use std::{
any::Any,
collections::{HashMap, VecDeque},
sync::Arc,
};
use bytes::Bytes;
use crate::{
FilterPipeline, StreamTermination,
context::PendingStreamChunks,
extensions::{RequestExtensions, SelectedClusterApplication},
results::{FilterResultSet, RetainedFilterResults},
};
pub(crate) struct SubrequestCompletion {
pub(crate) extensions: RequestExtensions,
#[cfg_attr(
not(feature = "iterative-request-router"),
expect(dead_code, reason = "read only by the iterative_request_router continuation path")
)]
pub(crate) filter_results: HashMap<&'static str, FilterResultSet>,
pub(crate) pending_chunks: VecDeque<Bytes>,
pub(crate) termination: Option<StreamTermination>,
}
pub(crate) struct FilteredSubrequestContinuation {
pub(super) pipeline: Arc<FilterPipeline>,
pub(super) request_snapshot: crate::Request,
pub(super) response_snapshot: crate::Response,
pub(super) extensions: RequestExtensions,
pub(super) filter_state: HashMap<usize, Box<dyn Any + Send + Sync>>,
pub(super) filter_results: HashMap<&'static str, FilterResultSet>,
pub(super) filter_metadata: HashMap<String, String>,
pub(super) structured_metadata: HashMap<String, serde_json::Value>,
pub(super) executed_filter_indices: Vec<bool>,
pub(super) body_done_indices: Vec<bool>,
pub(super) response_body_bytes: u64,
pub(super) response_body_mode: crate::body::BodyMode,
pub(super) completed: bool,
pub(super) client_addr: Option<std::net::IpAddr>,
pub(super) downstream_tls: bool,
pub(super) request_start: std::time::Instant,
pub(super) step_deadline: std::time::Instant,
pub(super) peer_identity: Option<Arc<praxis_tls::TlsPeerIdentity>>,
}
impl FilteredSubrequestContinuation {
#[expect(
clippy::too_many_arguments,
reason = "capture owns the complete response continuation boundary"
)]
pub(super) fn capture(
pipeline: Arc<FilterPipeline>,
request_snapshot: crate::Request,
response_snapshot: crate::Response,
ctx: &mut crate::filter::HttpFilterContext<'_>,
completed: bool,
step_deadline: std::time::Instant,
) -> Self {
Self {
pipeline,
request_snapshot,
response_snapshot,
extensions: std::mem::take(&mut ctx.extensions),
filter_state: std::mem::take(&mut ctx.filter_state),
filter_results: std::mem::take(&mut ctx.filter_results),
filter_metadata: std::mem::take(&mut ctx.filter_metadata),
structured_metadata: std::mem::take(&mut ctx.structured_metadata),
executed_filter_indices: std::mem::take(&mut ctx.executed_filter_indices),
body_done_indices: std::mem::take(&mut ctx.body_done_indices),
response_body_bytes: ctx.response_body_bytes,
response_body_mode: ctx.response_body_mode,
completed,
client_addr: ctx.client_addr,
downstream_tls: ctx.downstream_tls,
request_start: ctx.request_start,
step_deadline,
peer_identity: ctx.peer_identity.clone(),
}
}
pub(crate) fn into_parent_extensions(mut self) -> RequestExtensions {
self.extensions.remove::<SelectedClusterApplication>();
self.extensions.remove::<PendingStreamChunks>();
self.extensions.remove::<RetainedFilterResults>();
self.extensions.remove::<StreamTermination>();
self.extensions
}
pub(crate) fn into_completion(mut self) -> SubrequestCompletion {
self.extensions.remove::<SelectedClusterApplication>();
let pending_chunks = self
.extensions
.remove::<PendingStreamChunks>()
.map_or_else(VecDeque::new, PendingStreamChunks::into_chunks);
let termination = self.extensions.remove::<StreamTermination>();
let mut filter_results = self.extensions.remove::<RetainedFilterResults>().unwrap_or_default().0;
filter_results.extend(self.filter_results);
SubrequestCompletion {
extensions: self.extensions,
filter_results,
pending_chunks,
termination,
}
}
#[cfg_attr(
not(feature = "iterative-request-router"),
expect(dead_code, reason = "used only by the iterative_request_router filter")
)]
pub(crate) fn extensions(&self) -> &RequestExtensions {
&self.extensions
}
}