use super::super::coordination::ConversationOrderGate;
use super::super::large_message::As4FragmentJoiner;
use super::super::types::{As4ReceiveOutcome, As4ReceivePushProgress, As4ReceivePushRequest};
use super::ordered::As4Ordered;
use super::ordered_async;
use super::{As4WsSecVerifier, EventBus, SessionContext, sync_core};
use crate::core::Result;
use crate::storage::DedupStorage;
use std::sync::Arc;
#[derive(Debug)]
pub struct As4ReceivePushOrderedRequest<'a> {
pub request: As4ReceivePushRequest,
pub dedup_backend: Arc<dyn crate::storage::DedupStorage>,
pub gate: &'a dyn ConversationOrderGate,
}
#[derive(Debug)]
pub struct As4ReceivePushOrderedFragmentAwareRequest<'a> {
pub request: As4ReceivePushRequest,
pub dedup_backend: Arc<dyn crate::storage::DedupStorage>,
pub gate: &'a dyn ConversationOrderGate,
pub fragment_joiner: Arc<std::sync::Mutex<As4FragmentJoiner>>,
}
#[derive(Debug)]
pub struct As4ReceivePushAsyncFragmentAwareRequest {
pub request: As4ReceivePushRequest,
pub dedup_backend: Arc<dyn DedupStorage>,
pub fragment_joiner: Arc<std::sync::Mutex<As4FragmentJoiner>>,
}
#[derive(Debug)]
pub struct As4ReceivePushSyncRequest<'a> {
pub request: As4ReceivePushRequest,
pub dedup_backend: &'a dyn DedupStorage,
}
#[derive(Debug)]
pub struct As4ReceivePushSyncFragmentAwareRequest<'a> {
pub request: As4ReceivePushRequest,
pub dedup_backend: &'a dyn DedupStorage,
pub fragment_joiner: &'a mut As4FragmentJoiner,
}
pub fn receive_push_with_dedup_sync(
session: &SessionContext,
event_bus: &EventBus,
request: As4ReceivePushSyncRequest<'_>,
) -> Result<As4ReceiveOutcome> {
let As4ReceivePushSyncRequest {
request,
dedup_backend,
} = request;
super::receive_push_with_dedup_sync_with_verifier(
session,
event_bus,
request,
dedup_backend,
&As4WsSecVerifier,
)
}
pub fn receive_push_with_dedup_sync_fragment_aware(
session: &SessionContext,
event_bus: &EventBus,
request: As4ReceivePushSyncFragmentAwareRequest<'_>,
) -> Result<As4ReceivePushProgress> {
let As4ReceivePushSyncFragmentAwareRequest {
request,
dedup_backend,
fragment_joiner,
} = request;
sync_core::receive_push_with_dedup_sync_fragment_aware_with_verifier(
session,
event_bus,
request,
dedup_backend,
&As4WsSecVerifier,
Some(fragment_joiner),
)
}
pub async fn receive_push_with_dedup_async(
session: &SessionContext,
event_bus: &EventBus,
request: As4ReceivePushRequest,
dedup_backend: Arc<dyn DedupStorage>,
) -> Result<As4ReceiveOutcome> {
super::receive_push_with_dedup_async_with_verifier(
session,
event_bus,
request,
dedup_backend,
As4WsSecVerifier,
)
.await
}
pub async fn receive_push_with_dedup_async_fragment_aware(
session: &SessionContext,
event_bus: &EventBus,
request: As4ReceivePushAsyncFragmentAwareRequest,
) -> Result<As4ReceivePushProgress> {
let As4ReceivePushAsyncFragmentAwareRequest {
request,
dedup_backend,
fragment_joiner,
} = request;
super::async_bridge::receive_push_with_dedup_async_fragment_aware_with_shared_verifier(
session,
event_bus,
request,
dedup_backend,
Arc::new(As4WsSecVerifier),
fragment_joiner,
)
.await
}
#[cfg_attr(feature = "trace", tracing::instrument(skip_all, fields(partner_id = %session.partner_id())))]
pub async fn receive_push_ordered(
session: &SessionContext,
event_bus: &EventBus,
request: As4ReceivePushOrderedRequest<'_>,
) -> crate::core::Result<As4Ordered<As4ReceiveOutcome>> {
let As4ReceivePushOrderedRequest {
request,
dedup_backend,
gate,
} = request;
ordered_async::receive_push_ordered_with_verifier(
session,
event_bus,
request,
dedup_backend,
gate,
Arc::new(As4WsSecVerifier),
)
.await
}
#[cfg_attr(feature = "trace", tracing::instrument(skip_all, fields(partner_id = %session.partner_id())))]
pub async fn receive_push_ordered_fragment_aware(
session: &SessionContext,
event_bus: &EventBus,
request: As4ReceivePushOrderedFragmentAwareRequest<'_>,
) -> crate::core::Result<As4Ordered<As4ReceivePushProgress>> {
let As4ReceivePushOrderedFragmentAwareRequest {
request,
dedup_backend,
gate,
fragment_joiner,
} = request;
ordered_async::receive_push_ordered_fragment_aware_with_verifier(
session,
event_bus,
request,
dedup_backend,
gate,
Arc::new(As4WsSecVerifier),
fragment_joiner,
)
.await
}
#[cfg(feature = "testing")]
pub fn receive_push_with_dedup_sync_with_custom_verifier(
session: &SessionContext,
event_bus: &EventBus,
request: As4ReceivePushSyncRequest<'_>,
verifier: &(dyn super::As4Verifier + Send + Sync),
) -> Result<As4ReceiveOutcome> {
super::receive_push_with_dedup_sync_with_verifier(
session,
event_bus,
request.request,
request.dedup_backend,
verifier,
)
}
#[cfg(feature = "testing")]
pub async fn receive_push_with_dedup_async_with_custom_verifier<V>(
session: &SessionContext,
event_bus: &EventBus,
request: As4ReceivePushRequest,
dedup_backend: Arc<dyn DedupStorage>,
verifier: V,
) -> Result<As4ReceiveOutcome>
where
V: super::As4Verifier + Send + Sync + 'static,
{
super::receive_push_with_dedup_async_with_verifier(
session,
event_bus,
request,
dedup_backend,
verifier,
)
.await
}