use super::*;
pub fn spawn_remote_dashboard_worker_poller(
workspace_id: String,
) -> Result<RemoteDashboardWorkerPoller> {
let channels = spawn_remote_session_manager()?;
let crate::session_manager::RemoteSessionManagerChannels {
targets,
control,
updates,
shutdown,
publisher,
mut requests,
} = channels;
let (state_tx, state_rx) = tokio::sync::watch::channel(RuntimeStateUpdate::default());
let (reviews_tx, reviews_rx) = tokio::sync::watch::channel(Vec::new());
let (notices_tx, notices_rx) = tokio::sync::watch::channel(Vec::new());
let (quotas_tx, quotas_rx) = tokio::sync::watch::channel(Default::default());
let (capabilities_tx, capabilities_rx) = tokio::sync::watch::channel(Default::default());
let (config_tx, config_rx) = tokio::sync::watch::channel(mj_core::config::Config::default());
let (health_tx, health_rx) = tokio::sync::watch::channel(RuntimeFeedHealth::default());
tokio::spawn(async move {
let replica = Arc::new(tokio::sync::Mutex::new(
mj_client::runtime_feed::RuntimeReplica::default(),
));
let tails = Arc::clone(&replica);
let mut feed = spawn_runtime_feed_with(
workspace_id,
move |_, revision| {
let replica = replica.clone();
async move { poll_daemon_runtime(replica, revision).await }
},
move |session_id| load_session_tail(Arc::clone(&tails), session_id),
);
let mut native = super::native_agents::NativeAgentLoader::default();
let mut request_order = crate::session_manager::SessionRequestOrder::new();
loop {
tokio::select! {
_ = state_tx.closed() => return,
() = native.next(), if native.has_work() => {
let views = native.views_snapshot();
state_tx.send_if_modified(|state| {
if state.native_agents == views { return false; }
state.native_agents = views;
true
});
health_tx.send_if_modified(|health| {
let error = native.error();
if health.native_error == error { false } else { health.native_error = error; true }
});
},
request = requests.recv() => {
let Some(request) = request else { return; };
request_order.dispatch(request, forward_remote_session_request);
}
update = feed.updates.recv() => {
match update {
Some(RuntimeFeedUpdate::Snapshot(snapshot)) => {
let metadata = snapshot.metadata;
send_if_changed(&config_tx, metadata.installed_config());
native.update_snapshot(snapshot.native_agents);
health_tx.send_if_modified(|health| {
let recovered = health.refresh_error.take().is_some();
let native_error = native.error();
let changed = health.native_error != native_error;
health.native_error = native_error;
recovered || changed
});
publish_runtime_state(&state_tx, RuntimeStateUpdate {
last_subagent_policy: metadata.last_subagent_policy,
native_agents: native.views_snapshot(),
session_cpu: snapshot.session_cpu,
workspace_names: metadata.workspace_names,
revision: snapshot.revision,
records: snapshot.records,
lifecycles: metadata.lifecycles,
moves: snapshot.moves,
subagents: snapshot.subagents,
launch_recency: metadata.launch_recency,
storage: metadata.storage,
});
send_if_changed(&reviews_tx, metadata.reviews);
send_if_changed(&capabilities_tx, metadata.profile_capabilities);
send_if_changed(¬ices_tx, metadata.notices);
send_if_changed("as_tx, metadata.quotas);
}
Some(RuntimeFeedUpdate::Session { session_id, view }) => {
if mirror_daemon_view(&targets, &publisher, session_id, *view).await.is_err() { return; }
}
Some(RuntimeFeedUpdate::SessionRemoved(session_id)) => {
mirror_daemon_removal(&targets, &session_id);
}
Some(RuntimeFeedUpdate::Error(error)) => {
health_tx.send_if_modified(|health| {
if health.refresh_error.as_ref() == Some(&error) { return false; }
tracing::warn!(%error, "could not refresh sessions from controller daemon");
health.refresh_error = Some(error);
true
});
}
None => {
health_tx.send_modify(|health| health.refresh_error = Some("Session updates stopped; reconnect to the controller.".into()));
return;
},
}
}
}
}
});
Ok(RemoteDashboardWorkerPoller {
updates,
control,
shutdown,
state: state_rx,
reviews: reviews_rx,
notices: notices_rx,
quotas: quotas_rx,
profile_capabilities: capabilities_rx,
config: config_rx,
health: health_rx,
})
}
pub(super) fn send_if_changed<T: PartialEq>(tx: &tokio::sync::watch::Sender<T>, value: T) -> bool {
tx.send_if_modified(|current| {
if *current == value {
return false;
}
*current = value;
true
})
}
pub(super) fn publish_runtime_state(
tx: &tokio::sync::watch::Sender<RuntimeStateUpdate>,
next: RuntimeStateUpdate,
) -> bool {
tx.send_if_modified(|current| {
current.revision = next.revision;
if *current == next {
return false;
}
*current = next;
true
})
}
pub(super) async fn mirror_daemon_view(
targets: &tokio::sync::watch::Sender<Vec<WorkerPollTarget>>,
publisher: &crate::session_manager::RemoteSessionPublisher,
session_id: String,
view: ManagedSessionView,
) -> Result<()> {
targets.send_if_modified(|targets| {
if targets.iter().any(|target| target.session_id == session_id) {
return false;
}
targets.push(WorkerPollTarget {
session_id: session_id.clone(),
spec: crate::targets::CommandSpec::new("daemon-owned", Vec::<String>::new()),
worker_recovery: None,
project_memory: None,
});
true
});
publisher.publish(session_id, view).await
}
pub(super) fn mirror_daemon_removal(
targets: &tokio::sync::watch::Sender<Vec<WorkerPollTarget>>,
session_id: &str,
) {
targets.send_if_modified(|targets| {
let before = targets.len();
targets.retain(|target| target.session_id != session_id);
targets.len() != before
});
}
pub(super) async fn load_session_tail(
replica: Arc<tokio::sync::Mutex<mj_client::runtime_feed::RuntimeReplica>>,
session_id: String,
) -> Result<StoredProjection> {
let cursor = {
let replica = replica.lock().await;
if let Some(tail) = replica.projection.transcripts.get(&session_id) {
return Ok(Some((tail.materialized(), tail.window.clone())));
}
replica
.cursor
.clone()
.context("the runtime feed has no cursor to fetch a tail at")?
};
let reply = mj_client::daemon::connect_existing()
.await?
.session_tail(session_id.clone(), cursor.clone())
.await?;
let mut replica = replica.lock().await;
match reply {
mj_client::runtime_feed::SessionTailReply::Tail {
header,
window,
items,
} => {
let tail = mj_client::runtime_feed::SessionTail::from_parts(*header, window, items);
let projection = (tail.materialized(), tail.window.clone());
if replica.cursor.as_ref() == Some(&cursor) {
replica.projection.transcripts.insert(session_id, tail);
}
Ok(Some(projection))
}
mj_client::runtime_feed::SessionTailReply::NoTail => Ok(None),
mj_client::runtime_feed::SessionTailReply::ResetRequired => {
replica.cursor = None;
Err(anyhow::Error::new(super::runtime_feed::TailCursorExpired))
}
}
}
pub(super) async fn poll_daemon_runtime(
replica: Arc<tokio::sync::Mutex<mj_client::runtime_feed::RuntimeReplica>>,
after_revision: u64,
) -> Result<mj_client::runtime_feed::RuntimeProjection> {
let mut replica = replica.lock().await;
loop {
let mut daemon = mj_client::daemon::connect_existing().await?;
let wait = after_revision != 0 && after_revision == replica.projection.revision;
let frame = daemon.runtime_changes(replica.cursor.clone(), wait).await?;
if let Err(error) = replica.apply(frame) {
replica.cursor = None;
return Err(error);
}
if replica.cursor.is_none() {
continue;
}
return Ok(replica.projection.clone());
}
}
#[cfg(test)]
mod expired_tail_cursor_tests {
use super::*;
#[tokio::test]
async fn expired_session_tail_cursor_clears_the_replica_cursor_before_retry() {
const CHILD: &str = "MJ_TEST_EXPIRED_SESSION_TAIL_CURSOR";
const TEST: &str = "expired_session_tail_cursor_clears_the_replica_cursor_before_retry";
if std::env::var_os(CHILD).is_none() {
let root = tempfile::tempdir().expect("isolated data root");
crate::controller::test_support::IsolatedTest::new(
crate::controller::test_support::test_name(module_path!(), TEST),
)
.env(CHILD, "1")
.env("MJ_INSTANCE", "expired-tail-cursor")
.isolated_store(root.path())
.run();
return;
}
use mj_client::daemon::{
DaemonAction, DaemonMetadata, DaemonReply, RequestEnvelope, ResponseEnvelope,
};
use mj_client::runtime_feed::{RuntimeCursor, RuntimeReplica, SessionTailReply};
use tokio::net::TcpListener;
let listener = TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0))
.await
.expect("bind fake daemon protocol endpoint");
let address = listener.local_addr().unwrap();
let metadata = DaemonMetadata {
protocol_version: mj_client::daemon::PROTOCOL_VERSION,
pid: std::process::id(),
address,
token: "tail-cursor-test-token".into(),
started_at: "2026-10-05T00:00:00Z".into(),
build_version: "test".into(),
};
let metadata_path = mj_client::daemon::metadata_path();
std::fs::create_dir_all(metadata_path.parent().unwrap()).expect("create isolated data dir");
std::fs::write(
&metadata_path,
serde_json::to_vec(&metadata).expect("serialize daemon metadata"),
)
.expect("publish fake daemon metadata");
let server = tokio::spawn(async move {
let (mut stream, _) = listener.accept().await.expect("accept tail request");
let request: RequestEnvelope = mj_client::daemon::read_frame(&mut stream)
.await
.expect("read tail request");
let DaemonAction::SessionTail { session_id, cursor } = request.action else {
panic!("expected SessionTail request");
};
let response = ResponseEnvelope {
protocol_version: request.protocol_version,
request_id: request.request_id,
result: Ok(DaemonReply::SessionTail(Box::new(
SessionTailReply::ResetRequired,
))),
};
mj_client::daemon::write_frame(&mut stream, &response)
.await
.expect("write expired cursor response");
(session_id, cursor)
});
let cursor = RuntimeCursor {
incarnation: "previous-daemon".into(),
sequence: 42,
};
let replica = Arc::new(tokio::sync::Mutex::new(RuntimeReplica {
cursor: Some(cursor.clone()),
..RuntimeReplica::default()
}));
let error = load_session_tail(Arc::clone(&replica), "session-7".into())
.await
.expect_err("an expired cursor requires a fresh runtime snapshot");
assert_eq!(
error.to_string(),
"the daemon no longer holds this runtime cursor"
);
assert!(replica.lock().await.cursor.is_none());
assert_eq!(
server.await.expect("fake daemon task"),
("session-7".to_owned(), cursor)
);
}
}
fn remote_submit_failure(error: &anyhow::Error) -> mj_client::session::SubmitFailure {
let unconfirmed = error
.downcast_ref::<mj_client::daemon::DaemonRefusal>()
.is_none_or(mj_client::daemon::DaemonRefusal::delivery_unconfirmed);
mj_client::session::SubmitFailure {
unconfirmed,
refused: false,
message: format!("{error:#}"),
}
}
pub(super) async fn forward_remote_session_request(request: RemoteSessionRequest) {
match request {
RemoteSessionRequest::Submit {
session_id,
command_id,
command,
admission,
reply,
} => {
if admission.is_some() {
let _ = reply.send(Err(
"review delivery admissions cannot cross the daemon request bridge".into(),
));
return;
}
let result = async {
mj_client::daemon::connect_existing()
.await?
.submit_session_command(session_id, command_id, command, None)
.await
}
.await
.map_err(|error| remote_submit_failure(&error));
let _ = reply.send(result);
}
RemoteSessionRequest::Sync { session_id, reply } => {
let result = async {
mj_client::daemon::connect_existing()
.await?
.sync_session(session_id)
.await
}
.await
.map_err(|error| format!("{error:#}"));
let _ = reply.send(result);
}
RemoteSessionRequest::RespondElicitation {
session_id,
elicitation_id,
response,
reply,
} => {
let result = async {
mj_client::daemon::connect_existing()
.await?
.respond_elicitation(session_id, elicitation_id, response)
.await
}
.await
.map_err(|error| format!("{error:#}"));
let _ = reply.send(result);
}
RemoteSessionRequest::StopBackgroundTask {
session_id,
background_task_id,
reply,
} => {
let result = async {
mj_client::daemon::connect_existing()
.await?
.stop_background_task(session_id, background_task_id)
.await
}
.await
.map_err(|error| format!("{error:#}"));
let _ = reply.send(result);
}
RemoteSessionRequest::Reviewer {
session_id,
role,
action,
mut reply,
} => {
let result = tokio::select! {
_ = reply.closed() => return,
result = async {
mj_client::daemon::connect_existing()
.await?
.reviewer_action(session_id, role, action)
.await
} => result,
}
.map_err(|error| format!("{error:#}"));
let _ = reply.send(result);
}
}
}
pub fn queued_prompt_projection(
session: &MaterializedSession,
) -> Vec<mj_core::relay::QueuedPrompt> {
queued_prompt_entries(&session.queued_prompts)
}
pub(super) fn queued_prompt_entries(
prompts: &[mj_core::state::MaterializedQueuedPrompt],
) -> Vec<mj_core::relay::QueuedPrompt> {
prompts
.iter()
.map(|prompt| mj_core::relay::QueuedPrompt {
id: prompt.command_id.clone(),
text: mj_core::transcript::materialized_content_text(&prompt.content),
attachments: Vec::new(),
created_at_ms: prompt.queued_at_ms,
})
.collect()
}
#[cfg(test)]
mod tests {
use super::{publish_runtime_state, remote_submit_failure, send_if_changed};
use crate::pollers::{Feed, RuntimeStateUpdate};
use mj_client::daemon::{DaemonRefusal, RuntimeNotice};
#[test]
fn a_republished_snapshot_wakes_the_surface_only_for_what_changed() {
let (state_tx, state_rx) = tokio::sync::watch::channel(RuntimeStateUpdate::default());
let (notices_tx, notices_rx) = tokio::sync::watch::channel(Vec::<RuntimeNotice>::new());
let mut state = Feed::new(state_rx);
let mut notices = Feed::new(notices_rx);
let first = RuntimeStateUpdate {
revision: 1,
workspace_names: [("w".to_owned(), "work".to_owned())].into(),
..Default::default()
};
let notice = RuntimeNotice {
id: 1,
session_id: "s".into(),
text: "checkpoint saved".into(),
};
publish_runtime_state(&state_tx, first.clone());
send_if_changed(¬ices_tx, vec![notice.clone()]);
assert_eq!(state.next_ready().map(|state| state.revision), Some(1));
assert_eq!(notices.next_ready(), Some(vec![notice.clone()]));
publish_runtime_state(
&state_tx,
RuntimeStateUpdate {
revision: 2,
..first.clone()
},
);
send_if_changed(¬ices_tx, vec![notice.clone()]);
assert!(
state.next_ready().is_none(),
"an equal snapshot woke the surface"
);
assert!(
notices.next_ready().is_none(),
"equal notices woke the surface"
);
assert_eq!(
state_tx.borrow().revision,
2,
"the revision is still current"
);
let mut renamed = first;
renamed.revision = 3;
renamed
.workspace_names
.insert("w".to_owned(), "renamed".to_owned());
publish_runtime_state(&state_tx, renamed.clone());
assert_eq!(state.next_ready(), Some(renamed));
}
#[test]
fn a_daemon_refusal_is_a_rejection_and_a_lost_exchange_is_unconfirmed() {
let refused = remote_submit_failure(&anyhow::Error::new(DaemonRefusal(
"/clear requires an idle session".into(),
)));
assert!(!refused.unconfirmed);
assert_eq!(refused.message, "/clear requires an idle session");
let unconfirmed = remote_submit_failure(&anyhow::Error::new(DaemonRefusal(
"delivery unconfirmed: channel closed".into(),
)));
assert!(unconfirmed.unconfirmed);
let lost = remote_submit_failure(&anyhow::anyhow!("daemon connection reset"));
assert!(lost.unconfirmed);
}
}