use super::*;
impl HostState {
pub(super) async fn observe(
&mut self,
session_id: String,
snapshot: Option<Box<MaterializedSession>>,
prompt_driven: bool,
) {
let execution = snapshot
.as_ref()
.map_or(MaterializedExecutionState::Idle, |snapshot| {
snapshot.execution
});
let previous = self.sessions.insert(
session_id.clone(),
SessionWatch {
execution,
prompt_driven,
materialized: snapshot,
},
);
if self.recovery_candidates.contains(&session_id)
&& matches!(execution, MaterializedExecutionState::Idle)
{
self.begin_recovery(&session_id);
return;
}
let finished_turn = previous.as_ref().is_some_and(|watch| {
watch.prompt_driven
&& matches!(watch.execution, MaterializedExecutionState::Running { .. })
}) && matches!(execution, MaterializedExecutionState::Idle);
if !finished_turn || !(self.config)().enabled {
return;
}
self.begin(session_id, false, None);
}
pub(super) fn begin(
&mut self,
session_id: String,
manual: bool,
reply: Option<oneshot::Sender<Result<(), StartRefusal>>>,
) {
if crate::controller::move_session::move_owns_session(&session_id) {
answer(reply, Err(StartRefusal("session is moving".to_owned())));
return;
}
if let Some(refusal) = self.refuse_start(&session_id) {
answer(reply, Err(refusal));
return;
}
if self.preparing.contains(&session_id) {
answer(
reply,
Err(StartRefusal("a review is already starting".to_owned())),
);
return;
}
let config = (self.config)();
let Some(profile) = config.reviewer_profile().map(str::to_owned) else {
let refusal = StartRefusal(
"turn review needs a reviewer: set [review] profile in config.toml".to_owned(),
);
if self.missing_reviewer_reported.insert(session_id.clone()) {
self.record_notice(&session_id, refusal.0.clone());
}
answer(reply, Err(refusal));
return;
};
let reviewer = ReviewerIdentity {
profile,
model: config.model.clone(),
effort: config.effort.clone(),
};
let tier = config.tier;
let control = self.control.clone();
let events = self.events.clone();
let prepare_session = session_id.clone();
let environment = self.environment.clone();
hold_prompts(&session_id);
self.preparing.insert(session_id);
tokio::spawn(async move {
let prepared = prepare(&control, &environment, &prepare_session, &reviewer, tier).await;
let _ = events.send(HostEvent::Prepared {
session_id: prepare_session,
manual,
reply,
prepared,
});
});
}
pub(super) fn begin_recovery(&mut self, session_id: &str) {
if !self.recovery_candidates.contains(session_id)
|| self.recovery_in_flight.contains(session_id)
|| self.reviews.contains_key(session_id)
|| self.preparing.contains(session_id)
{
return;
}
let Some(watch) = self.sessions.get(session_id) else {
return;
};
if watch
.materialized
.as_ref()
.is_none_or(|snapshot| !snapshot.queued_prompts.is_empty())
{
return;
}
hold_prompts(session_id);
self.recovery_in_flight.insert(session_id.to_owned());
let control = self.control.clone();
let environment = self.environment.clone();
let events = self.events.clone();
let session_id = session_id.to_owned();
tokio::spawn(async move {
let prepared = prepare_recovery(&control, &environment, &session_id).await;
let _ = events.send(HostEvent::RecoveryPrepared {
session_id,
prepared,
});
});
}
pub(super) fn recovery_prepared(
&mut self,
session_id: String,
prepared: Result<Option<Prepared>, String>,
) {
self.recovery_in_flight.remove(&session_id);
match prepared {
Ok(Some(prepared)) => {
self.recovery_candidates.remove(&session_id);
self.preparing.insert(session_id.clone());
self.prepared(session_id, false, None, Ok(prepared));
}
Ok(None) => {
self.recovery_candidates.remove(&session_id);
release_prompts(&session_id);
self.record_notice(
&session_id,
"Turn review was cancelled when Mjolnir restarted; the next review covers the same changes".to_owned(),
);
}
Err(error) => {
tracing::warn!(session_id = %session_id, %error, "could not reconcile an interrupted review handoff");
release_prompts(&session_id);
}
}
}
pub(super) fn refuse_start(&self, session_id: &str) -> Option<StartRefusal> {
if self.reviews.contains_key(session_id) {
return Some(StartRefusal("a review is already open".to_owned()));
}
if self.recovery_candidates.contains(session_id)
|| self.recovery_in_flight.contains(session_id)
{
return Some(StartRefusal(
"an interrupted review handoff is being reconciled".to_owned(),
));
}
let Some(watch) = self.sessions.get(session_id) else {
return Some(StartRefusal("this session is not connected".to_owned()));
};
if !matches!(watch.execution, MaterializedExecutionState::Idle) {
return Some(StartRefusal(
"a review runs between turns; this one is still working".to_owned(),
));
}
let queued = watch
.materialized
.as_ref()
.is_some_and(|materialized| !materialized.queued_prompts.is_empty());
if queued {
return Some(StartRefusal(
"prompts are queued; the review waits for them".to_owned(),
));
}
None
}
}