use super::*;
impl RuntimeState {
pub(super) fn new(
session_manager: SessionManagerControl,
controller: Controller,
recovery_observer: RecoveryObserver,
worker_upgrade_observer: WorkerUpgradeObserver,
workspaces: Vec<WorkspaceRecord>,
) -> Self {
Self::new_with_controller_loader(
session_manager,
controller,
recovery_observer,
worker_upgrade_observer,
workspaces,
Controller::load,
)
}
pub(super) fn new_with_controller_loader(
session_manager: SessionManagerControl,
controller: Controller,
recovery_observer: RecoveryObserver,
worker_upgrade_observer: WorkerUpgradeObserver,
workspaces: Vec<WorkspaceRecord>,
controller_loader: fn() -> Result<Controller>,
) -> Self {
let initial_revision = u64::try_from(chrono::Utc::now().timestamp_micros()).unwrap_or(1);
let revisions = RuntimeRevisions::new(initial_revision);
let (workspaces_tx, _) = tokio::sync::watch::channel(workspaces);
let review_config = Arc::new(Mutex::new(controller.config.review.clone()));
let review_host = TurnReviewHost::spawn_notifying(
session_manager.clone(),
{
let installed = review_config.clone();
Arc::new(move || {
installed
.lock()
.unwrap_or_else(PoisonError::into_inner)
.clone()
})
},
revisions.notifier(),
Some(recovery_observer.gate.clone()),
);
Self {
attachments: Mutex::new(BTreeMap::new()),
phone_status: Mutex::new(WebViewerStatus::Starting),
web_viewer: crate::web_viewer::ViewerControl::new(),
ever_attached: AtomicBool::new(false),
sessions: Mutex::new(BTreeMap::new()),
background_policies: Mutex::new(BTreeMap::new()),
revisions,
workspaces_tx,
workspace_refresh: tokio::sync::Mutex::new(()),
session_manager,
lifecycle: Mutex::new(BTreeMap::new()),
workspace_closes: Mutex::new(BTreeMap::new()),
workspace_resume_admission: Mutex::new(BTreeMap::new()),
harness_readiness: Mutex::new(HarnessReadinessWatch::default()),
startup_prompts: Mutex::new(BTreeMap::new()),
close_requested: Mutex::new(BTreeSet::new()),
controller: Mutex::new(controller),
controller_loader,
config_mutation: tokio::sync::Mutex::new(()),
recovery_observer,
worker_upgrade_observer,
notices: Mutex::new(VecDeque::new()),
next_notice_id: AtomicU64::new(1),
review_config,
review_host,
wiki: crate::sessionwiki::WikiIndexer::spawn(),
}
}
pub fn review_host(&self) -> &TurnReviewHost {
&self.review_host
}
pub fn wiki(&self) -> &crate::sessionwiki::WikiIndexer {
&self.wiki
}
pub async fn wiki_search(&self, query: String, limit: usize) -> Result<WikiSearchPage> {
if crate::sessionwiki::sync_is_stale(self.wiki.last_success()) {
self.wiki.request_sync(false);
}
let live = self.live_session_ids();
let rows =
blocking(move || crate::sessionwiki::query_rows(&query, limit, &live, false)).await?;
Ok(WikiSearchPage {
rows,
status: self.wiki.status(),
})
}
pub async fn wiki_brief(&self, wiki_id: String, max_chars: usize) -> Result<Option<String>> {
blocking(move || crate::sessionwiki::brief(&wiki_id, max_chars)).await
}
pub async fn wiki_hits(
&self,
wiki_id: String,
query: String,
context_messages: usize,
per_message_chars: usize,
) -> Result<Option<WikiHitTranscript>> {
blocking(move || {
crate::sessionwiki::transcript_hits(
&wiki_id,
&query,
context_messages,
per_message_chars,
)
})
.await
}
pub async fn wiki_session(
&self,
wiki_id: String,
) -> Result<Option<mj_client::daemon::WikiSessionInfo>> {
let known = self.live_session_ids();
blocking(move || crate::sessionwiki::wiki_session(&wiki_id, &known)).await
}
pub async fn restore_wiki_session(
self: &Arc<Self>,
request: WikiRestoreRequest,
cancellation: &CancellationToken,
) -> Result<Option<RegisteredSession>> {
let wiki_id = request.wiki_id.clone();
let Some(archived) =
blocking(move || crate::sessionwiki::archived_session(&wiki_id)).await?
else {
return Ok(None);
};
let project_directory = request
.project_directory
.clone()
.or_else(|| archived.project_directory.clone())
.context(
"name a project directory: the archived session's own project is no longer on this machine",
)?;
let source = project_directory.display().to_string();
let bundle_id = blocking(move || {
crate::controller::create_bundle_from_sources(&[source])
.map(|created| created.bundle_id)
.map_err(anyhow::Error::new)
})
.await
.context("find or create a bundle for the restored session's project")?;
let registered = self
.start_create_session(CreateSessionRequest {
launch_base: None,
launch_branch: None,
create_managed_worktree: None,
mjolnir_subagents: None,
initial_prompt: None,
workspace_id: request.workspace_id,
profile_id: request.profile_id,
bundle_id,
project_directory: Some(project_directory),
target_template_id: request.target_template_id,
additional_mounts: request.additional_mounts,
resource_allocation: request.resource_allocation,
title: archived.title.clone(),
session_title_override: Some(archived.title.clone()),
})
.await?;
let session_id = registered.session.id.clone();
self.queue_startup_step(
&session_id,
StartupStep::InstallHandoff(Box::new(archived.snapshot)),
cancellation,
)?;
Ok(Some(registered))
}
pub(super) fn live_session_ids(&self) -> BTreeSet<String> {
self.controller
.lock()
.unwrap_or_else(PoisonError::into_inner)
.state
.sessions
.keys()
.cloned()
.collect()
}
async fn install_archive_handoff(
&self,
session_id: &str,
handle: &crate::session_manager::ManagedSessionHandle,
snapshot: &mj_core::archive::CanonicalSessionSnapshot,
) -> Result<()> {
let (config, profile_id) = {
let controller = self
.controller
.lock()
.unwrap_or_else(PoisonError::into_inner);
let profile_id = controller
.state
.sessions
.get(session_id)
.map(|record| record.last_profile.clone());
(controller.config.clone(), profile_id)
};
let context_bytes = crate::handoff::profile_handoff_bytes(
profile_id.and_then(|id| config.profiles.get(&id)),
);
let cancel = CancellationToken::new();
let handoff = crate::handoff::build_handoff_context(
session_id,
&config,
snapshot,
context_bytes,
&cancel,
)
.await
.context("compact the archived transcript")?;
handle
.install_prompt_context(format!(
"{} {handoff}",
crate::compaction::ARCHIVE_HANDOFF_PREAMBLE
))
.await
.context("install the archived hand-off")?;
tracing::info!(
session_id,
bytes = handoff.len(),
"installed the restored archive's hand-off"
);
Ok(())
}
pub(super) async fn wait_for_ready_session(
&self,
session_id: &str,
) -> Result<crate::session_manager::ManagedSessionHandle> {
const POLL: Duration = Duration::from_millis(250);
let deadline = tokio::time::Instant::now() + Duration::from_secs(30 * 60);
loop {
match self.session_state(session_id) {
Some(
SessionState::Provisioning
| SessionState::Running
| SessionState::Disconnected
| SessionState::Checkpointing,
) => {}
Some(state) => bail!("session {session_id} is {state:?} before its hand-off"),
None => bail!("session {session_id} disappeared before its hand-off"),
}
if let Ok(handle) = self.session_manager.session(session_id).await {
let view = handle.view();
if let Some(ViewError::TargetMissing(detail)) = &view.error {
bail!("session {session_id} lost its target: {detail}");
}
if view.connected
&& view
.snapshot
.is_some_and(|snapshot| snapshot.operational.native_session_is_ready())
{
return Ok(handle);
}
}
ensure!(
tokio::time::Instant::now() < deadline,
"session {session_id} was not ready for its hand-off within 30 minutes"
);
tokio::time::sleep(POLL).await;
}
}
pub(super) async fn fail_unfinished_provisioning(
self: &Arc<Self>,
session_id: &str,
error: &str,
) {
let provisioning = {
let controller = self
.controller
.lock()
.unwrap_or_else(PoisonError::into_inner);
durable_session_state(&controller, session_id) == Some(SessionState::Provisioning)
};
if !provisioning {
return;
}
let cause = format!("session provisioning ended without finishing: {error}");
let applied = blocking({
let session_id = session_id.to_owned();
move || {
let mut controller = Controller::load()?;
controller.fail_interrupted_lifecycle(&session_id, &cause)
}
})
.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 recording a failed provision");
}
}
Ok(false) => {}
Err(error) => tracing::warn!(
%session_id,
error = format!("{error:#}"),
"could not record that provisioning ended without finishing"
),
}
}
pub(super) async fn record_failed_close(
self: &Arc<Self>,
session_id: &str,
reference: &str,
failure: &LifecycleFailure,
) {
self.record_lifecycle_failure(
session_id,
reference,
failure,
mj_core::state::CLOSE_FAILURE_PREFIX,
)
.await;
}
pub(crate) async fn record_lifecycle_failure(
self: &Arc<Self>,
session_id: &str,
reference: &str,
failure: &LifecycleFailure,
prefix: &str,
) {
let cause = match &failure.refusal {
Some(refusal) => format!("{prefix}: {refusal}"),
None => {
format!("{prefix}; the daemon log records the reason under reference {reference}")
}
};
let applied = blocking({
let session_id = session_id.to_owned();
let cause = cause.clone();
move || {
let mut controller = Controller::load()?;
controller.record_failed_close(&session_id, &cause)
}
})
.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 recording a failed close");
}
self.publish_revision();
}
Ok(false) => {}
Err(error) => tracing::warn!(
%session_id,
error = format!("{error:#}"),
"could not record why a close failed"
),
}
}
pub async fn clear_recorded_close_failure(self: &Arc<Self>, session_id: &str) {
let recorded = self
.controller
.lock()
.unwrap_or_else(PoisonError::into_inner)
.state
.sessions
.get(session_id)
.is_some_and(|record| record.public_error().is_some());
if !recorded {
return;
}
let cleared = blocking({
let session_id = session_id.to_owned();
move || {
let mut controller = Controller::load()?;
controller.clear_recorded_close_failure(&session_id)
}
})
.await;
match cleared {
Ok(true) => {
if let Err(error) = self.reload_controller().await {
tracing::warn!(%session_id, error = format!("{error:#}"), "could not reload state after clearing a recorded close failure");
}
self.publish_revision();
}
Ok(false) => {}
Err(error) => tracing::warn!(
%session_id,
error = format!("{error:#}"),
"could not clear a recorded close failure"
),
}
}
pub(super) fn note_lifecycle_outcome(&self, session_id: &str) {
let stopped = {
let controller = self
.controller
.lock()
.unwrap_or_else(PoisonError::into_inner);
durable_session_state(&controller, session_id) == Some(SessionState::Stopped)
};
if stopped {
self.wiki.request_sync(false);
}
}
pub fn allocate_revision(&self) -> u64 {
self.revisions.allocate()
}
pub(super) fn publish_revision(&self) -> u64 {
self.revisions.publish()
}
pub(super) fn attachments(&self) -> std::sync::MutexGuard<'_, BTreeMap<String, Attachment>> {
self.attachments
.lock()
.unwrap_or_else(PoisonError::into_inner)
}
pub(super) fn prune_dead_clients(&self) {
self.attachments()
.retain(|_, attachment| process_is_alive(attachment.pid));
}
pub(super) fn workspace_has_active_resume(&self, workspace_id: &str) -> bool {
self.lifecycle
.lock()
.unwrap_or_else(PoisonError::into_inner)
.values()
.any(|active| {
active.result.borrow().is_none()
&& active.resume_workspace_id.as_deref() == Some(workspace_id)
})
}
pub fn publish_web_access(&self, access: crate::server::WebViewerAccess) {
use crate::server::WebViewerAccess;
let status = match &access {
WebViewerAccess::Starting => WebViewerStatus::Starting,
WebViewerAccess::Ready {
viewer_url,
viewer_code,
qr_login_url,
fallback_reason,
..
} => WebViewerStatus::Ready {
viewer_url: viewer_url.clone(),
viewer_code: viewer_code.clone(),
qr_login_url: qr_login_url.clone(),
fallback_reason: fallback_reason.clone(),
},
WebViewerAccess::Failed {
address, message, ..
} => WebViewerStatus::Error {
message: format!("{message} Address: {address}"),
},
WebViewerAccess::Unavailable(message) => WebViewerStatus::Error {
message: message.clone(),
},
};
self.web_viewer.publish(access);
self.set_phone_status(status);
}
pub(super) fn set_phone_status(&self, status: WebViewerStatus) {
*self
.phone_status
.lock()
.unwrap_or_else(PoisonError::into_inner) = status;
}
pub(super) fn phone_status(&self) -> WebViewerStatus {
self.phone_status
.lock()
.unwrap_or_else(PoisonError::into_inner)
.clone()
}
pub(super) fn workspaces(&self) -> tokio::sync::watch::Receiver<Vec<WorkspaceRecord>> {
self.workspaces_tx.subscribe()
}
pub(super) fn worker_poll_exclusion_session_ids(
&self,
controller: &Controller,
) -> BTreeSet<String> {
self.lifecycle
.lock()
.unwrap_or_else(PoisonError::into_inner)
.iter()
.filter(|(session_id, active)| {
active.result.borrow().is_none()
&& (active.move_source_closed
|| lifecycle_owns_worker_target(
active.kind,
controller
.state
.sessions
.get(*session_id)
.map(|session| session.state),
))
})
.map(|(session_id, _)| session_id.clone())
.collect()
}
pub fn revisions(&self) -> tokio::sync::watch::Receiver<u64> {
self.revisions.subscribe()
}
pub fn with_config<T>(&self, read: impl FnOnce(&Config) -> T) -> T {
read(
&self
.controller
.lock()
.unwrap_or_else(PoisonError::into_inner)
.config,
)
}
pub async fn create_quick_bundle(
&self,
source: String,
) -> std::result::Result<
crate::controller::QuickBundleCreation,
crate::controller::QuickBundleFailure,
> {
let _mutation = self.config_mutation.lock().await;
tokio::task::spawn_blocking(move || crate::controller::create_quick_bundle(&source))
.await
.map_err(|error| {
crate::controller::QuickBundleFailure::Persistence(anyhow!(
"bundle creation task panicked: {error}"
))
})?
}
pub async fn create_bundle_from_sources(
&self,
sources: Vec<String>,
) -> std::result::Result<
crate::controller::QuickBundleCreation,
crate::controller::QuickBundleFailure,
> {
let _mutation = self.config_mutation.lock().await;
tokio::task::spawn_blocking(move || crate::controller::create_bundle_from_sources(&sources))
.await
.map_err(|error| {
crate::controller::QuickBundleFailure::Persistence(anyhow!(
"bundle creation task panicked: {error}"
))
})?
}
pub(crate) fn publish_workspaces(&self, workspaces: Vec<WorkspaceRecord>) {
self.workspaces_tx.send_replace(workspaces);
self.publish_revision();
}
pub(crate) async fn refresh_workspaces(&self) -> Result<()> {
let _refresh = self.workspace_refresh.lock().await;
let workspaces = tokio::task::spawn_blocking(crate::database::list_workspaces)
.await
.context("daemon workspace refresh task panicked")??;
self.publish_workspaces(workspaces);
Ok(())
}
pub(crate) fn queue_startup_step(
self: &Arc<Self>,
session_id: &str,
step: StartupStep,
cancellation: &CancellationToken,
) -> Result<()> {
match self.session_state(session_id) {
Some(
SessionState::Provisioning
| SessionState::Running
| SessionState::Disconnected
| SessionState::Checkpointing,
) => {}
Some(state) => {
bail!("session {session_id} is {state:?}; it cannot take a queued prompt")
}
None => bail!("unknown session {session_id}"),
}
let mut queues = self
.startup_prompts
.lock()
.unwrap_or_else(PoisonError::into_inner);
if let Some(queue) = queues.get_mut(session_id) {
queue.pending.push_back(step);
return Ok(());
}
let cancel = cancellation.child_token();
let upgrade_work = crate::upgrade::activity("startup prompt delivery")?;
queues.insert(
session_id.to_owned(),
StartupQueue {
pending: VecDeque::from([step]),
in_flight: false,
cancel: cancel.clone(),
task: None,
},
);
let runtime = Arc::clone(self);
let drain_session = session_id.to_owned();
let task = tokio::spawn(async move {
let _upgrade_work = upgrade_work;
let supervised = {
let runtime = Arc::clone(&runtime);
let session_id = drain_session.clone();
let cancel = cancel.clone();
tokio::spawn(async move { runtime.drain_startup_queue(&session_id, &cancel).await })
};
if let Err(error) = supervised.await {
runtime
.fail_startup_queue(
&drain_session,
None,
&format!("the daemon's delivery task failed: {error}"),
)
.await;
}
});
if let Some(queue) = queues.get_mut(session_id) {
queue.task = Some(task);
}
Ok(())
}
async fn drain_startup_queue(self: Arc<Self>, session_id: &str, cancel: &CancellationToken) {
let handle = tokio::select! {
() = cancel.cancelled() => {
self.fail_startup_queue(
session_id,
None,
"the daemon stopped before the session was ready",
)
.await;
return;
}
ready = self.wait_for_ready_session(session_id) => match ready {
Ok(handle) => handle,
Err(error) => {
self.fail_startup_queue(session_id, None, &format!("{error:#}"))
.await;
return;
}
},
};
loop {
let step = {
let mut queues = self
.startup_prompts
.lock()
.unwrap_or_else(PoisonError::into_inner);
let Some(queue) = queues.get_mut(session_id) else {
return;
};
match queue.pending.pop_front() {
Some(step) => {
queue.in_flight = true;
step
}
None => {
queues.remove(session_id);
return;
}
}
};
let outcome = tokio::select! {
() = cancel.cancelled() => {
Err(anyhow!("the daemon stopped before the prompt was sent"))
}
result = self.run_startup_step(session_id, &handle, &step) => result,
};
if let Err(error) = outcome {
self.fail_startup_queue(session_id, Some(step), &format!("{error:#}"))
.await;
return;
}
let mut queues = self
.startup_prompts
.lock()
.unwrap_or_else(PoisonError::into_inner);
let Some(queue) = queues.get_mut(session_id) else {
return;
};
queue.in_flight = false;
if queue.pending.is_empty() {
queues.remove(session_id);
return;
}
}
}
async fn run_startup_step(
&self,
session_id: &str,
handle: &crate::session_manager::ManagedSessionHandle,
step: &StartupStep,
) -> Result<()> {
match step {
StartupStep::InstallHandoff(snapshot) => {
self.install_archive_handoff(session_id, handle, snapshot)
.await
}
StartupStep::Prompt {
text,
inherited_draft,
} => {
self.submit_startup_prompt(session_id, handle, text, inherited_draft.as_deref())
.await
}
}
}
async fn submit_startup_prompt(
&self,
session_id: &str,
handle: &crate::session_manager::ManagedSessionHandle,
text: &str,
inherited_draft: Option<&str>,
) -> Result<()> {
let bundle_id = self
.controller
.lock()
.unwrap_or_else(PoisonError::into_inner)
.state
.sessions
.get(session_id)
.map(|record| record.bundle_id.clone());
let ordinal = handle
.submit(
new_command_id("startup")?,
RelayCommand::Prompt {
prompt: vec![ContentBlock::Text(TextContent::new(text.to_owned()))],
},
)
.await?;
if let Some(expected) = inherited_draft {
let persisted_id = session_id.to_owned();
let persisted_expected = expected.to_owned();
if let Err(error) = blocking(move || {
crate::database::clear_session_draft_input_if_matches(
&persisted_id,
&persisted_expected,
)
})
.await
{
tracing::warn!(
session_id,
error = format!("{error:#}"),
"the delivered prompt's draft could not be cleared"
);
}
if let Some(record) = self
.controller
.lock()
.unwrap_or_else(PoisonError::into_inner)
.state
.sessions
.get_mut(session_id)
&& record.draft_input == expected
{
record.draft_input.clear();
}
self.publish_revision();
}
if let Some(bundle_id) = bundle_id {
let history_id = session_id.to_owned();
let history_text = text.to_owned();
if let Err(error) = blocking(move || {
crate::database::record_prompt(
&history_id,
&bundle_id,
ordinal,
None,
&history_text,
)
})
.await
{
tracing::warn!(
session_id,
error = format!("{error:#}"),
"the queued prompt was accepted but its history could not be stored"
);
}
}
Ok(())
}
async fn fail_startup_queue(
&self,
session_id: &str,
failed: Option<StartupStep>,
reason: &str,
) {
let remaining = self
.startup_prompts
.lock()
.unwrap_or_else(PoisonError::into_inner)
.remove(session_id)
.map(|queue| queue.pending)
.unwrap_or_default();
let mut texts = Vec::new();
let mut dropped_handoff = false;
for step in failed.into_iter().chain(remaining) {
match step {
StartupStep::Prompt { text, .. } => texts.push(text),
StartupStep::InstallHandoff(_) => dropped_handoff = true,
}
}
if dropped_handoff {
tracing::warn!(
session_id,
reason,
"could not install the restored archive's hand-off"
);
self.push_notice(
session_id,
format!("The restored session started without its archived hand-off: {reason}"),
);
}
if texts.is_empty() {
return;
}
let restored = texts.join("\n\n");
if let Err(error) = self.append_draft_input(session_id, &restored).await {
tracing::warn!(
session_id,
error = format!("{error:#}"),
"a queued prompt could not be saved back into the session's draft"
);
}
self.push_notice(
session_id,
format!(
"Your prompt could not be sent to session {} ({reason}); it is back in the composer draft.",
mj_core::state::short_id(session_id)
),
);
tracing::warn!(
session_id,
reason,
"a queued startup prompt could not be delivered"
);
}
pub(super) async fn append_draft_input(&self, session_id: &str, text: &str) -> Result<()> {
let existing = self
.session_record(session_id)
.map(|record| record.draft_input)
.unwrap_or_default();
let combined = [existing.as_str(), text]
.into_iter()
.filter(|part| !part.is_empty())
.collect::<Vec<_>>()
.join("\n\n");
let persisted_id = session_id.to_owned();
let persisted = combined.clone();
let stored =
blocking(move || crate::database::set_session_draft_input(&persisted_id, &persisted))
.await;
if let Some(record) = self
.controller
.lock()
.unwrap_or_else(PoisonError::into_inner)
.state
.sessions
.get_mut(session_id)
{
record.draft_input = combined;
}
self.publish_revision();
stored
}
pub(crate) async fn cancel_and_join_startup_prompts(&self) -> Result<()> {
let tasks = {
let mut queues = self
.startup_prompts
.lock()
.unwrap_or_else(PoisonError::into_inner);
queues
.values_mut()
.filter_map(|queue| {
queue.cancel.cancel();
queue.task.take()
})
.collect::<Vec<_>>()
};
let deadline = tokio::time::Instant::now() + Duration::from_secs(1);
let mut outcome = Ok(());
for task in tasks {
let joined = match tokio::time::timeout_at(deadline, task).await {
Ok(Ok(())) => Ok(()),
Ok(Err(error)) => Err(anyhow!("startup prompt delivery task failed: {error}")),
Err(_) => Err(anyhow!(
"a startup prompt delivery task did not stop within 1s"
)),
};
if outcome.is_ok() {
outcome = joined;
} else if let Err(error) = joined {
tracing::warn!(%error, "another startup prompt drain did not stop cleanly");
}
}
outcome
}
}