use super::*;
pub(super) struct BackgroundPolicyState {
quiet: bool,
checkpoint_wait: Option<mj_core::activity::CheckpointWait>,
latest_completed_turn_ordinal: Option<u64>,
worker_build: Option<String>,
}
impl RuntimeState {
fn observe_background_policy(
&self,
session: &SessionRecord,
config: &Config,
policy: &BackgroundPolicyState,
) {
if policy.quiet {
self.worker_upgrade_observer
.observe(WorkerUpgradeObservation {
session: session.clone(),
config: config.clone(),
worker_build: policy.worker_build.clone(),
quiet: true,
});
}
self.recovery_observer.observe(RecoveryObservation {
checkpoint_wait: policy.checkpoint_wait,
session: session.clone(),
config: config.clone(),
latest_completed_turn_ordinal: policy.latest_completed_turn_ordinal,
});
}
pub(super) fn refresh_background_policies(&self) {
let mut owner = self.owner();
let records = owner.controller().state.sessions.clone();
let config = owner.controller().config.clone();
owner.background_policies.retain(|session_id, policy| {
let Some(session) = records
.get(session_id)
.filter(|session| session.state.has_live_worker())
else {
return false;
};
self.observe_background_policy(session, &config, policy);
true
});
}
pub async fn reload_controller(&self) -> Result<()> {
let _mutation = self.config_mutation.lock().await;
if self.committed.is_some() {
let config = tokio::task::spawn_blocking(Config::load)
.await
.context("daemon configuration reader panicked")??;
self.owner().install_config(config);
} else {
let controller_loader = self.controller_loader;
let controller = tokio::task::spawn_blocking(controller_loader)
.await
.context("daemon controller loader panicked")??;
self.owner().install_controller(controller);
}
self.publish_revision();
Ok(())
}
pub(super) fn missing_target_record(
&self,
session_id: &str,
view: &ManagedSessionView,
) -> Option<(String, String)> {
let Some(ViewError::TargetMissing(detail)) = &view.error else {
return None;
};
if view.connected {
return None;
}
let controller_owner = self.owner();
if controller_owner
.lifecycle
.get(session_id)
.is_some_and(|active| active.is_running())
{
return None;
}
let controller = controller_owner.controller();
let session = controller.state.sessions.get(session_id)?;
if !matches!(
session.state,
SessionState::Running | SessionState::Disconnected
) {
return None;
}
Some((detail.clone(), session.updated_at.clone()))
}
pub(super) async fn persist_missing_target(
self: &Arc<Self>,
session_id: &str,
detail: String,
observed_updated_at: String,
) -> Result<()> {
let changed = blocking({
let session_id = session_id.to_owned();
let detail = detail.clone();
move || {
crate::database::mark_session_target_missing_if_current(
&session_id,
&detail,
&chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
&observed_updated_at,
)
}
})
.await?;
if changed.is_some() {
self.reload_controller().await?;
if changed == Some(SessionState::Lost) {
self.discard_lost_session(session_id.to_owned()).await;
} else {
self.push_notice(session_id, detail);
}
}
Ok(())
}
pub(super) fn sessions_without_a_usable_harness(&self) -> Vec<UnreadySession> {
let now = chrono::Utc::now();
let observations = {
let owner = self.owner();
owner
.indexes
.running
.keys()
.filter(|id| {
!owner
.lifecycle
.get(*id)
.is_some_and(|active| active.is_running())
})
.filter(|id| !owner.close_requested.contains(*id))
.filter_map(|id| owner.controller().state.sessions.get(id))
.map(|record| {
let view = owner.sessions.get(&record.id);
ReadinessObservation {
harness_ready: view.is_some_and(|view| {
view.operational.as_ref().is_some_and(
mj_core::relay::RelayOperationalState::native_session_is_ready,
)
}),
record_age: record_age(&record.updated_at, now),
updated_at: record.updated_at.clone(),
detail: view
.and_then(|view| view.error.as_ref())
.map(|error| error.detail().to_owned()),
session_id: record.id.clone(),
}
})
.collect::<Vec<_>>()
};
self.harness_readiness
.lock()
.unwrap_or_else(PoisonError::into_inner)
.observe(std::time::Instant::now(), observations)
}
pub(super) async fn fail_unready_session(&self, unready: UnreadySession) {
let cause = unready.cause();
let session_id = unready.session_id.clone();
tracing::warn!(%session_id, waited_seconds = unready.waited.as_secs(), "{cause}");
let applied = blocking({
let session_id = session_id.clone();
let cause = cause.clone();
move || {
let mut controller = Controller::load()?;
controller.fail_unready_session(&session_id, &cause, &unready.observed_updated_at)
}
})
.await;
match applied {
Ok(true) => {
if let Err(error) = self.reload_controller().await {
tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after failing an unready session");
}
self.push_notice(&session_id, cause);
self.publish_revision();
}
Ok(false) => {}
Err(error) => tracing::warn!(
%session_id,
error = format!("{error:#}"),
"could not record that a session's harness never became usable"
),
}
}
pub(super) async fn publish_session(
&self,
session_id: String,
view: ManagedSessionView,
) -> Result<()> {
let connected = view.connected;
let has_snapshot = view.snapshot.is_some();
tracing::debug!(
%session_id,
connected,
has_snapshot,
"daemon received a session view"
);
{
let mut owner = self.owner();
let controller = owner.controller();
if !controller.state.sessions.contains_key(&session_id) {
return Ok(());
}
if view.connected
&& let Some(snapshot) = view.snapshot.as_ref()
&& let Some(session) = controller.state.sessions.get(&session_id)
&& session.state.has_live_worker()
{
let policy = BackgroundPolicyState {
quiet: snapshot.operational.safe_to_replace(session.harness_kind),
checkpoint_wait: snapshot
.operational
.routine_checkpoint_wait(session.harness_kind),
latest_completed_turn_ordinal: snapshot.latest_completed_turn_ordinal(),
worker_build: snapshot.worker_build.clone(),
};
self.observe_background_policy(session, &controller.config, &policy);
owner.background_policies.insert(session_id.clone(), policy);
} else {
owner.background_policies.remove(&session_id);
}
owner.sessions.insert(
session_id.clone(),
RuntimeSessionView::from_managed(session_id, view),
);
}
reach_test_hook("relay_projection_before_revision_publication").await?;
self.publish_revision();
Ok(())
}
pub(super) async fn runtime_snapshot(
&self,
workspace_id: &str,
after_revision: u64,
all_workspaces: bool,
) -> Result<RuntimeSnapshot> {
let mut revisions = self.revisions.subscribe();
if *revisions.borrow_and_update() <= after_revision {
let _ = tokio::time::timeout(Duration::from_secs(30), revisions.changed()).await;
}
let revision = self.revisions.current();
let moves = blocking(crate::database::load_move_operations).await?;
let workspace_names = self
.workspaces()
.borrow()
.clone()
.into_iter()
.map(|workspace| (workspace.id, workspace.name))
.collect();
let session_ids = if all_workspaces {
self.owner()
.controller()
.state
.sessions
.keys()
.cloned()
.collect::<BTreeSet<_>>()
} else {
let workspace_id = workspace_id.to_owned();
blocking(move || crate::database::session_ids_for_workspace(&workspace_id))
.await?
.into_iter()
.collect()
};
let native_owners = session_ids.clone();
let native_agents = blocking(move || {
let mut agents = Vec::new();
for owner in native_owners {
agents.extend(crate::database::load_native_agent_summaries(&owner)?);
}
Ok(agents)
})
.await?;
let controller_owner = self.owner();
let sessions = controller_owner
.sessions
.iter()
.filter(|(session_id, _)| session_ids.contains(*session_id))
.map(|(_, view)| view.clone())
.collect();
let controller = controller_owner.controller();
let lifecycles = controller_owner
.lifecycle
.iter()
.filter(|(session_id, active)| {
(all_workspaces
|| session_ids.contains(*session_id)
|| active.resume_workspace_id.as_deref() == Some(workspace_id))
&& active.is_visible()
})
.map(|(session_id, active)| RuntimeLifecycleView {
operation_id: active.operation_id.clone(),
cancellable: active.is_cancellable()
&& lifecycle_cancellable(
active.kind,
durable_session_state(controller, session_id),
),
session_id: session_id.clone(),
kind: active.kind.into(),
started_at_epoch_seconds: active.started_at_epoch_seconds,
active_stages: active
.active_stages
.iter()
.map(|(stage, (_, started_at))| (*stage, *started_at))
.collect(),
resume_destination: active.resume_destination.clone(),
notice: active.notice.clone(),
})
.collect();
let reviews = self
.review_host
.views()
.into_iter()
.filter(|review| session_ids.contains(&review.session_id))
.collect();
let notices = self
.notices
.lock()
.unwrap_or_else(PoisonError::into_inner)
.iter()
.filter(|notice| notice_reaches_workspace(notice, &session_ids))
.cloned()
.collect();
let records = runtime_records_for_workspace(controller, &session_ids);
let subagents = runtime_subagents_for_workspace(controller, &records);
Ok(RuntimeSnapshot {
last_subagent_policy: controller.state.last_subagent_policy.clone(),
native_agents,
workspace_names,
moves: moves
.into_iter()
.filter(|operation| session_ids.contains(&operation.selection.session_id))
.collect(),
revision,
config: controller.config.clone(),
records,
sessions,
lifecycles,
reviews,
notices,
subagents,
})
}
}
pub(super) fn notice_reaches_workspace(
notice: &RuntimeNotice,
session_ids: &BTreeSet<String>,
) -> bool {
notice.session_id.is_empty() || session_ids.contains(¬ice.session_id)
}
fn record_age(updated_at: &str, now: chrono::DateTime<chrono::Utc>) -> Option<Duration> {
let written = chrono::DateTime::parse_from_rfc3339(updated_at).ok()?;
now.signed_duration_since(written.with_timezone(&chrono::Utc))
.to_std()
.ok()
}