REST relay loop: parse request -> evaluate -> allow/deny -> relay response -> repeat.
(
config: &L7EndpointConfig,
engine: &TunnelPolicyEngine,
client: &mut C,
upstream: &mut U,
ctx: &L7EvalContext,
)
| 736 | |
| 737 | /// REST relay loop: parse request -> evaluate -> allow/deny -> relay response -> repeat. |
| 738 | async fn relay_rest<C, U>( |
| 739 | config: &L7EndpointConfig, |
| 740 | engine: &TunnelPolicyEngine, |
| 741 | client: &mut C, |
| 742 | upstream: &mut U, |
| 743 | ctx: &L7EvalContext, |
| 744 | ) -> Result<()> |
| 745 | where |
| 746 | C: AsyncRead + AsyncWrite + Unpin + Send, |
| 747 | U: AsyncRead + AsyncWrite + Unpin + Send, |
| 748 | { |
| 749 | // Build a provider carrying the per-endpoint canonicalization options so |
| 750 | // request parsing honors the endpoint's `allow_encoded_slash` setting |
| 751 | // (e.g. APIs like GitLab that embed `%2F` in path segments). |
| 752 | let provider = |
| 753 | crate::l7::rest::RestProvider::with_options(crate::l7::path::CanonicalizeOptions { |
| 754 | allow_encoded_slash: config.allow_encoded_slash, |
| 755 | ..Default::default() |
| 756 | }); |
| 757 | loop { |
| 758 | if close_if_stale(engine.generation_guard(), ctx) { |
| 759 | return Ok(()); |
| 760 | } |
| 761 | |
| 762 | // Parse one HTTP request from client |
| 763 | let req = match provider.parse_request(client).await { |
| 764 | Ok(Some(req)) => req, |
| 765 | Ok(None) => return Ok(()), // Client closed connection |
| 766 | Err(e) => { |
| 767 | if is_benign_connection_error(&e) { |
| 768 | debug!( |
| 769 | host = %ctx.host, |
| 770 | port = ctx.port, |
| 771 | error = %e, |
| 772 | "L7 connection closed" |
| 773 | ); |
| 774 | } else { |
| 775 | let detail = |
| 776 | parse_rejection_detail(&e.to_string(), ParseRejectionMode::L7Endpoint); |
| 777 | emit_parse_rejection(ctx, &detail, "l7"); |
| 778 | } |
| 779 | return Ok(()); // Close connection on parse error |
| 780 | } |
| 781 | }; |
| 782 | |
| 783 | if deny_h2c_upgrade_if_requested(&req, config, ctx, client).await? { |
| 784 | return Ok(()); |
| 785 | } |
| 786 | |
| 787 | if close_if_stale(engine.generation_guard(), ctx) { |
| 788 | return Ok(()); |
| 789 | } |
| 790 | |
| 791 | // Rewrite credential placeholders in the request target BEFORE OPA |
| 792 | // evaluation. OPA sees the redacted path; the resolved path goes only |
| 793 | // to the upstream write. |
| 794 | let (eval_target, redacted_target) = if let Some(ref resolver) = ctx.secret_resolver { |
| 795 | match secrets::rewrite_target_for_eval(&req.target, resolver) { |
no test coverage detected