use super::*;
use mj_client::runtime_feed::{
RuntimeCursor, RuntimeDelta, RuntimeFrame, RuntimeMetadata, RuntimeProjection, SessionTailReply,
};
use mj_core::native_agent::NativeAgentSummary;
use mj_core::snapshot_map::SnapshotMap;
type NativeOwners = SnapshotMap<String, SnapshotMap<String, NativeAgentSummary>>;
type Children = SnapshotMap<String, SnapshotMap<String, ()>>;
#[derive(Default)]
pub(super) struct RuntimeHistory {
incarnation: String,
sequence: u64,
snapshots: VecDeque<(u64, RuntimeProjection, usize)>,
bytes: usize,
native_owners: NativeOwners,
full: RuntimeProjection,
operations: BTreeSet<String>,
}
impl RuntimeHistory {
fn cursor(&self) -> RuntimeCursor {
RuntimeCursor {
incarnation: self.incarnation.clone(),
sequence: self.sequence,
}
}
fn publish(&mut self, next: RuntimeProjection) -> Result<()> {
if self.incarnation.is_empty() {
self.incarnation = new_command_id("runtime-feed")?;
}
let size = if let Some((_, previous, _)) = self.snapshots.back() {
let delta = RuntimeDelta::between(previous, &next);
if delta.is_empty() {
return Ok(());
}
serde_json::to_vec(&delta)?.len()
} else {
0
};
self.sequence += 1;
self.snapshots.push_back((self.sequence, next, size));
self.bytes += size;
while self.snapshots.len() > 1
&& (self.snapshots.len() > 4096 || self.bytes > 16 * 1024 * 1024)
{
let (_, _, size) = self.snapshots.pop_front().expect("history entry");
self.bytes -= size;
}
if self.bytes > 16 * 1024 * 1024 {
self.snapshots.back_mut().expect("current projection").2 = 0;
self.bytes = 0;
}
Ok(())
}
fn live_projection(
&self,
full: &RuntimeProjection,
native_owners: &NativeOwners,
children: &Children,
operations: &BTreeSet<String>,
) -> RuntimeProjection {
let before = &self.full;
let mut live = self
.snapshots
.back()
.map(|(_, projection, _)| projection.clone())
.unwrap_or_default();
live.revision = full.revision;
live.sessions = full.sessions.clone();
live.transcripts = full.transcripts.clone();
live.metadata = full.metadata.clone();
let mut pending = before
.records
.changes(&full.records)
.map(|(id, _)| id.clone())
.chain(
before
.subagents
.changes(&full.subagents)
.map(|(id, _)| id.clone()),
)
.chain(before.moves.changes(&full.moves).map(|(id, _)| id.clone()))
.chain(self.operations.symmetric_difference(operations).cloned())
.collect::<Vec<_>>();
while let Some(id) = pending.pop() {
let was = live.records.contains_key(&id);
let record = full.records.get(&id);
let now = record.is_some_and(|record| {
let parent_live = full
.subagents
.get(&id)
.is_some_and(|relation| live.records.contains_key(&relation.parent_session_id));
mj_core::state::session_is_live(record, operations.contains(&id), parent_live)
});
sync_key(&mut live.records, &id, record.filter(|_| now));
sync_key(
&mut live.subagents,
&id,
full.subagents.get(&id).filter(|_| now),
);
sync_key(&mut live.moves, &id, full.moves.get(&id).filter(|_| now));
sync_key(
&mut live.session_cpu,
&id,
full.session_cpu.get(&id).filter(|_| now),
);
if was == now {
continue;
}
for summary in native_owners
.get(&id)
.into_iter()
.flat_map(SnapshotMap::values)
{
let view_id = summary.agent.view_id();
sync_key(
&mut live.native_agents,
&view_id,
Some(summary).filter(|_| now),
);
}
pending.extend(
children
.get(&id)
.into_iter()
.flat_map(SnapshotMap::keys)
.cloned(),
);
}
for (view_id, summary) in before.native_agents.changes(&full.native_agents) {
let owner = summary
.or_else(|| before.native_agents.get(view_id))
.map(|summary| summary.agent.owner_session_id.as_str());
if owner.is_some_and(|owner| live.records.contains_key(owner)) {
sync_key(&mut live.native_agents, view_id, summary);
}
}
for (id, value) in before.session_cpu.changes(&full.session_cpu) {
sync_key(
&mut live.session_cpu,
id,
value.filter(|_| live.records.contains_key(id)),
);
}
debug_assert_eq!(
live.records.keys().cloned().collect::<BTreeSet<_>>(),
mj_core::state::live_session_ids(&full.records, &full.subagents, operations),
"the incremental live set must equal a full evaluation"
);
live
}
fn session_tail(&self, cursor: &RuntimeCursor, session_id: &str) -> SessionTailReply {
if cursor.incarnation != self.incarnation {
return SessionTailReply::ResetRequired;
}
let Some((_, projection, _)) = self
.snapshots
.iter()
.find(|(sequence, _, _)| *sequence == cursor.sequence)
else {
return SessionTailReply::ResetRequired;
};
projection
.transcripts
.get(session_id)
.map_or(SessionTailReply::NoTail, SessionTailReply::of)
}
fn frame(&self, requested: Option<&RuntimeCursor>) -> RuntimeFrame {
let (_, current, _) = self.snapshots.back().expect("captured projection");
let Some(requested) = requested else {
return RuntimeFrame::Snapshot {
cursor: self.cursor(),
projection: Box::new(current.clone()),
};
};
if requested.incarnation == self.incarnation
&& let Some((_, before, _)) = self
.snapshots
.iter()
.find(|(sequence, _, _)| *sequence == requested.sequence)
{
return RuntimeFrame::Delta {
from: requested.clone(),
cursor: self.cursor(),
changes: Box::new(RuntimeDelta::between(before, current)),
};
}
RuntimeFrame::ResetRequired
}
}
fn launch_inputs_changed(before: &RuntimeProjection, after: &RuntimeProjection) -> bool {
let inputs = |record: &SessionRecord| {
(
record.bundle_id.clone(),
record.last_profile.clone(),
record.target_template_id.clone(),
record.created_at.clone(),
)
};
before
.records
.changes(&after.records)
.any(|(id, record)| before.records.get(id).map(inputs) != record.map(inputs))
}
fn sync_key<V: Clone + PartialEq>(map: &mut SnapshotMap<String, V>, id: &str, value: Option<&V>) {
match value {
Some(value) if map.get(id) != Some(value) => {
map.insert(id.to_owned(), value.clone());
}
Some(_) => {}
None => {
map.remove(id);
}
}
}
impl RuntimeState {
pub(crate) fn runtime_publication(&self) -> Result<RuntimeProjection> {
self.capture_runtime()?;
Ok(self
.feed
.lock()
.unwrap_or_else(PoisonError::into_inner)
.full
.clone())
}
fn capture_runtime(&self) -> Result<RuntimeCursor> {
let mut history = self.feed.lock().unwrap_or_else(PoisonError::into_inner);
let (mut full, native_owners, children, operations) = {
let owner = self.owner();
owner.ensure_available()?;
let controller = owner.controller();
let moves = owner
.committed()
.map(|state| state.moves.clone())
.unwrap_or_default();
let lifecycles = Self::active_lifecycles_with(&owner);
let operations = lifecycles
.iter()
.map(|lifecycle| lifecycle.session_id.clone())
.chain(
moves
.iter()
.filter(|(_, operation)| operation.is_active())
.map(|(id, _)| id.clone()),
)
.collect::<BTreeSet<_>>();
(
RuntimeProjection {
revision: self.revisions.current(),
records: owner.projected_records(),
subagents: controller.state.subagents.clone(),
sessions: owner.sessions.clone(),
transcripts: owner.transcripts.clone(),
session_cpu: self.session_manager.session_cpu.borrow().clone(),
moves,
metadata: RuntimeMetadata {
config: controller.config.clone(),
last_subagent_policy: controller.state.last_subagent_policy.clone(),
lifecycles,
..Default::default()
},
..Default::default()
},
owner
.committed()
.map(|state| state.native_agents.clone())
.unwrap_or_default(),
owner.indexes.children.clone(),
operations,
)
};
full.native_agents = history.full.native_agents.clone();
for (owner, children) in history.native_owners.changes(&native_owners) {
let empty = SnapshotMap::new();
let before = history.native_owners.get(owner).unwrap_or(&empty);
for (child, value) in before.changes(children.unwrap_or(&empty)) {
if let Some(old) = before.get(child) {
full.native_agents.remove(&old.agent.view_id());
}
if let Some(value) = value {
full.native_agents
.insert(value.agent.view_id(), value.clone());
}
}
}
full.metadata.workspace_names = self
.workspaces()
.borrow()
.iter()
.map(|w| (w.id.clone(), w.name.clone()))
.collect();
full.metadata.reviews = self.review_host.views();
full.metadata.notices = self
.notices
.lock()
.unwrap_or_else(PoisonError::into_inner)
.iter()
.cloned()
.collect();
full.metadata.quotas = self
.quota
.lock()
.unwrap_or_else(PoisonError::into_inner)
.snapshot
.clone();
full.metadata.profile_capabilities = self
.profile_catalog
.get()
.map(|catalog| catalog.snapshot())
.unwrap_or_default();
full.metadata.launch_recency = if launch_inputs_changed(&history.full, &full) {
mj_client::runtime_feed::launch_recency(&full.records)
} else {
history.full.metadata.launch_recency.clone()
};
let live = history.live_projection(&full, &native_owners, &children, &operations);
history.publish(live)?;
history.native_owners = native_owners;
history.full = full;
history.operations = operations;
Ok(history.cursor())
}
pub(super) fn session_tail(
&self,
session_id: &str,
cursor: &RuntimeCursor,
) -> SessionTailReply {
self.feed
.lock()
.unwrap_or_else(PoisonError::into_inner)
.session_tail(cursor, session_id)
}
pub(super) async fn runtime_changes(
&self,
cursor: Option<RuntimeCursor>,
wait: bool,
) -> Result<RuntimeFrame> {
let mut revisions = self.revisions.subscribe();
revisions.borrow_and_update();
let current = self.capture_runtime()?;
if wait && cursor.as_ref() == Some(¤t) {
let _ = tokio::time::timeout(Duration::from_secs(30), revisions.changed()).await;
self.capture_runtime()?;
}
Ok(self
.feed
.lock()
.unwrap_or_else(PoisonError::into_inner)
.frame(cursor.as_ref()))
}
}
#[cfg(test)]
mod tests {
use super::*;
use mj_client::runtime_feed::RuntimeReplica;
#[test]
fn retained_cursors_replay_and_slow_or_replaced_clients_reset() {
let mut history = RuntimeHistory::default();
history.publish(RuntimeProjection::default()).unwrap();
let mut replica = RuntimeReplica::default();
replica.apply(history.frame(None)).unwrap();
let original = replica.cursor.clone().unwrap();
let mut next = replica.projection.clone();
next.metadata
.workspace_names
.insert("workspace".into(), "Renamed".into());
history.publish(next.clone()).unwrap();
replica.apply(history.frame(Some(&original))).unwrap();
assert_eq!(replica.projection, next);
let wrong = RuntimeCursor {
incarnation: "previous-daemon".into(),
sequence: 1,
};
assert!(matches!(
history.frame(Some(&wrong)),
RuntimeFrame::ResetRequired
));
for revision in 1..=4096 {
next.revision = revision;
history.publish(next.clone()).unwrap();
}
assert_eq!(history.snapshots.len(), 4096);
assert!(matches!(
history.frame(Some(&original)),
RuntimeFrame::ResetRequired
));
assert!(matches!(history.frame(None), RuntimeFrame::Snapshot { .. }));
}
#[test]
fn a_tail_is_served_at_the_requested_cursor() {
let item = |position: u64, text: &str| {
Arc::new(mj_core::transcript::TranscriptItem {
stable_id: format!("agent:{position}"),
position,
latest_content_event_ordinal: Some(position),
created_at_ms: 0,
last_changed_at_ms: 0,
body: mj_core::transcript::TranscriptBody::Agent {
chunks: vec![serde_json::json!({"content": {"type": "text", "text": text}})],
streaming: true,
},
})
};
let tail_of = |items: Vec<Arc<mj_core::transcript::TranscriptItem>>| {
let mut materialized = mj_core::state::MaterializedSession::empty("s");
materialized.applied_event_ordinal = items.len() as u64;
materialized.transcript = items;
mj_client::runtime_feed::SessionTail::of(
&materialized,
&mj_core::state::ProjectionWindow::default(),
1_024,
)
};
let mut history = RuntimeHistory::default();
let first = tail_of(vec![item(1, "one")]);
let mut projection = RuntimeProjection::default();
projection.transcripts.insert("s".into(), first.clone());
history.publish(projection.clone()).unwrap();
let held = history.cursor();
projection
.transcripts
.insert("s".into(), tail_of(vec![item(1, "one, then more")]));
history.publish(projection).unwrap();
assert_eq!(
history.session_tail(&held, "s"),
SessionTailReply::of(&first),
"the tail at the held cursor, not the newer one"
);
assert_eq!(
history.session_tail(&held, "other"),
SessionTailReply::NoTail
);
let replaced = RuntimeCursor {
incarnation: "previous-daemon".into(),
sequence: held.sequence,
};
assert_eq!(
history.session_tail(&replaced, "s"),
SessionTailReply::ResetRequired
);
let pruned = RuntimeCursor {
sequence: 0,
..held
};
assert_eq!(
history.session_tail(&pruned, "s"),
SessionTailReply::ResetRequired
);
}
#[test]
fn an_oversized_batch_drops_history_but_preserves_the_current_snapshot() {
let mut history = RuntimeHistory::default();
history.publish(RuntimeProjection::default()).unwrap();
let original = history.cursor();
let mut next = RuntimeProjection::default();
next.metadata.notices.push(RuntimeNotice {
id: 1,
session_id: String::new(),
text: "x".repeat(16 * 1024 * 1024),
});
history.publish(next).unwrap();
assert_eq!(history.snapshots.len(), 1);
assert_eq!(history.bytes, 0);
assert!(matches!(
history.frame(Some(&original)),
RuntimeFrame::ResetRequired
));
}
#[tokio::test]
async fn draft_hydration_publishes_shared_capabilities_without_adopting_the_draft() {
use crate::server_runtime::profile_catalog::{ProfileCatalog, counting_probe, test_config};
use mj_core::config::HarnessKind;
use mj_core::profile_capabilities::CapabilityState;
use std::sync::atomic::{AtomicUsize, Ordering};
let state = super::super::tests::test_runtime_state();
let calls = Arc::new(AtomicUsize::new(0));
let catalog = ProfileCatalog::with_probe(counting_probe(calls.clone()));
let weak = Arc::downgrade(&state);
catalog.set_publisher(Arc::new(move || {
if let Some(state) = weak.upgrade() {
state.publish_revision();
}
}));
assert!(state.profile_catalog.set(catalog.clone()).is_ok());
let config = test_config(&[("draft", HarnessKind::Codex)], &[]);
let key = config.profiles["draft"].capabilities_key("draft");
let mut replica = RuntimeReplica::default();
replica
.apply(state.runtime_changes(None, false).await.unwrap())
.unwrap();
let cursor = replica.cursor.clone();
let metadata = super::super::tests::test_metadata("127.0.0.1:1".parse().unwrap());
let cancellation = CancellationToken::new();
for _ in 0..2 {
let reply = super::super::actions::handle_action(
DaemonAction::WarmProfileCapabilities {
config: Box::new(config.clone()),
},
&metadata,
&state,
&cancellation,
)
.await
.unwrap();
assert!(matches!(reply, DaemonReply::Done));
catalog.options_for(&config, "draft", None).await.unwrap();
}
let frame =
tokio::time::timeout(Duration::from_secs(1), state.runtime_changes(cursor, true))
.await
.unwrap()
.unwrap();
replica.apply(frame).unwrap();
assert!(matches!(
replica.projection.metadata.profile_capabilities.profiles[&key].choices,
CapabilityState::Ready(_)
));
assert_eq!(calls.load(Ordering::SeqCst), 1);
assert_eq!(replica.projection.metadata.config, Config::default());
}
#[tokio::test]
async fn attaching_and_then_changing_owner_state_yields_an_incremental_frame() {
let state = super::super::tests::test_runtime_state();
let initial = state.runtime_changes(None, false).await.unwrap();
let mut replica = RuntimeReplica::default();
replica.apply(initial).unwrap();
let cursor = replica.cursor.clone();
state.owner().edit_sessions(|sessions| {
sessions.insert(
"new".into(),
super::super::tests::runtime_test_session(
"new",
"workspace",
SessionState::Running,
),
);
});
state.publish_revision();
let frame = state.runtime_changes(cursor, true).await.unwrap();
let RuntimeFrame::Delta { changes, .. } = &frame else {
panic!("expected delta")
};
assert_eq!(changes.records.len(), 1);
replica.apply(frame).unwrap();
assert!(replica.projection.records.contains_key("new"));
}
#[test]
fn cpu_measurements_follow_live_membership_and_replay_by_key() {
let mut full = RuntimeProjection::default();
let mut history = RuntimeHistory::default();
let value = mj_client::runtime_feed::SessionCpuView::Measured {
usage: mj_core::cpu_usage::SessionCpuUsage {
recent_permille: 230,
hourly_permille: 100,
hourly_covered_secs: 30,
online_cpus: 8,
},
};
for (id, state) in [
("live", SessionState::Running),
("stopped", SessionState::Stopped),
] {
full.records.insert(id.into(), record(id, state));
full.session_cpu.insert(id.into(), value.clone());
}
capture(&mut history, &full, &NativeOwners::new(), &[]);
let mut replica = RuntimeReplica::default();
replica.apply(history.frame(None)).unwrap();
assert_eq!(replica.projection.session_cpu.len(), 1);
assert_eq!(replica.projection.session_cpu.get("live"), Some(&value));
let cursor = replica.cursor.clone();
full.records.get_mut("live").unwrap().state = SessionState::Stopped;
capture(&mut history, &full, &NativeOwners::new(), &[]);
replica.apply(history.frame(cursor.as_ref())).unwrap();
assert!(replica.projection.session_cpu.is_empty());
let cursor = replica.cursor.clone();
capture(&mut history, &full, &NativeOwners::new(), &["live"]);
replica.apply(history.frame(cursor.as_ref())).unwrap();
assert_eq!(replica.projection.session_cpu.get("live"), Some(&value));
}
fn record(id: &str, state: SessionState) -> SessionRecord {
super::super::tests::runtime_test_session(id, "workspace", state)
}
fn capture(
history: &mut RuntimeHistory,
full: &RuntimeProjection,
native_owners: &NativeOwners,
operations: &[&str],
) -> Vec<String> {
let mut children = Children::new();
for (child, relation) in &full.subagents {
let mut members = children
.get(&relation.parent_session_id)
.cloned()
.unwrap_or_default();
members.insert(child.clone(), ());
children.insert(relation.parent_session_id.clone(), members);
}
let operations = operations.iter().map(|id| (*id).to_owned()).collect();
let live = history.live_projection(full, native_owners, &children, &operations);
history.publish(live).unwrap();
history.full = full.clone();
history.native_owners = native_owners.clone();
history.operations = operations;
history
.snapshots
.back()
.unwrap()
.1
.records
.keys()
.cloned()
.collect()
}
#[test]
fn the_terminal_feed_follows_live_sessions_and_drops_stopped_ones() {
let mut full = RuntimeProjection::default();
for (id, state) in [
("parent", SessionState::Running),
("stopped", SessionState::Stopped),
("lost", SessionState::Lost),
("child", SessionState::Stopped),
] {
full.records.insert(id.into(), record(id, state));
}
full.subagents.insert(
"child".into(),
super::super::tests::runtime_test_subagent("child", "parent"),
);
let agent = mj_core::native_agent::NativeAgent {
owner_session_id: "stopped".into(),
session_id: "explore".into(),
parent_session_id: None,
name: "Explore".into(),
task: "Map the code".into(),
capabilities: Default::default(),
state: mj_core::native_agent::NativeAgentState::Completed,
availability: Default::default(),
availability_reason: None,
stable_id: None,
};
let view_id = agent.view_id();
let summary = NativeAgentSummary {
generation_ordinal: 1,
agent,
projection_ordinal: 0,
projection_digest: String::new(),
};
full.native_agents.insert(view_id.clone(), summary.clone());
let mut native_owners = NativeOwners::new();
native_owners.insert(
"stopped".into(),
SnapshotMap::from([(view_id.clone(), summary)]),
);
let mut history = RuntimeHistory::default();
assert_eq!(
capture(&mut history, &full, &native_owners, &[]),
["child", "lost", "parent"]
);
let first = history.cursor();
let current = &history.snapshots.back().unwrap().1;
assert!(current.native_agents.is_empty());
assert!(current.subagents.contains_key("child"));
assert_eq!(
capture(&mut history, &full, &native_owners, &["stopped"]),
["child", "lost", "parent", "stopped"]
);
assert!(
history
.snapshots
.back()
.unwrap()
.1
.native_agents
.contains_key(&view_id)
);
full.records
.insert("parent".into(), record("parent", SessionState::Stopped));
assert_eq!(
capture(&mut history, &full, &native_owners, &["stopped"]),
["lost", "stopped"]
);
assert_eq!(capture(&mut history, &full, &native_owners, &[]), ["lost"]);
let current = &history.snapshots.back().unwrap().1;
assert!(current.native_agents.is_empty());
assert!(current.subagents.is_empty());
let RuntimeFrame::Delta { changes, .. } = history.frame(Some(&first)) else {
panic!("expected a delta from the first capture");
};
let removed = changes
.records
.iter()
.filter(|(_, record)| record.is_none())
.map(|(id, _)| id.as_str())
.collect::<BTreeSet<_>>();
assert_eq!(removed, BTreeSet::from(["child", "parent"]));
assert_eq!(
history.full.records.len(),
4,
"the web server keeps every record"
);
}
}