use std::sync::Arc;
use std::time::Duration;
use meerkat_contracts::wire::{WireMemberHistoryPageBody, WireProjectionProvenance};
use meerkat_core::comms::TrustedPeerDescriptor;
use super::MobSupervisorBridge;
use super::bridge_protocol::{
BridgeCommand, BridgeMemberHistoryPage, BridgeProtocolVersion, BridgeReadHistoryPayload,
};
use crate::MobError;
use crate::machines::mob_machine as mob_dsl;
pub(crate) const HISTORY_BRIDGE_TIMEOUT: Duration = Duration::from_secs(15);
#[derive(Debug, Clone)]
pub struct MemberHistoryPageDomain {
pub generation: u64,
pub page: WireMemberHistoryPageBody,
pub placement: Option<mob_dsl::HostId>,
pub provenance: WireProjectionProvenance,
}
async fn read_remote_page(
bridge: &Arc<MobSupervisorBridge>,
peer: &TrustedPeerDescriptor,
placement: mob_dsl::HostId,
expected_member: &super::bridge_protocol::BridgeMemberIncarnation,
from_index: Option<u64>,
limit: Option<u32>,
) -> Result<MemberHistoryPageDomain, MobError> {
let authority = bridge.authority().await;
let sup_spec = bridge.supervisor_spec_for_recipient(peer).await?;
let command = BridgeCommand::ReadMemberHistory(BridgeReadHistoryPayload {
supervisor: sup_spec.into(),
epoch: authority.epoch,
protocol_version: BridgeProtocolVersion::V4,
expected_member: expected_member.clone(),
from_index,
limit,
});
let _ = bridge.trust_recipient(peer).await?;
let value = bridge
.send_bridge_command(peer, &command, HISTORY_BRIDGE_TIMEOUT)
.await?;
let page: BridgeMemberHistoryPage =
super::bridge_protocol::decode_bridge_payload(&command, value, "read member history")?;
Ok(MemberHistoryPageDomain {
generation: page.generation,
page: page.page,
placement: Some(placement),
provenance: WireProjectionProvenance::HostClaimed,
})
}
pub(crate) async fn read_remote_member_history_page(
bridge: &Arc<MobSupervisorBridge>,
peer: &TrustedPeerDescriptor,
placement: mob_dsl::HostId,
expected_member: super::bridge_protocol::BridgeMemberIncarnation,
from_index: Option<u64>,
limit: Option<u32>,
) -> Result<MemberHistoryPageDomain, MobError> {
read_remote_page(bridge, peer, placement, &expected_member, from_index, limit).await
}
fn checked_page_end(page: &WireMemberHistoryPageBody) -> Result<u64, MobError> {
let served = u64::try_from(page.messages.len()).map_err(|_| {
MobError::Internal("member history page row count does not fit u64".to_string())
})?;
let end = page.from_index.checked_add(served).ok_or_else(|| {
MobError::Internal("member history page index overflowed u64".to_string())
})?;
if end > page.message_count {
return Err(MobError::Internal(format!(
"member history page ends at {end} beyond its message_count {}",
page.message_count
)));
}
match (page.complete, page.next_index) {
(true, None) => {}
(true, Some(next)) => {
return Err(MobError::Internal(format!(
"complete member history page unexpectedly carries next_index {next}"
)));
}
(false, Some(next)) if next == end => {}
(false, Some(next)) => {
return Err(MobError::Internal(format!(
"member history page next_index {next} does not equal served end {end}"
)));
}
(false, None) => {
return Err(MobError::Internal(
"incomplete member history page carries no next_index".to_string(),
));
}
}
Ok(end)
}
async fn finish_remote_history_snapshot(
bridge: &Arc<MobSupervisorBridge>,
peer: &TrustedPeerDescriptor,
placement: mob_dsl::HostId,
expected_member: &super::bridge_protocol::BridgeMemberIncarnation,
mut aggregate: MemberHistoryPageDomain,
) -> Result<MemberHistoryPageDomain, MobError> {
let generation = aggregate.generation;
let snapshot_end = aggregate.page.message_count;
loop {
let end = checked_page_end(&aggregate.page)?;
if end == snapshot_end {
aggregate.page.message_count = snapshot_end;
aggregate.page.next_index = None;
aggregate.page.complete = true;
return Ok(aggregate);
}
if aggregate.page.messages.is_empty() {
return Err(MobError::Internal(format!(
"member history page made no progress at index {} before snapshot end {snapshot_end}",
aggregate.page.from_index
)));
}
let next_index = aggregate.page.next_index.ok_or_else(|| {
MobError::Internal(format!(
"member history snapshot ended at {end} before its declared message_count {snapshot_end}"
))
})?;
let remaining = snapshot_end - end;
let limit = u32::try_from(remaining).unwrap_or(u32::MAX).max(1);
let page = read_remote_page(
bridge,
peer,
placement.clone(),
expected_member,
Some(next_index),
Some(limit),
)
.await?;
if page.generation != generation {
return Err(MobError::Internal(format!(
"member transcript generation changed mid-read ({generation} -> {}); retry the read",
page.generation
)));
}
if page.page.from_index != next_index {
return Err(MobError::Internal(format!(
"member history continuation requested index {next_index} but host served {}",
page.page.from_index
)));
}
if page.page.message_count < snapshot_end {
return Err(MobError::Internal(format!(
"member transcript shrank mid-read ({snapshot_end} -> {}); retry the read",
page.page.message_count
)));
}
let page_end = checked_page_end(&page.page)?;
if page_end > snapshot_end {
return Err(MobError::Internal(format!(
"member history continuation crossed pinned snapshot end {snapshot_end} (served through {page_end})"
)));
}
if page.page.messages.is_empty() && page_end < snapshot_end {
return Err(MobError::Internal(format!(
"member history continuation made no progress at index {next_index} before snapshot end {snapshot_end}"
)));
}
aggregate.page.messages.extend(page.page.messages);
aggregate.page.next_index = page.page.next_index;
aggregate.page.complete = page.page.complete;
}
}
pub(crate) async fn read_remote_member_full_history(
bridge: &Arc<MobSupervisorBridge>,
peer: &TrustedPeerDescriptor,
placement: mob_dsl::HostId,
expected_member: super::bridge_protocol::BridgeMemberIncarnation,
) -> Result<MemberHistoryPageDomain, MobError> {
let first = read_remote_page(
bridge,
peer,
placement.clone(),
&expected_member,
Some(0),
None,
)
.await?;
if first.page.from_index != 0 {
return Err(MobError::Internal(format!(
"full member history read requested index 0 but host served {}",
first.page.from_index
)));
}
finish_remote_history_snapshot(bridge, peer, placement, &expected_member, first).await
}
pub(crate) async fn read_remote_member_history_tail(
bridge: &Arc<MobSupervisorBridge>,
peer: &TrustedPeerDescriptor,
placement: mob_dsl::HostId,
expected_member: super::bridge_protocol::BridgeMemberIncarnation,
count: u32,
) -> Result<MemberHistoryPageDomain, MobError> {
let first = if count == 0 {
read_remote_page(
bridge,
peer,
placement.clone(),
&expected_member,
Some(u64::MAX),
Some(1),
)
.await?
} else {
read_remote_page(
bridge,
peer,
placement.clone(),
&expected_member,
None,
Some(count),
)
.await?
};
if count == 0 {
if first.page.from_index != first.page.message_count || !first.page.messages.is_empty() {
return Err(MobError::Internal(
"zero-count member history tail did not resolve to an empty end page".to_string(),
));
}
return Ok(first);
}
finish_remote_history_snapshot(bridge, peer, placement, &expected_member, first).await
}