use super::*;
impl HostState {
pub(super) fn record_notice(&self, session_id: &str, text: String) {
let control = self.control.clone();
let session_id = session_id.to_owned();
tokio::spawn(async move {
let recorded = async {
let handle = control
.session(session_id.clone())
.await
.map_err(|error| format!("{error:#}"))?;
let command_id =
new_command_id("turn-review-notice").map_err(|error| format!("{error:#}"))?;
handle
.submit(command_id, RelayCommand::RecordNotice { text })
.await
.map(|_| ())
.map_err(|error| format!("{error:#}"))
}
.await;
if let Err(error) = recorded {
tracing::debug!(
session_id = %session_id,
%error,
"could not record a review notice in the conversation"
);
}
});
}
pub(super) fn review_step(
&mut self,
session_id: &str,
action: ReviewerAction,
into_step: impl FnOnce(Result<ReviewerOutcome, String>) -> ReviewStep + Send + 'static,
) {
let Some(epoch) = self.reviews.get(session_id).map(|slot| slot.epoch) else {
return;
};
let owner = session_id.to_owned();
self.spawn_reviewer(session_id.to_owned(), None, action, move |outcome| {
Some(HostEvent::Step {
session_id: owner,
epoch,
step: into_step(outcome),
})
});
}
pub(super) fn spawn_reviewer(
&self,
session_id: String,
role: Option<String>,
action: ReviewerAction,
into_event: impl FnOnce(Result<ReviewerOutcome, String>) -> Option<HostEvent> + Send + 'static,
) {
let control = self.control.clone();
let events = self.events.clone();
tokio::spawn(async move {
let retry_backoff = matches!(
action,
ReviewerAction::AckLaneDispatches { .. } | ReviewerAction::ReadLaneDispatches
);
let outcome = reviewer_action(&control, &session_id, role, action).await;
if retry_backoff && outcome.is_err() {
tokio::time::sleep(Duration::from_secs(1)).await;
}
if let Some(event) = into_event(outcome) {
let _ = events.send(event).await;
}
});
}
pub(super) fn publish(&self, session_id: &str) {
let mut views = self
.shared
.views
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let changed = match self.reviews.get(session_id) {
Some(slot) => {
let mut next = slot.view(session_id);
if let Some(error) = slot.delivery_errors.values().next() {
next.status = format!("Review delivery is being reconciled: {error}");
}
if let Some(error) = self.persistence_errors.get(session_id) {
next.status = format!("Review is waiting for durable storage: {error}");
}
if views.get(session_id) == Some(&next) {
false
} else {
views.insert(session_id.to_owned(), next);
true
}
}
None if self.preparing.contains(session_id) => {
let next = RuntimeReviewView {
session_id: session_id.into(),
tier: (self.config)().tier,
phase: TurnReviewPhase::LaunchingReviewer,
roles: Vec::new(),
status: "Preparing reviewer…".into(),
verdict: None,
};
let changed = views.get(session_id) != Some(&next);
views.insert(session_id.into(), next);
changed
}
None => views.remove(session_id).is_some(),
};
drop(views);
if changed {
(self.shared.changed)();
}
}
pub(super) async fn shutdown(&mut self) -> Result<(), String> {
self.shared.stop_delivery.cancel();
for reply in std::mem::take(&mut self.start_replies).into_values() {
let _ = reply.send(Err(StartRefusal(
"daemon shutdown interrupted review acknowledgement".into(),
)));
}
for reply in std::mem::take(&mut self.resolve_replies).into_values() {
let _ = reply.send(Err(
"daemon shutdown interrupted review acknowledgement".into()
));
}
let session_ids = self
.preparing
.iter()
.chain(self.reviews.keys())
.chain(self.recovery_in_flight.iter())
.chain(self.recovery_candidates.iter())
.cloned()
.collect::<BTreeSet<_>>();
for pending in std::mem::take(&mut self.pending_open).into_values() {
answer(
pending.reply,
Err(StartRefusal("the daemon is shutting down".to_owned())),
);
}
let lane = self.persistence.take();
for (session_id, slot) in &mut self.reviews {
if !self.closing.contains(session_id) {
slot.state.orchestration = Some(slot.checkpoint()?);
}
let queued = lane
.as_ref()
.ok_or_else(|| "the review persistence lane stopped".to_owned())
.and_then(|persistence| {
persistence
.send(PersistenceRequest::Save {
session_id: session_id.clone(),
state: Box::new(slot.state.clone()),
completion: None,
})
.map_err(|_| "the review persistence lane stopped".to_owned())
});
if let Err(error) = queued {
tracing::warn!(session_id, %error, "could not queue review shutdown persistence");
}
}
for flag in self.preparation_cancellation.values() {
flag.store(true, std::sync::atomic::Ordering::Release);
}
self.preparation_cancellation.clear();
self.preparing.clear();
self.closing.clear();
self.reviews.clear();
for session_id in &session_ids {
release_prompts(session_id);
self.publish(session_id);
}
drop(lane);
let clear_result: Result<(), String> = Ok(());
let task_result = match self.persistence_task.take() {
Some(task) => task
.await
.map_err(|error| format!("review persistence lane panicked: {error}")),
None => Ok(()),
};
clear_result.and(task_result)
}
}