use std::collections::{HashMap, HashSet};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex, Weak};
use async_trait::async_trait;
use super::contracts::{
AgentCustomizer, RosterProvider, SessionSnapshotMatchCandidate, TopologyProvider,
};
use super::types::{
AgentAddressability, AgentBuildContext, AgentBuildDraft, AgentIdentity, ContinuityStoreError,
CustomizerError, DurableAgentSpec, ManagedPeerEdge, RosterContext, RosterError,
TopologyContext, TopologyError,
};
use crate::mob_handle_runtime::{SessionCreatedContext, SessionHook};
use crate::types::AgentDiscoverySpec;
use crate::unified_runtime::edge_types::{Discovery, EdgeDiscovery};
#[derive(Default)]
pub struct MutableRosterProvider {
roster: std::sync::RwLock<Vec<DurableAgentSpec>>,
}
impl MutableRosterProvider {
pub fn new(initial: Vec<DurableAgentSpec>) -> Self {
Self {
roster: std::sync::RwLock::new(initial),
}
}
pub fn upsert(&self, spec: DurableAgentSpec) {
let mut roster = self
.roster
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
match roster
.iter_mut()
.find(|entry| entry.identity == spec.identity)
{
Some(entry) => *entry = spec,
None => roster.push(spec),
}
}
pub fn remove(&self, identity: &AgentIdentity) -> bool {
let mut roster = self
.roster
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let before = roster.len();
roster.retain(|entry| &entry.identity != identity);
roster.len() != before
}
pub fn snapshot(&self) -> Vec<DurableAgentSpec> {
self.roster
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
}
#[async_trait::async_trait]
impl RosterProvider for MutableRosterProvider {
async fn roster(&self, _context: &RosterContext) -> Result<Vec<DurableAgentSpec>, RosterError> {
Ok(self.snapshot())
}
}
pub struct DiscoveryRosterAdapter {
inner: Box<dyn Discovery>,
}
impl DiscoveryRosterAdapter {
pub fn new(discovery: impl Discovery + 'static) -> Self {
Self {
inner: Box::new(discovery),
}
}
}
pub fn agent_discovery_to_durable(
spec: &AgentDiscoverySpec,
) -> Result<DurableAgentSpec, RosterError> {
let identity = AgentIdentity::parse(&spec.meerkat_id)
.map_err(|e| RosterError::Io(format!("invalid meerkat_id: {e}")))?;
Ok(DurableAgentSpec {
identity,
profile: meerkat_mob::ProfileName::from(spec.profile.as_str()),
addressability: AgentAddressability::Addressable,
display_name: None,
labels: spec.labels.clone().unwrap_or_default(),
context: spec.context.clone(),
additional_instructions: spec.additional_instructions.clone(),
initial_message: None,
runtime_mode_override: None,
backend: None,
binding: None,
})
}
#[async_trait]
impl RosterProvider for DiscoveryRosterAdapter {
async fn roster(&self, _context: &RosterContext) -> Result<Vec<DurableAgentSpec>, RosterError> {
let specs = self.inner.discover(serde_json::Value::Null).await;
specs.iter().map(agent_discovery_to_durable).collect()
}
}
pub struct EdgeDiscoveryTopologyAdapter {
inner: Box<dyn EdgeDiscovery>,
}
impl EdgeDiscoveryTopologyAdapter {
pub fn new(edge_discovery: impl EdgeDiscovery + 'static) -> Self {
Self {
inner: Box::new(edge_discovery),
}
}
}
#[async_trait]
impl TopologyProvider for EdgeDiscoveryTopologyAdapter {
async fn compute_edges(
&self,
_target_identities: &[AgentIdentity],
context: &TopologyContext,
) -> Result<Vec<ManagedPeerEdge>, TopologyError> {
let member_views: Vec<crate::unified_runtime::edge_types::EdgeMemberView> = context
.roster
.iter()
.map(|spec| crate::unified_runtime::edge_types::EdgeMemberView {
agent_identity: spec.identity.as_str().to_string(),
role: spec.profile.as_str().to_string(),
wired_to: std::collections::BTreeSet::new(),
labels: spec.labels.clone(),
})
.collect();
let desired_edges = self.inner.discover_edges(member_views).await;
let mut edges = Vec::with_capacity(desired_edges.len());
for edge in &desired_edges {
let (a_str, b_str) = edge.endpoints();
let a = AgentIdentity::parse(a_str)
.map_err(|e| TopologyError::InvalidEdge(format!("endpoint {a_str:?}: {e}")))?;
let b = AgentIdentity::parse(b_str)
.map_err(|e| TopologyError::InvalidEdge(format!("endpoint {b_str:?}: {e}")))?;
let managed = ManagedPeerEdge::new(a, b)
.map_err(|e| TopologyError::InvalidEdge(format!("{e}")))?;
edges.push(managed);
}
Ok(edges)
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct SessionRuntimeState {
pub identity: AgentIdentity,
pub generation: super::types::ContinuityGeneration,
pub fencing_token: super::types::FencingToken,
pub checkpoint_version: super::types::CheckpointVersion,
}
struct ParkedDeltas {
store: meerkat_store::MemoryStore,
sessions: Mutex<HashMap<String, ParkedFootprint>>,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
struct ParkedFootprint {
rows: u64,
}
impl ParkedDeltas {
fn new() -> Self {
Self {
store: meerkat_store::MemoryStore::new(),
sessions: Mutex::new(HashMap::new()),
}
}
fn is_parked(&self, session_id: &meerkat_core::types::SessionId) -> bool {
self.sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.contains_key(&session_id.to_string())
}
fn reads(&self) -> &meerkat_store::MemoryStore {
&self.store
}
fn mark_parked(&self, session_id: &meerkat_core::types::SessionId, rows: u64) {
let mut sessions = self
.sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let footprint = sessions.entry(session_id.to_string()).or_default();
footprint.rows = footprint.rows.saturating_add(rows);
}
async fn park_append(
&self,
session_id: &meerkat_core::types::SessionId,
strand: &meerkat_core::session_store::TranscriptStrandId,
base_seq: u64,
messages: &[meerkat_core::Message],
) -> Result<(), meerkat_store::SessionStoreError> {
meerkat_core::session_store::IncrementalSessionStore::append_messages(
&self.store,
session_id,
strand,
base_seq,
messages,
)
.await?;
self.mark_parked(session_id, messages.len() as u64);
Ok(())
}
async fn park_rewrite(
&self,
session_id: &meerkat_core::types::SessionId,
record: &meerkat_core::TranscriptRewriteRecord,
expected: meerkat_core::session_store::SessionHeadCas,
) -> Result<meerkat_core::session_store::SessionHead, meerkat_store::SessionStoreError> {
let head = meerkat_core::session_store::IncrementalSessionStore::commit_rewrite(
&self.store,
session_id,
record,
expected,
)
.await?;
self.mark_parked(session_id, record.revision_body.messages.len() as u64);
Ok(head)
}
async fn park_head(
&self,
head: &meerkat_core::session_store::SessionHead,
expected: meerkat_core::session_store::SessionHeadCas,
) -> Result<(), meerkat_store::SessionStoreError> {
meerkat_core::session_store::IncrementalSessionStore::save_head(
&self.store,
head,
expected,
)
.await?;
self.mark_parked(&head.id, 0);
Ok(())
}
fn footprint(&self, session_id: &meerkat_core::types::SessionId) -> Option<ParkedFootprint> {
self.sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(&session_id.to_string())
.copied()
}
fn unmark(&self, session_id: &meerkat_core::types::SessionId) -> bool {
self.sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&session_id.to_string())
.is_some()
}
async fn purge(&self, session_id: &meerkat_core::types::SessionId) {
self.unmark(session_id);
let _ = meerkat::SessionStore::delete(&self.store, session_id).await;
}
}
enum ParkedFlush {
Adopted(super::types::CheckpointVersion),
Empty,
}
fn session_is_legacy_unverified(session: &meerkat_core::Session) -> bool {
matches!(
meerkat_core::session_checkpoint_metadata_state(session.id(), session.metadata()),
Ok(meerkat_core::SessionCheckpointMetadataState::LegacyUnverified { .. })
)
}
pub struct ContinuitySessionStoreAdapter {
store: Arc<dyn super::contracts::ContinuityStore>,
incremental: Option<Arc<dyn super::contracts::ContinuityIncrementalSessions>>,
parked_deltas: ParkedDeltas,
versions: Mutex<HashMap<String, AtomicU64>>,
session_registry: Mutex<HashMap<String, SessionRuntimeState>>,
pending_unregistered: Mutex<HashMap<String, Vec<u8>>>,
unregistered_sessions: Mutex<HashSet<String>>,
suspended_sessions: Mutex<HashSet<String>>,
superseded_sessions: Mutex<HashSet<String>>,
session_locks: Mutex<HashMap<String, Weak<tokio::sync::Mutex<()>>>>,
lazy_checkpoint_adoption: bool,
}
impl ContinuitySessionStoreAdapter {
pub fn new(store: Arc<dyn super::contracts::ContinuityStore>) -> Self {
let incremental = store.as_incremental_sessions();
Self {
store,
incremental,
parked_deltas: ParkedDeltas::new(),
versions: Mutex::new(HashMap::new()),
session_registry: Mutex::new(HashMap::new()),
pending_unregistered: Mutex::new(HashMap::new()),
unregistered_sessions: Mutex::new(HashSet::new()),
suspended_sessions: Mutex::new(HashSet::new()),
superseded_sessions: Mutex::new(HashSet::new()),
session_locks: Mutex::new(HashMap::new()),
lazy_checkpoint_adoption: false,
}
}
#[must_use]
pub fn with_lazy_checkpoint_adoption(mut self, enabled: bool) -> Self {
self.lazy_checkpoint_adoption = enabled;
self
}
fn session_lock(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Arc<tokio::sync::Mutex<()>> {
const PRUNE_THRESHOLD: usize = 1_024;
let key = session_id.to_string();
let mut locks = self
.session_locks
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(lock) = locks.get(&key).and_then(Weak::upgrade) {
return lock;
}
if locks.len() >= PRUNE_THRESHOLD {
locks.retain(|_, lock| lock.strong_count() > 0);
}
let lock = Arc::new(tokio::sync::Mutex::new(()));
locks.insert(key, Arc::downgrade(&lock));
lock
}
async fn lock_session(
&self,
session_id: &meerkat_core::types::SessionId,
) -> tokio::sync::OwnedMutexGuard<()> {
self.session_lock(session_id).lock_owned().await
}
#[allow(dead_code)]
pub async fn register_session(
&self,
session_id: &meerkat_core::types::SessionId,
state: SessionRuntimeState,
) -> Result<super::types::CheckpointVersion, meerkat_store::SessionStoreError> {
let _guard = self.lock_session(session_id).await;
let session_key = session_id.to_string();
let checkpoint_version = state.checkpoint_version.get();
if let Some(existing) = self.lookup_session(&session_key) {
if existing.identity != state.identity || existing.generation != state.generation {
return Err(meerkat_store::SessionStoreError::Internal(format!(
"session ownership conflict for {session_id}: registered owner {}/generation {} cannot be replaced by {}/generation {}",
existing.identity, existing.generation, state.identity, state.generation
)));
}
if state.fencing_token < existing.fencing_token {
return Err(meerkat_store::SessionStoreError::Internal(format!(
"session ownership conflict for {session_id}: fencing token {} cannot regress to {}",
existing.fencing_token, state.fencing_token
)));
}
}
let was_unregistered = self
.unregistered_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&session_key);
let was_suspended = self
.suspended_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&session_key);
let was_superseded = self
.superseded_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&session_key);
let previous_registry = {
let mut registry = self
.session_registry
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
registry.insert(session_key.clone(), state.clone())
};
{
let mut versions = self
.versions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let counter = versions
.entry(session_key.clone())
.or_insert_with(|| AtomicU64::new(checkpoint_version));
counter.fetch_max(checkpoint_version, Ordering::Relaxed);
}
let pending = {
let pending = self
.pending_unregistered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
pending.get(&session_id.to_string()).cloned()
};
let mut effective_checkpoint_version = self.current_version(session_id);
let restore_markers = |adapter: &Self| {
adapter.restore_registration_state(session_id, previous_registry.clone());
if was_unregistered {
adapter
.unregistered_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(session_key.clone());
}
if was_suspended {
adapter
.suspended_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(session_key.clone());
}
if was_superseded {
adapter
.superseded_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(session_key.clone());
}
};
if let Some(data) = pending {
let flush_result = self
.save_registered_snapshot(session_id, data, state.clone())
.await;
match flush_result {
Ok(version) => {
self.pending_unregistered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&session_id.to_string());
effective_checkpoint_version = version;
}
Err(err) => {
restore_markers(self);
return Err(err);
}
}
}
if self.parked_deltas.is_parked(session_id) {
match self.flush_parked_deltas(session_id, &state).await {
Ok(ParkedFlush::Adopted(version)) => {
effective_checkpoint_version = version;
self.parked_deltas.purge(session_id).await;
}
Ok(ParkedFlush::Empty) => {
self.parked_deltas.purge(session_id).await;
}
Err(err) => {
restore_markers(self);
return Err(err);
}
}
}
Ok(effective_checkpoint_version)
}
async fn flush_parked_deltas(
&self,
session_id: &meerkat_core::types::SessionId,
state: &SessionRuntimeState,
) -> Result<ParkedFlush, meerkat_store::SessionStoreError> {
let Some(incremental) = self.incremental.as_ref() else {
return Err(meerkat_store::SessionStoreError::Internal(format!(
"session {session_id} parked incremental writes but the continuity substrate \
advertises no delta channel"
)));
};
let parked = self.parked_deltas.reads();
let Some(head) =
meerkat_core::session_store::IncrementalSessionStore::load_head(parked, session_id)
.await?
else {
let parked_rows = self
.parked_deltas
.footprint(session_id)
.unwrap_or_default()
.rows;
if parked_rows == 0 {
let parked_rewrites =
meerkat_core::session_store::IncrementalSessionStore::load_rewrites(
parked, session_id,
)
.await?;
if !parked_rewrites.is_empty() {
return Err(meerkat_store::SessionStoreError::Internal(format!(
"session {session_id} reports an empty parked footprint but the parked \
store holds {} rewrite record(s) and no head; refusing the registration \
instead of purging state the footprint cannot account for",
parked_rewrites.len()
)));
}
return Ok(ParkedFlush::Empty);
}
tracing::warn!(
session_id = %session_id,
parked_rows,
"registration refused: parked delta rows have no adopting head yet; \
retaining them for the retry"
);
return Err(meerkat_store::SessionStoreError::Internal(format!(
"session {session_id} parked {parked_rows} delta message row(s) that no parked \
head adopts yet; refusing the registration instead of dropping them — the \
parked state is retained, so retry once the adopting head write lands"
)));
};
if head.rewrite_count > 0
|| !meerkat_core::session_store::IncrementalSessionStore::load_rewrites(
parked, session_id,
)
.await?
.is_empty()
{
return Err(meerkat_store::SessionStoreError::Internal(format!(
"session {session_id} parked a transcript rewrite before its owning identity was \
registered; refusing to flush an unauditable rewrite chain"
)));
}
let messages = meerkat_core::session_store::IncrementalSessionStore::load_messages(
parked,
session_id,
&head.strand,
0..head.message_count,
)
.await?;
if incremental.load_canonical_head(session_id).await?.is_some() {
let session = head.clone().into_session(messages)?;
let data = serde_json::to_vec(&session)
.map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))?;
let version = self
.save_registered_snapshot(session_id, data, state.clone())
.await?;
return Ok(ParkedFlush::Adopted(version));
}
let append_cursor = self.write_cursor(session_id, state);
incremental
.append_messages(&append_cursor, session_id, &head.strand, 0, &messages)
.await?;
let head_cursor = self.write_cursor(session_id, state);
let committed = head_cursor.checkpoint_version;
incremental
.save_head(
&head_cursor,
&head,
meerkat_core::session_store::SessionHeadCas::Create,
)
.await?;
Ok(ParkedFlush::Adopted(committed))
}
async fn head_canonical_previous(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<
Option<(
meerkat_core::Session,
Vec<meerkat_core::TranscriptRewriteCommit>,
)>,
meerkat_store::SessionStoreError,
> {
let Some(incremental) = self.incremental.as_ref() else {
return Ok(None);
};
if self.parked_deltas.is_parked(id) {
return Ok(None);
}
incremental.load_canonical_previous(id).await
}
async fn save_head_canonical_transcript_rewrite(
&self,
session: &meerkat_core::Session,
commit: &meerkat_core::TranscriptRewriteCommit,
) -> Result<bool, meerkat_store::SessionStoreError> {
let Some(incremental) = self.incremental.as_ref() else {
return Ok(false);
};
let id = session.id();
if self.parked_deltas.is_parked(id) {
return Ok(false);
}
let Some(state) = self.lookup_session(&id.to_string()) else {
return Ok(false);
};
let Some(stored) = incremental.load_canonical_head(id).await? else {
return Ok(false);
};
let incoming_revision = meerkat_core::transcript_messages_digest(session.messages())?;
if incoming_revision != commit.revision {
return Err(meerkat_store::SessionStoreError::InvalidTranscriptRewrite {
id: id.clone(),
reason: format!(
"incoming current transcript digest {incoming_revision} does not match \
commit revision {}",
commit.revision
),
});
}
let record = rewrite_record_from_session_bodies(session, commit)?;
let token = meerkat_core::session_store::session_head_cas_token(&stored)?;
let next = incremental
.commit_rewrite(
&self.write_cursor(id, &state),
id,
&record,
meerkat_core::session_store::SessionHeadCas::IfToken(token.clone()),
)
.await?;
let adopted = meerkat_core::session_store::SessionHead::from_session(
session,
next.strand.clone(),
next.rewrite_count,
)?;
incremental
.save_head(
&self.write_cursor(id, &state),
&adopted,
meerkat_core::session_store::SessionHeadCas::IfToken(token),
)
.await?;
Ok(true)
}
fn write_cursor(
&self,
session_id: &meerkat_core::types::SessionId,
state: &SessionRuntimeState,
) -> super::contracts::ContinuityWriteCursor {
super::contracts::ContinuityWriteCursor {
identity: state.identity.clone(),
generation: state.generation,
checkpoint_version: super::types::CheckpointVersion::new(
self.next_version(&session_id.to_string()),
),
fencing_token: state.fencing_token,
}
}
async fn forget_session(&self, session_id: &meerkat_core::types::SessionId) {
let key = session_id.to_string();
self.parked_deltas.purge(session_id).await;
self.session_registry
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&key);
self.pending_unregistered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&key);
self.versions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&key);
self.suspended_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&key);
}
pub(crate) async fn suspend_session(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<(), ContinuityStoreError> {
let _guard = self.lock_session(session_id).await;
self.suspended_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(session_id.to_string());
Ok(())
}
pub(crate) async fn unregister_session(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<(), ContinuityStoreError> {
let _guard = self.lock_session(session_id).await;
self.forget_session(session_id).await;
self.superseded_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(&session_id.to_string());
self.unregistered_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(session_id.to_string());
Ok(())
}
pub(crate) async fn abandon_superseded_session(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<(), meerkat_store::SessionStoreError> {
let _guard = self.lock_session(session_id).await;
let session_key = session_id.to_string();
self.suspended_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(session_key.clone());
if let Some(session) = self.load_persisted_session(session_id).await? {
let current_revision =
meerkat_core::session_store::session_projection_cas_token(&session)?;
let deleted = self
.store
.delete_session_snapshot_if_current_revision(session_id, ¤t_revision)
.await
.map_err(|error| {
meerkat_store::SessionStoreError::Internal(format!(
"continuity abandon superseded session: {error}"
))
})?;
if !deleted {
return Err(meerkat_store::SessionStoreError::Internal(format!(
"continuity abandon did not delete superseded session snapshot {session_id}"
)));
}
}
self.forget_session(session_id).await;
self.superseded_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(session_key);
Ok(())
}
fn session_was_superseded(&self, session_id: &meerkat_core::types::SessionId) -> bool {
self.superseded_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.contains(&session_id.to_string())
}
fn session_was_unregistered(&self, session_id: &meerkat_core::types::SessionId) -> bool {
self.unregistered_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.contains(&session_id.to_string())
}
fn session_was_suspended(&self, session_id: &meerkat_core::types::SessionId) -> bool {
self.suspended_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.contains(&session_id.to_string())
}
fn ensure_session_mutation_allowed(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<(), meerkat_store::SessionStoreError> {
if self.session_was_unregistered(session_id) {
return Err(meerkat_store::SessionStoreError::Internal(format!(
"session {session_id} was unregistered from identity runtime state"
)));
}
if self.session_was_suspended(session_id) {
return Err(meerkat_store::SessionStoreError::Internal(format!(
"session {session_id} persistence is suspended during identity authority rotation"
)));
}
Ok(())
}
fn next_version(&self, session_id: &str) -> u64 {
let mut map = self
.versions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let counter = map
.entry(session_id.to_string())
.or_insert_with(|| AtomicU64::new(0));
counter.fetch_add(1, Ordering::Relaxed) + 1
}
fn restore_registration_state(
&self,
session_id: &meerkat_core::types::SessionId,
previous_registry: Option<SessionRuntimeState>,
) {
let key = session_id.to_string();
{
let mut registry = self
.session_registry
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
match previous_registry {
Some(state) => {
registry.insert(key, state);
}
None => {
registry.remove(&key);
}
}
}
}
fn current_version(
&self,
session_id: &meerkat_core::types::SessionId,
) -> super::types::CheckpointVersion {
let map = self
.versions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let version = map
.get(&session_id.to_string())
.map(|counter| counter.load(Ordering::Relaxed))
.unwrap_or(0);
super::types::CheckpointVersion::new(version)
}
fn lookup_session(&self, session_id: &str) -> Option<SessionRuntimeState> {
let registry = self
.session_registry
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
registry.get(session_id).cloned()
}
async fn load_persisted_session(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<Option<meerkat_core::Session>, meerkat_store::SessionStoreError> {
if let Some(session) = self.load_head_canonical_session(id).await? {
return Ok(Some(session));
}
Ok(self
.load_persisted_session_with_bytes(id)
.await?
.map(|(session, _)| session))
}
async fn load_head_canonical_session(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<Option<meerkat_core::Session>, meerkat_store::SessionStoreError> {
let Some(incremental) = self.incremental.as_ref() else {
return Ok(None);
};
if self.parked_deltas.is_parked(id) {
return Ok(None);
}
let Some(session) = incremental.load_canonical_session(id).await? else {
return Ok(None);
};
if session.id() != id {
return Err(meerkat_store::SessionStoreError::Serialization(format!(
"continuity head row {id} materializes session {}",
session.id()
)));
}
Ok(Some(session))
}
async fn load_persisted_session_with_bytes(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<Option<(meerkat_core::Session, Vec<u8>)>, meerkat_store::SessionStoreError> {
let snapshot = self.store.load_session_snapshot(id).await.map_err(|e| {
meerkat_store::SessionStoreError::Internal(format!("continuity load: {e}"))
})?;
match snapshot {
Some(snap) => {
let session: meerkat_core::Session = serde_json::from_slice(&snap.data)
.map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))?;
if session.id() != id {
return Err(meerkat_store::SessionStoreError::Serialization(format!(
"continuity snapshot key {id} contains session {}",
session.id()
)));
}
Ok(Some((session, snap.data)))
}
None => Ok(None),
}
}
async fn lazy_adopt_legacy_snapshot(
&self,
id: &meerkat_core::types::SessionId,
session: meerkat_core::Session,
raw: &[u8],
) -> meerkat_core::Session {
if !session_is_legacy_unverified(&session) {
return session;
}
let Some(state) = self.lookup_session(&id.to_string()) else {
return session;
};
let observed_generation = meerkat_core::SessionGeneration::new(state.generation.get());
let observed_revision =
meerkat_core::SessionCheckpointRevision::new(state.checkpoint_version.get());
let adopted =
match meerkat_core::adopt_legacy_session(raw, observed_generation, observed_revision) {
Ok(adopted) => adopted,
Err(error) => {
tracing::warn!(
session_id = %id,
%error,
"lazy checkpoint adoption refused; passing the legacy document through"
);
return session;
}
};
match self
.save_registered_snapshot(id, adopted.serialized, state)
.await
{
Ok(committed_version) => {
tracing::info!(
session_id = %id,
observed_generation = observed_generation.get(),
observed_checkpoint_revision = observed_revision.get(),
committed_version = committed_version.get(),
"lazy checkpoint adoption stamped a legacy continuity snapshot at restore"
);
adopted.session
}
Err(error) => {
tracing::warn!(
session_id = %id,
%error,
"lazy checkpoint adoption could not persist the adopted bytes; \
passing the legacy document through"
);
session
}
}
}
async fn load_previous_session_for_save(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<Option<meerkat_core::Session>, meerkat_store::SessionStoreError> {
if let Some(session) = self.load_persisted_session(id).await? {
return Ok(Some(session));
}
let pending = self
.pending_unregistered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(&id.to_string())
.cloned();
pending
.map(|data| {
serde_json::from_slice(&data)
.map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))
})
.transpose()
}
async fn save_registered_snapshot(
&self,
session_id: &meerkat_core::types::SessionId,
data: Vec<u8>,
state: SessionRuntimeState,
) -> Result<super::types::CheckpointVersion, meerkat_store::SessionStoreError> {
let version = self.next_version(&session_id.to_string());
let checkpoint_version = super::types::CheckpointVersion::new(version);
let snapshot = super::types::SessionSnapshot { data };
self.store
.save_session_snapshot_owned(
state.identity,
session_id.clone(),
state.generation,
checkpoint_version,
state.fencing_token,
snapshot,
)
.await
.map_err(|e| {
meerkat_store::SessionStoreError::Internal(format!("continuity save: {e}"))
})?;
Ok(checkpoint_version)
}
}
#[async_trait]
impl meerkat::SessionStore for ContinuitySessionStoreAdapter {
async fn save(
&self,
session: &meerkat_core::Session,
) -> Result<(), meerkat_store::SessionStoreError> {
let _guard = self.lock_session(session.id()).await;
if self.session_was_superseded(session.id()) {
return Ok(());
}
self.ensure_session_mutation_allowed(session.id())?;
let snapshot = Arc::new(super::types::SessionSnapshot {
data: serde_json::to_vec(session)
.map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))?,
});
let sid_str = session.id().to_string();
let state = self.lookup_session(&sid_str);
if let Some(state) = state.as_ref() {
let candidate = SessionSnapshotMatchCandidate {
identity: state.identity.clone(),
session_id: session.id().clone(),
generation: state.generation,
checkpoint_version: self.current_version(session.id()),
fencing_token: state.fencing_token,
snapshot: snapshot.clone(),
};
let matches = self
.store
.session_snapshot_matches_current(candidate)
.await
.map_err(|e| {
meerkat_store::SessionStoreError::Internal(format!(
"continuity snapshot match: {e}"
))
})?;
if matches {
meerkat_core::session_store::append_only_save_guard(session, Some(session))?;
return Ok(());
}
}
match self.head_canonical_previous(session.id()).await? {
Some((previous_slim, stored_commits)) => {
meerkat_core::session_store::head_canonical_plain_save_guard(
session,
&previous_slim,
&stored_commits,
)?;
}
None => {
let previous = self.load_previous_session_for_save(session.id()).await?;
meerkat_core::session_store::append_only_save_guard(session, previous.as_ref())?;
}
}
let snapshot = Arc::try_unwrap(snapshot).unwrap_or_else(|snapshot| (*snapshot).clone());
let data = snapshot.data;
match state {
Some(state) => {
self.save_registered_snapshot(session.id(), data, state)
.await?;
}
None => {
tracing::warn!(
session_id = %sid_str,
"ContinuitySessionStoreAdapter: delaying save until runtime state is registered"
);
let mut pending = self
.pending_unregistered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
pending.insert(sid_str, data);
}
}
Ok(())
}
async fn save_transcript_rewrite(
&self,
session: &meerkat_core::Session,
commit: &meerkat_core::TranscriptRewriteCommit,
) -> Result<(), meerkat_store::SessionStoreError> {
let _guard = self.lock_session(session.id()).await;
if self.session_was_superseded(session.id()) {
return Ok(());
}
self.ensure_session_mutation_allowed(session.id())?;
if self
.save_head_canonical_transcript_rewrite(session, commit)
.await?
{
return Ok(());
}
let previous = self.load_previous_session_for_save(session.id()).await?;
meerkat_core::session_store::transcript_rewrite_save_guard(
session,
previous.as_ref(),
commit,
)?;
let data = serde_json::to_vec(session)
.map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))?;
let sid_str = session.id().to_string();
match self.lookup_session(&sid_str) {
Some(state) => {
self.save_registered_snapshot(session.id(), data, state)
.await?;
}
None => {
let mut pending = self
.pending_unregistered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
pending.insert(sid_str, data);
}
}
Ok(())
}
async fn save_authoritative_projection(
&self,
session: &meerkat_core::Session,
) -> Result<(), meerkat_store::SessionStoreError> {
let _guard = self.lock_session(session.id()).await;
if self.session_was_superseded(session.id()) {
return Ok(());
}
self.ensure_session_mutation_allowed(session.id())?;
let data = serde_json::to_vec(session)
.map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))?;
let sid_str = session.id().to_string();
match self.lookup_session(&sid_str) {
Some(state) => {
self.save_registered_snapshot(session.id(), data, state)
.await?;
}
None => {
let mut pending = self
.pending_unregistered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
pending.insert(sid_str, data);
}
}
Ok(())
}
async fn save_authoritative_projection_if_current_revision(
&self,
session: &meerkat_core::Session,
expected_current_revision: Option<String>,
) -> Result<(), meerkat_store::SessionStoreError> {
let _guard = self.lock_session(session.id()).await;
if self.session_was_superseded(session.id()) {
return Ok(());
}
self.ensure_session_mutation_allowed(session.id())?;
let previous = self.load_persisted_session(session.id()).await?;
meerkat_core::session_store::authoritative_projection_current_revision_guard(
session,
previous.as_ref(),
expected_current_revision.as_deref(),
)?;
let data = serde_json::to_vec(session)
.map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))?;
let sid_str = session.id().to_string();
match self.lookup_session(&sid_str) {
Some(state) => {
self.save_registered_snapshot(session.id(), data, state)
.await?;
Ok(())
}
None => {
let mut pending = self
.pending_unregistered
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
pending.insert(sid_str, data);
Ok(())
}
}
}
async fn load(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<Option<meerkat_core::Session>, meerkat_store::SessionStoreError> {
let _guard = self.lock_session(id).await;
if self.session_was_superseded(id) {
return Ok(None);
}
if let Some(session) = self.load_head_canonical_session(id).await? {
if !self.lazy_checkpoint_adoption || !session_is_legacy_unverified(&session) {
return Ok(Some(session));
}
let raw = serde_json::to_vec(&session)
.map_err(|e| meerkat_store::SessionStoreError::Serialization(e.to_string()))?;
return Ok(Some(
self.lazy_adopt_legacy_snapshot(id, session, &raw).await,
));
}
let Some((session, raw)) = self.load_persisted_session_with_bytes(id).await? else {
return Ok(None);
};
if !self.lazy_checkpoint_adoption {
return Ok(Some(session));
}
Ok(Some(
self.lazy_adopt_legacy_snapshot(id, session, &raw).await,
))
}
async fn list(
&self,
_filter: meerkat_store::SessionFilter,
) -> Result<Vec<meerkat_core::SessionMeta>, meerkat_store::SessionStoreError> {
Ok(Vec::new())
}
async fn delete(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<(), meerkat_store::SessionStoreError> {
let _guard = self.lock_session(id).await;
if self.session_was_superseded(id) {
return Ok(());
}
self.ensure_session_mutation_allowed(id)?;
let Some(session) = self.load_persisted_session(id).await? else {
self.forget_session(id).await;
return Ok(());
};
let current_revision = meerkat_core::session_store::session_projection_cas_token(&session)?;
let deleted = self
.store
.delete_session_snapshot_if_current_revision(id, ¤t_revision)
.await
.map_err(|e| {
meerkat_store::SessionStoreError::Internal(format!("continuity delete: {e}"))
})?;
if !deleted {
return Err(meerkat_store::SessionStoreError::Internal(format!(
"continuity delete did not remove session snapshot {id}"
)));
}
self.forget_session(id).await;
Ok(())
}
async fn delete_if_current_revision(
&self,
id: &meerkat_core::types::SessionId,
expected_current_revision: &str,
) -> Result<bool, meerkat_store::SessionStoreError> {
let _guard = self.lock_session(id).await;
if self.session_was_superseded(id) {
return Ok(false);
}
self.ensure_session_mutation_allowed(id)?;
let Some(session) = self.load_persisted_session(id).await? else {
self.forget_session(id).await;
return Ok(false);
};
let current_revision = meerkat_core::session_store::session_projection_cas_token(&session)?;
if current_revision != expected_current_revision {
return Ok(false);
}
let deleted = self
.store
.delete_session_snapshot_if_current_revision(id, expected_current_revision)
.await
.map_err(|e| {
meerkat_store::SessionStoreError::Internal(format!(
"continuity delete_if_current_revision: {e}"
))
})?;
if deleted {
self.forget_session(id).await;
}
Ok(deleted)
}
fn as_incremental(
self: Arc<Self>,
) -> Option<Arc<dyn meerkat_core::session_store::IncrementalSessionStore>> {
let inner = self.incremental.clone()?;
Some(Arc::new(ContinuityIncrementalSessionStore {
adapter: self,
inner,
}))
}
}
pub struct ContinuityIncrementalSessionStore {
adapter: Arc<ContinuitySessionStoreAdapter>,
inner: Arc<dyn super::contracts::ContinuityIncrementalSessions>,
}
fn rewrite_record_from_session_bodies(
session: &meerkat_core::Session,
commit: &meerkat_core::TranscriptRewriteCommit,
) -> Result<meerkat_core::TranscriptRewriteRecord, meerkat_store::SessionStoreError> {
let parent_body = session
.transcript_revision_body(&commit.parent_revision)?
.ok_or_else(
|| meerkat_store::SessionStoreError::InvalidTranscriptRewrite {
id: session.id().clone(),
reason: format!(
"incoming rewrite omitted parent revision body {}",
commit.parent_revision
),
},
)?;
let revision_body = session
.transcript_revision_body(&commit.revision)?
.ok_or_else(
|| meerkat_store::SessionStoreError::InvalidTranscriptRewrite {
id: session.id().clone(),
reason: format!(
"incoming rewrite omitted new revision body {}",
commit.revision
),
},
)?;
meerkat_core::TranscriptRewriteRecord::new(commit.clone(), parent_body, revision_body).map_err(
|err| meerkat_store::SessionStoreError::InvalidTranscriptRewrite {
id: session.id().clone(),
reason: format!("transcript rewrite record failed validation: {err}"),
},
)
}
enum DeltaRoute {
Durable(super::contracts::ContinuityWriteCursor),
Park,
SupersededNoOp,
}
impl ContinuityIncrementalSessionStore {
fn route_delta_write(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<DeltaRoute, meerkat_store::SessionStoreError> {
if self.adapter.session_was_superseded(id) {
return Ok(DeltaRoute::SupersededNoOp);
}
self.adapter.ensure_session_mutation_allowed(id)?;
match self.adapter.lookup_session(&id.to_string()) {
Some(state) => Ok(DeltaRoute::Durable(self.adapter.write_cursor(id, &state))),
None => Ok(DeltaRoute::Park),
}
}
fn parked(&self) -> &meerkat_store::MemoryStore {
self.adapter.parked_deltas.reads()
}
fn reads_parked(&self, id: &meerkat_core::types::SessionId) -> bool {
self.adapter.parked_deltas.is_parked(id)
}
}
#[async_trait]
impl meerkat::SessionStore for ContinuityIncrementalSessionStore {
async fn save(
&self,
session: &meerkat_core::Session,
) -> Result<(), meerkat_store::SessionStoreError> {
self.adapter.save(session).await
}
async fn save_transcript_rewrite(
&self,
session: &meerkat_core::Session,
commit: &meerkat_core::TranscriptRewriteCommit,
) -> Result<(), meerkat_store::SessionStoreError> {
self.adapter.save_transcript_rewrite(session, commit).await
}
async fn save_authoritative_projection(
&self,
session: &meerkat_core::Session,
) -> Result<(), meerkat_store::SessionStoreError> {
self.adapter.save_authoritative_projection(session).await
}
async fn save_authoritative_projection_if_current_revision(
&self,
session: &meerkat_core::Session,
expected_current_revision: Option<String>,
) -> Result<(), meerkat_store::SessionStoreError> {
self.adapter
.save_authoritative_projection_if_current_revision(session, expected_current_revision)
.await
}
async fn load(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<Option<meerkat_core::Session>, meerkat_store::SessionStoreError> {
self.adapter.load(id).await
}
async fn list(
&self,
filter: meerkat_store::SessionFilter,
) -> Result<Vec<meerkat_core::SessionMeta>, meerkat_store::SessionStoreError> {
self.adapter.list(filter).await
}
async fn delete(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<(), meerkat_store::SessionStoreError> {
self.adapter.delete(id).await
}
async fn delete_if_current_revision(
&self,
id: &meerkat_core::types::SessionId,
expected_current_revision: &str,
) -> Result<bool, meerkat_store::SessionStoreError> {
self.adapter
.delete_if_current_revision(id, expected_current_revision)
.await
}
fn as_incremental(
self: Arc<Self>,
) -> Option<Arc<dyn meerkat_core::session_store::IncrementalSessionStore>> {
Some(self)
}
}
#[async_trait]
impl meerkat_core::session_store::IncrementalSessionStore for ContinuityIncrementalSessionStore {
async fn append_messages(
&self,
id: &meerkat_core::types::SessionId,
strand: &meerkat_core::session_store::TranscriptStrandId,
base_seq: u64,
messages: &[meerkat_core::Message],
) -> Result<(), meerkat_store::SessionStoreError> {
let _guard = self.adapter.lock_session(id).await;
match self.route_delta_write(id)? {
DeltaRoute::Durable(cursor) => {
self.inner
.append_messages(&cursor, id, strand, base_seq, messages)
.await
}
DeltaRoute::Park => {
self.adapter
.parked_deltas
.park_append(id, strand, base_seq, messages)
.await
}
DeltaRoute::SupersededNoOp => Ok(()),
}
}
async fn commit_rewrite(
&self,
id: &meerkat_core::types::SessionId,
record: &meerkat_core::TranscriptRewriteRecord,
expected: meerkat_core::session_store::SessionHeadCas,
) -> Result<meerkat_core::session_store::SessionHead, meerkat_store::SessionStoreError> {
let _guard = self.adapter.lock_session(id).await;
match self.route_delta_write(id)? {
DeltaRoute::Durable(cursor) => {
self.inner
.commit_rewrite(&cursor, id, record, expected)
.await
}
DeltaRoute::Park => {
self.adapter
.parked_deltas
.park_rewrite(id, record, expected)
.await
}
DeltaRoute::SupersededNoOp => Err(meerkat_store::SessionStoreError::Internal(format!(
"session {id} was superseded by a committed continuity reset; \
rewrite commits are refused"
))),
}
}
async fn save_head(
&self,
head: &meerkat_core::session_store::SessionHead,
expected: meerkat_core::session_store::SessionHeadCas,
) -> Result<(), meerkat_store::SessionStoreError> {
let _guard = self.adapter.lock_session(&head.id).await;
match self.route_delta_write(&head.id)? {
DeltaRoute::Durable(cursor) => self.inner.save_head(&cursor, head, expected).await,
DeltaRoute::Park => self.adapter.parked_deltas.park_head(head, expected).await,
DeltaRoute::SupersededNoOp => Ok(()),
}
}
async fn load_head(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<Option<meerkat_core::session_store::SessionHead>, meerkat_store::SessionStoreError>
{
if self.adapter.session_was_superseded(id) {
return Ok(None);
}
if self.reads_parked(id) {
return meerkat_core::session_store::IncrementalSessionStore::load_head(
self.parked(),
id,
)
.await;
}
self.inner.load_head(id).await
}
async fn load_messages(
&self,
id: &meerkat_core::types::SessionId,
strand: &meerkat_core::session_store::TranscriptStrandId,
range: std::ops::Range<u64>,
) -> Result<Vec<meerkat_core::Message>, meerkat_store::SessionStoreError> {
if self.adapter.session_was_superseded(id) {
return Ok(Vec::new());
}
if self.reads_parked(id) {
return meerkat_core::session_store::IncrementalSessionStore::load_messages(
self.parked(),
id,
strand,
range,
)
.await;
}
self.inner.load_messages(id, strand, range).await
}
async fn load_rewrites(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<Vec<meerkat_core::TranscriptRewriteRecord>, meerkat_store::SessionStoreError> {
if self.adapter.session_was_superseded(id) {
return Ok(Vec::new());
}
if self.reads_parked(id) {
return meerkat_core::session_store::IncrementalSessionStore::load_rewrites(
self.parked(),
id,
)
.await;
}
self.inner.load_rewrites(id).await
}
}
pub struct SessionHookCustomizerAdapter {
hook: Arc<dyn SessionHook>,
}
impl SessionHookCustomizerAdapter {
pub fn new(hook: Arc<dyn SessionHook>) -> Self {
Self { hook }
}
}
#[async_trait]
impl AgentCustomizer for SessionHookCustomizerAdapter {
async fn customize_build(
&self,
_context: &AgentBuildContext,
spec: &DurableAgentSpec,
draft: &mut AgentBuildDraft,
) -> Result<(), CustomizerError> {
let mut req = meerkat_core::service::CreateSessionRequest {
model: draft.model.clone().unwrap_or_default(),
prompt: meerkat_core::ContentInput::Text(String::new()),
system_prompt: match draft.system_prompt.clone() {
Some(prompt) => meerkat_core::config::SystemPromptOverride::Set(prompt),
None => meerkat_core::config::SystemPromptOverride::Inherit,
},
max_tokens: None,
event_tx: None,
initial_turn: meerkat_core::service::InitialTurnPolicy::Defer,
build: None,
labels: if draft.labels.is_empty() {
None
} else {
Some(draft.labels.clone())
},
deferred_prompt_policy: meerkat_core::service::DeferredPromptPolicy::default(),
injected_context: Vec::new(),
};
let prompt_before = req.prompt.clone();
let max_tokens_before = req.max_tokens;
let event_tx_was_some = req.event_tx.is_some();
let initial_turn_before = req.initial_turn;
let build_before_is_none = req.build.is_none();
self.hook
.before_create(&mut req)
.await
.map_err(|e| CustomizerError::BuildFailed(format!("session hook: {e}")))?;
let mut unsupported_mutations: Vec<&str> = Vec::new();
if req.prompt != prompt_before {
unsupported_mutations.push("prompt");
}
if req.max_tokens != max_tokens_before {
unsupported_mutations.push("max_tokens");
}
if req.event_tx.is_some() != event_tx_was_some {
unsupported_mutations.push("event_tx");
}
if req.initial_turn != initial_turn_before {
unsupported_mutations.push("initial_turn");
}
if let Some(ref build) = req.build {
if build_before_is_none {
unsupported_mutations.push("build");
if build.resume_session.is_some() {
unsupported_mutations.push("build.resume_session");
}
} else if build.resume_session.is_some() {
unsupported_mutations.push("build.resume_session");
}
}
if !unsupported_mutations.is_empty() {
tracing::warn!(
identity = %spec.identity,
fields = ?unsupported_mutations,
"SessionHook mutated unsupported CreateSessionRequest fields — \
these mutations are NOT applied in the identity-first model. \
Migrate to AgentCustomizer."
);
}
if !req.model.is_empty() {
draft.model = Some(req.model);
}
draft.system_prompt = req.system_prompt.as_set_prompt().map(ToString::to_string);
draft.labels = req.labels.unwrap_or_default();
Ok(())
}
async fn after_create(
&self,
_identity: &AgentIdentity,
session_id: &meerkat_core::types::SessionId,
context: &SessionCreatedContext,
) -> Result<(), CustomizerError> {
self.hook.after_create(session_id, context).await;
Ok(())
}
}
#[cfg(test)]
#[allow(clippy::expect_used, clippy::panic)]
mod tests {
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering as AtomicOrdering};
use std::time::Duration;
use serde_json::json;
use super::super::contracts::ContinuityStore;
use super::super::local_store::LocalContinuityStore;
use super::super::types::{
AgentIdentity, AgentRuntimeId, CheckpointVersion, ContinuityGeneration, ContinuityRecord,
ContinuityResolveState, ContinuityStoreError, FencingToken, SessionSnapshot,
};
use super::*;
struct FailSaveContinuityStore {
inner: Arc<LocalContinuityStore>,
fail_save: AtomicBool,
commit_then_fail_save: AtomicBool,
fail_delete_once: AtomicBool,
block_next_save: AtomicBool,
save_entered: tokio::sync::Semaphore,
release_save: tokio::sync::Semaphore,
}
impl FailSaveContinuityStore {
fn new(inner: Arc<LocalContinuityStore>) -> Self {
Self {
inner,
fail_save: AtomicBool::new(false),
commit_then_fail_save: AtomicBool::new(false),
fail_delete_once: AtomicBool::new(false),
block_next_save: AtomicBool::new(false),
save_entered: tokio::sync::Semaphore::new(0),
release_save: tokio::sync::Semaphore::new(0),
}
}
fn fail_saves(&self, fail: bool) {
self.fail_save.store(fail, AtomicOrdering::SeqCst);
}
fn commit_then_fail_next_save(&self) {
self.commit_then_fail_save
.store(true, AtomicOrdering::SeqCst);
}
fn fail_next_delete(&self) {
self.fail_delete_once.store(true, AtomicOrdering::SeqCst);
}
fn block_one_save(&self) {
self.block_next_save.store(true, AtomicOrdering::SeqCst);
}
async fn wait_for_blocked_save(&self) {
self.save_entered
.acquire()
.await
.expect("save-entered semaphore remains open")
.forget();
}
fn release_blocked_save(&self) {
self.release_save.add_permits(1);
}
}
struct ConcurrentLoadStore {
in_flight: AtomicUsize,
max_in_flight: AtomicUsize,
rendezvous: tokio::sync::Barrier,
}
impl ConcurrentLoadStore {
fn new(expected_concurrent_loads: usize) -> Self {
Self {
in_flight: AtomicUsize::new(0),
max_in_flight: AtomicUsize::new(0),
rendezvous: tokio::sync::Barrier::new(expected_concurrent_loads),
}
}
}
#[async_trait]
impl ContinuityStore for FailSaveContinuityStore {
async fn resolve_many(
&self,
identities: &[AgentIdentity],
) -> Result<
std::collections::BTreeMap<AgentIdentity, ContinuityResolveState>,
ContinuityStoreError,
> {
self.inner.resolve_many(identities).await
}
async fn load_session_snapshot(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<Option<SessionSnapshot>, ContinuityStoreError> {
self.inner.load_session_snapshot(session_id).await
}
async fn delete_session_snapshot_if_current_revision(
&self,
session_id: &meerkat_core::types::SessionId,
expected_current_revision: &str,
) -> Result<bool, ContinuityStoreError> {
if self.fail_delete_once.swap(false, AtomicOrdering::SeqCst) {
return Err(ContinuityStoreError::Io(
"synthetic superseded snapshot delete failure".to_string(),
));
}
self.inner
.delete_session_snapshot_if_current_revision(session_id, expected_current_revision)
.await
}
async fn save_session_snapshot(
&self,
identity: &AgentIdentity,
session_id: &meerkat_core::types::SessionId,
generation: ContinuityGeneration,
version: CheckpointVersion,
fencing_token: FencingToken,
snapshot: &SessionSnapshot,
) -> Result<(), ContinuityStoreError> {
if self.block_next_save.swap(false, AtomicOrdering::SeqCst) {
self.save_entered.add_permits(1);
self.release_save
.acquire()
.await
.expect("release-save semaphore remains open")
.forget();
}
if self.fail_save.load(AtomicOrdering::SeqCst) {
return Err(ContinuityStoreError::Io("forced save failure".to_string()));
}
self.inner
.save_session_snapshot(
identity,
session_id,
generation,
version,
fencing_token,
snapshot,
)
.await?;
if self
.commit_then_fail_save
.swap(false, AtomicOrdering::SeqCst)
{
return Err(ContinuityStoreError::Io(
"synthetic lost save acknowledgement".to_string(),
));
}
Ok(())
}
async fn upsert_continuity_record(
&self,
record: &ContinuityRecord,
fencing_token: FencingToken,
) -> Result<(), ContinuityStoreError> {
self.inner
.upsert_continuity_record(record, fencing_token)
.await
}
async fn delete_continuity_record(
&self,
identity: &AgentIdentity,
fencing_token: FencingToken,
) -> Result<(), ContinuityStoreError> {
self.inner
.delete_continuity_record(identity, fencing_token)
.await
}
}
#[async_trait]
impl ContinuityStore for ConcurrentLoadStore {
async fn resolve_many(
&self,
identities: &[AgentIdentity],
) -> Result<
std::collections::BTreeMap<AgentIdentity, ContinuityResolveState>,
ContinuityStoreError,
> {
Ok(identities
.iter()
.cloned()
.map(|identity| (identity, ContinuityResolveState::Uninitialized))
.collect())
}
async fn load_session_snapshot(
&self,
_session_id: &meerkat_core::types::SessionId,
) -> Result<Option<SessionSnapshot>, ContinuityStoreError> {
let now = self.in_flight.fetch_add(1, AtomicOrdering::SeqCst) + 1;
self.max_in_flight.fetch_max(now, AtomicOrdering::SeqCst);
self.rendezvous.wait().await;
self.in_flight.fetch_sub(1, AtomicOrdering::SeqCst);
Ok(None)
}
async fn save_session_snapshot(
&self,
_identity: &AgentIdentity,
_session_id: &meerkat_core::types::SessionId,
_generation: ContinuityGeneration,
_version: CheckpointVersion,
_fencing_token: FencingToken,
_snapshot: &SessionSnapshot,
) -> Result<(), ContinuityStoreError> {
Ok(())
}
async fn upsert_continuity_record(
&self,
_record: &ContinuityRecord,
_fencing_token: FencingToken,
) -> Result<(), ContinuityStoreError> {
Ok(())
}
async fn delete_continuity_record(
&self,
_identity: &AgentIdentity,
_fencing_token: FencingToken,
) -> Result<(), ContinuityStoreError> {
Ok(())
}
}
#[tokio::test]
async fn continuity_session_store_adapter_parallelizes_different_sessions() {
let store = Arc::new(ConcurrentLoadStore::new(2));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let first = meerkat_core::Session::new();
let second = meerkat_core::Session::new();
assert_ne!(first.id(), second.id());
tokio::time::timeout(std::time::Duration::from_secs(1), async {
let (first_result, second_result) = tokio::join!(
meerkat::SessionStore::save(&adapter, &first),
meerkat::SessionStore::save(&adapter, &second),
);
first_result.expect("first save");
second_result.expect("second save");
})
.await
.expect("different session IDs must not share one global save lock");
assert_eq!(
store.max_in_flight.load(AtomicOrdering::SeqCst),
2,
"both independent session loads should overlap"
);
}
#[tokio::test]
async fn continuity_session_store_adapter_exact_resave_is_a_noop() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let session = meerkat_core::Session::new();
let identity = AgentIdentity::parse("agent:exact-resave").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:exact-resave:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(3);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity: identity.clone(),
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("initial save");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("exact resave");
let resolved = store
.resolve_many(std::slice::from_ref(&identity))
.await
.expect("resolve");
let ContinuityResolveState::Ready { record } = resolved.get(&identity).expect("record")
else {
panic!("expected ready record");
};
assert_eq!(
record.checkpoint_version,
CheckpointVersion::new(1),
"an exact durable resave must not manufacture a new checkpoint"
);
}
#[tokio::test]
async fn continuity_session_store_adapter_abandons_only_superseded_snapshot() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let identity = AgentIdentity::parse("agent:reset-abandon").expect("identity");
let old_session = meerkat_core::Session::new();
let old_record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:reset-abandon:0")
.expect("old runtime id"),
session_id: old_session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let old_fence = FencingToken::new(3);
store
.upsert_continuity_record(&old_record, old_fence)
.await
.expect("seed old record");
adapter
.register_session(
old_session.id(),
SessionRuntimeState {
identity: identity.clone(),
generation: old_record.generation,
fencing_token: old_fence,
checkpoint_version: old_record.checkpoint_version,
},
)
.await
.expect("register old session");
meerkat::SessionStore::save(&adapter, &old_session)
.await
.expect("persist old session");
let replacement_session = meerkat_core::types::SessionId::new();
let replacement_record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:reset-abandon:1")
.expect("replacement runtime id"),
session_id: replacement_session.clone(),
generation: ContinuityGeneration::new(1),
checkpoint_version: CheckpointVersion::new(0),
};
let replacement_fence = FencingToken::new(4);
store
.upsert_continuity_record(&replacement_record, replacement_fence)
.await
.expect("commit replacement record");
adapter
.abandon_superseded_session(old_session.id())
.await
.expect("abandon old projection");
assert!(
store
.load_session_snapshot(old_session.id())
.await
.expect("load old snapshot")
.is_none(),
"the exact superseded snapshot must be CAS-deleted"
);
assert_eq!(
store
.resolve_many(std::slice::from_ref(&identity))
.await
.expect("resolve replacement")
.get(&identity),
Some(&ContinuityResolveState::Ready {
record: replacement_record
}),
"abandonment must not disturb the replacement continuity head"
);
meerkat::SessionStore::save_authoritative_projection(&adapter, &old_session)
.await
.expect("terminal superseded projection is acknowledged");
assert!(
store
.load_session_snapshot(old_session.id())
.await
.expect("reload old snapshot")
.is_none()
);
adapter
.unregister_session(old_session.id())
.await
.expect("finalize old session authority");
assert!(
meerkat::SessionStore::save(&adapter, &old_session)
.await
.is_err(),
"late writes must fail closed after structural member absence"
);
}
#[tokio::test]
async fn continuity_session_store_adapter_retries_failed_superseded_snapshot_cas() {
let inner = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let store = Arc::new(FailSaveContinuityStore::new(inner.clone()));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let identity = AgentIdentity::parse("agent:reset-abandon-retry").expect("identity");
let session = meerkat_core::Session::new();
let old_record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:reset-abandon-retry:0")
.expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let old_fence = FencingToken::new(7);
store
.upsert_continuity_record(&old_record, old_fence)
.await
.expect("seed old record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity: identity.clone(),
generation: old_record.generation,
fencing_token: old_fence,
checkpoint_version: old_record.checkpoint_version,
},
)
.await
.expect("register old session");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("persist old session");
let replacement = ContinuityRecord {
identity,
agent_runtime_id: AgentRuntimeId::parse("rt:agent:reset-abandon-retry:1")
.expect("replacement runtime id"),
session_id: meerkat_core::types::SessionId::new(),
generation: ContinuityGeneration::new(1),
checkpoint_version: CheckpointVersion::new(0),
};
store
.upsert_continuity_record(&replacement, FencingToken::new(8))
.await
.expect("commit replacement");
store.fail_next_delete();
assert!(
adapter
.abandon_superseded_session(session.id())
.await
.is_err(),
"the injected CAS failure must remain visible"
);
assert!(adapter.session_was_suspended(session.id()));
assert!(
inner
.load_session_snapshot(session.id())
.await
.expect("snapshot after failed abandon")
.is_some()
);
adapter
.abandon_superseded_session(session.id())
.await
.expect("retry exact CAS abandon");
assert!(adapter.session_was_superseded(session.id()));
assert!(
inner
.load_session_snapshot(session.id())
.await
.expect("snapshot after retry")
.is_none()
);
}
#[tokio::test]
async fn continuity_session_store_adapter_exact_bytes_do_not_mask_a_newer_fence() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let session = meerkat_core::Session::new();
let identity = AgentIdentity::parse("agent:exact-stale-fence").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:exact-stale-fence:0")
.expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(3);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity: identity.clone(),
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("initial save");
store
.upsert_continuity_record(&record, FencingToken::new(4))
.await
.expect("advance durable fence");
let error = meerkat::SessionStore::save(&adapter, &session)
.await
.expect_err("the stale registered fence must still be rejected");
assert!(
error.to_string().contains("stale fencing token"),
"unexpected error: {error}"
);
}
#[tokio::test]
async fn continuity_session_store_adapter_seeds_registered_checkpoint_version() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let session = meerkat_core::Session::new();
let identity = AgentIdentity::parse("agent:restored").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:restored:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(2),
checkpoint_version: CheckpointVersion::new(5),
};
let fencing_token = FencingToken::new(9);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity: identity.clone(),
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("save should advance from restored checkpoint");
let effective_version = adapter
.register_session(
session.id(),
SessionRuntimeState {
identity: identity.clone(),
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("post-save register should report advanced version");
assert_eq!(effective_version, CheckpointVersion::new(6));
let resolved = store
.resolve_many(std::slice::from_ref(&identity))
.await
.expect("resolve");
let ContinuityResolveState::Ready { record } = resolved.get(&identity).expect("record")
else {
panic!("expected ready record");
};
assert_eq!(record.checkpoint_version, CheckpointVersion::new(6));
}
#[tokio::test]
async fn continuity_session_store_adapter_flushes_pending_save_under_registered_identity() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let session = meerkat_core::Session::new();
let identity = AgentIdentity::parse("agent:fresh").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:fresh:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(3);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("unregistered save should be delayed, not written under fallback identity");
assert!(
store
.load_session_snapshot(session.id())
.await
.expect("load before register")
.is_none(),
"unregistered save must not be visible in continuity store"
);
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity: identity.clone(),
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register flushes pending");
assert!(
store
.load_session_snapshot(session.id())
.await
.expect("load after register")
.is_some(),
"pending save should flush under the registered identity"
);
let resolved = store
.resolve_many(std::slice::from_ref(&identity))
.await
.expect("resolve");
let ContinuityResolveState::Ready { record } = resolved.get(&identity).expect("record")
else {
panic!("expected ready record");
};
assert_eq!(record.checkpoint_version, CheckpointVersion::new(1));
}
#[tokio::test]
async fn continuity_session_store_adapter_rejects_saves_after_unregister() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let session = meerkat_core::Session::new();
let identity = AgentIdentity::parse("agent:retired").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:retired:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(9);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity: identity.clone(),
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
adapter
.unregister_session(session.id())
.await
.expect("unregister");
let err = meerkat::SessionStore::save(&adapter, &session)
.await
.expect_err("post-unregister save must fail closed");
assert!(
err.to_string().contains("was unregistered"),
"unexpected error: {err}"
);
assert!(
store
.load_session_snapshot(session.id())
.await
.expect("load")
.is_none(),
"post-unregister save must not be queued as pending"
);
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity,
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("registering the same id later should not flush stale pending data");
assert!(
store
.load_session_snapshot(session.id())
.await
.expect("load after re-register")
.is_none(),
"stale post-unregister save must not flush on a later registration"
);
}
#[tokio::test]
async fn continuity_session_store_adapter_suspension_blocks_every_mutation_until_reregister() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let mut session = meerkat_core::Session::new();
session.append_external_user_content(meerkat_core::ContentInput::Text(
"before rotation".to_string(),
));
let identity = AgentIdentity::parse("agent:suspended").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:suspended:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let old_token = FencingToken::new(7);
store
.upsert_continuity_record(&record, old_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity: identity.clone(),
generation: record.generation,
fencing_token: old_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("seed snapshot");
let parent_revision = session.transcript_revision().expect("parent revision");
let mut rewritten = session.clone();
let rewrite_commit = rewritten
.commit_transcript_rewrite(
meerkat_core::TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
vec![meerkat_core::Message::User(
meerkat_core::UserMessage::text("rewritten".to_string()),
)],
meerkat_core::TranscriptRewriteReason::new("rotation-test"),
Some("mobkit-test".to_string()),
Some(parent_revision),
)
.expect("rewrite commit");
adapter
.suspend_session(session.id())
.await
.expect("suspend");
let save_error = meerkat::SessionStore::save(&adapter, &session)
.await
.expect_err("ordinary save must fail while suspended");
let rewrite_error =
meerkat::SessionStore::save_transcript_rewrite(&adapter, &rewritten, &rewrite_commit)
.await
.expect_err("transcript rewrite must fail while suspended");
let projection_error =
meerkat::SessionStore::save_authoritative_projection(&adapter, &session)
.await
.expect_err("authoritative projection must fail while suspended");
let projection_cas_error =
meerkat::SessionStore::save_authoritative_projection_if_current_revision(
&adapter, &session, None,
)
.await
.expect_err("authoritative projection CAS must fail while suspended");
let delete_error = meerkat::SessionStore::delete(&adapter, session.id())
.await
.expect_err("delete must fail while suspended");
let delete_cas_error = meerkat::SessionStore::delete_if_current_revision(
&adapter,
session.id(),
"row-sha256:any",
)
.await
.expect_err("delete CAS must fail while suspended");
for error in [
save_error,
rewrite_error,
projection_error,
projection_cas_error,
delete_error,
delete_cas_error,
] {
assert!(
error.to_string().contains("persistence is suspended"),
"unexpected suspension error: {error}"
);
}
let new_token = FencingToken::new(8);
store
.upsert_continuity_record(&record, new_token)
.await
.expect("publish replacement authority");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity,
generation: record.generation,
fencing_token: new_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("re-register replacement authority");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("writes resume only after exact replacement registration");
}
#[tokio::test]
async fn continuity_session_store_adapter_suspension_drains_admitted_save() {
let inner = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let store = Arc::new(FailSaveContinuityStore::new(inner.clone()));
let adapter = Arc::new(ContinuitySessionStoreAdapter::new(store.clone()));
let session = meerkat_core::Session::new();
let identity = AgentIdentity::parse("agent:suspend-drain").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:suspend-drain:0")
.expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(17);
inner
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity,
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
store.block_one_save();
let save_adapter = adapter.clone();
let save_session = session.clone();
let save_task = tokio::spawn(async move {
meerkat::SessionStore::save(save_adapter.as_ref(), &save_session).await
});
store.wait_for_blocked_save().await;
let suspend_adapter = adapter.clone();
let session_id = session.id().clone();
let mut suspend_task =
tokio::spawn(async move { suspend_adapter.suspend_session(&session_id).await });
assert!(
tokio::time::timeout(Duration::from_millis(20), &mut suspend_task)
.await
.is_err(),
"suspension must wait for the already-admitted save to leave the session lock"
);
store.release_blocked_save();
save_task
.await
.expect("save task joins")
.expect("admitted save completes before suspension");
suspend_task
.await
.expect("suspend task joins")
.expect("suspension completes after drain");
let error = meerkat::SessionStore::save(adapter.as_ref(), &session)
.await
.expect_err("later saves must see the suspension barrier");
assert!(error.to_string().contains("persistence is suspended"));
}
#[tokio::test]
async fn continuity_session_store_adapter_register_keeps_pending_snapshot_on_flush_failure() {
let inner = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let fail_store = Arc::new(FailSaveContinuityStore::new(inner.clone()));
let adapter = ContinuitySessionStoreAdapter::new(fail_store.clone());
let mut session = meerkat_core::Session::new();
session.set_metadata("pending", json!(true));
let identity = AgentIdentity::parse("agent:pending-fail").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:pending-fail:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(14);
inner
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("pending save");
fail_store.fail_saves(true);
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity: identity.clone(),
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect_err("forced pending flush failure");
assert!(
meerkat::SessionStore::load(&adapter, session.id())
.await
.expect("load after failed register")
.is_none(),
"failed register must not leave a synthetic registered session"
);
fail_store.fail_saves(false);
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity: identity.clone(),
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("retry register should flush preserved pending snapshot");
let loaded = meerkat::SessionStore::load(&adapter, session.id())
.await
.expect("load after retry")
.expect("snapshot");
assert_eq!(loaded.metadata().get("pending"), Some(&json!(true)));
}
#[tokio::test]
async fn continuity_session_store_adapter_lost_ack_consumes_checkpoint_version() {
let inner = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let fail_store = Arc::new(FailSaveContinuityStore::new(inner.clone()));
let adapter = ContinuitySessionStoreAdapter::new(fail_store.clone());
let session = meerkat_core::Session::new();
let identity = AgentIdentity::parse("agent:lost-ack").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:lost-ack:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(23);
inner
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("queue pre-registration snapshot");
let state = SessionRuntimeState {
identity: identity.clone(),
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
};
fail_store.commit_then_fail_next_save();
adapter
.register_session(session.id(), state.clone())
.await
.expect_err("first flush commits version 1 but loses its acknowledgement");
let effective = adapter
.register_session(session.id(), state)
.await
.expect("retry must allocate a fresh checkpoint version");
assert_eq!(effective, CheckpointVersion::new(2));
let resolved = inner
.resolve_many(std::slice::from_ref(&identity))
.await
.expect("resolve");
let ContinuityResolveState::Ready { record } = resolved.get(&identity).expect("record")
else {
panic!("expected ready record");
};
assert_eq!(record.checkpoint_version, CheckpointVersion::new(2));
}
#[tokio::test]
async fn lazy_adoption_failed_save_is_not_observable_and_retry_stamps_durable_cursor() {
let inner = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let fail_store = Arc::new(FailSaveContinuityStore::new(inner.clone()));
let adapter = ContinuitySessionStoreAdapter::new(fail_store.clone())
.with_lazy_checkpoint_adoption(true);
let session = meerkat_core::Session::new();
let sid = session.id().clone();
let legacy = serde_json::to_vec(&session).expect("serialize legacy session");
let identity = AgentIdentity::parse("agent:lazy-fail").expect("identity");
let fencing_token = FencingToken::new(5);
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:lazy-fail:0").expect("runtime id"),
session_id: sid.clone(),
generation: ContinuityGeneration::new(3),
checkpoint_version: CheckpointVersion::new(0),
};
inner
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
inner
.save_session_snapshot(
&identity,
&sid,
ContinuityGeneration::new(3),
CheckpointVersion::new(4),
fencing_token,
&SessionSnapshot {
data: legacy.clone(),
},
)
.await
.expect("seed legacy snapshot");
adapter
.register_session(
&sid,
SessionRuntimeState {
identity: identity.clone(),
generation: ContinuityGeneration::new(3),
fencing_token,
checkpoint_version: CheckpointVersion::new(4),
},
)
.await
.expect("register session");
fail_store.fail_saves(true);
let loaded = meerkat::SessionStore::load(&adapter, &sid)
.await
.expect("load under failing save")
.expect("session present");
assert!(
matches!(
loaded.try_checkpoint_state().expect("checkpoint state"),
meerkat_core::SessionCheckpointState::LegacyUnverified { .. }
),
"a failed adoption save must pass the legacy document through"
);
let durable = inner
.load_session_snapshot(&sid)
.await
.expect("durable load")
.expect("snapshot present");
assert_eq!(
durable.data, legacy,
"a failed adoption save must leave the durable legacy bytes untouched"
);
fail_store.fail_saves(false);
let adopted = meerkat::SessionStore::load(&adapter, &sid)
.await
.expect("load after recovery")
.expect("session present");
let stamp = match adopted.try_checkpoint_state().expect("checkpoint state") {
meerkat_core::SessionCheckpointState::Verified(stamp) => stamp,
other => panic!("expected an adopted document, got {other:?}"),
};
assert_eq!(stamp.generation(), meerkat_core::SessionGeneration::new(3));
assert_eq!(
stamp.checkpoint_revision(),
meerkat_core::SessionCheckpointRevision::new(4),
"the stamp must bind the durable cursor, never a revision that only \
the in-memory allocator ever saw"
);
let durable = inner
.load_session_snapshot(&sid)
.await
.expect("durable load after adoption")
.expect("snapshot present");
let durable_session: meerkat_core::Session =
serde_json::from_slice(&durable.data).expect("decode adopted snapshot");
match durable_session
.try_checkpoint_state()
.expect("durable checkpoint state")
{
meerkat_core::SessionCheckpointState::Verified(durable_stamp) => {
assert_eq!(durable_stamp, stamp);
}
other => panic!("expected a verified durable document, got {other:?}"),
}
}
#[tokio::test]
async fn continuity_session_store_adapter_rejects_owner_generation_and_fence_regression() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store);
let session = meerkat_core::Session::new();
let first = SessionRuntimeState {
identity: AgentIdentity::parse("agent:owner-a").expect("identity"),
generation: ContinuityGeneration::new(4),
fencing_token: FencingToken::new(10),
checkpoint_version: CheckpointVersion::new(7),
};
adapter
.register_session(session.id(), first.clone())
.await
.expect("initial registration");
let foreign_owner = SessionRuntimeState {
identity: AgentIdentity::parse("agent:owner-b").expect("identity"),
..first.clone()
};
adapter
.register_session(session.id(), foreign_owner)
.await
.expect_err("a session id cannot be rebound to another identity");
let foreign_generation = SessionRuntimeState {
generation: ContinuityGeneration::new(5),
..first.clone()
};
adapter
.register_session(session.id(), foreign_generation)
.await
.expect_err("a session id cannot be rebound to another generation");
let greater = SessionRuntimeState {
fencing_token: FencingToken::new(11),
..first.clone()
};
adapter
.register_session(session.id(), greater.clone())
.await
.expect("a monotonic fence may replace the prior write authority");
adapter
.suspend_session(session.id())
.await
.expect("suspend replacement authority");
let regressed = SessionRuntimeState {
fencing_token: FencingToken::new(9),
..first.clone()
};
adapter
.register_session(session.id(), regressed)
.await
.expect_err("suspension must never authorize fence regression");
adapter
.register_session(session.id(), greater.clone())
.await
.expect("the same current fence resumes suspended persistence");
assert_eq!(
adapter.lookup_session(&session.id().to_string()),
Some(greater),
"rejected registrations must not replace the owner"
);
}
#[tokio::test]
async fn continuity_session_store_adapter_delete_if_current_revision_removes_matching_snapshot()
{
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let session = meerkat_core::Session::new();
let identity = AgentIdentity::parse("agent:quarantine").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:quarantine:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(4);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity,
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("save snapshot");
let stale_revision = "row-sha256:not-current".to_string();
assert!(
!meerkat::SessionStore::delete_if_current_revision(
&adapter,
session.id(),
&stale_revision
)
.await
.expect("stale delete should be clean"),
"stale revision must not delete"
);
assert!(
store
.load_session_snapshot(session.id())
.await
.expect("load after stale")
.is_some(),
"stale CAS delete must leave snapshot in place"
);
let current_revision =
meerkat_core::session_store::session_projection_cas_token(&session).expect("revision");
assert!(
meerkat::SessionStore::delete_if_current_revision(
&adapter,
session.id(),
¤t_revision
)
.await
.expect("matching delete should succeed"),
"matching revision should delete"
);
assert!(
store
.load_session_snapshot(session.id())
.await
.expect("load after delete")
.is_none(),
"matching CAS delete must remove the continuity snapshot"
);
assert!(
meerkat::SessionStore::load(&adapter, session.id())
.await
.expect("adapter load after delete")
.is_none(),
"adapter must not synthesize a session after successful CAS delete"
);
}
#[tokio::test]
async fn continuity_session_store_adapter_save_rejects_transcript_shrink() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let mut session = meerkat_core::Session::new();
session.append_external_user_content(meerkat_core::ContentInput::Text("first".to_string()));
session
.append_external_user_content(meerkat_core::ContentInput::Text("second".to_string()));
let identity = AgentIdentity::parse("agent:append-only").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:append-only:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(12);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity,
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("initial save");
let mut stale = meerkat_core::Session::with_id(session.id().clone());
stale.append_external_user_content(meerkat_core::ContentInput::Text("first".to_string()));
let err = meerkat::SessionStore::save(&adapter, &stale)
.await
.expect_err("plain save must reject transcript shrink");
assert!(
err.to_string().contains("transcript")
|| err.to_string().contains("monotonicity")
|| err.to_string().contains("continuity"),
"unexpected shrink error: {err}"
);
}
#[tokio::test]
#[allow(deprecated)] async fn raw_save_rejects_stamped_checkpoint_residue_rollback() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let mut authority = meerkat_core::Session::new();
authority
.append_external_user_content(meerkat_core::ContentInput::Text("first".to_string()));
authority
.append_external_user_content(meerkat_core::ContentInput::Text("second".to_string()));
let identity = AgentIdentity::parse("agent:torn-shutdown").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:torn-shutdown:0")
.expect("runtime id"),
session_id: authority.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(21);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
authority.id(),
SessionRuntimeState {
identity,
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
meerkat::SessionStore::save(&adapter, &authority)
.await
.expect("boundary save of the authority");
let mut stamped_head = authority.clone();
stamped_head
.append_external_user_content(meerkat_core::ContentInput::Text("mid-turn".to_string()));
stamped_head
.set_runtime_checkpoint_provenance()
.expect("stamp legacy runtime checkpoint provenance");
meerkat::SessionStore::save(&adapter, &stamped_head)
.await
.expect("checkpointer save of the stamped head");
let error = meerkat::SessionStore::save(&adapter, &authority)
.await
.expect_err("raw save must reject transcript rollback");
assert!(
error.to_string().contains("transcript")
|| error.to_string().contains("monotonicity")
|| error.to_string().contains("continuation"),
"unexpected rollback error: {error}"
);
let loaded = meerkat::SessionStore::load(&adapter, authority.id())
.await
.expect("load")
.expect("session");
assert_eq!(
loaded.messages().len(),
stamped_head.messages().len(),
"rejected raw rollback must leave the durable stamped row unchanged"
);
}
#[tokio::test]
async fn save_keeps_rejecting_unstamped_longer_heads() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let mut authority = meerkat_core::Session::new();
authority
.append_external_user_content(meerkat_core::ContentInput::Text("first".to_string()));
let identity = AgentIdentity::parse("agent:unstamped").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:unstamped:0").expect("runtime id"),
session_id: authority.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(22);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
authority.id(),
SessionRuntimeState {
identity,
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
meerkat::SessionStore::save(&adapter, &authority)
.await
.expect("initial save");
let mut unstamped_head = authority.clone();
unstamped_head
.append_external_user_content(meerkat_core::ContentInput::Text("tail".to_string()));
meerkat::SessionStore::save(&adapter, &unstamped_head)
.await
.expect("longer head appends fine");
meerkat::SessionStore::save(&adapter, &authority)
.await
.expect_err("unstamped longer head must keep failing closed");
}
#[tokio::test]
#[allow(deprecated)] async fn save_keeps_rejecting_stamped_forked_heads() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let mut base = meerkat_core::Session::new();
base.append_external_user_content(meerkat_core::ContentInput::Text("first".to_string()));
let identity = AgentIdentity::parse("agent:forked").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:forked:0").expect("runtime id"),
session_id: base.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(23);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
base.id(),
SessionRuntimeState {
identity,
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
let mut stamped_fork = base.clone();
stamped_fork
.append_external_user_content(meerkat_core::ContentInput::Text("forked".to_string()));
stamped_fork
.set_runtime_checkpoint_provenance()
.expect("stamp legacy runtime checkpoint provenance");
meerkat::SessionStore::save(&adapter, &stamped_fork)
.await
.expect("seed the stamped head");
let mut diverged = base.clone();
diverged
.append_external_user_content(meerkat_core::ContentInput::Text("other".to_string()));
diverged
.append_external_user_content(meerkat_core::ContentInput::Text("branch".to_string()));
meerkat::SessionStore::save(&adapter, &diverged)
.await
.expect_err("stamped but forked head must keep failing closed");
}
#[tokio::test]
async fn continuity_session_store_adapter_saves_transcript_rewrite() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let mut session = meerkat_core::Session::new();
session.append_external_user_content(meerkat_core::ContentInput::Text("first".to_string()));
session
.append_external_user_content(meerkat_core::ContentInput::Text("second".to_string()));
let identity = AgentIdentity::parse("agent:rewrite").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:rewrite:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(13);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity,
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("initial save");
let parent_revision = session.transcript_revision().expect("parent revision");
let mut rewritten = session.clone();
let commit = rewritten
.commit_transcript_rewrite(
meerkat_core::TranscriptRewriteSelection::MessageRange { start: 0, end: 1 },
vec![meerkat_core::Message::User(
meerkat_core::UserMessage::text("compacted first".to_string()),
)],
meerkat_core::TranscriptRewriteReason::new("test"),
Some("mobkit-test".to_string()),
Some(parent_revision),
)
.expect("rewrite commit");
meerkat::SessionStore::save_transcript_rewrite(&adapter, &rewritten, &commit)
.await
.expect("rewrite save should be supported");
let loaded = meerkat::SessionStore::load(&adapter, session.id())
.await
.expect("load rewritten")
.expect("rewritten session");
assert_eq!(loaded.messages().len(), rewritten.messages().len());
assert_eq!(
loaded.transcript_revision().expect("loaded revision"),
commit.revision
);
}
#[tokio::test]
async fn continuity_session_store_adapter_delete_removes_current_snapshot() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let session = meerkat_core::Session::new();
let identity = AgentIdentity::parse("agent:delete").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:delete:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(7);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity,
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
meerkat::SessionStore::save(&adapter, &session)
.await
.expect("save snapshot");
meerkat::SessionStore::delete(&adapter, session.id())
.await
.expect("delete should remove current snapshot");
assert!(
store
.load_session_snapshot(session.id())
.await
.expect("load after delete")
.is_none(),
"delete must not be a successful no-op"
);
assert!(
meerkat::SessionStore::load(&adapter, session.id())
.await
.expect("adapter load after delete")
.is_none(),
"adapter must forget registry state after delete"
);
}
#[tokio::test]
async fn continuity_session_store_adapter_queues_unregistered_authoritative_projection() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let session = meerkat_core::Session::new();
meerkat::SessionStore::save_authoritative_projection(&adapter, &session)
.await
.expect("create-time authoritative projection should queue before registration");
assert!(
meerkat::SessionStore::load(&adapter, session.id())
.await
.expect("load")
.is_none(),
"pending authoritative projection must stay invisible until registration"
);
let identity = AgentIdentity::parse("agent:queued").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:queued:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(7);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity,
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register flushes pending authoritative projection");
assert!(
meerkat::SessionStore::load(&adapter, session.id())
.await
.expect("load after register")
.is_some(),
"registration must flush the pending authoritative projection"
);
}
#[tokio::test]
async fn continuity_session_store_adapter_delete_forgets_registered_session_without_snapshot() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let session = meerkat_core::Session::new();
let identity = AgentIdentity::parse("agent:delete-empty").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:delete-empty:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(11);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity,
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
assert!(
meerkat::SessionStore::load(&adapter, session.id())
.await
.expect("load before delete")
.is_none(),
"registration alone must not fabricate a session document"
);
meerkat::SessionStore::delete(&adapter, session.id())
.await
.expect("delete with no persisted snapshot should be idempotent");
assert!(
meerkat::SessionStore::load(&adapter, session.id())
.await
.expect("load after delete")
.is_none(),
"delete must forget registry state when no persisted row exists"
);
}
#[tokio::test]
async fn continuity_session_store_adapter_projection_cas_ignores_pending_bytes() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let mut session = meerkat_core::Session::new();
session.set_metadata("projection", json!("first-pending"));
meerkat::SessionStore::save_authoritative_projection_if_current_revision(
&adapter, &session, None,
)
.await
.expect("first missing-row CAS");
session.set_metadata("projection", json!("latest-pending"));
meerkat::SessionStore::save_authoritative_projection_if_current_revision(
&adapter, &session, None,
)
.await
.expect("pending bytes are not a visible durable revision");
let identity = AgentIdentity::parse("agent:projection-pending").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:projection-pending:0")
.expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(31);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity,
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("flush latest pending projection");
let loaded = meerkat::SessionStore::load(&adapter, session.id())
.await
.expect("load")
.expect("snapshot");
assert_eq!(
loaded.metadata().get("projection"),
Some(&json!("latest-pending"))
);
}
#[tokio::test]
async fn continuity_session_store_adapter_rejects_snapshot_with_foreign_embedded_id() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let requested = meerkat_core::Session::new();
let foreign = meerkat_core::Session::new();
let identity = AgentIdentity::parse("agent:foreign-snapshot-id").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:foreign-snapshot-id:0")
.expect("runtime id"),
session_id: requested.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(41);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
let bytes = serde_json::to_vec(&foreign).expect("serialize foreign session");
store
.save_session_snapshot(
&identity,
requested.id(),
record.generation,
CheckpointVersion::new(1),
fencing_token,
&SessionSnapshot { data: bytes },
)
.await
.expect("seed corrupt keyed snapshot");
let error = meerkat::SessionStore::load(&adapter, requested.id())
.await
.expect_err("embedded foreign session id must be explicit corruption");
assert!(error.to_string().contains("contains session"));
}
#[tokio::test]
async fn continuity_session_store_adapter_authoritative_projection_cas_guards_rewrites() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = ContinuitySessionStoreAdapter::new(store.clone());
let mut session = meerkat_core::Session::new();
let identity = AgentIdentity::parse("agent:projection").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:projection:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
let fencing_token = FencingToken::new(5);
store
.upsert_continuity_record(&record, fencing_token)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity: identity.clone(),
generation: record.generation,
fencing_token,
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register");
meerkat::SessionStore::save_authoritative_projection_if_current_revision(
&adapter, &session, None,
)
.await
.expect("initial projection should accept missing current revision");
let original_revision =
meerkat_core::session_store::session_projection_cas_token(&session).expect("revision");
let mut stale_rewrite = session.clone();
stale_rewrite.set_metadata("projection", json!("stale"));
let stale_error = meerkat::SessionStore::save_authoritative_projection_if_current_revision(
&adapter,
&stale_rewrite,
Some("row-sha256:not-current".to_string()),
)
.await
.expect_err("stale CAS projection must reject");
assert!(
stale_error.to_string().contains("not a continuation"),
"unexpected stale error: {stale_error}"
);
let loaded = meerkat::SessionStore::load(&adapter, session.id())
.await
.expect("load")
.expect("snapshot");
assert_eq!(
meerkat_core::session_store::session_projection_cas_token(&loaded).expect("revision"),
original_revision,
"stale authoritative projection must leave stored row unchanged"
);
session.set_metadata("projection", json!("current"));
meerkat::SessionStore::save_authoritative_projection_if_current_revision(
&adapter,
&session,
Some(original_revision),
)
.await
.expect("matching CAS projection should save");
let loaded = meerkat::SessionStore::load(&adapter, session.id())
.await
.expect("load after save")
.expect("snapshot after save");
assert_eq!(loaded.metadata().get("projection"), Some(&json!("current")));
}
struct RecordingIncrementalSessions {
inner: Arc<meerkat_store::MemoryStore>,
cursors: Mutex<Vec<super::super::contracts::ContinuityWriteCursor>>,
refuse_writes: AtomicBool,
}
impl RecordingIncrementalSessions {
fn refuse_writes(&self, refuse: bool) {
self.refuse_writes.store(refuse, AtomicOrdering::SeqCst);
}
fn injected_failure(&self) -> Option<meerkat_core::SessionStoreError> {
self.refuse_writes
.load(AtomicOrdering::SeqCst)
.then(|| meerkat_core::SessionStoreError::Internal("injected substrate".into()))
}
fn new(inner: Arc<meerkat_store::MemoryStore>) -> Self {
Self {
inner,
refuse_writes: AtomicBool::new(false),
cursors: Mutex::new(Vec::new()),
}
}
fn record(&self, cursor: &super::super::contracts::ContinuityWriteCursor) {
self.cursors
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(cursor.clone());
}
fn recorded(&self) -> Vec<super::super::contracts::ContinuityWriteCursor> {
self.cursors
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
}
#[async_trait]
impl super::super::contracts::ContinuityIncrementalSessions for RecordingIncrementalSessions {
async fn append_messages(
&self,
cursor: &super::super::contracts::ContinuityWriteCursor,
id: &meerkat_core::types::SessionId,
strand: &meerkat_core::session_store::TranscriptStrandId,
base_seq: u64,
messages: &[meerkat_core::Message],
) -> Result<(), meerkat_core::SessionStoreError> {
self.record(cursor);
if let Some(failure) = self.injected_failure() {
return Err(failure);
}
meerkat_core::session_store::IncrementalSessionStore::append_messages(
self.inner.as_ref(),
id,
strand,
base_seq,
messages,
)
.await
}
async fn commit_rewrite(
&self,
cursor: &super::super::contracts::ContinuityWriteCursor,
id: &meerkat_core::types::SessionId,
record: &meerkat_core::TranscriptRewriteRecord,
expected: meerkat_core::session_store::SessionHeadCas,
) -> Result<meerkat_core::session_store::SessionHead, meerkat_core::SessionStoreError>
{
self.record(cursor);
if let Some(failure) = self.injected_failure() {
return Err(failure);
}
meerkat_core::session_store::IncrementalSessionStore::commit_rewrite(
self.inner.as_ref(),
id,
record,
expected,
)
.await
}
async fn save_head(
&self,
cursor: &super::super::contracts::ContinuityWriteCursor,
head: &meerkat_core::session_store::SessionHead,
expected: meerkat_core::session_store::SessionHeadCas,
) -> Result<(), meerkat_core::SessionStoreError> {
self.record(cursor);
if let Some(failure) = self.injected_failure() {
return Err(failure);
}
meerkat_core::session_store::IncrementalSessionStore::save_head(
self.inner.as_ref(),
head,
expected,
)
.await
}
async fn load_head(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<Option<meerkat_core::session_store::SessionHead>, meerkat_core::SessionStoreError>
{
meerkat_core::session_store::IncrementalSessionStore::load_head(self.inner.as_ref(), id)
.await
}
async fn load_canonical_head(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<Option<meerkat_core::session_store::SessionHead>, meerkat_core::SessionStoreError>
{
meerkat_core::session_store::IncrementalSessionStore::load_head(self.inner.as_ref(), id)
.await
}
async fn load_canonical_session(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<Option<meerkat_core::Session>, meerkat_core::SessionStoreError> {
if meerkat_core::session_store::IncrementalSessionStore::load_head(
self.inner.as_ref(),
id,
)
.await?
.is_none()
{
return Ok(None);
}
meerkat::SessionStore::load(self.inner.as_ref(), id).await
}
async fn load_canonical_previous(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<
Option<(
meerkat_core::Session,
Vec<meerkat_core::TranscriptRewriteCommit>,
)>,
meerkat_core::SessionStoreError,
> {
let Some(session) = self.load_canonical_session(id).await? else {
return Ok(None);
};
let adopted = meerkat_core::session_store::IncrementalSessionStore::load_rewrites(
self.inner.as_ref(),
id,
)
.await?
.into_iter()
.map(|record| record.commit)
.collect();
Ok(Some((session, adopted)))
}
async fn load_messages(
&self,
id: &meerkat_core::types::SessionId,
strand: &meerkat_core::session_store::TranscriptStrandId,
range: std::ops::Range<u64>,
) -> Result<Vec<meerkat_core::Message>, meerkat_core::SessionStoreError> {
meerkat_core::session_store::IncrementalSessionStore::load_messages(
self.inner.as_ref(),
id,
strand,
range,
)
.await
}
async fn load_rewrites(
&self,
id: &meerkat_core::types::SessionId,
) -> Result<Vec<meerkat_core::TranscriptRewriteRecord>, meerkat_core::SessionStoreError>
{
meerkat_core::session_store::IncrementalSessionStore::load_rewrites(
self.inner.as_ref(),
id,
)
.await
}
}
struct IncrementalCapableStore {
inner: Arc<LocalContinuityStore>,
incremental: Arc<RecordingIncrementalSessions>,
}
impl IncrementalCapableStore {
fn new() -> Self {
Self {
inner: Arc::new(LocalContinuityStore::in_memory().expect("store")),
incremental: Arc::new(RecordingIncrementalSessions::new(Arc::new(
meerkat_store::MemoryStore::new(),
))),
}
}
}
#[async_trait]
impl ContinuityStore for IncrementalCapableStore {
async fn resolve_many(
&self,
identities: &[AgentIdentity],
) -> Result<
std::collections::BTreeMap<AgentIdentity, ContinuityResolveState>,
ContinuityStoreError,
> {
self.inner.resolve_many(identities).await
}
async fn load_session_snapshot(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<Option<SessionSnapshot>, ContinuityStoreError> {
self.inner.load_session_snapshot(session_id).await
}
async fn save_session_snapshot(
&self,
identity: &AgentIdentity,
session_id: &meerkat_core::types::SessionId,
generation: ContinuityGeneration,
version: CheckpointVersion,
fencing_token: FencingToken,
snapshot: &SessionSnapshot,
) -> Result<(), ContinuityStoreError> {
self.inner
.save_session_snapshot(
identity,
session_id,
generation,
version,
fencing_token,
snapshot,
)
.await
}
async fn upsert_continuity_record(
&self,
record: &ContinuityRecord,
fencing_token: FencingToken,
) -> Result<(), ContinuityStoreError> {
self.inner
.upsert_continuity_record(record, fencing_token)
.await
}
async fn delete_continuity_record(
&self,
identity: &AgentIdentity,
fencing_token: FencingToken,
) -> Result<(), ContinuityStoreError> {
self.inner
.delete_continuity_record(identity, fencing_token)
.await
}
fn as_incremental_sessions(
&self,
) -> Option<Arc<dyn super::super::contracts::ContinuityIncrementalSessions>> {
Some(self.incremental.clone())
}
}
struct WholeSnapshotOnlyStore {
inner: Arc<LocalContinuityStore>,
}
#[async_trait]
impl ContinuityStore for WholeSnapshotOnlyStore {
async fn resolve_many(
&self,
identities: &[AgentIdentity],
) -> Result<
std::collections::BTreeMap<AgentIdentity, ContinuityResolveState>,
ContinuityStoreError,
> {
self.inner.resolve_many(identities).await
}
async fn load_session_snapshot(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<Option<SessionSnapshot>, ContinuityStoreError> {
self.inner.load_session_snapshot(session_id).await
}
async fn save_session_snapshot(
&self,
identity: &AgentIdentity,
session_id: &meerkat_core::types::SessionId,
generation: ContinuityGeneration,
version: CheckpointVersion,
fencing_token: FencingToken,
snapshot: &SessionSnapshot,
) -> Result<(), ContinuityStoreError> {
self.inner
.save_session_snapshot(
identity,
session_id,
generation,
version,
fencing_token,
snapshot,
)
.await
}
async fn upsert_continuity_record(
&self,
record: &ContinuityRecord,
fencing_token: FencingToken,
) -> Result<(), ContinuityStoreError> {
self.inner
.upsert_continuity_record(record, fencing_token)
.await
}
async fn delete_continuity_record(
&self,
identity: &AgentIdentity,
fencing_token: FencingToken,
) -> Result<(), ContinuityStoreError> {
self.inner
.delete_continuity_record(identity, fencing_token)
.await
}
}
#[tokio::test]
async fn adapter_forwards_incremental_capability_only_when_substrate_advertises() {
let bundled = Arc::new(ContinuitySessionStoreAdapter::new(Arc::new(
LocalContinuityStore::in_memory().expect("store"),
)));
assert!(
meerkat::SessionStore::as_incremental(bundled).is_some(),
"the bundled LocalContinuityStore ships the delta channel (M4b)"
);
let declining = Arc::new(ContinuitySessionStoreAdapter::new(Arc::new(
WholeSnapshotOnlyStore {
inner: Arc::new(LocalContinuityStore::in_memory().expect("store")),
},
)));
assert!(
meerkat::SessionStore::as_incremental(declining).is_none(),
"a whole-snapshot-only substrate must not surface a delta channel"
);
let capable = Arc::new(ContinuitySessionStoreAdapter::new(Arc::new(
IncrementalCapableStore::new(),
)));
assert!(
meerkat::SessionStore::as_incremental(capable).is_some(),
"an advertising substrate must surface through the adapter"
);
}
#[tokio::test]
async fn incremental_mutations_park_before_registration_and_flush_on_register() {
let store = Arc::new(IncrementalCapableStore::new());
let channel = store.incremental.clone();
let adapter = Arc::new(ContinuitySessionStoreAdapter::new(
store.clone() as Arc<dyn ContinuityStore>
));
let incremental = meerkat::SessionStore::as_incremental(adapter.clone())
.expect("advertising substrate forwards");
let session = meerkat_core::Session::new();
let root = meerkat_core::session_store::TranscriptStrandId::root();
let message =
meerkat_core::Message::User(meerkat_core::UserMessage::text("delta turn".to_string()));
let mut document = session.clone();
document.push(message.clone());
let head =
meerkat_core::session_store::SessionHead::from_session(&document, root.clone(), 0)
.expect("head from session");
incremental
.append_messages(session.id(), &root, 0, std::slice::from_ref(&message))
.await
.expect("pre-registration delta writes must park, not fail");
incremental
.save_head(&head, meerkat_core::session_store::SessionHeadCas::Create)
.await
.expect("pre-registration head writes must park, not fail");
assert!(
channel.recorded().is_empty(),
"parked writes must reach nothing durable"
);
let parked_head = incremental
.load_head(session.id())
.await
.expect("parked load_head")
.expect("parked head is visible to the preflight");
assert_eq!(parked_head.head_revision, head.head_revision);
assert_eq!(
incremental
.load_messages(session.id(), &root, 0..1)
.await
.expect("parked rows")
.len(),
1
);
let identity = AgentIdentity::parse("agent:incremental").expect("identity");
let record = ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:incremental:0").expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
};
store
.upsert_continuity_record(&record, FencingToken::new(1))
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity: identity.clone(),
generation: record.generation,
fencing_token: FencingToken::new(1),
checkpoint_version: record.checkpoint_version,
},
)
.await
.expect("register flushes the parked deltas");
let cursors = channel.recorded();
assert_eq!(
cursors.len(),
2,
"the flush replays one append + one head save under the real cursor"
);
assert!(
cursors.iter().all(|cursor| cursor.identity == identity
&& cursor.fencing_token == FencingToken::new(1)),
"every flushed mutation carries the registered continuity cursor"
);
assert!(
cursors[0].checkpoint_version < cursors[1].checkpoint_version,
"each mutation mints a strictly advancing checkpoint version"
);
let durable_rows = incremental
.load_messages(session.id(), &root, 0..1)
.await
.expect("read back after flush");
assert_eq!(
durable_rows.len(),
1,
"the flushed delta row must be durable in the channel"
);
let second = meerkat_core::Message::User(meerkat_core::UserMessage::text(
"second delta turn".to_string(),
));
incremental
.append_messages(session.id(), &root, 1, std::slice::from_ref(&second))
.await
.expect("registered delta writes delegate to the substrate channel");
assert_eq!(channel.recorded().len(), 3);
adapter
.suspend_session(session.id())
.await
.expect("suspend");
let suspended = incremental
.append_messages(session.id(), &root, 2, std::slice::from_ref(&message))
.await
.expect_err("suspended sessions must refuse delta writes");
assert!(
suspended.to_string().contains("suspended"),
"the refusal must name the suspension: {suspended}"
);
}
async fn seed_incremental_record(
store: &IncrementalCapableStore,
identity: &AgentIdentity,
session_id: &meerkat_core::types::SessionId,
) -> SessionRuntimeState {
store
.upsert_continuity_record(
&ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse(&format!("rt:{identity}:0"))
.expect("runtime id"),
session_id: session_id.clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
},
FencingToken::new(1),
)
.await
.expect("seed record");
SessionRuntimeState {
identity: identity.clone(),
generation: ContinuityGeneration::new(0),
fencing_token: FencingToken::new(1),
checkpoint_version: CheckpointVersion::new(0),
}
}
#[tokio::test]
async fn registration_never_purges_parked_rows_that_no_head_adopts_yet() {
let store = Arc::new(IncrementalCapableStore::new());
let channel = store.incremental.clone();
let adapter = Arc::new(ContinuitySessionStoreAdapter::new(
store.clone() as Arc<dyn ContinuityStore>
));
let incremental = meerkat::SessionStore::as_incremental(adapter.clone())
.expect("advertising substrate forwards");
let session = meerkat_core::Session::new();
let root = meerkat_core::session_store::TranscriptStrandId::root();
let message = meerkat_core::Message::User(meerkat_core::UserMessage::text(
"creation-window turn".to_string(),
));
let mut document = session.clone();
document.push(message.clone());
let head =
meerkat_core::session_store::SessionHead::from_session(&document, root.clone(), 0)
.expect("head from session");
incremental
.append_messages(session.id(), &root, 0, std::slice::from_ref(&message))
.await
.expect("pre-registration append must park");
let identity = AgentIdentity::parse("agent:racing").expect("identity");
let state = seed_incremental_record(&store, &identity, session.id()).await;
let refused = adapter
.register_session(session.id(), state.clone())
.await
.expect_err("an unadoptable parked state must refuse the registration");
assert!(
refused.to_string().contains("no parked head adopts"),
"the refusal must name the unadoptable parked state: {refused}"
);
assert!(
adapter.parked_deltas.is_parked(session.id()),
"the routing marker must survive a refused registration"
);
assert_eq!(
adapter
.parked_deltas
.footprint(session.id())
.expect("footprint")
.rows,
1,
"the parked footprint must still account for the parked row"
);
assert!(
adapter.lookup_session(&session.id().to_string()).is_none(),
"a refused registration must restore the registry entry"
);
assert!(
channel.recorded().is_empty(),
"an unadoptable parked state must reach nothing durable"
);
incremental
.save_head(&head, meerkat_core::session_store::SessionHeadCas::Create)
.await
.expect("the retained rows must still be there for the head to adopt");
assert_eq!(
incremental
.load_messages(session.id(), &root, 0..1)
.await
.expect("parked rows")
.len(),
1,
"the row parked before the refused registration must be the row the head adopts"
);
let committed = adapter
.register_session(session.id(), state)
.await
.expect("the retry must flush the now-adoptable parked document");
assert!(committed.get() > 0, "the flush must commit a version");
assert!(
!adapter.parked_deltas.is_parked(session.id()),
"an adopted parked state is purged"
);
assert_eq!(
incremental
.load_messages(session.id(), &root, 0..1)
.await
.expect("durable rows")
.len(),
1,
"the retained row must reach the durable channel on the retry"
);
}
#[tokio::test]
async fn failed_parked_flush_restores_registry_and_markers_and_retries_clean() {
let store = Arc::new(IncrementalCapableStore::new());
let channel = store.incremental.clone();
let adapter = Arc::new(ContinuitySessionStoreAdapter::new(
store.clone() as Arc<dyn ContinuityStore>
));
let incremental = meerkat::SessionStore::as_incremental(adapter.clone())
.expect("advertising substrate forwards");
let session = meerkat_core::Session::new();
let root = meerkat_core::session_store::TranscriptStrandId::root();
let message = meerkat_core::Message::User(meerkat_core::UserMessage::text(
"flush failure turn".to_string(),
));
let mut document = session.clone();
document.push(message.clone());
let head =
meerkat_core::session_store::SessionHead::from_session(&document, root.clone(), 0)
.expect("head from session");
incremental
.append_messages(session.id(), &root, 0, std::slice::from_ref(&message))
.await
.expect("pre-registration append must park");
incremental
.save_head(&head, meerkat_core::session_store::SessionHeadCas::Create)
.await
.expect("pre-registration head must park");
adapter
.suspend_session(session.id())
.await
.expect("suspend");
let identity = AgentIdentity::parse("agent:flushfail").expect("identity");
let state = seed_incremental_record(&store, &identity, session.id()).await;
channel.refuse_writes(true);
let refused = adapter
.register_session(session.id(), state.clone())
.await
.expect_err("a failing parked flush must refuse the registration");
assert!(
refused.to_string().contains("injected substrate"),
"the substrate failure must ride out typed: {refused}"
);
assert!(
adapter.lookup_session(&session.id().to_string()).is_none(),
"a failed flush must restore the registry entry"
);
let still_suspended = adapter
.ensure_session_mutation_allowed(session.id())
.expect_err("the suspension marker must be restored");
assert!(
still_suspended.to_string().contains("suspended"),
"the restored marker must be the suspension: {still_suspended}"
);
assert!(
adapter.parked_deltas.is_parked(session.id()),
"a failed flush must retain the routing marker"
);
assert_eq!(
incremental
.load_messages(session.id(), &root, 0..1)
.await
.expect("parked rows")
.len(),
1,
"a failed flush must retain the parked rows for the retry"
);
assert!(
meerkat_core::session_store::IncrementalSessionStore::load_head(
channel.inner.as_ref(),
session.id()
)
.await
.expect("durable head")
.is_none(),
"a failed flush must leave the durable channel empty"
);
channel.refuse_writes(false);
adapter
.register_session(session.id(), state)
.await
.expect("the retry must flush the retained parked document");
assert!(
!adapter.parked_deltas.is_parked(session.id()),
"an adopted parked state is purged"
);
assert_eq!(
incremental
.load_messages(session.id(), &root, 0..1)
.await
.expect("durable rows")
.len(),
1,
"the retained row must reach the durable channel on the retry"
);
assert_eq!(
meerkat_core::session_store::IncrementalSessionStore::load_head(
channel.inner.as_ref(),
session.id()
)
.await
.expect("durable head")
.expect("durable head after retry")
.head_revision,
head.head_revision,
"the retry must adopt the parked head"
);
}
#[tokio::test]
async fn transcript_rewrite_on_a_head_canonical_session_records_the_commit() {
let store = Arc::new(LocalContinuityStore::in_memory().expect("store"));
let adapter = Arc::new(ContinuitySessionStoreAdapter::new(
store.clone() as Arc<dyn ContinuityStore>
));
let incremental = meerkat::SessionStore::as_incremental(adapter.clone())
.expect("the bundled store advertises the delta channel");
let mut session = meerkat_core::Session::new();
session.push(meerkat_core::Message::User(
meerkat_core::UserMessage::text("one".to_string()),
));
session.push(meerkat_core::Message::User(
meerkat_core::UserMessage::text("two".to_string()),
));
let identity = AgentIdentity::parse("agent:rewrite").expect("identity");
store
.upsert_continuity_record(
&ContinuityRecord {
identity: identity.clone(),
agent_runtime_id: AgentRuntimeId::parse("rt:agent:rewrite:0")
.expect("runtime id"),
session_id: session.id().clone(),
generation: ContinuityGeneration::new(0),
checkpoint_version: CheckpointVersion::new(0),
},
FencingToken::new(1),
)
.await
.expect("seed record");
adapter
.register_session(
session.id(),
SessionRuntimeState {
identity,
generation: ContinuityGeneration::new(0),
fencing_token: FencingToken::new(1),
checkpoint_version: CheckpointVersion::new(0),
},
)
.await
.expect("register");
let root = meerkat_core::session_store::TranscriptStrandId::root();
incremental
.append_messages(session.id(), &root, 0, session.messages())
.await
.expect("seed rows");
let head =
meerkat_core::session_store::SessionHead::from_session(&session, root.clone(), 0)
.expect("head");
incremental
.save_head(&head, meerkat_core::session_store::SessionHeadCas::Create)
.await
.expect("adopt head");
let mut rewritten = session.clone();
let commit = rewritten
.commit_transcript_rewrite(
meerkat_core::TranscriptRewriteSelection::MessageRange { start: 0, end: 2 },
vec![meerkat_core::Message::User(
meerkat_core::UserMessage::text("[summary]".to_string()),
)],
meerkat_core::TranscriptRewriteReason::new("compaction"),
Some("adapter-test".to_string()),
None,
)
.expect("commit rewrite");
meerkat::SessionStore::save_transcript_rewrite(adapter.as_ref(), &rewritten, &commit)
.await
.expect("head-canonical rewrite must be admitted");
let records = incremental
.load_rewrites(session.id())
.await
.expect("load rewrites");
assert_eq!(
records.len(),
1,
"the rewrite must be recorded and adopted, not flattened into a rebase"
);
assert_eq!(records[0].commit.revision, commit.revision);
let loaded = meerkat::SessionStore::load(adapter.as_ref(), session.id())
.await
.expect("load")
.expect("session");
assert_eq!(
loaded.messages(),
rewritten.messages(),
"the adopted head must serve the rewritten transcript"
);
}
#[tokio::test]
async fn superseded_sessions_drop_terminal_writes_without_parking_them() {
let store = Arc::new(IncrementalCapableStore::new());
let adapter = Arc::new(ContinuitySessionStoreAdapter::new(
store.clone() as Arc<dyn ContinuityStore>
));
let incremental = meerkat::SessionStore::as_incremental(adapter.clone())
.expect("advertising substrate forwards");
let session = meerkat_core::Session::new();
adapter
.superseded_sessions
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(session.id().to_string());
let root = meerkat_core::session_store::TranscriptStrandId::root();
let message =
meerkat_core::Message::User(meerkat_core::UserMessage::text("orphan".to_string()));
incremental
.append_messages(session.id(), &root, 0, std::slice::from_ref(&message))
.await
.expect("a post-abandon terminal write must be acknowledged, not refused");
assert!(
meerkat_core::session_store::IncrementalSessionStore::load_head(
adapter.parked_deltas.reads(),
session.id(),
)
.await
.expect("parked head probe")
.is_none(),
"a superseded terminal write must not park anything"
);
assert!(
meerkat_core::session_store::IncrementalSessionStore::load_head(
incremental.as_ref(),
session.id(),
)
.await
.expect("substrate head probe")
.is_none(),
"a superseded terminal write must not reach the substrate"
);
}
#[tokio::test]
async fn registration_over_an_empty_parked_marker_clears_it_and_keeps_the_session_usable() {
let store = Arc::new(IncrementalCapableStore::new());
let channel = store.incremental.clone();
let adapter = Arc::new(ContinuitySessionStoreAdapter::new(
store.clone() as Arc<dyn ContinuityStore>
));
let incremental = meerkat::SessionStore::as_incremental(adapter.clone())
.expect("advertising substrate forwards");
let session = meerkat_core::Session::new();
let root = meerkat_core::session_store::TranscriptStrandId::root();
incremental
.append_messages(session.id(), &root, 0, &[])
.await
.expect("an empty pre-registration append must park, not fail");
assert!(
adapter.parked_deltas.is_parked(session.id()),
"an empty append still publishes the routing marker"
);
assert_eq!(
adapter
.parked_deltas
.footprint(session.id())
.expect("footprint")
.rows,
0,
"an empty append parks zero rows"
);
assert!(
meerkat_core::session_store::IncrementalSessionStore::load_head(
adapter.parked_deltas.reads(),
session.id(),
)
.await
.expect("parked head probe")
.is_none(),
"the Empty arm is only reachable with no parked head"
);
let identity = AgentIdentity::parse("agent:empty-marker").expect("identity");
let state = seed_incremental_record(&store, &identity, session.id()).await;
adapter
.register_session(session.id(), state)
.await
.expect("an empty parked marker must not refuse the registration");
assert!(
!adapter.parked_deltas.is_parked(session.id()),
"the empty marker must be cleared, or every later read for this session \
is served from a parked view that will never hold anything"
);
assert!(
channel.recorded().is_empty(),
"an empty parked marker replays nothing durable"
);
let message = meerkat_core::Message::User(meerkat_core::UserMessage::text(
"first real turn".to_string(),
));
let mut document = session.clone();
document.push(message.clone());
let head =
meerkat_core::session_store::SessionHead::from_session(&document, root.clone(), 0)
.expect("head from session");
incremental
.append_messages(session.id(), &root, 0, std::slice::from_ref(&message))
.await
.expect("registered append");
incremental
.save_head(&head, meerkat_core::session_store::SessionHeadCas::Create)
.await
.expect("registered head save");
assert_eq!(
channel.recorded().len(),
2,
"post-registration writes must reach the durable channel"
);
}
#[tokio::test]
async fn deleting_a_parked_session_reclaims_its_parked_rows() {
let store = Arc::new(IncrementalCapableStore::new());
let adapter = Arc::new(ContinuitySessionStoreAdapter::new(
store.clone() as Arc<dyn ContinuityStore>
));
let incremental = meerkat::SessionStore::as_incremental(adapter.clone())
.expect("advertising substrate forwards");
let session = meerkat_core::Session::new();
let root = meerkat_core::session_store::TranscriptStrandId::root();
let message = meerkat_core::Message::User(meerkat_core::UserMessage::text(
"parked transcript".to_string(),
));
let mut document = session.clone();
document.push(message.clone());
let head =
meerkat_core::session_store::SessionHead::from_session(&document, root.clone(), 0)
.expect("head from session");
incremental
.append_messages(session.id(), &root, 0, std::slice::from_ref(&message))
.await
.expect("pre-registration append parks");
incremental
.save_head(&head, meerkat_core::session_store::SessionHeadCas::Create)
.await
.expect("pre-registration head parks");
assert!(
meerkat_core::session_store::IncrementalSessionStore::load_head(
adapter.parked_deltas.reads(),
session.id(),
)
.await
.expect("parked head probe")
.is_some(),
"the parked document must exist before the delete"
);
meerkat::SessionStore::delete(adapter.as_ref(), session.id())
.await
.expect("deleting a session with no durable document succeeds");
assert!(
!adapter.parked_deltas.is_parked(session.id()),
"delete must clear the routing marker"
);
assert!(
adapter.parked_deltas.footprint(session.id()).is_none(),
"delete must clear the footprint"
);
assert!(
meerkat_core::session_store::IncrementalSessionStore::load_head(
adapter.parked_deltas.reads(),
session.id(),
)
.await
.expect("parked head probe")
.is_none(),
"delete must reclaim the parked ROWS, not merely stop routing to them — \
an unmarked-but-retained parked document is unreachable memory held \
for the lifetime of the process"
);
}
#[tokio::test]
async fn head_canonical_load_takes_one_substrate_snapshot() {
struct TearingSubstrate {
settled: meerkat_core::Session,
stale_head: meerkat_core::session_store::SessionHead,
}
#[async_trait]
impl super::super::contracts::ContinuityIncrementalSessions for TearingSubstrate {
async fn append_messages(
&self,
_cursor: &super::super::contracts::ContinuityWriteCursor,
_id: &meerkat_core::types::SessionId,
_strand: &meerkat_core::session_store::TranscriptStrandId,
_base_seq: u64,
_messages: &[meerkat_core::Message],
) -> Result<(), meerkat_core::SessionStoreError> {
Ok(())
}
async fn commit_rewrite(
&self,
_cursor: &super::super::contracts::ContinuityWriteCursor,
_id: &meerkat_core::types::SessionId,
_record: &meerkat_core::TranscriptRewriteRecord,
_expected: meerkat_core::session_store::SessionHeadCas,
) -> Result<meerkat_core::session_store::SessionHead, meerkat_core::SessionStoreError>
{
Ok(self.stale_head.clone())
}
async fn save_head(
&self,
_cursor: &super::super::contracts::ContinuityWriteCursor,
_head: &meerkat_core::session_store::SessionHead,
_expected: meerkat_core::session_store::SessionHeadCas,
) -> Result<(), meerkat_core::SessionStoreError> {
Ok(())
}
async fn load_head(
&self,
_id: &meerkat_core::types::SessionId,
) -> Result<
Option<meerkat_core::session_store::SessionHead>,
meerkat_core::SessionStoreError,
> {
Ok(Some(self.stale_head.clone()))
}
async fn load_canonical_head(
&self,
_id: &meerkat_core::types::SessionId,
) -> Result<
Option<meerkat_core::session_store::SessionHead>,
meerkat_core::SessionStoreError,
> {
Ok(Some(self.stale_head.clone()))
}
async fn load_canonical_session(
&self,
_id: &meerkat_core::types::SessionId,
) -> Result<Option<meerkat_core::Session>, meerkat_core::SessionStoreError>
{
Ok(Some(self.settled.clone()))
}
async fn load_canonical_previous(
&self,
_id: &meerkat_core::types::SessionId,
) -> Result<
Option<(
meerkat_core::Session,
Vec<meerkat_core::TranscriptRewriteCommit>,
)>,
meerkat_core::SessionStoreError,
> {
Ok(Some((self.settled.clone(), Vec::new())))
}
async fn load_messages(
&self,
_id: &meerkat_core::types::SessionId,
_strand: &meerkat_core::session_store::TranscriptStrandId,
_range: std::ops::Range<u64>,
) -> Result<Vec<meerkat_core::Message>, meerkat_core::SessionStoreError> {
Ok(self.settled.messages().to_vec())
}
async fn load_rewrites(
&self,
_id: &meerkat_core::types::SessionId,
) -> Result<Vec<meerkat_core::TranscriptRewriteRecord>, meerkat_core::SessionStoreError>
{
Ok(Vec::new())
}
}
struct TearingStore {
inner: Arc<LocalContinuityStore>,
channel: Arc<TearingSubstrate>,
}
#[async_trait]
impl ContinuityStore for TearingStore {
async fn resolve_many(
&self,
identities: &[AgentIdentity],
) -> Result<
std::collections::BTreeMap<AgentIdentity, ContinuityResolveState>,
ContinuityStoreError,
> {
self.inner.resolve_many(identities).await
}
async fn load_session_snapshot(
&self,
session_id: &meerkat_core::types::SessionId,
) -> Result<Option<SessionSnapshot>, ContinuityStoreError> {
self.inner.load_session_snapshot(session_id).await
}
async fn save_session_snapshot(
&self,
identity: &AgentIdentity,
session_id: &meerkat_core::types::SessionId,
generation: ContinuityGeneration,
checkpoint_version: CheckpointVersion,
fencing_token: FencingToken,
snapshot: &SessionSnapshot,
) -> Result<(), ContinuityStoreError> {
self.inner
.save_session_snapshot(
identity,
session_id,
generation,
checkpoint_version,
fencing_token,
snapshot,
)
.await
}
async fn upsert_continuity_record(
&self,
record: &ContinuityRecord,
fencing_token: FencingToken,
) -> Result<(), ContinuityStoreError> {
self.inner
.upsert_continuity_record(record, fencing_token)
.await
}
async fn delete_continuity_record(
&self,
identity: &AgentIdentity,
fencing_token: FencingToken,
) -> Result<(), ContinuityStoreError> {
self.inner
.delete_continuity_record(identity, fencing_token)
.await
}
fn as_incremental_sessions(
&self,
) -> Option<Arc<dyn super::super::contracts::ContinuityIncrementalSessions>>
{
Some(self.channel.clone())
}
}
let root = meerkat_core::session_store::TranscriptStrandId::root();
let mut settled = meerkat_core::Session::new();
settled.push(meerkat_core::Message::User(
meerkat_core::UserMessage::text("settled turn".to_string()),
));
let mut earlier = settled.clone();
earlier.push(meerkat_core::Message::User(
meerkat_core::UserMessage::text("since rewritten away".to_string()),
));
let stale_head =
meerkat_core::session_store::SessionHead::from_session(&earlier, root.clone(), 0)
.expect("head from session");
assert_ne!(
stale_head.message_count,
settled.messages().len() as u64,
"the double must actually model a torn pair"
);
let adapter = Arc::new(ContinuitySessionStoreAdapter::new(Arc::new(TearingStore {
inner: Arc::new(LocalContinuityStore::in_memory().expect("store")),
channel: Arc::new(TearingSubstrate {
settled: settled.clone(),
stale_head,
}),
})));
let loaded = meerkat::SessionStore::load(adapter.as_ref(), settled.id())
.await
.expect(
"the head-first load must take ONE substrate snapshot; composing an \
independent head read with an independent rows read surfaces the torn pair",
)
.expect("the session is head-canonical");
assert_eq!(
loaded.messages().len(),
settled.messages().len(),
"the loaded document must be the substrate's single-snapshot materialization"
);
}
}