use std::net::{Ipv4Addr, SocketAddr};
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::{Duration, Instant};
use anyhow::{Context, Result, bail};
mod api;
mod api_activity;
use mj_core::config::{Config, HarnessProfile, PhoneConfig, is_bare_project_target};
use mj_core::remote_git::{default_branch, display_url, resolve_repository};
use mj_core::state::{MaterializedSession, ProjectSourceIdentity, SessionRecord, State};
use crate::controller::Controller;
use crate::quota::ProfileQuota;
use crate::server::{
ActionOutcome, BackgroundTaskStopFailure, BackgroundTaskStopRequest, BrowserTranscript,
ControllerAction, ControllerRequest, MovePreparationRequest, PreflightFailure,
ReadReceiptRequest, ResumeQueueDisposition, ServerOptions, ViewerActivityDetails,
ViewerActivityKind, ViewerBackgroundTask, ViewerMoveRecovery, ViewerQueuedPrompt, ViewerQuota,
ViewerSnapshot, ViewerUserShell,
};
use crate::session_manager::{SessionManagerChannels, SessionManagerControl, new_command_id};
use crate::tailscale::TailscaleTls;
#[cfg(test)]
use crate::targets::ProcessExecutor;
use crate::targets::{CancellableProcessExecutor, CommandExecutor};
use crate::worker_client::CredentialSyncCoordinator;
use mj_core::relay::RelayCommand;
use mj_core::workspace::WorkspaceRecord;
use crate::controller::config_only_controller;
use crate::daemon::{
CreateSessionControl, CreateSessionRequest, ResumeSessionRequest, RuntimeState,
};
use crate::pollers::{
CredentialSyncNotices, CredentialSyncSignalTracker, QUOTA_STALE_AFTER, QuotaRefreshBatch,
QuotaUpdate, apply_worker_record_update, credential_sync_targets, dashboard_worker_targets,
projected_queued_prompts, queued_prompt_projection, quota_refresh_profiles,
schedule_due_credential_syncs, spawn_quota_refresher,
};
#[derive(Debug, Clone)]
pub struct ServerArgs {
bind: String,
tailscale_detect: bool,
tls_cert: Option<PathBuf>,
tls_key: Option<PathBuf>,
}
impl From<&PhoneConfig> for ServerArgs {
fn from(config: &PhoneConfig) -> Self {
Self {
bind: config.bind.clone(),
tailscale_detect: config.tailscale_detect,
tls_cert: config.tls_cert.clone(),
tls_key: config.tls_key.clone(),
}
}
}
const TAILSCALE_COMMAND_TIMEOUT: Duration = Duration::from_secs(120);
const TAILSCALE_RENEW_INTERVAL: Duration = Duration::from_secs(24 * 60 * 60);
const MAX_CONCURRENT_CONVERSATION_PROJECTIONS: usize = 2;
const CONVERSATION_PROJECTION_CHANNEL_CAPACITY: usize = 16;
#[derive(Debug, Clone, PartialEq, Eq)]
struct ConversationProjectionKey {
ordinal: u64,
digest: String,
}
impl ConversationProjectionKey {
fn of(materialized: &MaterializedSession) -> Self {
Self {
ordinal: materialized.applied_event_ordinal,
digest: materialized.applied_event_digest.clone(),
}
}
fn is_newer_than(&self, other: &Self) -> bool {
self.ordinal > other.ordinal
|| (self.ordinal == other.ordinal && self.digest != other.digest)
}
}
struct ConversationProjectionRequest {
materialized: MaterializedSession,
key: ConversationProjectionKey,
generation: u64,
}
struct ConversationProjectionResult {
session_id: String,
key: ConversationProjectionKey,
generation: u64,
result: std::result::Result<BrowserTranscript, String>,
}
struct ConversationProjectionDispatcher {
in_flight: std::collections::BTreeMap<String, (ConversationProjectionKey, u64)>,
pending: std::collections::BTreeMap<String, ConversationProjectionRequest>,
completed: std::collections::BTreeMap<String, ConversationProjectionKey>,
generations: std::collections::BTreeMap<String, u64>,
permits: Arc<tokio::sync::Semaphore>,
results: tokio::sync::mpsc::Sender<ConversationProjectionResult>,
shutdown: tokio_util::sync::CancellationToken,
}
impl ConversationProjectionDispatcher {
fn new(
results: tokio::sync::mpsc::Sender<ConversationProjectionResult>,
shutdown: tokio_util::sync::CancellationToken,
) -> Self {
Self {
in_flight: std::collections::BTreeMap::new(),
pending: std::collections::BTreeMap::new(),
completed: std::collections::BTreeMap::new(),
generations: std::collections::BTreeMap::new(),
permits: Arc::new(tokio::sync::Semaphore::new(
MAX_CONCURRENT_CONVERSATION_PROJECTIONS,
)),
results,
shutdown,
}
}
#[cfg(test)]
fn with_permits(
results: tokio::sync::mpsc::Sender<ConversationProjectionResult>,
shutdown: tokio_util::sync::CancellationToken,
permits: usize,
) -> Self {
let mut dispatcher = Self::new(results, shutdown);
dispatcher.permits = Arc::new(tokio::sync::Semaphore::new(permits));
dispatcher
}
fn enqueue(&mut self, materialized: MaterializedSession) {
let session_id = materialized.session_id.clone();
let key = ConversationProjectionKey::of(&materialized);
if self
.completed
.get(&session_id)
.is_some_and(|completed| !key.is_newer_than(completed))
{
return;
}
let generation = *self.generations.entry(session_id.clone()).or_default();
if let Some((in_flight, in_flight_generation)) = self.in_flight.get(&session_id) {
if generation != *in_flight_generation || key.is_newer_than(in_flight) {
let replace = self.pending.get(&session_id).is_none_or(|pending| {
pending.generation != generation || key.is_newer_than(&pending.key)
});
if replace {
self.pending.insert(
session_id,
ConversationProjectionRequest {
materialized,
key,
generation,
},
);
}
}
return;
}
self.in_flight
.insert(session_id.clone(), (key.clone(), generation));
self.start(ConversationProjectionRequest {
materialized,
key,
generation,
});
}
fn finish(
&mut self,
result: ConversationProjectionResult,
session_active: bool,
) -> Option<(String, ConversationProjectionKey, BrowserTranscript)> {
let expected = self.in_flight.remove(&result.session_id);
if expected.as_ref() != Some(&(result.key.clone(), result.generation)) {
tracing::warn!(
session_id = %result.session_id,
"discarding an out-of-date browser transcript projection"
);
return None;
}
let session_id = result.session_id;
let key = result.key;
let current_generation = self
.generations
.get(&session_id)
.copied()
.unwrap_or_default();
let current = result.generation == current_generation;
let projected = match result.result {
Ok(transcript) if session_active && current => {
self.completed.insert(session_id.clone(), key.clone());
Some((session_id.clone(), key, transcript))
}
Ok(_) => {
self.completed.remove(&session_id);
None
}
Err(error) => {
tracing::warn!(
session_id = %session_id,
"browser transcript projection failed: {error}"
);
None
}
};
if session_active {
if let Some(pending) = self.pending.remove(&session_id) {
self.enqueue(pending.materialized);
}
} else {
self.pending.remove(&session_id);
self.completed.remove(&session_id);
}
projected
}
fn forget(&mut self, session_id: &str) {
self.pending.remove(session_id);
self.completed.remove(session_id);
let generation = self.generations.entry(session_id.to_owned()).or_default();
*generation = generation.wrapping_add(1);
}
fn session_ids(&self) -> std::collections::BTreeSet<String> {
self.in_flight
.keys()
.chain(self.pending.keys())
.chain(self.completed.keys())
.cloned()
.collect()
}
fn start(&self, request: ConversationProjectionRequest) {
let session_id = request.materialized.session_id.clone();
let key = request.key;
let generation = request.generation;
let permits = Arc::clone(&self.permits);
let results = self.results.clone();
let shutdown = self.shutdown.clone();
tokio::spawn(async move {
let result = match tokio::select! {
_ = shutdown.cancelled() => return,
result = permits.acquire_owned() => result,
} {
Ok(permit) => {
let projection = tokio::task::spawn_blocking(move || {
mj_client::transcript::materialized_browser_transcript(
&request.materialized,
)
})
.await;
drop(permit);
projection
.map_err(|error| format!("transcript projection task failed: {error}"))
}
Err(error) => Err(format!("transcript projection worker stopped: {error}")),
};
let message = ConversationProjectionResult {
session_id,
key,
generation,
result,
};
tokio::select! {
_ = shutdown.cancelled() => {}
result = results.send(message) => {
if let Err(error) = result {
tracing::debug!(%error, "browser transcript projection result dropped after server shutdown");
}
}
}
});
}
}
struct ResolvedServerArgs {
bind: SocketAddr,
viewer_url: String,
tls_files: Option<(PathBuf, PathBuf)>,
tailscale: Option<TailscaleTls>,
fallback_reason: Option<String>,
}
async fn resolve_server_args(
args: ServerArgs,
termination: tokio_util::sync::CancellationToken,
) -> Result<ResolvedServerArgs> {
let configured_bind: SocketAddr = args.bind.parse().context("parse web viewer bind address")?;
match (args.tls_cert, args.tls_key) {
(Some(cert), Some(key)) => {
let scheme = "https";
return Ok(ResolvedServerArgs {
bind: configured_bind,
viewer_url: format!("{scheme}://{configured_bind}"),
tls_files: Some((cert, key)),
tailscale: None,
fallback_reason: None,
});
}
(None, None) => {}
_ => bail!("web viewer TLS requires both a certificate and private key"),
}
if !args.tailscale_detect {
return Ok(loopback_server_args(
configured_bind,
Some("automatic Tailscale detection is disabled".into()),
));
}
let tls_root = mj_core::config::data_dir().join("viewer");
let prepared = run_tailscale_blocking(termination.clone(), move |executor| {
crate::tailscale::prepare_tailscale_tls(&tls_root, executor)
})
.await;
match prepared {
Ok(tailscale) => {
let bind = tailscale_bind(configured_bind);
let viewer_url = format!(
"https://{}:{}",
tailscale.cert_domain(),
configured_bind.port()
);
Ok(ResolvedServerArgs {
bind,
viewer_url,
tls_files: Some((
tailscale.cert_path().to_owned(),
tailscale.key_path().to_owned(),
)),
tailscale: Some(tailscale),
fallback_reason: None,
})
}
Err(error) if termination.is_cancelled() => Err(error),
Err(error) => {
let reason = format!("{error:#}");
tracing::debug!(error = reason, "Tailscale HTTPS unavailable for web viewer");
Ok(loopback_server_args(configured_bind, Some(reason)))
}
}
}
fn tailscale_bind(configured_bind: SocketAddr) -> SocketAddr {
SocketAddr::from((Ipv4Addr::UNSPECIFIED, configured_bind.port()))
}
fn loopback_server_args(bind: SocketAddr, fallback_reason: Option<String>) -> ResolvedServerArgs {
ResolvedServerArgs {
bind,
viewer_url: format!("http://{bind}"),
tls_files: None,
tailscale: None,
fallback_reason,
}
}
async fn run_tailscale_blocking<T>(
termination: tokio_util::sync::CancellationToken,
operation: impl FnOnce(&CancellableProcessExecutor) -> Result<T> + Send + 'static,
) -> Result<T>
where
T: Send + 'static,
{
let cancelled = Arc::new(AtomicBool::new(false));
let executor_cancelled = cancelled.clone();
let mut task = tokio::task::spawn_blocking(move || {
let executor = CancellableProcessExecutor::new(executor_cancelled)
.with_deadline(TAILSCALE_COMMAND_TIMEOUT);
operation(&executor)
});
tokio::select! {
result = &mut task => result.context("Tailscale background task panicked")?,
_ = termination.cancelled() => {
cancelled.store(true, Ordering::Release);
let _ = task.await;
bail!("Tailscale operation cancelled during web viewer shutdown")
}
}
}
fn spawn_tailscale_cert_renewer(
tailscale: TailscaleTls,
rustls: axum_server::tls_rustls::RustlsConfig,
termination: tokio_util::sync::CancellationToken,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut interval = tokio::time::interval(TAILSCALE_RENEW_INTERVAL);
interval.tick().await;
loop {
tokio::select! {
_ = termination.cancelled() => return,
_ = interval.tick() => {}
}
let renewing = tailscale.clone();
let result = run_tailscale_blocking(termination.clone(), move |executor| {
renewing.renew(executor)
})
.await;
if let Err(error) = result {
if !termination.is_cancelled() {
tracing::warn!(
error = format!("{error:#}"),
"Tailscale certificate renewal failed"
);
}
continue;
}
if let Err(error) = rustls
.reload_from_pem_file(tailscale.cert_path(), tailscale.key_path())
.await
{
tracing::warn!(%error, "could not activate renewed Tailscale certificate");
}
}
})
}
const MAX_CONCURRENT_PHONE_ACTIONS: usize = 4;
const MAX_CONCURRENT_BUNDLE_CREATIONS: usize = 4;
const MAX_CONCURRENT_PREFLIGHTS: usize = 4;
struct PhoneActionStarted {
action_id: u64,
session: SessionRecord,
published: tokio::sync::oneshot::Sender<std::result::Result<(), String>>,
}
#[derive(Default)]
struct PendingActionReplies(
std::collections::BTreeMap<u64, tokio::sync::oneshot::Sender<ActionOutcome>>,
);
impl PendingActionReplies {
fn accept(
&mut self,
action_id: u64,
action: &ControllerAction,
reply: tokio::sync::oneshot::Sender<ActionOutcome>,
) {
if matches!(action, ControllerAction::New { .. }) {
self.0.insert(action_id, reply);
} else {
if reply.send(ActionOutcome::accepted()).is_err() {
tracing::debug!(
action_id,
"phone action acceptance reply dropped after client disconnect"
);
}
}
}
fn resolve(&mut self, action_id: u64, outcome: ActionOutcome) {
if let Some(reply) = self.0.remove(&action_id)
&& reply.send(outcome).is_err()
{
tracing::debug!(
action_id,
"phone action completion reply dropped after client disconnect"
);
}
}
}
fn admit_phone_action(
action: &ControllerAction,
running_actions: usize,
active_sessions: &mut std::collections::BTreeSet<String>,
) -> std::result::Result<Option<String>, ActionOutcome> {
let closing = matches!(
action,
ControllerAction::Close { .. } | ControllerAction::ForceClose { .. }
);
if !closing && !phone_action_capacity_available(running_actions) {
return Err(ActionOutcome::Busy);
}
let session_id = controller_action_session_id(action);
if let Some(session_id) = &session_id
&& !active_sessions.insert(session_id.clone())
&& !closing
{
return Err(ActionOutcome::SessionBusy);
}
Ok(session_id)
}
struct ReadReceiptPersisted {
session_id: String,
result: std::result::Result<u64, String>,
reply: tokio::sync::oneshot::Sender<std::result::Result<(), String>>,
}
struct ControllerReloaded {
result: std::result::Result<Controller, String>,
}
struct BundleCreated {
result: std::result::Result<
crate::controller::QuickBundleCreation,
crate::controller::QuickBundleFailure,
>,
reply: tokio::sync::oneshot::Sender<std::result::Result<String, crate::server::BundleFailure>>,
}
struct MovePrepared {
result: std::result::Result<mj_core::state::MovePreparation, String>,
reply:
tokio::sync::oneshot::Sender<std::result::Result<mj_core::state::MovePreparation, String>>,
}
fn spawn_controller_reload(completed: tokio::sync::mpsc::UnboundedSender<ControllerReloaded>) {
spawn_controller_reload_with(completed, Controller::load);
}
fn spawn_controller_reload_with(
completed: tokio::sync::mpsc::UnboundedSender<ControllerReloaded>,
load: impl FnOnce() -> Result<Controller> + Send + 'static,
) {
tokio::spawn(async move {
let result = match tokio::task::spawn_blocking(load).await {
Ok(result) => result.map_err(|error| format!("{error:#}")),
Err(error) => Err(format!("controller reload task failed: {error}")),
};
if completed.send(ControllerReloaded { result }).is_err() {
tracing::debug!("controller reload completed after the phone control loop stopped");
}
});
}
fn request_controller_reload(
in_flight: &mut bool,
requested: &mut bool,
completed: &tokio::sync::mpsc::UnboundedSender<ControllerReloaded>,
) {
if *in_flight {
*requested = true;
} else {
*in_flight = true;
spawn_controller_reload(completed.clone());
}
}
fn request_daemon_controller_reload(daemon_runtime: Arc<RuntimeState>, reason: &'static str) {
tokio::spawn(async move {
if let Err(error) = daemon_runtime.reload_controller().await {
tracing::warn!(
error = format!("{error:#}"),
reason,
"phone operation could not refresh dashboard controller state"
);
}
});
}
#[derive(Debug, PartialEq, Eq)]
#[cfg(test)]
enum ReadReceiptPlan {
UnknownSession,
AlreadyRead,
Persist,
}
#[cfg(test)]
fn plan_read_receipt(state: &State, session_id: &str, through: u64) -> ReadReceiptPlan {
let Some(session) = state.sessions.get(session_id) else {
return ReadReceiptPlan::UnknownSession;
};
if through > session.viewed_through_event_ordinal {
ReadReceiptPlan::Persist
} else {
ReadReceiptPlan::AlreadyRead
}
}
#[cfg(test)]
fn apply_read_receipt(state: &mut State, session_id: &str, receipt: u64) -> bool {
let Some(session) = state.sessions.get_mut(session_id) else {
return false;
};
if receipt <= session.viewed_through_event_ordinal {
return false;
}
session.viewed_through_event_ordinal = receipt;
true
}
#[derive(Clone)]
struct PhoneActionControl {
cancelled: Arc<AtomicBool>,
create: Option<CreateSessionControl>,
}
impl PhoneActionControl {
fn for_action(action: &ControllerAction) -> Self {
let create =
matches!(action, ControllerAction::New { .. }).then(CreateSessionControl::default);
let cancelled = create.as_ref().map_or_else(
|| Arc::new(AtomicBool::new(false)),
|control| control.cancelled.clone(),
);
Self { cancelled, create }
}
fn request_cancel(&self) -> bool {
let accepted = self.create.as_ref().map_or_else(
|| {
self.cancelled
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
},
|control| control.request_cancel(),
);
if accepted {
self.cancelled.store(true, Ordering::Release);
}
accepted
}
#[cfg(test)]
fn grant_new_commit(&self) -> bool {
self.create
.as_ref()
.is_some_and(|control| control.grant_commit())
}
}
pub async fn run_server(
args: ServerArgs,
termination: tokio_util::sync::CancellationToken,
worker: SessionManagerChannels,
daemon_runtime: Arc<RuntimeState>,
mut workspace_updates: tokio::sync::watch::Receiver<Vec<WorkspaceRecord>>,
) -> Result<()> {
let resolved = resolve_server_args(args, termination.clone()).await?;
let bind = resolved.bind;
let mut controller = Controller::load()?;
let mut daemon_revisions = daemon_runtime.revisions();
daemon_revisions.borrow_and_update();
let mut phone_workspaces = workspace_updates.borrow_and_update().clone();
let mut quotas = std::collections::BTreeMap::new();
let subagent_quota_reports = Arc::new(std::sync::Mutex::new(quotas.clone()));
let (quota_profiles_tx, mut quota_updates_rx) = spawn_quota_refresher();
let mut quota_batch = QuotaRefreshBatch::default();
let mut published_quota_profiles = std::collections::BTreeMap::new();
republish_quota_profiles(
&controller,
&mut published_quota_profiles,
&mut quota_batch,
"a_profiles_tx,
);
let mut revision = daemon_runtime.allocate_revision();
let mut conversations = std::collections::BTreeMap::new();
let mut queued_prompts = projected_queued_prompts(&controller)?;
let mut active_user_shells = std::collections::BTreeMap::new();
let mut pending_elicitations = std::collections::BTreeMap::new();
let mut prompt_images = std::collections::BTreeSet::new();
let mut operational = std::collections::BTreeMap::new();
let mut materialized_activity = load_materialized_activity(&controller).await?;
let mut project_sources = PhoneProjectSources::default();
let (records, lifecycles) = daemon_runtime.session_projection();
controller.state.sessions = records;
let mut operations = lifecycles
.iter()
.map(|view| (view.session_id.clone(), viewer_operation(view)))
.collect::<std::collections::BTreeMap<_, _>>();
let mut move_recoveries = ViewerMoveRecoveries::new();
let (move_recovery_tx, mut move_recovery_rx) =
tokio::sync::mpsc::unbounded_channel::<Result<ViewerMoveRecoveries, String>>();
let mut move_recovery_load_in_flight = false;
let mut launch_failures = Vec::new();
let mut capacity_state: std::collections::BTreeMap<String, PhoneCapacity> =
std::collections::BTreeMap::new();
let (capacity_targets_tx, capacity_triggers_tx, mut capacity_updates_rx) =
crate::pollers::spawn_dashboard_capacity_poller();
let (snapshot_tx, snapshot_rx) = tokio::sync::watch::channel(viewer_snapshot(
&controller,
&phone_workspaces,
"as,
&PhoneSessionViews {
conversations: &conversations,
queued_prompts: &queued_prompts,
active_user_shells: &active_user_shells,
pending_elicitations: &pending_elicitations,
prompt_images: &prompt_images,
operational: &operational,
materialized_activity: &materialized_activity,
project_sources: &project_sources,
operations: &operations,
move_recoveries: &move_recoveries,
capacity: &viewer_capacity(&capacity_state),
launch_failures: &launch_failures,
reviews: &review_views(&daemon_runtime),
},
revision,
));
let (conversation_tx, conversation_rx) = tokio::sync::watch::channel(conversations.clone());
let (action_tx, mut action_rx) = tokio::sync::mpsc::channel(32);
let (bundle_tx, mut bundle_rx) = tokio::sync::mpsc::channel(16);
let (receipt_tx, mut receipt_rx) = tokio::sync::mpsc::channel(32);
let (preflight_tx, mut preflight_rx) = tokio::sync::mpsc::channel(32);
let (move_preparation_tx, mut move_preparation_rx) = tokio::sync::mpsc::channel(32);
let (client_state_tx, mut client_state_rx) = tokio::sync::mpsc::channel(64);
let (dictation_tx, mut dictation_rx) =
tokio::sync::mpsc::channel::<crate::dictation::DictationRequest>(8);
let (background_task_stop_tx, mut background_task_stop_rx) =
tokio::sync::mpsc::channel::<BackgroundTaskStopRequest>(32);
let SessionManagerChannels {
targets: worker_targets_tx,
control: worker_commands_tx,
updates: mut worker_updates_rx,
shutdown: worker_shutdown,
} = worker;
worker_targets_tx.send_replace(dashboard_worker_targets(&controller));
publish_capacity_targets(&controller, &capacity_targets_tx, &mut capacity_state);
let mut credential_sync = CredentialSyncCoordinator::spawn();
let credential_sync_handle = credential_sync.handle();
credential_sync_handle.set_targets(credential_sync_targets(&controller));
let mut credential_sync_signals = CredentialSyncSignalTracker::default();
let mut credential_sync_notices = CredentialSyncNotices::default();
let options_session_ttl = crate::server::default_session_ttl();
let activity_snapshots = snapshot_rx.clone();
let mut options = ServerOptions::new(
bind,
snapshot_rx,
conversation_rx,
crate::server::ServerRequests {
action_tx,
bundle_tx,
receipt_tx,
preflight_tx,
move_preparation_tx,
client_state_tx,
dictation_tx,
},
)?;
options.set_background_task_stop_tx(background_task_stop_tx);
options.shutdown = termination.clone();
let cookie_key_path = crate::server::cookie_key_path();
options.set_cookie_key(crate::server::load_or_create_cookie_key(&cookie_key_path)?)?;
options.set_api_token(crate::server::load_or_create_api_token(
&crate::server::api_token_path(),
)?);
let api_runtime = daemon_runtime.clone();
let api_backend = Arc::new(
api::ApiBackend::new(
worker_commands_tx.client(),
Arc::new(move |session_id: &str| api_runtime.session_state(session_id)),
daemon_runtime.clone(),
)
.with_quota_reports(subagent_quota_reports.clone()),
);
options.set_subagent_backend(api_backend.clone());
let renewal_cancellation = termination.child_token();
let mut renewal_task = None;
if let Some((cert, key)) = resolved.tls_files {
let rustls = axum_server::tls_rustls::RustlsConfig::from_pem_file(cert, key)
.await
.context("load web viewer TLS certificate")?;
options.set_tls_config(rustls.clone());
if let Some(tailscale) = resolved.tailscale {
renewal_task = Some(spawn_tailscale_cert_renewer(
tailscale,
rustls,
renewal_cancellation.clone(),
));
}
} else if bind.ip().is_loopback() {
options.secure_cookie = false;
} else {
anyhow::bail!("non-loopback web viewer requires TLS");
}
let fallback_reason = resolved.fallback_reason;
let qr_login_url = if fallback_reason.is_none() && resolved.viewer_url.starts_with("https://") {
let encoded = url::form_urlencoded::byte_serialize(options.login_token().as_bytes())
.collect::<String>();
Some(format!(
"{}/auth/login?token={encoded}",
resolved.viewer_url.trim_end_matches('/')
))
} else {
None
};
let ready = crate::server::WebViewerAccess::Ready {
viewer_url: resolved.viewer_url,
viewer_code: options.viewer_code().to_owned(),
qr_login_url,
fallback_reason,
};
let serve = crate::web_viewer::serve(options, ready, &daemon_runtime.web_viewer, |access| {
daemon_runtime.publish_web_access(access);
});
let conversation_projection_shutdown = termination.child_token();
let control = async {
let mut credential_tick = tokio::time::interval(Duration::from_millis(250));
let mut prune_tick = tokio::time::interval(Duration::from_secs(60 * 60));
let client_state_retention = options_session_ttl;
let (action_done_tx, mut action_done_rx) = tokio::sync::mpsc::unbounded_channel::<(
u64,
Option<String>,
std::result::Result<(), String>,
)>();
let (action_started_tx, mut action_started_rx) =
tokio::sync::mpsc::unbounded_channel::<PhoneActionStarted>();
let (receipt_done_tx, mut receipt_done_rx) =
tokio::sync::mpsc::unbounded_channel::<ReadReceiptPersisted>();
let (controller_reload_tx, mut controller_reload_rx) =
tokio::sync::mpsc::unbounded_channel::<ControllerReloaded>();
let (bundle_done_tx, mut bundle_done_rx) =
tokio::sync::mpsc::unbounded_channel::<BundleCreated>();
let (move_prepared_tx, mut move_prepared_rx) =
tokio::sync::mpsc::unbounded_channel::<MovePrepared>();
let mut dictation_jobs = tokio::task::JoinSet::new();
let mut bundle_jobs = tokio::task::JoinSet::new();
let mut preflight_jobs = tokio::task::JoinSet::new();
let mut move_preparation_jobs = tokio::task::JoinSet::new();
let mut move_recovery_jobs = tokio::task::JoinSet::new();
let mut background_task_stop_jobs = tokio::task::JoinSet::new();
let mut background_task_stop_open = true;
let mut controller_reload_in_flight = false;
let mut controller_reload_requested = false;
let mut controller_reload_invalidated = false;
let mut pending_action_errors = std::collections::BTreeMap::<String, String>::new();
let mut active_actions = std::collections::BTreeSet::new();
let mut closing_actions = std::collections::BTreeMap::<String, u64>::new();
let mut next_action_id = 0_u64;
let mut action_cancellations = std::collections::BTreeMap::<u64, PhoneActionControl>::new();
let mut action_sessions = std::collections::BTreeMap::<u64, String>::new();
let mut action_replies = PendingActionReplies::default();
let mut launch_workspaces = std::collections::BTreeMap::new();
let mut subagent_jobs = tokio::task::JoinSet::new();
let mut subagent_completion_jobs = tokio::task::JoinSet::new();
let mut active_subagent_requests = std::collections::BTreeSet::new();
let (conversation_projection_tx, mut conversation_projection_rx) =
tokio::sync::mpsc::channel(CONVERSATION_PROJECTION_CHANNEL_CAPACITY);
let mut conversation_projections = ConversationProjectionDispatcher::new(
conversation_projection_tx,
conversation_projection_shutdown.clone(),
);
let mut quota_updates_open = true;
let mut failure: Option<anyhow::Error> = None;
request_move_recovery_reload(
&move_recovery_tx,
&mut move_recovery_load_in_flight,
&mut move_recovery_jobs,
);
macro_rules! publish_snapshot {
($revision:expr) => {
let (records, lifecycles) = daemon_runtime.session_projection();
controller.state.sessions = records;
for (session_id, error) in &pending_action_errors {
if let Some(session) = controller.state.sessions.get_mut(session_id)
&& session.last_error.is_none()
{
session.last_error = Some(error.clone());
}
}
operations = lifecycles.iter()
.map(|view| (view.session_id.clone(), viewer_operation(view)))
.collect();
if let Err(error) = snapshot_tx.send(viewer_snapshot(
&controller,
&phone_workspaces,
"as,
&PhoneSessionViews {
conversations: &conversations,
queued_prompts: &queued_prompts,
active_user_shells: &active_user_shells,
pending_elicitations: &pending_elicitations,
prompt_images: &prompt_images,
operational: &operational,
materialized_activity: &materialized_activity,
project_sources: &project_sources,
operations: &operations,
move_recoveries: &move_recoveries,
capacity: &viewer_capacity(&capacity_state),
launch_failures: &launch_failures,
reviews: &review_views(&daemon_runtime),
},
$revision,
)) {
tracing::debug!(revision = $revision, %error, "phone snapshot delivery failed; no viewer is subscribed");
}
};
}
loop {
project_sources.synchronize(&controller);
tokio::select! {
_ = termination.cancelled() => break,
request = background_task_stop_rx.recv(), if background_task_stop_open => {
let Some(request) = request else {
background_task_stop_open = false;
tracing::warn!("background-task stop request feed closed while the phone server was running");
continue;
};
let session_control = worker_commands_tx.clone();
background_task_stop_jobs.spawn(async move {
let result = match session_control.session(&request.session_id).await {
Ok(session) => session
.client()
.stop_background_task(request.background_task_id.clone())
.await
.map_err(|error| {
tracing::warn!(
session_id = %request.session_id,
background_task_id = %request.background_task_id,
%error,
"provider rejected background-task stop"
);
BackgroundTaskStopFailure::Provider
}),
Err(error) => {
tracing::warn!(
session_id = %request.session_id,
%error,
"could not resolve live session for background-task stop"
);
Err(BackgroundTaskStopFailure::SessionUnavailable)
}
};
if request.reply.send(result).is_err() {
tracing::debug!(
session_id = %request.session_id,
background_task_id = %request.background_task_id,
"background-task stop result dropped after viewer disconnected"
);
}
});
}
completed = background_task_stop_jobs.join_next(), if !background_task_stop_jobs.is_empty() => {
if let Some(Err(error)) = completed {
tracing::error!(%error, "background-task stop task failed unexpectedly");
}
}
move_reloaded = move_recovery_rx.recv() => {
let Some(result) = move_reloaded else {
failure = feed_stopped(
termination.is_cancelled(),
"the Move recovery projection stopped while the phone server was running",
);
break;
};
move_recovery_load_in_flight = false;
match result {
Ok(recoveries) => {
move_recoveries = recoveries;
revision = daemon_runtime.allocate_revision();
publish_snapshot!(revision);
}
Err(error) => tracing::warn!(%error, "could not refresh Move recovery projection"),
}
}
resolved = project_sources.jobs.join_next(), if !project_sources.jobs.is_empty() => {
match resolved {
Some(Ok(resolved)) => project_sources.complete(resolved),
Some(Err(error)) => {
failure = Some(anyhow::anyhow!("web project source task failed: {error}"));
break;
}
None => unreachable!("project source jobs were not empty"),
}
revision = daemon_runtime.allocate_revision();
publish_snapshot!(revision);
}
changed = daemon_revisions.changed() => {
if changed.is_err() {
failure = feed_stopped(
termination.is_cancelled(),
"the daemon stopped publishing runtime revisions to the phone server",
);
break;
}
daemon_revisions.borrow_and_update();
revision = daemon_runtime.allocate_revision();
publish_snapshot!(revision);
request_controller_reload(
&mut controller_reload_in_flight,
&mut controller_reload_requested,
&controller_reload_tx,
);
}
changed = workspace_updates.changed() => {
if changed.is_err() {
failure = feed_stopped(
termination.is_cancelled(),
"the daemon stopped publishing workspaces to the phone server",
);
break;
}
phone_workspaces = workspace_updates.borrow_and_update().clone();
revision = daemon_runtime.allocate_revision();
publish_snapshot!(revision);
}
update = capacity_updates_rx.recv() => {
let Some(update) = update else {
failure = feed_stopped(termination.is_cancelled(), "the capacity poller stopped while the phone server was running");
break;
};
if let Some(entry) = capacity_state.get_mut(&update.target_id) {
entry.refreshing = false;
entry.sampled_at_epoch_seconds = Some(update.sampled_at_epoch_seconds);
match update.result {
Ok(usage) => {
entry.on_demand = usage.is_none();
entry.usage = usage;
entry.failed = false;
}
Err(_) => entry.failed = true,
}
}
revision = daemon_runtime.allocate_revision();
publish_snapshot!(revision);
}
update = quota_updates_rx.recv(), if quota_updates_open => {
match update {
Some(QuotaUpdate::Report(outcome)) => {
if outcome.credentials_changed {
credential_sync_handle
.sync_profile_now(&outcome.report.profile_id, None);
}
quotas.insert(outcome.report.profile_id.clone(), outcome.report.clone());
subagent_quota_reports
.lock()
.expect("sub-agent quota reports lock poisoned")
.insert(outcome.report.profile_id.clone(), outcome.report);
revision = daemon_runtime.allocate_revision();
publish_snapshot!(revision);
}
Some(QuotaUpdate::Refreshing { .. } | QuotaUpdate::Finished { .. }) => {}
None => {
quota_updates_open = false;
tracing::warn!("quota refresher stopped while the phone server is running");
}
}
}
projected = conversation_projection_rx.recv() => {
let Some(projected) = projected else {
failure = feed_stopped(
termination.is_cancelled(),
"the browser transcript projection feed stopped",
);
break;
};
let session_id = projected.session_id.clone();
let session_active = controller
.state
.sessions
.get(&session_id)
.is_some_and(|session| session.state.is_active());
if !session_active {
conversation_projections.forget(&session_id);
}
if let Some((session_id, _key, transcript)) =
conversation_projections.finish(projected, session_active)
{
conversations.insert(session_id, transcript);
revision = daemon_runtime.allocate_revision();
conversation_tx.send_replace(conversations.clone());
publish_snapshot!(revision);
} else if !session_active && conversations.remove(&session_id).is_some() {
revision = daemon_runtime.allocate_revision();
conversation_tx.send_replace(conversations.clone());
publish_snapshot!(revision);
}
tokio::task::yield_now().await;
}
update = worker_updates_rx.recv() => {
let Some(update) = update else {
failure = feed_stopped(termination.is_cancelled(), "the session manager stopped; the phone server can no longer follow sessions");
break;
};
if let Some(snapshot) = update.view.snapshot.as_ref()
&& let Some(session) = controller.state.sessions.get(&update.session_id)
&& let Some(signal) = snapshot.latest_credential_sync_signal.clone()
{
credential_sync_signals.observe(
&update.session_id,
&session.last_profile,
signal,
);
}
schedule_due_credential_syncs(
&mut credential_sync_signals,
&credential_sync_handle,
Instant::now(),
);
apply_worker_record_update(&mut controller, &update);
if let Some(snapshot) = update.view.snapshot {
for request in snapshot.subagent_requests.iter().cloned() {
let identity = (update.session_id.clone(), request.request_id.clone());
if !active_subagent_requests.insert(identity.clone()) {
continue;
}
let backend = api_backend.clone();
let runtime = daemon_runtime.clone();
let parent_session_id = update.session_id.clone();
subagent_jobs.spawn(async move {
let result = backend
.execute_subagent_tool(parent_session_id.clone(), request)
.await;
let outcome = async {
let handle = runtime
.workspace_session_handle(&parent_session_id)
.await?;
let mut lease = handle.lease_connection().await?;
lease
.connection_mut()
.complete_subagent_request(result)
.await?;
lease.release();
anyhow::Ok(())
}
.await;
(identity, outcome)
});
}
if let Some(relation) = controller.state.subagents.get_mut(&update.session_id)
&& matches!(snapshot.materialized.execution, mj_core::state::MaterializedExecutionState::Idle)
&& let Some(outcome) = snapshot.materialized.last_turn_outcome.as_ref()
&& relation.delivered_turn != Some(outcome.completed_ordinal)
{
let turn = outcome.completed_ordinal;
let output = snapshot
.materialized
.transcript
.iter()
.rev()
.find_map(|item| match &item.body {
mj_core::transcript::TranscriptBody::Agent { chunks, .. }
if item.position >= outcome.turn_start_position.unwrap_or(0) =>
{
Some(mj_core::transcript::materialized_chunks_text(chunks))
}
_ => None,
})
.unwrap_or_else(|| "The sub-agent completed without a final text response.".to_owned());
let child_id = relation.child_session_id.clone();
let parent_id = relation.parent_session_id.clone();
let task_name = relation.task_name.clone();
let outcome_name = format!("{:?}", outcome.outcome).to_lowercase();
relation.delivered_turn = Some(turn);
let backend = api_backend.clone();
subagent_completion_jobs.spawn(async move {
let result = async {
backend
.deliver_subagent_completion(
parent_id,
&child_id,
&task_name,
turn,
&outcome_name,
&output,
)
.await?;
tokio::task::spawn_blocking({
let child_id = child_id.clone();
move || crate::database::mark_subagent_turn_delivered(&child_id, turn)
})
.await??;
anyhow::Ok(())
}
.await;
(child_id, turn, result)
});
}
if snapshot.operational.native_session_is_ready()
&& operational.get(&update.session_id).is_none_or(|old: &mj_core::relay::RelayOperationalState| old.config_options != snapshot.operational.config_options)
&& let Some(session) = controller.state.sessions.get(&update.session_id)
&& matches!(session.target, Some(mj_core::state::TargetLocator::LocalBare { .. } | mj_core::state::TargetLocator::SshBare { .. } | mj_core::state::TargetLocator::AwsEc2 { .. }))
&& let Some(build) = snapshot.worker_build.clone()
{
let profile = session.last_profile.clone();
let state = snapshot.operational.clone();
tokio::spawn(async move {
if let Err(error) = crate::controller::profile_config::observe(profile, build, state).await {
tracing::warn!(%error, "could not cache observed profile choices");
}
});
}
let materialized = snapshot.materialized;
let operational_state = snapshot.operational;
materialized_activity.insert(
update.session_id.clone(),
materialized.last_activity_at_ms,
);
let queued = queued_prompt_projection(&materialized);
let pending = materialized.pending_elicitations.clone();
let active_shells = operational_state.active_user_shells.clone();
let prompt_images_supported =
agent_accepts_prompt_images(&operational_state);
active_user_shells.insert(
update.session_id.clone(),
active_shells,
);
conversation_projections.enqueue(materialized);
queued_prompts.insert(
update.session_id.clone(),
queued,
);
pending_elicitations.insert(
update.session_id.clone(),
pending,
);
if prompt_images_supported {
prompt_images.insert(update.session_id.clone());
} else {
prompt_images.remove(&update.session_id);
}
operational.insert(
update.session_id.clone(),
operational_state,
);
revision = daemon_runtime.allocate_revision();
conversation_tx.send_replace(conversations.clone());
publish_snapshot!(revision);
}
tokio::task::yield_now().await;
}
completed = subagent_jobs.join_next(), if !subagent_jobs.is_empty() => {
match completed {
Some(Ok((identity, Ok(())))) => {
active_subagent_requests.remove(&identity);
}
Some(Ok((identity, Err(error)))) => {
active_subagent_requests.remove(&identity);
tracing::warn!(
parent_session_id = %identity.0,
request_id = %identity.1,
error = %format!("{error:#}"),
"sub-agent tool request failed"
);
}
Some(Err(error)) => tracing::warn!(%error, "sub-agent tool task panicked"),
None => {}
}
}
completed = subagent_completion_jobs.join_next(), if !subagent_completion_jobs.is_empty() => {
match completed {
Some(Ok((_, _, Ok(())))) => {}
Some(Ok((child_id, turn, Err(error)))) => {
if let Some(relation) = controller.state.subagents.get_mut(&child_id)
&& relation.delivered_turn == Some(turn)
{
relation.delivered_turn = None;
}
tracing::warn!(%child_id, turn, error = %format!("{error:#}"), "could not deliver sub-agent completion");
}
Some(Err(error)) => tracing::warn!(%error, "sub-agent completion task panicked"),
None => {}
}
}
_ = prune_tick.tick() => {
tokio::spawn(async move {
let pruned = tokio::task::spawn_blocking(move || {
crate::database::prune_phone_client_state(client_state_retention)
})
.await;
match pruned {
Ok(Ok(0)) => {}
Ok(Ok(rows)) => tracing::debug!(rows, "pruned expired phone viewer state"),
Ok(Err(error)) => tracing::warn!(%error, "could not prune phone viewer state"),
Err(error) => tracing::warn!(%error, "phone viewer state pruning task failed"),
}
});
}
_ = credential_tick.tick() => {
schedule_due_credential_syncs(
&mut credential_sync_signals,
&credential_sync_handle,
Instant::now(),
);
while let Some(result) = credential_sync.try_result() {
crate::pollers::log_credential_sync_actions(&result);
let harness = controller
.config
.profiles
.get(&result.profile_id)
.map(|profile| profile.kind);
if let Some(notice) = credential_sync_notices.notice(&result, harness) {
eprintln!("Mjolnir: {notice}");
}
}
}
request = dictation_rx.recv() => {
let Some(request) = request else {
failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering dictation requests");
break;
};
let paths = controller.state.sessions.get(&request.session_id).map(|session| {
crate::dictation::auth_paths(&controller.config, &session.last_profile)
});
dictation_jobs.spawn(crate::dictation::execute(
request, paths, termination.clone(),
));
}
job = dictation_jobs.join_next(), if !dictation_jobs.is_empty() => {
if let Some(Err(error)) = job {
tracing::warn!(%error, "web dictation task failed");
}
}
stored = client_state_rx.recv() => {
let Some(stored) = stored else {
failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering viewer state requests");
break;
};
let workspace_of = |session_id: &str| {
controller
.state
.sessions
.get(session_id)
.map(|session| session.workspace_id.clone())
};
let bundle_of = |session_id: &str| {
controller
.state
.sessions
.get(session_id)
.map(|session| session.bundle_id.clone())
};
match stored {
crate::server::ClientStateRequest::Read { client_id, session_id, reply } => {
let workspace = workspace_of(&session_id);
tokio::spawn(async move {
let answer = tokio::task::spawn_blocking(move || {
let workspace = workspace.context("unknown session")?;
let state = crate::database::client_session_state(
&client_id, &workspace, &session_id,
)?;
anyhow::Ok(crate::server::ViewerClientState {
draft: state.draft,
through_event_ordinal: state.through_event_ordinal,
})
})
.await;
reply.send(flatten_stored(answer)).ok();
});
}
crate::server::ClientStateRequest::SaveDraft { client_id, session_id, draft, reply } => {
let workspace = workspace_of(&session_id);
tokio::spawn(async move {
let answer = tokio::task::spawn_blocking(move || {
let workspace = workspace.context("unknown session")?;
crate::database::persist_client_draft(
&client_id, &workspace, &session_id, &draft,
)
})
.await;
reply.send(flatten_stored(answer)).ok();
});
}
crate::server::ClientStateRequest::MarkWorkspaceRead { client_id, workspace_id, reply } => {
let sessions = controller
.state
.sessions
.values()
.filter(|session| session.workspace_id == workspace_id)
.map(|session| (session.id.clone(), session.viewed_through_event_ordinal))
.collect::<Vec<_>>();
tokio::spawn(async move {
let answer = tokio::task::spawn_blocking(move || {
for (session_id, through) in sessions {
crate::database::persist_read_receipt(
&client_id, &workspace_id, &session_id, through,
)
.ok();
}
anyhow::Ok(())
})
.await;
reply.send(flatten_stored(answer)).ok();
});
}
crate::server::ClientStateRequest::History { session_id, query, scope, reply } => {
let bundle = bundle_of(&session_id);
tokio::spawn(async move {
let answer = tokio::task::spawn_blocking(move || {
let bundle = bundle.context("unknown session")?;
let scope = match scope.as_str() {
"session" => crate::database::HistoryScope::Session,
"all" => crate::database::HistoryScope::All,
_ => crate::database::HistoryScope::Project,
};
let found = crate::database::search_prompts_bounded(
&session_id,
&bundle,
scope,
&query,
crate::server::MAX_HISTORY_MATCHES,
)?;
anyhow::Ok(crate::server::ViewerPromptHistory {
entries: found
.entries
.into_iter()
.map(|entry| entry.text)
.collect(),
truncated: found.truncated,
})
})
.await;
reply.send(flatten_stored(answer)).ok();
});
}
}
}
bundle = bundle_rx.recv(), if bundle_jobs.len() < MAX_CONCURRENT_BUNDLE_CREATIONS => {
let Some(crate::server::BundleRequest { source, reply }) = bundle else {
failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering bundle requests");
break;
};
let done = bundle_done_tx.clone();
let daemon_runtime = daemon_runtime.clone();
bundle_jobs.spawn(async move {
let result = daemon_runtime
.create_quick_bundle(source)
.await;
if let Err(error) = done.send(BundleCreated { result, reply }) {
tracing::debug!(%error, "bundle creation finished after the server stopped");
}
});
}
bundle_done = bundle_done_rx.recv() => {
let Some(BundleCreated { result, reply }) = bundle_done else {
failure = feed_stopped(termination.is_cancelled(), "the bundle creation pipeline stopped while the phone server was running");
break;
};
match result {
Ok(created) => {
let bundle_id = created.bundle_id;
let Some(bundle) = created.config.bundles.get(&bundle_id) else {
tracing::error!(%bundle_id, "bundle creation returned a config without its bundle");
if reply.send(Err(crate::server::BundleFailure::Controller)).is_err() {
tracing::debug!("bundle creation failure reply dropped after client disconnect");
}
continue;
};
controller
.config
.bundles
.insert(bundle_id.clone(), bundle.clone());
controller_reload_invalidated |= controller_reload_in_flight;
revision = daemon_runtime.allocate_revision();
publish_snapshot!(revision);
request_daemon_controller_reload(
daemon_runtime.clone(),
"new bundle publication",
);
if reply.send(Ok(bundle_id)).is_err() {
tracing::debug!("bundle creation reply dropped after client disconnect");
}
}
Err(error) => {
let failure = match error {
crate::controller::QuickBundleFailure::InvalidSource(
detail,
) => {
tracing::debug!(error = %detail, "phone bundle source was invalid");
crate::server::BundleFailure::InvalidSource
}
crate::controller::QuickBundleFailure::Persistence(
detail,
) => {
tracing::warn!(error = %detail, "phone bundle creation failed");
crate::server::BundleFailure::Controller
}
};
if reply.send(Err(failure)).is_err() {
tracing::debug!("bundle creation failure reply dropped after client disconnect");
}
}
}
}
bundle_job = bundle_jobs.join_next(), if !bundle_jobs.is_empty() => {
if let Some(Err(error)) = bundle_job {
tracing::warn!(%error, "bundle creation task panicked");
}
}
preflight = preflight_rx.recv(), if preflight_jobs.len() < MAX_CONCURRENT_PREFLIGHTS => {
let Some(preflight) = preflight else {
failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering preflight requests");
break;
};
let crate::server::NewPreflightRequest {
bundle_id,
target_id,
project_directory,
mut reply,
remote_repairs,
} = match preflight {
crate::server::PreflightRequest::New(request) => request,
crate::server::PreflightRequest::Resume(request) => {
spawn_resume_preflight(
&mut preflight_jobs,
&controller,
request,
&termination,
);
continue;
}
};
let config = controller.config.clone();
let project_validation = project_directory.is_some();
let task_termination = termination.clone();
preflight_jobs.spawn(async move {
let cancelled = Arc::new(AtomicBool::new(false));
let cancellation_guard = ProcessCancellationGuard(cancelled.clone());
let mut blocking = tokio::task::spawn_blocking(move || {
run_new_preflight_with_cancellation(
config,
bundle_id,
target_id,
project_directory,
cancelled,
remote_repairs,
)
});
let answer = tokio::select! {
biased;
_ = task_termination.cancelled() => None,
_ = reply.closed() => None,
answer = &mut blocking => Some(answer),
};
let Some(answer) = answer else {
drop(cancellation_guard);
match blocking.await {
Err(error) => tracing::warn!(%error, "cancelled phone preflight task failed"),
Ok(Err(error)) => tracing::debug!(%error, "phone preflight cancelled"),
Ok(Ok(_)) => {}
}
return;
};
let answer = match answer {
Ok(Ok(answer)) => Ok(answer),
Ok(Err(error)) => {
tracing::debug!(
error = %error,
project_validation,
"phone preflight check failed"
);
Err(if project_validation {
PreflightFailure::Validation
} else {
PreflightFailure::InvalidRepository(format!("{error:#}"))
})
}
Err(error) => {
tracing::warn!(%error, "phone preflight task failed");
Err(PreflightFailure::Controller(format!(
"preflight task failed: {error}"
)))
}
};
if reply.send(answer).is_err() {
tracing::debug!("phone preflight reply dropped after client disconnect");
}
});
}
preflight_job = preflight_jobs.join_next(), if !preflight_jobs.is_empty() => {
if let Some(Err(error)) = preflight_job {
tracing::warn!(%error, "phone preflight task panicked");
}
}
preparation = move_preparation_rx.recv() => {
let Some(MovePreparationRequest { selection, reply }) = preparation else {
failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering move preparation requests");
break;
};
let done = move_prepared_tx.clone();
let daemon_runtime = daemon_runtime.clone();
move_preparation_jobs.spawn(async move {
let result = daemon_runtime
.prepare_move_session(selection)
.await
.map_err(|error| format!("{error:#}"));
if let Err(error) = done.send(MovePrepared { result, reply }) {
tracing::debug!(%error, "move preparation finished after the server stopped");
}
});
}
prepared = move_prepared_rx.recv() => {
let Some(MovePrepared { result, reply }) = prepared else {
failure = feed_stopped(termination.is_cancelled(), "the move preparation pipeline stopped while the phone server was running");
break;
};
if reply.send(result).is_err() {
tracing::debug!("move preparation reply dropped after client disconnect");
}
}
move_preparation_job = move_preparation_jobs.join_next(), if !move_preparation_jobs.is_empty() => {
if let Some(Err(error)) = move_preparation_job {
tracing::warn!(%error, "move preparation task failed");
}
}
move_recovery_job = move_recovery_jobs.join_next(), if !move_recovery_jobs.is_empty() => {
if let Some(Err(error)) = move_recovery_job {
move_recovery_load_in_flight = false;
tracing::warn!(%error, "Move recovery projection task failed");
}
}
receipt = receipt_rx.recv() => {
let Some(ReadReceiptRequest { client_id, session_id, through, reply }) = receipt else {
failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering read receipts");
break;
};
match controller.state.sessions.get(&session_id) {
None => {
if reply.send(Err("unknown session".into())).is_err() {
tracing::debug!(%session_id, "unknown-session read receipt reply dropped after client disconnect");
}
}
Some(session) => {
let workspace_id = session.workspace_id.clone();
let done = receipt_done_tx.clone();
let persisted_session_id = session_id.clone();
tokio::spawn(async move {
let joined = tokio::task::spawn_blocking(move || {
crate::database::persist_read_receipt(
&client_id,
&workspace_id,
&persisted_session_id,
through,
)
})
.await;
let result = match joined {
Ok(result) => result.map_err(|error| format!("{error:#}")),
Err(error) => Err(format!("phone read receipt task failed: {error}")),
};
if let Err(error) = done.send(ReadReceiptPersisted { session_id, result, reply }) {
tracing::debug!(%error, "phone read receipt finished after the server stopped");
}
});
}
}
}
persisted = receipt_done_rx.recv() => {
let Some(ReadReceiptPersisted { session_id, result, reply }) = persisted else { continue };
match result {
Ok(receipt) => {
let _ = receipt;
if reply.send(Ok(())).is_err() {
tracing::debug!(%session_id, "phone read receipt reply dropped after client disconnect");
}
}
Err(error) => {
tracing::warn!(%session_id, "could not persist a phone read receipt: {error}");
if reply.send(Err(error)).is_err() {
tracing::debug!(%session_id, "failed phone read receipt reply dropped after client disconnect");
}
}
}
}
action = action_rx.recv() => {
let Some(request) = action else {
failure = feed_stopped(termination.is_cancelled(), "the phone HTTP server stopped delivering actions");
break;
};
match &request.action {
ControllerAction::RefreshCapacity { target_id } => {
let known = capacity_state.contains_key(target_id);
let accepted = known && match capacity_triggers_tx.try_send(()) {
Ok(()) | Err(tokio::sync::mpsc::error::TrySendError::Full(())) => true,
Err(tokio::sync::mpsc::error::TrySendError::Closed(())) => {
tracing::warn!("phone capacity refresh rejected: poller stopped");
false
}
};
if known {
if let Some(entry) = capacity_state.get_mut(target_id) {
entry.refreshing = accepted;
if !accepted {
entry.failed = true;
}
}
revision = daemon_runtime.allocate_revision();
publish_snapshot!(revision);
}
let outcome = if accepted {
ActionOutcome::accepted()
} else {
ActionOutcome::Failed
};
if request.reply.send(outcome).is_err() {
tracing::debug!(%target_id, "phone capacity refresh reply dropped after client disconnect");
}
tokio::task::yield_now().await;
continue;
}
ControllerAction::RefreshQuota { profile_id } => {
let known = controller.config.enabled_profile(profile_id).is_some();
if known {
quota_batch.generation = quota_batch.generation.saturating_add(1);
quota_batch.profiles = quota_refresh_profiles(&controller);
quota_profiles_tx.send_replace(quota_batch.clone());
}
let outcome = if known {
ActionOutcome::accepted()
} else {
ActionOutcome::Failed
};
if request.reply.send(outcome).is_err() {
tracing::debug!(%profile_id, "phone quota refresh reply dropped after client disconnect");
}
tokio::task::yield_now().await;
continue;
}
_ => {}
}
if let ControllerAction::Cancel { session_id } = &request.action {
let outcome = if request_phone_action_cancellation(
session_id,
&action_sessions,
&action_cancellations,
) {
daemon_runtime.cancel_lifecycle_if_active(session_id);
ActionOutcome::accepted()
} else {
ActionOutcome::NotCancellable
};
if request.reply.send(outcome).is_err() {
tracing::debug!(%session_id, "phone cancellation reply dropped after client disconnect");
}
tokio::task::yield_now().await;
continue;
}
if let ControllerAction::Close { session_id } = &request.action {
if closing_actions.contains_key(session_id) {
if request.reply.send(ActionOutcome::accepted()).is_err() { tracing::debug!(%session_id, "repeated close reply dropped"); }
continue;
}
request_phone_action_cancellation(session_id, &action_sessions, &action_cancellations);
daemon_runtime.request_close(session_id);
}
if let ControllerAction::ForceClose { session_id } = &request.action {
request_phone_action_cancellation(session_id, &action_sessions, &action_cancellations);
daemon_runtime.request_close(session_id);
}
let session_id = match admit_phone_action(
&request.action,
action_cancellations.len(),
&mut active_actions,
) {
Ok(session_id) => session_id,
Err(refusal) => {
if request.reply.send(refusal).is_err() {
tracing::debug!("phone action refusal reply dropped after client disconnect");
}
tokio::task::yield_now().await;
continue;
}
};
let ControllerRequest { action, reply } = request;
let done = action_done_tx.clone();
let session_control = worker_commands_tx.clone();
let daemon_runtime = daemon_runtime.clone();
let started = action_started_tx.clone();
next_action_id = next_action_id.wrapping_add(1).max(1);
let action_id = next_action_id;
if let ControllerAction::Close { session_id } | ControllerAction::ForceClose { session_id } = &action { closing_actions.insert(session_id.clone(), action_id); }
if let ControllerAction::New { workspace_id, .. } = &action {
let workspace_id = if workspace_id.is_empty() && phone_workspaces.len() == 1 {
phone_workspaces[0].id.clone()
} else {
workspace_id.clone()
};
launch_workspaces.insert(action_id, workspace_id);
}
let control = PhoneActionControl::for_action(&action);
action_cancellations.insert(action_id, control.clone());
if let Some(session_id) = &session_id {
action_sessions.insert(action_id, session_id.clone());
}
action_replies.accept(action_id, &action, reply);
let runtime = tokio::runtime::Handle::current();
tokio::spawn(async move {
let joined = tokio::task::spawn_blocking(move || {
let result = (|| -> Result<()> {
if control.cancelled.load(Ordering::Acquire) {
bail!("phone action cancelled");
}
let mut operation_controller = Controller::load()?;
let executor =
CancellableProcessExecutor::new(control.cancelled.clone());
runtime.block_on(apply_phone_action(
&mut operation_controller,
PhoneActionServices {
sessions: &session_control,
daemon_runtime: &daemon_runtime,
},
action,
&executor,
action_id,
&started,
&control,
))
})();
result.map_err(|error| format!("{error:#}"))
})
.await;
let result = match joined {
Ok(result) => result,
Err(error) => Err(format!("phone action task failed: {error}")),
};
if let Err(error) = done.send((action_id, session_id, result)) {
tracing::debug!(action_id, %error, "phone action finished after the server stopped");
}
});
}
started = action_started_rx.recv() => {
let Some(started) = started else {
tokio::task::yield_now().await;
continue;
};
let started_session_id = started.session.id.clone();
let publication = if !action_cancellations.contains_key(&started.action_id) {
Err("phone action completed before its provisional session was published".into())
} else {
track_started_phone_session(
&mut controller.state,
&mut active_actions,
&mut action_sessions,
started.action_id,
started.session,
)
};
if publication.is_ok() {
revision = daemon_runtime.allocate_revision();
publish_snapshot!(revision);
request_daemon_controller_reload(
daemon_runtime.clone(),
"new session publication",
);
};
if publication.is_err()
&& let Some(control) = action_cancellations.get(&started.action_id)
{
control.request_cancel();
}
action_replies.resolve(
started.action_id,
if publication.is_ok() {
ActionOutcome::Accepted {
session_id: Some(started_session_id),
}
} else {
ActionOutcome::Failed
},
);
if started.published.send(publication).is_err() {
tracing::debug!(action_id = started.action_id, "phone new-session publication reply dropped after client disconnect");
}
}
completed = action_done_rx.recv() => {
let Some((action_id, session_id, result)) = completed else {
failure = feed_stopped(termination.is_cancelled(), "the phone action pipeline stopped reporting completions");
break;
};
action_cancellations.remove(&action_id);
let session_id = action_sessions.remove(&action_id).or(session_id);
if closing_actions.values().any(|closing_id| *closing_id == action_id) && let Some(id) = &session_id { daemon_runtime.clear_close_request(id); }
closing_actions.retain(|_, closing_id| *closing_id != action_id);
if let Some(session_id) = &session_id && !action_sessions.values().any(|active| active == session_id) {
active_actions.remove(session_id);
}
action_replies.resolve(action_id, ActionOutcome::Failed);
if let Some(workspace_id) = launch_workspaces.remove(&action_id)
&& result.is_err()
&& !session_id.as_ref().is_some_and(|id| closing_actions.contains_key(id))
{
record_launch_failure(
&mut launch_failures,
action_id,
workspace_id,
session_id.clone(),
result.as_ref().err().cloned(),
);
revision = daemon_runtime.allocate_revision();
publish_snapshot!(revision);
}
if let Err(error) = &result {
tracing::warn!(
action_id,
session_id = session_id.as_deref(),
%error,
"phone action failed"
);
}
record_action_result(
&mut pending_action_errors,
session_id.as_deref(),
&result,
);
request_controller_reload(
&mut controller_reload_in_flight,
&mut controller_reload_requested,
&controller_reload_tx,
);
request_move_recovery_reload(
&move_recovery_tx,
&mut move_recovery_load_in_flight,
&mut move_recovery_jobs,
);
request_daemon_controller_reload(
daemon_runtime.clone(),
"phone action completion",
);
}
reloaded = controller_reload_rx.recv() => {
let Some(ControllerReloaded { result }) = reloaded else {
failure = feed_stopped(
termination.is_cancelled(),
"the controller reload pipeline stopped while the phone server was running",
);
break;
};
controller_reload_in_flight = false;
if std::mem::take(&mut controller_reload_invalidated) {
if let Err(error) = &result {
tracing::warn!(%error, "superseded controller reload failed");
}
controller_reload_requested = false;
request_controller_reload(
&mut controller_reload_in_flight,
&mut controller_reload_requested,
&controller_reload_tx,
);
continue;
}
match result {
Ok(mut reloaded) => {
for (session_id, error) in &pending_action_errors {
if let Some(session) = reloaded.state.sessions.get_mut(session_id)
&& session.last_error.is_none()
{
session.last_error = Some(error.clone());
}
}
controller = reloaded;
quotas.retain(|id, _| controller.config.enabled_profile(id).is_some());
subagent_quota_reports
.lock()
.expect("sub-agent quota reports lock poisoned")
.retain(|id, _| controller.config.enabled_profile(id).is_some());
worker_targets_tx.send_replace(dashboard_worker_targets(&controller));
publish_capacity_targets(
&controller,
&capacity_targets_tx,
&mut capacity_state,
);
credential_sync_handle.set_targets(credential_sync_targets(&controller));
republish_quota_profiles(
&controller,
&mut published_quota_profiles,
&mut quota_batch,
"a_profiles_tx,
);
queued_prompts.retain(|session_id, _| {
controller.state.sessions.contains_key(session_id)
});
pending_elicitations.retain(|session_id, _| {
controller.state.sessions.contains_key(session_id)
});
prompt_images.retain(|session_id| {
controller.state.sessions.contains_key(session_id)
});
operational.retain(|session_id, _| {
controller.state.sessions.contains_key(session_id)
});
materialized_activity.retain(|session_id, _| {
controller.state.sessions.contains_key(session_id)
});
request_move_recovery_reload(
&move_recovery_tx,
&mut move_recovery_load_in_flight,
&mut move_recovery_jobs,
);
conversations.retain(|id, _| {
controller.state.sessions.get(id).is_some_and(|session| session.state.is_active())
});
for session_id in conversation_projections.session_ids() {
if !controller
.state
.sessions
.get(&session_id)
.is_some_and(|session| session.state.is_active())
{
conversation_projections.forget(&session_id);
}
}
revision = daemon_runtime.allocate_revision();
conversation_tx.send_replace(conversations.clone());
publish_snapshot!(revision);
}
Err(error) => {
tracing::warn!(%error, "completed phone operation could not reload controller state");
}
}
if controller_reload_requested {
controller_reload_requested = false;
controller_reload_in_flight = true;
spawn_controller_reload(controller_reload_tx.clone());
}
}
}
}
dictation_jobs.shutdown().await;
bundle_jobs.shutdown().await;
preflight_jobs.shutdown().await;
move_preparation_jobs.shutdown().await;
move_recovery_jobs.shutdown().await;
crate::controller::profile_config::cancel_all();
for control in action_cancellations.values() {
control.request_cancel();
}
match failure {
Some(failure) => Err(failure),
None => Ok::<(), anyhow::Error>(()),
}
};
let result = tokio::select! {
result = api_activity::record_activity_stream(activity_snapshots) => result.context("native API activity recorder stopped"),
result = serve => result,
result = control => result,
};
conversation_projection_shutdown.cancel();
renewal_cancellation.cancel();
if let Some(task) = renewal_task
&& let Err(error) = task.await
{
tracing::warn!(%error, "Tailscale certificate renewal task failed");
}
worker_shutdown
.shutdown()
.await
.context("shut down phone server session manager")?;
result?;
Ok(())
}
fn feed_stopped(shutting_down: bool, reason: &'static str) -> Option<anyhow::Error> {
(!shutting_down).then(|| anyhow::anyhow!(reason))
}
fn agent_accepts_prompt_images(operational: &mj_core::relay::RelayOperationalState) -> bool {
operational
.agent_capabilities
.as_ref()
.is_some_and(|capabilities| capabilities.prompt_capabilities.image)
}
fn controller_action_session_id(action: &ControllerAction) -> Option<String> {
match action {
ControllerAction::New { .. } => None,
ControllerAction::Prompt { session_id, .. }
| ControllerAction::RunShell { session_id, .. }
| ControllerAction::CancelShell { session_id, .. }
| ControllerAction::Close { session_id }
| ControllerAction::ForceClose { session_id }
| ControllerAction::Resume { session_id, .. }
| ControllerAction::Open { session_id }
| ControllerAction::Cancel { session_id }
| ControllerAction::RemoveQueuedPrompt { session_id, .. }
| ControllerAction::RespondElicitation { session_id, .. }
| ControllerAction::Rename { session_id, .. }
| ControllerAction::CancelTurn { session_id }
| ControllerAction::SetConfig { session_id, .. }
| ControllerAction::SetPlanMode { session_id, .. }
| ControllerAction::StartReview { session_id }
| ControllerAction::ResolveReview { session_id, .. } => Some(session_id.clone()),
ControllerAction::Move { request } => {
Some(request.preparation.selection.session_id.clone())
}
ControllerAction::RefreshQuota { .. } | ControllerAction::RefreshCapacity { .. } => None,
}
}
fn flatten_stored<T>(
joined: std::result::Result<Result<T>, tokio::task::JoinError>,
) -> std::result::Result<T, String> {
match joined {
Ok(Ok(value)) => Ok(value),
Ok(Err(error)) => Err(format!("{error:#}")),
Err(error) => Err(format!("viewer state task failed: {error}")),
}
}
struct ProcessCancellationGuard(Arc<AtomicBool>);
impl Drop for ProcessCancellationGuard {
fn drop(&mut self) {
self.0.store(true, Ordering::Release);
}
}
#[cfg(test)]
fn run_new_preflight(
config: Config,
bundle_id: String,
target_id: String,
project_directory: Option<PathBuf>,
) -> Result<crate::server::PreflightNew> {
run_new_preflight_with_cancellation(
config,
bundle_id,
target_id,
project_directory,
Arc::new(AtomicBool::new(false)),
Vec::new(),
)
}
fn spawn_resume_preflight(
jobs: &mut tokio::task::JoinSet<()>,
controller: &Controller,
request: crate::server::ResumePreflightRequest,
termination: &tokio_util::sync::CancellationToken,
) {
let crate::server::ResumePreflightRequest {
session_id,
target_id,
mut reply,
} = request;
let config = controller.config.clone();
let session = controller.state.sessions.get(&session_id).cloned();
let termination = termination.clone();
jobs.spawn(async move {
let cancelled = Arc::new(AtomicBool::new(false));
let cancellation_guard = ProcessCancellationGuard(cancelled.clone());
let mut blocking = tokio::task::spawn_blocking(move || {
run_resume_preflight(config, session, &target_id, cancelled)
});
let answer = tokio::select! {
biased;
_ = termination.cancelled() => None,
_ = reply.closed() => None,
answer = &mut blocking => Some(answer),
};
let Some(answer) = answer else {
drop(cancellation_guard);
if let Err(error) = blocking.await {
tracing::warn!(%error, "cancelled phone resume preflight task failed");
}
return;
};
let answer = answer.map_err(|error| {
tracing::warn!(%error, "phone resume preflight task failed");
PreflightFailure::Controller(format!("resume preflight task failed: {error}"))
});
if reply.send(answer).is_err() {
tracing::debug!("phone resume preflight reply dropped after client disconnect");
}
});
}
fn run_resume_preflight(
config: Config,
session: Option<mj_core::state::SessionRecord>,
target_id: &str,
cancelled: Arc<AtomicBool>,
) -> crate::server::PreflightResume {
let Some(session) = session else {
return crate::server::PreflightResume::Unavailable {
detail: "this session is no longer available".to_owned(),
};
};
match crate::controller::resume_compatibility(&session, &config, target_id) {
Err(reason) => crate::server::PreflightResume::Unavailable { detail: reason },
Ok(plan) if plan != crate::controller::ResumePlan::RawToWorkspace => {
crate::server::PreflightResume::Ready
}
Ok(_) => {
let executor =
CancellableProcessExecutor::new(cancelled).with_deadline(Duration::from_secs(30));
match crate::controller::raw_conversion_preview_for(&session, &config, &executor) {
Err(error) => crate::server::PreflightResume::Unavailable {
detail: format!("{error:#}"),
},
Ok(mut preview) => {
preview.fetch_url = display_url(&preview.fetch_url);
preview.push_urls = preview
.push_urls
.iter()
.map(|url| display_url(url))
.collect();
crate::server::PreflightResume::ConvertingRawCheckout {
preview: Box::new(preview),
}
}
}
}
}
}
fn run_new_preflight_with_cancellation(
config: Config,
bundle_id: String,
target_id: String,
project_directory: Option<PathBuf>,
cancelled: Arc<AtomicBool>,
remote_repairs: Vec<mj_core::local_git::LocalRemoteRepair>,
) -> Result<crate::server::PreflightNew> {
let executor =
CancellableProcessExecutor::new(cancelled).with_deadline(Duration::from_secs(30));
if !remote_repairs.is_empty() {
let target = config.targets.get(&target_id).context("unknown target")?;
anyhow::ensure!(
!is_bare_project_target(target) && project_directory.is_none(),
"remote repair requires an isolated target"
);
let bundle = config.bundles.get(&bundle_id).context("unknown bundle")?;
mj_core::local_git::apply_repository_remote_repairs(bundle, &remote_repairs, &executor)?;
}
run_new_preflight_with_executor(config, bundle_id, target_id, project_directory, &executor)
}
fn run_new_preflight_with_executor(
config: Config,
bundle_id: String,
target_id: String,
project_directory: Option<PathBuf>,
executor: &impl CommandExecutor,
) -> Result<crate::server::PreflightNew> {
let target_is_bare = config
.targets
.get(&target_id)
.with_context(|| format!("unknown target template {target_id:?}"))
.map(is_bare_project_target)?;
if target_is_bare {
let directory =
project_directory.context("project directory is required for a bare target")?;
let controller = config_only_controller(config);
let directory = controller.resolve_project_directory(&target_id, &directory, executor)?;
let managed_worktree =
controller.managed_worktree_options(&target_id, &directory, executor)?;
return Ok(crate::server::PreflightNew {
managed_worktree,
project_directory: Some(directory),
remote_repairs: Vec::new(),
dirty_repositories: Vec::new(),
remote_repositories: Vec::new(),
local_changes_excluded: false,
});
}
if project_directory.is_some() {
bail!("project directory is unsupported for this target");
}
let bundle = config.bundles.get(&bundle_id).context("unknown bundle")?;
let repairs = mj_core::local_git::repository_remote_repairs(bundle, executor)?;
if !repairs.is_empty() {
return Ok(crate::server::PreflightNew {
managed_worktree: Default::default(),
project_directory: None,
remote_repairs: repairs,
dirty_repositories: Vec::new(),
remote_repositories: Vec::new(),
local_changes_excluded: true,
});
}
let remote_repositories = bundle
.repositories
.iter()
.map(|repository| {
let source = resolve_repository(repository, executor)
.with_context(|| format!("repository {:?}", repository.id))?;
let default_branch = default_branch(&source, executor)
.with_context(|| format!("repository {:?}", repository.id))?;
Ok(crate::server::PreflightRepository {
id: repository.id.clone(),
fetch_url: display_url(&source.fetch_url),
default_branch,
push_urls: source
.push_urls
.iter()
.map(|url| display_url(url))
.collect(),
})
})
.collect::<Result<Vec<_>>>()?;
Ok(crate::server::PreflightNew {
managed_worktree: Default::default(),
project_directory: None,
remote_repairs: Vec::new(),
dirty_repositories: Vec::new(),
remote_repositories,
local_changes_excluded: true,
})
}
fn phone_action_capacity_available(active_actions: usize) -> bool {
active_actions < MAX_CONCURRENT_PHONE_ACTIONS
}
fn republish_quota_profiles(
controller: &Controller,
published: &mut std::collections::BTreeMap<String, HarnessProfile>,
batch: &mut QuotaRefreshBatch,
profiles_tx: &tokio::sync::watch::Sender<QuotaRefreshBatch>,
) -> bool {
if *published == controller.config.profiles {
return false;
}
published.clone_from(&controller.config.profiles);
batch.generation = batch.generation.saturating_add(1);
batch.profiles = quota_refresh_profiles(controller);
profiles_tx.send_replace(batch.clone());
true
}
fn request_phone_action_cancellation(
session_id: &str,
action_sessions: &std::collections::BTreeMap<u64, String>,
action_cancellations: &std::collections::BTreeMap<u64, PhoneActionControl>,
) -> bool {
let control = action_sessions
.iter()
.find_map(|(action_id, active_session_id)| {
(active_session_id == session_id)
.then(|| action_cancellations.get(action_id))
.flatten()
});
if let Some(control) = control {
return control.request_cancel();
}
false
}
fn track_started_phone_session(
state: &mut State,
active_actions: &mut std::collections::BTreeSet<String>,
action_sessions: &mut std::collections::BTreeMap<u64, String>,
action_id: u64,
session: SessionRecord,
) -> std::result::Result<(), String> {
let session_id = session.id.clone();
if !active_actions.insert(session_id.clone()) {
return Err("another operation is already running for the new session".into());
}
action_sessions.insert(action_id, session_id.clone());
state.sessions.insert(session_id, session);
Ok(())
}
fn record_action_result(
pending_action_errors: &mut std::collections::BTreeMap<String, String>,
session_id: Option<&str>,
result: &std::result::Result<(), String>,
) {
let Some(session_id) = session_id else {
return;
};
match result {
Err(error) => {
pending_action_errors.insert(session_id.to_owned(), error.clone());
}
Ok(()) => {
pending_action_errors.remove(session_id);
}
}
}
fn record_launch_failure(
failures: &mut Vec<crate::server::ViewerLaunchFailure>,
action_id: u64,
workspace_id: String,
session_id: Option<String>,
error: Option<String>,
) {
failures.push(crate::server::ViewerLaunchFailure {
id: format!("{}-{action_id}", std::process::id()),
workspace_id,
session_id,
error,
});
if failures.len() > 16 {
failures.remove(0);
}
}
struct PhoneActionServices<'a> {
sessions: &'a SessionManagerControl,
daemon_runtime: &'a Arc<RuntimeState>,
}
async fn apply_phone_action(
controller: &mut Controller,
services: PhoneActionServices<'_>,
action: ControllerAction,
_executor: &(impl CommandExecutor + Sync),
action_id: u64,
started: &tokio::sync::mpsc::UnboundedSender<PhoneActionStarted>,
control: &PhoneActionControl,
) -> Result<()> {
match action {
ControllerAction::New {
workspace_id,
profile_id,
bundle_id,
target_id,
title,
project_directory,
create_managed_worktree,
mjolnir_subagents,
dirty_ack: _dirty_ack,
} => {
let workspace_id = if workspace_id.is_empty() {
let workspaces = crate::database::list_workspaces()?;
match workspaces.as_slice() {
[workspace] => workspace.id.clone(),
[] => bail!("create a workspace before starting a phone session"),
_ => bail!("phone session creation requires a workspace_id"),
}
} else {
workspace_id
};
let title = title.unwrap_or_else(|| {
let project = project_directory
.as_ref()
.and_then(|path| path.file_name())
.map(|name| name.to_string_lossy().into_owned())
.unwrap_or_else(|| bundle_id.clone());
format!("{project} via {profile_id}")
});
let session_title_override = Some(title.clone());
let allow_dirty_local = false;
let (published, publication) = tokio::sync::oneshot::channel();
let registered = services
.daemon_runtime
.start_create_session_controlled(
CreateSessionRequest {
create_managed_worktree,
mjolnir_subagents,
initial_prompt: None,
workspace_id,
profile_id,
bundle_id,
project_directory,
target_template_id: target_id,
additional_mounts: Vec::new(),
allow_dirty_local,
resource_allocation: None,
title,
session_title_override,
},
control
.create
.clone()
.expect("New action has a daemon create control"),
publication,
)
.await?;
let registered_session_id = registered.session.id.clone();
started
.send(PhoneActionStarted {
action_id,
session: registered.session,
published,
})
.map_err(|_| anyhow::anyhow!("phone server stopped before publishing session"))?;
services
.daemon_runtime
.wait_create_session(®istered_session_id)
.await
}
ControllerAction::Prompt {
session_id,
text,
images,
} => {
services
.sessions
.wait_for_session(&session_id, Duration::from_secs(5))
.await?
.submit(
new_command_id("phone-prompt")?,
RelayCommand::Prompt {
prompt: phone_prompt_blocks(text, images),
},
)
.await?;
Ok(())
}
ControllerAction::RunShell {
session_id,
command,
} => {
services
.sessions
.wait_for_session(&session_id, Duration::from_secs(5))
.await?
.submit(
new_command_id("phone-shell")?,
RelayCommand::RunUserShell { command },
)
.await?;
Ok(())
}
ControllerAction::CancelShell {
session_id,
shell_command_id,
} => {
services
.sessions
.wait_for_session(&session_id, Duration::from_secs(5))
.await?
.submit(
new_command_id("phone-cancel-shell")?,
RelayCommand::CancelUserShell { shell_command_id },
)
.await?;
Ok(())
}
ControllerAction::Close { session_id } => {
services.daemon_runtime.close_session(session_id).await
}
ControllerAction::ForceClose { session_id } => {
services
.daemon_runtime
.force_destroy_session(session_id)
.await
}
ControllerAction::Resume {
session_id,
workspace_id,
profile_id,
target_id,
queue,
additional_mounts,
resource_allocation,
} => services
.daemon_runtime
.resume_session(ResumeSessionRequest {
session_id,
workspace_id,
profile_id,
target_template_id: target_id,
additional_mounts,
resource_allocation,
discard_queue: queue == ResumeQueueDisposition::Discard,
repository_preflight: None,
})
.await
.map(|_| ()),
ControllerAction::Move { request } => {
let outcome = services.daemon_runtime.move_session(request).await?;
match outcome.outcome.as_str() {
"completed" | "unchanged" => Ok(()),
"cancelled" | "failed" => {
let recovery = outcome.recovery.unwrap_or_else(|| {
"Inspect the session and retry Move or Resume with previous settings."
.into()
});
bail!("Move {}: {recovery}", outcome.outcome)
}
status => bail!("Move returned an unknown outcome: {status}"),
}
}
ControllerAction::Open { .. } => Ok(()),
ControllerAction::Cancel { .. } => {
bail!("cancel actions must be handled by the phone control loop")
}
ControllerAction::RemoveQueuedPrompt {
session_id,
queue_id,
} => {
services
.sessions
.session(&session_id)
.await?
.submit(
new_command_id("phone-remove-prompt")?,
RelayCommand::RemoveQueuedPrompt {
queued_command_id: queue_id,
},
)
.await?;
Ok(())
}
ControllerAction::RespondElicitation {
session_id,
elicitation_id,
response,
} => {
services
.sessions
.session(&session_id)
.await?
.respond_elicitation(elicitation_id, response)
.await
}
ControllerAction::Rename { session_id, title } => {
controller.rename_session(&session_id, &title)?;
Ok(())
}
ControllerAction::StartReview { session_id } => {
services
.daemon_runtime
.review_host()
.start(&session_id, true)
.await
.map_err(|refusal| anyhow::anyhow!("{refusal}"))?;
Ok(())
}
ControllerAction::ResolveReview {
session_id,
resolution,
} => {
let resolution = crate::server::resolution_from_name(&resolution)
.context("a review is resolved by forward, dismiss, or cancel")?;
services
.daemon_runtime
.review_host()
.resolve(&session_id, resolution)
.await
.map_err(|error| anyhow::anyhow!("{error}"))?;
Ok(())
}
ControllerAction::CancelTurn { session_id } => {
services
.sessions
.session(&session_id)
.await?
.submit(new_command_id("phone-cancel-turn")?, RelayCommand::Cancel)
.await?;
Ok(())
}
ControllerAction::SetConfig {
session_id,
key,
value,
} => {
services
.sessions
.session(&session_id)
.await?
.submit(
new_command_id("phone-set-config")?,
RelayCommand::SetConfig { key, value },
)
.await?;
Ok(())
}
ControllerAction::SetPlanMode { session_id, active } => {
let harness_kind = controller
.state
.sessions
.get(&session_id)
.with_context(|| format!("unknown session {session_id}"))?
.harness_kind;
let handle = services.sessions.session(&session_id).await?;
let operational = handle
.view()
.snapshot
.map(|snapshot| snapshot.operational)
.context("the session has not reported what it supports yet")?;
let facts = mj_core::acp::AcpSessionFacts::from_operational(
harness_kind,
&operational.config,
&operational.config_options,
operational.modes.as_ref(),
);
let command = match facts.plan_control(active) {
Ok(mj_core::acp::PlanControl::SetConfig { key, value }) => {
RelayCommand::SetConfig { key, value }
}
Ok(mj_core::acp::PlanControl::SetSessionMode { mode_id }) => {
RelayCommand::SetSessionMode { mode_id }
}
Err(reason) => bail!("{reason}"),
};
handle
.submit(new_command_id("phone-plan-mode")?, command)
.await?;
Ok(())
}
ControllerAction::RefreshQuota { .. } | ControllerAction::RefreshCapacity { .. } => {
bail!("refresh actions must be handled by the phone control loop")
}
}
}
#[derive(Clone, PartialEq, Eq)]
struct ProjectSourceKey {
directory: Option<PathBuf>,
worktree: Option<mj_core::state::ManagedWorktree>,
target: Option<mj_core::config::TargetTemplate>,
fallback: ProjectSourceIdentity,
}
impl ProjectSourceKey {
fn of(session: &SessionRecord, config: &Config) -> Self {
Self {
directory: session.project_directory.clone(),
worktree: session.managed_worktree.clone(),
target: config.targets.get(&session.target_template_id).cloned(),
fallback: session.project_source(config),
}
}
}
struct ProjectSourceEntry {
key: ProjectSourceKey,
source: Option<ProjectSourceIdentity>,
retry_at: Option<Instant>,
cancelled: Arc<AtomicBool>,
}
struct ProjectSourceResolved {
cancelled: Arc<AtomicBool>,
session_id: String,
key: ProjectSourceKey,
result: Result<ProjectSourceIdentity, String>,
}
#[derive(Default)]
struct PhoneProjectSources {
entries: std::collections::BTreeMap<String, ProjectSourceEntry>,
jobs: tokio::task::JoinSet<ProjectSourceResolved>,
}
impl PhoneProjectSources {
fn synchronize(&mut self, controller: &Controller) {
self.entries.retain(|id, entry| {
let keep = controller.state.sessions.get(id).is_some_and(|session| {
session.project_directory.is_some()
&& entry.key == ProjectSourceKey::of(session, &controller.config)
});
if !keep {
entry.cancelled.store(true, Ordering::Release);
}
keep
});
for session in controller.state.sessions.values() {
if self.jobs.len() >= 8 {
break;
}
if session.project_directory.is_none()
|| self.entries.get(&session.id).is_some_and(|entry| {
entry
.retry_at
.is_none_or(|deadline| Instant::now() < deadline)
})
{
continue;
}
let key = ProjectSourceKey::of(session, &controller.config);
let cancelled = Arc::new(AtomicBool::new(false));
self.entries.insert(
session.id.clone(),
ProjectSourceEntry {
key: key.clone(),
source: None,
retry_at: None,
cancelled: cancelled.clone(),
},
);
let source_controller = Controller {
config: controller.config.clone(),
state: State {
sessions: [(session.id.clone(), session.clone())]
.into_iter()
.collect(),
..State::default()
},
};
let session_id = session.id.clone();
self.jobs.spawn_blocking(move || {
let executor = CancellableProcessExecutor::new(cancelled.clone())
.with_deadline(Duration::from_secs(8));
let result = source_controller
.resolve_session_project_source(&session_id, &executor)
.map_err(|error| format!("{error:#}"));
ProjectSourceResolved {
cancelled,
session_id,
key,
result,
}
});
}
}
fn complete(&mut self, resolved: ProjectSourceResolved) {
let Some(entry) = self.entries.get_mut(&resolved.session_id) else {
return;
};
if entry.key != resolved.key || !Arc::ptr_eq(&entry.cancelled, &resolved.cancelled) {
return;
}
match resolved.result {
Ok(source) => entry.source = Some(source),
Err(error) => {
tracing::warn!(session_id = %resolved.session_id, %error, "could not resolve web project source");
entry.retry_at = Some(Instant::now() + Duration::from_secs(30));
}
}
}
fn source(&self, session: &SessionRecord, config: &Config) -> Option<&ProjectSourceIdentity> {
self.entries
.get(&session.id)
.filter(|entry| entry.key == ProjectSourceKey::of(session, config))
.and_then(|entry| entry.source.as_ref())
}
}
impl Drop for PhoneProjectSources {
fn drop(&mut self) {
for entry in self.entries.values() {
entry.cancelled.store(true, Ordering::Release);
}
}
}
struct PhoneSessionViews<'a> {
conversations: &'a std::collections::BTreeMap<String, crate::server::BrowserTranscript>,
queued_prompts: &'a std::collections::BTreeMap<String, Vec<mj_core::relay::QueuedPrompt>>,
active_user_shells:
&'a std::collections::BTreeMap<String, Vec<mj_core::relay::ActiveUserShell>>,
pending_elicitations:
&'a std::collections::BTreeMap<String, Vec<mj_core::elicitation::ElicitationRequest>>,
prompt_images: &'a std::collections::BTreeSet<String>,
operational: &'a std::collections::BTreeMap<String, mj_core::relay::RelayOperationalState>,
materialized_activity: &'a std::collections::BTreeMap<String, Option<i64>>,
project_sources: &'a PhoneProjectSources,
operations: &'a std::collections::BTreeMap<String, crate::server::ViewerOperation>,
move_recoveries: &'a std::collections::BTreeMap<String, ViewerMoveRecovery>,
capacity: &'a [crate::server::ViewerTargetCapacity],
launch_failures: &'a [crate::server::ViewerLaunchFailure],
reviews: &'a std::collections::BTreeMap<String, crate::review_host::RuntimeReviewView>,
}
#[derive(Debug, Clone)]
struct PhoneCapacity {
target: crate::targets::DeploymentCapacityTarget,
usage: Option<crate::targets::DeploymentCapacityUsage>,
on_demand: bool,
sampled_at_epoch_seconds: Option<u64>,
refreshing: bool,
failed: bool,
}
const CAPACITY_STALE_AFTER: Duration = Duration::from_secs(120);
fn publish_capacity_targets(
controller: &Controller,
targets_tx: &tokio::sync::watch::Sender<Vec<crate::targets::DeploymentCapacityTarget>>,
state: &mut std::collections::BTreeMap<String, PhoneCapacity>,
) {
let targets = controller.deployment_capacity_targets();
state.retain(|id, _| targets.iter().any(|target| target.id == *id));
for target in &targets {
state
.entry(target.id.clone())
.and_modify(|entry| entry.target = target.clone())
.or_insert_with(|| PhoneCapacity {
target: target.clone(),
usage: None,
on_demand: false,
sampled_at_epoch_seconds: None,
refreshing: true,
failed: false,
});
}
if targets_tx.borrow().as_slice() != targets.as_slice() {
targets_tx.send_replace(targets);
}
}
fn viewer_capacity(
state: &std::collections::BTreeMap<String, PhoneCapacity>,
) -> Vec<crate::server::ViewerTargetCapacity> {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
state
.values()
.map(|entry| {
let usage = entry.usage.as_ref();
crate::server::ViewerTargetCapacity {
id: entry.target.id.clone(),
label: entry.target.host.clone(),
target_ids: entry.target.target_ids.clone(),
cpu_percent: usage.and_then(|usage| usage.cpu_percent),
memory_used_bytes: usage.map(|usage| usage.memory_used_bytes),
memory_total_bytes: usage.map(|usage| usage.memory_total_bytes),
logical_cores: usage.map(|usage| usage.logical_cores),
disk_total_bytes: usage.and_then(|usage| usage.disk_total_bytes),
virtual_machines: matches!(
entry.target.kind,
crate::targets::DeploymentCapacityKind::AwsFleet
)
.then(|| u64::from(!entry.on_demand)),
sampled_at_epoch_seconds: entry.sampled_at_epoch_seconds,
refreshing: entry.refreshing,
stale: entry.sampled_at_epoch_seconds.is_some_and(|sampled| {
now.saturating_sub(sampled) > CAPACITY_STALE_AFTER.as_secs()
}),
has_error: entry.failed,
}
})
.collect()
}
fn viewer_operation(view: &crate::daemon::RuntimeLifecycleView) -> crate::server::ViewerOperation {
use crate::server::{ViewerOperationKind, ViewerOperationStage};
crate::server::ViewerOperation {
id: view.operation_id.clone(),
session_id: view.session_id.clone(),
kind: match view.kind {
crate::daemon::RuntimeLifecycleKind::Create => ViewerOperationKind::Create,
crate::daemon::RuntimeLifecycleKind::Resume => ViewerOperationKind::Resume,
crate::daemon::RuntimeLifecycleKind::Move => ViewerOperationKind::Move,
crate::daemon::RuntimeLifecycleKind::Close
| crate::daemon::RuntimeLifecycleKind::ForceStop => ViewerOperationKind::Stop,
crate::daemon::RuntimeLifecycleKind::DestroyStopped
| crate::daemon::RuntimeLifecycleKind::ForceDestroy => ViewerOperationKind::Destroy,
crate::daemon::RuntimeLifecycleKind::Cleanup => ViewerOperationKind::Cleanup,
},
started_at_epoch_seconds: view.started_at_epoch_seconds,
stages: view
.active_stages
.iter()
.map(|(stage, started_at)| ViewerOperationStage {
label: stage.label(),
started_at_epoch_seconds: *started_at,
})
.collect(),
notice: view.notice.clone(),
cancellable: view.cancellable,
}
}
fn session_capabilities(
session: &crate::server::ViewerSession,
operational: Option<&mj_core::relay::RelayOperationalState>,
operation: Option<&crate::server::ViewerOperation>,
facts: Option<&mj_core::acp::AcpSessionFacts>,
) -> crate::server::ViewerSessionCapabilities {
use crate::server::ViewerLifecycleCategory;
let live = session.lifecycle == ViewerLifecycleCategory::Live;
let attached = operational.is_some();
let transition_busy = session.transitioning && !session.has_error;
let busy = operation.is_some() || transition_busy;
let partial_move_queue = session.move_recovery.as_ref().is_some_and(|recovery| {
recovery.queue_admission_started && !recovery.queue_admission_finished
});
let mutation_busy = busy || partial_move_queue;
let idle = operational
.is_some_and(|state| state.execution == mj_core::relay::RelayExecutionState::Idle);
crate::server::ViewerSessionCapabilities {
open: session.conversation_available
&& !session.transitioning
&& session.lifecycle == ViewerLifecycleCategory::Live,
prompt: live && attached && !mutation_busy,
run_shell: live && attached && !mutation_busy,
cancel_turn: live
&& !mutation_busy
&& operational.is_some_and(|state| {
state.active_prompt.is_some() || state.capacity_retry.is_some()
}),
cancel_operation: operation.is_some_and(|operation| operation.cancellable),
stop: session.lifecycle.is_dashboard_visible() && !mutation_busy,
rename: !session.transitioning,
resume: !session.lifecycle.is_dashboard_visible() && !mutation_busy,
move_session: live && !busy,
set_config: live && attached && facts.is_some() && !mutation_busy,
set_plan_mode: live
&& !mutation_busy
&& idle
&& facts.is_some_and(mj_core::acp::AcpSessionFacts::supports_plan_mode),
}
}
fn phone_prompt_blocks(
text: String,
images: Vec<crate::server::ViewerPromptImage>,
) -> Vec<agent_client_protocol::schema::v1::ContentBlock> {
use agent_client_protocol::schema::v1::{ContentBlock, ImageContent, TextContent};
let mut prompt = Vec::with_capacity(images.len() + 1);
if !text.is_empty() {
prompt.push(ContentBlock::Text(TextContent::new(text)));
}
prompt.extend(images.into_iter().map(|image| match image.attachment {
Some(reference) => reference.content_block(),
None => ContentBlock::Image(ImageContent::new(image.data_base64, image.mime_type)),
}));
prompt
}
fn phone_commands(
session: &crate::server::ViewerSession,
operational: Option<&mj_core::relay::RelayOperationalState>,
) -> Vec<crate::server::ViewerMjCommand> {
use crate::server::ViewerCommandSource;
use agent_client_protocol::schema::v1::AvailableCommandInput;
let command =
|name: &str, description: &str, argument: Option<&str>| crate::server::ViewerMjCommand {
name: name.to_owned(),
description: description.to_owned(),
source: ViewerCommandSource::Mj,
argument: argument.map(str::to_owned),
};
let mut commands = vec![
command("help", "show available Mjolnir and agent commands", None),
command(
"detach",
"leave the conversation without stopping the worker",
None,
),
];
let option = |key: &str| {
session
.config_options
.iter()
.any(|option| option.key == key)
};
if option("model") {
commands.push(command("model", "change the active model", Some("value")));
commands.push(command("fast", "toggle Codex Fast mode", None));
}
if option("effort") {
commands.push(command(
"effort",
"change the active reasoning effort",
Some("value"),
));
}
if session.plan_mode_active.is_some() && session.capabilities.set_plan_mode {
commands.push(command("plan", "toggle plan mode", Some("message")));
commands.push(command(
"implement",
"leave plan mode and implement",
Some("instruction"),
));
}
if session.capabilities.prompt || session.turn_review.is_some() {
commands.push(command(
"review",
"review the finished turn now, or report how review is configured",
Some("status"),
));
}
let reserved = [
"help",
"detach",
"model",
"fast",
"effort",
"plan",
"implement",
"review",
];
for advertised in operational
.into_iter()
.flat_map(|state| state.available_commands.iter())
{
let name = advertised.name.trim();
if name.is_empty()
|| reserved
.iter()
.any(|local| name.eq_ignore_ascii_case(local))
|| commands
.iter()
.any(|existing| existing.name.eq_ignore_ascii_case(name))
{
continue;
}
let argument = advertised.input.as_ref().and_then(|input| match input {
AvailableCommandInput::Unstructured(input) => {
let hint = input.hint.trim();
(!hint.is_empty()).then(|| hint.to_owned())
}
_ => None,
});
commands.push(crate::server::ViewerMjCommand {
name: name.to_owned(),
description: advertised.description.trim().to_owned(),
source: ViewerCommandSource::Agent,
argument,
});
}
commands
}
fn review_views(
daemon_runtime: &Arc<RuntimeState>,
) -> std::collections::BTreeMap<String, crate::review_host::RuntimeReviewView> {
daemon_runtime
.review_host()
.views()
.into_iter()
.map(|review| (review.session_id.clone(), review))
.collect()
}
type ViewerMoveRecoveries = std::collections::BTreeMap<String, crate::server::ViewerMoveRecovery>;
fn request_move_recovery_reload(
completed: &tokio::sync::mpsc::UnboundedSender<Result<ViewerMoveRecoveries, String>>,
in_flight: &mut bool,
jobs: &mut tokio::task::JoinSet<()>,
) {
if *in_flight {
return;
}
*in_flight = true;
let completed = completed.clone();
jobs.spawn(async move {
let result = match tokio::task::spawn_blocking(|| {
crate::database::load_move_operations()
.map(|operations| {
operations
.into_iter()
.filter_map(|operation| {
crate::server::ViewerMoveRecovery::from_operation(&operation)
.map(|recovery| (operation.selection.session_id.clone(), recovery))
})
.collect()
})
.map_err(|error| format!("{error:#}"))
})
.await
{
Ok(result) => result,
Err(error) => Err(format!("move recovery projection task failed: {error}")),
};
let _ = completed.send(result);
});
}
fn viewer_snapshot(
controller: &Controller,
workspaces: &[mj_core::workspace::WorkspaceRecord],
quotas: &std::collections::BTreeMap<String, ProfileQuota>,
views: &PhoneSessionViews<'_>,
revision: u64,
) -> ViewerSnapshot {
let PhoneSessionViews {
conversations,
reviews,
queued_prompts,
active_user_shells,
pending_elicitations,
prompt_images,
operational,
materialized_activity,
project_sources,
operations,
move_recoveries,
capacity,
launch_failures,
} = views;
let mut snapshot =
ViewerSnapshot::from_config_state(&controller.config, &controller.state, revision);
snapshot.launch_failures = launch_failures.to_vec();
snapshot.workspaces = workspaces
.iter()
.map(|workspace| crate::server::ViewerWorkspace {
id: workspace.id.clone(),
name: workspace.name.clone(),
})
.collect();
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
for profile in &mut snapshot.profiles {
let Some(quota) = quotas.get(&profile.id) else {
continue;
};
profile.quota = Some(ViewerQuota {
summary: quota.compact(),
windows: quota
.windows
.iter()
.map(|window| crate::server::ViewerQuotaWindow {
label: window.label.clone(),
percent_used: window
.remaining_percent
.map(|left| 100_u8.saturating_sub(left)),
resets_at: window.resets.clone(),
projects_exhaustion_before_reset: crate::quota::projects_exhaustion(
window,
quota.refreshed_at_epoch_seconds,
),
})
.collect(),
resets_at: quota
.windows
.iter()
.find_map(|window| window.resets.clone()),
stale: now.saturating_sub(quota.refreshed_at_epoch_seconds)
> QUOTA_STALE_AFTER.as_secs(),
refreshed_at_epoch_seconds: quota.refreshed_at_epoch_seconds,
has_error: quota.error.is_some(),
});
}
for session in &mut snapshot.sessions {
session.move_recovery = move_recoveries.get(&session.id).cloned();
if let Some(record) = controller.state.sessions.get(&session.id)
&& let Some(source) = project_sources.source(record, &controller.config)
{
session.set_project_source(source);
}
session.last_activity_at_ms = materialized_activity.get(&session.id).copied().flatten();
session.queued_prompts = queued_prompts
.get(&session.id)
.into_iter()
.flatten()
.map(|prompt| ViewerQueuedPrompt {
id: prompt.id.clone(),
text: prompt.text.clone(),
created_at: prompt.created_at_ms.to_string(),
})
.collect();
session.active_user_shells = active_user_shells
.get(&session.id)
.into_iter()
.flatten()
.map(|shell| ViewerUserShell {
id: shell.command_id.clone(),
command: shell.command.clone(),
started_at_ms: shell.started_at_ms,
})
.collect();
session.background_tasks = operational
.get(&session.id)
.map(|state| {
state
.background_commands
.iter()
.map(|task| ViewerBackgroundTask {
id: task.id.clone(),
command: task.command.clone(),
started_at_ms: task.started_at_ms,
can_stop: task.can_stop,
})
.collect()
})
.unwrap_or_default();
session.pending_elicitations = pending_elicitations
.get(&session.id)
.cloned()
.unwrap_or_default();
session.prompt_images_supported = prompt_images.contains(&session.id);
session.operation = operations.get(&session.id).cloned();
if let Some(operation) = session.operation.as_ref()
&& operation.kind.transition_kind().is_some()
{
session.transitioning = true;
}
let live = operational.get(&session.id);
let facts = live.map(|state| {
mj_core::acp::AcpSessionFacts::from_operational(
controller
.state
.sessions
.get(&session.id)
.map_or(mj_core::config::HarnessKind::Codex, |record| {
record.harness_kind
}),
&state.config,
&state.config_options,
state.modes.as_ref(),
)
});
if let Some(state) = live {
session.latest_event_ordinal = state.latest_ordinal;
session.chat_phase = match state.execution {
mj_core::relay::RelayExecutionState::Idle => crate::server::ViewerChatPhase::Idle,
mj_core::relay::RelayExecutionState::Running => {
crate::server::ViewerChatPhase::Running
}
mj_core::relay::RelayExecutionState::Closing => {
crate::server::ViewerChatPhase::Closing
}
mj_core::relay::RelayExecutionState::Closed => {
crate::server::ViewerChatPhase::Closed
}
};
session.config_options = crate::server::viewer_config_options(
&state.config_options,
facts
.as_ref()
.expect("live operational state always has ACP session facts"),
);
let turn_started_at_ms = state
.active_prompt
.as_ref()
.map(|prompt| prompt.started_at_ms)
.or_else(|| state.harness_turn.map(|turn| turn.started_at_ms));
let turn_started_at = turn_started_at_ms
.and_then(|started_at_ms| u64::try_from(started_at_ms).ok())
.map(|started_at_ms| started_at_ms / 1_000);
session.capacity_retry = state.capacity_retry.clone();
let activity = mj_client::usage_format::SessionActivity::of(state);
let activity_details =
activity.details(turn_started_at_ms, state.current_step_started_at_ms);
session.activity_details = Some(viewer_activity_details(&activity_details));
session.is_idle = controller
.state
.sessions
.get(&session.id)
.is_some_and(|record| record.state == mj_core::state::SessionState::Running)
&& session.operation.is_none()
&& activity.is_idle(turn_started_at);
session.activity = mj_client::usage_format::format_activity_columns(
now,
turn_started_at,
state
.current_step_started_at_ms
.and_then(|value| u64::try_from(value).ok()),
&activity,
)
.join(" ")
.trim()
.to_owned();
}
session.plan_mode_active = facts
.as_ref()
.filter(|facts| facts.supports_plan_mode())
.map(mj_core::acp::AcpSessionFacts::plan_mode_active);
session.turn_review = reviews
.get(&session.id)
.map(crate::server::ViewerTurnReview::from_runtime);
session.capabilities =
session_capabilities(session, live, operations.get(&session.id), facts.as_ref());
session.available_commands = phone_commands(session, live);
if let Some(transcript) = conversations.get(&session.id) {
session.conversation_available = true;
if !session.transitioning {
let mut lines = transcript
.entries
.iter()
.flat_map(|entry| {
entry
.lines
.iter()
.enumerate()
.filter_map(move |(index, line)| {
let line = line.trim();
(!line.is_empty()).then(|| {
if index == 0 {
format!("{}: {line}", entry.label)
} else {
line.to_owned()
}
})
})
})
.collect::<Vec<_>>();
session.preview = lines.split_off(lines.len().saturating_sub(4));
}
}
session.capabilities.open = session.conversation_available && !session.transitioning;
}
snapshot.capacity = capacity.to_vec();
snapshot
}
fn viewer_activity_details(
details: &mj_client::usage_format::SessionActivityDetails,
) -> ViewerActivityDetails {
ViewerActivityDetails {
kind: match details.kind {
mj_client::usage_format::SessionActivityKind::Turn => ViewerActivityKind::Turn,
mj_client::usage_format::SessionActivityKind::Step => ViewerActivityKind::Step,
mj_client::usage_format::SessionActivityKind::Background => {
ViewerActivityKind::Background
}
mj_client::usage_format::SessionActivityKind::Idle => ViewerActivityKind::Idle,
mj_client::usage_format::SessionActivityKind::Lifecycle
| mj_client::usage_format::SessionActivityKind::Goal => ViewerActivityKind::Lifecycle,
},
turn_started_at_ms: details.turn_started_at_ms,
step_started_at_ms: details.step_started_at_ms,
background_started_at_ms: details.background_started_at_ms,
idle_since_ms: details.idle_since_ms,
label: details.label.clone(),
}
}
async fn load_materialized_activity(
controller: &Controller,
) -> Result<std::collections::BTreeMap<String, Option<i64>>> {
let session_ids = controller
.state
.sessions
.keys()
.cloned()
.collect::<Vec<_>>();
tokio::task::spawn_blocking(move || {
session_ids
.into_iter()
.map(|session_id| {
crate::database::load_materialized_session_summary(&session_id).map(|summary| {
(
session_id,
summary.and_then(|summary| summary.last_activity_at_ms),
)
})
})
.collect()
})
.await
.context("materialized activity startup task failed")?
}
#[cfg(test)]
mod tests {
use super::*;
use crate::pollers::QUOTA_REFRESH_INTERVAL;
use std::collections::BTreeMap;
use agent_client_protocol::schema::v1::{
SessionConfigOption, SessionConfigOptionCategory, SessionConfigSelectOption,
SessionConfigSelectOptions,
};
use mj_core::config::{
CONFIG_VERSION, Config, HarnessKind, ProjectBundle, ProjectRepository, TargetTemplate,
};
use mj_core::state::SessionState;
#[test]
fn viewer_config_options_publish_current_advertised_values() {
let make_options = |model, effort| {
vec![
SessionConfigOption::select(
"model_selector",
"Model",
model,
SessionConfigSelectOptions::Ungrouped(vec![
SessionConfigSelectOption::new("sonnet", "Claude Sonnet"),
SessionConfigSelectOption::new("opus", "Claude Opus"),
]),
)
.category(SessionConfigOptionCategory::Model),
SessionConfigOption::select(
"reasoning_effort",
"Effort",
effort,
SessionConfigSelectOptions::Ungrouped(vec![
SessionConfigSelectOption::new("high", "High"),
SessionConfigSelectOption::new("max", "Maximum"),
]),
),
]
};
let options = make_options("sonnet", "high");
let defaults = mj_core::acp::AcpSessionFacts::from_operational(
HarnessKind::Claude,
&BTreeMap::new(),
&options,
None,
);
let projected = crate::server::viewer_config_options(&options, &defaults);
assert_eq!(
projected
.iter()
.map(|option| (option.key.as_str(), option.current.as_deref()))
.collect::<Vec<_>>(),
[("model", Some("sonnet")), ("effort", Some("high"))]
);
let options = make_options("opus", "max");
let updated = mj_core::acp::AcpSessionFacts::from_operational(
HarnessKind::Claude,
&BTreeMap::new(),
&options,
None,
);
let projected = crate::server::viewer_config_options(&options, &updated);
assert_eq!(projected[0].current.as_deref(), Some("opus"));
assert_eq!(projected[1].current.as_deref(), Some("max"));
}
#[tokio::test]
async fn explicit_tls_takes_precedence_over_tailscale_detection() {
let resolved = resolve_server_args(
ServerArgs {
bind: "0.0.0.0:4443".into(),
tailscale_detect: true,
tls_cert: Some(PathBuf::from("configured-cert.pem")),
tls_key: Some(PathBuf::from("configured-key.pem")),
},
tokio_util::sync::CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(resolved.bind, "0.0.0.0:4443".parse().unwrap());
assert_eq!(resolved.viewer_url, "https://0.0.0.0:4443");
assert_eq!(
resolved.tls_files,
Some((
PathBuf::from("configured-cert.pem"),
PathBuf::from("configured-key.pem")
))
);
assert!(resolved.tailscale.is_none());
assert!(resolved.fallback_reason.is_none());
}
#[tokio::test]
async fn disabling_detection_keeps_the_viewer_on_loopback() {
let resolved = resolve_server_args(
ServerArgs {
bind: "127.0.0.1:4765".into(),
tailscale_detect: false,
tls_cert: None,
tls_key: None,
},
tokio_util::sync::CancellationToken::new(),
)
.await
.unwrap();
assert_eq!(resolved.bind, "127.0.0.1:4765".parse().unwrap());
assert_eq!(resolved.viewer_url, "http://127.0.0.1:4765");
assert!(
resolved
.fallback_reason
.unwrap()
.contains("detection is disabled")
);
}
fn bare_preflight_config() -> Config {
let mut config = Config::default();
config
.targets
.insert("raw".into(), TargetTemplate::LocalBare);
config
}
#[test]
fn move_recovery_projection_exposes_safe_retry_settings_only() {
let operation = mj_core::state::MoveOperation {
source_checkpoint_only: false,
operation_id: "move-1".into(),
selection: mj_core::state::MoveSelection {
clear_resource_allocation: true,
session_id: "session-1".into(),
profile_id: Some("destination-profile".into()),
target_template_id: Some("destination-target".into()),
additional_mounts: Some(vec![crate::targets::AdditionalMount {
source: "/destination/source".into(),
destination: "/destination/target".into(),
read_only: true,
}]),
resource_allocation: None,
},
source_profile_id: "source-profile".into(),
source_target_template_id: "source-target".into(),
source_target: None,
source_native_session_id: Some("private-native-id".into()),
source_additional_mounts: vec![crate::targets::AdditionalMount {
source: "/source/source".into(),
destination: "/source/target".into(),
read_only: false,
}],
source_resource_allocation: Some(
mj_core::state::SessionResourceAllocation::Container {
cpus: 2,
memory_bytes: 4096,
},
),
destination_target: None,
destination_native_session_id: None,
destination_store_id: None,
configuration_fingerprint: "private-fingerprint".into(),
checkpoint: None,
recovery_session: None,
queue: mj_core::state::ResumeQueueDisposition::Start,
phase: mj_core::state::MovePhase::Cancelled,
queue_admission_started: false,
queue_admission_finished: false,
cancellation_requested: true,
created_at: "now".into(),
updated_at: "now".into(),
error: Some("private path and token".into()),
};
let recovery = ViewerMoveRecovery::from_operation(&operation).unwrap();
assert_eq!(recovery.phase, "cancelled");
assert_eq!(recovery.source_profile_id, "source-profile");
assert!(recovery.clear_resource_allocation);
assert_eq!(recovery.source_additional_mounts.len(), 1);
assert!(recovery.source_resource_allocation.is_some());
assert_eq!(recovery.destination_additional_mounts.len(), 1);
assert!(recovery.destination_resource_allocation.is_none());
let json = serde_json::to_string(&recovery).unwrap();
assert!(!json.contains("private-native-id"));
assert!(!json.contains("private-fingerprint"));
assert!(!json.contains("private path and token"));
}
#[test]
fn new_preflight_rejects_a_bare_project_without_a_git_head() {
let error = run_new_preflight(
bare_preflight_config(),
"hel".into(),
"raw".into(),
Some(PathBuf::from("/definitely/not/a/project")),
)
.expect_err("a missing project directory must fail preflight");
assert!(
error
.to_string()
.contains("project directory does not exist or is not a directory")
);
}
#[test]
fn new_preflight_accepts_a_git_project_for_a_bare_target() {
let directory = std::env::current_dir().expect("the test has a working directory");
let answer = run_new_preflight(
bare_preflight_config(),
"hel".into(),
"raw".into(),
Some(directory.clone()),
)
.expect("the repository running the test has a valid Git HEAD");
assert!(answer.dirty_repositories.is_empty());
assert_eq!(answer.project_directory, Some(directory));
}
#[test]
fn new_preflight_requires_network_sources_for_isolated_targets() {
let mut config = Config::default();
config.targets.insert(
"podman".into(),
TargetTemplate::LocalPodman {
container: mj_core::config::ContainerTemplate {
image: "test-image".into(),
pull_policy: Default::default(),
platform: None,
cpus: None,
memory: None,
environment: Default::default(),
workspace_storage: Default::default(),
},
},
);
config.bundles.insert(
"hel".into(),
ProjectBundle {
primary_repo: "hel".into(),
repositories: vec![ProjectRepository {
id: "hel".into(),
github: None,
local: Some(PathBuf::from("/definitely/not/a/repository")),
destination: "hel".into(),
git_ref: None,
}],
},
);
let error = run_new_preflight(config, "hel".into(), "podman".into(), None)
.expect_err("an isolated bundle cannot use a local source");
assert!(error.to_string().contains("repository"));
}
#[test]
fn a_phone_prompt_becomes_its_text_then_its_images() {
use agent_client_protocol::schema::v1::ContentBlock;
let image = |data: &str| crate::server::ViewerPromptImage {
attachment: None,
data_base64: data.into(),
mime_type: "image/png".into(),
width: 32,
height: 24,
};
let blocks = phone_prompt_blocks(
"look at this".into(),
vec![image("aW1hZ2U="), image("c2Vjb25k")],
);
let ContentBlock::Text(text) = &blocks[0] else {
panic!("the prompt leads with its text");
};
assert_eq!(text.text, "look at this");
let ContentBlock::Image(first) = &blocks[1] else {
panic!("each attachment travels as an image block");
};
assert_eq!(first.data, "aW1hZ2U=");
assert_eq!(first.mime_type, "image/png");
assert!(matches!(blocks[2], ContentBlock::Image(_)));
assert_eq!(blocks.len(), 3);
let images_only = phone_prompt_blocks(String::new(), vec![image("aW1hZ2U=")]);
assert_eq!(images_only.len(), 1);
assert!(matches!(images_only[0], ContentBlock::Image(_)));
}
#[test]
fn image_prompts_are_offered_only_after_the_agent_advertises_them() {
use agent_client_protocol::schema::v1::AgentCapabilities;
use mj_core::relay::{RelayExecutionState, RelayOperationalState};
let operational = |agent_capabilities| RelayOperationalState {
goal: Default::default(),
capacity_retry: None,
activity_turn_started_at_ms: None,
session_id: "session-1".into(),
store_id: None,
idle_since_ms: None,
execution: RelayExecutionState::Idle,
latest_ordinal: 0,
latest_digest: String::new(),
acknowledged_through: 0,
acknowledged_digest: String::new(),
recovery_floor_ordinal: 0,
recovery_floor_digest: String::new(),
native_session_id: None,
native_continuity_lost: false,
checkpoint_only: false,
acp_ready: None,
agent_capabilities,
agent_info: None,
steering_supported: None,
config_options: Vec::new(),
modes: None,
available_commands: Vec::new(),
config: std::collections::BTreeMap::new(),
active_prompt: None,
queued_prompts: Vec::new(),
active_user_shells: Vec::new(),
active_agent_terminals: Vec::new(),
checkpoint_barrier: None,
checkpoint_ready: None,
last_acp_activity_at_ms: None,
current_step_started_at_ms: None,
foreground_tool_started_at_ms: None,
harness_turn: None,
last_harness_turn_started_ordinal: None,
background_commands: Vec::new(),
background_work_known: None,
};
assert!(!agent_accepts_prompt_images(&operational(None)));
assert!(!agent_accepts_prompt_images(&operational(Some(Box::new(
AgentCapabilities::default()
)))));
let mut capabilities = AgentCapabilities::default();
capabilities.prompt_capabilities.image = true;
assert!(agent_accepts_prompt_images(&operational(Some(Box::new(
capabilities
)))));
}
#[test]
fn phone_snapshot_projects_capability_gated_and_agent_commands_with_provenance() {
use agent_client_protocol::schema::v1::{
AvailableCommand, AvailableCommandInput, SessionMode, SessionModeState,
UnstructuredCommandInput,
};
use mj_core::relay::{RelayExecutionState, RelayOperationalState};
use crate::server::ViewerCommandSource;
let mut controller = controller_with_profiles(&["claude"]);
let mut record = phone_session("session-1", 0);
record.harness_kind = HarnessKind::Claude;
record.last_profile = "claude".into();
record.state = SessionState::Running;
controller.state.sessions.insert(record.id.clone(), record);
let operational = RelayOperationalState {
goal: Default::default(),
capacity_retry: None,
activity_turn_started_at_ms: None,
session_id: "session-1".into(),
store_id: None,
idle_since_ms: None,
execution: RelayExecutionState::Idle,
latest_ordinal: 0,
latest_digest: String::new(),
acknowledged_through: 0,
acknowledged_digest: String::new(),
recovery_floor_ordinal: 0,
recovery_floor_digest: String::new(),
native_session_id: None,
native_continuity_lost: false,
checkpoint_only: false,
acp_ready: None,
agent_capabilities: None,
agent_info: None,
steering_supported: None,
config_options: Vec::new(),
modes: Some(SessionModeState::new(
"default",
vec![
SessionMode::new("default", "Default"),
SessionMode::new("plan", "Plan"),
],
)),
available_commands: vec![
AvailableCommand::new("inspect", " Inspect the workspace ").input(
AvailableCommandInput::Unstructured(UnstructuredCommandInput::new(" query ")),
),
AvailableCommand::new("Review", "agent collision"),
AvailableCommand::new("INSPECT", "duplicate agent command"),
],
config: std::collections::BTreeMap::new(),
active_prompt: None,
queued_prompts: Vec::new(),
active_user_shells: Vec::new(),
active_agent_terminals: Vec::new(),
checkpoint_barrier: None,
checkpoint_ready: None,
last_acp_activity_at_ms: None,
current_step_started_at_ms: None,
foreground_tool_started_at_ms: None,
harness_turn: None,
last_harness_turn_started_ordinal: None,
background_commands: Vec::new(),
background_work_known: None,
};
let mut operational = std::collections::BTreeMap::from([("session-1".into(), operational)]);
let materialized_activity =
std::collections::BTreeMap::from([("session-1".into(), Some(7_777_i64))]);
let project = |operational: &std::collections::BTreeMap<String, RelayOperationalState>| {
viewer_snapshot(
&controller,
&[],
&std::collections::BTreeMap::new(),
&PhoneSessionViews {
conversations: &std::collections::BTreeMap::new(),
queued_prompts: &std::collections::BTreeMap::new(),
active_user_shells: &std::collections::BTreeMap::new(),
pending_elicitations: &std::collections::BTreeMap::new(),
prompt_images: &std::collections::BTreeSet::new(),
operational,
materialized_activity: &materialized_activity,
project_sources: &PhoneProjectSources::default(),
operations: &std::collections::BTreeMap::new(),
move_recoveries: &std::collections::BTreeMap::new(),
capacity: &[],
launch_failures: &[],
reviews: &std::collections::BTreeMap::new(),
},
1,
)
};
let snapshot = project(&operational);
let session = &snapshot.sessions[0];
assert_eq!(session.display_location, "podman");
assert_eq!(session.last_activity_at_ms, Some(7_777));
assert_eq!(
session.activity_details,
Some(crate::server::ViewerActivityDetails {
kind: ViewerActivityKind::Idle,
turn_started_at_ms: None,
step_started_at_ms: None,
background_started_at_ms: None,
idle_since_ms: None,
label: None,
})
);
assert!(session.capabilities.prompt);
assert!(session.capabilities.set_plan_mode);
assert_eq!(
session
.available_commands
.iter()
.map(|command| (command.name.as_str(), command.source))
.collect::<Vec<_>>(),
vec![
("help", ViewerCommandSource::Mj),
("detach", ViewerCommandSource::Mj),
("plan", ViewerCommandSource::Mj),
("implement", ViewerCommandSource::Mj),
("review", ViewerCommandSource::Mj),
("inspect", ViewerCommandSource::Agent),
]
);
let inspect = session.available_commands.last().unwrap();
assert_eq!(inspect.description, "Inspect the workspace");
assert_eq!(inspect.argument.as_deref(), Some("query"));
assert!(session.is_idle);
assert_eq!(session.activity, "[idle]");
let state = operational.get_mut("session-1").unwrap();
state
.background_commands
.push(mj_core::relay::BackgroundCommand {
id: "background-1".into(),
started_at_ms: 1_000,
command: "background check".into(),
can_stop: true,
});
let background = project(&operational);
assert_eq!(
background.sessions[0].background_tasks,
vec![ViewerBackgroundTask {
id: "background-1".into(),
command: "background check".into(),
started_at_ms: 1_000,
can_stop: true,
}]
);
assert!(!background.sessions[0].is_idle);
assert!(background.sessions[0].activity.starts_with("BG "));
let state = operational.get_mut("session-1").unwrap();
state.background_commands.clear();
state.execution = RelayExecutionState::Running;
let stale_running = project(&operational);
assert!(stale_running.sessions[0].is_idle);
operational
.get_mut("session-1")
.unwrap()
.current_step_started_at_ms = Some(1_000);
let running = project(&operational);
assert!(!running.sessions[0].is_idle);
assert_ne!(running.sessions[0].activity, "[idle]");
operational.get_mut("session-1").unwrap().execution = RelayExecutionState::Idle;
let idle_again = project(&operational);
assert!(idle_again.sessions[0].is_idle);
assert_eq!(idle_again.sessions[0].activity, "[idle]");
let unknown = project(&std::collections::BTreeMap::new());
assert!(!unknown.sessions[0].is_idle);
assert!(unknown.sessions[0].activity.is_empty());
assert!(unknown.sessions[0].activity_details.is_none());
}
#[test]
fn tailscale_listener_preserves_the_configured_port() {
assert_eq!(
tailscale_bind("127.0.0.1:4765".parse().unwrap()),
"0.0.0.0:4765".parse().unwrap()
);
}
fn controller_with_profiles(ids: &[&str]) -> Controller {
Controller {
config: Config {
subagents: Default::default(),
version: CONFIG_VERSION,
sessions_side: Default::default(),
advanced: Default::default(),
show_stopped_sessions: false,
newer_config_version: None,
spinner: Default::default(),
theme: Default::default(),
phone: Default::default(),
review: Default::default(),
legacy_startup: (),
profiles: ids
.iter()
.map(|id| {
(
(*id).to_owned(),
HarnessProfile {
enabled: true,
context_window_bytes: None,
kind: HarnessKind::Codex,
home: PathBuf::from("/home/agent").join(id),
environment: std::collections::BTreeMap::new(),
},
)
})
.collect(),
bundles: std::collections::BTreeMap::new(),
targets: std::collections::BTreeMap::new(),
},
state: State::default(),
}
}
fn snapshot_with_project_sources(
controller: &Controller,
sources: &PhoneProjectSources,
) -> ViewerSnapshot {
viewer_snapshot(
controller,
&[],
&Default::default(),
&PhoneSessionViews {
conversations: &Default::default(),
queued_prompts: &Default::default(),
active_user_shells: &Default::default(),
pending_elicitations: &Default::default(),
prompt_images: &Default::default(),
operational: &Default::default(),
materialized_activity: &Default::default(),
project_sources: sources,
operations: &Default::default(),
move_recoveries: &Default::default(),
capacity: &[],
launch_failures: &[],
reviews: &Default::default(),
},
1,
)
}
#[tokio::test]
async fn phone_projects_resolve_origins_and_discard_results_after_location_changes() {
let root = tempfile::tempdir().unwrap();
let mut controller = controller_with_profiles(&["codex"]);
controller
.config
.targets
.insert("local".into(), TargetTemplate::LocalBare);
for (id, origin) in [
("first-checkout", "git@github.com:BrokkAi/hel.git"),
("second-checkout", "https://github.com/BrokkAi/hel.git"),
] {
let directory = root.path().join(id);
std::fs::create_dir(&directory).unwrap();
for args in [vec!["init"], vec!["remote", "add", "origin", origin]] {
let command = crate::targets::CommandSpec::new(
"git",
["-C".to_owned(), directory.to_string_lossy().into_owned()]
.into_iter()
.chain(args.into_iter().map(str::to_owned)),
);
assert_eq!(ProcessExecutor.execute(&command).unwrap().status, 0);
}
let mut record = phone_session(id, 0);
record.project_directory = Some(directory);
record.target_template_id = "local".into();
controller.state.sessions.insert(id.into(), record);
}
let mut sources = PhoneProjectSources::default();
sources.synchronize(&controller);
assert_eq!(sources.jobs.len(), 2);
tokio::time::timeout(Duration::from_secs(10), async {
while let Some(result) = sources.jobs.join_next().await {
sources.complete(result.unwrap());
}
})
.await
.unwrap();
let snapshot = snapshot_with_project_sources(&controller, &sources);
assert_eq!(
snapshot.sessions[0].project_key,
snapshot.sessions[1].project_key
);
assert!(
snapshot
.sessions
.iter()
.all(|session| session.project_label == "hel")
);
assert!(
!serde_json::to_string(&snapshot)
.unwrap()
.contains(&root.path().to_string_lossy().to_string())
);
sources.synchronize(&controller);
assert!(
sources.jobs.is_empty(),
"unchanged inputs reuse the resolved origin"
);
let previous = &sources.entries["first-checkout"];
let late = ProjectSourceResolved {
cancelled: previous.cancelled.clone(),
session_id: "first-checkout".into(),
key: previous.key.clone(),
result: Ok(ProjectSourceIdentity::git_remote("old/wrong").unwrap()),
};
controller
.state
.sessions
.get_mut("first-checkout")
.unwrap()
.project_directory = None;
let snapshot = snapshot_with_project_sources(&controller, &sources);
assert_ne!(
snapshot.sessions[0].project_key,
snapshot.sessions[1].project_key
);
sources.synchronize(&controller);
assert!(late.cancelled.load(Ordering::Acquire));
sources.complete(late);
assert!(!sources.entries.contains_key("first-checkout"));
assert_eq!(snapshot.sessions[0].project_label, "project");
}
#[test]
fn capacity_target_publication_skips_unchanged_targets_and_preserves_readings() {
let mut controller = controller_with_profiles(&[]);
controller
.config
.targets
.insert("raw".into(), TargetTemplate::LocalBare);
let (targets_tx, mut targets_rx) = tokio::sync::watch::channel(Vec::new());
let mut state = std::collections::BTreeMap::new();
publish_capacity_targets(&controller, &targets_tx, &mut state);
assert!(targets_rx.has_changed().expect("target sender is alive"));
assert_eq!(targets_rx.borrow_and_update().len(), 1);
let usage = crate::targets::DeploymentCapacityUsage {
cpu_percent: Some(37),
memory_used_bytes: 3,
memory_total_bytes: 4,
logical_cores: 8,
disk_total_bytes: Some(5),
};
let local = state.get_mut("local").expect("local capacity state");
local.usage = Some(usage.clone());
local.on_demand = true;
local.sampled_at_epoch_seconds = Some(42);
local.refreshing = false;
publish_capacity_targets(&controller, &targets_tx, &mut state);
assert!(!targets_rx.has_changed().expect("target sender is alive"));
let local_capacity = viewer_capacity(&state)
.into_iter()
.find(|capacity| capacity.id == "local")
.expect("local viewer capacity");
assert_eq!(local_capacity.cpu_percent, usage.cpu_percent);
assert_eq!(
local_capacity.memory_used_bytes,
Some(usage.memory_used_bytes)
);
assert_eq!(local_capacity.logical_cores, Some(usage.logical_cores));
assert_eq!(local_capacity.sampled_at_epoch_seconds, Some(42));
controller
.config
.targets
.insert("second-local".into(), TargetTemplate::LocalBare);
publish_capacity_targets(&controller, &targets_tx, &mut state);
assert!(targets_rx.has_changed().expect("target sender is alive"));
assert_eq!(targets_rx.borrow_and_update().len(), 1);
assert_eq!(state["local"].usage, Some(usage.clone()));
controller.config.targets.insert(
"fleet".into(),
TargetTemplate::AwsEc2 {
aws_profile: None,
region: "us-east-1".into(),
launch_template: "hel-runson".into(),
launch_template_version: None,
ssh_user: "ubuntu".into(),
address_source: Default::default(),
identity_file: None,
ssh_args: Vec::new(),
},
);
publish_capacity_targets(&controller, &targets_tx, &mut state);
assert!(targets_rx.has_changed().expect("target sender is alive"));
assert_eq!(targets_rx.borrow_and_update().len(), 2);
assert!(state.contains_key("aws:fleet"));
assert_eq!(state["local"].usage, Some(usage.clone()));
controller.config.targets.remove("fleet");
publish_capacity_targets(&controller, &targets_tx, &mut state);
assert!(targets_rx.has_changed().expect("target sender is alive"));
assert_eq!(targets_rx.borrow_and_update().len(), 1);
assert!(!state.contains_key("aws:fleet"));
assert_eq!(state["local"].usage, Some(usage));
}
fn prompt_action() -> ControllerAction {
ControllerAction::Prompt {
session_id: "session-1".into(),
text: "ship it".into(),
images: Vec::new(),
}
}
fn new_action() -> ControllerAction {
ControllerAction::New {
mjolnir_subagents: None,
create_managed_worktree: None,
workspace_id: String::new(),
profile_id: "codex".into(),
bundle_id: "project".into(),
target_id: "podman".into(),
title: Some("Phone launch".into()),
project_directory: None,
dirty_ack: Vec::new(),
}
}
fn phone_session(id: &str, viewed_through_event_ordinal: u64) -> SessionRecord {
SessionRecord {
mjolnir_subagents: None,
create_managed_worktree: None,
workspace_id: mj_core::workspace::DEFAULT_WORKSPACE_ID.to_owned(),
archived: false,
container_cpus: None,
container_memory: None,
id: id.into(),
title: "Phone launch".into(),
harness_kind: mj_core::config::HarnessKind::Codex,
last_profile: "codex".into(),
bundle_id: "project".into(),
project_directory: None,
managed_worktree: None,
target_template_id: "podman".into(),
resource_allocation: None,
additional_mounts: Vec::new(),
state: SessionState::Provisioning,
target: None,
native_session_id: None,
acp_session_title: None,
session_title_override: Some("Phone launch".into()),
created_at: "2026-08-14T00:00:00Z".into(),
updated_at: "2026-08-14T00:00:00Z".into(),
viewed_through_event_ordinal,
draft_input: String::new(),
last_error: None,
last_checkpoint_error: None,
checkpoint: None,
}
}
#[tokio::test]
async fn controller_reload_does_not_block_the_phone_control_loop() {
let (completed_tx, mut completed_rx) = tokio::sync::mpsc::unbounded_channel();
let (started_tx, started_rx) = std::sync::mpsc::channel();
let (release_tx, release_rx) = std::sync::mpsc::channel();
let loaded = controller_with_profiles(&[]);
let release = std::thread::spawn(move || {
started_rx.recv().unwrap();
std::thread::sleep(Duration::from_millis(250));
release_tx.send(()).unwrap();
});
let started = Instant::now();
spawn_controller_reload_with(completed_tx, move || {
started_tx.send(()).unwrap();
release_rx.recv().unwrap();
Ok(loaded)
});
assert!(
started.elapsed() < Duration::from_millis(100),
"scheduling a controller reload occupied the control loop for {:?}",
started.elapsed()
);
let completed = tokio::time::timeout(Duration::from_secs(1), completed_rx.recv())
.await
.expect("background reload timed out")
.expect("background reload channel closed");
assert!(completed.result.is_ok());
release.join().unwrap();
}
#[test]
fn read_receipt_only_persists_and_refreshes_when_the_cursor_advances() {
let session_id = "0123456789abcdef0123456789abcdef";
let mut state = State::default();
state
.sessions
.insert(session_id.into(), phone_session(session_id, 5));
assert_eq!(
plan_read_receipt(&state, session_id, 5),
ReadReceiptPlan::AlreadyRead
);
assert_eq!(
plan_read_receipt(&state, session_id, 4),
ReadReceiptPlan::AlreadyRead
);
assert_eq!(
plan_read_receipt(&state, "missing", 9),
ReadReceiptPlan::UnknownSession
);
assert_eq!(
plan_read_receipt(&state, session_id, 9),
ReadReceiptPlan::Persist
);
assert!(apply_read_receipt(&mut state, session_id, 9));
assert_eq!(state.sessions[session_id].viewed_through_event_ordinal, 9);
assert!(!apply_read_receipt(&mut state, session_id, 9));
assert!(!apply_read_receipt(&mut state, session_id, 7));
assert!(!apply_read_receipt(&mut state, "missing", 9));
assert_eq!(
plan_read_receipt(&state, session_id, 9),
ReadReceiptPlan::AlreadyRead
);
}
#[tokio::test]
async fn an_admitted_action_answers_its_phone_before_the_work_runs() {
let mut replies = PendingActionReplies::default();
let (reply, answer) = tokio::sync::oneshot::channel();
replies.accept(1, &prompt_action(), reply);
assert_eq!(answer.await.unwrap(), ActionOutcome::accepted());
}
#[tokio::test]
async fn a_new_action_answers_once_its_provisional_session_is_published() {
let mut replies = PendingActionReplies::default();
let (reply, mut answer) = tokio::sync::oneshot::channel();
replies.accept(7, &new_action(), reply);
assert!(
matches!(
answer.try_recv(),
Err(tokio::sync::oneshot::error::TryRecvError::Empty)
),
"a new session has no id to report before it is published"
);
replies.resolve(7, ActionOutcome::accepted());
assert_eq!(answer.await.unwrap(), ActionOutcome::accepted());
}
#[tokio::test]
async fn a_new_action_that_never_publishes_still_answers_its_phone() {
let mut replies = PendingActionReplies::default();
let (reply, answer) = tokio::sync::oneshot::channel();
replies.accept(7, &new_action(), reply);
replies.resolve(7, ActionOutcome::Failed);
assert_eq!(answer.await.unwrap(), ActionOutcome::Failed);
replies.resolve(7, ActionOutcome::accepted());
}
#[test]
fn force_close_is_admitted_like_close() {
let mut active = std::collections::BTreeSet::from(["session-1".to_owned()]);
let force_close = ControllerAction::ForceClose {
session_id: "session-1".into(),
};
assert_eq!(
admit_phone_action(&force_close, MAX_CONCURRENT_PHONE_ACTIONS, &mut active),
Ok(Some("session-1".into()))
);
assert_eq!(
controller_action_session_id(&force_close),
Some("session-1".to_owned())
);
}
#[test]
fn close_is_admitted_while_provisioning_occupies_a_full_action_pool() {
let mut active = std::collections::BTreeSet::from(["session-1".to_owned()]);
let close = ControllerAction::Close {
session_id: "session-1".into(),
};
assert_eq!(
admit_phone_action(&close, MAX_CONCURRENT_PHONE_ACTIONS, &mut active),
Ok(Some("session-1".into()))
);
assert_eq!(
admit_phone_action(&prompt_action(), 0, &mut active),
Err(ActionOutcome::SessionBusy)
);
}
#[test]
fn a_refused_action_reports_the_reason_the_phone_can_act_on() {
let mut active = std::collections::BTreeSet::new();
assert_eq!(
admit_phone_action(&prompt_action(), 0, &mut active),
Ok(Some("session-1".to_owned()))
);
assert_eq!(
admit_phone_action(&prompt_action(), 1, &mut active),
Err(ActionOutcome::SessionBusy)
);
assert_eq!(
admit_phone_action(&new_action(), MAX_CONCURRENT_PHONE_ACTIONS, &mut active),
Err(ActionOutcome::Busy)
);
assert_eq!(active.len(), 1);
assert_eq!(admit_phone_action(&new_action(), 1, &mut active), Ok(None));
}
#[test]
fn a_feed_that_ends_outside_shutdown_names_the_failure() {
assert!(feed_stopped(true, "the session manager stopped").is_none());
let failure = feed_stopped(false, "the session manager stopped").expect("named failure");
assert!(failure.to_string().contains("session manager"));
}
#[test]
fn a_profile_added_while_the_server_runs_reaches_the_quota_refresher() {
let (profiles_tx, profiles_rx) = tokio::sync::watch::channel(QuotaRefreshBatch::default());
let mut published = std::collections::BTreeMap::new();
let mut batch = QuotaRefreshBatch::default();
let controller = controller_with_profiles(&["codex"]);
assert!(republish_quota_profiles(
&controller,
&mut published,
&mut batch,
&profiles_tx
));
assert_eq!(
profiles_rx
.borrow()
.profiles
.iter()
.map(|profile| profile.profile_id.clone())
.collect::<Vec<_>>(),
vec!["codex".to_owned()]
);
let first_generation = profiles_rx.borrow().generation;
assert!(!republish_quota_profiles(
&controller,
&mut published,
&mut batch,
&profiles_tx
));
assert_eq!(profiles_rx.borrow().generation, first_generation);
let grown = controller_with_profiles(&["claude", "codex"]);
assert!(republish_quota_profiles(
&grown,
&mut published,
&mut batch,
&profiles_tx
));
assert_eq!(
profiles_rx
.borrow()
.profiles
.iter()
.map(|profile| profile.profile_id.clone())
.collect::<Vec<_>>(),
vec!["claude".to_owned(), "codex".to_owned()]
);
assert!(profiles_rx.borrow().generation > first_generation);
}
#[test]
fn a_quota_reads_stale_only_once_its_next_refresh_is_overdue() {
let controller = controller_with_profiles(&["codex"]);
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs();
let quota_refreshed = |age: Duration| {
let quotas = std::collections::BTreeMap::from([(
"codex".to_owned(),
ProfileQuota {
profile_id: "codex".into(),
harness: HarnessKind::Codex,
windows: Vec::new(),
extra: None,
error: None,
refreshed_at_epoch_seconds: now - age.as_secs(),
},
)]);
viewer_snapshot(
&controller,
&[],
"as,
&PhoneSessionViews {
conversations: &std::collections::BTreeMap::new(),
queued_prompts: &std::collections::BTreeMap::new(),
active_user_shells: &std::collections::BTreeMap::new(),
pending_elicitations: &std::collections::BTreeMap::new(),
prompt_images: &std::collections::BTreeSet::new(),
operational: &std::collections::BTreeMap::new(),
materialized_activity: &std::collections::BTreeMap::new(),
project_sources: &PhoneProjectSources::default(),
operations: &std::collections::BTreeMap::new(),
move_recoveries: &std::collections::BTreeMap::new(),
capacity: &[],
launch_failures: &[],
reviews: &std::collections::BTreeMap::new(),
},
1,
)
.profiles[0]
.quota
.as_ref()
.expect("the profile carries its quota")
.stale
};
assert!(!quota_refreshed(QUOTA_REFRESH_INTERVAL));
assert!(!quota_refreshed(QUOTA_STALE_AFTER));
assert!(quota_refreshed(QUOTA_STALE_AFTER + Duration::from_secs(1)));
}
#[test]
fn phone_action_capacity_is_bounded() {
assert!(phone_action_capacity_available(
MAX_CONCURRENT_PHONE_ACTIONS - 1
));
assert!(!phone_action_capacity_available(
MAX_CONCURRENT_PHONE_ACTIONS
));
}
#[test]
fn started_phone_session_is_visible_and_mapped_before_provisioning() {
let session_id = "0123456789abcdef0123456789abcdef";
let session = phone_session(session_id, 0);
let mut state = State::default();
let mut active_actions = std::collections::BTreeSet::new();
let mut action_sessions = std::collections::BTreeMap::new();
track_started_phone_session(
&mut state,
&mut active_actions,
&mut action_sessions,
7,
session,
)
.unwrap();
assert_eq!(state.sessions[session_id].state, SessionState::Provisioning);
assert_eq!(state.sessions[session_id].display_title(), "Phone launch");
assert!(active_actions.contains(session_id));
assert_eq!(
action_sessions.get(&7).map(String::as_str),
Some(session_id)
);
}
#[test]
fn failed_launch_notice_survives_session_rollback_and_history_is_bounded() {
let controller = controller_with_profiles(&["codex"]);
let mut failures = Vec::new();
for index in 0..20 {
record_launch_failure(
&mut failures,
index,
format!("workspace-{index}"),
Some(format!("session-{index}")),
Some(format!("worker bootstrap failed for {index}")),
);
}
let snapshot = viewer_snapshot(
&controller,
&[],
&std::collections::BTreeMap::new(),
&PhoneSessionViews {
conversations: &std::collections::BTreeMap::new(),
queued_prompts: &std::collections::BTreeMap::new(),
active_user_shells: &std::collections::BTreeMap::new(),
pending_elicitations: &std::collections::BTreeMap::new(),
prompt_images: &std::collections::BTreeSet::new(),
operational: &std::collections::BTreeMap::new(),
materialized_activity: &std::collections::BTreeMap::new(),
project_sources: &PhoneProjectSources::default(),
operations: &std::collections::BTreeMap::new(),
move_recoveries: &std::collections::BTreeMap::new(),
capacity: &[],
launch_failures: &failures,
reviews: &std::collections::BTreeMap::new(),
},
1,
);
assert!(snapshot.sessions.is_empty());
assert_eq!(failures.len(), 16);
assert_eq!(failures[0].workspace_id, "workspace-4");
assert_eq!(failures[15].workspace_id, "workspace-19");
assert_ne!(failures[0].id, failures[1].id);
assert_eq!(
failures[15].session_id.as_deref(),
Some("session-19"),
"a wait on that session has to be able to recognize its own launch failure"
);
let json = serde_json::to_value(snapshot).unwrap();
assert_eq!(json["launch_failures"][15]["workspace_id"], "workspace-19");
assert_eq!(
json["launch_failures"][15]["error"], "worker bootstrap failed for 19",
"the recorded failure carries its reason so a client can show it"
);
assert_eq!(json["launch_failures"][15].as_object().unwrap().len(), 4);
}
#[test]
fn a_later_successful_action_clears_a_session_s_recorded_failure() {
let mut pending = std::collections::BTreeMap::new();
record_action_result(&mut pending, Some("session-1"), &Err("relay hiccup".into()));
assert_eq!(
pending.get("session-1").map(String::as_str),
Some("relay hiccup")
);
record_action_result(&mut pending, Some("session-2"), &Ok(()));
record_action_result(&mut pending, Some("session-1"), &Ok(()));
assert!(
pending.is_empty(),
"the overlay has no other expiry, so a stale error would badge the session forever"
);
record_action_result(&mut pending, None, &Err("orphaned".into()));
assert!(pending.is_empty());
}
#[test]
fn phone_cancel_targets_the_matching_background_action() {
let first = PhoneActionControl {
cancelled: Arc::new(AtomicBool::new(false)),
create: None,
};
let second = PhoneActionControl {
cancelled: Arc::new(AtomicBool::new(false)),
create: None,
};
let action_sessions =
std::collections::BTreeMap::from([(1, "session-1".into()), (2, "session-2".into())]);
let cancellations =
std::collections::BTreeMap::from([(1, first.clone()), (2, second.clone())]);
assert!(request_phone_action_cancellation(
"session-2",
&action_sessions,
&cancellations,
));
assert!(!first.cancelled.load(Ordering::Acquire));
assert!(second.cancelled.load(Ordering::Acquire));
assert!(!request_phone_action_cancellation(
"missing",
&action_sessions,
&cancellations,
));
}
#[test]
fn phone_new_cancel_and_running_commit_have_one_atomic_winner() {
for _ in 0..100 {
let create = CreateSessionControl::default();
let control = PhoneActionControl {
cancelled: create.cancelled.clone(),
create: Some(create),
};
let cancelling = control.clone();
let committing = control.clone();
let (cancelled, committed) = std::thread::scope(|scope| {
let cancel = scope.spawn(move || cancelling.request_cancel());
let commit = scope.spawn(move || committing.grant_new_commit());
(cancel.join().unwrap(), commit.join().unwrap())
});
assert_ne!(cancelled, committed);
assert_eq!(control.cancelled.load(Ordering::Acquire), cancelled);
assert!(!control.request_cancel());
assert!(!control.grant_new_commit());
}
}
fn materialized_at(ordinal: u64) -> MaterializedSession {
let mut materialized = MaterializedSession::empty("session-1");
materialized.applied_event_ordinal = ordinal;
materialized.applied_event_digest = format!("digest-{ordinal}");
materialized
}
#[tokio::test(flavor = "multi_thread", worker_threads = 1)]
async fn browser_projection_enqueue_does_not_wait_for_a_blocking_permit() {
let (results, mut completed) = tokio::sync::mpsc::channel(1);
let shutdown = tokio_util::sync::CancellationToken::new();
let mut dispatcher = ConversationProjectionDispatcher::with_permits(results, shutdown, 0);
dispatcher.enqueue(materialized_at(1));
let (progress, progress_done) = tokio::sync::oneshot::channel();
tokio::spawn(async move {
progress.send(()).expect("control task is still alive");
});
tokio::time::timeout(Duration::from_millis(100), progress_done)
.await
.expect("control progress was starved by transcript projection")
.expect("progress task did not report");
assert!(completed.try_recv().is_err());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 1)]
async fn browser_projection_keeps_only_the_newest_pending_snapshot() {
let (results, mut completed) = tokio::sync::mpsc::channel(4);
let shutdown = tokio_util::sync::CancellationToken::new();
let mut dispatcher = ConversationProjectionDispatcher::with_permits(results, shutdown, 0);
dispatcher.enqueue(materialized_at(1));
dispatcher.enqueue(materialized_at(3));
dispatcher.enqueue(materialized_at(2));
assert_eq!(dispatcher.pending["session-1"].key.ordinal, 3);
dispatcher.permits.add_permits(1);
let first = tokio::time::timeout(Duration::from_secs(1), completed.recv())
.await
.expect("first projection did not complete")
.expect("projection result channel closed");
assert_eq!(first.key.ordinal, 1);
assert!(dispatcher.finish(first, true).is_some());
let second = tokio::time::timeout(Duration::from_secs(1), completed.recv())
.await
.expect("coalesced projection did not complete")
.expect("projection result channel closed");
assert_eq!(second.key.ordinal, 3);
assert!(dispatcher.finish(second, true).is_some());
assert!(dispatcher.pending.is_empty());
assert!(dispatcher.in_flight.is_empty());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 1)]
async fn forgotten_projection_cannot_repopulate_after_a_same_cursor_resume() {
let (results, mut completed) = tokio::sync::mpsc::channel(4);
let shutdown = tokio_util::sync::CancellationToken::new();
let mut dispatcher = ConversationProjectionDispatcher::with_permits(results, shutdown, 0);
let snapshot = materialized_at(7);
dispatcher.enqueue(snapshot.clone());
dispatcher.forget("session-1");
dispatcher.enqueue(snapshot);
dispatcher.permits.add_permits(1);
let old = tokio::time::timeout(Duration::from_secs(1), completed.recv())
.await
.expect("old projection did not complete")
.expect("projection result channel closed");
assert!(dispatcher.finish(old, true).is_none());
let current = tokio::time::timeout(Duration::from_secs(1), completed.recv())
.await
.expect("resumed projection did not complete")
.expect("projection result channel closed");
assert!(dispatcher.finish(current, true).is_some());
assert!(dispatcher.in_flight.is_empty());
}
}