use super::super::coordination::ConversationOrderGate;
use super::super::large_message::As4FragmentJoiner;
use super::super::types::{As4ReceiveOutcome, As4ReceivePushProgress, As4ReceivePushRequest};
use super::ordered::{self, As4Ordered};
use super::{As4Verifier, EventBus, SessionContext};
use crate::core::Result;
use std::sync::Arc;
async fn reserve_turn_for_payload(
session: &SessionContext,
gate: &dyn ConversationOrderGate,
payload: &[u8],
) -> Result<(
String,
Box<dyn super::super::coordination::ConversationTurnHandle>,
)> {
let conversation_id = super::super::parser::extract_conversation_id_for_gate(payload)
.ok_or_else(|| ordered::ordered_missing_conversation_id_error(session))?;
let reservation = gate.reserve_ordered_turn(&conversation_id, session).await?;
Ok((conversation_id, reservation))
}
pub(super) async fn receive_push_ordered_with_verifier(
session: &SessionContext,
event_bus: &EventBus,
request: As4ReceivePushRequest,
dedup_backend: Arc<dyn crate::storage::DedupStorage>,
gate: &dyn ConversationOrderGate,
verifier: Arc<dyn As4Verifier + Send + Sync>,
) -> Result<As4Ordered<As4ReceiveOutcome>> {
let (conversation_id, reservation) =
reserve_turn_for_payload(session, gate, &request.payload).await?;
let outcome = super::async_completion::receive_push_with_dedup_async_with_shared_verifier(
session,
event_bus,
request,
dedup_backend,
verifier,
)
.await?;
match outcome {
As4ReceiveOutcome::FirstSeen(output) => {
let turn = reservation.wait_for_turn().await?;
ordered::confirm_and_record_turn(session, gate, &conversation_id, &output).await?;
Ok(As4Ordered::held(As4ReceiveOutcome::FirstSeen(output), turn))
}
duplicate @ As4ReceiveOutcome::Duplicate { .. } => Ok(As4Ordered::untimed(duplicate)),
}
}
pub(super) async fn receive_push_ordered_fragment_aware_with_verifier(
session: &SessionContext,
event_bus: &EventBus,
request: As4ReceivePushRequest,
dedup_backend: Arc<dyn crate::storage::DedupStorage>,
gate: &dyn ConversationOrderGate,
verifier: Arc<dyn As4Verifier + Send + Sync>,
fragment_joiner: Arc<std::sync::Mutex<As4FragmentJoiner>>,
) -> Result<As4Ordered<As4ReceivePushProgress>> {
let (conversation_id, reservation) =
reserve_turn_for_payload(session, gate, &request.payload).await?;
let progress =
super::async_bridge::receive_push_with_dedup_async_fragment_aware_with_shared_verifier(
session,
event_bus,
request,
dedup_backend,
verifier,
fragment_joiner,
)
.await?;
match progress {
As4ReceivePushProgress::Complete(output) => {
let turn = reservation.wait_for_turn().await?;
ordered::confirm_and_record_turn(session, gate, &conversation_id, &output).await?;
Ok(As4Ordered::held(
As4ReceivePushProgress::Complete(output),
turn,
))
}
other => Ok(As4Ordered::untimed(other)),
}
}