use super::large_message::As4FragmentJoiner;
mod async_bridge;
mod async_completion;
#[doc(hidden)]
mod entrypoints;
mod metadata;
mod ordered;
pub use ordered::As4Ordered;
mod ordered_async;
mod payload;
mod processing;
mod receipt;
mod runtime_guards;
mod sync_core;
mod verifier;
mod verifier_wrappers;
use super::types::{As4PushPolicy, As4ReceivePushProgress};
use crate::core::{PayloadInput, Result, SessionContext};
use crate::observability::EventBus;
use crate::storage::DedupStorage;
pub use entrypoints::{
As4ReceivePushAsyncFragmentAwareRequest, As4ReceivePushOrderedFragmentAwareRequest,
As4ReceivePushOrderedRequest, As4ReceivePushSyncFragmentAwareRequest,
As4ReceivePushSyncRequest, receive_push_ordered, receive_push_ordered_fragment_aware,
receive_push_with_dedup_async, receive_push_with_dedup_async_fragment_aware,
receive_push_with_dedup_sync, receive_push_with_dedup_sync_fragment_aware,
};
#[cfg(feature = "testing")]
pub use entrypoints::{
receive_push_with_dedup_async_with_custom_verifier,
receive_push_with_dedup_sync_with_custom_verifier,
};
#[cfg(feature = "testing")]
pub use verifier::InsecureBypassAs4Verifier;
#[cfg(all(test, not(feature = "testing")))]
pub(crate) use verifier::private;
#[cfg(feature = "testing")]
pub use verifier::private;
pub use verifier::{As4Verifier, As4WsSecVerifier};
use verifier_wrappers::{
receive_push_with_dedup_async_with_verifier, receive_push_with_dedup_sync_with_verifier,
};
pub(super) struct As4PushReceiveCtx<'a> {
pub(super) session: &'a SessionContext,
pub(super) event_bus: &'a EventBus,
pub(super) policy: &'a As4PushPolicy,
pub(super) dedup_backend: &'a dyn DedupStorage,
pub(super) verifier: &'a (dyn As4Verifier + Send + Sync),
}
fn receive_push_with_dedup_inner(
ctx: &As4PushReceiveCtx<'_>,
payload_input: PayloadInput<'_>,
receipt_payload: Option<&[u8]>,
http_content_type: &str,
fragment_joiner: Option<&mut As4FragmentJoiner>,
authenticated_sender_scope: Option<&str>,
) -> Result<As4ReceivePushProgress> {
runtime_guards::enforce_receive_push_runtime_guards(
ctx.session,
ctx.event_bus,
ctx.policy,
ctx.dedup_backend,
)?;
let payload_bytes = payload_input.as_slice();
if let Some(fragment_progress) = processing::maybe_handle_fragment_message(
ctx,
payload_bytes,
receipt_payload,
http_content_type,
fragment_joiner,
authenticated_sender_scope,
)? {
return Ok(fragment_progress);
}
processing::process_non_fragment_push(ctx, payload_bytes, receipt_payload, http_content_type)
}
#[cfg(test)]
#[cfg_attr(not(feature = "interop-relaxed"), allow(unused_imports))]
mod tests;