use rusqlite::{params, Connection, OptionalExtension, Transaction, TransactionBehavior};
use time::OffsetDateTime;
use crate::child::ChildRef;
use crate::durable::{
AdvanceReceipt, AgentInvocation, AgentInvocationId, Ask, AskBody, AskClaim, AskId, AskOrigin,
AskResult, AskState, AskTarget, Author, Basis, BoundarySeed, BoundaryState, Containment,
ContainmentObservation, DoneProposal, DoneProposalId, Epoch, EpochId, EpochReceipt, EpochState,
FlowPosition, Home, HomeId, InterruptReceipt, InvocationRoute, InvocationSurface, Placement,
ProjectId, Run, RunAdvance, RunId, RunLease, RunLeaseToken, RunState, RunTrigger, Send, SendId,
SendState, SendVia, Steer, SteerId, SteerReceipt, StopCause, StopReceipt, TaskId,
ToolResponseId, ToolResponseReceipt, ToolResponseWrite, Turn, TurnId, Wait, WaitId, WaitOn,
WorkRef, WorkStatus,
};
use crate::id::WaveId;
use crate::project::Project;
use crate::store::durable::{AskCommentTransition, AskCommentWrite};
use crate::store::rows::now_unix;
use crate::store::{StoreError, StoreResult};
use crate::task::Task;
use super::SqliteStore;
const HAS_PENDING_USER_ASK_FOR_WORK_SQL: &str = "SELECT EXISTS(
SELECT 1 FROM ask_exchanges a
JOIN epochs e ON e.id=a.epoch_id
WHERE a.target_kind='user' AND a.state IN ('queued', 'claimed')
AND e.state='open'
AND (e.wave_id=?1 OR e.project_id=?2 OR e.task_id=?3)
)";
impl SqliteStore {
pub(crate) fn task_writer_state(
&self,
external_issue_id: &str,
) -> StoreResult<Option<crate::store::durable::TaskWriterState>> {
let task = {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.query_row(
"SELECT id, issue_identifier FROM tasks WHERE external_issue_id=?1",
[external_issue_id],
|row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
)
.optional()?
};
let Some((id, identifier)) = task else {
return Ok(None);
};
let work = WorkRef::Task(TaskId::parse(&id).map_err(invalid_durable)?);
let run = self.current_run(&work)?;
Ok(Some(crate::store::durable::TaskWriterState {
work,
identifier,
run,
}))
}
pub fn home_by_id(&self, home_id: &HomeId) -> StoreResult<Option<Home>> {
let conn = self.conn.lock().expect("store mutex poisoned");
map_home_by_id(&conn, home_id)
}
pub fn local_home(&self) -> StoreResult<Home> {
let conn = self.conn.lock().expect("store mutex poisoned");
map_local_home(&conn)
}
pub fn observe_home(&self, home_id: &HomeId, route: &str) -> StoreResult<Home> {
let route = route.trim();
if route.is_empty() {
return Err(StoreError::InvalidData(
"Home route cannot be empty".to_string(),
));
}
let conn = self.conn.lock().expect("store mutex poisoned");
if map_home_by_id(&conn, home_id)?
.is_some_and(|home| home.route == "local" && route != "local")
{
return Err(StoreError::InvalidData(format!(
"cannot replace local Home {home_id} with remote route {route:?}"
)));
}
let existing_id = conn
.query_row("SELECT id FROM homes WHERE route=?1", [route], |row| {
row.get::<_, String>(0)
})
.optional()?;
if existing_id
.as_deref()
.is_some_and(|id| id != home_id.as_str())
{
return Err(StoreError::InvalidData(format!(
"Home route {route:?} is already observed for {}",
existing_id.expect("checked as present")
)));
}
let now = now_unix();
conn.execute(
"INSERT INTO homes (id, route, created_at, observed_at)
VALUES (?1, ?2, ?3, ?3)
ON CONFLICT(id) DO UPDATE SET
route=excluded.route, observed_at=excluded.observed_at",
params![home_id.as_str(), route, now],
)?;
map_home_by_id(&conn, home_id)?.ok_or(StoreError::NotFound)
}
pub fn placement(&self, work: &WorkRef) -> StoreResult<Placement> {
let conn = self.conn.lock().expect("store mutex poisoned");
placement_in(&conn, work)
}
pub fn set_work_enabled(&self, work: &WorkRef, enabled: bool) -> StoreResult<Placement> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
placement_in(&tx, work)?;
let enabled = if enabled { 1_i64 } else { 0_i64 };
match work {
WorkRef::Wave(id) => tx.execute(
"UPDATE work_placements SET enabled=?2 WHERE wave_id=?1",
params![id.as_str(), enabled],
)?,
WorkRef::Project(id) => tx.execute(
"UPDATE work_placements SET enabled=?2 WHERE project_id=?1",
params![id.as_str(), enabled],
)?,
WorkRef::Task(id) => tx.execute(
"UPDATE work_placements SET enabled=?2 WHERE task_id=?1",
params![id.as_str(), enabled],
)?,
};
let placement = placement_in(&tx, work)?;
tx.commit()?;
Ok(placement)
}
pub(crate) fn place_work(&self, work: &WorkRef, home_id: &HomeId) -> StoreResult<Placement> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
current_epoch_in(&tx, work)?;
tx.query_row(
"SELECT 1 FROM homes WHERE id=?1",
[home_id.as_str()],
|_| Ok(()),
)?;
if let Some(current) = find_placement_in(&tx, work)? {
if current.home_id == *home_id {
tx.commit()?;
return Ok(current);
}
}
if let Some(run) = current_run_for_work_in(&tx, work)? {
return Err(StoreError::InvalidData(format!(
"cannot move {} {} while Run {} is {:?}",
work.kind(),
work.id(),
run.id,
run.state
)));
}
write_placement(&tx, work, home_id, now_unix())?;
let placement = placement_in(&tx, work)?;
tx.commit()?;
Ok(placement)
}
pub fn reserve_run(
&self,
work: &WorkRef,
trigger: &RunTrigger,
) -> StoreResult<(Run, RunLease)> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let receipt = reserve_run_in(&tx, work, trigger)?;
tx.commit()?;
Ok(receipt)
}
pub(crate) fn reserve_child_run(
&self,
caller: &RunLease,
work: &WorkRef,
trigger: &RunTrigger,
) -> StoreResult<(Run, RunLease)> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
validate_control_caller(&tx, Some(caller), work)?;
let receipt = reserve_run_in(&tx, work, trigger)?;
tx.commit()?;
Ok(receipt)
}
pub(crate) fn reserve_recovery_run(&self, lease: &RunLease) -> StoreResult<(Run, RunLease)> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let prior = validate_run_lease(&tx, lease)?;
if prior.state != RunState::Active {
return Err(StoreError::InvalidData(format!(
"Run {} cannot hand off recovery while {:?}",
prior.id, prior.state
)));
}
let now = now_unix();
end_open_turns_for_run(&tx, &prior.id, now, "failed")?;
tx.execute(
"UPDATE agent_invocations
SET ended_at=COALESCE(ended_at, ?2),
outcome=CASE WHEN outcome='running' THEN 'failed' ELSE outcome END,
handback_state=COALESCE(handback_state, 'unknown')
WHERE supervising_run_id=?1 AND ended_at IS NULL",
params![prior.id.as_str(), now],
)?;
let stop_reason =
serde_json::to_string(&StopCause::Recovery).expect("Stop cause must serialize");
tx.execute(
"UPDATE runs SET state='ended', ended_at=?2, stop_reason=?3
WHERE id=?1 AND state='active'",
params![prior.id.as_str(), now, stop_reason],
)?;
let trigger = RunTrigger::Recovery {
prior_run_id: prior.id.clone(),
};
let token = RunLeaseToken::new();
let runtime_generation = runtime_generation_in(&tx, &prior.home_id)?;
let run = Run {
id: RunId::new(),
work: prior.work.clone(),
epoch_id: prior.epoch_id.clone(),
home_id: prior.home_id,
runtime_generation,
state: RunState::Reserved,
trigger: trigger.clone(),
retry_of: Some(prior.id),
containment: None,
cwd: None,
created_at: OffsetDateTime::now_utc(),
started_at: None,
ended_at: None,
};
if run.runtime_generation.is_some() {
let runtime_generation = runtime_generation_sql(run.runtime_generation)?;
tx.execute(
"INSERT INTO runs (
id, epoch_id, home_id, runtime_generation, state, trigger_json, retry_of,
lease_hash, lease_generation, source_kind, source_id, created_at, ended_at,
stop_reason
) VALUES (?1, ?2, ?3, ?4, 'reserved', ?5, ?6, ?7, NULL, ?8, ?9, ?10, NULL, NULL)",
params![
run.id.as_str(),
run.epoch_id.as_str(),
run.home_id.as_str(),
runtime_generation,
serde_json::to_string(&trigger).expect("Run trigger must serialize"),
run.retry_of.as_ref().map(RunId::as_str),
token.hash(),
run.work.kind(),
run.work.id(),
run.created_at.unix_timestamp(),
],
)?;
} else {
tx.execute(
"INSERT INTO runs (
id, epoch_id, home_id, state, trigger_json, retry_of, lease_hash,
lease_generation, source_kind, source_id, created_at, ended_at, stop_reason
) VALUES (?1, ?2, ?3, 'reserved', ?4, ?5, ?6, NULL, ?7, ?8, ?9, NULL, NULL)",
params![
run.id.as_str(),
run.epoch_id.as_str(),
run.home_id.as_str(),
serde_json::to_string(&trigger).expect("Run trigger must serialize"),
run.retry_of.as_ref().map(RunId::as_str),
token.hash(),
run.work.kind(),
run.work.id(),
run.created_at.unix_timestamp(),
],
)?;
}
let recovery_lease = RunLease::new(
run.id.clone(),
run.work.clone(),
current_epoch_in(&tx, &run.work)?.current_basis,
token,
);
tx.commit()?;
Ok((run, recovery_lease))
}
pub fn current_run(&self, work: &WorkRef) -> StoreResult<Option<Run>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let Ok(epoch) = current_epoch_in(&conn, work) else {
return Ok(None);
};
run_for_epoch_in(&conn, &epoch.id)
}
pub(crate) fn run_by_id(&self, run_id: &RunId) -> StoreResult<Run> {
let conn = self.conn.lock().expect("store mutex poisoned");
run_by_id_in(&conn, run_id)
}
pub(crate) fn latest_run(&self, work: &WorkRef) -> StoreResult<Option<Run>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let epoch = current_epoch_in(&conn, work)?;
latest_run_for_epoch_in(&conn, &epoch.id)
}
pub(crate) fn resolve_run_lease(&self, token: &RunLeaseToken) -> StoreResult<RunLease> {
let conn = self.conn.lock().expect("store mutex poisoned");
let id = conn
.query_row(
"SELECT id FROM runs
WHERE lease_hash=?1 AND state IN ('reserved', 'active')",
[token.hash()],
|row| row.get::<_, String>(0),
)
.optional()?
.ok_or_else(|| {
StoreError::InvalidAuthority(
"Run lease is malformed, stale, stopped, or unknown".to_string(),
)
})?;
let run = run_by_id_in(&conn, &RunId::parse(&id).map_err(invalid_durable)?)?;
let basis = current_epoch_in(&conn, &run.work)?.current_basis;
Ok(RunLease::new(run.id, run.work, basis, token.clone()))
}
pub(crate) fn validate_run_lease(&self, lease: &RunLease) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
validate_run_lease(&conn, lease)?;
Ok(())
}
pub(crate) fn record_first_material_at(
&self,
lease: &RunLease,
observed_at: OffsetDateTime,
) -> StoreResult<OffsetDateTime> {
let conn = self.conn.lock().expect("store mutex poisoned");
let run = validate_run_lease(&conn, lease)?;
if !matches!(run.work, WorkRef::Task(_)) || run.state != RunState::Active {
return Err(StoreError::InvalidAuthority(format!(
"Run {} does not hold active Task progress authority",
run.id
)));
}
let started_at = run.started_at.ok_or_else(|| {
StoreError::InvalidData(format!("active Run {} has no start time", run.id))
})?;
if observed_at < started_at {
return Err(StoreError::InvalidData(format!(
"Run {} first material event precedes its start",
run.id
)));
}
conn.execute(
"UPDATE runs
SET first_material_at=COALESCE(first_material_at, ?2)
WHERE id=?1 AND state='active'",
params![run.id.as_str(), observed_at.unix_timestamp()],
)?;
let accepted: i64 = conn.query_row(
"SELECT first_material_at FROM runs WHERE id=?1",
[run.id.as_str()],
|row| row.get(0),
)?;
OffsetDateTime::from_unix_timestamp(accepted)
.map_err(|error| StoreError::InvalidData(error.to_string()))
}
pub fn advance_run(
&self,
lease: &RunLease,
advance: &RunAdvance,
) -> StoreResult<AdvanceReceipt> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let mut run = validate_run_lease(&tx, lease)?;
let receipt = match advance {
RunAdvance::RunStarting { containment, cwd } => {
if run.state != RunState::Reserved {
return Err(StoreError::InvalidData(format!(
"Run {} cannot start while {:?}",
run.id, run.state
)));
}
if !cwd.is_absolute() {
return Err(StoreError::InvalidData(
"Run cwd must be absolute".to_string(),
));
}
let (kind, id) = containment.parts();
if id.trim().is_empty() {
return Err(StoreError::InvalidData(
"Run containment identity cannot be empty".to_string(),
));
}
let started_at = now_unix();
tx.execute(
"UPDATE runs
SET state='active', containment_kind=?2, containment_id=?3,
cwd=?4, started_at=?5
WHERE id=?1 AND state='reserved'",
params![
run.id.as_str(),
kind,
id,
cwd.display().to_string(),
started_at,
],
)?;
run = run_by_id_in(&tx, &run.id)?;
AdvanceReceipt::Run(run.clone())
}
RunAdvance::InvocationStarting {
route,
surface,
resume_token,
answer_ask_id,
} => {
if run.state != RunState::Active {
return Err(StoreError::InvalidData(format!(
"Run {} cannot supervise an Invocation while {:?}",
run.id, run.state
)));
}
if route.provider.trim().is_empty() || surface.trim().is_empty() {
return Err(StoreError::InvalidData(
"Invocation provider and surface cannot be empty".to_string(),
));
}
let invocation = AgentInvocation {
id: AgentInvocationId::new(),
supervising_run_id: Some(run.id.clone()),
answer_ask_id: answer_ask_id.clone(),
route: route.clone(),
surface: surface.clone(),
resume_token: resume_token.clone(),
started_at: OffsetDateTime::now_utc(),
ended_at: None,
};
insert_supervised_invocation(&tx, &run, &invocation)?;
AdvanceReceipt::Invocation(invocation)
}
RunAdvance::InvocationEnded {
invocation_id,
outcome,
} => {
if !outcome.is_terminal() {
return Err(StoreError::InvalidData(
"Invocation outcome must be terminal".to_string(),
));
}
require_invocation_for_run(&tx, invocation_id, &run.id)?;
let now = now_unix();
tx.execute(
"UPDATE agent_invocations
SET ended_at=COALESCE(ended_at, ?2), outcome=?3,
handback_state=CASE
WHEN surface IN ('tui', 'ide') AND answer_ask_id IS NULL
THEN handback_state
ELSE ?4
END
WHERE id=?1 AND ended_at IS NULL",
params![
invocation_id.as_str(),
now,
outcome.as_invocation_outcome(),
handback_state(*outcome),
],
)?;
AdvanceReceipt::Invocation(supervised_invocation_in(&tx, invocation_id)?)
}
RunAdvance::TurnStarting { invocation_id } => {
require_open_invocation_for_run(&tx, invocation_id, &run.id)?;
let basis = current_epoch_in(&tx, &run.work)?.current_basis;
let ordinal: i64 = tx.query_row(
"SELECT COALESCE(MAX(ordinal), 0) + 1 FROM agent_turns WHERE invocation_id=?1",
[invocation_id.as_str()],
|row| row.get(0),
)?;
let turn = Turn {
id: TurnId::new(),
invocation_id: invocation_id.clone(),
basis: basis.clone(),
state: BoundaryState::Starting,
provider_turn_id: None,
root_output: None,
started_at: OffsetDateTime::now_utc(),
ended_at: None,
};
tx.execute(
"INSERT INTO agent_turns (
id, invocation_id, ordinal, provider_turn_id, started_at, ended_at,
status, input_op, context_coverage, tokenizer, system_prompt_path,
task_prompt_path, system_tokens, task_tokens, supplied_context_tokens,
context_gather_ms, context_render_ms, context_persist_ms,
first_event_seq, last_event_seq, epoch_id, basis_rev
) VALUES (
?1, ?2, ?3, NULL, ?4, NULL, 'running', 'initial', 'unknown',
'cl100k_base', NULL, '', 0, 0, 0, 0, 0, 0, NULL, NULL, ?5, ?6
)",
params![
turn.id.as_str(),
turn.invocation_id.as_str(),
ordinal,
turn.started_at.unix_timestamp(),
basis.epoch_id.as_str(),
basis.revision as i64,
],
)?;
insert_seed_sends_for_turn(&tx, turn.id.as_str(), &basis)?;
AdvanceReceipt::Turn(turn)
}
RunAdvance::TurnActive {
turn_id,
provider_turn_id,
} => {
require_turn_for_run(&tx, turn_id, &run.id)?;
tx.execute(
"UPDATE agent_turns SET provider_turn_id=?2
WHERE id=?1 AND status='running'",
params![turn_id.as_str(), provider_turn_id],
)?;
AdvanceReceipt::Turn(control_turn_in(&tx, turn_id)?)
}
RunAdvance::TurnEnded { turn_id, outcome } => {
if !outcome.is_terminal() {
return Err(StoreError::InvalidData(
"Turn outcome must be terminal".to_string(),
));
}
require_turn_for_run(&tx, turn_id, &run.id)?;
tx.execute(
"UPDATE agent_turns SET status=?2, ended_at=COALESCE(ended_at, ?3)
WHERE id=?1 AND status='running'",
params![turn_id.as_str(), outcome.as_turn_status(), now_unix()],
)?;
AdvanceReceipt::Turn(control_turn_in(&tx, turn_id)?)
}
RunAdvance::Wait { on } => {
let open: bool = tx.query_row(
"SELECT EXISTS(
SELECT 1 FROM agent_invocations
WHERE supervising_run_id=?1 AND ended_at IS NULL
)",
[run.id.as_str()],
|row| row.get(0),
)?;
if open {
return Err(StoreError::InvalidData(
"Run cannot wait while an Invocation is open".to_string(),
));
}
let wait = Wait {
id: WaitId::new(),
work: run.work.clone(),
epoch_id: run.epoch_id.clone(),
on: on.clone(),
created_at: OffsetDateTime::now_utc(),
resolved_at: None,
};
tx.execute(
"INSERT INTO waits (id, epoch_id, on_json, created_at, resolved_at)
VALUES (?1, ?2, ?3, ?4, NULL)",
params![
wait.id.as_str(),
wait.epoch_id.as_str(),
serde_json::to_string(on).expect("Wait must serialize"),
wait.created_at.unix_timestamp(),
],
)?;
tx.execute(
"UPDATE runs SET state='ended', ended_at=?2 WHERE id=?1",
params![run.id.as_str(), now_unix()],
)?;
AdvanceReceipt::Wait(wait)
}
};
run = run_by_id_in(&tx, &run.id)?;
let _ = run;
tx.commit()?;
Ok(receipt)
}
pub fn stop_run(
&self,
lease: &RunLease,
cause: &StopCause,
containment: ContainmentObservation,
) -> StoreResult<StopReceipt> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let run = validate_stop_lease(&tx, lease)?;
let cause_json = serde_json::to_string(cause).expect("Stop cause must serialize");
match containment {
ContainmentObservation::Absent => {
let now = now_unix();
end_open_turns_for_run(&tx, &run.id, now, "failed")?;
tx.execute(
"UPDATE agent_invocations
SET ended_at=COALESCE(ended_at, ?2),
outcome=CASE WHEN outcome='running' THEN 'failed' ELSE outcome END,
handback_state=COALESCE(handback_state, 'unknown')
WHERE supervising_run_id=?1 AND ended_at IS NULL",
params![run.id.as_str(), now],
)?;
tx.execute(
"UPDATE runs SET state='ended', ended_at=?2, stop_reason=?3 WHERE id=?1",
params![run.id.as_str(), now, cause_json],
)?;
}
ContainmentObservation::Present | ContainmentObservation::Unprovable => {
tx.execute(
"UPDATE runs
SET state=CASE WHEN state='reserved' THEN 'reserved' ELSE 'stopping' END,
stop_reason=?2
WHERE id=?1",
params![run.id.as_str(), cause_json],
)?;
}
}
let run = run_by_id_in(&tx, &run.id)?;
tx.commit()?;
Ok(StopReceipt { run, containment })
}
pub(crate) fn finish_run(&self, lease: &RunLease, outcome: BoundaryState) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
end_run_for_lease(&conn, lease, outcome)
}
pub fn run_control(
&self,
lease: &RunLease,
active_turn_id: Option<&str>,
) -> StoreResult<Option<crate::durable::RunControl>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let run = validate_stop_lease(&conn, lease)?;
let epoch_state: String = conn.query_row(
"SELECT state FROM epochs WHERE id=?1",
[run.epoch_id.as_str()],
|row| row.get(0),
)?;
if epoch_state == "abandoned" {
let reason = conn
.query_row(
"SELECT stop_reason FROM runs WHERE id=?1",
[run.id.as_str()],
|row| row.get::<_, Option<String>>(0),
)?
.unwrap_or_else(|| "Work was abandoned".to_string());
return Ok(Some(crate::durable::RunControl::Abandon { reason }));
}
if run.state == RunState::Stopping {
let cause = conn
.query_row(
"SELECT stop_reason FROM runs WHERE id=?1",
[run.id.as_str()],
|row| row.get::<_, Option<String>>(0),
)?
.and_then(|value| serde_json::from_str::<StopCause>(&value).ok());
if let Some(StopCause::HomeUpgrade {
upgrade_id,
deadline,
}) = cause
{
return Ok(Some(crate::durable::RunControl::Quiesce {
upgrade_id,
deadline,
}));
}
return Ok(Some(crate::durable::RunControl::Interrupt));
}
let Some(turn_id) = active_turn_id else {
return Ok(None);
};
let interrupted = conn
.query_row(
"SELECT t.status='interrupted'
FROM agent_turns t
JOIN agent_invocations l ON l.id=t.invocation_id
WHERE t.id=?1 AND l.supervising_run_id=?2",
params![turn_id, run.id.as_str()],
|row| row.get::<_, bool>(0),
)
.optional()?
.unwrap_or(false);
Ok(interrupted.then_some(crate::durable::RunControl::Interrupt))
}
pub fn set_flow_position(
&self,
lease: &RunLease,
position: &FlowPosition,
) -> StoreResult<FlowPosition> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let run = validate_run_lease(&tx, lease)?;
if position.work != run.work || position.epoch_id != run.epoch_id {
return Err(StoreError::InvalidAuthority(
"flow position does not belong to the active Run".to_string(),
));
}
if position.flow.trim().is_empty() || position.step.trim().is_empty() {
return Err(StoreError::InvalidData(
"flow and step cannot be empty".to_string(),
));
}
if position.human && position.node_id.is_none() {
return Err(StoreError::InvalidData(
"human flow positions require a stable node id".to_string(),
));
}
tx.execute(
"INSERT INTO work_flow_positions (
epoch_id, flow, step, node_id, human, step_index, iteration, updated_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
ON CONFLICT(epoch_id) DO UPDATE SET
flow=excluded.flow, step=excluded.step, node_id=excluded.node_id,
human=excluded.human, step_index=excluded.step_index,
iteration=excluded.iteration,
updated_at=excluded.updated_at",
params![
position.epoch_id.as_str(),
position.flow,
position.step,
position.node_id,
position.human,
i64::from(position.step_index),
i64::from(position.iteration),
position.updated_at.unix_timestamp(),
],
)?;
tx.commit()?;
Ok(position.clone())
}
pub fn create_ask(
&self,
lease: &RunLease,
origin: &AskOrigin,
request: &AskBody,
target: &AskTarget,
) -> StoreResult<Ask> {
validate_ask_body(request)?;
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let run = validate_ask_origin(&tx, lease, origin)?;
validate_requested_target(&tx, &run.work, target)?;
if let AskBody::FlowStep {
flow,
node_id,
skill,
iteration,
} = request
{
if let Some(existing) =
flow_ask_in(&tx, &run.epoch_id, flow, node_id, skill, *iteration, target)?
{
tx.commit()?;
return Ok(existing);
}
}
let ask = Ask {
id: AskId::new(),
origin: origin.clone(),
target: target.clone(),
request: request.clone(),
state: AskState::Queued,
active_invocation_id: None,
result: None,
terminal_author: None,
asked_at: OffsetDateTime::from_unix_timestamp(now_unix()).map_err(invalid_durable)?,
terminal_at: None,
};
insert_ask(&tx, &run.epoch_id, &ask)?;
enqueue_ask_comment(&tx, &run.epoch_id, &ask)?;
tx.commit()?;
Ok(ask)
}
pub fn ask_by_id(&self, ask_id: &AskId) -> StoreResult<Ask> {
let conn = self.conn.lock().expect("store mutex poisoned");
ask_by_id_in(&conn, ask_id)
}
pub fn pending_asks(
&self,
caller: Option<&RunLease>,
target: &AskTarget,
) -> StoreResult<Vec<Ask>> {
let conn = self.conn.lock().expect("store mutex poisoned");
validate_target_caller(&conn, caller, target)?;
query_asks(&conn, AskScope::Target(target))
}
pub fn claim_ask(
&self,
caller: Option<&RunLease>,
ask_id: &AskId,
route: &InvocationRoute,
surface: &str,
) -> StoreResult<AskClaim> {
if route.provider.trim().is_empty() || surface.trim().is_empty() {
return Err(StoreError::InvalidData(
"Ask provider and surface cannot be empty".to_string(),
));
}
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let ask = ask_by_id_in(&tx, ask_id)?;
validate_target_caller(&tx, caller, &ask.target)?;
require_open_ask_epoch(&tx, ask_id)?;
if ask.state.is_terminal() {
return Err(StoreError::InvalidAuthority(format!(
"Ask {ask_id} is already terminal"
)));
}
if ask.state == AskState::Claimed {
let invocation_id = ask.active_invocation_id.as_ref().ok_or_else(|| {
StoreError::InvalidData(format!("claimed Ask {ask_id} has no active session"))
})?;
tx.commit()?;
return Ok(AskClaim {
invocation_id: invocation_id.clone(),
needs_launch: false,
});
}
let invocation_id = AgentInvocationId::new();
let run = ask_run_in(&tx, &ask)?;
let invocation = AgentInvocation {
id: invocation_id,
supervising_run_id: Some(run.id.clone()),
answer_ask_id: Some(ask.id.clone()),
route: route.clone(),
surface: surface.to_string(),
resume_token: None,
started_at: OffsetDateTime::now_utc(),
ended_at: None,
};
insert_ask_invocation(
&tx,
&run,
&ask,
&invocation,
&ask_containment_id(&invocation.id),
)?;
if tx.execute(
"UPDATE ask_exchanges SET state='claimed', active_invocation_id=?2
WHERE id=?1 AND state='queued' AND active_invocation_id IS NULL",
params![ask_id.as_str(), invocation.id.as_str()],
)? != 1
{
return Err(StoreError::InvalidAuthority(format!(
"Ask {ask_id} was claimed concurrently"
)));
}
tx.commit()?;
Ok(AskClaim {
invocation_id: invocation.id,
needs_launch: true,
})
}
pub(crate) fn claim_flow_step_run_lease(
&self,
ask_id: &AskId,
invocation_id: &AgentInvocationId,
) -> StoreResult<Option<RunLease>> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let invocation = validate_ask_invocation(&tx, ask_id, invocation_id)?;
let ask = ask_by_id_in(&tx, ask_id)?;
if !matches!(ask.request, AskBody::FlowStep { .. }) {
tx.commit()?;
return Ok(None);
}
if ask.state != AskState::Claimed
|| ask.active_invocation_id.as_ref() != Some(&invocation.id)
|| invocation.ended_at.is_some()
|| ask_presentation_in(&tx, &invocation.id)?.0
{
return Err(StoreError::InvalidAuthority(format!(
"Ask Invocation {} cannot claim flow-step writer authority",
invocation.id
)));
}
validate_flow_step_position(&tx, &ask)?;
let run = current_flow_step_run(&tx, &ask)?.ok_or_else(|| {
StoreError::InvalidAuthority(format!("Ask {} has no active flow-step Run", ask.id))
})?;
if run.state != RunState::Active || run.cwd.as_ref() != Some(&ask.origin.cwd) {
return Err(StoreError::InvalidAuthority(format!(
"Ask {} current Run is not the active flow-step writer in its captured cwd",
ask.id
)));
}
let epoch = current_epoch_in(&tx, &run.work)?;
if epoch.id != run.epoch_id {
return Err(StoreError::InvalidAuthority(format!(
"Ask {} origin Epoch is no longer current",
ask.id
)));
}
let token = RunLeaseToken::new();
if tx.execute(
"UPDATE runs SET lease_hash=?2 WHERE id=?1 AND state='active'",
params![run.id.as_str(), token.hash()],
)? != 1
{
return Err(StoreError::InvalidAuthority(format!(
"Ask {} lost flow-step writer authority",
ask.id
)));
}
let run_lease = RunLease::new(run.id, run.work, epoch.current_basis, token);
tx.commit()?;
Ok(Some(run_lease))
}
pub fn mark_ask_ready(
&self,
ask_id: &AskId,
invocation_id: &AgentInvocationId,
) -> StoreResult<AgentInvocation> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let invocation = validate_active_ask_invocation(&tx, ask_id, invocation_id)?;
if invocation.ended_at.is_some() {
return Err(StoreError::InvalidAuthority(format!(
"Ask Invocation {} is already closed",
invocation.id
)));
}
if tx.execute(
"UPDATE agent_invocations
SET ask_ready_at=COALESCE(ask_ready_at, ?2)
WHERE id=?1 AND ended_at IS NULL",
params![invocation.id.as_str(), now_unix()],
)? != 1
{
return Err(StoreError::NotFound);
}
let invocation = supervised_invocation_in(&tx, &invocation.id)?;
tx.commit()?;
Ok(invocation)
}
pub fn mark_ask_presented(
&self,
ask_id: &AskId,
invocation_id: &AgentInvocationId,
) -> StoreResult<AgentInvocation> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let invocation = validate_active_ask_invocation(&tx, ask_id, invocation_id)?;
if !ask_presentation_in(&tx, &invocation.id)?.0 || invocation.ended_at.is_some() {
return Err(StoreError::InvalidAuthority(format!(
"Ask Invocation {} is not attachable",
invocation.id
)));
}
present_ask_invocation(&tx, &invocation.id)?;
let invocation = supervised_invocation_in(&tx, &invocation.id)?;
tx.commit()?;
Ok(invocation)
}
pub fn mark_presented_by_target(
&self,
caller: Option<&RunLease>,
ask_id: &AskId,
expected_invocation_id: &AgentInvocationId,
) -> StoreResult<AgentInvocation> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let ask = ask_by_id_in(&tx, ask_id)?;
validate_target_caller(&tx, caller, &ask.target)?;
let invocation_id = ask.active_invocation_id.ok_or_else(|| {
StoreError::InvalidData(format!("Ask {ask_id} has no active Invocation"))
})?;
if &invocation_id != expected_invocation_id {
return Err(StoreError::InvalidAuthority(format!(
"Ask {ask_id} owns Invocation {invocation_id}, not {expected_invocation_id}"
)));
}
let invocation = supervised_invocation_in(&tx, &invocation_id)?;
if !ask_presentation_in(&tx, &invocation.id)?.0 || invocation.ended_at.is_some() {
return Err(StoreError::InvalidAuthority(format!(
"Ask Invocation {} is not attachable",
invocation.id
)));
}
present_ask_invocation(&tx, &invocation.id)?;
let invocation = supervised_invocation_in(&tx, &invocation.id)?;
tx.commit()?;
Ok(invocation)
}
pub fn interrupt_ask_on_interrupt(
&self,
ask_id: &AskId,
invocation_id: &AgentInvocationId,
) -> StoreResult<Ask> {
self.close_ask_invocation(
ask_id,
invocation_id,
Some("Ask process interrupted"),
BoundaryState::Interrupted,
)
}
pub fn settle_ask(
&self,
ask_id: &AskId,
invocation_id: &AgentInvocationId,
result: &AskResult,
) -> StoreResult<Ask> {
if matches!(result, AskResult::Cancelled { .. }) {
return Err(StoreError::InvalidAuthority(
"an Ask Invocation cannot cancel its Ask".to_string(),
));
}
validate_ask_result(result)?;
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let invocation = validate_ask_invocation(&tx, ask_id, invocation_id)?;
let ask = ask_by_id_in(&tx, ask_id)?;
if ask.state.is_terminal() {
if ask_invocation_is_latest_in(&tx, &ask.id, &invocation.id)?
&& ask.result.as_ref() == Some(result)
{
tx.commit()?;
return Ok(ask);
}
return Err(StoreError::InvalidAuthority(format!(
"Ask {ask_id} was already settled"
)));
}
if ask.state != AskState::Claimed
|| ask.active_invocation_id.as_ref() != Some(&invocation.id)
{
return Err(StoreError::InvalidAuthority(format!(
"Invocation {} no longer owns Ask {ask_id}",
invocation.id
)));
}
if !ask_presentation_in(&tx, &invocation.id)?.1 || invocation.ended_at.is_some() {
return Err(StoreError::InvalidAuthority(format!(
"Ask Invocation {} is not ready to settle",
invocation.id
)));
}
validate_flow_step_position(&tx, &ask)?;
let now = now_unix();
let author = ask_author_in(&tx, &ask)?;
finish_ask_invocation_in(&tx, &invocation.id, BoundaryState::Succeeded, None)?;
let settled = write_terminal_ask_in(&tx, &ask, result, &author, now)?;
enqueue_ask_result_comment(&tx, &settled)?;
end_flow_step_run(&tx, &settled)?;
tx.commit()?;
Ok(settled)
}
pub fn release_ask(
&self,
ask_id: &AskId,
invocation_id: &AgentInvocationId,
reason: Option<&str>,
) -> StoreResult<Ask> {
self.close_ask_invocation(ask_id, invocation_id, reason, BoundaryState::Unknown)
}
pub fn close_ask_invocation(
&self,
ask_id: &AskId,
invocation_id: &AgentInvocationId,
reason: Option<&str>,
outcome: BoundaryState,
) -> StoreResult<Ask> {
if !outcome.is_terminal() {
return Err(StoreError::InvalidData(
"Ask Invocation outcome must be terminal".to_string(),
));
}
let reason = normalize_optional_reason(reason)?;
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let invocation = validate_ask_invocation(&tx, ask_id, invocation_id)?;
let ask = ask_by_id_in(&tx, ask_id)?;
if ask.state != AskState::Claimed
|| ask.active_invocation_id.as_ref() != Some(&invocation.id)
{
return Err(StoreError::InvalidAuthority(format!(
"Invocation {} no longer owns Ask {ask_id}",
invocation.id
)));
}
requeue_ask_in(&tx, &ask, &invocation.id, outcome, reason.as_deref())?;
let ask = ask_by_id_in(&tx, &ask.id)?;
tx.commit()?;
Ok(ask)
}
pub fn escalate_ask(
&self,
ask_id: &AskId,
invocation_id: &AgentInvocationId,
) -> StoreResult<Ask> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let invocation = validate_ask_invocation(&tx, ask_id, invocation_id)?;
let ask = ask_by_id_in(&tx, ask_id)?;
if !matches!(ask.target, AskTarget::Parent(_)) {
return Err(StoreError::InvalidData(format!(
"Ask {} already targets the User",
ask.id
)));
}
if ask.state != AskState::Claimed
|| ask.active_invocation_id.as_ref() != Some(&invocation.id)
{
return Err(StoreError::InvalidAuthority(format!(
"Invocation {} no longer owns Ask {}",
invocation.id, ask.id
)));
}
finish_ask_invocation_in(
&tx,
&invocation.id,
BoundaryState::Unknown,
Some("escalated to User"),
)?;
tx.execute(
"UPDATE ask_exchanges
SET target_kind='user', target_work_kind=NULL, target_work_id=NULL,
state='queued', active_invocation_id=NULL
WHERE id=?1 AND state='claimed' AND active_invocation_id=?2",
params![ask.id.as_str(), invocation.id.as_str()],
)?;
let ask = ask_by_id_in(&tx, &ask.id)?;
tx.commit()?;
Ok(ask)
}
pub fn escalate_queued_ask(
&self,
caller: Option<&RunLease>,
ask_id: &AskId,
) -> StoreResult<Ask> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let ask = ask_by_id_in(&tx, ask_id)?;
validate_target_caller(&tx, caller, &ask.target)?;
if !matches!(ask.target, AskTarget::Parent(_)) || ask.state != AskState::Queued {
return Err(StoreError::InvalidAuthority(format!(
"Ask {ask_id} is not a queued parent request"
)));
}
tx.execute(
"UPDATE ask_exchanges
SET target_kind='user', target_work_kind=NULL, target_work_id=NULL
WHERE id=?1 AND state='queued'",
[ask_id.as_str()],
)?;
let ask = ask_by_id_in(&tx, ask_id)?;
tx.commit()?;
Ok(ask)
}
pub fn cancel_ask(
&self,
caller: Option<&RunLease>,
ask_id: &AskId,
reason: &str,
) -> StoreResult<Ask> {
let reason = normalize_reason(reason, "Ask cancellation reason")?;
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let ask = ask_by_id_in(&tx, ask_id)?;
let author = match caller {
None => Author::User,
Some(lease) => {
let run = validate_run_lease(&tx, lease)?;
let epoch_id = ask_epoch_id_in(&tx, ask_id)?;
if run.work != ask.origin.work || run.epoch_id != epoch_id {
return Err(StoreError::InvalidAuthority(
"Run may cancel only an Ask from its current Work Epoch".to_string(),
));
}
Author::Run(run.id)
}
};
if ask.state.is_terminal() {
if ask.result
== Some(AskResult::Cancelled {
reason: reason.clone(),
})
{
tx.commit()?;
return Ok(ask);
}
return Err(StoreError::InvalidAuthority(format!(
"Ask {ask_id} is already terminal"
)));
}
let now = now_unix();
if let Some(invocation_id) = ask.active_invocation_id.as_ref() {
finish_ask_invocation_in(&tx, invocation_id, BoundaryState::Unknown, Some(&reason))?;
}
let ask = write_terminal_ask_in(&tx, &ask, &AskResult::Cancelled { reason }, &author, now)?;
enqueue_ask_result_comment(&tx, &ask)?;
end_flow_step_run(&tx, &ask)?;
tx.commit()?;
Ok(ask)
}
pub fn reconcile_ask(
&self,
invocation_id: &AgentInvocationId,
observation: ContainmentObservation,
) -> StoreResult<Ask> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let invocation = supervised_invocation_in(&tx, invocation_id)?;
let ask_id = invocation.answer_ask_id.as_ref().ok_or_else(|| {
StoreError::InvalidData(format!(
"Invocation {invocation_id} does not belong to an Ask"
))
})?;
let ask = ask_by_id_in(&tx, ask_id)?;
if invocation.ended_at.is_some()
|| ask.state.is_terminal()
|| observation != ContainmentObservation::Absent
|| ask.active_invocation_id.as_ref() != Some(invocation_id)
{
tx.commit()?;
return Ok(ask);
}
requeue_ask_in(
&tx,
&ask,
invocation_id,
BoundaryState::Unknown,
Some("Ask session disappeared"),
)?;
let ask = ask_by_id_in(&tx, &ask.id)?;
tx.commit()?;
Ok(ask)
}
pub fn ask_invocations(&self, ask_id: &AskId) -> StoreResult<Vec<AgentInvocation>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut statement = conn.prepare(
"SELECT id FROM agent_invocations WHERE answer_ask_id=?1
ORDER BY started_at, rowid",
)?;
let ids = statement
.query_map([ask_id.as_str()], |row| row.get::<_, String>(0))?
.collect::<Result<Vec<_>, _>>()?;
ids.into_iter()
.map(|id| {
let id = AgentInvocationId::parse(&id).map_err(invalid_durable)?;
supervised_invocation_in(&conn, &id)
})
.collect()
}
pub fn ask_presentation(&self, invocation_id: &AgentInvocationId) -> StoreResult<(bool, bool)> {
let conn = self.conn.lock().expect("store mutex poisoned");
ask_presentation_in(&conn, invocation_id)
}
pub fn request_intervention(
&self,
lease: &RunLease,
invocation_id: &AgentInvocationId,
prompt: &str,
user: bool,
) -> StoreResult<Ask> {
let prompt = prompt.trim();
if prompt.is_empty() {
return Err(StoreError::InvalidData(
"Ask request cannot be empty".to_string(),
));
}
let conn = self.conn.lock().expect("store mutex poisoned");
let run = validate_run_lease(&conn, lease)?;
require_open_invocation_for_run(&conn, invocation_id, &run.id)?;
let turn_id = current_turn_for_invocation_in(&conn, invocation_id)?;
let request = AskBody::Intervention {
prompt: prompt.to_string(),
};
let target = if user {
AskTarget::User
} else {
let Some(parent) = parent_work(&conn, &run.work)? else {
return Err(StoreError::InvalidData(format!(
"root {} {} has no parent; use `lf ask --user` for genuine User intervention",
run.work.kind(),
run.work.id()
)));
};
AskTarget::Parent(parent)
};
let origin = AskOrigin {
work: run.work.clone(),
run_id: run.id,
turn_id: Some(turn_id),
invocation_id: Some(invocation_id.clone()),
home_id: run.home_id,
cwd: run.cwd.ok_or_else(|| {
StoreError::InvalidData("Ask origin Run has no execution cwd".to_string())
})?,
};
drop(conn);
self.create_ask(lease, &origin, &request, &target)
}
pub fn asks_for_work_epoch(&self, lease: &RunLease) -> StoreResult<Vec<Ask>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let run = validate_run_lease(&conn, lease)?;
query_asks(
&conn,
AskScope::OriginEpoch {
work: &run.work,
epoch_id: &run.epoch_id,
},
)
}
pub(crate) fn pending_ask_comment_writes(&self) -> StoreResult<Vec<AskCommentWrite>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut statement = conn.prepare(
"SELECT o.ask_id, o.transition, o.issue_id, o.body,
w.repo, w.name, o.attempt_count, o.attempt_started_at,
o.last_error, o.linear_comment_id, o.delivered_at
FROM ask_linear_comment_outbox o
JOIN tasks t ON t.id=o.task_id
JOIN projects p ON p.id=t.project_id
JOIN waves w ON w.id=p.wave_id
WHERE o.delivered_at IS NULL
ORDER BY o.created_at, o.ask_id,
CASE o.transition WHEN 'ask' THEN 0 ELSE 1 END",
)?;
let rows = statement
.query_map([], read_ask_comment_write_row)?
.collect::<Result<Vec<_>, _>>()?;
rows.into_iter().map(parse_ask_comment_write_row).collect()
}
pub(crate) fn claim_ask_comment_write(
&self,
ask_id: &AskId,
transition: AskCommentTransition,
attempted_at: i64,
stale_before: i64,
) -> StoreResult<Option<AskCommentWrite>> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let changed = tx.execute(
"UPDATE ask_linear_comment_outbox
SET attempt_count=attempt_count + 1, attempt_started_at=?3
WHERE ask_id=?1 AND transition=?2 AND delivered_at IS NULL
AND (attempt_started_at IS NULL OR attempt_started_at <= ?4)",
params![
ask_id.as_str(),
transition.as_str(),
attempted_at,
stale_before
],
)?;
let write = if changed == 1 {
Some(ask_comment_write_in(&tx, ask_id, transition)?)
} else {
None
};
tx.commit()?;
Ok(write)
}
pub(crate) fn complete_ask_comment_write(
&self,
ask_id: &AskId,
transition: AskCommentTransition,
comment_id: &str,
delivered_at: i64,
) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let changed = conn.execute(
"UPDATE ask_linear_comment_outbox
SET linear_comment_id=?3, delivered_at=?4,
attempt_started_at=NULL, last_error=NULL
WHERE ask_id=?1 AND transition=?2 AND delivered_at IS NULL",
params![
ask_id.as_str(),
transition.as_str(),
comment_id,
delivered_at
],
)?;
if changed == 0 {
let existing = ask_comment_write_in(&conn, ask_id, transition)?;
if existing.linear_comment_id.as_deref() == Some(comment_id) {
return Ok(());
}
return Err(StoreError::InvalidData(format!(
"Ask {ask_id} {} comment write completed concurrently",
transition.as_str()
)));
}
Ok(())
}
pub(crate) fn fail_ask_comment_write(
&self,
ask_id: &AskId,
transition: AskCommentTransition,
error: &str,
) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
if conn.execute(
"UPDATE ask_linear_comment_outbox
SET attempt_started_at=NULL, last_error=?3
WHERE ask_id=?1 AND transition=?2 AND delivered_at IS NULL",
params![ask_id.as_str(), transition.as_str(), error],
)? == 0
{
return Err(StoreError::NotFound);
}
Ok(())
}
pub fn has_pending_user_ask_for_work(&self, work: &WorkRef) -> StoreResult<bool> {
let conn = self.conn.lock().expect("store mutex poisoned");
let (wave_id, project_id, task_id) = match work {
WorkRef::Wave(id) => (Some(id.as_str()), None, None),
WorkRef::Project(id) => (None, Some(id.as_str()), None),
WorkRef::Task(id) => (None, None, Some(id.as_str())),
};
conn.query_row(
HAS_PENDING_USER_ASK_FOR_WORK_SQL,
params![wave_id, project_id, task_id],
|row| row.get(0),
)
.map_err(StoreError::from)
}
pub fn invocation_surface(
&self,
invocation_id: &AgentInvocationId,
) -> StoreResult<Option<InvocationSurface>> {
let conn = self.conn.lock().expect("store mutex poisoned");
invocation_surface_in(&conn, invocation_id)
}
pub(crate) fn open_invocation_for_run(
&self,
run_id: &RunId,
) -> StoreResult<Option<AgentInvocation>> {
let conn = self.conn.lock().expect("store mutex poisoned");
open_invocation_for_run_in(&conn, run_id)
}
pub(crate) fn open_invocation_for_run_by_id(
&self,
run_id: &RunId,
invocation_id: &AgentInvocationId,
) -> StoreResult<Option<AgentInvocation>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let exists = conn
.query_row(
"SELECT 1 FROM agent_invocations
WHERE id=?1 AND supervising_run_id=?2 AND ended_at IS NULL",
params![invocation_id.as_str(), run_id.as_str()],
|_| Ok(()),
)
.optional()?
.is_some();
exists
.then(|| supervised_invocation_in(&conn, invocation_id))
.transpose()
}
pub(crate) fn invocations_for_run(&self, run_id: &RunId) -> StoreResult<Vec<AgentInvocation>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut statement = conn.prepare(
"SELECT id FROM agent_invocations
WHERE supervising_run_id=?1 ORDER BY started_at, rowid",
)?;
let ids = statement
.query_map([run_id.as_str()], |row| row.get::<_, String>(0))?
.collect::<Result<Vec<_>, _>>()?;
ids.into_iter()
.map(|id| {
let id = AgentInvocationId::parse(&id).map_err(invalid_durable)?;
supervised_invocation_in(&conn, &id)
})
.collect()
}
pub(crate) fn recover_run(
&self,
run_id: &RunId,
containment: ContainmentObservation,
) -> StoreResult<StopReceipt> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let run = run_by_id_in(&tx, run_id)?;
if run.state == RunState::Ended {
tx.commit()?;
return Ok(StopReceipt { run, containment });
}
let cause = serde_json::to_string(&StopCause::Recovery).expect("Stop cause must serialize");
match containment {
ContainmentObservation::Absent => {
let now = now_unix();
end_open_turns_for_run(&tx, run_id, now, "failed")?;
tx.execute(
"UPDATE agent_invocations
SET ended_at=COALESCE(ended_at, ?2),
outcome=CASE WHEN outcome='running' THEN 'failed' ELSE outcome END,
handback_state=COALESCE(handback_state, 'unknown')
WHERE supervising_run_id=?1 AND ended_at IS NULL",
params![run_id.as_str(), now],
)?;
tx.execute(
"UPDATE runs SET state='ended', ended_at=?2, stop_reason=?3
WHERE id=?1 AND state != 'ended'",
params![run_id.as_str(), now, cause],
)?;
}
ContainmentObservation::Present | ContainmentObservation::Unprovable => {
tx.execute(
"UPDATE runs
SET state=CASE WHEN state='reserved' THEN 'reserved' ELSE 'stopping' END,
stop_reason=?2
WHERE id=?1 AND state IN ('reserved', 'active', 'stopping')",
params![run_id.as_str(), cause],
)?;
}
}
let run = run_by_id_in(&tx, run_id)?;
tx.commit()?;
Ok(StopReceipt { run, containment })
}
pub fn invocation_surfaces(&self, active_only: bool) -> StoreResult<Vec<InvocationSurface>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let sql = if active_only {
"SELECT i.id FROM agent_invocations i
JOIN runs r ON r.id=i.supervising_run_id
WHERE i.ended_at IS NULL AND r.state IN ('active', 'stopping')
ORDER BY i.started_at, i.id"
} else {
"SELECT id FROM agent_invocations
WHERE supervising_run_id IS NOT NULL ORDER BY started_at, id"
};
let mut statement = conn.prepare(sql)?;
let ids = statement
.query_map([], |row| row.get::<_, String>(0))?
.collect::<Result<Vec<_>, _>>()?;
ids.into_iter()
.map(|id| {
let id = AgentInvocationId::parse(&id).map_err(invalid_durable)?;
invocation_surface_in(&conn, &id)?.ok_or(StoreError::NotFound)
})
.collect()
}
pub fn observe_invocation_provider(
&self,
lease: &RunLease,
invocation_id: &AgentInvocationId,
account_id: Option<&crate::store::ProviderAccountId>,
resume_token: Option<&str>,
) -> StoreResult<AgentInvocation> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let run = validate_run_lease(&tx, lease)?;
if !matches!(run.state, RunState::Reserved | RunState::Active) {
return Err(StoreError::InvalidData(format!(
"Run {} cannot observe a provider while {:?}",
run.id, run.state
)));
}
if tx.execute(
"UPDATE agent_invocations
SET account_id=COALESCE(account_id, ?3),
resume_token=COALESCE(?4, resume_token),
provider_session_id=COALESCE(?4, provider_session_id)
WHERE id=?1 AND supervising_run_id=?2 AND ended_at IS NULL
AND (?3 IS NULL OR account_id IS NULL OR account_id=?3)",
params![
invocation_id.as_str(),
lease.run_id.as_str(),
account_id.map(crate::store::ProviderAccountId::as_str),
resume_token,
],
)? == 0
{
return Err(StoreError::NotFound);
}
let invocation = supervised_invocation_in(&tx, invocation_id)?;
tx.commit()?;
Ok(invocation)
}
pub fn handback_invocation(
&self,
invocation_id: &AgentInvocationId,
outcome: BoundaryState,
) -> StoreResult<InvocationSurface> {
if !outcome.is_terminal() {
return Err(StoreError::InvalidData(
"Invocation handback outcome must be terminal".to_string(),
));
}
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
if tx.execute(
"UPDATE agent_invocations
SET ended_at=COALESCE(ended_at, ?2),
outcome=?3, handback_state=?4
WHERE id=?1 AND supervising_run_id IS NOT NULL AND ended_at IS NULL
AND surface IN ('tui', 'ide') AND answer_ask_id IS NULL",
params![
invocation_id.as_str(),
now_unix(),
outcome.as_invocation_outcome(),
handback_state(outcome),
],
)? == 0
{
return Err(StoreError::NotFound);
}
let surface = invocation_surface_in(&tx, invocation_id)?.ok_or(StoreError::NotFound)?;
tx.commit()?;
Ok(surface)
}
pub fn interrupt(
&self,
caller: Option<&RunLease>,
work: &WorkRef,
if_run: &RunId,
) -> StoreResult<InterruptReceipt> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
validate_control_caller(&tx, caller, work)?;
let run = current_run_for_work_in(&tx, work)?.ok_or(StoreError::NotFound)?;
if &run.id != if_run {
return Err(StoreError::InvalidAuthority(format!(
"Run {if_run} is not current for {}",
work.id()
)));
}
let mut statement = tx.prepare(
"SELECT t.id FROM agent_turns t
JOIN agent_invocations i ON i.id=t.invocation_id
WHERE i.supervising_run_id=?1 AND t.status='running'
ORDER BY t.started_at, t.ordinal",
)?;
let turn_ids = statement
.query_map([run.id.as_str()], |row| row.get::<_, String>(0))?
.collect::<Result<Vec<_>, _>>()?;
drop(statement);
let now = now_unix();
tx.execute(
"UPDATE agent_turns SET status='interrupted', ended_at=?2
WHERE status='running' AND invocation_id IN (
SELECT id FROM agent_invocations WHERE supervising_run_id=?1
)",
params![run.id.as_str(), now],
)?;
let cause =
serde_json::to_string(&StopCause::Interrupted).expect("interrupt cause must serialize");
if run.state == RunState::Reserved {
tx.execute(
"UPDATE runs SET state='ended', ended_at=?2, stop_reason=?3 WHERE id=?1",
params![run.id.as_str(), now, cause],
)?;
} else {
tx.execute(
"UPDATE runs SET state='stopping', stop_reason=?2 WHERE id=?1",
params![run.id.as_str(), cause],
)?;
}
let receipt = InterruptReceipt {
run_id: run.id,
turn_ids: turn_ids
.into_iter()
.map(|id| TurnId::parse(&id).map_err(invalid_durable))
.collect::<StoreResult<Vec<_>>>()?,
};
tx.commit()?;
Ok(receipt)
}
pub fn done(&self, lease: &RunLease, basis: &Basis) -> StoreResult<DoneProposal> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let run = validate_run_lease(&tx, lease)?;
let epoch = current_epoch_in(&tx, &run.work)?;
validate_basis(&epoch.current_basis, basis)?;
let applied = applied_basis_in(&tx, &epoch.id)?.ok_or_else(|| {
StoreError::InvalidData("no successful boundary can complete Work".to_string())
})?;
validate_basis(&applied, basis)?;
validate_completion_readiness_in(&tx, &run)?;
let proposal = DoneProposal {
id: DoneProposalId::new(),
run_id: run.id.clone(),
basis: basis.clone(),
proposed_at: OffsetDateTime::now_utc(),
};
tx.execute(
"INSERT INTO done_proposals (id, run_id, epoch_id, basis_rev, proposed_at)
VALUES (?1, ?2, ?3, ?4, ?5)",
params![
proposal.id.as_str(),
proposal.run_id.as_str(),
proposal.basis.epoch_id.as_str(),
proposal.basis.revision as i64,
proposal.proposed_at.unix_timestamp(),
],
)?;
let now = now_unix();
cancel_pending_asks_for_epoch(&tx, &epoch.id, "owning Work Epoch completed", now)?;
tx.execute(
"UPDATE epochs SET state='done', terminal_at=?2
WHERE id=?1 AND state='open' AND current_rev=?3",
params![epoch.id.as_str(), now, basis.revision as i64],
)?;
tx.execute(
"UPDATE runs SET state='ended', ended_at=?2 WHERE id=?1",
params![run.id.as_str(), now],
)?;
tx.commit()?;
Ok(proposal)
}
pub fn abandon(
&self,
work: &WorkRef,
reason: &str,
if_basis: &Basis,
) -> StoreResult<EpochReceipt> {
let reason = reason.trim();
if reason.is_empty() {
return Err(StoreError::InvalidData(
"abandon reason cannot be empty".to_string(),
));
}
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let mut epoch = current_epoch_in(&tx, work)?;
validate_basis(&epoch.current_basis, if_basis)?;
let now = now_unix();
cancel_pending_asks_for_epoch(&tx, &epoch.id, "owning Work Epoch abandoned", now)?;
tx.execute(
"UPDATE runs
SET state=CASE WHEN state='reserved' THEN 'ended' ELSE 'stopping' END,
ended_at=CASE WHEN state='reserved' THEN ?3 ELSE ended_at END,
stop_reason=?2
WHERE epoch_id=?1 AND state != 'ended'",
params![epoch.id.as_str(), reason, now],
)?;
tx.execute(
"UPDATE epochs SET state='abandoned', terminal_at=?2 WHERE id=?1 AND state='open'",
params![epoch.id.as_str(), now],
)?;
epoch.state = EpochState::Abandoned;
epoch.terminal_at = Some(
OffsetDateTime::from_unix_timestamp(now).expect("current Unix timestamp must be valid"),
);
tx.commit()?;
Ok(EpochReceipt { epoch })
}
pub fn work_status(&self, work: &WorkRef) -> StoreResult<WorkStatus> {
let conn = self.conn.lock().expect("store mutex poisoned");
work_status_in(&conn, work)
}
pub fn work_for_child(&self, target: &ChildRef) -> StoreResult<WorkRef> {
let conn = self.conn.lock().expect("store mutex poisoned");
work_for_child_in(&conn, target)
}
pub fn current_epoch(&self, work: &WorkRef) -> StoreResult<Epoch> {
let conn = self.conn.lock().expect("store mutex poisoned");
current_epoch_in(&conn, work)
}
pub fn boundary_seed(&self, work: &WorkRef) -> StoreResult<BoundarySeed> {
let conn = self.conn.lock().expect("store mutex poisoned");
boundary_seed_in(&conn, work)
}
pub(crate) fn boundary_seed_for_child(&self, target: &ChildRef) -> StoreResult<BoundarySeed> {
let conn = self.conn.lock().expect("store mutex poisoned");
let work = work_for_child_in(&conn, target)?;
let (column, id) = match &work {
WorkRef::Project(id) => ("project_id", id.as_str()),
WorkRef::Task(id) => ("task_id", id.as_str()),
WorkRef::Wave(_) => unreachable!("a child is Project or Task Work"),
};
let (epoch_id, revision) = conn.query_row(
&format!(
"SELECT id, current_rev FROM epochs WHERE {column}=?1 ORDER BY number DESC LIMIT 1"
),
[id],
|row| Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)?)),
)?;
let epoch_id = EpochId::parse(&epoch_id).map_err(|error| {
StoreError::InvalidData(format!("invalid stored Epoch id: {error}"))
})?;
boundary_seed_for_epoch_in(&conn, &work, epoch_id, revision as u64)
}
pub fn append_steer(
&self,
work: &WorkRef,
author: &Author,
text: &str,
if_basis: Option<&Basis>,
) -> StoreResult<SteerReceipt> {
let text = text.trim();
if text.is_empty() {
return Err(StoreError::InvalidData(
"Steer text cannot be empty".to_string(),
));
}
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let epoch = current_epoch_in(&tx, work)?;
if let Some(expected) = if_basis {
validate_basis(&epoch.current_basis, expected)?;
}
let receipt = Self::append_steer_in(&tx, work, author, text)?;
tx.commit()?;
Ok(receipt)
}
pub fn steer(
&self,
caller: Option<&RunLease>,
work: &WorkRef,
text: &str,
if_basis: Option<&Basis>,
) -> StoreResult<SteerReceipt> {
let text = text.trim();
if text.is_empty() {
return Err(StoreError::InvalidData(
"Steer text cannot be empty".to_string(),
));
}
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let epoch = current_epoch_in(&tx, work)?;
if let Some(expected) = if_basis {
validate_basis(&epoch.current_basis, expected)?;
}
validate_control_caller(&tx, caller, work)?;
let author = caller.map_or(Author::User, |lease| Author::Run(lease.run_id.clone()));
let receipt = Self::append_steer_in(&tx, work, &author, text)?;
tx.commit()?;
Ok(receipt)
}
pub fn list_steers_since(&self, since: i64) -> StoreResult<Vec<Steer>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut statement = conn.prepare(
"SELECT s.id, s.epoch_id, s.rev, s.author_kind, s.author_run_id,
s.text, s.issued_at, e.wave_id, e.project_id, e.task_id
FROM steers s
JOIN epochs e ON e.id=s.epoch_id
WHERE s.issued_at >= ?1
ORDER BY s.issued_at DESC, s.id DESC",
)?;
let rows = statement.query_map([since], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, i64>(2)?,
row.get::<_, String>(3)?,
row.get::<_, Option<String>>(4)?,
row.get::<_, String>(5)?,
row.get::<_, i64>(6)?,
row.get::<_, Option<String>>(7)?,
row.get::<_, Option<String>>(8)?,
row.get::<_, Option<String>>(9)?,
))
})?;
let mut steers = Vec::new();
for row in rows {
let (
id,
epoch_id,
revision,
author_kind,
author_run_id,
text,
issued_at,
wave_id,
project_id,
task_id,
) = row?;
let work = parse_work_columns(wave_id, project_id, task_id)?;
steers.push(decode_steer(
(
id,
epoch_id,
revision,
author_kind,
author_run_id,
text,
issued_at,
),
work,
)?);
}
Ok(steers)
}
pub(crate) fn append_steer_in(
tx: &Transaction<'_>,
work: &WorkRef,
author: &Author,
text: &str,
) -> StoreResult<SteerReceipt> {
let text = text.trim();
if text.is_empty() {
return Err(StoreError::InvalidData(
"Steer text cannot be empty".to_string(),
));
}
let epoch = current_epoch_in(tx, work)?;
validate_author(tx, work, author)?;
let revision = epoch.current_basis.revision + 1;
let steer = Steer {
id: SteerId::new(),
work: work.clone(),
basis: Basis {
epoch_id: epoch.id.clone(),
revision,
},
author: author.clone(),
text: text.to_string(),
issued_at: OffsetDateTime::now_utc(),
};
tx.execute(
"INSERT INTO epoch_revisions (epoch_id, rev, kind, source_id, created_at)
VALUES (?1, ?2, 'steer', ?3, ?4)",
params![
steer.basis.epoch_id.as_str(),
steer.basis.revision as i64,
steer.id.as_str(),
steer.issued_at.unix_timestamp(),
],
)?;
let (author_kind, author_run_id) = match &steer.author {
Author::User => ("user", None),
Author::Run(run_id) => ("run", Some(run_id.as_str())),
};
tx.execute(
"INSERT INTO steers (
id, epoch_id, rev, author_kind, author_run_id, text, issued_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
params![
steer.id.as_str(),
steer.basis.epoch_id.as_str(),
steer.basis.revision as i64,
author_kind,
author_run_id,
steer.text,
steer.issued_at.unix_timestamp(),
],
)?;
tx.execute(
"UPDATE epochs SET current_rev=?2 WHERE id=?1 AND state='open'",
params![steer.basis.epoch_id.as_str(), revision as i64],
)?;
Ok(SteerReceipt {
steer,
sends: Vec::new(),
applied_by: None,
})
}
pub fn write_tool_response(
&self,
work: &WorkRef,
write: &ToolResponseWrite,
if_basis: Option<&Basis>,
) -> StoreResult<(ToolResponseReceipt, bool)> {
let choice = write.choice.trim();
if choice.is_empty() {
return Err(StoreError::InvalidData(
"tool response choice cannot be empty".to_string(),
));
}
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let epoch = current_epoch_in(&tx, work)?;
if let Some(expected) = if_basis {
validate_basis(&epoch.current_basis, expected)?;
}
if let Some(existing) = tool_response_in(&tx, &epoch.id, &write.request_id)? {
if existing.choice != choice {
return Err(StoreError::InvalidData(format!(
"tool response {} is already resolved as {:?}",
write.request_id, existing.choice
)));
}
return Ok((existing, false));
}
let revision = epoch.current_basis.revision + 1;
let receipt = ToolResponseReceipt {
id: ToolResponseId::new(),
work: work.clone(),
basis: Basis {
epoch_id: epoch.id.clone(),
revision,
},
request_id: write.request_id.clone(),
choice: choice.to_string(),
responded_at: OffsetDateTime::now_utc(),
};
tx.execute(
"INSERT INTO epoch_revisions (epoch_id, rev, kind, source_id, created_at)
VALUES (?1, ?2, 'tool_response', ?3, ?4)",
params![
receipt.basis.epoch_id.as_str(),
receipt.basis.revision as i64,
receipt.id.as_str(),
receipt.responded_at.unix_timestamp(),
],
)?;
tx.execute(
"INSERT INTO tool_responses (id, epoch_id, rev, request_id, choice, responded_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
params![
receipt.id.as_str(),
receipt.basis.epoch_id.as_str(),
receipt.basis.revision as i64,
receipt.request_id,
receipt.choice,
receipt.responded_at.unix_timestamp(),
],
)?;
tx.execute(
"UPDATE epochs SET current_rev=?2 WHERE id=?1 AND state='open'",
params![receipt.basis.epoch_id.as_str(), revision as i64],
)?;
tx.commit()?;
Ok((receipt, true))
}
pub fn tool_response(
&self,
work: &WorkRef,
request_id: &str,
) -> StoreResult<Option<ToolResponseReceipt>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let epoch = current_epoch_in(&conn, work)?;
tool_response_in(&conn, &epoch.id, request_id)
}
pub fn begin_live_send(&self, steer_id: &SteerId, turn_id: &str) -> StoreResult<Option<Send>> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
if let Some(existing) = send_for(&tx, steer_id, turn_id, SendVia::Live)? {
return Ok(Some(existing));
}
let eligible: bool = tx.query_row(
"SELECT EXISTS(
SELECT 1
FROM steers s
JOIN agent_turns t ON t.id=?2
WHERE s.id=?1
AND s.epoch_id=t.epoch_id
AND t.status='running'
AND s.rev > COALESCE((
SELECT MAX(done.basis_rev)
FROM agent_turns done
WHERE done.epoch_id=s.epoch_id AND done.status='completed'
), -1)
)",
params![steer_id.as_str(), turn_id],
|row| row.get(0),
)?;
if !eligible {
tx.commit()?;
return Ok(None);
}
let send = Send {
id: SendId::new(),
steer_id: steer_id.clone(),
turn_id: turn_id.to_string(),
via: SendVia::Live,
state: SendState::Sending,
provider_turn_id: None,
reason: None,
attempted_at: OffsetDateTime::now_utc(),
finished_at: None,
};
insert_send(&tx, &send)?;
tx.commit()?;
Ok(Some(send))
}
pub fn finish_send(
&self,
send_id: &SendId,
state: SendState,
provider_turn_id: Option<&str>,
reason: Option<&str>,
) -> StoreResult<Send> {
if state == SendState::Sending {
return Err(StoreError::InvalidData(
"a Send cannot finish as sending".to_string(),
));
}
let now = OffsetDateTime::now_utc();
let conn = self.conn.lock().expect("store mutex poisoned");
let changed = conn.execute(
"UPDATE sends
SET state=?2, provider_turn_id=?3, reason=?4, finished_at=?5
WHERE id=?1 AND state='sending'",
params![
send_id.as_str(),
state.as_str(),
provider_turn_id,
reason,
now.unix_timestamp(),
],
)?;
if changed == 0 {
return send_by_id(&conn, send_id)?.ok_or(StoreError::NotFound);
}
send_by_id(&conn, send_id)?.ok_or(StoreError::NotFound)
}
pub fn validate_completion_basis(&self, work: &WorkRef, proposed: &Basis) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let epoch = current_epoch_in(&conn, work)?;
validate_basis(&epoch.current_basis, proposed)?;
let applied = applied_basis_in(&conn, &epoch.id)?.ok_or_else(|| {
StoreError::InvalidData("no successful root boundary can complete Work".to_string())
})?;
validate_basis(&applied, proposed)
}
}
fn map_local_home(conn: &Connection) -> StoreResult<Home> {
conn.query_row(
"SELECT id, route, created_at, observed_at FROM homes WHERE route='local'",
[],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, i64>(2)?,
row.get::<_, i64>(3)?,
))
},
)
.map_err(StoreError::from)
.and_then(|(id, route, created_at, observed_at)| {
Ok(Home {
id: HomeId::parse(&id).map_err(invalid_durable)?,
route,
created_at: OffsetDateTime::from_unix_timestamp(created_at).map_err(invalid_durable)?,
observed_at: OffsetDateTime::from_unix_timestamp(observed_at)
.map_err(invalid_durable)?,
})
})
}
fn map_home_by_id(conn: &Connection, home_id: &HomeId) -> StoreResult<Option<Home>> {
conn.query_row(
"SELECT id, route, created_at, observed_at FROM homes WHERE id=?1",
[home_id.as_str()],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, i64>(2)?,
row.get::<_, i64>(3)?,
))
},
)
.optional()
.map_err(StoreError::from)?
.map(|(id, route, created_at, observed_at)| {
Ok(Home {
id: HomeId::parse(&id).map_err(invalid_durable)?,
route,
created_at: OffsetDateTime::from_unix_timestamp(created_at).map_err(invalid_durable)?,
observed_at: OffsetDateTime::from_unix_timestamp(observed_at)
.map_err(invalid_durable)?,
})
})
.transpose()
}
fn placement_in(conn: &Connection, work: &WorkRef) -> StoreResult<Placement> {
find_placement_in(conn, work)?.ok_or_else(|| {
StoreError::InvalidData(format!(
"{} {} has no Home placement",
work.kind(),
work.id()
))
})
}
fn reserving_home_in(conn: &Connection, work: &WorkRef) -> StoreResult<HomeId> {
let placement = placement_in(conn, work)?;
if !placement.enabled {
return Err(StoreError::InvalidData(format!(
"cannot reserve {} {}; it is disabled",
work.kind(),
work.id()
)));
}
let placed = placement.home_id;
let local = map_local_home(conn)?.id;
if placed != local {
return Err(StoreError::InvalidData(format!(
"cannot reserve {} {} on local Home {local}; it is placed on {placed}",
work.kind(),
work.id()
)));
}
Ok(local)
}
fn find_placement_in(conn: &Connection, work: &WorkRef) -> StoreResult<Option<Placement>> {
let row = match work {
WorkRef::Wave(id) => conn.query_row(
"SELECT home_id, enabled, placed_at FROM work_placements WHERE wave_id=?1",
[id.as_str()],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, bool>(1)?,
row.get::<_, i64>(2)?,
))
},
),
WorkRef::Project(id) => conn.query_row(
"SELECT home_id, enabled, placed_at FROM work_placements WHERE project_id=?1",
[id.as_str()],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, bool>(1)?,
row.get::<_, i64>(2)?,
))
},
),
WorkRef::Task(id) => conn.query_row(
"SELECT home_id, enabled, placed_at FROM work_placements WHERE task_id=?1",
[id.as_str()],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, bool>(1)?,
row.get::<_, i64>(2)?,
))
},
),
}
.optional()?;
row.map(|(home_id, enabled, placed_at)| {
Ok(Placement {
work: work.clone(),
home_id: HomeId::parse(&home_id).map_err(invalid_durable)?,
enabled,
placed_at: OffsetDateTime::from_unix_timestamp(placed_at).map_err(invalid_durable)?,
})
})
.transpose()
}
fn write_placement(
tx: &Transaction<'_>,
work: &WorkRef,
home_id: &HomeId,
placed_at: i64,
) -> StoreResult<()> {
match work {
WorkRef::Wave(id) => tx.execute(
"INSERT INTO work_placements (wave_id, home_id, enabled, placed_at)
VALUES (?1, ?2, 1, ?3)
ON CONFLICT(wave_id) DO UPDATE SET
home_id=excluded.home_id, placed_at=excluded.placed_at",
params![id.as_str(), home_id.as_str(), placed_at],
)?,
WorkRef::Project(id) => tx.execute(
"INSERT INTO work_placements (project_id, home_id, enabled, placed_at)
VALUES (?1, ?2, 1, ?3)
ON CONFLICT(project_id) DO UPDATE SET
home_id=excluded.home_id, placed_at=excluded.placed_at",
params![id.as_str(), home_id.as_str(), placed_at],
)?,
WorkRef::Task(id) => tx.execute(
"INSERT INTO work_placements (task_id, home_id, enabled, placed_at)
VALUES (?1, ?2, 1, ?3)
ON CONFLICT(task_id) DO UPDATE SET
home_id=excluded.home_id, placed_at=excluded.placed_at",
params![id.as_str(), home_id.as_str(), placed_at],
)?,
};
Ok(())
}
fn has_runtime_generations(conn: &Connection) -> StoreResult<bool> {
conn.query_row(
"SELECT EXISTS (
SELECT 1 FROM sqlite_master
WHERE type='table' AND name='home_runtime_generations'
)",
[],
|row| row.get::<_, bool>(0),
)
.map_err(StoreError::from)
}
fn runtime_generation_in(conn: &Connection, home_id: &HomeId) -> StoreResult<Option<u64>> {
if !has_runtime_generations(conn)? {
return Ok(None);
}
let generation = conn.query_row(
"SELECT MAX(generation) FROM home_runtime_generations WHERE home_id=?1",
[home_id.as_str()],
|row| row.get::<_, Option<i64>>(0),
)?;
generation
.map(|value| {
u64::try_from(value).map_err(|_| {
StoreError::InvalidData(format!(
"Home {} has invalid runtime generation {value}",
home_id.as_str()
))
})
})
.transpose()
}
fn runtime_generation_sql(generation: Option<u64>) -> StoreResult<Option<i64>> {
generation
.map(|value| {
i64::try_from(value).map_err(|_| {
StoreError::InvalidData(format!("runtime generation {value} exceeds SQLite"))
})
})
.transpose()
}
fn validate_upgrade_reservation(
tx: &Transaction<'_>,
home_id: &HomeId,
runtime_generation: Option<u64>,
) -> StoreResult<()> {
let has_upgrades = tx.query_row(
"SELECT EXISTS (
SELECT 1 FROM sqlite_master
WHERE type='table' AND name='home_upgrades'
)",
[],
|row| row.get::<_, bool>(0),
)?;
if !has_upgrades {
return Ok(());
}
let active = tx
.query_row(
"SELECT id, target_generation, artifacts_activated
FROM home_upgrades
WHERE home_id=?1
AND phase NOT IN ('completed', 'failed', 'rolled_back')
ORDER BY target_generation DESC, started_at DESC, id DESC
LIMIT 1",
[home_id.as_str()],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, i64>(1)?,
row.get::<_, bool>(2)?,
))
},
)
.optional()?;
let Some((upgrade_id, target_generation, artifacts_activated)) = active else {
return Ok(());
};
let target_generation = u64::try_from(target_generation).map_err(|_| {
StoreError::InvalidData(format!(
"Home upgrade {upgrade_id} has invalid target generation {target_generation}"
))
})?;
if artifacts_activated && runtime_generation == Some(target_generation) {
return Ok(());
}
Err(StoreError::HomeUpgradeFenced {
upgrade_id,
runtime_generation,
})
}
pub(super) fn reserve_run_in(
tx: &Transaction<'_>,
work: &WorkRef,
trigger: &RunTrigger,
) -> StoreResult<(Run, RunLease)> {
let epoch = current_epoch_in(tx, work)?;
let home_id = reserving_home_in(tx, work)?;
if let Some(run) = run_for_epoch_in(tx, &epoch.id)? {
return Err(StoreError::RunFenced {
target: format!("{} {}", work.kind(), work.id()),
run_id: run.id,
state: run.state,
});
}
resolve_wait_for_trigger(tx, &epoch, trigger)?;
let token = RunLeaseToken::new();
let runtime_generation = runtime_generation_in(tx, &home_id)?;
validate_upgrade_reservation(tx, &home_id, runtime_generation)?;
let run = Run {
id: RunId::new(),
work: work.clone(),
epoch_id: epoch.id.clone(),
home_id,
runtime_generation,
state: RunState::Reserved,
trigger: trigger.clone(),
retry_of: match trigger {
RunTrigger::Recovery { prior_run_id } => Some(prior_run_id.clone()),
_ => None,
},
containment: None,
cwd: None,
created_at: OffsetDateTime::now_utc(),
started_at: None,
ended_at: None,
};
if run.runtime_generation.is_some() {
let runtime_generation = runtime_generation_sql(run.runtime_generation)?;
tx.execute(
"INSERT INTO runs (
id, epoch_id, home_id, runtime_generation, state, trigger_json, retry_of,
lease_hash, lease_generation, source_kind, source_id, created_at, ended_at,
stop_reason
) VALUES (?1, ?2, ?3, ?4, 'reserved', ?5, ?6, ?7, NULL, ?8, ?9, ?10, NULL, NULL)",
params![
run.id.as_str(),
run.epoch_id.as_str(),
run.home_id.as_str(),
runtime_generation,
serde_json::to_string(trigger).expect("Run trigger must serialize"),
run.retry_of.as_ref().map(RunId::as_str),
token.hash(),
work.kind(),
work.id(),
run.created_at.unix_timestamp(),
],
)?;
} else {
tx.execute(
"INSERT INTO runs (
id, epoch_id, home_id, state, trigger_json, retry_of, lease_hash,
lease_generation, source_kind, source_id, created_at, ended_at, stop_reason
) VALUES (?1, ?2, ?3, 'reserved', ?4, ?5, ?6, NULL, ?7, ?8, ?9, NULL, NULL)",
params![
run.id.as_str(),
run.epoch_id.as_str(),
run.home_id.as_str(),
serde_json::to_string(trigger).expect("Run trigger must serialize"),
run.retry_of.as_ref().map(RunId::as_str),
token.hash(),
work.kind(),
work.id(),
run.created_at.unix_timestamp(),
],
)?;
}
let lease = RunLease::new(run.id.clone(), work.clone(), epoch.current_basis, token);
Ok((run, lease))
}
fn inherit_placement(
tx: &Transaction<'_>,
work: &WorkRef,
parent: Option<&WorkRef>,
placed_at: i64,
) -> StoreResult<()> {
if find_placement_in(tx, work)?.is_some() {
return Ok(());
}
let home_id = match parent {
Some(parent) => placement_in(tx, parent)?.home_id,
None => map_local_home(tx)?.id,
};
write_placement(tx, work, &home_id, placed_at)
}
fn resolve_wait_for_trigger(
tx: &Transaction<'_>,
epoch: &Epoch,
trigger: &RunTrigger,
) -> StoreResult<()> {
let row = tx
.query_row(
"SELECT id, on_json FROM waits WHERE epoch_id=?1 AND resolved_at IS NULL",
[epoch.id.as_str()],
|row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
)
.optional()?;
let Some((wait_id, on_json)) = row else {
return Ok(());
};
let wait_on: WaitOn = serde_json::from_str(&on_json)?;
let resolves = match (&wait_on, trigger) {
(WaitOn::Input { after }, RunTrigger::Input { basis }) => {
basis.epoch_id == after.epoch_id && basis.revision > after.revision
}
(WaitOn::Input { .. }, RunTrigger::User) => true,
(WaitOn::Time { not_before }, RunTrigger::Time { scheduled_at }) => {
scheduled_at >= not_before
}
(WaitOn::Event { event }, RunTrigger::Event { event: observed }) => event == observed,
(WaitOn::Child { work }, RunTrigger::Child { work: observed }) => work == observed,
(WaitOn::Capability { .. }, RunTrigger::Recovery { .. })
| (WaitOn::Effect { .. }, RunTrigger::Recovery { .. }) => true,
_ => false,
};
if !resolves {
return Err(StoreError::InvalidData(format!(
"Run trigger does not resolve Wait {wait_id}"
)));
}
tx.execute(
"UPDATE waits SET resolved_at=?2 WHERE id=?1 AND resolved_at IS NULL",
params![wait_id, now_unix()],
)?;
Ok(())
}
pub(crate) fn validate_run_lease(conn: &Connection, lease: &RunLease) -> StoreResult<Run> {
let run = run_by_id_in(conn, &lease.run_id)?;
if run.work != lease.work || !matches!(run.state, RunState::Reserved | RunState::Active) {
return Err(StoreError::InvalidAuthority(format!(
"Run {} no longer holds execution authority",
lease.run_id
)));
}
let stored_hash: String = conn.query_row(
"SELECT lease_hash FROM runs WHERE id=?1",
[lease.run_id.as_str()],
|row| row.get(0),
)?;
if stored_hash != lease.token_hash() {
return Err(StoreError::InvalidAuthority(format!(
"Run {} lease token does not match",
lease.run_id
)));
}
Ok(run)
}
pub(crate) fn validate_stop_lease(conn: &Connection, lease: &RunLease) -> StoreResult<Run> {
let run = run_by_id_in(conn, &lease.run_id)?;
if run.work != lease.work
|| !matches!(
run.state,
RunState::Reserved | RunState::Active | RunState::Stopping
)
{
return Err(StoreError::InvalidAuthority(format!(
"Run {} no longer owns cleanup authority",
lease.run_id
)));
}
let stored_hash: String = conn.query_row(
"SELECT lease_hash FROM runs WHERE id=?1",
[lease.run_id.as_str()],
|row| row.get(0),
)?;
if stored_hash != lease.token_hash() {
return Err(StoreError::InvalidAuthority(format!(
"Run {} lease token does not match",
lease.run_id
)));
}
Ok(run)
}
fn run_for_epoch_in(conn: &Connection, epoch_id: &EpochId) -> StoreResult<Option<Run>> {
let id = conn
.query_row(
"SELECT id FROM runs WHERE epoch_id=?1 AND state != 'ended'",
[epoch_id.as_str()],
|row| row.get::<_, String>(0),
)
.optional()?;
id.map(|id| {
RunId::parse(&id)
.map_err(invalid_durable)
.and_then(|id| run_by_id_in(conn, &id))
})
.transpose()
}
fn latest_run_for_epoch_in(conn: &Connection, epoch_id: &EpochId) -> StoreResult<Option<Run>> {
let id = conn
.query_row(
"SELECT id FROM runs WHERE epoch_id=?1 ORDER BY created_at DESC, rowid DESC LIMIT 1",
[epoch_id.as_str()],
|row| row.get::<_, String>(0),
)
.optional()?;
id.map(|id| {
RunId::parse(&id)
.map_err(invalid_durable)
.and_then(|id| run_by_id_in(conn, &id))
})
.transpose()
}
fn current_run_for_work_in(conn: &Connection, work: &WorkRef) -> StoreResult<Option<Run>> {
let epoch = current_epoch_in(conn, work)?;
run_for_epoch_in(conn, &epoch.id)
}
fn run_by_id_in(conn: &Connection, run_id: &RunId) -> StoreResult<Run> {
let generation_column = if has_runtime_generations(conn)? {
"r.runtime_generation"
} else {
"NULL"
};
let query = format!(
"SELECT r.epoch_id, r.home_id, r.state, r.trigger_json, r.retry_of,
r.created_at, r.ended_at, e.wave_id, e.project_id, e.task_id,
r.containment_kind, r.containment_id, r.cwd, r.started_at,
{generation_column}
FROM runs r JOIN epochs e ON e.id=r.epoch_id WHERE r.id=?1"
);
let row = conn.query_row(&query, [run_id.as_str()], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
row.get::<_, String>(3)?,
row.get::<_, Option<String>>(4)?,
row.get::<_, i64>(5)?,
row.get::<_, Option<i64>>(6)?,
row.get::<_, Option<String>>(7)?,
row.get::<_, Option<String>>(8)?,
row.get::<_, Option<String>>(9)?,
row.get::<_, Option<String>>(10)?,
row.get::<_, Option<String>>(11)?,
row.get::<_, Option<String>>(12)?,
row.get::<_, Option<i64>>(13)?,
row.get::<_, Option<i64>>(14)?,
))
})?;
let work = work_from_parts((row.7, row.8, row.9))?;
Ok(Run {
id: run_id.clone(),
work,
epoch_id: EpochId::parse(&row.0).map_err(invalid_durable)?,
home_id: HomeId::parse(&row.1).map_err(invalid_durable)?,
runtime_generation: row
.14
.map(u64::try_from)
.transpose()
.map_err(|_| StoreError::InvalidData("Run has a negative runtime generation".into()))?,
state: RunState::parse(&row.2).map_err(invalid_durable)?,
trigger: serde_json::from_str(&row.3)?,
retry_of: row
.4
.map(|id| RunId::parse(&id).map_err(invalid_durable))
.transpose()?,
containment: match (row.10, row.11) {
(Some(kind), Some(id)) => Some(Containment::parse(&kind, id).map_err(invalid_durable)?),
(None, None) => None,
_ => {
return Err(StoreError::InvalidData(
"stored Run containment is incomplete".to_string(),
))
}
},
cwd: row.12.map(Into::into),
created_at: OffsetDateTime::from_unix_timestamp(row.5).map_err(invalid_durable)?,
started_at: row
.13
.map(OffsetDateTime::from_unix_timestamp)
.transpose()
.map_err(invalid_durable)?,
ended_at: row
.6
.map(OffsetDateTime::from_unix_timestamp)
.transpose()
.map_err(invalid_durable)?,
})
}
fn insert_supervised_invocation(
tx: &Transaction<'_>,
run: &Run,
invocation: &AgentInvocation,
) -> StoreResult<()> {
if let Some(ask_id) = invocation.answer_ask_id.as_ref() {
let ask = ask_by_id_in(tx, ask_id)?;
if ask.target != AskTarget::Parent(run.work.clone()) {
return Err(StoreError::InvalidAuthority(format!(
"Run {} does not own Ask {ask_id}'s answer route",
run.id
)));
}
if ask.state.is_terminal() || !ask_epoch_is_open_in(tx, ask_id)? {
return Err(StoreError::InvalidAuthority(format!(
"Ask {ask_id} is no longer answerable"
)));
}
}
let labels = work_labels(tx, &run.work)?;
let cwd = run
.cwd
.as_ref()
.expect("an active Run has cwd")
.display()
.to_string();
let (_, containment_id) = run
.containment
.as_ref()
.expect("an active Run has containment")
.parts();
tx.execute(
"INSERT INTO agent_invocations (
id, run_id, process_id, started_at, ended_at, repo, worktree, wave,
flow, skill, project, task, provider, model, surface, capture_status,
incomplete_reason, outcome, artifact_dir, conversation_path,
provider_events_path, provider_session_id, provider_session_path,
conversation_event_count, conversation_bytes, supervising_run_id,
account_id, resume_token, answer_ask_id
) VALUES (
?1, ?2, ?3, ?4, NULL, ?5, ?6, ?7, NULL, NULL, ?8, ?9, ?10, ?11,
?12, 'prompt_only', NULL, 'running', '', '', NULL, NULL, NULL, 0, 0,
?2, ?13, ?14, ?15
)",
params![
invocation.id.as_str(),
run.id.as_str(),
containment_id,
invocation.started_at.unix_timestamp(),
labels.repo,
cwd,
labels.wave,
labels.project,
labels.task,
invocation.route.provider,
invocation.route.model,
invocation.surface,
invocation.route.account_id,
invocation.resume_token,
invocation.answer_ask_id.as_ref().map(AskId::as_str),
],
)?;
Ok(())
}
fn insert_ask_invocation(
tx: &Transaction<'_>,
run: &Run,
ask: &Ask,
invocation: &AgentInvocation,
containment_id: &str,
) -> StoreResult<()> {
let labels = work_labels(tx, &ask.origin.work)?;
tx.execute(
"INSERT INTO agent_invocations (
id, run_id, process_id, started_at, ended_at, repo, worktree, wave,
flow, skill, project, task, provider, model, surface, capture_status,
incomplete_reason, outcome, artifact_dir, conversation_path,
provider_events_path, provider_session_id, provider_session_path,
conversation_event_count, conversation_bytes, supervising_run_id,
account_id, resume_token, answer_ask_id, ask_ready_at,
ask_presented_at
) VALUES (
?1, ?2, ?3, ?4, NULL, ?5, ?6, ?7, NULL, 'ask', ?8, ?9, ?10, ?11,
?12, 'prompt_only', NULL, 'running', '', '', NULL, NULL, NULL, 0, 0,
?2, ?13, NULL, ?14, NULL, NULL
)",
params![
invocation.id.as_str(),
run.id.as_str(),
containment_id,
invocation.started_at.unix_timestamp(),
labels.repo,
ask.origin.cwd.display().to_string(),
labels.wave,
labels.project,
labels.task,
invocation.route.provider,
invocation.route.model,
invocation.surface,
invocation.route.account_id,
ask.id.as_str(),
],
)?;
Ok(())
}
fn ask_containment_id(invocation_id: &AgentInvocationId) -> String {
format!("lf-ask-{}", &invocation_id.as_str()[11..23])
}
struct WorkLabels {
wave: Option<String>,
project: Option<String>,
task: Option<String>,
repo: String,
}
fn work_labels(conn: &Connection, work: &WorkRef) -> StoreResult<WorkLabels> {
match work {
WorkRef::Wave(id) => conn
.query_row(
"SELECT name, repo FROM waves WHERE id=?1",
[id.as_str()],
|row| {
Ok(WorkLabels {
wave: Some(row.get(0)?),
project: None,
task: None,
repo: row.get(1)?,
})
},
)
.map_err(StoreError::from),
WorkRef::Project(id) => conn
.query_row(
"SELECT w.name, p.external_project_id, w.repo
FROM projects p JOIN waves w ON w.id=p.wave_id WHERE p.id=?1",
[id.as_str()],
|row| {
Ok(WorkLabels {
wave: Some(row.get(0)?),
project: Some(row.get(1)?),
task: None,
repo: row.get(2)?,
})
},
)
.map_err(StoreError::from),
WorkRef::Task(id) => conn
.query_row(
"SELECT w.name, p.external_project_id, t.issue_identifier, w.repo
FROM tasks t
JOIN projects p ON p.id=t.project_id
JOIN waves w ON w.id=p.wave_id
WHERE t.id=?1",
[id.as_str()],
|row| {
Ok(WorkLabels {
wave: Some(row.get(0)?),
project: Some(row.get(1)?),
task: Some(row.get(2)?),
repo: row.get(3)?,
})
},
)
.map_err(StoreError::from),
}
}
fn require_invocation_for_run(
conn: &Connection,
invocation_id: &AgentInvocationId,
run_id: &RunId,
) -> StoreResult<()> {
conn.query_row(
"SELECT 1 FROM agent_invocations WHERE id=?1 AND supervising_run_id=?2",
params![invocation_id.as_str(), run_id.as_str()],
|_| Ok(()),
)
.map_err(StoreError::from)
}
fn require_open_invocation_for_run(
conn: &Connection,
invocation_id: &AgentInvocationId,
run_id: &RunId,
) -> StoreResult<()> {
conn.query_row(
"SELECT 1 FROM agent_invocations
WHERE id=?1 AND supervising_run_id=?2 AND ended_at IS NULL",
params![invocation_id.as_str(), run_id.as_str()],
|_| Ok(()),
)
.map_err(StoreError::from)
}
fn open_invocation_for_run_in(
conn: &Connection,
run_id: &RunId,
) -> StoreResult<Option<AgentInvocation>> {
let invocation_id = conn
.query_row(
"SELECT id FROM agent_invocations
WHERE supervising_run_id=?1 AND ended_at IS NULL
ORDER BY started_at, rowid LIMIT 1",
[run_id.as_str()],
|row| row.get::<_, String>(0),
)
.optional()?;
invocation_id
.map(|id| {
let id = AgentInvocationId::parse(&id).map_err(invalid_durable)?;
supervised_invocation_in(conn, &id)
})
.transpose()
}
fn supervised_invocation_in(
conn: &Connection,
invocation_id: &AgentInvocationId,
) -> StoreResult<AgentInvocation> {
let row = conn.query_row(
"SELECT supervising_run_id, provider, model, account_id, surface,
resume_token, started_at, ended_at, answer_ask_id
FROM agent_invocations WHERE id=?1 AND supervising_run_id IS NOT NULL",
[invocation_id.as_str()],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, Option<String>>(2)?,
row.get::<_, Option<String>>(3)?,
row.get::<_, String>(4)?,
row.get::<_, Option<String>>(5)?,
row.get::<_, i64>(6)?,
row.get::<_, Option<i64>>(7)?,
row.get::<_, Option<String>>(8)?,
))
},
)?;
Ok(AgentInvocation {
id: invocation_id.clone(),
supervising_run_id: Some(RunId::parse(&row.0).map_err(invalid_durable)?),
answer_ask_id: row
.8
.map(|id| AskId::parse(&id).map_err(invalid_durable))
.transpose()?,
route: InvocationRoute {
provider: row.1,
model: row.2,
account_id: row.3,
},
surface: row.4,
resume_token: row.5,
started_at: OffsetDateTime::from_unix_timestamp(row.6).map_err(invalid_durable)?,
ended_at: row
.7
.map(OffsetDateTime::from_unix_timestamp)
.transpose()
.map_err(invalid_durable)?,
})
}
fn invocation_surface_in(
conn: &Connection,
invocation_id: &AgentInvocationId,
) -> StoreResult<Option<InvocationSurface>> {
let row = conn
.query_row(
"SELECT r.id, e.wave_id, e.project_id, e.task_id, h.route,
l.handback_state, l.process_id, l.answer_ask_id, l.surface
FROM agent_invocations l
JOIN runs r ON r.id=l.supervising_run_id
JOIN epochs e ON e.id=r.epoch_id
JOIN homes h ON h.id=r.home_id
WHERE l.id=?1 AND l.supervising_run_id IS NOT NULL",
[invocation_id.as_str()],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, Option<String>>(1)?,
row.get::<_, Option<String>>(2)?,
row.get::<_, Option<String>>(3)?,
row.get::<_, String>(4)?,
row.get::<_, Option<String>>(5)?,
row.get::<_, String>(6)?,
row.get::<_, Option<String>>(7)?,
row.get::<_, String>(8)?,
))
},
)
.optional()?;
let Some(row) = row else {
return Ok(None);
};
let work = work_from_parts((row.1, row.2, row.3))?;
let handback = row
.5
.as_deref()
.map(BoundaryState::parse_handback)
.transpose()
.map_err(invalid_durable)?;
let invocation = supervised_invocation_in(conn, invocation_id)?;
let run_id = RunId::parse(&row.0).map_err(invalid_durable)?;
let run = run_by_id_in(conn, &run_id)?;
let ask_tmux = row.7.is_some() && row.8.starts_with("ask_");
let attach_argv = if ask_tmux {
Some(vec![
"tmux".to_string(),
"attach-session".to_string(),
"-t".to_string(),
row.6.clone(),
])
} else {
match &run.containment {
Some(Containment::Tmux { name }) if row.7.is_none() => Some(vec![
"tmux".to_string(),
"attach-session".to_string(),
"-t".to_string(),
name.clone(),
]),
Some(Containment::Tmux { .. }) | Some(Containment::ProcessGroup { .. }) | None => None,
}
};
debug_assert_eq!(invocation.supervising_run_id.as_ref(), Some(&run.id));
let wave_id = match &work {
WorkRef::Wave(id) => id.clone(),
WorkRef::Project(id) => {
let value: String = conn.query_row(
"SELECT wave_id FROM projects WHERE id=?1",
[id.as_str()],
|row| row.get(0),
)?;
WaveId::parse(&value).map_err(invalid_durable)?
}
WorkRef::Task(id) => {
let value: String = conn.query_row(
"SELECT p.wave_id FROM tasks t JOIN projects p ON p.id=t.project_id
WHERE t.id=?1",
[id.as_str()],
|row| row.get(0),
)?;
WaveId::parse(&value).map_err(invalid_durable)?
}
};
Ok(Some(InvocationSurface {
invocation,
run,
work,
wave_id,
home_route: row.4,
handback,
attach_argv,
}))
}
fn require_turn_for_run(conn: &Connection, turn_id: &TurnId, run_id: &RunId) -> StoreResult<()> {
conn.query_row(
"SELECT 1 FROM agent_turns t
JOIN agent_invocations l ON l.id=t.invocation_id
WHERE t.id=?1 AND l.supervising_run_id=?2",
params![turn_id.as_str(), run_id.as_str()],
|_| Ok(()),
)
.map_err(StoreError::from)
}
fn control_turn_in(conn: &Connection, turn_id: &TurnId) -> StoreResult<Turn> {
let row = conn.query_row(
"SELECT invocation_id, epoch_id, basis_rev, status, provider_turn_id,
root_output, started_at, ended_at FROM agent_turns WHERE id=?1",
[turn_id.as_str()],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, i64>(2)?,
row.get::<_, String>(3)?,
row.get::<_, Option<String>>(4)?,
row.get::<_, Option<String>>(5)?,
row.get::<_, i64>(6)?,
row.get::<_, Option<i64>>(7)?,
))
},
)?;
Ok(Turn {
id: turn_id.clone(),
invocation_id: AgentInvocationId::parse(&row.0).map_err(invalid_durable)?,
basis: Basis {
epoch_id: EpochId::parse(&row.1).map_err(invalid_durable)?,
revision: row.2 as u64,
},
state: BoundaryState::parse_turn(&row.3).map_err(invalid_durable)?,
provider_turn_id: row.4,
root_output: row.5,
started_at: OffsetDateTime::from_unix_timestamp(row.6).map_err(invalid_durable)?,
ended_at: row
.7
.map(OffsetDateTime::from_unix_timestamp)
.transpose()
.map_err(invalid_durable)?,
})
}
fn handback_state(state: BoundaryState) -> &'static str {
match state {
BoundaryState::Succeeded => "succeeded",
BoundaryState::Failed => "failed",
BoundaryState::Interrupted => "interrupted",
BoundaryState::Unknown => "unknown",
BoundaryState::Starting | BoundaryState::Active => {
unreachable!("terminal outcome validated before handback mapping")
}
}
}
fn parent_work(conn: &Connection, work: &WorkRef) -> StoreResult<Option<WorkRef>> {
match work {
WorkRef::Wave(_) => Ok(None),
WorkRef::Project(id) => {
let wave_id: String = conn.query_row(
"SELECT wave_id FROM projects WHERE id=?1",
[id.as_str()],
|row| row.get(0),
)?;
Ok(Some(WorkRef::Wave(
WaveId::parse(&wave_id).map_err(invalid_durable)?,
)))
}
WorkRef::Task(id) => {
let project_id: String = conn.query_row(
"SELECT project_id FROM tasks WHERE id=?1",
[id.as_str()],
|row| row.get(0),
)?;
Ok(Some(WorkRef::Project(
ProjectId::parse(&project_id).map_err(invalid_durable)?,
)))
}
}
}
fn current_turn_for_invocation_in(
conn: &Connection,
invocation_id: &AgentInvocationId,
) -> StoreResult<TurnId> {
let id = conn
.query_row(
"SELECT id FROM agent_turns
WHERE invocation_id=?1 AND status='running'
ORDER BY ordinal DESC LIMIT 1",
[invocation_id.as_str()],
|row| row.get::<_, String>(0),
)
.optional()?
.ok_or_else(|| {
StoreError::InvalidAuthority(format!(
"AgentInvocation {invocation_id} has no active Turn"
))
})?;
TurnId::parse(&id).map_err(invalid_durable)
}
fn insert_ask(conn: &Connection, epoch_id: &EpochId, ask: &Ask) -> StoreResult<()> {
let (target_kind, target_work_kind, target_work_id) = match &ask.target {
AskTarget::User => ("user", None, None),
AskTarget::Parent(work) => ("parent", Some(work.kind()), Some(work.id())),
};
let (request_kind, prompt, flow, node_id, skill, iteration) = match &ask.request {
AskBody::Intervention { prompt } => (
"intervention",
Some(prompt.as_str()),
None,
None,
None,
None,
),
AskBody::FlowStep {
flow,
node_id,
skill,
iteration,
} => (
"flow_step",
None,
Some(flow.as_str()),
Some(node_id.as_str()),
Some(skill.as_str()),
Some(i64::from(*iteration)),
),
};
let (result_kind, result_text) = ask
.result
.as_ref()
.map(|result| (Some(result.state().as_str()), Some(result.text())))
.unwrap_or((None, None));
let (author_kind, author_id) = ask
.terminal_author
.as_ref()
.map(author_parts)
.unwrap_or((None, None));
conn.execute(
"INSERT INTO ask_exchanges (
id, epoch_id, origin_work_kind, origin_work_id, origin_run_id,
origin_turn_id, origin_invocation_id, origin_home_id, origin_cwd,
target_kind, target_work_kind, target_work_id,
request_kind, request_prompt, request_flow, request_node_id,
request_skill, request_iteration, state, active_invocation_id,
result_kind, result_text, terminal_author_kind, terminal_author_id,
asked_at, terminal_at
) VALUES (
?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13,
?14, ?15, ?16, ?17, ?18, ?19, ?20, ?21, ?22, ?23, ?24, ?25, ?26
)",
params![
ask.id.as_str(),
epoch_id.as_str(),
ask.origin.work.kind(),
ask.origin.work.id(),
ask.origin.run_id.as_str(),
ask.origin.turn_id.as_ref().map(TurnId::as_str),
ask.origin
.invocation_id
.as_ref()
.map(AgentInvocationId::as_str),
ask.origin.home_id.as_str(),
ask.origin.cwd.display().to_string(),
target_kind,
target_work_kind,
target_work_id,
request_kind,
prompt,
flow,
node_id,
skill,
iteration,
ask.state.as_str(),
ask.active_invocation_id
.as_ref()
.map(AgentInvocationId::as_str),
result_kind,
result_text,
author_kind,
author_id,
ask.asked_at.unix_timestamp(),
ask.terminal_at.map(OffsetDateTime::unix_timestamp),
],
)?;
Ok(())
}
fn enqueue_ask_comment(conn: &Connection, epoch_id: &EpochId, ask: &Ask) -> StoreResult<()> {
let route = match &ask.target {
AskTarget::User => "User".to_string(),
AskTarget::Parent(work) => format!("{} `{}`", work.kind(), work.id()),
};
let request = match &ask.request {
AskBody::Intervention { prompt } => prompt.clone(),
AskBody::FlowStep {
flow,
node_id,
skill,
..
} => format!("Run `{flow}` node `{node_id}` with `{skill}`."),
};
let transition = AskCommentTransition::Requested;
let body = format!(
"### Loopflow Ask\n\n**Route:** {route}\n\n{}\n\n{}",
request,
transition.marker(&ask.id)
);
enqueue_ask_comment_write(
conn,
&ask.id,
epoch_id,
&ask.origin.work,
transition,
&body,
ask.asked_at.unix_timestamp(),
)
}
fn enqueue_ask_result_comment(conn: &Connection, ask: &Ask) -> StoreResult<()> {
let result = ask
.result
.as_ref()
.ok_or_else(|| StoreError::InvalidData(format!("terminal Ask {} has no result", ask.id)))?;
let author = match ask.terminal_author.as_ref() {
Some(Author::User) => "User".to_string(),
Some(Author::Run(run_id)) => format!("Run `{run_id}`"),
None => "Loopflow".to_string(),
};
let heading = match result {
AskResult::Resolved { .. } => "Loopflow Ask Resolved",
AskResult::Declined { .. } => "Loopflow Ask Declined",
AskResult::Cancelled { .. } => "Loopflow Ask Cancelled",
};
let transition = AskCommentTransition::Result;
let body = format!(
"### {heading}\n\n**Author:** {author}\n\n{}\n\n{}",
result.text(),
transition.marker(&ask.id)
);
enqueue_ask_comment_write(
conn,
&ask.id,
&ask_epoch_id_in(conn, &ask.id)?,
&ask.origin.work,
transition,
&body,
ask.terminal_at
.expect("a terminal Ask has terminal_at")
.unix_timestamp(),
)
}
fn enqueue_ask_comment_write(
conn: &Connection,
ask_id: &AskId,
epoch_id: &EpochId,
work: &WorkRef,
transition: AskCommentTransition,
body: &str,
created_at: i64,
) -> StoreResult<()> {
let WorkRef::Task(task_id) = work else {
return Ok(());
};
let current_epoch = conn.query_row(
"SELECT id FROM epochs WHERE id=?1 AND task_id=?2",
params![epoch_id.as_str(), task_id.as_str()],
|row| row.get::<_, String>(0),
)?;
if current_epoch != epoch_id.as_str() {
return Err(StoreError::InvalidData(
"Ask Task attribution does not match its Epoch".to_string(),
));
}
let issue_id: String = conn.query_row(
"SELECT external_issue_id FROM tasks WHERE id=?1",
[task_id.as_str()],
|row| row.get(0),
)?;
conn.execute(
"INSERT OR IGNORE INTO ask_linear_comment_outbox (
ask_id, transition, task_id, issue_id, body, created_at,
attempt_count, attempt_started_at, last_error,
linear_comment_id, delivered_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, 0, NULL, NULL, NULL, NULL)",
params![
ask_id.as_str(),
transition.as_str(),
task_id.as_str(),
issue_id,
body,
created_at,
],
)?;
Ok(())
}
type AskCommentWriteRow = (
String,
String,
String,
String,
String,
String,
i64,
Option<i64>,
Option<String>,
Option<String>,
Option<i64>,
);
fn read_ask_comment_write_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<AskCommentWriteRow> {
Ok((
row.get(0)?,
row.get(1)?,
row.get(2)?,
row.get(3)?,
row.get(4)?,
row.get(5)?,
row.get(6)?,
row.get(7)?,
row.get(8)?,
row.get(9)?,
row.get(10)?,
))
}
fn parse_ask_comment_write_row(row: AskCommentWriteRow) -> StoreResult<AskCommentWrite> {
let transition = match row.1.as_str() {
"ask" => AskCommentTransition::Requested,
"answer" => AskCommentTransition::Result,
value => {
return Err(StoreError::InvalidData(format!(
"invalid Ask comment transition {value:?}"
)))
}
};
Ok(AskCommentWrite {
ask_id: AskId::parse(&row.0).map_err(invalid_durable)?,
transition,
issue_id: row.2,
body: row.3,
repo: row.4,
wave: row.5,
attempt_count: u32::try_from(row.6).map_err(|_| {
StoreError::InvalidData("invalid Ask comment attempt count".to_string())
})?,
attempt_started_at: row.7,
last_error: row.8,
linear_comment_id: row.9,
delivered_at: row.10,
})
}
fn ask_comment_write_in(
conn: &Connection,
ask_id: &AskId,
transition: AskCommentTransition,
) -> StoreResult<AskCommentWrite> {
let row = conn.query_row(
"SELECT o.ask_id, o.transition, o.issue_id, o.body,
w.repo, w.name, o.attempt_count, o.attempt_started_at,
o.last_error, o.linear_comment_id, o.delivered_at
FROM ask_linear_comment_outbox o
JOIN tasks t ON t.id=o.task_id
JOIN projects p ON p.id=t.project_id
JOIN waves w ON w.id=p.wave_id
WHERE o.ask_id=?1 AND o.transition=?2",
params![ask_id.as_str(), transition.as_str()],
read_ask_comment_write_row,
)?;
parse_ask_comment_write_row(row)
}
pub(super) fn ask_by_id_in(conn: &Connection, ask_id: &AskId) -> StoreResult<Ask> {
conn.query_row(
"SELECT origin_work_kind, origin_work_id, origin_run_id, origin_turn_id,
origin_invocation_id, origin_home_id, origin_cwd,
target_kind, target_work_kind, target_work_id,
request_kind, request_prompt, request_flow, request_node_id,
request_skill, request_iteration, state, active_invocation_id,
result_kind, result_text, terminal_author_kind, terminal_author_id,
asked_at, terminal_at
FROM ask_exchanges WHERE id=?1",
[ask_id.as_str()],
|row| {
map_ask_row(
ask_id.clone(),
row.get(0)?,
row.get(1)?,
row.get(2)?,
row.get(3)?,
row.get(4)?,
row.get(5)?,
row.get(6)?,
row.get(7)?,
row.get(8)?,
row.get(9)?,
row.get(10)?,
row.get(11)?,
row.get(12)?,
row.get(13)?,
row.get(14)?,
row.get(15)?,
row.get(16)?,
row.get(17)?,
row.get(18)?,
row.get(19)?,
row.get(20)?,
row.get(21)?,
row.get(22)?,
row.get(23)?,
)
.map_err(to_sqlite_conversion_error)
},
)
.map_err(StoreError::from)
}
#[allow(clippy::too_many_arguments)]
fn map_ask_row(
id: AskId,
origin_work_kind: String,
origin_work_id: String,
origin_run_id: String,
origin_turn_id: Option<String>,
origin_invocation_id: Option<String>,
origin_home_id: String,
origin_cwd: String,
target_kind: String,
target_work_kind: Option<String>,
target_work_id: Option<String>,
request_kind: String,
request_prompt: Option<String>,
request_flow: Option<String>,
request_node_id: Option<String>,
request_skill: Option<String>,
request_iteration: Option<i64>,
state: String,
active_invocation_id: Option<String>,
result_kind: Option<String>,
result_text: Option<String>,
terminal_author_kind: Option<String>,
terminal_author_id: Option<String>,
asked_at: i64,
terminal_at: Option<i64>,
) -> StoreResult<Ask> {
let target = match (target_kind.as_str(), target_work_kind, target_work_id) {
("user", None, None) => AskTarget::User,
("parent", Some(kind), Some(id)) => AskTarget::Parent(parse_work_ref(&kind, &id)?),
_ => {
return Err(StoreError::InvalidData(
"stored Ask route is inconsistent".to_string(),
))
}
};
let request = match (
request_kind.as_str(),
request_prompt,
request_flow,
request_node_id,
request_skill,
request_iteration,
) {
("intervention", Some(prompt), None, None, None, None) => AskBody::Intervention { prompt },
("flow_step", None, Some(flow), Some(node_id), Some(skill), Some(iteration)) => {
AskBody::FlowStep {
flow,
node_id,
skill,
iteration: u32::try_from(iteration).map_err(|_| {
StoreError::InvalidData(
"stored flow-step iteration is outside the u32 range".to_string(),
)
})?,
}
}
_ => {
return Err(StoreError::InvalidData(
"stored Ask body is inconsistent".to_string(),
))
}
};
let result = match (result_kind.as_deref(), result_text) {
(None, None) => None,
(Some("resolved"), Some(summary)) => Some(AskResult::Resolved { summary }),
(Some("declined"), Some(reason)) => Some(AskResult::Declined { reason }),
(Some("cancelled"), Some(reason)) => Some(AskResult::Cancelled { reason }),
_ => {
return Err(StoreError::InvalidData(
"stored Ask result is inconsistent".to_string(),
))
}
};
let terminal_author = match (terminal_author_kind.as_deref(), terminal_author_id) {
(None, None) => None,
(Some("user"), None) => Some(Author::User),
(Some("run"), Some(run_id)) => {
Some(Author::Run(RunId::parse(&run_id).map_err(invalid_durable)?))
}
_ => {
return Err(StoreError::InvalidData(
"stored Ask terminal author is inconsistent".to_string(),
))
}
};
Ok(Ask {
id,
origin: AskOrigin {
work: parse_work_ref(&origin_work_kind, &origin_work_id)?,
run_id: RunId::parse(&origin_run_id).map_err(invalid_durable)?,
turn_id: origin_turn_id
.map(|turn_id| TurnId::parse(&turn_id).map_err(invalid_durable))
.transpose()?,
invocation_id: origin_invocation_id
.map(|invocation_id| {
AgentInvocationId::parse(&invocation_id).map_err(invalid_durable)
})
.transpose()?,
home_id: HomeId::parse(&origin_home_id).map_err(invalid_durable)?,
cwd: origin_cwd.into(),
},
target,
request,
state: AskState::parse(&state).map_err(invalid_durable)?,
active_invocation_id: active_invocation_id
.map(|invocation_id| AgentInvocationId::parse(&invocation_id).map_err(invalid_durable))
.transpose()?,
result,
terminal_author,
asked_at: OffsetDateTime::from_unix_timestamp(asked_at).map_err(invalid_durable)?,
terminal_at: terminal_at
.map(OffsetDateTime::from_unix_timestamp)
.transpose()
.map_err(invalid_durable)?,
})
}
fn validate_target_caller(
conn: &Connection,
caller: Option<&RunLease>,
target: &AskTarget,
) -> StoreResult<()> {
match (target, caller) {
(AskTarget::User, None) => Ok(()),
(AskTarget::Parent(parent), Some(lease)) => {
let run = validate_run_lease(conn, lease)?;
if &run.work == parent {
Ok(())
} else {
Err(StoreError::InvalidAuthority(
"Run does not own this Ask target".to_string(),
))
}
}
_ => Err(StoreError::InvalidAuthority(
"caller does not own this Ask target".to_string(),
)),
}
}
fn ask_epoch_is_open_in(conn: &Connection, ask_id: &AskId) -> StoreResult<bool> {
conn.query_row(
"SELECT EXISTS(
SELECT 1 FROM ask_exchanges a
JOIN epochs e ON e.id=a.epoch_id
WHERE a.id=?1 AND e.state='open'
)",
[ask_id.as_str()],
|row| row.get(0),
)
.map_err(StoreError::from)
}
enum AskScope<'a> {
Target(&'a AskTarget),
OriginEpoch {
work: &'a WorkRef,
epoch_id: &'a EpochId,
},
}
fn query_asks(conn: &Connection, scope: AskScope<'_>) -> StoreResult<Vec<Ask>> {
let ids = match scope {
AskScope::Target(target) => {
let (predicate, kind, id) = match target {
AskTarget::User => ("a.target_kind='user'", None, None),
AskTarget::Parent(work) => (
"a.target_kind='parent' AND a.target_work_kind=?1 AND a.target_work_id=?2",
Some(work.kind()),
Some(work.id()),
),
};
let sql = format!(
"SELECT a.id FROM ask_exchanges a
JOIN epochs e ON e.id=a.epoch_id
WHERE a.state IN ('queued', 'claimed') AND e.state='open'
AND {predicate}
ORDER BY a.asked_at, a.rowid"
);
let mut statement = conn.prepare(&sql)?;
match (kind, id) {
(Some(kind), Some(id)) => statement
.query_map(params![kind, id], |row| row.get::<_, String>(0))?
.collect::<Result<Vec<_>, _>>()?,
(None, None) => statement
.query_map([], |row| row.get::<_, String>(0))?
.collect::<Result<Vec<_>, _>>()?,
_ => unreachable!("Ask target parts are complete"),
}
}
AskScope::OriginEpoch { work, epoch_id } => {
let mut statement = conn.prepare(
"SELECT id FROM ask_exchanges
WHERE epoch_id=?1 AND origin_work_kind=?2 AND origin_work_id=?3
ORDER BY asked_at DESC, rowid DESC",
)?;
let ids = statement
.query_map(params![epoch_id.as_str(), work.kind(), work.id()], |row| {
row.get::<_, String>(0)
})?
.collect::<Result<Vec<_>, _>>()?;
ids
}
};
ids.into_iter()
.map(|id| {
let id = AskId::parse(&id).map_err(invalid_durable)?;
ask_by_id_in(conn, &id)
})
.collect()
}
fn validate_ask_origin(
conn: &Connection,
lease: &RunLease,
origin: &AskOrigin,
) -> StoreResult<Run> {
let run = validate_run_lease(conn, lease)?;
if origin.work != run.work
|| origin.run_id != run.id
|| origin.home_id != run.home_id
|| run.cwd.as_ref() != Some(&origin.cwd)
{
return Err(StoreError::InvalidAuthority(
"Ask origin does not match the active Run lease".to_string(),
));
}
match (&origin.turn_id, &origin.invocation_id) {
(None, None) => {}
(Some(turn_id), Some(invocation_id)) => {
require_open_invocation_for_run(conn, invocation_id, &run.id)?;
let valid: bool = conn.query_row(
"SELECT EXISTS(
SELECT 1 FROM agent_turns
WHERE id=?1 AND invocation_id=?2 AND epoch_id=?3 AND status='running'
)",
params![
turn_id.as_str(),
invocation_id.as_str(),
run.epoch_id.as_str()
],
|row| row.get(0),
)?;
if !valid {
return Err(StoreError::InvalidAuthority(
"Ask Turn does not belong to the active origin Invocation".to_string(),
));
}
}
_ => {
return Err(StoreError::InvalidData(
"Ask origin Turn and Invocation must be both present or both absent".to_string(),
))
}
}
Ok(run)
}
fn validate_requested_target(
conn: &Connection,
origin: &WorkRef,
target: &AskTarget,
) -> StoreResult<()> {
if let AskTarget::Parent(parent) = target {
if parent_work(conn, origin)?.as_ref() != Some(parent) {
return Err(StoreError::InvalidAuthority(
"Ask parent target is not the origin Work's immediate parent".to_string(),
));
}
}
Ok(())
}
fn validate_ask_body(request: &AskBody) -> StoreResult<()> {
let empty = |value: &str| value.trim().is_empty();
match request {
AskBody::Intervention { prompt } if empty(prompt) => Err(StoreError::InvalidData(
"Ask intervention prompt cannot be empty".to_string(),
)),
AskBody::FlowStep {
flow,
node_id,
skill,
..
} if [flow, node_id, skill].into_iter().any(|value| empty(value)) => Err(
StoreError::InvalidData("Ask flow, node, and skill cannot be empty".to_string()),
),
_ => Ok(()),
}
}
fn validate_flow_step_position(conn: &Connection, ask: &Ask) -> StoreResult<()> {
let AskBody::FlowStep {
flow,
node_id,
skill,
iteration,
} = &ask.request
else {
return Ok(());
};
let epoch_id = ask_epoch_id_in(conn, &ask.id)?;
let position = conn
.query_row(
"SELECT flow, step, node_id, human, iteration
FROM work_flow_positions WHERE epoch_id=?1",
[epoch_id.as_str()],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, Option<String>>(2)?,
row.get::<_, bool>(3)?,
row.get::<_, i64>(4)?,
))
},
)
.optional()?;
if position.as_ref()
!= Some(&(
flow.clone(),
skill.clone(),
Some(node_id.clone()),
true,
i64::from(*iteration),
))
{
return Err(StoreError::InvalidAuthority(format!(
"Ask {} no longer matches the persisted human flow position",
ask.id
)));
}
Ok(())
}
fn current_flow_step_run(conn: &Connection, ask: &Ask) -> StoreResult<Option<Run>> {
let ask_epoch_id = ask_epoch_id_in(conn, &ask.id)?;
let Some(run) = current_run_for_work_in(conn, &ask.origin.work)? else {
return Ok(None);
};
if run.epoch_id != ask_epoch_id {
return Err(StoreError::InvalidAuthority(format!(
"Ask {} current Run does not belong to its flow-step Epoch",
ask.id
)));
}
Ok(Some(run))
}
fn end_flow_step_run(conn: &Connection, ask: &Ask) -> StoreResult<()> {
if !matches!(ask.request, AskBody::FlowStep { .. }) {
return Ok(());
}
if let Some(run) = current_flow_step_run(conn, ask)? {
end_run_in(conn, &run, BoundaryState::Succeeded)?;
}
Ok(())
}
fn validate_ask_result(result: &AskResult) -> StoreResult<()> {
if result.text().trim().is_empty() {
return Err(StoreError::InvalidData(
"Ask result cannot be empty".to_string(),
));
}
Ok(())
}
fn flow_ask_in(
conn: &Connection,
epoch_id: &EpochId,
flow: &str,
node_id: &str,
skill: &str,
iteration: u32,
target: &AskTarget,
) -> StoreResult<Option<Ask>> {
let mut statement = conn.prepare(
"SELECT id FROM ask_exchanges
WHERE epoch_id=?1 AND request_kind='flow_step'
AND request_flow=?2 AND request_node_id=?3
AND request_skill=?4 AND request_iteration=?5
ORDER BY asked_at DESC, rowid DESC",
)?;
let ids = statement
.query_map(
params![epoch_id.as_str(), flow, node_id, skill, iteration],
|row| row.get::<_, String>(0),
)?
.collect::<Result<Vec<_>, _>>()?;
for id in ids {
let id = AskId::parse(&id).map_err(invalid_durable)?;
let ask = ask_by_id_in(conn, &id)?;
if &ask.target == target {
return Ok(Some(ask));
}
}
Ok(None)
}
fn ask_epoch_id_in(conn: &Connection, ask_id: &AskId) -> StoreResult<EpochId> {
let epoch_id = conn.query_row(
"SELECT epoch_id FROM ask_exchanges WHERE id=?1",
[ask_id.as_str()],
|row| row.get::<_, String>(0),
)?;
EpochId::parse(&epoch_id).map_err(invalid_durable)
}
fn require_open_ask_epoch(conn: &Connection, ask_id: &AskId) -> StoreResult<()> {
if ask_epoch_is_open_in(conn, ask_id)? {
Ok(())
} else {
Err(StoreError::InvalidAuthority(format!(
"Ask {ask_id} no longer belongs to an open Work Epoch"
)))
}
}
fn ask_run_in(conn: &Connection, ask: &Ask) -> StoreResult<Run> {
let run = if matches!(ask.request, AskBody::FlowStep { .. }) {
validate_flow_step_position(conn, ask)?;
current_flow_step_run(conn, ask)?.ok_or_else(|| {
StoreError::InvalidAuthority(format!("Ask {} has no active flow-step Run", ask.id))
})?
} else {
run_by_id_in(conn, &ask.origin.run_id)?
};
if run.work != ask.origin.work {
return Err(StoreError::InvalidData(format!(
"Ask {} supervising Run no longer maps to its Work",
ask.id
)));
}
if matches!(ask.request, AskBody::FlowStep { .. })
&& (run.state != RunState::Active || run.cwd.as_ref() != Some(&ask.origin.cwd))
{
return Err(StoreError::InvalidAuthority(format!(
"Ask {} current Run cannot supervise the flow step in its captured cwd",
ask.id
)));
}
Ok(run)
}
fn validate_ask_invocation(
conn: &Connection,
ask_id: &AskId,
invocation_id: &AgentInvocationId,
) -> StoreResult<AgentInvocation> {
let owned_ask_id = conn.query_row(
"SELECT answer_ask_id FROM agent_invocations WHERE id=?1",
[invocation_id.as_str()],
|row| row.get::<_, Option<String>>(0),
)?;
if owned_ask_id.as_deref() != Some(ask_id.as_str()) {
return Err(StoreError::InvalidAuthority(format!(
"Invocation {invocation_id} does not belong to Ask {ask_id}"
)));
}
supervised_invocation_in(conn, invocation_id)
}
fn validate_active_ask_invocation(
conn: &Connection,
ask_id: &AskId,
invocation_id: &AgentInvocationId,
) -> StoreResult<AgentInvocation> {
let invocation = validate_ask_invocation(conn, ask_id, invocation_id)?;
let ask = ask_by_id_in(conn, ask_id)?;
if ask.state != AskState::Claimed || ask.active_invocation_id.as_ref() != Some(&invocation.id) {
return Err(StoreError::InvalidAuthority(format!(
"Invocation {} no longer owns Ask {}",
invocation.id, ask.id
)));
}
Ok(invocation)
}
fn ask_presentation_in(
conn: &Connection,
invocation_id: &AgentInvocationId,
) -> StoreResult<(bool, bool)> {
conn.query_row(
"SELECT ask_ready_at IS NOT NULL, ask_presented_at IS NOT NULL
FROM agent_invocations WHERE id=?1 AND answer_ask_id IS NOT NULL",
[invocation_id.as_str()],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.map_err(StoreError::from)
}
fn present_ask_invocation(conn: &Connection, invocation_id: &AgentInvocationId) -> StoreResult<()> {
if conn.execute(
"UPDATE agent_invocations
SET ask_presented_at=COALESCE(ask_presented_at, ?2)
WHERE id=?1 AND ask_ready_at IS NOT NULL AND ended_at IS NULL",
params![invocation_id.as_str(), now_unix()],
)? != 1
{
return Err(StoreError::InvalidAuthority(format!(
"Ask Invocation {invocation_id} is not attachable"
)));
}
Ok(())
}
fn finish_ask_invocation_in(
conn: &Connection,
invocation_id: &AgentInvocationId,
outcome: BoundaryState,
reason: Option<&str>,
) -> StoreResult<()> {
if conn.execute(
"UPDATE agent_invocations
SET ended_at=COALESCE(ended_at, ?2), outcome=?3, handback_state=?4,
incomplete_reason=COALESCE(?5, incomplete_reason)
WHERE id=?1 AND answer_ask_id IS NOT NULL",
params![
invocation_id.as_str(),
now_unix(),
outcome.as_invocation_outcome(),
handback_state(outcome),
reason,
],
)? != 1
{
return Err(StoreError::NotFound);
}
Ok(())
}
fn ask_invocation_is_latest_in(
conn: &Connection,
ask_id: &AskId,
invocation_id: &AgentInvocationId,
) -> StoreResult<bool> {
let latest = conn
.query_row(
"SELECT id FROM agent_invocations WHERE answer_ask_id=?1
ORDER BY started_at DESC, rowid DESC LIMIT 1",
[ask_id.as_str()],
|row| row.get::<_, String>(0),
)
.optional()?;
Ok(latest.as_deref() == Some(invocation_id.as_str()))
}
fn requeue_ask_in(
conn: &Connection,
ask: &Ask,
invocation_id: &AgentInvocationId,
outcome: BoundaryState,
reason: Option<&str>,
) -> StoreResult<()> {
finish_ask_invocation_in(conn, invocation_id, outcome, reason)?;
if conn.execute(
"UPDATE ask_exchanges SET state='queued', active_invocation_id=NULL
WHERE id=?1 AND state='claimed' AND active_invocation_id=?2",
params![ask.id.as_str(), invocation_id.as_str()],
)? != 1
{
return Err(StoreError::InvalidAuthority(format!(
"Invocation {} no longer owns Ask {}",
invocation_id, ask.id
)));
}
Ok(())
}
fn write_terminal_ask_in(
conn: &Connection,
ask: &Ask,
result: &AskResult,
author: &Author,
terminal_at: i64,
) -> StoreResult<Ask> {
let (author_kind, author_id) = author_parts(author);
if conn.execute(
"UPDATE ask_exchanges
SET state=?2, active_invocation_id=NULL, result_kind=?2, result_text=?3,
terminal_author_kind=?4, terminal_author_id=?5, terminal_at=?6
WHERE id=?1 AND state=?7 AND active_invocation_id IS ?8",
params![
ask.id.as_str(),
result.state().as_str(),
result.text(),
author_kind,
author_id,
terminal_at,
ask.state.as_str(),
ask.active_invocation_id
.as_ref()
.map(AgentInvocationId::as_str),
],
)? != 1
{
return Err(StoreError::InvalidAuthority(format!(
"Ask {} changed before terminal settlement",
ask.id
)));
}
ask_by_id_in(conn, &ask.id)
}
fn ask_author_in(conn: &Connection, ask: &Ask) -> StoreResult<Author> {
match &ask.target {
AskTarget::User => Ok(Author::User),
AskTarget::Parent(work) => {
let epoch = current_epoch_in(conn, work)?;
current_run_for_work_in(conn, work)?
.or(latest_run_for_epoch_in(conn, &epoch.id)?)
.map(|run| Author::Run(run.id))
.ok_or_else(|| {
StoreError::InvalidAuthority(format!(
"parent {} {} has no Run to author Ask {} settlement",
work.kind(),
work.id(),
ask.id
))
})
}
}
}
fn author_parts(author: &Author) -> (Option<&'static str>, Option<&str>) {
match author {
Author::User => (Some("user"), None),
Author::Run(run_id) => (Some("run"), Some(run_id.as_str())),
}
}
fn normalize_reason(value: &str, label: &str) -> StoreResult<String> {
let value = value.trim();
if value.is_empty() {
return Err(StoreError::InvalidData(format!("{label} cannot be empty")));
}
Ok(value.to_string())
}
fn normalize_optional_reason(value: Option<&str>) -> StoreResult<Option<String>> {
value
.map(|value| normalize_reason(value, "Ask release reason"))
.transpose()
}
fn cancel_pending_asks_for_epoch(
conn: &Connection,
epoch_id: &EpochId,
reason: &str,
terminal_at: i64,
) -> StoreResult<()> {
let mut statement = conn.prepare(
"SELECT id FROM ask_exchanges
WHERE epoch_id=?1 AND state IN ('queued', 'claimed')
ORDER BY asked_at, rowid",
)?;
let ids = statement
.query_map([epoch_id.as_str()], |row| row.get::<_, String>(0))?
.collect::<Result<Vec<_>, _>>()?;
drop(statement);
conn.execute(
"UPDATE agent_invocations
SET ended_at=COALESCE(ended_at, ?3), outcome='failed',
handback_state='unknown', incomplete_reason=?2
WHERE answer_ask_id IN (
SELECT id FROM ask_exchanges
WHERE epoch_id=?1 AND state='claimed'
) AND ended_at IS NULL",
params![epoch_id.as_str(), reason, terminal_at],
)?;
conn.execute(
"UPDATE ask_exchanges
SET state='cancelled', active_invocation_id=NULL,
result_kind='cancelled', result_text=?2,
terminal_author_kind=NULL, terminal_author_id=NULL, terminal_at=?3
WHERE epoch_id=?1 AND state IN ('queued', 'claimed')",
params![epoch_id.as_str(), reason, terminal_at],
)?;
for id in ids {
let id = AskId::parse(&id).map_err(invalid_durable)?;
let ask = ask_by_id_in(conn, &id)?;
enqueue_ask_result_comment(conn, &ask)?;
}
Ok(())
}
pub(super) fn validate_control_caller(
conn: &Connection,
caller: Option<&RunLease>,
target: &WorkRef,
) -> StoreResult<()> {
let Some(lease) = caller else {
return Ok(());
};
let run = validate_run_lease(conn, lease)?;
if parent_work(conn, target)?.as_ref() == Some(&run.work) {
let parent_epoch = current_epoch_in(conn, &run.work)?;
let selected = conn
.query_row(
"SELECT t.epoch_id, t.basis_rev
FROM agent_turns t
JOIN agent_invocations i ON i.id=t.invocation_id
WHERE i.supervising_run_id=?1
ORDER BY (t.status='running') DESC, t.started_at DESC, t.rowid DESC
LIMIT 1",
[run.id.as_str()],
|row| Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)?)),
)
.optional()?
.map(|(epoch_id, revision)| {
Ok::<Basis, StoreError>(Basis {
epoch_id: EpochId::parse(&epoch_id).map_err(invalid_durable)?,
revision: revision as u64,
})
})
.transpose()?
.ok_or_else(|| {
StoreError::InvalidAuthority(format!(
"Run {} has no durable Turn Basis for child control",
run.id
))
})?;
validate_basis(&parent_epoch.current_basis, &selected)
} else {
Err(StoreError::InvalidAuthority(
"Run may control only immediate child Work".to_string(),
))
}
}
pub(crate) fn work_status_in(conn: &Connection, work: &WorkRef) -> StoreResult<WorkStatus> {
let epoch = latest_epoch_in(conn, work)?;
match epoch.state {
EpochState::Done => return Ok(WorkStatus::Done),
EpochState::Abandoned => return Ok(WorkStatus::Abandoned),
EpochState::Open => {}
}
if let Some(run) = run_for_epoch_in(conn, &epoch.id)? {
return Ok(WorkStatus::Running { run_id: run.id });
}
let wait = conn
.query_row(
"SELECT id, on_json, created_at, resolved_at FROM waits
WHERE epoch_id=?1 AND resolved_at IS NULL",
[epoch.id.as_str()],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, i64>(2)?,
row.get::<_, Option<i64>>(3)?,
))
},
)
.optional()?;
if let Some((id, on_json, created_at, resolved_at)) = wait {
return Ok(WorkStatus::Waiting {
wait: Wait {
id: WaitId::parse(&id).map_err(invalid_durable)?,
work: work.clone(),
epoch_id: epoch.id,
on: serde_json::from_str(&on_json)?,
created_at: OffsetDateTime::from_unix_timestamp(created_at)
.map_err(invalid_durable)?,
resolved_at: resolved_at
.map(OffsetDateTime::from_unix_timestamp)
.transpose()
.map_err(invalid_durable)?,
},
});
}
Ok(WorkStatus::Ready)
}
fn latest_epoch_in(conn: &Connection, work: &WorkRef) -> StoreResult<Epoch> {
let (column, id) = match work {
WorkRef::Wave(id) => ("wave_id", id.as_str()),
WorkRef::Project(id) => ("project_id", id.as_str()),
WorkRef::Task(id) => ("task_id", id.as_str()),
};
let sql = format!(
"SELECT id, number, state, current_rev, created_at, terminal_at
FROM epochs WHERE {column}=?1 ORDER BY number DESC LIMIT 1"
);
let row = conn.query_row(&sql, [id], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, i64>(1)?,
row.get::<_, String>(2)?,
row.get::<_, i64>(3)?,
row.get::<_, i64>(4)?,
row.get::<_, Option<i64>>(5)?,
))
})?;
let epoch_id = EpochId::parse(&row.0).map_err(invalid_durable)?;
Ok(Epoch {
id: epoch_id.clone(),
work: work.clone(),
number: row.1 as u32,
state: EpochState::parse(&row.2).map_err(invalid_durable)?,
current_basis: Basis {
epoch_id,
revision: row.3 as u64,
},
created_at: OffsetDateTime::from_unix_timestamp(row.4).map_err(invalid_durable)?,
terminal_at: row
.5
.map(OffsetDateTime::from_unix_timestamp)
.transpose()
.map_err(invalid_durable)?,
})
}
fn parse_work_ref(kind: &str, id: &str) -> StoreResult<WorkRef> {
match kind {
"wave" => WaveId::parse(id)
.map(WorkRef::Wave)
.map_err(invalid_durable),
"project" => ProjectId::parse(id)
.map(WorkRef::Project)
.map_err(invalid_durable),
"task" => TaskId::parse(id)
.map(WorkRef::Task)
.map_err(invalid_durable),
value => Err(StoreError::InvalidData(format!(
"invalid Work kind: {value}"
))),
}
}
fn work_from_parts(
parts: (Option<String>, Option<String>, Option<String>),
) -> StoreResult<WorkRef> {
match parts {
(Some(id), None, None) => parse_work_ref("wave", &id),
(None, Some(id), None) => parse_work_ref("project", &id),
(None, None, Some(id)) => parse_work_ref("task", &id),
_ => Err(StoreError::InvalidData(
"stored Epoch owns an invalid Work reference".to_string(),
)),
}
}
fn invalid_durable(error: impl std::fmt::Display) -> StoreError {
StoreError::InvalidData(error.to_string())
}
fn to_sqlite_conversion_error(error: impl std::fmt::Display) -> rusqlite::Error {
rusqlite::Error::FromSqlConversionFailure(
0,
rusqlite::types::Type::Text,
Box::new(std::io::Error::new(
std::io::ErrorKind::InvalidData,
error.to_string(),
)),
)
}
pub(crate) fn create_wave_spine(
tx: &Transaction<'_>,
wave_id: &WaveId,
name: &str,
repo: &str,
created_at: i64,
) -> StoreResult<()> {
let work = WorkRef::Wave(wave_id.clone());
inherit_placement(tx, &work, None, created_at)?;
let exists: bool = tx.query_row(
"SELECT EXISTS(SELECT 1 FROM epochs WHERE wave_id=?1)",
[wave_id.as_str()],
|row| row.get(0),
)?;
if exists {
return Ok(());
}
let epoch_id = EpochId::new();
tx.execute(
"INSERT INTO epochs (
id, number, wave_id, project_id, task_id, state, current_rev,
created_at, terminal_at
) VALUES (?1, 1, ?2, NULL, NULL, 'open', 0, ?3, NULL)",
params![epoch_id.as_str(), wave_id.as_str(), created_at],
)?;
insert_truth(
tx,
&epoch_id,
serde_json::json!({"name": name, "repo": repo}),
OffsetDateTime::from_unix_timestamp(created_at).map_err(invalid_durable)?,
)
}
pub(crate) fn create_project_spine(tx: &Transaction<'_>, project: &Project) -> StoreResult<()> {
let project_id = tx
.query_row(
"SELECT id FROM projects WHERE external_project_id=?1",
[project.plan.id.as_str()],
|row| row.get::<_, String>(0),
)
.optional()?
.unwrap_or_else(|| ProjectId::new().to_string());
tx.execute(
"INSERT OR IGNORE INTO projects (
id, wave_id, external_project_id, created_at
) VALUES (?1, ?2, ?3, ?4)",
params![
project_id,
project.wave_id.as_str(),
project.plan.id.as_str(),
project.created_at.unix_timestamp(),
],
)?;
let work = WorkRef::Project(ProjectId::parse(&project_id).map_err(invalid_durable)?);
let parent = WorkRef::Wave(project.wave_id.clone());
inherit_placement(
tx,
&work,
Some(&parent),
project.created_at.unix_timestamp(),
)?;
let epoch_id = EpochId::new();
let number: i64 = tx.query_row(
"SELECT COALESCE(MAX(number), 0) + 1 FROM epochs WHERE project_id=?1",
[&project_id],
|row| row.get(0),
)?;
tx.execute(
"INSERT INTO epochs (
id, number, wave_id, project_id, task_id, state, current_rev,
created_at, terminal_at
) VALUES (?1, ?2, NULL, ?3, NULL, 'open', 0, ?4, NULL)",
params![
epoch_id.as_str(),
number,
project_id,
project.updated_at.unix_timestamp(),
],
)?;
insert_truth(
tx,
&epoch_id,
serde_json::json!({
"external_project_id": project.plan.id.as_str(),
"slug": project.plan.slug,
"name": project.plan.name,
"prompt_context": project.plan.prompt_context,
"pm_snapshot_synced_at": project.plan.pm_snapshot_synced_at,
}),
project.updated_at,
)?;
Ok(())
}
pub(crate) fn create_task_spine(tx: &Transaction<'_>, task: &Task) -> StoreResult<()> {
let project_id = task.project_id.as_str().to_string();
let task_id = tx
.query_row(
"SELECT id FROM tasks WHERE external_issue_id=?1",
[task.plan.id.as_str()],
|row| row.get::<_, String>(0),
)
.optional()?
.unwrap_or_else(|| TaskId::new().to_string());
tx.execute(
"INSERT OR IGNORE INTO tasks (
id, project_id, external_issue_id, issue_identifier, created_at
) VALUES (?1, ?2, ?3, ?4, ?5)",
params![
task_id,
project_id,
task.plan.id.as_str(),
task.plan.identifier,
task.created_at.unix_timestamp(),
],
)?;
let work = WorkRef::Task(TaskId::parse(&task_id).map_err(invalid_durable)?);
let parent = WorkRef::Project(ProjectId::parse(&project_id).map_err(invalid_durable)?);
inherit_placement(tx, &work, Some(&parent), task.created_at.unix_timestamp())?;
let epoch_id = EpochId::new();
let number: i64 = tx.query_row(
"SELECT COALESCE(MAX(number), 0) + 1 FROM epochs WHERE task_id=?1",
[&task_id],
|row| row.get(0),
)?;
tx.execute(
"INSERT INTO epochs (
id, number, wave_id, project_id, task_id, state, current_rev,
created_at, terminal_at
) VALUES (?1, ?2, NULL, NULL, ?3, 'open', 0, ?4, NULL)",
params![
epoch_id.as_str(),
number,
task_id,
task.updated_at.unix_timestamp(),
],
)?;
insert_truth(
tx,
&epoch_id,
serde_json::json!({
"external_issue_id": task.plan.id.as_str(),
"identifier": task.plan.identifier,
"title": task.plan.title,
"description": task.plan.description,
"pm_snapshot_synced_at": task.plan.pm_snapshot_synced_at,
}),
task.updated_at,
)?;
Ok(())
}
pub(crate) fn end_run_for_lease(
conn: &Connection,
lease: &RunLease,
outcome: BoundaryState,
) -> StoreResult<()> {
let run = validate_stop_lease(conn, lease)?;
end_run_in(conn, &run, outcome)
}
fn end_run_in(conn: &Connection, run: &Run, outcome: BoundaryState) -> StoreResult<()> {
if !outcome.is_terminal() {
return Err(StoreError::InvalidData(
"Run finish outcome must be terminal".to_string(),
));
}
let now = now_unix();
let turn_outcome = if outcome == BoundaryState::Interrupted {
"interrupted"
} else {
"failed"
};
end_open_turns_for_run(conn, &run.id, now, turn_outcome)?;
conn.execute(
"UPDATE agent_invocations SET
ended_at=COALESCE(ended_at, ?2),
outcome=?3, handback_state=?4
WHERE supervising_run_id=?1 AND ended_at IS NULL",
params![
run.id.as_str(),
now,
outcome.as_invocation_outcome(),
handback_state(outcome)
],
)?;
conn.execute(
"UPDATE runs SET state='ended', ended_at=?2
WHERE id=?1 AND state != 'ended'",
params![run.id.as_str(), now],
)?;
Ok(())
}
fn end_open_turns_for_run(
conn: &Connection,
run_id: &RunId,
ended_at: i64,
outcome: &str,
) -> StoreResult<()> {
conn.execute(
"UPDATE agent_turns SET status=?3, ended_at=COALESCE(ended_at, ?2)
WHERE status='running' AND invocation_id IN (
SELECT id FROM agent_invocations WHERE supervising_run_id=?1
)",
params![run_id.as_str(), ended_at, outcome],
)?;
Ok(())
}
pub(crate) fn insert_seed_sends_for_turn(
tx: &Transaction<'_>,
turn_id: &str,
basis: &Basis,
) -> StoreResult<()> {
let current: i64 = tx.query_row(
"SELECT current_rev FROM epochs WHERE id=?1 AND state='open'",
[basis.epoch_id.as_str()],
|row| row.get(0),
)?;
if current != basis.revision as i64 {
return Err(StoreError::StaleBasis {
expected: format!("{}:{}", basis.epoch_id, basis.revision),
current: format!("{}:{current}", basis.epoch_id),
});
}
let applied: i64 = tx.query_row(
"SELECT COALESCE(MAX(basis_rev), -1)
FROM agent_turns WHERE epoch_id=?1 AND status='completed'",
[basis.epoch_id.as_str()],
|row| row.get(0),
)?;
let mut statement = tx.prepare(
"SELECT id FROM steers
WHERE epoch_id=?1 AND rev > ?2 AND rev <= ?3
ORDER BY rev",
)?;
let steer_ids = statement
.query_map(
params![basis.epoch_id.as_str(), applied, basis.revision as i64],
|row| row.get::<_, String>(0),
)?
.collect::<Result<Vec<_>, _>>()?;
drop(statement);
let now = now_unix();
for steer_id in steer_ids {
tx.execute(
"INSERT INTO sends (
id, steer_id, turn_id, via, state, provider_turn_id,
reason, attempted_at, finished_at
) VALUES (?1, ?2, ?3, 'seed', 'sent', NULL, NULL, ?4, ?4)",
params![SendId::new().as_str(), steer_id, turn_id, now],
)?;
}
Ok(())
}
fn insert_truth(
tx: &Connection,
epoch_id: &EpochId,
payload: serde_json::Value,
at: OffsetDateTime,
) -> StoreResult<()> {
let source_id = format!("truth:{}:0", epoch_id);
tx.execute(
"INSERT INTO epoch_revisions (epoch_id, rev, kind, source_id, created_at)
VALUES (?1, 0, 'truth', ?2, ?3)",
params![epoch_id.as_str(), source_id, at.unix_timestamp()],
)?;
tx.execute(
"INSERT INTO work_truth (epoch_id, rev, payload_json, created_at)
VALUES (?1, 0, ?2, ?3)",
params![epoch_id.as_str(), payload.to_string(), at.unix_timestamp()],
)?;
Ok(())
}
pub(crate) fn validate_basis(current: &Basis, expected: &Basis) -> StoreResult<()> {
if current == expected {
return Ok(());
}
Err(StoreError::StaleBasis {
expected: format!("{}:{}", expected.epoch_id, expected.revision),
current: format!("{}:{}", current.epoch_id, current.revision),
})
}
pub(crate) fn validate_completion_readiness_in(conn: &Connection, run: &Run) -> StoreResult<()> {
let invocation_open: bool = conn.query_row(
"SELECT EXISTS(
SELECT 1 FROM agent_invocations
WHERE supervising_run_id=?1 AND ended_at IS NULL
)",
[run.id.as_str()],
|row| row.get(0),
)?;
if invocation_open {
return Err(StoreError::InvalidData(
"Run has an open Invocation".to_string(),
));
}
let child_ask_open: bool = conn.query_row(
"SELECT EXISTS(
SELECT 1 FROM ask_exchanges a
JOIN epochs child_epoch ON child_epoch.id=a.epoch_id
WHERE a.state IN ('queued', 'claimed') AND child_epoch.state='open'
AND a.target_kind='parent' AND a.target_work_kind=?1
AND a.target_work_id=?2
)",
params![run.work.kind(), run.work.id()],
|row| row.get(0),
)?;
if child_ask_open {
return Err(StoreError::InvalidData(
"Run cannot complete while a child Ask is unanswered".to_string(),
));
}
Ok(())
}
fn validate_author(tx: &Transaction<'_>, target: &WorkRef, author: &Author) -> StoreResult<()> {
let Author::Run(run_id) = author else {
return Ok(());
};
let source = tx
.query_row(
"SELECT e.wave_id, e.project_id, e.task_id
FROM runs r JOIN epochs e ON e.id=r.epoch_id
WHERE r.id=?1 AND r.state IN ('reserved', 'active')",
[run_id.as_str()],
|row| {
Ok((
row.get::<_, Option<String>>(0)?,
row.get::<_, Option<String>>(1)?,
row.get::<_, Option<String>>(2)?,
))
},
)
.optional()?
.ok_or_else(|| {
StoreError::InvalidAuthority("Run no longer holds execution authority".to_string())
})?;
let allowed = match target {
WorkRef::Project(project_id) => {
let parent: String = tx.query_row(
"SELECT wave_id FROM projects WHERE id=?1",
[project_id.as_str()],
|row| row.get(0),
)?;
source.0.as_deref() == Some(parent.as_str())
}
WorkRef::Task(task_id) => {
let parent: String = tx.query_row(
"SELECT project_id FROM tasks WHERE id=?1",
[task_id.as_str()],
|row| row.get(0),
)?;
source.1.as_deref() == Some(parent.as_str())
}
WorkRef::Wave(_) => false,
};
if allowed {
Ok(())
} else {
Err(StoreError::InvalidAuthority(format!(
"Run {run_id} may steer only immediate child Work"
)))
}
}
pub(crate) fn work_for_child_in(conn: &Connection, target: &ChildRef) -> StoreResult<WorkRef> {
match target {
ChildRef::Project(project_id) => {
conn.query_row(
"SELECT 1 FROM projects WHERE id=?1",
[project_id.as_str()],
|_| Ok(()),
)?;
Ok(WorkRef::Project(project_id.clone()))
}
ChildRef::Task(task_id) => {
conn.query_row(
"SELECT 1 FROM tasks WHERE id=?1",
[task_id.as_str()],
|_| Ok(()),
)?;
Ok(WorkRef::Task(task_id.clone()))
}
}
}
pub(crate) fn current_epoch_in(conn: &Connection, work: &WorkRef) -> StoreResult<Epoch> {
let (column, id) = match work {
WorkRef::Wave(id) => ("wave_id", id.as_str()),
WorkRef::Project(id) => ("project_id", id.as_str()),
WorkRef::Task(id) => ("task_id", id.as_str()),
};
let sql = format!(
"SELECT id, number, state, current_rev, created_at, terminal_at
FROM epochs WHERE {column}=?1 AND state='open'"
);
conn.query_row(&sql, [id], |row| {
let epoch_id = row.get::<_, String>(0)?;
let state = row.get::<_, String>(2)?;
Ok((
epoch_id,
row.get::<_, i64>(1)?,
state,
row.get::<_, i64>(3)?,
row.get::<_, i64>(4)?,
row.get::<_, Option<i64>>(5)?,
))
})
.optional()?
.map(
|(epoch_id, number, state, revision, created_at, terminal_at)| -> StoreResult<Epoch> {
let id = EpochId::parse(&epoch_id).map_err(|error| {
StoreError::InvalidData(format!("invalid stored Epoch id: {error}"))
})?;
Ok(Epoch {
id: id.clone(),
work: work.clone(),
number: number as u32,
state: EpochState::parse(&state)
.map_err(|error| StoreError::InvalidData(error.to_string()))?,
current_basis: Basis {
epoch_id: id,
revision: revision as u64,
},
created_at: OffsetDateTime::from_unix_timestamp(created_at).map_err(|error| {
StoreError::InvalidData(format!("invalid Epoch timestamp: {error}"))
})?,
terminal_at: terminal_at
.map(OffsetDateTime::from_unix_timestamp)
.transpose()
.map_err(|error| {
StoreError::InvalidData(format!(
"invalid Epoch terminal timestamp: {error}"
))
})?,
})
},
)
.transpose()?
.ok_or(StoreError::NotFound)
}
fn applied_basis_in(conn: &Connection, epoch_id: &EpochId) -> StoreResult<Option<Basis>> {
let revision = conn.query_row(
"SELECT MAX(basis_rev) FROM agent_turns
WHERE epoch_id=?1 AND status='completed'",
[epoch_id.as_str()],
|row| row.get::<_, Option<i64>>(0),
)?;
Ok(revision.map(|revision| Basis {
epoch_id: epoch_id.clone(),
revision: revision as u64,
}))
}
fn boundary_seed_in(conn: &Connection, work: &WorkRef) -> StoreResult<BoundarySeed> {
let epoch = current_epoch_in(conn, work)?;
boundary_seed_for_epoch_in(
conn,
work,
epoch.current_basis.epoch_id,
epoch.current_basis.revision,
)
}
fn boundary_seed_for_epoch_in(
conn: &Connection,
work: &WorkRef,
epoch_id: EpochId,
revision: u64,
) -> StoreResult<BoundarySeed> {
let applied = applied_basis_in(conn, &epoch_id)?
.map(|basis| basis.revision as i64)
.unwrap_or(-1);
let mut statement = conn.prepare(
"SELECT id, epoch_id, rev, author_kind, author_run_id, text, issued_at
FROM steers WHERE epoch_id=?1 AND rev > ?2 AND rev <= ?3 ORDER BY rev",
)?;
let rows = statement.query_map(
params![epoch_id.as_str(), applied, revision as i64],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, i64>(2)?,
row.get::<_, String>(3)?,
row.get::<_, Option<String>>(4)?,
row.get::<_, String>(5)?,
row.get::<_, i64>(6)?,
))
},
)?;
let mut steers = Vec::new();
for row in rows {
steers.push(decode_steer(row?, work.clone())?);
}
Ok(BoundarySeed {
basis: Basis { epoch_id, revision },
steers,
})
}
fn tool_response_in(
conn: &Connection,
epoch_id: &EpochId,
request_id: &str,
) -> StoreResult<Option<ToolResponseReceipt>> {
let row = conn
.query_row(
"SELECT id, rev, choice, responded_at FROM tool_responses
WHERE epoch_id=?1 AND request_id=?2",
params![epoch_id.as_str(), request_id],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, i64>(1)?,
row.get::<_, String>(2)?,
row.get::<_, i64>(3)?,
))
},
)
.optional()?;
let Some((id, revision, choice, responded_at)) = row else {
return Ok(None);
};
let work = work_for_epoch(conn, epoch_id)?;
Ok(Some(ToolResponseReceipt {
id: ToolResponseId::parse(&id).map_err(|error| {
StoreError::InvalidData(format!("invalid stored ToolResponse id: {error}"))
})?,
work,
basis: Basis {
epoch_id: epoch_id.clone(),
revision: revision as u64,
},
request_id: request_id.to_string(),
choice,
responded_at: OffsetDateTime::from_unix_timestamp(responded_at).map_err(|error| {
StoreError::InvalidData(format!("invalid Decision timestamp: {error}"))
})?,
}))
}
fn work_for_epoch(conn: &Connection, epoch_id: &EpochId) -> StoreResult<WorkRef> {
let row = conn.query_row(
"SELECT wave_id, project_id, task_id FROM epochs WHERE id=?1",
[epoch_id.as_str()],
|row| {
Ok((
row.get::<_, Option<String>>(0)?,
row.get::<_, Option<String>>(1)?,
row.get::<_, Option<String>>(2)?,
))
},
)?;
parse_work_columns(row.0, row.1, row.2)
}
type SteerFields = (String, String, i64, String, Option<String>, String, i64);
fn decode_steer(fields: SteerFields, work: WorkRef) -> StoreResult<Steer> {
let (id, epoch_id, revision, author_kind, author_run_id, text, issued_at) = fields;
let author = match (author_kind.as_str(), author_run_id) {
("user", None) => Author::User,
("run", Some(id)) => Author::Run(RunId::parse(&id).map_err(invalid_durable)?),
_ => {
return Err(StoreError::InvalidData(
"stored Steer author is inconsistent".to_string(),
))
}
};
Ok(Steer {
id: SteerId::parse(&id).map_err(invalid_durable)?,
work,
basis: Basis {
epoch_id: EpochId::parse(&epoch_id).map_err(invalid_durable)?,
revision: u64::try_from(revision).map_err(|_| {
StoreError::InvalidData(format!("invalid stored Steer revision: {revision}"))
})?,
},
author,
text,
issued_at: OffsetDateTime::from_unix_timestamp(issued_at).map_err(|error| {
StoreError::InvalidData(format!("invalid Steer timestamp: {error}"))
})?,
})
}
fn parse_work_columns(
wave_id: Option<String>,
project_id: Option<String>,
task_id: Option<String>,
) -> StoreResult<WorkRef> {
match (wave_id, project_id, task_id) {
(Some(id), None, None) => Ok(WorkRef::Wave(WaveId::parse(&id).map_err(invalid_durable)?)),
(None, Some(id), None) => Ok(WorkRef::Project(
ProjectId::parse(&id).map_err(invalid_durable)?,
)),
(None, None, Some(id)) => Ok(WorkRef::Task(TaskId::parse(&id).map_err(invalid_durable)?)),
_ => Err(StoreError::InvalidData(
"stored Epoch Work identity is inconsistent".to_string(),
)),
}
}
fn insert_send(conn: &Connection, send: &Send) -> StoreResult<()> {
conn.execute(
"INSERT INTO sends (
id, steer_id, turn_id, via, state, provider_turn_id, reason,
attempted_at, finished_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
params![
send.id.as_str(),
send.steer_id.as_str(),
send.turn_id.as_str(),
send.via.as_str(),
send.state.as_str(),
send.provider_turn_id,
send.reason,
send.attempted_at.unix_timestamp(),
send.finished_at.map(|at| at.unix_timestamp()),
],
)?;
Ok(())
}
fn send_for(
conn: &Connection,
steer_id: &SteerId,
turn_id: &str,
via: SendVia,
) -> StoreResult<Option<Send>> {
conn.query_row(
"SELECT id, steer_id, turn_id, via, state, provider_turn_id, reason,
attempted_at, finished_at
FROM sends WHERE steer_id=?1 AND turn_id=?2 AND via=?3",
params![steer_id.as_str(), turn_id, via.as_str()],
map_send,
)
.optional()
.map_err(StoreError::from)
}
fn send_by_id(conn: &Connection, send_id: &SendId) -> StoreResult<Option<Send>> {
conn.query_row(
"SELECT id, steer_id, turn_id, via, state, provider_turn_id, reason,
attempted_at, finished_at
FROM sends WHERE id=?1",
[send_id.as_str()],
map_send,
)
.optional()
.map_err(StoreError::from)
}
fn map_send(row: &rusqlite::Row<'_>) -> rusqlite::Result<Send> {
let id = row.get::<_, String>(0)?;
let steer_id = row.get::<_, String>(1)?;
let via = row.get::<_, String>(3)?;
let state = row.get::<_, String>(4)?;
let attempted_at = row.get::<_, i64>(7)?;
let finished_at = row.get::<_, Option<i64>>(8)?;
Ok(Send {
id: SendId::parse(&id).map_err(to_sql_error)?,
steer_id: SteerId::parse(&steer_id).map_err(to_sql_error)?,
turn_id: row.get(2)?,
via: match via.as_str() {
"live" => SendVia::Live,
"seed" => SendVia::Seed,
_ => return Err(to_sql_error(format!("invalid send via: {via}"))),
},
state: SendState::parse(&state).map_err(to_sql_error)?,
provider_turn_id: row.get(5)?,
reason: row.get(6)?,
attempted_at: OffsetDateTime::from_unix_timestamp(attempted_at).map_err(to_sql_error)?,
finished_at: finished_at
.map(OffsetDateTime::from_unix_timestamp)
.transpose()
.map_err(to_sql_error)?,
})
}
fn to_sql_error(error: impl std::fmt::Display) -> rusqlite::Error {
rusqlite::Error::FromSqlConversionFailure(
0,
rusqlite::types::Type::Text,
Box::new(std::io::Error::new(
std::io::ErrorKind::InvalidData,
error.to_string(),
)),
)
}
#[cfg(test)]
mod durable_store_tests {
use crate::durable::{
AdvanceReceipt, Ask, AskBody, AskId, AskOrigin, AskState, AskTarget, Author, BoundaryState,
Containment, ContainmentObservation, EpochId, InvocationRoute, RunAdvance, RunControl,
RunState, RunTrigger, StopCause, WorkRef,
};
use crate::id::WaveId;
use crate::project::ProjectId;
use crate::store::sqlite::SqliteStore;
use crate::store::StoreError;
use crate::task::TaskId;
use std::path::{Path, PathBuf};
use time::OffsetDateTime;
fn table_exists(conn: &rusqlite::Connection, table: &str) -> bool {
conn.query_row(
"SELECT EXISTS(
SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1
)",
[table],
|row| row.get(0),
)
.unwrap()
}
fn column_exists(conn: &rusqlite::Connection, table: &str, column: &str) -> bool {
conn.query_row(
"SELECT EXISTS(
SELECT 1 FROM pragma_table_info(?1) WHERE name=?2
)",
rusqlite::params![table, column],
|row| row.get(0),
)
.unwrap()
}
fn store_with_wave() -> (tempfile::TempDir, SqliteStore, WorkRef) {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("loopflow.db");
let store = SqliteStore::new(&path).expect("open a fresh store");
let wave_id = WaveId::new();
let conn = rusqlite::Connection::open(&path).unwrap();
if !column_exists(&conn, "work_placements", "enabled") {
conn.execute_batch(&crate::store::migrations::migration_sql_for_test(
"work_enablement",
))
.unwrap();
}
conn.execute(
"INSERT INTO waves (id, name, repo, created_at, parent_wave_id)
VALUES (?1, 'probe', '/repo', 1700000000, NULL)",
[wave_id.as_str()],
)
.unwrap();
drop(conn);
{
let mut raw = rusqlite::Connection::open(&path).unwrap();
let tx = raw.transaction().unwrap();
super::create_wave_spine(&tx, &wave_id, "probe", "/repo", 1_700_000_000).unwrap();
tx.commit().unwrap();
}
let work = WorkRef::Wave(wave_id);
(dir, store, work)
}
#[test]
fn activity_steers_span_historical_epochs() {
let (dir, store, work) = store_with_wave();
let first = store
.append_steer(&work, &Author::User, "first direction", None)
.unwrap();
let second_epoch = EpochId::new();
let conn = rusqlite::Connection::open(dir.path().join("loopflow.db")).unwrap();
conn.execute(
"UPDATE epochs SET state='done', terminal_at=1700000010
WHERE id=?1",
[first.steer.basis.epoch_id.as_str()],
)
.unwrap();
conn.execute(
"INSERT INTO epochs (
id, number, wave_id, project_id, task_id, state, current_rev,
created_at, terminal_at
) VALUES (?1, 2, ?2, NULL, NULL, 'open', 0, 1700000020, NULL)",
rusqlite::params![second_epoch.as_str(), work.id()],
)
.unwrap();
drop(conn);
let second = store
.append_steer(&work, &Author::User, "second direction", None)
.unwrap();
let conn = rusqlite::Connection::open(dir.path().join("loopflow.db")).unwrap();
conn.execute(
"UPDATE steers SET issued_at=1700000001 WHERE id=?1",
[first.steer.id.as_str()],
)
.unwrap();
conn.execute(
"UPDATE steers SET issued_at=1700000021 WHERE id=?1",
[second.steer.id.as_str()],
)
.unwrap();
let steers = store.list_steers_since(0).unwrap();
assert_eq!(
steers
.iter()
.map(|steer| steer.text.as_str())
.collect::<Vec<_>>(),
["second direction", "first direction"]
);
assert_eq!(steers[0].basis.epoch_id, second.steer.basis.epoch_id);
assert_eq!(steers[1].basis.epoch_id, first.steer.basis.epoch_id);
assert!(steers.iter().all(|steer| steer.work == work));
assert_eq!(
store
.list_steers_since(1_700_000_010)
.unwrap()
.into_iter()
.map(|steer| steer.id)
.collect::<Vec<_>>(),
[second.steer.id]
);
}
fn start_invocation(
store: &SqliteStore,
work: &WorkRef,
) -> (crate::durable::RunLease, crate::durable::AgentInvocation) {
let (_, lease) = store
.reserve_run(work, &RunTrigger::User)
.expect("reserve a Run");
let cwd = PathBuf::from("/repo");
store
.advance_run(
&lease,
&RunAdvance::RunStarting {
containment: Containment::Tmux {
name: "probe".to_string(),
},
cwd: cwd.clone(),
},
)
.expect("start a Run");
let receipt = store
.advance_run(
&lease,
&RunAdvance::InvocationStarting {
route: InvocationRoute {
provider: "codex".to_string(),
model: None,
account_id: None,
},
surface: "headless".to_string(),
resume_token: None,
answer_ask_id: None,
},
)
.expect("start an AgentInvocation");
let AdvanceReceipt::Invocation(invocation) = receipt else {
panic!("expected AgentInvocation receipt")
};
(lease, invocation)
}
fn store_with_work_hierarchy() -> (tempfile::TempDir, SqliteStore, Vec<(WorkRef, WorkRef)>) {
let (directory, store, wave) = store_with_wave();
let path = directory.path().join("loopflow.db");
let project_id = ProjectId::new();
let project = WorkRef::Project(project_id.clone());
let task_id = TaskId::new();
let task = WorkRef::Task(task_id.clone());
let mut conn = rusqlite::Connection::open(path).unwrap();
let tx = conn.transaction().unwrap();
tx.execute(
"INSERT INTO projects (id, wave_id, external_project_id, created_at)
VALUES (?1, ?2, 'linear-project', 1700000001)",
rusqlite::params![project_id.as_str(), wave.id()],
)
.unwrap();
super::inherit_placement(&tx, &project, Some(&wave), 1_700_000_001).unwrap();
let project_epoch_id = EpochId::new();
tx.execute(
"INSERT INTO epochs (
id, number, wave_id, project_id, task_id, state, current_rev,
created_at, terminal_at
) VALUES (?1, 1, NULL, ?2, NULL, 'open', 0, 1700000001, NULL)",
rusqlite::params![project_epoch_id.as_str(), project_id.as_str()],
)
.unwrap();
super::insert_truth(
&tx,
&project_epoch_id,
serde_json::json!({"external_project_id": "linear-project"}),
OffsetDateTime::from_unix_timestamp(1_700_000_001).unwrap(),
)
.unwrap();
tx.execute(
"INSERT INTO tasks (
id, project_id, external_issue_id, issue_identifier, created_at
) VALUES (?1, ?2, 'linear-issue', 'ENG-1', 1700000002)",
rusqlite::params![task_id.as_str(), project_id.as_str()],
)
.unwrap();
super::inherit_placement(&tx, &task, Some(&project), 1_700_000_002).unwrap();
let task_epoch_id = EpochId::new();
tx.execute(
"INSERT INTO epochs (
id, number, wave_id, project_id, task_id, state, current_rev,
created_at, terminal_at
) VALUES (?1, 1, NULL, NULL, ?2, 'open', 0, 1700000002, NULL)",
rusqlite::params![task_epoch_id.as_str(), task_id.as_str()],
)
.unwrap();
super::insert_truth(
&tx,
&task_epoch_id,
serde_json::json!({"issue_identifier": "ENG-1"}),
OffsetDateTime::from_unix_timestamp(1_700_000_002).unwrap(),
)
.unwrap();
tx.commit().unwrap();
(
directory,
store,
vec![
(wave.clone(), wave.clone()),
(project.clone(), wave),
(task, project),
],
)
}
#[test]
fn disabled_ancestor_does_not_block_a_manual_descendant_run() {
let (_directory, store, work_routes) = store_with_work_hierarchy();
let wave = &work_routes[0].0;
let task = &work_routes[2].0;
store.set_work_enabled(wave, false).unwrap();
store.reserve_run(task, &RunTrigger::User).unwrap();
}
fn start_turn(
store: &SqliteStore,
work: &WorkRef,
) -> (
crate::durable::RunLease,
crate::durable::AgentInvocation,
crate::durable::Turn,
) {
let (lease, invocation) = start_invocation(store, work);
let receipt = store
.advance_run(
&lease,
&RunAdvance::TurnStarting {
invocation_id: invocation.id.clone(),
},
)
.unwrap();
let AdvanceReceipt::Turn(turn) = receipt else {
panic!("expected Turn receipt")
};
(lease, invocation, turn)
}
#[test]
fn status_ask_sql_prepares_against_the_migrated_schema() {
let (directory, _, _) = store_with_wave();
let conn = rusqlite::Connection::open(directory.path().join("loopflow.db")).unwrap();
conn.prepare(super::HAS_PENDING_USER_ASK_FOR_WORK_SQL)
.expect("runtime status SQL must prepare against the migration head");
}
#[test]
fn pending_user_ask_status_resolves_every_work_kind_and_excludes_parent_routes() {
let (directory, store, work_routes) = store_with_work_hierarchy();
let path = directory.path().join("loopflow.db");
for (work, parent) in &work_routes {
let (lease, invocation, turn) = start_turn(&store, work);
let run = store.run_by_id(&lease.run_id).unwrap();
let ask = Ask {
id: AskId::new(),
origin: AskOrigin {
work: work.clone(),
run_id: run.id,
turn_id: Some(turn.id),
invocation_id: Some(invocation.id),
home_id: run.home_id,
cwd: run.cwd.unwrap(),
},
target: AskTarget::User,
request: AskBody::Intervention {
prompt: format!("What blocks {}?", work.kind()),
},
state: AskState::Queued,
active_invocation_id: None,
result: None,
terminal_author: None,
asked_at: OffsetDateTime::now_utc(),
terminal_at: None,
};
let conn = rusqlite::Connection::open(&path).unwrap();
super::insert_ask(&conn, &lease.basis.epoch_id, &ask).unwrap();
for (candidate, _) in &work_routes {
assert_eq!(
store.has_pending_user_ask_for_work(candidate).unwrap(),
candidate == work,
"the User Ask must belong only to its {} Epoch",
work.kind()
);
}
conn.execute(
"UPDATE ask_exchanges
SET target_kind='parent', target_work_kind=?2, target_work_id=?3
WHERE id=?1",
rusqlite::params![ask.id.as_str(), parent.kind(), parent.id()],
)
.unwrap();
assert!(!store.has_pending_user_ask_for_work(work).unwrap());
let conn = store.conn.lock().expect("store mutex poisoned");
let pending = super::query_asks(
&conn,
super::AskScope::Target(&AskTarget::Parent(parent.clone())),
)
.unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].target, AskTarget::Parent(parent.clone()));
conn.execute("DELETE FROM ask_exchanges WHERE id=?1", [ask.id.as_str()])
.unwrap();
}
}
#[test]
fn run_execution_shape_is_enforced_and_containment_is_immutable() {
let (dir, store, work) = store_with_wave();
let path = dir.path().join("loopflow.db");
let (run, lease) = store
.reserve_run(&work, &RunTrigger::User)
.expect("reserve a Run");
let conn = rusqlite::Connection::open(&path).unwrap();
assert!(conn
.execute(
"UPDATE runs SET state='active' WHERE id=?1",
[run.id.as_str()],
)
.is_err());
store
.advance_run(
&lease,
&RunAdvance::RunStarting {
containment: Containment::Tmux {
name: "probe".to_string(),
},
cwd: PathBuf::from("/repo"),
},
)
.expect("start a Run with complete containment");
assert!(conn
.execute(
"UPDATE runs SET containment_id='replacement' WHERE id=?1",
[run.id.as_str()],
)
.is_err());
assert!(conn
.execute(
"UPDATE runs SET containment_kind=NULL WHERE id=?1",
[run.id.as_str()],
)
.is_err());
store
.stop_run(
&lease,
&StopCause::Requested,
ContainmentObservation::Absent,
)
.expect("end the contained Run");
let ended = store.run_by_id(&run.id).unwrap();
assert_eq!(ended.state, RunState::Ended);
assert_eq!(
ended.containment,
Some(Containment::Tmux {
name: "probe".to_string()
})
);
assert_eq!(ended.cwd, Some(PathBuf::from("/repo")));
assert!(ended.started_at.is_some());
}
fn live_invocation_id(store: &SqliteStore, path: &Path) -> crate::durable::AgentInvocationId {
let _ = store;
let conn = rusqlite::Connection::open(path).unwrap();
let id: String = conn
.query_row("SELECT id FROM agent_invocations LIMIT 1", [], |row| {
row.get(0)
})
.unwrap();
crate::durable::AgentInvocationId::parse(&id).unwrap()
}
#[test]
fn a_running_invocation_records_its_provider_continuity() {
let (dir, store, work) = store_with_wave();
let path = dir.path().join("loopflow.db");
let (lease, _) = start_invocation(&store, &work);
let invocation_id = live_invocation_id(&store, &path);
let invocation = store
.observe_invocation_provider(&lease, &invocation_id, None, Some("thread_abc"))
.expect("record the observed provider");
assert_eq!(invocation.resume_token.as_deref(), Some("thread_abc"));
assert_eq!(
store.current_run(&work).unwrap().unwrap().containment,
Some(Containment::Tmux {
name: "probe".to_string()
}),
"containment is spawn-time fencing evidence and must survive a provider observation"
);
}
#[test]
fn an_empty_observation_never_erases_recorded_continuity() {
let (dir, store, work) = store_with_wave();
let path = dir.path().join("loopflow.db");
let (lease, _) = start_invocation(&store, &work);
let invocation_id = live_invocation_id(&store, &path);
store
.observe_invocation_provider(&lease, &invocation_id, None, Some("thread_abc"))
.unwrap();
let invocation = store
.observe_invocation_provider(&lease, &invocation_id, None, None)
.unwrap();
assert_eq!(invocation.resume_token.as_deref(), Some("thread_abc"));
}
#[test]
fn an_invocation_records_one_exact_account_route_and_rejects_route_drift() {
let (dir, store, work) = store_with_wave();
let path = dir.path().join("loopflow.db");
let (lease, _) = start_invocation(&store, &work);
let invocation_id = live_invocation_id(&store, &path);
let work_account = crate::store::ProviderAccountId::parse("work").unwrap();
let personal_account = crate::store::ProviderAccountId::parse("personal").unwrap();
let invocation = store
.observe_invocation_provider(&lease, &invocation_id, Some(&work_account), None)
.unwrap();
assert_eq!(invocation.route.account_id.as_deref(), Some("work"));
assert!(store
.observe_invocation_provider(&lease, &invocation_id, Some(&personal_account), None)
.is_err());
}
#[test]
fn a_stopped_run_cannot_record_a_provider_observation() {
let (dir, store, work) = store_with_wave();
let path = dir.path().join("loopflow.db");
let (lease, _) = start_invocation(&store, &work);
let invocation_id = live_invocation_id(&store, &path);
store
.stop_run(
&lease,
&StopCause::Requested,
crate::durable::ContainmentObservation::Absent,
)
.expect("stop the Run with proven containment absence");
let error = store
.observe_invocation_provider(&lease, &invocation_id, None, Some("thread_zzz"))
.expect_err("a stopped Run must not remain a writer");
assert!(
matches!(
error,
StoreError::InvalidAuthority(_) | StoreError::InvalidData(_)
),
"expected an authority refusal, got {error:?}"
);
}
#[test]
fn an_interrupted_run_keeps_only_cleanup_authority() {
let (dir, store, work) = store_with_wave();
let (lease, _) = start_invocation(&store, &work);
store
.interrupt(None, &work, &lease.run_id)
.expect("mark the Run for interruption");
assert!(store.validate_run_lease(&lease).is_err());
assert_eq!(
store.run_control(&lease, None).unwrap(),
Some(RunControl::Interrupt)
);
let conn = rusqlite::Connection::open(dir.path().join("loopflow.db")).unwrap();
super::end_run_for_lease(&conn, &lease, BoundaryState::Interrupted)
.expect("the stopped runner can finish cleanup");
assert_eq!(
store.run_by_id(&lease.run_id).unwrap().state,
RunState::Ended
);
}
#[test]
fn a_home_upgrade_quiesces_at_the_provider_boundary() {
let (_directory, store, work) = store_with_wave();
let (lease, _) = start_invocation(&store, &work);
store
.stop_run(
&lease,
&StopCause::HomeUpgrade {
upgrade_id: "upgrade-one".to_string(),
deadline: 1_900_000_000,
},
crate::durable::ContainmentObservation::Present,
)
.unwrap();
assert_eq!(
store.run_control(&lease, None).unwrap(),
Some(RunControl::Quiesce {
upgrade_id: "upgrade-one".to_string(),
deadline: 1_900_000_000,
})
);
}
#[test]
fn a_new_run_records_the_active_home_runtime_generation() {
let (directory, store, work) = store_with_wave();
let connection = rusqlite::Connection::open(directory.path().join("loopflow.db")).unwrap();
if !table_exists(&connection, "home_runtime_generations") {
connection
.execute_batch(&crate::store::migrations::migration_sql_for_test(
"home_runtime_generation",
))
.unwrap();
}
let home_id: String = connection
.query_row("SELECT id FROM homes WHERE route='local'", [], |row| {
row.get(0)
})
.unwrap();
connection
.execute(
"INSERT INTO home_runtime_generations (
home_id, generation, build_version, source_revision,
migration_frontier, activated_at
) VALUES (?1, 12, '0.12.6', 'abc', '0.12.6.001_release', 1)",
[home_id],
)
.unwrap();
let (run, _) = store.reserve_run(&work, &RunTrigger::User).unwrap();
assert_eq!(run.runtime_generation, Some(12));
assert_eq!(
store.run_by_id(&run.id).unwrap().runtime_generation,
Some(12)
);
}
#[test]
fn an_upgrade_fences_old_generations_inside_run_reservation() {
let (directory, store, work) = store_with_wave();
let connection = rusqlite::Connection::open(directory.path().join("loopflow.db")).unwrap();
if !table_exists(&connection, "home_runtime_generations") {
connection
.execute_batch(&crate::store::migrations::migration_sql_for_test(
"home_runtime_generation",
))
.unwrap();
}
let home_id: String = connection
.query_row("SELECT id FROM homes WHERE route='local'", [], |row| {
row.get(0)
})
.unwrap();
connection
.execute(
"INSERT INTO home_runtime_generations (
home_id, generation, build_version, source_revision,
migration_frontier, activated_at
) VALUES (?1, 12, '0.12.5', 'old', '0.12.5.001_release', 1)",
[&home_id],
)
.unwrap();
connection
.execute(
"INSERT INTO home_upgrades (
id, home_id, source_revision, source_identity,
migration_authority, package_version, latest_known_migration,
prior_generation, target_generation, phase, keeper_mode,
migration_required, started_at, artifacts_activated,
migration_applied, daemon_restarted, drain_timed_out,
coordinator_started_at
) VALUES (
'upgrade-next', ?1, 'new', 'release', 'published', '0.12.6',
'0.12.6.001_release', 12, 13, 'draining', 'none',
1, 1, 0, 0, 0, 0, 1
)",
[&home_id],
)
.unwrap();
let error = store.reserve_run(&work, &RunTrigger::User).unwrap_err();
assert!(matches!(
error,
StoreError::HomeUpgradeFenced {
runtime_generation: Some(12),
..
}
));
connection
.execute(
"UPDATE home_upgrades SET artifacts_activated=1, phase='reconciling'
WHERE id='upgrade-next'",
[],
)
.unwrap();
let error = store.reserve_run(&work, &RunTrigger::User).unwrap_err();
assert!(matches!(error, StoreError::HomeUpgradeFenced { .. }));
connection
.execute(
"INSERT INTO home_runtime_generations (
home_id, generation, build_version, source_revision,
migration_frontier, activated_at
) VALUES (?1, 13, '0.12.6', 'new', '0.12.6.001_release', 2)",
[&home_id],
)
.unwrap();
let (run, _) = store.reserve_run(&work, &RunTrigger::User).unwrap();
assert_eq!(run.runtime_generation, Some(13));
}
#[test]
fn proven_runner_loss_ends_its_incomplete_turns() {
let (dir, store, work) = store_with_wave();
let (lease, _) = start_invocation(&store, &work);
let invocation_id = live_invocation_id(&store, &dir.path().join("loopflow.db"));
let receipt = store
.advance_run(
&lease,
&RunAdvance::TurnStarting {
invocation_id: invocation_id.clone(),
},
)
.unwrap();
let crate::durable::AdvanceReceipt::Turn(turn) = receipt else {
panic!("expected Turn receipt")
};
store
.recover_run(
&lease.run_id,
crate::durable::ContainmentObservation::Absent,
)
.expect("recover proven missing containment");
let conn = rusqlite::Connection::open(dir.path().join("loopflow.db")).unwrap();
let (status, ended_at): (String, Option<i64>) = conn
.query_row(
"SELECT status, ended_at FROM agent_turns WHERE id=?1",
[turn.id.as_str()],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.unwrap();
assert_eq!(status, "failed");
assert!(ended_at.is_some());
}
}