use std::collections::{BTreeMap, BTreeSet};
use std::path::PathBuf;
use std::sync::{Arc, LazyLock, Mutex};
use std::time::Duration;
use tokio::sync::{mpsc, oneshot};
use crate::hel_session_manager::{
ManagedSessionHandle, ManagedSessionView, ReviewDeliveryAdmission, ReviewerAction,
ReviewerOutcome, SessionManagerControl, new_command_id,
};
use hel::hel_config::ReviewConfig;
use hel::hel_database::TurnReviewState;
use hel::hel_state::{MaterializedExecutionState, MaterializedSession};
use hel::hel_worker::{RelayCommand, RelayEvent, RelayObservation};
use hel::hel_review::driver::{
INTENT_ROLE, PendingForward, Resolution, ReviewRequest, RoleState, RoleStatus, SUPERVISOR_ROLE,
TurnReviewDriver, TurnReviewPhase, TurnReviewSeed,
};
use hel::hel_review::lanes::{ReviewTier, UserMessage};
use hel::hel_review::verdict::ReviewVerdict;
const ROLE_POLL_IDLE_INTERVAL: Duration = Duration::from_millis(200);
#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct RuntimeReviewView {
pub session_id: String,
pub tier: ReviewTier,
pub phase: TurnReviewPhase,
pub roles: Vec<RoleStatus>,
pub status: String,
pub verdict: Option<VerdictView>,
}
impl RuntimeReviewView {
#[must_use]
pub fn activity_label(&self) -> Option<&'static str> {
match &self.phase {
TurnReviewPhase::Resolved(_) => None,
TurnReviewPhase::Forwarding { error: None, .. } => Some("Sending findings"),
TurnReviewPhase::Forwarding { error: Some(_), .. } => Some("Forward failed"),
TurnReviewPhase::Verdict(verdict) => Some(match verdict {
ReviewVerdict::Findings { .. } => "Findings",
ReviewVerdict::Failed { .. } => "Review failed",
ReviewVerdict::Clean => "Review complete",
}),
TurnReviewPhase::Running { roles }
if roles.iter().any(|role| {
role.role == hel::hel_review::driver::VALIDATOR_ROLE
&& matches!(role.state, RoleState::Pending | RoleState::Running)
}) =>
{
Some("Validating")
}
_ => Some("Reviewing"),
}
}
}
#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct VerdictView {
pub kind: VerdictKind,
pub text: String,
pub allowed: Vec<Resolution>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum VerdictKind {
Clean,
Findings,
Failed,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StartRefusal(pub String);
impl std::fmt::Display for StartRefusal {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str(&self.0)
}
}
pub const PROMPT_HELD_MESSAGE: &str =
"a review of the last turn is open; forward, dismiss or cancel it first";
static PROMPT_LOCK: LazyLock<Mutex<BTreeMap<String, PromptHold>>> = LazyLock::new(Mutex::default);
#[derive(Debug, Default)]
struct PromptHold {
delivery_epoch: Option<u64>,
delivery_command_id: Option<String>,
}
pub(crate) fn next_review_generation() -> Result<u64, String> {
let mut random = [0_u8; 8];
getrandom::fill(&mut random)
.map_err(|error| format!("generate reviewer generation: {error}"))?;
let generation = u64::from_le_bytes(random);
if generation == 0 {
return Err("generate reviewer generation: random nonce was zero".to_owned());
}
Ok(generation)
}
#[must_use]
pub fn prompt_refusal(session_id: &str) -> Option<&'static str> {
PROMPT_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.contains_key(session_id)
.then_some(PROMPT_HELD_MESSAGE)
}
fn hold_prompts(session_id: &str) {
PROMPT_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(session_id.to_owned(), PromptHold::default());
}
fn release_prompts(session_id: &str) {
PROMPT_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(session_id);
}
fn admit_review_delivery(
session_id: &str,
epoch: u64,
command_id: &str,
) -> Option<ReviewDeliveryAdmission> {
let mut locks = PROMPT_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let hold = locks.get_mut(session_id)?;
match (hold.delivery_epoch, hold.delivery_command_id.as_deref()) {
(Some(existing_epoch), Some(existing_command))
if existing_epoch != epoch || existing_command != command_id =>
{
None
}
_ => {
hold.delivery_epoch = Some(epoch);
hold.delivery_command_id = Some(command_id.to_owned());
Some(ReviewDeliveryAdmission::new(
session_id.to_owned(),
epoch,
command_id.to_owned(),
))
}
}
}
pub(crate) fn review_delivery_admitted(
session_id: &str,
admission: &ReviewDeliveryAdmission,
) -> bool {
let locks = PROMPT_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
locks.get(session_id).is_some_and(|hold| {
admission.session_id() == session_id
&& hold.delivery_epoch == Some(admission.epoch())
&& hold.delivery_command_id.as_deref() == Some(admission.command_id())
})
}
pub type ReviewConfigSource = Arc<dyn Fn() -> ReviewConfig + Send + Sync>;
pub trait ReviewEnvironment: Send + Sync {
fn check(&self, session_id: &str, profile: &str) -> Result<(), String>;
fn stage(
&self,
session_id: &str,
profile: &str,
generation: u64,
mcp_servers: &[hel::hel_worker_launch::ReviewMcpServer],
dispatch_tool: bool,
) -> Result<hel::hel_worker_launch::ReviewerLaunchConfig, String>;
fn load_state(&self, session_id: &str) -> Result<TurnReviewState, String>;
fn save_state(&self, session_id: &str, state: &TurnReviewState) -> Result<(), String>;
fn clear_interrupted(&self) -> Result<Vec<String>, String>;
}
#[derive(Debug, Default)]
pub struct ControllerEnvironment;
impl ReviewEnvironment for ControllerEnvironment {
fn check(&self, session_id: &str, profile: &str) -> Result<(), String> {
let controller =
crate::hel_controller::Controller::load().map_err(|error| format!("{error:#}"))?;
if !controller.config.profiles.contains_key(profile) {
return Err(format!(
"turn review needs a reviewer: [review] profile {profile:?} is not a profile in config.toml"
));
}
validate_reviewer_assignment(
session_id,
controller.state.sessions.get(session_id),
profile,
)
}
fn stage(
&self,
session_id: &str,
profile: &str,
generation: u64,
mcp_servers: &[hel::hel_worker_launch::ReviewMcpServer],
dispatch_tool: bool,
) -> Result<hel::hel_worker_launch::ReviewerLaunchConfig, String> {
let controller =
crate::hel_controller::Controller::load().map_err(|error| format!("{error:#}"))?;
controller
.stage_reviewer_profile_with_mcp(
session_id,
profile,
generation,
mcp_servers,
dispatch_tool,
)
.map_err(|error| format!("{error:#}"))
}
fn load_state(&self, session_id: &str) -> Result<TurnReviewState, String> {
hel::hel_database::turn_review_state(session_id).map_err(|error| format!("{error:#}"))
}
fn save_state(&self, session_id: &str, state: &TurnReviewState) -> Result<(), String> {
hel::hel_database::save_turn_review_state(session_id, state)
.map_err(|error| format!("{error:#}"))
}
fn clear_interrupted(&self) -> Result<Vec<String>, String> {
hel::hel_database::clear_interrupted_turn_reviews().map_err(|error| format!("{error:#}"))
}
}
pub(crate) fn validate_reviewer_assignment(
session_id: &str,
session: Option<&hel::hel_state::SessionRecord>,
profile: &str,
) -> Result<(), String> {
let Some(session) = session else {
return Err(format!(
"session {session_id:?} is not in the controller store"
));
};
if session.archived {
return Err("this session is archived".to_owned());
}
if session.last_profile == profile {
return Err(format!(
"turn review profile {profile:?} is also this session's primary profile; choose a different [review] profile"
));
}
Ok(())
}
#[derive(Clone)]
pub struct TurnReviewHost {
events: mpsc::UnboundedSender<HostEvent>,
shared: Arc<HostShared>,
}
struct HostShared {
views: Mutex<BTreeMap<String, RuntimeReviewView>>,
changed: Arc<dyn Fn() + Send + Sync>,
shutdown: tokio::sync::OnceCell<Result<(), String>>,
}
impl std::fmt::Debug for TurnReviewHost {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("TurnReviewHost")
}
}
impl TurnReviewHost {
#[must_use]
pub fn spawn(control: SessionManagerControl, config: ReviewConfigSource) -> Self {
Self::spawn_notifying(control, config, Arc::new(|| {}))
}
#[must_use]
pub fn spawn_notifying(
control: SessionManagerControl,
config: ReviewConfigSource,
changed: Arc<dyn Fn() + Send + Sync>,
) -> Self {
Self::spawn_in_notifying(control, config, Arc::new(ControllerEnvironment), changed)
}
#[must_use]
pub fn spawn_in(
control: SessionManagerControl,
config: ReviewConfigSource,
environment: Arc<dyn ReviewEnvironment>,
) -> Self {
Self::spawn_in_notifying(control, config, environment, Arc::new(|| {}))
}
#[must_use]
fn spawn_in_notifying(
control: SessionManagerControl,
config: ReviewConfigSource,
environment: Arc<dyn ReviewEnvironment>,
changed: Arc<dyn Fn() + Send + Sync>,
) -> Self {
let (events, receiver) = mpsc::unbounded_channel();
let (persistence, persistence_receiver) = mpsc::unbounded_channel();
let shared = Arc::new(HostShared {
views: Mutex::default(),
changed,
shutdown: tokio::sync::OnceCell::new(),
});
let host = Self {
events: events.clone(),
shared: shared.clone(),
};
let persistence_task = tokio::spawn(persistence_loop(
environment.clone(),
events.clone(),
persistence_receiver,
));
persistence
.send(PersistenceRequest::SweepInterrupted)
.expect("new review persistence lane accepts its initial sweep");
tokio::spawn(host_loop(
HostState {
control,
config,
environment,
shared,
events,
persistence: Some(persistence),
persistence_task: Some(persistence_task),
reviews: BTreeMap::new(),
preparing: BTreeSet::new(),
pending_open: BTreeMap::new(),
closing: BTreeSet::new(),
awaiting_forward_persistence: BTreeMap::new(),
next_epoch: 0,
sessions: BTreeMap::new(),
missing_reviewer_reported: BTreeSet::new(),
recovery_candidates: BTreeSet::new(),
recovery_in_flight: BTreeSet::new(),
},
receiver,
));
host
}
pub fn observe(&self, session_id: &str, view: &ManagedSessionView) {
let _ = self.events.send(HostEvent::View {
session_id: session_id.to_owned(),
snapshot: view
.snapshot
.as_ref()
.map(|snapshot| Box::new(snapshot.materialized.clone())),
prompt_driven: view
.snapshot
.as_ref()
.is_some_and(|snapshot| snapshot.operational.active_prompt.is_some()),
});
}
pub async fn start(&self, session_id: &str, manual: bool) -> Result<(), StartRefusal> {
let (reply, answer) = oneshot::channel();
self.events
.send(HostEvent::Start {
session_id: session_id.to_owned(),
manual,
reply: Some(reply),
})
.map_err(|_| StartRefusal("the review host stopped".to_owned()))?;
answer
.await
.map_err(|_| StartRefusal("the review host stopped".to_owned()))?
}
pub async fn resolve(&self, session_id: &str, resolution: Resolution) -> Result<(), String> {
let (reply, answer) = oneshot::channel();
self.events
.send(HostEvent::Resolve {
session_id: session_id.to_owned(),
resolution,
reply,
})
.map_err(|_| "the review host stopped".to_owned())?;
answer
.await
.map_err(|_| "the review host stopped".to_owned())?
}
pub async fn shutdown(&self) -> Result<(), String> {
self.shared
.shutdown
.get_or_init(|| async {
let (reply, answer) = oneshot::channel();
self.events
.send(HostEvent::Shutdown { reply })
.map_err(|_| "the review host stopped".to_owned())?;
answer
.await
.map_err(|_| "the review host stopped during shutdown".to_owned())?
})
.await
.clone()
}
#[must_use]
pub fn views(&self) -> Vec<RuntimeReviewView> {
self.shared
.views
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.values()
.cloned()
.collect()
}
#[must_use]
pub fn view(&self, session_id: &str) -> Option<RuntimeReviewView> {
self.shared
.views
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(session_id)
.cloned()
}
#[must_use]
pub fn refuses_prompt(&self, session_id: &str) -> bool {
prompt_refusal(session_id).is_some()
}
}
enum HostEvent {
View {
session_id: String,
snapshot: Option<Box<MaterializedSession>>,
prompt_driven: bool,
},
Start {
session_id: String,
manual: bool,
reply: Option<oneshot::Sender<Result<(), StartRefusal>>>,
},
Prepared {
session_id: String,
manual: bool,
reply: Option<oneshot::Sender<Result<(), StartRefusal>>>,
prepared: Result<Prepared, StartRefusal>,
},
RecoveryPrepared {
session_id: String,
prepared: Result<Option<Prepared>, String>,
},
StateSaved {
session_id: String,
completion: PersistenceCompletion,
result: Result<(), String>,
},
Step {
session_id: String,
epoch: u64,
step: ReviewStep,
},
Resolve {
session_id: String,
resolution: Resolution,
reply: oneshot::Sender<Result<(), String>>,
},
Interrupted { interrupted: Vec<String> },
Shutdown {
reply: oneshot::Sender<Result<(), String>>,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum PersistenceCompletion {
Open,
Forward,
Close,
}
enum PersistenceRequest {
SweepInterrupted,
Save {
session_id: String,
state: Box<TurnReviewState>,
completion: Option<PersistenceCompletion>,
},
ClearActive {
reply: oneshot::Sender<Result<(), String>>,
},
}
enum ReviewStep {
Delta(Result<Vec<hel::hel_worker::RepoDelta>, String>),
Analysis(Result<String, String>),
RoleStarted {
role: String,
result: Result<(), String>,
},
RolePrompted {
role: String,
result: Result<(), String>,
},
PrimaryPrompted(Result<(), String>),
RoleEvents {
role: String,
result: Result<Vec<RelayEvent>, String>,
},
Dispatches(Result<Vec<hel::hel_review::lanes::ReviewSubagentRequest>, String>),
}
struct Prepared {
state: TurnReviewState,
reviewer: ReviewerIdentity,
tier: ReviewTier,
materialized: Box<MaterializedSession>,
resume_forward: Option<PendingForward>,
}
struct PendingOpen {
epoch: u64,
manual: bool,
reply: Option<oneshot::Sender<Result<(), StartRefusal>>>,
prepared: Prepared,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct ReviewerIdentity {
profile: String,
model: Option<String>,
effort: Option<String>,
}
struct ReviewSlot {
epoch: u64,
driver: TurnReviewDriver,
roles: BTreeMap<String, RoleTranscript>,
reviewer: ReviewerIdentity,
state: TurnReviewState,
generation: u64,
}
#[derive(Default)]
struct RoleTranscript {
session: Option<MaterializedSession>,
cursor_ordinal: u64,
cursor_digest: String,
}
impl RoleTranscript {
fn apply(&mut self, session_id: &str, events: &[RelayEvent]) {
let session = self
.session
.get_or_insert_with(|| MaterializedSession::empty(session_id));
for event in events {
let Ok(projected) = hel::hel_projection::project_relay_event(session, event) else {
continue;
};
if hel::hel_projection::apply_committed_projection_event(
session,
event,
projected.mutation,
)
.is_err()
{
continue;
}
self.cursor_ordinal = event.ordinal;
self.cursor_digest.clone_from(&event.digest);
}
}
fn latest_answer(&self) -> Option<String> {
let session = self.session.as_ref()?;
session
.transcript
.iter()
.rev()
.find(|item| item.is_nonempty_agent_message())
.and_then(|item| {
let hel::hel_state::TranscriptBody::Agent { chunks, .. } = &item.body else {
return None;
};
Some(hel::hel_transcript::materialized_chunks_text(chunks))
})
.filter(|text| !text.trim().is_empty())
}
}
struct HostState {
control: SessionManagerControl,
config: ReviewConfigSource,
environment: Arc<dyn ReviewEnvironment>,
shared: Arc<HostShared>,
events: mpsc::UnboundedSender<HostEvent>,
persistence: Option<mpsc::UnboundedSender<PersistenceRequest>>,
persistence_task: Option<tokio::task::JoinHandle<()>>,
reviews: BTreeMap<String, ReviewSlot>,
preparing: BTreeSet<String>,
pending_open: BTreeMap<String, PendingOpen>,
closing: BTreeSet<String>,
awaiting_forward_persistence: BTreeMap<String, Vec<ReviewRequest>>,
next_epoch: u64,
sessions: BTreeMap<String, SessionWatch>,
missing_reviewer_reported: BTreeSet<String>,
recovery_candidates: BTreeSet<String>,
recovery_in_flight: BTreeSet<String>,
}
struct SessionWatch {
execution: MaterializedExecutionState,
prompt_driven: bool,
materialized: Option<Box<MaterializedSession>>,
}
async fn host_loop(mut state: HostState, mut events: mpsc::UnboundedReceiver<HostEvent>) {
while let Some(event) = events.recv().await {
if state.handle(event).await {
break;
}
}
}
async fn persistence_loop(
environment: Arc<dyn ReviewEnvironment>,
events: mpsc::UnboundedSender<HostEvent>,
mut requests: mpsc::UnboundedReceiver<PersistenceRequest>,
) {
while let Some(request) = requests.recv().await {
match request {
PersistenceRequest::SweepInterrupted => {
let environment = environment.clone();
match tokio::task::spawn_blocking(move || environment.clear_interrupted()).await {
Ok(Ok(interrupted)) if !interrupted.is_empty() => {
let _ = events.send(HostEvent::Interrupted { interrupted });
}
Ok(Ok(_)) => {}
Ok(Err(error)) => {
tracing::warn!(%error, "could not clear interrupted reviews");
}
Err(error) => {
tracing::warn!(%error, "the interrupted-review sweep did not run");
}
}
}
PersistenceRequest::Save {
session_id,
state,
completion,
} => {
let environment = environment.clone();
let owner = session_id.clone();
let result =
tokio::task::spawn_blocking(move || environment.save_state(&owner, &state))
.await
.map_err(|error| format!("review state persistence task stopped: {error}"))
.and_then(|result| result);
if let Some(completion) = completion {
let _ = events.send(HostEvent::StateSaved {
session_id,
completion,
result,
});
} else if let Err(error) = result {
tracing::warn!(
session_id = %session_id,
%error,
"could not record how far this session has been reviewed"
);
}
}
PersistenceRequest::ClearActive { reply } => {
let environment = environment.clone();
let result = tokio::task::spawn_blocking(move || {
environment.clear_interrupted().map(|_| ())
})
.await
.map_err(|error| format!("review shutdown persistence task stopped: {error}"))
.and_then(|result| result);
let _ = reply.send(result);
}
}
}
}
impl HostState {
async fn handle(&mut self, event: HostEvent) -> bool {
match event {
HostEvent::View {
session_id,
snapshot,
prompt_driven,
} => self.observe(session_id, snapshot, prompt_driven).await,
HostEvent::Start {
session_id,
manual,
reply,
} => self.begin(session_id, manual, reply),
HostEvent::Prepared {
session_id,
manual,
reply,
prepared,
} => self.prepared(session_id, manual, reply, prepared),
HostEvent::RecoveryPrepared {
session_id,
prepared,
} => self.recovery_prepared(session_id, prepared),
HostEvent::StateSaved {
session_id,
completion,
result,
} => self.state_saved(session_id, completion, result),
HostEvent::Resolve {
session_id,
resolution,
reply,
} => {
let answer = self.resolve(&session_id, resolution);
let _ = reply.send(answer);
}
HostEvent::Step {
session_id,
epoch,
step,
} => self.step(session_id, epoch, step),
HostEvent::Interrupted { interrupted } => {
for session_id in interrupted {
self.recovery_candidates.insert(session_id.clone());
if self.sessions.get(&session_id).is_some_and(|watch| {
matches!(watch.execution, MaterializedExecutionState::Idle)
}) {
self.begin_recovery(&session_id);
}
}
}
HostEvent::Shutdown { reply } => {
let result = self.shutdown().await;
let _ = reply.send(result);
return true;
}
}
false
}
fn step(&mut self, session_id: String, epoch: u64, step: ReviewStep) {
if self.reviews.get(&session_id).map(|slot| slot.epoch) != Some(epoch) {
return;
}
match step {
ReviewStep::Delta(result) => {
let requests = match result {
Ok(deltas) => self
.reviews
.get_mut(&session_id)
.map(|slot| slot.driver.delta_captured(deltas))
.unwrap_or_default(),
Err(error) => {
self.fail(
&session_id,
format!("the change could not be captured: {error}"),
);
return;
}
};
self.run(&session_id, requests);
}
ReviewStep::Analysis(result) => {
let requests = self
.reviews
.get_mut(&session_id)
.map(|slot| slot.driver.analysis_completed(result))
.unwrap_or_default();
self.run(&session_id, requests);
}
ReviewStep::RoleStarted { role, result } => {
let requests = match self.reviews.get_mut(&session_id) {
Some(slot) => match result {
Ok(()) => slot.driver.role_started(&role),
Err(error) if hel::hel_review::lanes::lane_by_id(&role).is_some() => {
slot.driver.lane_failed(&role, error)
}
Err(error)
if matches!(
slot.driver.phase(),
TurnReviewPhase::Forwarding { .. }
) =>
{
tracing::debug!(
session_id = %session_id,
role = %role,
%error,
"ignoring a late reviewer start result during primary handoff"
);
return;
}
Err(error) => {
self.fail(&session_id, error);
return;
}
},
None => return,
};
self.run(&session_id, requests);
}
ReviewStep::RolePrompted { role, result } => {
if let Err(error) = result {
if self.reviews.get(&session_id).is_some_and(|slot| {
matches!(slot.driver.phase(), TurnReviewPhase::Forwarding { .. })
}) {
tracing::debug!(
session_id = %session_id,
role = %role,
%error,
"ignoring a late reviewer-role result during primary handoff"
);
return;
}
self.fail(
&session_id,
format!("reviewing role {role:?} could not be prompted: {error}"),
);
}
}
ReviewStep::PrimaryPrompted(result) => {
let requests = match result {
Ok(()) => self
.reviews
.get_mut(&session_id)
.map(|slot| slot.driver.forward_succeeded())
.unwrap_or_default(),
Err(error) => self
.reviews
.get_mut(&session_id)
.map(|slot| slot.driver.forward_failed(error))
.unwrap_or_default(),
};
self.run(&session_id, requests);
}
ReviewStep::RoleEvents { role, result } => self.role_events(session_id, role, result),
ReviewStep::Dispatches(result) => {
let requests = match result {
Ok(requests) => self
.reviews
.get_mut(&session_id)
.map(|slot| slot.driver.lanes_dispatched(requests))
.unwrap_or_default(),
Err(error) => {
if self.reviews.get(&session_id).is_some_and(|slot| {
matches!(slot.driver.phase(), TurnReviewPhase::Forwarding { .. })
}) {
tracing::debug!(
session_id = %session_id,
%error,
"ignoring a late lane dispatch result during primary handoff"
);
return;
}
self.fail(
&session_id,
format!(
"the review could not collect the supervisor's specialists: {error}"
),
);
return;
}
};
self.run(&session_id, requests);
}
}
}
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);
}
fn begin(
&mut self,
session_id: String,
manual: bool,
reply: Option<oneshot::Sender<Result<(), StartRefusal>>>,
) {
if crate::hel_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,
});
});
}
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,
});
});
}
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);
}
}
}
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
}
fn prepared(
&mut self,
session_id: String,
manual: bool,
reply: Option<oneshot::Sender<Result<(), StartRefusal>>>,
prepared: Result<Prepared, StartRefusal>,
) {
let prepared = match prepared {
Ok(prepared) => prepared,
Err(refusal) => {
self.preparing.remove(&session_id);
release_prompts(&session_id);
answer(reply, Err(refusal));
return;
}
};
if self.reviews.contains_key(&session_id) {
self.preparing.remove(&session_id);
release_prompts(&session_id);
let refusal = StartRefusal("a review is already open".to_owned());
answer(reply, Err(refusal));
return;
}
self.next_epoch = self.next_epoch.saturating_add(1);
let epoch = self.next_epoch;
let mut prepared = prepared;
prepared.state.active = Some(format!("review-{epoch}"));
let state = prepared.state.clone();
self.pending_open.insert(
session_id.clone(),
PendingOpen {
epoch,
manual,
reply,
prepared,
},
);
if let Err(error) =
self.persist(session_id.clone(), state, Some(PersistenceCompletion::Open))
{
let pending = self
.pending_open
.remove(&session_id)
.expect("pending review was just inserted");
let retry_recovery = pending.prepared.resume_forward.is_some();
self.preparing.remove(&session_id);
if retry_recovery {
self.recovery_candidates.insert(session_id.clone());
}
release_prompts(&session_id);
answer(
pending.reply,
Err(StartRefusal(format!(
"could not record the active review: {error}"
))),
);
}
}
fn state_saved(
&mut self,
session_id: String,
completion: PersistenceCompletion,
result: Result<(), String>,
) {
match completion {
PersistenceCompletion::Open => {
let Some(pending) = self.pending_open.remove(&session_id) else {
return;
};
self.preparing.remove(&session_id);
if let Err(error) = result {
if pending.prepared.resume_forward.is_some() {
self.recovery_candidates.insert(session_id.clone());
}
release_prompts(&session_id);
answer(
pending.reply,
Err(StartRefusal(format!(
"could not record the active review: {error}"
))),
);
return;
}
let seed = seed_from_session(
&pending.prepared.materialized,
pending.prepared.tier,
&pending.prepared.state,
if pending.manual {
"manual"
} else {
"automatic"
},
);
let (driver, requests) =
if let Some(pending) = pending.prepared.resume_forward.clone() {
let command_id = pending.command_id.clone();
let (mut driver, _) = TurnReviewDriver::resume_forward(seed, pending);
let requests = driver.forward(command_id);
(driver, requests)
} else {
TurnReviewDriver::start(seed)
};
self.reviews.insert(
session_id.clone(),
ReviewSlot {
epoch: pending.epoch,
driver,
roles: BTreeMap::new(),
reviewer: pending.prepared.reviewer,
state: pending.prepared.state,
generation: 0,
},
);
answer(pending.reply, Ok(()));
self.run(&session_id, requests);
}
PersistenceCompletion::Forward => {
let Some(requests) = self.awaiting_forward_persistence.remove(&session_id) else {
return;
};
if let Err(error) = result {
if let Some(slot) = self.reviews.get_mut(&session_id) {
slot.driver.forward_failed(format!(
"the handoff could not be recorded durably: {error}"
));
}
self.publish(&session_id);
return;
}
self.run(&session_id, requests);
}
PersistenceCompletion::Close => {
self.closing.remove(&session_id);
if let Err(error) = result {
tracing::warn!(
session_id = %session_id,
%error,
"could not clear the active review marker"
);
}
let notice = self.reviews.get(&session_id).and_then(|slot| {
resolution_notice(slot.driver.phase(), slot.driver.last_verdict())
});
self.reviews.remove(&session_id);
release_prompts(&session_id);
if let Some(notice) = notice {
self.record_notice(&session_id, notice);
}
self.publish(&session_id);
}
}
}
fn persist(
&self,
session_id: String,
state: TurnReviewState,
completion: Option<PersistenceCompletion>,
) -> Result<(), String> {
self.persistence
.as_ref()
.ok_or_else(|| "the review persistence lane stopped".to_owned())?
.send(PersistenceRequest::Save {
session_id,
state: Box::new(state),
completion,
})
.map_err(|_| "the review persistence lane stopped".to_owned())
}
fn resolve(&mut self, session_id: &str, resolution: Resolution) -> Result<(), String> {
let (requests, pending_state) = {
let Some(slot) = self.reviews.get_mut(session_id) else {
return Err("no review is open for that session".to_owned());
};
let requests = match resolution {
Resolution::Forwarded => {
if !slot.driver.can_forward() {
return Err("there are no findings to forward".to_owned());
}
slot.driver.forward(
new_command_id("review-forward").map_err(|error| format!("{error:#}"))?,
)
}
Resolution::Dismissed => {
if slot.driver.verdict().is_none() {
return Err("the review has not reached a verdict yet".to_owned());
}
slot.driver.dismiss()
}
Resolution::Cancelled => slot.driver.cancel(),
Resolution::NothingToReview | Resolution::CoverageStarted => {
return Err("that is not a resolution a surface can ask for".to_owned());
}
};
if requests.is_empty() {
return Err("the review could not be resolved that way".to_owned());
}
let pending_state = if resolution == Resolution::Forwarded {
let Some(pending) = slot.driver.pending_forward() else {
return Err("the review handoff has no durable findings".to_owned());
};
slot.state.pending_forward = Some(pending);
Some(slot.state.clone())
} else {
None
};
(requests, pending_state)
};
if let Some(state) = pending_state {
self.awaiting_forward_persistence
.insert(session_id.to_owned(), requests);
if let Err(error) = self.persist(
session_id.to_owned(),
state,
Some(PersistenceCompletion::Forward),
) {
self.awaiting_forward_persistence.remove(session_id);
if let Some(slot) = self.reviews.get_mut(session_id) {
slot.driver.forward_failed(format!(
"the handoff could not be recorded durably: {error}"
));
}
self.publish(session_id);
return Err(error);
}
self.publish(session_id);
return Ok(());
}
self.run(session_id, requests);
Ok(())
}
fn fail(&mut self, session_id: &str, message: impl Into<String>) {
let Some(slot) = self.reviews.get_mut(session_id) else {
return;
};
slot.state.active = None;
let state = slot.state.clone();
let requests = slot.driver.request_failed(message);
release_prompts(session_id);
if let Err(error) = self.persist(session_id.to_owned(), state, None) {
tracing::warn!(session_id, %error, "could not queue failed review persistence");
}
self.run(session_id, requests);
}
fn run(&mut self, session_id: &str, requests: Vec<ReviewRequest>) {
for request in requests {
self.run_one(session_id, request);
}
self.publish(session_id);
}
fn run_one(&mut self, session_id: &str, request: ReviewRequest) {
match request {
ReviewRequest::CaptureDelta { baselines } => {
self.review_step(
session_id,
ReviewerAction::CaptureDelta { baselines },
|outcome| {
ReviewStep::Delta(match outcome {
Ok(ReviewerOutcome::Delta { repositories }) => Ok(repositories),
other => Err(unexpected(other)),
})
},
);
}
ReviewRequest::AnalyzeDelta { repositories } => {
self.review_step(
session_id,
ReviewerAction::AnalyzeDelta { repositories },
|outcome| {
ReviewStep::Analysis(match outcome {
Ok(ReviewerOutcome::ChangedFunctions { packet }) => Ok(packet),
other => Err(unexpected(other)),
})
},
);
}
ReviewRequest::StartRole { role, fresh } => self.start_role(session_id, role, fresh),
ReviewRequest::PromptRole {
role,
command_id,
prompt,
} => {
self.prompt_role(session_id, &role, command_id, prompt);
self.poll_role(session_id, &role, Duration::ZERO);
}
ReviewRequest::PromptPrimary { command_id, prompt } => {
self.prompt_primary(session_id, command_id, prompt);
}
ReviewRequest::PauseRole { role } => {
let session_id = session_id.to_owned();
self.spawn_reviewer(
session_id.clone(),
Some(role),
ReviewerAction::Pause,
move |outcome| {
if let Err(error) = outcome {
tracing::debug!(
session_id = %session_id,
%error,
"pausing a review role failed"
);
}
None
},
);
}
ReviewRequest::AdvanceBaseline {
trees,
reviewed_through_ordinal,
} => {
if let Some(slot) = self.reviews.get_mut(session_id) {
slot.state.baselines = trees.clone();
slot.state.reviewed_through_ordinal = reviewed_through_ordinal;
slot.state.pending_forward = None;
let state = slot.state.clone();
if let Err(error) = self.persist(session_id.to_owned(), state, None) {
tracing::warn!(session_id, %error, "could not queue review baseline persistence");
}
}
let session_id = session_id.to_owned();
self.spawn_reviewer(
session_id.clone(),
None,
ReviewerAction::AdvanceBaseline { trees },
move |outcome| {
if let Err(error) = outcome {
tracing::debug!(
session_id = %session_id,
%error,
"the review baseline ref could not be pinned"
);
}
None
},
);
}
ReviewRequest::RecordPriorReview { prior } => {
if let Some(slot) = self.reviews.get_mut(session_id) {
slot.state.prior_review = Some(prior);
}
}
ReviewRequest::ClearPriorReview => {
if let Some(slot) = self.reviews.get_mut(session_id) {
slot.state.prior_review = None;
let state = slot.state.clone();
if let Err(error) = self.persist(session_id.to_owned(), state, None) {
tracing::warn!(session_id, %error, "could not queue prior review cleanup");
}
}
}
ReviewRequest::Close => {
if self.closing.contains(session_id) {
return;
}
if let Some(slot) = self.reviews.get_mut(session_id) {
slot.state.active = None;
if matches!(
slot.driver.phase(),
TurnReviewPhase::Resolved(Resolution::Cancelled)
) {
slot.state.pending_forward = None;
}
let state = slot.state.clone();
self.closing.insert(session_id.to_owned());
release_prompts(session_id);
if let Err(error) = self.persist(
session_id.to_owned(),
state,
Some(PersistenceCompletion::Close),
) {
tracing::warn!(session_id, %error, "could not queue review close persistence");
self.closing.remove(session_id);
self.reviews.remove(session_id);
release_prompts(session_id);
}
}
}
}
}
fn start_role(&mut self, session_id: &str, role: String, fresh: bool) {
let fresh_generation = if fresh {
match next_review_generation() {
Ok(generation) => Some(generation),
Err(error) => {
self.fail(
session_id,
format!("the reviewer could not allocate a fresh conversation: {error}"),
);
return;
}
}
} else {
None
};
let Some(slot) = self.reviews.get_mut(session_id) else {
return;
};
if let Some(generation) = fresh_generation {
slot.generation = generation;
}
let generation = slot.generation;
let epoch = slot.epoch;
let reviewer = slot.reviewer.clone();
let repositories = slot.driver.repository_roots();
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 result = launch_role(
&control,
&environment,
&session_id,
&role,
&reviewer,
generation,
&repositories,
)
.await;
let _ = events.send(HostEvent::Step {
session_id,
epoch,
step: ReviewStep::RoleStarted { role, result },
});
});
}
fn prompt_role(&mut self, session_id: &str, role: &str, command_id: String, prompt: String) {
let Some(epoch) = self.reviews.get(session_id).map(|slot| slot.epoch) else {
return;
};
let owner = session_id.to_owned();
let result_role = role.to_owned();
self.spawn_reviewer(
session_id.to_owned(),
Some(role.to_owned()),
ReviewerAction::Submit {
command_id,
command: prompt_command(prompt),
},
move |outcome| {
let result = match outcome {
Ok(ReviewerOutcome::Accepted { .. }) => Ok(()),
other => Err(unexpected(other)),
};
Some(HostEvent::Step {
session_id: owner,
epoch,
step: ReviewStep::RolePrompted {
role: result_role,
result,
},
})
},
);
}
fn prompt_primary(&mut self, session_id: &str, command_id: String, prompt: String) {
let Some(epoch) = self.reviews.get(session_id).map(|slot| slot.epoch) else {
return;
};
let Some(admission) = admit_review_delivery(session_id, epoch, &command_id) else {
self.step(
session_id.to_owned(),
epoch,
ReviewStep::PrimaryPrompted(Err(
"the review handoff admission is no longer valid".to_owned()
)),
);
return;
};
let control = self.control.clone();
let events = self.events.clone();
let session_id = session_id.to_owned();
tokio::spawn(async move {
let submitted = async {
let handle = control
.session(session_id.clone())
.await
.map_err(|error| format!("{error:#}"))?;
handle
.submit_review_delivery(admission, prompt_command(prompt))
.await
.map(|_| ())
.map_err(|error| format!("{error:#}"))
}
.await;
let _ = events.send(HostEvent::Step {
session_id,
epoch,
step: ReviewStep::PrimaryPrompted(submitted),
});
});
}
fn poll_role(&mut self, session_id: &str, role: &str, delay: Duration) {
let Some(slot) = self.reviews.get_mut(session_id) else {
return;
};
let transcript = slot.roles.entry(role.to_owned()).or_default();
let after_ordinal = transcript.cursor_ordinal;
let after_digest = if transcript.cursor_digest.is_empty() {
hel::hel_worker::RELAY_EVENT_GENESIS_DIGEST.to_owned()
} else {
transcript.cursor_digest.clone()
};
let epoch = slot.epoch;
let control = self.control.clone();
let events = self.events.clone();
let session_id = session_id.to_owned();
let role = role.to_owned();
tokio::spawn(async move {
if !delay.is_zero() {
tokio::time::sleep(delay).await;
}
let result = reviewer_action(
&control,
&session_id,
Some(role.clone()),
ReviewerAction::Attach {
after_ordinal,
after_digest,
},
)
.await;
let result = match result {
Ok(ReviewerOutcome::Attached(attachment)) => Ok(attachment.events),
other => Err(unexpected(other)),
};
let _ = events.send(HostEvent::Step {
session_id,
epoch,
step: ReviewStep::RoleEvents { role, result },
});
});
}
fn role_events(
&mut self,
session_id: String,
role: String,
result: Result<Vec<RelayEvent>, String>,
) {
let events = match result {
Ok(events) => events,
Err(error) => {
if self.reviews.get(&session_id).is_some_and(|slot| {
matches!(slot.driver.phase(), TurnReviewPhase::Forwarding { .. })
}) {
tracing::debug!(
session_id = %session_id,
role = %role,
%error,
"ignoring a late reviewer-role poll during primary handoff"
);
return;
}
self.fail(&session_id, error);
return;
}
};
let Some(slot) = self.reviews.get_mut(&session_id) else {
return;
};
let idle = events.is_empty();
let relay_session = role_session_id(&session_id, &role);
let transcript = slot.roles.entry(role.clone()).or_default();
transcript.apply(&relay_session, &events);
let awaited = slot
.driver
.awaited_commands()
.into_iter()
.find(|(awaited_role, _)| *awaited_role == role)
.map(|(_, command_id)| command_id);
let completed = awaited.as_ref().is_some_and(|awaited| {
events.iter().any(|event| {
matches!(
&event.observation,
RelayObservation::CommandCompleted { command_id, outcome }
if command_id == awaited
&& matches!(
outcome,
hel::hel_worker::RelayCommandOutcome::Prompt { .. }
)
)
})
});
let requests = match (completed, awaited) {
(true, Some(awaited)) => {
let answer = slot
.roles
.get(&role)
.and_then(RoleTranscript::latest_answer)
.unwrap_or_default();
let slot = self.reviews.get_mut(&session_id).expect("the slot is open");
slot.driver.role_turn_completed(&awaited, &answer)
}
_ => Vec::new(),
};
self.run(&session_id, requests);
let Some(slot) = self.reviews.get(&session_id) else {
return;
};
if slot.driver.active_roles().contains(&role) {
self.poll_role(
&session_id,
&role,
if idle {
ROLE_POLL_IDLE_INTERVAL
} else {
Duration::ZERO
},
);
}
if role == SUPERVISOR_ROLE
&& self
.reviews
.get(&session_id)
.is_some_and(|slot| slot.driver.supervisor_running())
{
self.poll_dispatches(&session_id);
}
}
fn poll_dispatches(&mut self, session_id: &str) {
self.review_step(session_id, ReviewerAction::TakeLaneDispatches, |outcome| {
ReviewStep::Dispatches(match outcome {
Ok(ReviewerOutcome::LaneDispatches { requests }) => Ok(requests),
other => Err(unexpected(other)),
})
});
}
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"
);
}
});
}
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),
})
});
}
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 outcome = reviewer_action(&control, &session_id, role, action).await;
if let Some(event) = into_event(outcome) {
let _ = events.send(event);
}
});
}
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 next = slot.view(session_id);
if views.get(session_id) == Some(&next) {
false
} else {
views.insert(session_id.to_owned(), next);
true
}
}
None => views.remove(session_id).is_some(),
};
drop(views);
if changed {
(self.shared.changed)();
}
}
async fn shutdown(&mut self) -> Result<(), String> {
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 {
slot.state.active = None;
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");
}
}
self.preparing.clear();
self.closing.clear();
self.reviews.clear();
for session_id in &session_ids {
release_prompts(session_id);
self.publish(session_id);
}
let clear_result = match lane {
Some(lane) => {
let (reply, cleared) = oneshot::channel();
let sent = lane
.send(PersistenceRequest::ClearActive { reply })
.map_err(|_| "the review persistence lane stopped during shutdown".to_owned());
drop(lane);
match sent {
Ok(()) => cleared.await.map_err(|_| {
"the review persistence lane stopped before cleanup".to_owned()
})?,
Err(error) => Err(error),
}
}
None => Err("the review persistence lane already stopped".to_owned()),
};
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)
}
}
impl ReviewSlot {
fn view(&self, session_id: &str) -> RuntimeReviewView {
let verdict = match self.driver.phase() {
TurnReviewPhase::Forwarding { synthesis, .. } => Some(VerdictView {
kind: VerdictKind::Findings,
text: synthesis.clone(),
allowed: if matches!(
self.driver.phase(),
TurnReviewPhase::Forwarding { error: Some(_), .. }
) {
vec![Resolution::Forwarded, Resolution::Cancelled]
} else {
Vec::new()
},
}),
_ => self.driver.verdict().map(|verdict| match verdict {
ReviewVerdict::Clean => VerdictView {
kind: VerdictKind::Clean,
text: String::new(),
allowed: Vec::new(),
},
ReviewVerdict::Findings { synthesis, .. } => VerdictView {
kind: VerdictKind::Findings,
text: synthesis.clone(),
allowed: vec![
Resolution::Forwarded,
Resolution::Dismissed,
Resolution::Cancelled,
],
},
ReviewVerdict::Failed { reason } => VerdictView {
kind: VerdictKind::Failed,
text: reason.clone(),
allowed: vec![Resolution::Dismissed, Resolution::Cancelled],
},
}),
};
RuntimeReviewView {
session_id: session_id.to_owned(),
tier: self.driver.tier(),
phase: self.driver.phase().clone(),
roles: self.driver.roles(),
status: self.driver.status().to_owned(),
verdict,
}
}
}
fn answer(
reply: Option<oneshot::Sender<Result<(), StartRefusal>>>,
result: Result<(), StartRefusal>,
) {
if let Some(reply) = reply {
let _ = reply.send(result);
}
}
fn unexpected(outcome: Result<ReviewerOutcome, String>) -> String {
match outcome {
Ok(other) => format!("unexpected reviewer response {other:?}"),
Err(error) => error,
}
}
fn prompt_command(prompt: String) -> RelayCommand {
RelayCommand::Prompt {
prompt: vec![agent_client_protocol::schema::v1::ContentBlock::Text(
agent_client_protocol::schema::v1::TextContent::new(prompt),
)],
}
}
#[must_use]
pub fn role_session_id(primary_session_id: &str, role: &str) -> String {
if role == hel::hel_review::driver::REVIEWER_ROLE {
format!("{primary_session_id}-reviewer")
} else {
format!("{primary_session_id}-review-{role}")
}
}
async fn reviewer_action(
control: &SessionManagerControl,
session_id: &str,
role: Option<String>,
action: ReviewerAction,
) -> Result<ReviewerOutcome, String> {
let handle: ManagedSessionHandle = control
.session(session_id.to_owned())
.await
.map_err(|error| format!("{error:#}"))?;
handle
.reviewer_as(role, action)
.await
.map_err(|error| format!("{error:#}"))
}
async fn launch_role(
control: &SessionManagerControl,
environment: &Arc<dyn ReviewEnvironment>,
session_id: &str,
role: &str,
reviewer: &ReviewerIdentity,
generation: u64,
repositories: &[PathBuf],
) -> Result<(), String> {
let lane = hel::hel_review::lanes::lane_by_id(role).is_some();
let mcp_servers = if role == INTENT_ROLE {
Vec::new()
} else {
hel::hel_review::bifrost::review_mcp_servers(
repositories,
if lane {
hel::hel_review::lanes::LANE_BIFROST_TOOLSET
} else {
hel::hel_review::lanes::SUPERVISOR_BIFROST_TOOLSET
},
)
};
let dispatch_tool = role == SUPERVISOR_ROLE;
let staged = {
let session_id = session_id.to_owned();
let profile = reviewer.profile.clone();
let environment = environment.clone();
tokio::task::spawn_blocking(move || {
environment.stage(
&session_id,
&profile,
generation,
&mcp_servers,
dispatch_tool,
)
})
.await
.map_err(|error| format!("staging the reviewer stopped: {error}"))??
};
let mut config = staged;
config.model = reviewer.model.clone();
config.effort = reviewer.effort.clone();
match reviewer_action(
control,
session_id,
Some(role.to_owned()),
ReviewerAction::Start {
config: Box::new(config),
},
)
.await
{
Ok(ReviewerOutcome::Started(_)) => Ok(()),
other => Err(unexpected(other)),
}
}
async fn prepare(
control: &SessionManagerControl,
environment: &Arc<dyn ReviewEnvironment>,
session_id: &str,
reviewer: &ReviewerIdentity,
tier: ReviewTier,
) -> Result<Prepared, StartRefusal> {
let profile = reviewer.profile.clone();
let session = session_id.to_owned();
let environment = environment.clone();
let checked = tokio::task::spawn_blocking(move || -> Result<TurnReviewState, String> {
environment.check(&session, &profile)?;
environment.load_state(&session)
})
.await
.map_err(|error| StartRefusal(format!("preparing the review stopped: {error}")))?;
let state = checked.map_err(StartRefusal)?;
let handle = control
.session(session_id.to_owned())
.await
.map_err(|error| StartRefusal(format!("{error:#}")))?;
match handle.reviewer(ReviewerAction::Status).await {
Ok(ReviewerOutcome::Status(state)) if state.active_prompt.is_some() => {
return Err(StartRefusal(
"the reviewer is busy with a second opinion".to_owned(),
));
}
Ok(_) => {}
Err(error) => return Err(StartRefusal(format!("{error:#}"))),
}
let view = handle.view();
if !view.connected {
return Err(StartRefusal("this session is not connected".to_owned()));
}
let Some(snapshot) = view.snapshot else {
return Err(StartRefusal(
"this session has no transcript yet".to_owned(),
));
};
if !matches!(
snapshot.materialized.execution,
MaterializedExecutionState::Idle
) {
return Err(StartRefusal(
"a review runs between turns; this one is still working".to_owned(),
));
}
if !snapshot.materialized.queued_prompts.is_empty() {
return Err(StartRefusal(
"prompts are queued; the review waits for them".to_owned(),
));
}
Ok(Prepared {
state,
reviewer: reviewer.clone(),
tier,
materialized: Box::new(snapshot.materialized),
resume_forward: None,
})
}
async fn prepare_recovery(
control: &SessionManagerControl,
environment: &Arc<dyn ReviewEnvironment>,
session_id: &str,
) -> Result<Option<Prepared>, String> {
let session = session_id.to_owned();
let environment = environment.clone();
let state = tokio::task::spawn_blocking(move || environment.load_state(&session))
.await
.map_err(|error| format!("loading the pending review handoff stopped: {error}"))??;
let Some(pending) = state.pending_forward.clone() else {
return Ok(None);
};
let handle = control
.session(session_id.to_owned())
.await
.map_err(|error| format!("{error:#}"))?;
let view = handle.view();
if !view.connected {
return Err("the primary session is not connected".to_owned());
}
let Some(snapshot) = view.snapshot else {
return Err("the primary session has no transcript yet".to_owned());
};
if !matches!(
snapshot.materialized.execution,
MaterializedExecutionState::Idle
) {
return Err("the primary session is still working".to_owned());
}
if !snapshot.materialized.queued_prompts.is_empty() {
return Err("prompts are queued; the pending handoff waits for them".to_owned());
}
Ok(Some(Prepared {
state,
reviewer: ReviewerIdentity {
profile: String::new(),
model: None,
effort: None,
},
tier: ReviewTier::Quick,
materialized: Box::new(snapshot.materialized),
resume_forward: Some(pending),
}))
}
#[must_use]
pub fn resolution_notice(
phase: &TurnReviewPhase,
last_verdict: Option<&ReviewVerdict>,
) -> Option<String> {
let TurnReviewPhase::Resolved(resolution) = phase else {
return None;
};
Some(match resolution {
Resolution::Forwarded => "Review findings sent to the agent".to_owned(),
Resolution::Dismissed => match last_verdict {
Some(ReviewVerdict::Clean) => "Review complete: no material findings".to_owned(),
Some(ReviewVerdict::Failed { .. }) => {
"Review failed; the change stays unreviewed".to_owned()
}
_ => "Review dismissed".to_owned(),
},
Resolution::Cancelled => match last_verdict {
Some(ReviewVerdict::Failed { .. }) => {
"Review failed; the change stays unreviewed".to_owned()
}
_ => "Review cancelled".to_owned(),
},
Resolution::NothingToReview => "Nothing to review: the turn changed no files".to_owned(),
Resolution::CoverageStarted => {
"Review coverage starts here; the next completed turn is reviewed".to_owned()
}
})
}
fn seed_from_session(
session: &MaterializedSession,
tier: ReviewTier,
state: &TurnReviewState,
_trigger: &str,
) -> TurnReviewSeed {
let reviewed_through = state.reviewed_through_ordinal;
let mut task = String::new();
let mut user_messages = Vec::new();
let mut initial_result = String::new();
let mut trajectory = Vec::new();
for item in &session.transcript {
match &item.body {
hel::hel_state::TranscriptBody::User { content } => {
let text = hel::hel_transcript::materialized_content_text(content);
let text = text.trim();
if text.is_empty() {
continue;
}
if hel::hel_second_opinion::is_control_origin_prompt(text) {
continue;
}
task = text.to_owned();
user_messages.push(UserMessage::prompt(text));
if item.position > reviewed_through {
trajectory.push(format!("user: {text}"));
}
}
hel::hel_state::TranscriptBody::Agent { chunks, .. } => {
if !item.is_nonempty_agent_message() {
continue;
}
let text = hel::hel_transcript::materialized_chunks_text(chunks);
let text = text.trim();
if text.is_empty() {
continue;
}
initial_result = text.to_owned();
if item.position > reviewed_through {
trajectory.push(format!("agent: {text}"));
}
}
hel::hel_state::TranscriptBody::Tool { call, .. } => {
let title = call
.get("title")
.and_then(serde_json::Value::as_str)
.unwrap_or_default()
.trim();
if item.position > reviewed_through && !title.is_empty() {
trajectory.push(format!("tool: {title}"));
}
}
_ => {}
}
}
TurnReviewSeed {
tier,
task,
user_messages,
initial_result,
trajectory: trajectory.join("\n"),
baselines: state.baselines.clone(),
through_ordinal: session.applied_event_ordinal,
prior_review: state.prior_review.clone(),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::hel_session_manager::{
RelaySessionTarget, RemoteSessionRequest, RemoteSessionRequests,
spawn_remote_session_manager,
};
use hel::hel_state::{ManagedSessionSnapshot, MaterializedSession};
use hel::hel_worker::{
RELAY_EVENT_FORMAT_V1, RelayCommandOutcome, RelayOperationalState, relay_event_digest,
};
fn session_id(test: &str) -> String {
format!("018f9dd2-a3b4-7c8d-9000-{test}")
}
#[test]
fn review_activity_follows_typed_transitions_without_reading_progress_prose() {
let mut view = RuntimeReviewView {
session_id: "activity".to_owned(),
tier: ReviewTier::Quick,
phase: TurnReviewPhase::LaunchingReviewer,
roles: Vec::new(),
status: "validating configuration".to_owned(),
verdict: None,
};
assert_eq!(view.activity_label(), Some("Reviewing"));
view.phase = TurnReviewPhase::Running {
roles: vec![RoleStatus {
role: hel::hel_review::driver::VALIDATOR_ROLE.to_owned(),
label: "Validator".to_owned(),
state: RoleState::Running,
}],
};
view.status = "checking source".to_owned();
assert_eq!(view.activity_label(), Some("Validating"));
view.phase = TurnReviewPhase::Verdict(ReviewVerdict::Findings {
synthesis: "[P2] app.py:1 -- incorrect bounds".to_owned(),
evidence: Default::default(),
});
assert_eq!(view.activity_label(), Some("Findings"));
view.phase = TurnReviewPhase::Verdict(ReviewVerdict::Failed {
reason: "reviewer unavailable".to_owned(),
});
assert_eq!(view.activity_label(), Some("Review failed"));
view.phase = TurnReviewPhase::Resolved(Resolution::Cancelled);
assert_eq!(view.activity_label(), None);
}
#[test]
fn resolution_notices_keep_the_verdict_context_after_close() {
let resolved_dismissed = TurnReviewPhase::Resolved(Resolution::Dismissed);
assert_eq!(
resolution_notice(&resolved_dismissed, Some(&ReviewVerdict::Clean)),
Some("Review complete: no material findings".to_owned())
);
assert_eq!(
resolution_notice(
&TurnReviewPhase::Resolved(Resolution::Cancelled),
Some(&ReviewVerdict::Failed {
reason: "harness failed".to_owned(),
}),
),
Some("Review failed; the change stays unreviewed".to_owned())
);
assert_eq!(
resolution_notice(
&resolved_dismissed,
Some(&ReviewVerdict::Findings {
synthesis: "[P1] broken".to_owned(),
evidence: Default::default(),
}),
),
Some("Review dismissed".to_owned())
);
}
fn user_prompt(position: u64, text: &str) -> Arc<hel::hel_state::TranscriptItem> {
Arc::new(hel::hel_state::TranscriptItem {
stable_id: format!("user:{position}"),
position,
latest_content_event_ordinal: None,
created_at_ms: 0,
last_changed_at_ms: 0,
body: hel::hel_state::TranscriptBody::User {
content: vec![serde_json::json!({
"type": "text",
"text": text,
})],
},
})
}
#[test]
fn seed_uses_the_latest_real_prompt_and_keeps_history_for_intent() {
let mut session = MaterializedSession::empty("seed-prompts");
session.applied_event_ordinal = 5;
session.transcript = vec![
user_prompt(1, "implement the old parser"),
user_prompt(2, "support parse_range"),
user_prompt(3, "[HARNESS NOTE: review the parser]"),
user_prompt(4, "also finish the parser error path"),
user_prompt(5, "[HARNESS NOTE: forwarded findings]"),
];
let mut state = TurnReviewState {
reviewed_through_ordinal: 1,
..TurnReviewState::default()
};
let seed = seed_from_session(&session, ReviewTier::Extended, &state, "manual");
assert_eq!(seed.task, "also finish the parser error path");
assert_eq!(
seed.user_messages
.iter()
.map(|message| message.text.as_str())
.collect::<Vec<_>>(),
vec![
"implement the old parser",
"support parse_range",
"also finish the parser error path",
],
"intent receives real prompts in chronological order"
);
assert!(!seed.trajectory.contains("HARNESS NOTE"));
state.reviewed_through_ordinal = 5;
let corrective = seed_from_session(&session, ReviewTier::Extended, &state, "manual");
assert_eq!(corrective.task, "also finish the parser error path");
assert_eq!(corrective.user_messages.len(), 3);
}
fn operational() -> RelayOperationalState {
serde_json::from_value(serde_json::json!({
"session_id": "reviewer",
"execution": "idle",
"latest_ordinal": 0,
"latest_digest": hel::hel_worker::RELAY_EVENT_GENESIS_DIGEST,
"acknowledged_through": 0,
"acknowledged_digest": hel::hel_worker::RELAY_EVENT_GENESIS_DIGEST,
"recovery_floor_ordinal": 0,
"recovery_floor_digest": hel::hel_worker::RELAY_EVENT_GENESIS_DIGEST,
"native_session_id": null,
"agent_capabilities": null,
"agent_info": null,
"config_options": [],
"available_commands": [],
"config": {},
"active_prompt": null,
"queued_prompts": [],
"checkpoint_barrier": null,
"checkpoint_ready": null,
}))
.expect("the operational state fixture matches its schema")
}
struct FakeManager {
session: String,
control: SessionManagerControl,
requests: RemoteSessionRequests,
publisher: crate::hel_session_manager::RemoteSessionPublisher,
_shutdown: crate::hel_session_manager::SessionManagerShutdown,
_targets: tokio::sync::watch::Sender<Vec<RelaySessionTarget>>,
}
impl FakeManager {
async fn new(session: &str) -> Self {
let channels = spawn_remote_session_manager().expect("remote manager");
channels.targets.send_replace(vec![RelaySessionTarget {
session_id: session.to_owned(),
spec: hel::hel_targets::CommandSpec::new("true", Vec::<String>::new()),
worker_recovery: None,
project_memory: None,
}]);
let manager = Self {
session: session.to_owned(),
control: channels.control,
requests: channels.requests,
publisher: channels.publisher,
_shutdown: channels.shutdown,
_targets: channels.targets,
};
manager
.publisher
.publish(
session.to_owned(),
view(session, hel::hel_state::MaterializedExecutionState::Idle),
)
.await
.expect("publish the first view");
manager
.control
.wait_for_session(session, Duration::from_secs(5))
.await
.expect("the fake manager manages the session");
manager
}
async fn next(&mut self) -> RemoteSessionRequest {
tokio::time::timeout(Duration::from_secs(5), self.requests.recv())
.await
.expect("the host makes a request")
.expect("the manager is still running")
}
async fn next_reviewer(
&mut self,
wanted: impl Fn(&Option<String>, &ReviewerAction) -> bool,
) -> (
Option<String>,
ReviewerAction,
oneshot::Sender<Result<ReviewerOutcome, String>>,
) {
loop {
match self.next().await {
RemoteSessionRequest::Reviewer {
role,
action,
reply,
..
} => {
if wanted(&role, &action) {
return (role, action, reply);
}
let _ = reply.send(answer_for(&action));
}
RemoteSessionRequest::Submit { reply, .. } => {
let _ = reply.send(Ok(1));
}
other => panic!("unexpected request {}", other.session_id()),
}
}
}
}
fn answer_for(action: &ReviewerAction) -> Result<ReviewerOutcome, String> {
match action {
ReviewerAction::Status => Ok(ReviewerOutcome::Status(Box::new(operational()))),
ReviewerAction::CaptureDelta { .. } => Ok(ReviewerOutcome::Delta {
repositories: Vec::new(),
}),
ReviewerAction::AnalyzeDelta { .. } => Ok(ReviewerOutcome::ChangedFunctions {
packet: "- edited retry()".to_owned(),
}),
ReviewerAction::AdvanceBaseline { .. } => Ok(ReviewerOutcome::BaselineAdvanced),
ReviewerAction::TakeLaneDispatches => Ok(ReviewerOutcome::LaneDispatches {
requests: Vec::new(),
}),
ReviewerAction::Attach { .. } => Ok(ReviewerOutcome::Attached(Box::new(
crate::hel_worker_client::RelayAttachment {
state: operational(),
events: Vec::new(),
through_ordinal: 0,
through_digest: hel::hel_worker::RELAY_EVENT_GENESIS_DIGEST.to_owned(),
},
))),
ReviewerAction::Pause => Ok(ReviewerOutcome::Paused),
ReviewerAction::Submit { .. } => Ok(ReviewerOutcome::Accepted { ordinal: 1 }),
ReviewerAction::Start { .. } => Err("no harness in this test".to_owned()),
ReviewerAction::RespondElicitation { .. } => Ok(ReviewerOutcome::ElicitationResolved),
ReviewerAction::Acknowledge { .. } => Ok(ReviewerOutcome::Acknowledged(
hel::hel_worker::RelayCursor {
ordinal: 0,
digest: hel::hel_worker::RELAY_EVENT_GENESIS_DIGEST.to_owned(),
},
)),
}
}
fn view(
session: &str,
execution: hel::hel_state::MaterializedExecutionState,
) -> ManagedSessionView {
view_of(session, execution, true)
}
fn view_of(
session: &str,
execution: hel::hel_state::MaterializedExecutionState,
prompt_driven: bool,
) -> ManagedSessionView {
let mut materialized = MaterializedSession::empty(session);
materialized.execution = execution;
materialized.applied_event_ordinal = 12;
let mut operational = operational();
if prompt_driven
&& matches!(
execution,
hel::hel_state::MaterializedExecutionState::Running { .. }
)
{
operational.active_prompt = Some(hel::hel_worker::ActiveRelayPrompt {
command_id: "prompt-1".to_owned(),
created_at_ms: 0,
started_at_ms: 0,
});
}
ManagedSessionView {
snapshot: Some(ManagedSessionSnapshot {
window: hel::hel_state::ProjectionWindow::of(&materialized),
materialized,
operational,
latest_credential_sync_signal: None,
worker_build: None,
}),
connected: true,
error: None,
}
}
struct FakeEnvironment {
staged: Mutex<Vec<(String, u64, bool)>>,
state: Mutex<TurnReviewState>,
writes: Mutex<Vec<(TurnReviewState, std::thread::ThreadId)>>,
save_gate: Mutex<Option<Arc<SaveGate>>>,
}
struct SaveGate {
entered: tokio::sync::Notify,
released: Mutex<bool>,
released_changed: std::sync::Condvar,
}
impl SaveGate {
fn new() -> Arc<Self> {
Arc::new(Self {
entered: tokio::sync::Notify::new(),
released: Mutex::new(false),
released_changed: std::sync::Condvar::new(),
})
}
async fn entered(&self) {
self.entered.notified().await;
}
fn wait(&self) {
self.entered.notify_one();
let released = self
.released
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
drop(
self.released_changed
.wait_while(released, |released| !*released)
.unwrap_or_else(std::sync::PoisonError::into_inner),
);
}
fn release(&self) {
*self
.released
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = true;
self.released_changed.notify_all();
}
}
impl FakeEnvironment {
fn new() -> Arc<Self> {
Arc::new(Self {
staged: Mutex::new(Vec::new()),
state: Mutex::new(TurnReviewState::default()),
writes: Mutex::new(Vec::new()),
save_gate: Mutex::new(None),
})
}
fn state(&self) -> TurnReviewState {
self.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
fn staged_roles(&self) -> Vec<(String, u64, bool)> {
self.staged
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
fn writes(&self) -> Vec<(TurnReviewState, std::thread::ThreadId)> {
self.writes
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
fn block_saves(&self) -> Arc<SaveGate> {
let gate = SaveGate::new();
*self
.save_gate
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(gate.clone());
gate
}
}
impl ReviewEnvironment for FakeEnvironment {
fn check(&self, _session_id: &str, _profile: &str) -> Result<(), String> {
Ok(())
}
fn stage(
&self,
_session_id: &str,
profile: &str,
generation: u64,
mcp_servers: &[hel::hel_worker_launch::ReviewMcpServer],
dispatch_tool: bool,
) -> Result<hel::hel_worker_launch::ReviewerLaunchConfig, String> {
self.staged
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push((profile.to_owned(), generation, dispatch_tool));
Ok(hel::hel_worker_launch::ReviewerLaunchConfig {
profile_id: profile.to_owned(),
harness: hel::hel_config::HarnessKind::Claude,
bridge_command: std::path::PathBuf::from("/bin/false"),
bridge_args: Vec::new(),
environment: Default::default(),
execution_policy: hel::hel_config::ExecutionPolicy::ConfiguredApprovals,
model: None,
effort: None,
generation,
mcp_servers: mcp_servers.to_vec(),
})
}
fn load_state(&self, _session_id: &str) -> Result<TurnReviewState, String> {
Ok(self.state())
}
fn save_state(&self, _session_id: &str, state: &TurnReviewState) -> Result<(), String> {
let gate = self
.save_gate
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
if let Some(gate) = gate {
gate.wait();
}
*self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = state.clone();
self.writes
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push((state.clone(), std::thread::current().id()));
Ok(())
}
fn clear_interrupted(&self) -> Result<Vec<String>, String> {
self.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.active = None;
Ok(Vec::new())
}
}
fn armed(profile: Option<&str>) -> ReviewConfigSource {
let profile = profile.map(str::to_owned);
Arc::new(move || ReviewConfig {
enabled: true,
tier: ReviewTier::Quick,
profile: profile.clone(),
model: None,
effort: None,
})
}
#[test]
fn the_reviewer_profile_must_be_separate_from_the_primary_profile() {
let session = hel::hel_state::SessionRecord {
id: "session-1".to_owned(),
workspace_id: hel::hel_workspace::DEFAULT_WORKSPACE_ID.to_owned(),
title: "task".to_owned(),
harness_kind: hel::hel_config::HarnessKind::Codex,
last_profile: "primary".to_owned(),
bundle_id: "bundle".to_owned(),
project_directory: None,
managed_worktree: None,
target_template_id: "local".to_owned(),
resource_allocation: None,
additional_mounts: Vec::new(),
container_cpus: None,
container_memory: None,
state: hel::hel_state::SessionState::Running,
archived: false,
target: None,
native_session_id: None,
acp_session_title: None,
session_title_override: None,
created_at: "2026-01-01T00:00:00Z".to_owned(),
updated_at: "2026-01-01T00:00:00Z".to_owned(),
viewed_through_event_ordinal: 0,
draft_input: String::new(),
last_error: None,
last_checkpoint_error: None,
checkpoint: None,
};
let refusal = validate_reviewer_assignment("session-1", Some(&session), "primary")
.expect_err("one harness profile cannot review its own output independently");
assert!(refusal.contains("primary profile"), "{refusal}");
validate_reviewer_assignment("session-1", Some(&session), "reviewer")
.expect("a separate reviewer profile is accepted");
}
#[test]
fn delivery_admission_bypasses_only_the_matching_held_prompt() {
let session = session_id("admission");
hold_prompts(&session);
let admission = admit_review_delivery(&session, 7, "forward-7")
.expect("the live review hold grants its own corrective command");
assert!(review_delivery_admitted(&session, &admission));
assert!(!review_delivery_admitted(
&session,
&ReviewDeliveryAdmission::new(session.clone(), 7, "other-command".to_owned())
));
assert!(!review_delivery_admitted(
&session,
&ReviewDeliveryAdmission::new(session.clone(), 8, "forward-7".to_owned())
));
assert!(prompt_refusal(&session).is_some());
release_prompts(&session);
}
async fn finish_a_turn(manager: &FakeManager, host: &TurnReviewHost) {
for execution in [
hel::hel_state::MaterializedExecutionState::Running { started_at_ms: 0 },
hel::hel_state::MaterializedExecutionState::Idle,
] {
let published = view(&manager.session, execution);
let _ = manager
.publisher
.publish(manager.session.clone(), published.clone())
.await;
host.observe(&manager.session, &published);
}
}
#[tokio::test]
async fn an_unconfigured_reviewer_is_reported_once_per_session() {
let session = session_id("unreviewable");
let session = session.as_str();
let mut manager = FakeManager::new(session).await;
let environment = FakeEnvironment::new();
let host =
TurnReviewHost::spawn_in(manager.control.clone(), armed(None), environment.clone());
finish_a_turn(&manager, &host).await;
let request = manager.next().await;
let RemoteSessionRequest::Submit { command, reply, .. } = request else {
panic!("the only thing an unreviewable turn does is say so");
};
let RelayCommand::RecordNotice { text } = command else {
panic!("the notice is a controller-authored conversation line");
};
assert!(
text.contains("[review] profile"),
"the notice names the key that fixes it: {text}"
);
let _ = reply.send(Ok(1));
finish_a_turn(&manager, &host).await;
assert!(
tokio::time::timeout(Duration::from_millis(300), manager.requests.recv())
.await
.is_err(),
"a second unreviewable turn is silent"
);
assert!(!host.refuses_prompt(session));
}
#[tokio::test]
async fn a_self_started_turn_does_not_arm_an_automatic_review() {
let session = session_id("selfstarted0");
let session = session.as_str();
let mut manager = FakeManager::new(session).await;
let environment = FakeEnvironment::new();
let host =
TurnReviewHost::spawn_in(manager.control.clone(), armed(None), environment.clone());
for execution in [
MaterializedExecutionState::Running { started_at_ms: 0 },
MaterializedExecutionState::Idle,
] {
host.observe(session, &view_of(session, execution, false));
}
assert!(
tokio::time::timeout(Duration::from_millis(300), manager.requests.recv())
.await
.is_err(),
"a turn the harness started on its own arms nothing"
);
finish_a_turn(&manager, &host).await;
assert!(
matches!(manager.next().await, RemoteSessionRequest::Submit { .. }),
"a prompt-driven turn still reaches the automatic edge"
);
host.shutdown().await.expect("shutdown the host");
}
#[tokio::test]
async fn a_headless_turn_is_reviewed_and_resolves_itself() {
let session = session_id("headless000");
let session = session.as_str();
let mut manager = FakeManager::new(session).await;
let environment = FakeEnvironment::new();
let host = TurnReviewHost::spawn_in(
manager.control.clone(),
armed(Some("reviewer")),
environment.clone(),
);
finish_a_turn(&manager, &host).await;
let (_, action, reply) = manager
.next_reviewer(|_, action| matches!(action, ReviewerAction::Status))
.await;
assert!(matches!(action, ReviewerAction::Status));
assert!(
host.refuses_prompt(session),
"admission holds prompts before preparation waits on the session actor"
);
let _ = reply.send(Ok(ReviewerOutcome::Status(Box::new(operational()))));
let (_, action, reply) = manager
.next_reviewer(|_, action| matches!(action, ReviewerAction::CaptureDelta { .. }))
.await;
assert!(matches!(action, ReviewerAction::CaptureDelta { .. }));
assert!(
host.refuses_prompt(session),
"the review holds the session's prompts from the moment it opens"
);
assert!(
environment.state().active.is_some(),
"the active marker is durable before review work starts"
);
let _ = reply.send(Ok(ReviewerOutcome::Delta {
repositories: vec![hel::hel_worker::RepoDelta {
root: std::path::PathBuf::from("/workspace/app"),
baseline_tree: None,
current_tree: "first-tree".to_owned(),
patch: String::new(),
diffstat: "0 files changed".to_owned(),
changed_lines: 0,
}],
}));
let (_, action, reply) = manager
.next_reviewer(|_, action| matches!(action, ReviewerAction::AdvanceBaseline { .. }))
.await;
let ReviewerAction::AdvanceBaseline { trees } = action else {
unreachable!("matched above");
};
assert_eq!(
trees
.get(std::path::Path::new("/workspace/app"))
.map(String::as_str),
Some("first-tree"),
"the capture becomes the baseline the next review measures from"
);
let _ = reply.send(Ok(ReviewerOutcome::BaselineAdvanced));
tokio::time::timeout(Duration::from_secs(5), async {
while host.refuses_prompt(session)
|| host.view(session).is_some()
|| environment.state().active.is_some()
{
tokio::task::yield_now().await;
}
})
.await
.expect("a resolved review releases prompts and drains its durable close");
assert!(host.view(session).is_none(), "the review is over");
assert_eq!(environment.state().active, None);
}
#[tokio::test]
async fn preparation_rechecks_the_live_actor_after_installing_the_prompt_hold() {
let session = session_id("preparelive");
let session = session.as_str();
let mut manager = FakeManager::new(session).await;
let environment = FakeEnvironment::new();
let host = TurnReviewHost::spawn_in(
manager.control.clone(),
armed(Some("reviewer")),
environment.clone(),
);
finish_a_turn(&manager, &host).await;
let (_, _, reply) = manager
.next_reviewer(|_, action| matches!(action, ReviewerAction::Status))
.await;
assert!(host.refuses_prompt(session));
manager
.publisher
.publish(
session.to_owned(),
view(
session,
MaterializedExecutionState::Running { started_at_ms: 1 },
),
)
.await
.expect("publish the command that won the admission race");
let _ = reply.send(Ok(ReviewerOutcome::Status(Box::new(operational()))));
tokio::time::timeout(Duration::from_secs(5), async {
while host.refuses_prompt(session) {
tokio::task::yield_now().await;
}
})
.await
.expect("a refused preparation releases its prompt hold");
assert!(host.view(session).is_none());
assert_eq!(environment.state().active, None);
assert!(
tokio::time::timeout(Duration::from_millis(300), manager.requests.recv())
.await
.is_err(),
"stale preparation never starts capture"
);
host.shutdown().await.expect("shutdown the host");
}
#[tokio::test]
async fn observation_bursts_do_not_drop_the_final_idle_edge() {
let session = session_id("losslessobs");
let session = session.as_str();
let mut manager = FakeManager::new(session).await;
let environment = FakeEnvironment::new();
let config: ReviewConfigSource = Arc::new(|| ReviewConfig {
enabled: false,
tier: ReviewTier::Quick,
profile: Some("reviewer".to_owned()),
model: None,
effort: None,
});
let host = TurnReviewHost::spawn_in(manager.control.clone(), config, environment);
let running = view(
session,
MaterializedExecutionState::Running { started_at_ms: 1 },
);
for _ in 0..256 {
host.observe(session, &running);
}
host.observe(session, &view(session, MaterializedExecutionState::Idle));
let starting_host = host.clone();
let session_owned = session.to_owned();
let starting = tokio::spawn(async move { starting_host.start(&session_owned, true).await });
let (_, _, reply) = manager
.next_reviewer(|_, action| matches!(action, ReviewerAction::Status))
.await;
let _ = reply.send(Ok(ReviewerOutcome::Status(Box::new(operational()))));
starting
.await
.expect("start task")
.expect("the retained idle edge admits the review");
assert!(host.refuses_prompt(session));
host.shutdown().await.expect("shutdown the host");
}
#[tokio::test]
async fn persistence_is_nonblocking_ordered_and_drained_on_shutdown() {
let session = session_id("persistlane");
let session = session.as_str();
let mut manager = FakeManager::new(session).await;
let environment = FakeEnvironment::new();
let host = TurnReviewHost::spawn_in(
manager.control.clone(),
armed(Some("reviewer")),
environment.clone(),
);
finish_a_turn(&manager, &host).await;
let (_, _, reply) = manager
.next_reviewer(|_, action| matches!(action, ReviewerAction::Status))
.await;
let open_gate = environment.block_saves();
let _ = reply.send(Ok(ReviewerOutcome::Status(Box::new(operational()))));
tokio::time::timeout(Duration::from_secs(5), open_gate.entered())
.await
.expect("the active write reaches the blocking lane");
let refusal = tokio::time::timeout(Duration::from_secs(1), host.start(session, true))
.await
.expect("the host loop remains responsive while persistence blocks")
.expect_err("the same review is already starting");
assert!(refusal.0.contains("already starting"), "{refusal}");
assert!(host.view(session).is_none(), "open is not exposed early");
open_gate.release();
let (_, _, _capture_reply) = manager
.next_reviewer(|_, action| matches!(action, ReviewerAction::CaptureDelta { .. }))
.await;
assert!(environment.state().active.is_some());
assert!(host.view(session).is_some());
let close_gate = environment.block_saves();
let first_host = host.clone();
let second_host = host.clone();
let first = tokio::spawn(async move { first_host.shutdown().await });
let second = tokio::spawn(async move { second_host.shutdown().await });
tokio::time::timeout(Duration::from_secs(5), close_gate.entered())
.await
.expect("shutdown queues the final inactive state");
assert!(!first.is_finished(), "shutdown drains the blocked write");
assert!(
!second.is_finished(),
"concurrent shutdown joins the same drain"
);
assert!(
!host.refuses_prompt(session),
"logical shutdown releases prompts before persistence finishes"
);
close_gate.release();
first
.await
.expect("first shutdown task")
.expect("first drain");
second
.await
.expect("second shutdown task")
.expect("shared drain");
host.shutdown().await.expect("shutdown stays idempotent");
assert_eq!(environment.state().active, None);
assert!(host.view(session).is_none());
let writes = environment.writes();
assert!(
writes
.first()
.is_some_and(|(state, _)| state.active.is_some())
);
assert!(
writes
.last()
.is_some_and(|(state, _)| state.active.is_none())
);
let test_thread = std::thread::current().id();
assert!(
writes.iter().all(|(_, writer)| *writer != test_thread),
"synchronous database writes run off the Tokio host thread"
);
}
#[tokio::test]
async fn an_interrupted_handoff_retains_findings_until_acceptance_and_retries_the_same_id() {
let session = session_id("handoff0000");
let mut manager = FakeManager::new(&session).await;
let environment = FakeEnvironment::new();
let pending = PendingForward {
synthesis: "[P2] src/lib.rs:1 -- incorrect boundary".to_owned(),
evidence: Default::default(),
command_id: "durable-forward-id".to_owned(),
trees: BTreeMap::from([(PathBuf::from("/workspace/app"), "new".to_owned())]),
reviewed_through_ordinal: 12,
};
{
let mut state = environment.state.lock().unwrap();
state
.baselines
.insert("/workspace/app".into(), "old".into());
state.pending_forward = Some(pending.clone());
}
let host = TurnReviewHost::spawn_in(
manager.control.clone(),
armed(Some("reviewer")),
environment.clone(),
);
host.observe(&session, &view(&session, MaterializedExecutionState::Idle));
host.events
.send(HostEvent::Interrupted {
interrupted: vec![session.clone()],
})
.unwrap();
let RemoteSessionRequest::Submit {
command_id,
admission,
reply,
..
} = manager.next().await
else {
panic!("recovery submits the pending handoff directly, without starting a reviewer");
};
assert_eq!(command_id, pending.command_id);
assert!(review_delivery_admitted(&session, &admission.unwrap()));
assert!(host.refuses_prompt(&session));
assert_eq!(environment.state().pending_forward, Some(pending.clone()));
assert_eq!(
environment.state().baselines[&PathBuf::from("/workspace/app")],
"old"
);
assert!(
host.resolve(&session, Resolution::Forwarded).await.is_err(),
"duplicate Forward is not another submission"
);
assert!(
host.resolve(&session, Resolution::Cancelled).await.is_err(),
"an unknown delivery cannot be undone"
);
reply
.send(Err("primary temporarily unavailable".to_owned()))
.unwrap();
tokio::time::timeout(Duration::from_secs(2), async {
while !host.view(&session).is_some_and(|view| {
matches!(
view.phase,
TurnReviewPhase::Forwarding { error: Some(_), .. }
)
}) {
tokio::task::yield_now().await;
}
})
.await
.expect("rejection remains actionable");
assert_eq!(environment.state().pending_forward, Some(pending.clone()));
host.resolve(&session, Resolution::Forwarded).await.unwrap();
let RemoteSessionRequest::Submit {
command_id,
admission,
reply,
..
} = manager.next().await
else {
panic!("retry submits the same handoff");
};
assert_eq!(command_id, pending.command_id);
let epoch = admission.unwrap().epoch();
host.events
.send(HostEvent::Step {
session_id: session.clone(),
epoch,
step: ReviewStep::RoleEvents {
role: "reviewer".to_owned(),
result: Err("late reviewer disconnect".to_owned()),
},
})
.unwrap();
let gate = environment.block_saves();
reply.send(Ok(42)).unwrap();
tokio::time::timeout(Duration::from_secs(2), gate.entered())
.await
.unwrap();
assert_eq!(
environment.state().pending_forward,
Some(pending),
"the durable pending record remains until the complete accepted outcome is written"
);
gate.release();
tokio::time::timeout(Duration::from_secs(2), async {
while host.view(&session).is_some() {
tokio::task::yield_now().await;
}
})
.await
.expect("accepted handoff closes");
let state = environment.state();
assert!(state.pending_forward.is_none());
assert!(state.prior_review.is_some());
assert_eq!(state.baselines[&PathBuf::from("/workspace/app")], "new");
assert!(!host.refuses_prompt(&session));
assert!(environment.staged_roles().is_empty());
host.shutdown().await.unwrap();
}
#[tokio::test]
async fn queued_prompts_hold_a_review_back() {
let session = session_id("queued00000");
let session = session.as_str();
let mut manager = FakeManager::new(session).await;
let environment = FakeEnvironment::new();
let host = TurnReviewHost::spawn_in(
manager.control.clone(),
armed(Some("reviewer")),
environment.clone(),
);
let mut queued = view(session, hel::hel_state::MaterializedExecutionState::Idle);
if let Some(snapshot) = queued.snapshot.as_mut() {
snapshot.materialized.queued_prompts = vec![hel::hel_state::MaterializedQueuedPrompt {
command_id: "queued-1".to_owned(),
kind: hel::hel_state::QueuedCommandKind::Prompt,
content: vec![serde_json::json!({"type": "text", "text": "next"})],
queued_at_ms: 0,
}];
}
host.observe(
session,
&view(
session,
hel::hel_state::MaterializedExecutionState::Running { started_at_ms: 0 },
),
);
host.observe(session, &queued);
assert!(
tokio::time::timeout(Duration::from_millis(300), manager.requests.recv())
.await
.is_err(),
"no review starts while prompts are queued"
);
let refusal = host
.start(session, true)
.await
.expect_err("a manual review is refused for the same reason");
assert!(refusal.0.contains("queued"), "{refusal}");
}
#[tokio::test]
async fn resolving_a_review_that_has_no_verdict_is_refused() {
let session = session_id("resolution0");
let session = session.as_str();
let mut manager = FakeManager::new(session).await;
let environment = FakeEnvironment::new();
let host = TurnReviewHost::spawn_in(
manager.control.clone(),
armed(Some("reviewer")),
environment.clone(),
);
let error = host
.resolve(session, Resolution::Forwarded)
.await
.expect_err("there is no review at all");
assert!(error.contains("no review is open"), "{error}");
finish_a_turn(&manager, &host).await;
let (_, _, reply) = manager
.next_reviewer(|_, action| matches!(action, ReviewerAction::CaptureDelta { .. }))
.await;
let _ = reply.send(Ok(ReviewerOutcome::Delta {
repositories: vec![hel::hel_worker::RepoDelta {
root: std::path::PathBuf::from("/workspace/app"),
baseline_tree: Some("base".to_owned()),
current_tree: "new".to_owned(),
patch: "diff --git a/a b/a\n@@\n+one\n".to_owned(),
diffstat: "1 file changed, 1 insertion(+)".to_owned(),
changed_lines: 1,
}],
}));
tokio::time::timeout(Duration::from_secs(5), async {
while host.view(session).is_none() {
tokio::task::yield_now().await;
}
})
.await
.expect("the review is open");
let error = host
.resolve(session, Resolution::Forwarded)
.await
.expect_err("nothing has been found yet");
assert!(error.contains("no findings"), "{error}");
let error = host
.resolve(session, Resolution::Dismissed)
.await
.expect_err("nothing has been decided yet");
assert!(error.contains("verdict"), "{error}");
host.resolve(session, Resolution::Cancelled)
.await
.expect("cancel needs no verdict");
tokio::time::timeout(Duration::from_secs(5), async {
while host.refuses_prompt(session) {
tokio::task::yield_now().await;
}
})
.await
.expect("cancelling releases the prompts");
}
#[tokio::test]
async fn a_failed_review_clears_durable_active_state_and_the_prompt_hold() {
let session = session_id("failedrole0");
let session = session.as_str();
let mut manager = FakeManager::new(session).await;
let environment = FakeEnvironment::new();
let host = TurnReviewHost::spawn_in(
manager.control.clone(),
armed(Some("reviewer")),
environment.clone(),
);
finish_a_turn(&manager, &host).await;
let (_, _, reply) = manager
.next_reviewer(|_, action| matches!(action, ReviewerAction::CaptureDelta { .. }))
.await;
assert!(environment.state().active.is_some());
let _ = reply.send(Ok(ReviewerOutcome::Delta {
repositories: vec![hel::hel_worker::RepoDelta {
root: PathBuf::from("/workspace/app"),
baseline_tree: Some("base".to_owned()),
current_tree: "new".to_owned(),
patch: "diff --git a/a b/a\n@@\n+one\n".to_owned(),
diffstat: "1 file changed, 1 insertion(+)".to_owned(),
changed_lines: 1,
}],
}));
let (_, _, reply) = manager
.next_reviewer(|_, action| matches!(action, ReviewerAction::Start { .. }))
.await;
let _ = reply.send(Err("review harness failed to launch".to_owned()));
tokio::time::timeout(Duration::from_secs(5), async {
loop {
let failed = host.view(session).is_some_and(|view| {
matches!(
view.verdict,
Some(VerdictView {
kind: VerdictKind::Failed,
..
})
)
});
if failed && !host.refuses_prompt(session) && environment.state().active.is_none() {
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("failure releases and persists the turn");
host.resolve(session, Resolution::Dismissed)
.await
.expect("the visible failure can be dismissed");
tokio::time::timeout(Duration::from_secs(5), async {
while host.view(session).is_some() {
tokio::task::yield_now().await;
}
})
.await
.expect("dismissal closes the failed review");
host.shutdown().await.expect("shutdown the host");
}
fn agent_event(ordinal: u64, previous_digest: &str, text: &str) -> RelayEvent {
let mut event = RelayEvent {
format: RELAY_EVENT_FORMAT_V1,
ordinal,
previous_digest: previous_digest.to_owned(),
digest: String::new(),
recorded_at_ms: i64::try_from(ordinal).unwrap_or_default() * 100,
command_id: None,
observation: RelayObservation::SessionUpdate {
update: Box::new(
agent_client_protocol::schema::v1::SessionUpdate::AgentMessageChunk(
agent_client_protocol::schema::v1::ContentChunk::new(
agent_client_protocol::schema::v1::ContentBlock::Text(
agent_client_protocol::schema::v1::TextContent::new(text),
),
),
),
),
},
};
event.digest = relay_event_digest(&event).expect("digest");
event
}
fn completion_event(ordinal: u64, previous_digest: &str, command_id: &str) -> RelayEvent {
let mut event = RelayEvent {
format: RELAY_EVENT_FORMAT_V1,
ordinal,
previous_digest: previous_digest.to_owned(),
digest: String::new(),
recorded_at_ms: i64::try_from(ordinal).unwrap_or_default() * 100,
command_id: Some(command_id.to_owned()),
observation: RelayObservation::CommandCompleted {
command_id: command_id.to_owned(),
outcome: RelayCommandOutcome::Prompt {
stop_reason: "end_turn".to_owned(),
},
},
};
event.digest = relay_event_digest(&event).expect("digest");
event
}
#[tokio::test]
async fn a_clean_reviewer_report_resolves_the_review() {
let session = session_id("cleanreport");
let session = session.as_str();
let mut manager = FakeManager::new(session).await;
let environment = FakeEnvironment::new();
let publications = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let published = publications.clone();
let host = TurnReviewHost::spawn_in_notifying(
manager.control.clone(),
armed(Some("reviewer")),
environment.clone(),
Arc::new(move || {
published.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
}),
);
finish_a_turn(&manager, &host).await;
let (_, _, reply) = manager
.next_reviewer(|_, action| matches!(action, ReviewerAction::CaptureDelta { .. }))
.await;
let after_add = publications.load(std::sync::atomic::Ordering::SeqCst);
assert!(after_add > 0, "opening publishes and wakes surfaces");
let _ = reply.send(Ok(ReviewerOutcome::Delta {
repositories: vec![hel::hel_worker::RepoDelta {
root: std::path::PathBuf::from("/workspace/app"),
baseline_tree: Some("base".to_owned()),
current_tree: "new".to_owned(),
patch: "diff --git a/a b/a\n@@\n+one\n".to_owned(),
diffstat: "1 file changed, 1 insertion(+)".to_owned(),
changed_lines: 1,
}],
}));
let (role, _, reply) = manager
.next_reviewer(|_, action| matches!(action, ReviewerAction::Start { .. }))
.await;
let after_change = publications.load(std::sync::atomic::Ordering::SeqCst);
assert!(
after_change > after_add,
"a projected review state change wakes surfaces"
);
assert_eq!(
role.as_deref(),
Some(hel::hel_review::driver::REVIEWER_ROLE)
);
let _ = reply.send(Ok(ReviewerOutcome::Started(Box::new(
crate::hel_worker_client::StartedReviewer {
native_session_id: None,
config_options: Vec::new(),
reused: false,
state: operational(),
},
))));
let (_, action, reply) = manager
.next_reviewer(|role, action| {
role.as_deref() == Some(hel::hel_review::driver::REVIEWER_ROLE)
&& matches!(action, ReviewerAction::Submit { .. })
})
.await;
let ReviewerAction::Submit {
command_id,
command,
} = action
else {
unreachable!("matched above");
};
let RelayCommand::Prompt { prompt } = command else {
panic!("a reviewing role is prompted");
};
assert!(
format!("{prompt:?}").contains("+one"),
"the reviewer is given the captured change"
);
let _ = reply.send(Ok(ReviewerOutcome::Accepted { ordinal: 1 }));
let (_, _, reply) = manager
.next_reviewer(|role, action| {
role.as_deref() == Some(hel::hel_review::driver::REVIEWER_ROLE)
&& matches!(action, ReviewerAction::Attach { .. })
})
.await;
assert!(
tokio::time::timeout(Duration::from_millis(100), manager.requests.recv())
.await
.is_err(),
"one role prompt has only one attachment poll in flight"
);
let before_identical = publications.load(std::sync::atomic::Ordering::SeqCst);
let _ = reply.send(Ok(ReviewerOutcome::Attached(Box::new(
crate::hel_worker_client::RelayAttachment {
state: operational(),
events: Vec::new(),
through_ordinal: 0,
through_digest: hel::hel_worker::RELAY_EVENT_GENESIS_DIGEST.to_owned(),
},
))));
let (_, _, reply) = manager
.next_reviewer(|role, action| {
role.as_deref() == Some(hel::hel_review::driver::REVIEWER_ROLE)
&& matches!(action, ReviewerAction::Attach { .. })
})
.await;
assert_eq!(
publications.load(std::sync::atomic::Ordering::SeqCst),
before_identical,
"an identical projection does not wake surfaces"
);
let answer = agent_event(
1,
hel::hel_worker::RELAY_EVENT_GENESIS_DIGEST,
"No findings.",
);
let completion = completion_event(2, &answer.digest, &command_id);
let through_digest = completion.digest.clone();
let _ = reply.send(Ok(ReviewerOutcome::Attached(Box::new(
crate::hel_worker_client::RelayAttachment {
state: operational(),
events: vec![answer, completion],
through_ordinal: 2,
through_digest,
},
))));
tokio::time::timeout(Duration::from_secs(5), async {
while host.refuses_prompt(session)
|| host.view(session).is_some()
|| environment.state().active.is_some()
{
tokio::task::yield_now().await;
}
})
.await
.expect("a clean review releases and durably closes the turn by itself");
assert!(host.view(session).is_none());
assert!(
publications.load(std::sync::atomic::Ordering::SeqCst) > before_identical,
"closing removes the view and wakes surfaces"
);
let staged = environment.staged_roles();
assert_eq!(staged.len(), 1);
assert_eq!(staged[0].0, "reviewer");
assert_ne!(staged[0].1, 0, "fresh reviewer generation");
assert!(!staged[0].2);
let recorded = environment.state();
assert_eq!(
recorded
.baselines
.get(std::path::Path::new("/workspace/app"))
.map(String::as_str),
Some("new")
);
assert_eq!(recorded.reviewed_through_ordinal, 12);
assert_eq!(recorded.active, None);
host.shutdown().await.expect("shutdown the host");
}
}