use std::sync::Arc;
use ag_session::IntegrationApproach;
use async_trait::async_trait;
use sqlx::SqlitePool;
use super::status;
use crate::timestamp::TimestampSource;
use crate::{DbError, DbResultExt};
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct SessionOrchestrationRow {
pub controller_project_id: i64,
pub controller_session_id: String,
pub goal_statement: String,
pub id: i64,
pub max_parallelism: i64,
pub relayed_question_task_id: Option<i64>,
pub status: String,
pub verification_generation: i64,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct SessionOrchestrationTaskRow {
pub acceptance_criteria: String,
pub area_violations: String,
pub areas_compliant: Option<bool>,
pub attempt_count: i64,
pub child_added_lines: i64,
pub child_answer: Option<String>,
pub child_deleted_lines: i64,
pub child_focused_review_status: Option<String>,
pub child_focused_review_text: Option<String>,
pub child_has_diff: Option<bool>,
pub child_input_tokens: i64,
pub child_output_tokens: i64,
pub child_questions: Option<String>,
pub child_session_id: Option<String>,
pub child_status: Option<String>,
pub child_summary: Option<String>,
pub continuation_generation: i64,
pub continuation_prompt: Option<String>,
pub id: i64,
pub infrastructure_retry_count: i64,
pub kind: String,
pub last_error: Option<String>,
pub merge_position: i64,
pub prompt: String,
pub research_report: Option<String>,
pub result_summary: Option<String>,
pub review_iteration: i64,
pub status: String,
pub task_key: String,
pub touched_areas: String,
pub title: String,
pub verification_reason: Option<String>,
pub verification_verdict: Option<String>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct OrchestrationTaskScopeRow {
pub base_branch: String,
pub id: i64,
pub touched_areas: String,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct SessionOrchestrationMetadataRow {
pub controller_session_id: Option<String>,
pub orchestration_status: Option<String>,
pub running_task_count: i64,
pub session_id: String,
pub waiting_task_count: i64,
}
pub struct PersistedOrchestrationTask {
pub acceptance_criteria: String,
pub kind: String,
pub merge_position: i64,
pub prompt: String,
pub session_orchestration_id: i64,
pub task_key: String,
pub title: String,
pub touched_areas: String,
}
#[cfg_attr(any(test, feature = "test-utils"), mockall::automock)]
#[async_trait]
pub trait OrchestrationRepository: Send + Sync {
async fn insert_orchestration(
&self,
controller_session_id: &str,
status: &str,
max_parallelism: i64,
) -> Result<i64, DbError>;
async fn upsert_orchestration_task(
&self,
task: PersistedOrchestrationTask,
) -> Result<i64, DbError>;
async fn load_orchestration_for_controller(
&self,
controller_session_id: &str,
) -> Result<Option<SessionOrchestrationRow>, DbError>;
async fn load_active_orchestrations(&self) -> Result<Vec<SessionOrchestrationRow>, DbError>;
async fn load_recoverable_focused_review_session_ids(
&self,
project_id: i64,
) -> Result<Vec<String>, DbError>;
async fn load_session_metadata_for_project(
&self,
project_id: i64,
) -> Result<Vec<SessionOrchestrationMetadataRow>, DbError>;
async fn load_orchestration_tasks(
&self,
session_orchestration_id: i64,
) -> Result<Vec<SessionOrchestrationTaskRow>, DbError>;
async fn load_orchestration_integration_approach(&self, id: i64) -> Result<String, DbError>;
async fn load_orchestration_task_scope_for_child(
&self,
child_session_id: &str,
) -> Result<Option<OrchestrationTaskScopeRow>, DbError>;
async fn load_child_session_id_for_task(&self, task_id: i64)
-> Result<Option<String>, DbError>;
async fn begin_orchestration_cancellation(&self, id: i64) -> Result<bool, DbError>;
async fn claim_orchestration_task(&self, id: i64) -> Result<bool, DbError>;
async fn claim_orchestration_review_application(
&self,
id: i64,
prompt: &str,
iteration_limit: i64,
) -> Result<bool, DbError>;
async fn claim_orchestration_rollup(&self, id: i64) -> Result<bool, DbError>;
async fn complete_orchestration_rollup(&self, id: i64) -> Result<bool, DbError>;
async fn record_orchestration_verdict(
&self,
id: i64,
task_key: &str,
is_pass: bool,
reason: &str,
) -> Result<bool, DbError>;
async fn complete_orchestration_campaign(&self, id: i64) -> Result<bool, DbError>;
async fn load_rollup_operation_status(
&self,
operation_id: &str,
) -> Result<Option<String>, DbError>;
async fn update_orchestration_status(&self, id: i64, status: &str) -> Result<(), DbError>;
async fn approve_orchestration_plan(&self, id: i64) -> Result<bool, DbError>;
async fn approve_orchestration_integration(
&self,
id: i64,
approach: IntegrationApproach,
) -> Result<bool, DbError>;
async fn update_orchestration_plan(
&self,
id: i64,
goal_statement: &str,
max_parallelism: i64,
) -> Result<(), DbError>;
async fn queue_orchestration_continuation(
&self,
id: i64,
prompt: &str,
acceptance_criteria: &str,
touched_areas: &str,
) -> Result<bool, DbError>;
async fn reset_orchestration_verification(&self, id: i64) -> Result<(), DbError>;
async fn record_orchestration_spawn_failure(
&self,
id: i64,
error: &str,
retry_limit: i64,
) -> Result<String, DbError>;
async fn detach_orchestration_child(&self, child_session_id: &str) -> Result<bool, DbError>;
async fn surface_orchestration_questions(
&self,
session_orchestration_id: i64,
task_id: i64,
questions: &str,
) -> Result<bool, DbError>;
async fn clear_orchestration_questions(
&self,
session_orchestration_id: i64,
) -> Result<(), DbError>;
async fn link_orchestration_task_child(
&self,
id: i64,
child_session_id: &str,
) -> Result<bool, DbError>;
async fn update_orchestration_task_status(
&self,
id: i64,
status: &str,
last_error: Option<String>,
) -> Result<(), DbError>;
async fn update_orchestration_task_result_summary(
&self,
id: i64,
result_summary: &str,
) -> Result<(), DbError>;
async fn update_orchestration_task_research_report(
&self,
id: i64,
research_report: &str,
) -> Result<(), DbError>;
async fn update_orchestration_task_area_compliance(
&self,
id: i64,
areas_compliant: Option<bool>,
area_violations: &str,
) -> Result<(), DbError>;
}
#[derive(Clone)]
pub(crate) struct SqliteOrchestrationRepository(SqlitePool, Arc<dyn TimestampSource>);
struct OrchestrationTransition {
context: &'static str,
from_orchestration_status: &'static str,
from_task_status: &'static str,
require_pass_verdict: bool,
to_orchestration_status: &'static str,
to_task_status: &'static str,
}
impl SqliteOrchestrationRepository {
pub(crate) fn new(pool: SqlitePool, timestamp_source: Arc<dyn TimestampSource>) -> Self {
Self(pool, timestamp_source)
}
fn now(&self) -> i64 {
self.1.now_timestamp_seconds()
}
async fn transition_orchestration_and_tasks(
&self,
id: i64,
transition: OrchestrationTransition,
) -> Result<bool, DbError> {
let now = self.now();
let mut transaction = self.0.begin().await.db_context(transition.context)?;
let result = sqlx::query(
r"
UPDATE session_orchestration
SET status = ?,
updated_at = ?
WHERE id = ?
AND status = ?
",
)
.bind(transition.to_orchestration_status)
.bind(now)
.bind(id)
.bind(transition.from_orchestration_status)
.execute(&mut *transaction)
.await
.db_context(transition.context)?;
if result.rows_affected() == 1 {
sqlx::query(
r"
UPDATE session_orchestration_task
SET status = ?,
updated_at = ?
WHERE session_orchestration_id = ?
AND status = ?
AND (? = 0 OR verification_verdict = 'Pass')
",
)
.bind(transition.to_task_status)
.bind(now)
.bind(id)
.bind(transition.from_task_status)
.bind(i64::from(transition.require_pass_verdict))
.execute(&mut *transaction)
.await
.db_context(transition.context)?;
}
transaction.commit().await.db_context(transition.context)?;
Ok(result.rows_affected() == 1)
}
}
#[async_trait]
impl OrchestrationRepository for SqliteOrchestrationRepository {
async fn insert_orchestration(
&self,
controller_session_id: &str,
status: &str,
max_parallelism: i64,
) -> Result<i64, DbError> {
status::validate_orchestration(status)?;
let now = self.now();
let row = sqlx::query!(
r#"
INSERT INTO session_orchestration (
controller_session_id,
goal_statement,
status,
max_parallelism,
created_at,
updated_at
)
VALUES (?, '', ?, ?, ?, ?)
RETURNING id AS "id!: i64"
"#,
controller_session_id,
status,
max_parallelism,
now,
now
)
.fetch_one(&self.0)
.await?;
Ok(row.id)
}
async fn upsert_orchestration_task(
&self,
task: PersistedOrchestrationTask,
) -> Result<i64, DbError> {
let PersistedOrchestrationTask {
acceptance_criteria,
kind,
merge_position,
prompt,
session_orchestration_id,
task_key,
title,
touched_areas,
} = task;
kind.parse::<ag_session::OrchestrationTaskKind>()
.map_err(|_| DbError::InvalidData {
entity: "orchestration task kind",
reason: format!("unknown persisted kind `{kind}`"),
})?;
let now = self.now();
let mut transaction = self
.0
.begin()
.await
.db_context("upsert orchestration task")?;
let row = sqlx::query!(
r#"
INSERT INTO session_orchestration_task (
session_orchestration_id,
task_key,
title,
prompt,
touched_areas,
acceptance_criteria,
kind,
merge_position,
status,
created_at,
updated_at
)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, 'Planned', ?, ?)
ON CONFLICT(session_orchestration_id, task_key) DO UPDATE
SET title = excluded.title,
prompt = excluded.prompt,
kind = excluded.kind,
touched_areas = excluded.touched_areas,
acceptance_criteria = excluded.acceptance_criteria,
merge_position = excluded.merge_position,
status = 'Planned',
child_session_id = NULL,
continuation_prompt = NULL,
review_iteration = 0,
research_report = NULL,
result_summary = NULL,
verification_reason = NULL,
verification_verdict = NULL,
last_error = NULL,
updated_at = excluded.updated_at
RETURNING id AS "id!: i64"
"#,
session_orchestration_id,
task_key,
title,
prompt,
touched_areas,
acceptance_criteria,
kind,
merge_position,
now,
now
)
.fetch_one(&mut *transaction)
.await
.db_context("upsert orchestration task")?;
sqlx::query!(
r"
UPDATE session
SET orchestration_task_id = NULL,
updated_at = ?
WHERE orchestration_task_id = ?
",
now,
row.id
)
.execute(&mut *transaction)
.await
.db_context("upsert orchestration task")?;
transaction
.commit()
.await
.db_context("upsert orchestration task")?;
Ok(row.id)
}
async fn load_orchestration_for_controller(
&self,
controller_session_id: &str,
) -> Result<Option<SessionOrchestrationRow>, DbError> {
let row = sqlx::query_as!(
SessionOrchestrationRow,
r#"
SELECT orchestration.id AS "id!: i64",
session.project_id AS "controller_project_id!: i64",
orchestration.controller_session_id,
orchestration.goal_statement,
orchestration.relayed_question_task_id,
orchestration.status,
orchestration.max_parallelism,
orchestration.verification_generation
FROM session_orchestration AS orchestration
INNER JOIN session
ON session.id = orchestration.controller_session_id
WHERE orchestration.controller_session_id = ?
ORDER BY orchestration.id DESC
LIMIT 1
"#,
controller_session_id
)
.fetch_optional(&self.0)
.await?;
if let Some(row) = &row {
status::validate_orchestration(&row.status)?;
}
Ok(row)
}
async fn load_active_orchestrations(&self) -> Result<Vec<SessionOrchestrationRow>, DbError> {
let rows = sqlx::query_as!(
SessionOrchestrationRow,
r#"
SELECT orchestration.id AS "id!: i64",
session.project_id AS "controller_project_id!: i64",
orchestration.controller_session_id,
orchestration.goal_statement,
orchestration.relayed_question_task_id,
orchestration.status,
orchestration.max_parallelism,
orchestration.verification_generation
FROM session_orchestration AS orchestration
INNER JOIN session
ON session.id = orchestration.controller_session_id
WHERE orchestration.status IN (
'AwaitingApproval',
'Running',
'Verifying',
'AwaitingIntegration',
'Integrating',
'Canceling'
)
ORDER BY orchestration.id
"#
)
.fetch_all(&self.0)
.await?;
for row in &rows {
status::validate_orchestration(&row.status)?;
}
Ok(rows)
}
async fn load_recoverable_focused_review_session_ids(
&self,
project_id: i64,
) -> Result<Vec<String>, DbError> {
let session_ids = sqlx::query_scalar!(
r#"
SELECT child.id AS "id!: String"
FROM session_orchestration_task AS task
INNER JOIN session_orchestration AS orchestration
ON orchestration.id = task.session_orchestration_id
INNER JOIN session AS child
ON child.id = task.child_session_id
WHERE child.project_id = ?
AND orchestration.status IN ('AwaitingApproval', 'Running')
AND task.status = 'Reviewing'
AND child.status IN ('Review', 'AgentReview')
AND (
child.focused_review_status IS NULL
OR child.focused_review_status = 'Pending'
)
ORDER BY task.id
"#,
project_id
)
.fetch_all(&self.0)
.await?;
Ok(session_ids)
}
async fn load_session_metadata_for_project(
&self,
project_id: i64,
) -> Result<Vec<SessionOrchestrationMetadataRow>, DbError> {
let rows = sqlx::query_as!(
SessionOrchestrationMetadataRow,
r#"
WITH latest_orchestration_id AS (
SELECT controller_session_id,
MAX(id) AS orchestration_id
FROM session_orchestration
GROUP BY controller_session_id
),
latest_orchestration AS (
SELECT orchestration.id,
orchestration.controller_session_id,
orchestration.status
FROM session_orchestration AS orchestration
INNER JOIN latest_orchestration_id AS latest
ON latest.orchestration_id = orchestration.id
),
controller_metadata AS (
SELECT orchestration.controller_session_id AS session_id,
orchestration.status AS orchestration_status,
COALESCE(SUM(
CASE WHEN task.status IN (
'Creating',
'Running',
'Reviewing',
'ReviewApplying',
'ContinuationPending'
)
THEN 1 ELSE 0 END
), 0) AS running_task_count,
COALESCE(SUM(
CASE WHEN task.status = 'WaitingForInput' THEN 1 ELSE 0 END
), 0) AS waiting_task_count
FROM latest_orchestration AS orchestration
LEFT JOIN session_orchestration_task AS task
ON task.session_orchestration_id = orchestration.id
GROUP BY orchestration.id
),
child_metadata AS (
SELECT task.child_session_id AS session_id,
orchestration.controller_session_id
FROM session_orchestration_task AS task
INNER JOIN session_orchestration AS orchestration
ON orchestration.id = task.session_orchestration_id
WHERE task.child_session_id IS NOT NULL
)
SELECT session.id AS "session_id!: String",
child_metadata.controller_session_id,
controller_metadata.orchestration_status,
COALESCE(controller_metadata.running_task_count, 0) AS "running_task_count!: i64",
COALESCE(controller_metadata.waiting_task_count, 0) AS "waiting_task_count!: i64"
FROM session
LEFT JOIN controller_metadata
ON controller_metadata.session_id = session.id
LEFT JOIN child_metadata
ON child_metadata.session_id = session.id
WHERE session.project_id = ?
AND (
controller_metadata.session_id IS NOT NULL
OR child_metadata.session_id IS NOT NULL
)
ORDER BY session.id
"#,
project_id
)
.fetch_all(&self.0)
.await?;
for row in &rows {
if let Some(orchestration_status) = &row.orchestration_status {
status::validate_orchestration(orchestration_status)?;
}
}
Ok(rows)
}
async fn load_orchestration_tasks(
&self,
session_orchestration_id: i64,
) -> Result<Vec<SessionOrchestrationTaskRow>, DbError> {
let rows = sqlx::query_as!(
SessionOrchestrationTaskRow,
r#"
SELECT task.id AS "id!: i64",
task.acceptance_criteria,
task.area_violations,
task.areas_compliant AS "areas_compliant?: bool",
task.attempt_count,
COALESCE(child.added_lines, 0) AS "child_added_lines!: i64",
(
SELECT message.content
FROM session_message AS message
WHERE message.session_id = task.child_session_id
AND message.kind = 'assistant_answer'
ORDER BY message.position DESC
LIMIT 1
) AS child_answer,
COALESCE(child.deleted_lines, 0) AS "child_deleted_lines!: i64",
child.focused_review_status AS child_focused_review_status,
child.focused_review_text AS child_focused_review_text,
child.has_diff AS "child_has_diff?: bool",
COALESCE(child.input_tokens, 0) AS "child_input_tokens!: i64",
COALESCE(child.output_tokens, 0) AS "child_output_tokens!: i64",
task.child_session_id,
child.status AS child_status,
child.questions AS child_questions,
child.summary AS child_summary,
task.continuation_generation,
task.continuation_prompt,
task.infrastructure_retry_count,
task.kind,
task.last_error,
task.merge_position,
task.prompt,
task.research_report,
task.result_summary,
task.review_iteration,
task.status,
task.task_key,
task.touched_areas,
task.title,
task.verification_reason,
task.verification_verdict
FROM session_orchestration_task AS task
LEFT JOIN session AS child
ON child.id = task.child_session_id
WHERE task.session_orchestration_id = ?
ORDER BY task.merge_position, task.id
"#,
session_orchestration_id
)
.fetch_all(&self.0)
.await?;
for row in &rows {
status::validate_orchestration_task(&row.status)?;
row.kind
.parse::<ag_session::OrchestrationTaskKind>()
.map_err(|_| DbError::InvalidData {
entity: "orchestration task kind",
reason: format!("unknown persisted kind `{}`", row.kind),
})?;
if let Some(child_status) = &row.child_status {
status::validate_session(child_status)?;
}
}
Ok(rows)
}
async fn load_orchestration_integration_approach(&self, id: i64) -> Result<String, DbError> {
sqlx::query_scalar::<_, String>(
"SELECT integration_approach FROM session_orchestration WHERE id = ?",
)
.bind(id)
.fetch_one(&self.0)
.await
.db_context("load orchestration integration approach")
}
async fn load_orchestration_task_scope_for_child(
&self,
child_session_id: &str,
) -> Result<Option<OrchestrationTaskScopeRow>, DbError> {
sqlx::query_as!(
OrchestrationTaskScopeRow,
r#"
SELECT child.base_branch,
task.id AS "id!: i64",
task.touched_areas
FROM session_orchestration_task AS task
INNER JOIN session AS child
ON child.id = task.child_session_id
WHERE child.id = ?
AND child.role = 'OrchestrationWorker'
AND task.kind = 'Implementation'
"#,
child_session_id
)
.fetch_optional(&self.0)
.await
.db_context("load orchestration task scope for child")
}
async fn load_child_session_id_for_task(
&self,
task_id: i64,
) -> Result<Option<String>, DbError> {
let row = sqlx::query!(
r#"
SELECT id AS "id!: String"
FROM session
WHERE orchestration_task_id = ?
"#,
task_id
)
.fetch_optional(&self.0)
.await?;
Ok(row.map(|row| row.id))
}
async fn begin_orchestration_cancellation(&self, id: i64) -> Result<bool, DbError> {
let now = self.now();
let result = sqlx::query!(
r"
UPDATE session_orchestration
SET status = 'Canceling',
updated_at = ?
WHERE id = ?
AND status IN (
'AwaitingApproval',
'Running',
'Verifying',
'AwaitingIntegration',
'Integrating',
'Canceling'
)
",
now,
id
)
.execute(&self.0)
.await
.db_context("begin orchestration cancellation")?;
Ok(result.rows_affected() == 1)
}
async fn claim_orchestration_task(&self, id: i64) -> Result<bool, DbError> {
let now = self.now();
let result = sqlx::query!(
r"
UPDATE session_orchestration_task
SET status = 'Creating',
last_error = NULL,
updated_at = ?
WHERE id = ?
AND status = 'Planned'
AND EXISTS (
SELECT 1
FROM session_orchestration
WHERE session_orchestration.id = session_orchestration_task.session_orchestration_id
AND session_orchestration.status = 'Running'
)
",
now,
id
)
.execute(&self.0)
.await
.db_context("claim orchestration task")?;
Ok(result.rows_affected() == 1)
}
async fn claim_orchestration_review_application(
&self,
id: i64,
prompt: &str,
iteration_limit: i64,
) -> Result<bool, DbError> {
let now = self.now();
let mut transaction = self
.0
.begin()
.await
.db_context("claim orchestration review application")?;
let claim = sqlx::query!(
r"
UPDATE session_orchestration_task
SET continuation_generation = continuation_generation + 1,
continuation_prompt = ?,
review_iteration = review_iteration + 1,
status = 'ReviewApplying',
updated_at = ?
WHERE id = ?
AND status = 'Reviewing'
AND review_iteration < ?
AND child_session_id IS NOT NULL
",
prompt,
now,
id,
iteration_limit
)
.execute(&mut *transaction)
.await
.db_context("claim orchestration review application")?;
if claim.rows_affected() == 0 {
transaction
.rollback()
.await
.db_context("claim orchestration review application")?;
return Ok(false);
}
sqlx::query!(
r"
UPDATE session
SET focused_review_status = NULL,
focused_review_diff_hash = NULL,
focused_review_text = NULL,
updated_at = ?
WHERE id = (
SELECT child_session_id
FROM session_orchestration_task
WHERE id = ?
)
",
now,
id
)
.execute(&mut *transaction)
.await
.db_context("claim orchestration review application")?;
transaction
.commit()
.await
.db_context("claim orchestration review application")?;
Ok(true)
}
async fn claim_orchestration_rollup(&self, id: i64) -> Result<bool, DbError> {
let now = self.now();
let result = sqlx::query!(
r"
UPDATE session_orchestration
SET status = 'Verifying',
verification_generation = verification_generation + 1,
updated_at = ?
WHERE id = ?
AND status = 'Running'
",
now,
id
)
.execute(&self.0)
.await
.db_context("claim orchestration rollup")?;
Ok(result.rows_affected() == 1)
}
async fn complete_orchestration_rollup(&self, id: i64) -> Result<bool, DbError> {
self.transition_orchestration_and_tasks(
id,
OrchestrationTransition {
context: "complete orchestration rollup",
from_orchestration_status: "Verifying",
from_task_status: "Ready",
require_pass_verdict: true,
to_orchestration_status: "AwaitingIntegration",
to_task_status: "AwaitingIntegration",
},
)
.await
}
async fn record_orchestration_verdict(
&self,
id: i64,
task_key: &str,
is_pass: bool,
reason: &str,
) -> Result<bool, DbError> {
let now = self.now();
let verdict = if is_pass { "Pass" } else { "Flag" };
let result = sqlx::query!(
r"
UPDATE session_orchestration_task
SET verification_reason = ?,
verification_verdict = ?,
updated_at = ?
WHERE session_orchestration_id = ?
AND task_key = ?
AND status IN ('Ready', 'Reported')
AND EXISTS (
SELECT 1
FROM session_orchestration
WHERE id = ?
AND status = 'Verifying'
)
",
reason,
verdict,
now,
id,
task_key,
id
)
.execute(&self.0)
.await
.db_context("record orchestration verdict")?;
Ok(result.rows_affected() == 1)
}
async fn complete_orchestration_campaign(&self, id: i64) -> Result<bool, DbError> {
let now = self.now();
let mut transaction = self
.0
.begin()
.await
.db_context("complete orchestration campaign")?;
let orchestration = sqlx::query!(
r#"
UPDATE session_orchestration
SET status = 'Done',
updated_at = ?
WHERE id = ?
AND status IN ('AwaitingIntegration', 'Integrating')
RETURNING controller_session_id AS "controller_session_id!: String"
"#,
now,
id
)
.fetch_optional(&mut *transaction)
.await
.db_context("complete orchestration campaign")?;
if let Some(orchestration) = &orchestration {
sqlx::query!(
r"
UPDATE session
SET status = 'Done',
questions = '',
updated_at = ?
WHERE id = ?
AND role = 'Orchestrator'
AND status IN ('Review', 'Question')
",
now,
orchestration.controller_session_id
)
.execute(&mut *transaction)
.await
.db_context("complete orchestration campaign")?;
}
transaction
.commit()
.await
.db_context("complete orchestration campaign")?;
Ok(orchestration.is_some())
}
async fn load_rollup_operation_status(
&self,
operation_id: &str,
) -> Result<Option<String>, DbError> {
let row = sqlx::query!(
r#"
SELECT status AS "status!: String"
FROM session_operation
WHERE id = ?
"#,
operation_id
)
.fetch_optional(&self.0)
.await?;
if let Some(row) = &row {
status::validate_operation(&row.status)?;
}
Ok(row.map(|row| row.status))
}
async fn update_orchestration_status(&self, id: i64, status: &str) -> Result<(), DbError> {
status::validate_orchestration(status)?;
let now = self.now();
sqlx::query!(
r"
UPDATE session_orchestration
SET status = ?,
updated_at = ?
WHERE id = ?
",
status,
now,
id
)
.execute(&self.0)
.await?;
Ok(())
}
async fn approve_orchestration_plan(&self, id: i64) -> Result<bool, DbError> {
self.transition_orchestration_and_tasks(
id,
OrchestrationTransition {
context: "approve orchestration plan",
from_orchestration_status: "AwaitingApproval",
from_task_status: "Proposed",
require_pass_verdict: false,
to_orchestration_status: "Running",
to_task_status: "Planned",
},
)
.await
}
async fn approve_orchestration_integration(
&self,
id: i64,
approach: IntegrationApproach,
) -> Result<bool, DbError> {
let approach = approach.to_string();
let now = self.now();
let result = sqlx::query(
r"
UPDATE session_orchestration
SET integration_approach = ?,
status = 'Integrating',
updated_at = ?
WHERE id = ?
AND status = 'AwaitingIntegration'
",
)
.bind(approach)
.bind(now)
.bind(id)
.execute(&self.0)
.await
.db_context("approve orchestration integration")?;
Ok(result.rows_affected() == 1)
}
async fn update_orchestration_plan(
&self,
id: i64,
goal_statement: &str,
max_parallelism: i64,
) -> Result<(), DbError> {
let now = self.now();
sqlx::query!(
r"
UPDATE session_orchestration
SET goal_statement = ?,
max_parallelism = ?,
updated_at = ?
WHERE id = ?
AND status = 'AwaitingApproval'
",
goal_statement,
max_parallelism,
now,
id
)
.execute(&self.0)
.await?;
Ok(())
}
async fn queue_orchestration_continuation(
&self,
id: i64,
prompt: &str,
acceptance_criteria: &str,
touched_areas: &str,
) -> Result<bool, DbError> {
let now = self.now();
let mut transaction = self
.0
.begin()
.await
.db_context("queue orchestration continuation")?;
let result = sqlx::query!(
r"
UPDATE session_orchestration_task
SET acceptance_criteria = ?,
area_violations = '[]',
areas_compliant = NULL,
continuation_generation = continuation_generation + 1,
continuation_prompt = ?,
review_iteration = 0,
status = 'ContinuationPending',
touched_areas = ?,
result_summary = NULL,
verification_verdict = NULL,
verification_reason = NULL,
last_error = NULL,
updated_at = ?
WHERE id = ?
AND child_session_id IS NOT NULL
AND status IN ('Ready', 'AwaitingIntegration', 'IntegrationFailed')
",
acceptance_criteria,
prompt,
touched_areas,
now,
id
)
.execute(&mut *transaction)
.await?;
if result.rows_affected() == 1 {
sqlx::query!(
r"
UPDATE session
SET focused_review_status = NULL,
focused_review_diff_hash = NULL,
focused_review_text = NULL,
updated_at = ?
WHERE id = (
SELECT child_session_id
FROM session_orchestration_task
WHERE id = ?
)
",
now,
id
)
.execute(&mut *transaction)
.await?;
}
transaction.commit().await?;
Ok(result.rows_affected() == 1)
}
async fn reset_orchestration_verification(&self, id: i64) -> Result<(), DbError> {
let now = self.now();
sqlx::query!(
r"
UPDATE session_orchestration_task
SET status = 'Ready',
verification_reason = NULL,
verification_verdict = NULL,
updated_at = ?
WHERE session_orchestration_id = ?
AND status = 'AwaitingIntegration'
",
now,
id
)
.execute(&self.0)
.await
.db_context("reset orchestration verification")?;
Ok(())
}
async fn record_orchestration_spawn_failure(
&self,
id: i64,
error: &str,
retry_limit: i64,
) -> Result<String, DbError> {
let now = self.now();
let mut transaction = self
.0
.begin()
.await
.db_context("record orchestration spawn failure")?;
sqlx::query!(
r"
UPDATE session
SET orchestration_task_id = NULL,
updated_at = ?
WHERE orchestration_task_id = ?
",
now,
id
)
.execute(&mut *transaction)
.await
.db_context("record orchestration spawn failure")?;
let row = sqlx::query!(
r#"
UPDATE session_orchestration_task
SET child_session_id = NULL,
infrastructure_retry_count = infrastructure_retry_count + 1,
status = CASE
WHEN infrastructure_retry_count < ? THEN 'Planned'
ELSE 'Failed'
END,
last_error = ?,
updated_at = ?
WHERE id = ?
RETURNING status AS "status!: String"
"#,
retry_limit,
error,
now,
id
)
.fetch_one(&mut *transaction)
.await
.db_context("record orchestration spawn failure")?;
transaction
.commit()
.await
.db_context("record orchestration spawn failure")?;
status::validate_orchestration_task(&row.status)?;
Ok(row.status)
}
async fn detach_orchestration_child(&self, child_session_id: &str) -> Result<bool, DbError> {
let now = self.now();
let mut transaction = self
.0
.begin()
.await
.db_context("detach orchestration child")?;
let task = sqlx::query!(
r#"
SELECT orchestration_task_id AS "task_id!: i64"
FROM session
WHERE id = ?
AND role = 'OrchestrationWorker'
AND orchestration_task_id IS NOT NULL
"#,
child_session_id
)
.fetch_optional(&mut *transaction)
.await
.db_context("detach orchestration child")?;
let Some(task) = task else {
transaction
.commit()
.await
.db_context("detach orchestration child")?;
return Ok(false);
};
sqlx::query!(
r"
UPDATE session_orchestration_task
SET child_session_id = NULL,
status = 'Detached',
updated_at = ?
WHERE id = ?
",
now,
task.task_id
)
.execute(&mut *transaction)
.await
.db_context("detach orchestration child")?;
sqlx::query!(
r"
UPDATE session
SET orchestration_task_id = NULL,
role = 'Worker',
updated_at = ?
WHERE id = ?
",
now,
child_session_id
)
.execute(&mut *transaction)
.await
.db_context("detach orchestration child")?;
transaction
.commit()
.await
.db_context("detach orchestration child")?;
Ok(true)
}
async fn surface_orchestration_questions(
&self,
session_orchestration_id: i64,
task_id: i64,
questions: &str,
) -> Result<bool, DbError> {
let now = self.now();
let mut transaction = self
.0
.begin()
.await
.db_context("surface orchestration questions")?;
let claim = sqlx::query!(
r"
UPDATE session_orchestration
SET relayed_question_task_id = ?,
updated_at = ?
WHERE id = ?
AND relayed_question_task_id IS NULL
AND EXISTS (
SELECT 1
FROM session_orchestration_task AS task
WHERE task.id = ?
AND task.session_orchestration_id = session_orchestration.id
AND task.status = 'WaitingForInput'
AND task.child_session_id IS NOT NULL
)
AND EXISTS (
SELECT 1
FROM session AS controller
WHERE controller.id = session_orchestration.controller_session_id
AND controller.role = 'Orchestrator'
AND controller.status IN ('Review', 'Question')
AND COALESCE(controller.questions, '') = ''
)
",
task_id,
now,
session_orchestration_id,
task_id
)
.execute(&mut *transaction)
.await
.db_context("surface orchestration questions")?;
if claim.rows_affected() == 0 {
transaction
.rollback()
.await
.db_context("surface orchestration questions")?;
return Ok(false);
}
sqlx::query!(
r"
UPDATE session
SET questions = ?,
status = 'Question',
updated_at = ?
WHERE id = (
SELECT controller_session_id
FROM session_orchestration
WHERE id = ?
)
",
questions,
now,
session_orchestration_id
)
.execute(&mut *transaction)
.await
.db_context("surface orchestration questions")?;
transaction
.commit()
.await
.db_context("surface orchestration questions")?;
Ok(true)
}
async fn clear_orchestration_questions(
&self,
session_orchestration_id: i64,
) -> Result<(), DbError> {
let now = self.now();
let mut transaction = self
.0
.begin()
.await
.db_context("clear orchestration questions")?;
sqlx::query!(
r"
UPDATE session
SET questions = '',
status = 'Review',
updated_at = ?
WHERE id = (
SELECT controller_session_id
FROM session_orchestration
WHERE id = ?
AND relayed_question_task_id IS NOT NULL
)
AND role = 'Orchestrator'
AND status = 'Question'
",
now,
session_orchestration_id
)
.execute(&mut *transaction)
.await
.db_context("clear orchestration questions")?;
sqlx::query!(
r"
UPDATE session_orchestration
SET relayed_question_task_id = NULL,
updated_at = ?
WHERE id = ?
AND relayed_question_task_id IS NOT NULL
",
now,
session_orchestration_id
)
.execute(&mut *transaction)
.await
.db_context("clear orchestration questions")?;
transaction
.commit()
.await
.db_context("clear orchestration questions")?;
Ok(())
}
async fn link_orchestration_task_child(
&self,
id: i64,
child_session_id: &str,
) -> Result<bool, DbError> {
let now = self.now();
let result = sqlx::query!(
r"
UPDATE session_orchestration_task
SET child_session_id = ?,
status = 'Running',
attempt_count = attempt_count + 1,
updated_at = ?
WHERE id = ?
AND status = 'Creating'
AND EXISTS (
SELECT 1
FROM session_orchestration
WHERE session_orchestration.id = session_orchestration_task.session_orchestration_id
AND session_orchestration.status = 'Running'
)
",
child_session_id,
now,
id
)
.execute(&self.0)
.await?;
Ok(result.rows_affected() == 1)
}
async fn update_orchestration_task_status(
&self,
id: i64,
status: &str,
last_error: Option<String>,
) -> Result<(), DbError> {
status::validate_orchestration_task(status)?;
let now = self.now();
sqlx::query!(
r"
UPDATE session_orchestration_task
SET status = ?,
last_error = ?,
updated_at = ?
WHERE id = ?
",
status,
last_error,
now,
id
)
.execute(&self.0)
.await?;
Ok(())
}
async fn update_orchestration_task_result_summary(
&self,
id: i64,
result_summary: &str,
) -> Result<(), DbError> {
let now = self.now();
sqlx::query!(
r"
UPDATE session_orchestration_task
SET result_summary = ?,
updated_at = ?
WHERE id = ?
",
result_summary,
now,
id
)
.execute(&self.0)
.await?;
Ok(())
}
async fn update_orchestration_task_research_report(
&self,
id: i64,
research_report: &str,
) -> Result<(), DbError> {
let now = self.now();
sqlx::query!(
r"
UPDATE session_orchestration_task
SET research_report = ?,
updated_at = ?
WHERE id = ?
AND kind = 'Research'
",
research_report,
now,
id
)
.execute(&self.0)
.await?;
Ok(())
}
async fn update_orchestration_task_area_compliance(
&self,
id: i64,
areas_compliant: Option<bool>,
area_violations: &str,
) -> Result<(), DbError> {
let now = self.now();
sqlx::query!(
r"
UPDATE session_orchestration_task
SET area_violations = ?,
areas_compliant = ?,
updated_at = ?
WHERE id = ?
",
area_violations,
areas_compliant,
now,
id
)
.execute(&self.0)
.await?;
Ok(())
}
}
#[cfg(test)]
mod tests {
use ag_agent::{AgentKind, ReasoningLevel, SpeedMode};
use ag_session::{
FocusedReviewStatus, OrchestrationStatus, OrchestrationTaskKind, OrchestrationTaskStatus,
SessionMessageKind,
};
use super::*;
use crate::{AppRepositories, PersistedSessionCreation, SessionRow};
async fn controller_fixture() -> AppRepositories {
controller_fixture_with_pool().await.0
}
async fn controller_fixture_with_pool() -> (AppRepositories, SqlitePool) {
let (database, pool) = AppRepositories::in_memory_with_pool()
.await
.expect("db should open");
let project_id = database
.projects()
.upsert_project("/tmp/project", None)
.await
.expect("failed to upsert project");
database
.sessions()
.insert_session_with_agent(PersistedSessionCreation {
agent: "codex",
base_branch: "main",
id: "controller",
is_draft: false,
model: AgentKind::Codex.default_model().as_str(),
orchestration_task_id: None,
parent_session_id: None,
personality_id: None,
project_id,
reasoning_level: ReasoningLevel::default(),
role: Some("Orchestrator"),
speed_mode: SpeedMode::Normal,
status: "Review",
})
.await
.expect("failed to insert controller session");
(database, pool)
}
#[tokio::test]
async fn hydration_rejects_unknown_orchestration_status() {
let (database, pool) = AppRepositories::in_memory_with_pool()
.await
.expect("db should open");
let project_id = database
.projects()
.upsert_project("/tmp/invalid-orchestration", None)
.await
.expect("failed to upsert project");
database
.sessions()
.insert_session("controller", "gpt-5.6-sol", "main", "Draft", project_id)
.await
.expect("failed to insert controller session");
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", "Running", 1)
.await
.expect("failed to insert orchestration");
sqlx::query("UPDATE session_orchestration SET status = 'Unknown' WHERE id = ?")
.bind(orchestration_id)
.execute(&pool)
.await
.expect("failed to corrupt orchestration status");
let error = database
.orchestrations()
.load_orchestration_for_controller("controller")
.await
.expect_err("invalid status should fail hydration");
assert!(matches!(
error,
DbError::InvalidStatus {
entity: "orchestration",
value,
} if value == "Unknown"
));
}
#[tokio::test]
async fn rollup_recovery_failures_report_semantic_operation_context() {
let (database, pool) = AppRepositories::in_memory_with_pool()
.await
.expect("db should open");
sqlx::query("DROP TABLE session_orchestration")
.execute(&pool)
.await
.expect("failed to drop orchestration table");
let claim_error = database
.orchestrations()
.claim_orchestration_rollup(1)
.await
.expect_err("claim should fail");
let completion_error = database
.orchestrations()
.complete_orchestration_rollup(1)
.await
.expect_err("completion should fail");
assert!(matches!(
claim_error,
DbError::QueryContext {
operation: "claim orchestration rollup",
..
}
));
assert!(matches!(
completion_error,
DbError::QueryContext {
operation: "complete orchestration rollup",
..
}
));
}
fn planned_task(session_orchestration_id: i64, task_key: &str) -> PersistedOrchestrationTask {
PersistedOrchestrationTask {
acceptance_criteria: format!(r#"["Complete {task_key}"]"#),
kind: "Implementation".to_string(),
merge_position: 0,
prompt: format!("Complete {task_key}"),
session_orchestration_id,
task_key: task_key.to_string(),
title: format!("Task {task_key}"),
touched_areas: format!(r#"["crates/{task_key}/"]"#),
}
}
async fn insert_orchestration_child(
database: &AppRepositories,
project_id: i64,
session_id: &str,
task_id: i64,
) {
database
.sessions()
.insert_session_with_agent(PersistedSessionCreation {
agent: "codex",
base_branch: "main",
id: session_id,
is_draft: false,
model: AgentKind::Codex.default_model().as_str(),
orchestration_task_id: Some(task_id),
parent_session_id: None,
personality_id: None,
project_id,
reasoning_level: ReasoningLevel::default(),
role: Some("OrchestrationWorker"),
speed_mode: SpeedMode::Normal,
status: "Review",
})
.await
.expect("failed to insert orchestration child");
}
async fn insert_waiting_orchestration_task(
database: &AppRepositories,
project_id: i64,
session_orchestration_id: i64,
task_key: &str,
child_session_id: &str,
) -> i64 {
let task_id = database
.orchestrations()
.upsert_orchestration_task(planned_task(session_orchestration_id, task_key))
.await
.expect("failed to insert waiting task");
assert!(
database
.orchestrations()
.claim_orchestration_task(task_id)
.await
.expect("failed to claim waiting task")
);
insert_orchestration_child(database, project_id, child_session_id, task_id).await;
assert!(
database
.orchestrations()
.link_orchestration_task_child(task_id, child_session_id)
.await
.expect("failed to link waiting child")
);
database
.orchestrations()
.update_orchestration_task_status(
task_id,
&OrchestrationTaskStatus::WaitingForInput.to_string(),
None,
)
.await
.expect("failed to wait for child question");
task_id
}
async fn surface_and_clear_orchestration_questions(
database: &AppRepositories,
session_orchestration_id: i64,
task_id: i64,
questions: &str,
) -> bool {
let surfaced = database
.orchestrations()
.surface_orchestration_questions(session_orchestration_id, task_id, questions)
.await
.expect("failed to surface child question");
database
.orchestrations()
.clear_orchestration_questions(session_orchestration_id)
.await
.expect("failed to clear child question");
surfaced
}
async fn load_detached_campaign_state(
database: &AppRepositories,
orchestration_id: i64,
) -> (SessionOrchestrationTaskRow, SessionRow, SessionRow) {
let task = database
.orchestrations()
.load_orchestration_tasks(orchestration_id)
.await
.expect("failed to load detached task")
.remove(0);
let child = database
.sessions()
.load_session("child-alpha")
.await
.expect("failed to load detached child")
.expect("detached child should exist");
let controller = database
.sessions()
.load_session("controller")
.await
.expect("failed to load controller")
.expect("controller should exist");
(task, child, controller)
}
#[tokio::test]
async fn test_planned_orchestration_round_trips_before_any_child_exists() {
let database = controller_fixture().await;
let orchestration_id = database
.orchestrations()
.insert_orchestration(
"controller",
&OrchestrationStatus::AwaitingApproval.to_string(),
3,
)
.await
.expect("failed to insert orchestration");
database
.orchestrations()
.upsert_orchestration_task(planned_task(orchestration_id, "alpha"))
.await
.expect("failed to insert task");
let orchestration = database
.orchestrations()
.load_orchestration_for_controller("controller")
.await
.expect("failed to load orchestration")
.expect("orchestration should exist");
let tasks = database
.orchestrations()
.load_orchestration_tasks(orchestration_id)
.await
.expect("failed to load tasks");
assert_eq!(orchestration.controller_project_id, 1);
assert_eq!(orchestration.controller_session_id, "controller");
assert_eq!(orchestration.max_parallelism, 3);
assert_eq!(
orchestration.status,
OrchestrationStatus::AwaitingApproval.to_string()
);
assert_eq!(tasks.len(), 1);
assert_eq!(tasks[0].task_key, "alpha");
assert_eq!(tasks[0].child_session_id, None);
assert_eq!(tasks[0].attempt_count, 0);
assert_eq!(
tasks[0].kind,
OrchestrationTaskKind::Implementation.to_string()
);
assert_eq!(tasks[0].touched_areas, r#"["crates/alpha/"]"#);
assert!(tasks[0].child_answer.is_none());
assert!(tasks[0].research_report.is_none());
}
#[tokio::test]
async fn research_task_round_trips_report_and_latest_child_answer_without_scope() {
let database = controller_fixture().await;
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert orchestration");
let task_id = database
.orchestrations()
.upsert_orchestration_task(PersistedOrchestrationTask {
kind: OrchestrationTaskKind::Research.to_string(),
touched_areas: "[]".to_string(),
..planned_task(orchestration_id, "architecture")
})
.await
.expect("failed to insert research task");
assert!(
database
.orchestrations()
.claim_orchestration_task(task_id)
.await
.expect("failed to claim research task")
);
database
.sessions()
.insert_session_with_agent(PersistedSessionCreation {
agent: "codex",
base_branch: "main",
id: "research-child",
is_draft: false,
model: AgentKind::Codex.default_model().as_str(),
orchestration_task_id: Some(task_id),
parent_session_id: None,
personality_id: None,
project_id: 1,
reasoning_level: ReasoningLevel::default(),
role: Some("OrchestrationResearcher"),
speed_mode: SpeedMode::Normal,
status: "Review",
})
.await
.expect("failed to insert research child");
assert!(
database
.orchestrations()
.link_orchestration_task_child(task_id, "research-child")
.await
.expect("failed to link research child")
);
database
.sessions()
.append_session_message(
"research-child",
SessionMessageKind::AssistantAnswer,
"Full architecture report",
)
.await
.expect("failed to persist research answer");
database
.orchestrations()
.update_orchestration_task_research_report(task_id, "Full architecture report")
.await
.expect("failed to persist research report");
let task = database
.orchestrations()
.load_orchestration_tasks(orchestration_id)
.await
.expect("failed to load research task")
.remove(0);
let scope = database
.orchestrations()
.load_orchestration_task_scope_for_child("research-child")
.await
.expect("failed to inspect research child scope");
assert_eq!(task.kind, OrchestrationTaskKind::Research.to_string());
assert_eq!(
task.child_answer.as_deref(),
Some("Full architecture report")
);
assert_eq!(
task.research_report.as_deref(),
Some("Full architecture report")
);
assert!(scope.is_none());
}
#[tokio::test]
async fn orchestration_task_kind_validation_rejects_unknown_writes_and_hydration() {
let (database, pool) = controller_fixture_with_pool().await;
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert orchestration");
let invalid_task = PersistedOrchestrationTask {
kind: "Unknown".to_string(),
..planned_task(orchestration_id, "invalid-write")
};
let write_error = database
.orchestrations()
.upsert_orchestration_task(invalid_task)
.await
.expect_err("unknown task kind should be rejected before persistence");
sqlx::query(
"INSERT INTO session_orchestration_task (session_orchestration_id, task_key, title, \
prompt, status, kind) VALUES (?, 'invalid-read', 'Invalid', 'Inspect', 'Planned', \
'Unknown')",
)
.bind(orchestration_id)
.execute(&pool)
.await
.expect("failed to seed invalid persisted kind");
let read_error = database
.orchestrations()
.load_orchestration_tasks(orchestration_id)
.await
.expect_err("unknown persisted task kind should fail hydration");
assert!(matches!(
write_error,
DbError::InvalidData {
entity: "orchestration task kind",
reason,
} if reason == "unknown persisted kind `Unknown`"
));
assert!(matches!(
read_error,
DbError::InvalidData {
entity: "orchestration task kind",
reason,
} if reason == "unknown persisted kind `Unknown`"
));
}
#[tokio::test]
async fn test_retry_with_same_task_key_updates_the_existing_row() {
let database = controller_fixture().await;
let project_id = database
.projects()
.upsert_project("/tmp/project", None)
.await
.expect("failed to load project");
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert orchestration");
let first_id = database
.orchestrations()
.upsert_orchestration_task(planned_task(orchestration_id, "alpha"))
.await
.expect("failed to insert task");
let first_claim = database
.orchestrations()
.claim_orchestration_task(first_id)
.await
.expect("failed to claim first attempt");
insert_orchestration_child(&database, project_id, "child-old", first_id).await;
let first_link = database
.orchestrations()
.link_orchestration_task_child(first_id, "child-old")
.await
.expect("failed to link old child");
database
.orchestrations()
.update_orchestration_task_status(
first_id,
&OrchestrationTaskStatus::Failed.to_string(),
Some("agent crashed".to_string()),
)
.await
.expect("failed to fail task");
let retried_id = database
.orchestrations()
.upsert_orchestration_task(PersistedOrchestrationTask {
title: "Task alpha, retried".to_string(),
..planned_task(orchestration_id, "alpha")
})
.await
.expect("failed to retry task");
let detached_child = database
.orchestrations()
.load_child_session_id_for_task(first_id)
.await
.expect("failed to check detached child");
let retry_claim = database
.orchestrations()
.claim_orchestration_task(retried_id)
.await
.expect("failed to claim replacement attempt");
insert_orchestration_child(&database, project_id, "child-replacement", retried_id).await;
let replacement_link = database
.orchestrations()
.link_orchestration_task_child(retried_id, "child-replacement")
.await
.expect("failed to link replacement child");
let linked_child = database
.orchestrations()
.load_child_session_id_for_task(first_id)
.await
.expect("failed to check replacement child");
let tasks = database
.orchestrations()
.load_orchestration_tasks(orchestration_id)
.await
.expect("failed to load tasks");
assert!(first_claim);
assert!(first_link);
assert!(retry_claim);
assert!(replacement_link);
assert_eq!(retried_id, first_id);
assert_eq!(tasks.len(), 1);
assert_eq!(tasks[0].title, "Task alpha, retried");
assert_eq!(
tasks[0].status,
OrchestrationTaskStatus::Running.to_string()
);
assert_eq!(tasks[0].attempt_count, 2);
assert_eq!(
tasks[0].child_session_id.as_deref(),
Some("child-replacement")
);
assert_eq!(tasks[0].last_error, None);
assert_eq!(detached_child, None);
assert_eq!(linked_child.as_deref(), Some("child-replacement"));
}
#[tokio::test]
async fn test_child_linkage_counts_attempts_and_loads_observed_state() {
let database = controller_fixture().await;
let project_id = database
.projects()
.upsert_project("/tmp/project", None)
.await
.expect("failed to upsert project");
database
.sessions()
.insert_session("child-a", "gpt-5.6-sol", "main", "Draft", project_id)
.await
.expect("failed to insert child session");
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert orchestration");
let task_id = database
.orchestrations()
.upsert_orchestration_task(planned_task(orchestration_id, "alpha"))
.await
.expect("failed to insert task");
let claimed = database
.orchestrations()
.claim_orchestration_task(task_id)
.await
.expect("failed to claim task");
let linked = database
.orchestrations()
.link_orchestration_task_child(task_id, "child-a")
.await
.expect("failed to link child");
database
.orchestrations()
.update_orchestration_task_result_summary(task_id, "Added the parser")
.await
.expect("failed to record summary");
let task = database
.orchestrations()
.load_orchestration_tasks(orchestration_id)
.await
.expect("failed to load task")
.remove(0);
assert!(claimed);
assert!(linked);
assert_eq!(task.id, task_id);
assert_eq!(task.attempt_count, 1);
assert_eq!(task.child_session_id.as_deref(), Some("child-a"));
assert_eq!(task.result_summary.as_deref(), Some("Added the parser"));
assert_eq!(task.status, OrchestrationTaskStatus::Running.to_string());
assert_eq!(task.child_status.as_deref(), Some("Draft"));
assert_eq!(task.child_summary, None);
assert_eq!(task.child_input_tokens, 0);
assert_eq!(task.child_output_tokens, 0);
}
#[tokio::test]
async fn test_active_orchestration_load_includes_recoverable_states() {
let database = controller_fixture().await;
let running_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert running orchestration");
let submitting_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert submitting orchestration");
let canceling_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Canceling.to_string(), 2)
.await
.expect("failed to insert canceling orchestration");
let settled_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert settled orchestration");
database
.orchestrations()
.update_orchestration_status(settled_id, &OrchestrationStatus::Done.to_string())
.await
.expect("failed to settle orchestration");
let first_claim = database
.orchestrations()
.claim_orchestration_rollup(submitting_id)
.await
.expect("failed to claim roll-up");
let duplicate_claim = database
.orchestrations()
.claim_orchestration_rollup(submitting_id)
.await
.expect("failed to repeat roll-up claim");
let active = database
.orchestrations()
.load_active_orchestrations()
.await
.expect("failed to load active orchestrations");
assert!(first_claim);
assert!(!duplicate_claim);
assert_eq!(
active.iter().map(|row| row.id).collect::<Vec<_>>(),
vec![running_id, submitting_id, canceling_id]
);
assert_eq!(active[1].status, OrchestrationStatus::Verifying.to_string());
assert_eq!(active[1].verification_generation, 1);
assert_eq!(active[2].status, OrchestrationStatus::Canceling.to_string());
}
#[tokio::test]
async fn test_cancellation_barrier_blocks_fan_out_claims_and_links() {
let database = controller_fixture().await;
let project_id = database
.projects()
.upsert_project("/tmp/project", None)
.await
.expect("failed to reload project");
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert orchestration");
let late_task_id = database
.orchestrations()
.upsert_orchestration_task(planned_task(orchestration_id, "late"))
.await
.expect("failed to insert late task");
let claimed_task_id = database
.orchestrations()
.upsert_orchestration_task(planned_task(orchestration_id, "claimed"))
.await
.expect("failed to insert claimed task");
let initial_claim = database
.orchestrations()
.claim_orchestration_task(claimed_task_id)
.await
.expect("failed to claim task before cancellation");
insert_orchestration_child(&database, project_id, "child-claimed", claimed_task_id).await;
let cancellation_started = database
.orchestrations()
.begin_orchestration_cancellation(orchestration_id)
.await
.expect("failed to begin cancellation");
let cancellation_retried = database
.orchestrations()
.begin_orchestration_cancellation(orchestration_id)
.await
.expect("failed to retry cancellation");
let late_claim = database
.orchestrations()
.claim_orchestration_task(late_task_id)
.await
.expect("failed to inspect late claim");
let late_link = database
.orchestrations()
.link_orchestration_task_child(claimed_task_id, "child-claimed")
.await
.expect("failed to inspect late link");
let rollup_completion = database
.orchestrations()
.complete_orchestration_rollup(orchestration_id)
.await
.expect("failed to inspect late roll-up completion");
let orchestration = database
.orchestrations()
.load_orchestration_for_controller("controller")
.await
.expect("failed to load orchestration")
.expect("orchestration should exist");
assert!(initial_claim);
assert!(cancellation_started);
assert!(cancellation_retried);
assert!(!late_claim);
assert!(!late_link);
assert!(!rollup_completion);
assert_eq!(
orchestration.status,
OrchestrationStatus::Canceling.to_string()
);
}
#[tokio::test]
async fn test_rollup_operation_status_controls_orchestration_completion() {
let database = controller_fixture().await;
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert orchestration");
let claimed = database
.orchestrations()
.claim_orchestration_rollup(orchestration_id)
.await
.expect("failed to claim roll-up");
assert!(claimed);
let operation_id = format!("orchestration-rollup-{orchestration_id}-1");
let missing_status = database
.orchestrations()
.load_rollup_operation_status(&operation_id)
.await
.expect("failed to load missing operation");
database
.operations()
.claim_session_operation(&operation_id, "controller", "reply")
.await
.expect("failed to claim operation");
let queued_status = database
.orchestrations()
.load_rollup_operation_status(&operation_id)
.await
.expect("failed to load queued operation");
database
.operations()
.mark_session_operation_running(&operation_id)
.await
.expect("failed to run operation");
let running_status = database
.orchestrations()
.load_rollup_operation_status(&operation_id)
.await
.expect("failed to load running operation");
database
.operations()
.mark_session_operation_failed(&operation_id, "turn failed")
.await
.expect("failed to fail operation");
let failed_status = database
.orchestrations()
.load_rollup_operation_status(&operation_id)
.await
.expect("failed to load failed operation");
database
.operations()
.claim_session_operation(&operation_id, "controller", "reply")
.await
.expect("failed to reclaim operation");
database
.operations()
.mark_session_operation_done(&operation_id)
.await
.expect("failed to complete operation");
let done_status = database
.orchestrations()
.load_rollup_operation_status(&operation_id)
.await
.expect("failed to load completed operation");
let completed = database
.orchestrations()
.complete_orchestration_rollup(orchestration_id)
.await
.expect("failed to complete orchestration");
let duplicate_completion = database
.orchestrations()
.complete_orchestration_rollup(orchestration_id)
.await
.expect("failed to inspect duplicate completion");
let orchestration = database
.orchestrations()
.load_orchestration_for_controller("controller")
.await
.expect("failed to load orchestration")
.expect("orchestration should exist");
assert_eq!(missing_status, None);
assert_eq!(queued_status.as_deref(), Some("queued"));
assert_eq!(running_status.as_deref(), Some("running"));
assert_eq!(failed_status.as_deref(), Some("failed"));
assert_eq!(done_status.as_deref(), Some("done"));
assert!(completed);
assert!(!duplicate_completion);
assert_eq!(
orchestration.status,
OrchestrationStatus::AwaitingIntegration.to_string()
);
}
#[tokio::test]
async fn test_session_metadata_for_project_loads_controller_and_children_together() {
let database = controller_fixture().await;
let project_id = database
.projects()
.upsert_project("/tmp/project", None)
.await
.expect("failed to reload project");
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert orchestration");
let running_task_id = database
.orchestrations()
.upsert_orchestration_task(planned_task(orchestration_id, "running"))
.await
.expect("failed to insert running task");
let waiting_task_id = database
.orchestrations()
.upsert_orchestration_task(planned_task(orchestration_id, "waiting"))
.await
.expect("failed to insert waiting task");
insert_orchestration_child(&database, project_id, "child-running", running_task_id).await;
insert_orchestration_child(&database, project_id, "child-waiting", waiting_task_id).await;
let running_claim = database
.orchestrations()
.claim_orchestration_task(running_task_id)
.await
.expect("failed to claim running task");
let waiting_claim = database
.orchestrations()
.claim_orchestration_task(waiting_task_id)
.await
.expect("failed to claim waiting task");
let running_link = database
.orchestrations()
.link_orchestration_task_child(running_task_id, "child-running")
.await
.expect("failed to link running child");
let waiting_link = database
.orchestrations()
.link_orchestration_task_child(waiting_task_id, "child-waiting")
.await
.expect("failed to link waiting child");
database
.orchestrations()
.update_orchestration_task_status(
waiting_task_id,
&OrchestrationTaskStatus::WaitingForInput.to_string(),
None,
)
.await
.expect("failed to mark waiting child");
let metadata = database
.orchestrations()
.load_session_metadata_for_project(project_id)
.await
.expect("failed to load bulk orchestration metadata");
assert!(running_claim);
assert!(waiting_claim);
assert!(running_link);
assert!(waiting_link);
assert_eq!(
metadata,
vec![
SessionOrchestrationMetadataRow {
controller_session_id: Some("controller".to_string()),
orchestration_status: None,
running_task_count: 0,
session_id: "child-running".to_string(),
waiting_task_count: 0,
},
SessionOrchestrationMetadataRow {
controller_session_id: Some("controller".to_string()),
orchestration_status: None,
running_task_count: 0,
session_id: "child-waiting".to_string(),
waiting_task_count: 0,
},
SessionOrchestrationMetadataRow {
controller_session_id: None,
orchestration_status: Some(OrchestrationStatus::Running.to_string()),
running_task_count: 1,
session_id: "controller".to_string(),
waiting_task_count: 1,
},
]
);
}
#[tokio::test]
async fn approval_persists_plan_and_releases_proposed_tasks() {
let database = controller_fixture().await;
let orchestration_id = database
.orchestrations()
.insert_orchestration(
"controller",
&OrchestrationStatus::AwaitingApproval.to_string(),
2,
)
.await
.expect("failed to insert orchestration");
let task_id = database
.orchestrations()
.upsert_orchestration_task(planned_task(orchestration_id, "follow-up"))
.await
.expect("failed to insert proposed task");
database
.orchestrations()
.update_orchestration_task_status(
task_id,
&OrchestrationTaskStatus::Proposed.to_string(),
None,
)
.await
.expect("failed to propose task");
database
.orchestrations()
.update_orchestration_plan(orchestration_id, "Ship the campaign", 4)
.await
.expect("failed to update campaign plan");
let approved = database
.orchestrations()
.approve_orchestration_plan(orchestration_id)
.await
.expect("failed to approve campaign");
let duplicate_approval = database
.orchestrations()
.approve_orchestration_plan(orchestration_id)
.await
.expect("failed to inspect duplicate approval");
let orchestration = database
.orchestrations()
.load_orchestration_for_controller("controller")
.await
.expect("failed to load campaign")
.expect("campaign should exist");
let task = database
.orchestrations()
.load_orchestration_tasks(orchestration_id)
.await
.expect("failed to load campaign tasks")
.remove(0);
assert!(approved);
assert!(!duplicate_approval);
assert_eq!(orchestration.goal_statement, "Ship the campaign");
assert_eq!(orchestration.max_parallelism, 4);
assert_eq!(
orchestration.status,
OrchestrationStatus::Running.to_string()
);
assert_eq!(task.status, OrchestrationTaskStatus::Planned.to_string());
}
#[tokio::test]
async fn integration_approval_persists_selected_approach_atomically() {
let database = controller_fixture().await;
let orchestration_id = database
.orchestrations()
.insert_orchestration(
"controller",
&OrchestrationStatus::AwaitingIntegration.to_string(),
2,
)
.await
.expect("failed to insert orchestration");
let approved = database
.orchestrations()
.approve_orchestration_integration(orchestration_id, IntegrationApproach::ReviewRequest)
.await
.expect("failed to approve integration");
let duplicate_approval = database
.orchestrations()
.approve_orchestration_integration(orchestration_id, IntegrationApproach::LocalMerge)
.await
.expect("failed to inspect duplicate approval");
let approach = database
.orchestrations()
.load_orchestration_integration_approach(orchestration_id)
.await
.expect("failed to load integration approach");
let orchestration = database
.orchestrations()
.load_orchestration_for_controller("controller")
.await
.expect("failed to load orchestration")
.expect("orchestration should exist");
assert!(approved);
assert!(!duplicate_approval);
assert_eq!(approach, IntegrationApproach::ReviewRequest.to_string());
assert_eq!(
orchestration.status,
OrchestrationStatus::Integrating.to_string()
);
}
#[tokio::test]
async fn recoverable_focused_reviews_exclude_outstanding_continuations() {
let database = controller_fixture().await;
let project_id = database
.projects()
.upsert_project("/tmp/project", None)
.await
.expect("failed to reload project");
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert orchestration");
let task_id = database
.orchestrations()
.upsert_orchestration_task(planned_task(orchestration_id, "recoverable"))
.await
.expect("failed to insert recoverable task");
assert!(
database
.orchestrations()
.claim_orchestration_task(task_id)
.await
.expect("failed to claim recoverable task")
);
insert_orchestration_child(&database, project_id, "child-recoverable", task_id).await;
assert!(
database
.orchestrations()
.link_orchestration_task_child(task_id, "child-recoverable")
.await
.expect("failed to link recoverable child")
);
database
.orchestrations()
.update_orchestration_task_status(
task_id,
&OrchestrationTaskStatus::Reviewing.to_string(),
None,
)
.await
.expect("failed to mark child reviewing");
let incomplete = database
.orchestrations()
.load_recoverable_focused_review_session_ids(project_id)
.await
.expect("failed to load incomplete review");
database
.orchestrations()
.update_orchestration_task_status(
task_id,
&OrchestrationTaskStatus::ContinuationPending.to_string(),
None,
)
.await
.expect("failed to mark continuation pending");
let pending = database
.orchestrations()
.load_recoverable_focused_review_session_ids(project_id)
.await
.expect("failed to inspect pending continuation recovery");
database
.orchestrations()
.update_orchestration_task_status(
task_id,
&OrchestrationTaskStatus::ReviewApplying.to_string(),
None,
)
.await
.expect("failed to mark review application pending");
let applying = database
.orchestrations()
.load_recoverable_focused_review_session_ids(project_id)
.await
.expect("failed to inspect review application recovery");
database
.orchestrations()
.update_orchestration_task_status(
task_id,
&OrchestrationTaskStatus::Reviewing.to_string(),
None,
)
.await
.expect("failed to restore reviewing state");
database
.sessions()
.update_session_focused_review(
"child-recoverable",
Some(FocusedReviewStatus::Ready),
Some("42".to_string()),
Some("### Suggestions\n\n- None".to_string()),
)
.await
.expect("failed to complete focused review");
let completed = database
.orchestrations()
.load_recoverable_focused_review_session_ids(project_id)
.await
.expect("failed to reload completed review");
assert_eq!(incomplete, ["child-recoverable"]);
assert!(pending.is_empty() && applying.is_empty() && completed.is_empty());
}
#[tokio::test]
async fn review_application_claim_is_bounded_and_clears_consumed_review() {
let database = controller_fixture().await;
let project_id = database
.projects()
.upsert_project("/tmp/project", None)
.await
.expect("failed to reload project");
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert orchestration");
let task_id = database
.orchestrations()
.upsert_orchestration_task(planned_task(orchestration_id, "reviewed"))
.await
.expect("failed to insert reviewed task");
assert!(
database
.orchestrations()
.claim_orchestration_task(task_id)
.await
.expect("failed to claim reviewed task")
);
insert_orchestration_child(&database, project_id, "child-reviewed", task_id).await;
assert!(
database
.orchestrations()
.link_orchestration_task_child(task_id, "child-reviewed")
.await
.expect("failed to link reviewed child")
);
database
.orchestrations()
.update_orchestration_task_status(
task_id,
&OrchestrationTaskStatus::Reviewing.to_string(),
None,
)
.await
.expect("failed to mark child reviewing");
database
.sessions()
.update_session_focused_review(
"child-reviewed",
Some(FocusedReviewStatus::Ready),
Some("42".to_string()),
Some("### Suggestions\n\n- Fix it".to_string()),
)
.await
.expect("failed to seed focused review");
let claimed = database
.orchestrations()
.claim_orchestration_review_application(task_id, "Verify then apply", 3)
.await
.expect("failed to claim review application");
let duplicate_claim = database
.orchestrations()
.claim_orchestration_review_application(task_id, "Duplicate", 3)
.await
.expect("failed to inspect duplicate review application");
let task = database
.orchestrations()
.load_orchestration_tasks(orchestration_id)
.await
.expect("failed to load reviewed task")
.remove(0);
let review_cache = database
.sessions()
.load_session_focused_reviews_for_project(project_id)
.await
.expect("failed to load consumed review cache");
assert!(claimed);
assert!(!duplicate_claim);
assert_eq!(
task.status,
OrchestrationTaskStatus::ReviewApplying.to_string()
);
assert_eq!(task.continuation_generation, 1);
assert_eq!(
task.continuation_prompt.as_deref(),
Some("Verify then apply")
);
assert_eq!(task.review_iteration, 1);
assert_eq!(review_cache, [] as [crate::SessionFocusedReviewRow; 0]);
assert_eq!(task.child_focused_review_status, None);
assert_eq!(task.child_focused_review_text, None);
}
#[tokio::test]
async fn managed_child_continuation_questions_and_detach_are_durable() {
let database = controller_fixture().await;
let project_id = database
.projects()
.upsert_project("/tmp/project", None)
.await
.expect("failed to reload project");
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert orchestration");
let task_id = database
.orchestrations()
.upsert_orchestration_task(planned_task(orchestration_id, "alpha"))
.await
.expect("failed to insert task");
assert!(
database
.orchestrations()
.claim_orchestration_task(task_id)
.await
.expect("failed to claim task")
);
insert_orchestration_child(&database, project_id, "child-alpha", task_id).await;
assert!(
database
.orchestrations()
.link_orchestration_task_child(task_id, "child-alpha")
.await
.expect("failed to link child")
);
database
.orchestrations()
.update_orchestration_task_status(
task_id,
&OrchestrationTaskStatus::WaitingForInput.to_string(),
None,
)
.await
.expect("failed to wait for child question");
let surfaced = surface_and_clear_orchestration_questions(
&database,
orchestration_id,
task_id,
r#"[{"text":"Choose one"}]"#,
)
.await;
database
.orchestrations()
.update_orchestration_task_status(
task_id,
&OrchestrationTaskStatus::Ready.to_string(),
None,
)
.await
.expect("failed to settle child after its question");
let queued = database
.orchestrations()
.queue_orchestration_continuation(
task_id,
"Add the missing edge case",
r#"["The edge case is tested"]"#,
r#"["docs/"]"#,
)
.await
.expect("failed to queue continuation");
let duplicate_queue = database
.orchestrations()
.queue_orchestration_continuation(
task_id,
"Duplicate",
r#"["Duplicate"]"#,
r#"["ignored/"]"#,
)
.await
.expect("failed to inspect duplicate continuation");
let detached = database
.orchestrations()
.detach_orchestration_child("child-alpha")
.await
.expect("failed to detach child");
let duplicate_detach = database
.orchestrations()
.detach_orchestration_child("child-alpha")
.await
.expect("failed to inspect duplicate detach");
let (task, child, controller) =
load_detached_campaign_state(&database, orchestration_id).await;
let continuation_prompt = task.continuation_prompt.as_deref();
assert!(queued && surfaced && detached);
assert!(!duplicate_queue && !duplicate_detach);
assert_eq!(task.status, OrchestrationTaskStatus::Detached.to_string());
assert_eq!(task.child_session_id, None);
assert_eq!(task.continuation_generation, 1);
assert_eq!(continuation_prompt, Some("Add the missing edge case"));
assert_eq!(task.touched_areas, r#"["docs/"]"#);
assert_eq!(child.role.as_deref(), Some("Worker"));
assert_eq!(controller.status, "Review");
assert_eq!(controller.questions.as_deref(), Some(""));
}
#[tokio::test]
async fn question_relay_preserves_controller_questions() {
let database = controller_fixture().await;
let project_id = database
.projects()
.upsert_project("/tmp/project", None)
.await
.expect("failed to reload project");
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert orchestration");
let task_id = insert_waiting_orchestration_task(
&database,
project_id,
orchestration_id,
"worker",
"child-worker",
)
.await;
let controller_question = r#"[{"text":"Controller question"}]"#;
database
.sessions()
.update_session_questions("controller", controller_question)
.await
.expect("failed to seed controller question");
let blocked_by_controller = database
.orchestrations()
.surface_orchestration_questions(
orchestration_id,
task_id,
r#"[{"text":"Child question"}]"#,
)
.await
.expect("failed to inspect controller question");
let orchestration = database
.orchestrations()
.load_orchestration_for_controller("controller")
.await
.expect("failed to load orchestration")
.expect("orchestration should exist");
let controller = database
.sessions()
.load_session("controller")
.await
.expect("failed to load controller")
.expect("controller should exist");
assert!(!blocked_by_controller);
assert_eq!(orchestration.relayed_question_task_id, None);
assert_eq!(controller.questions.as_deref(), Some(controller_question));
}
#[tokio::test]
async fn question_relay_claims_one_exact_task_at_a_time() {
let database = controller_fixture().await;
let project_id = database
.projects()
.upsert_project("/tmp/project", None)
.await
.expect("failed to reload project");
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert orchestration");
let first_task_id = insert_waiting_orchestration_task(
&database,
project_id,
orchestration_id,
"first",
"child-first",
)
.await;
let second_task_id = insert_waiting_orchestration_task(
&database,
project_id,
orchestration_id,
"second",
"child-second",
)
.await;
let first_surfaced = database
.orchestrations()
.surface_orchestration_questions(
orchestration_id,
first_task_id,
r#"[{"text":"First child question"}]"#,
)
.await
.expect("failed to surface first child question");
let second_blocked = database
.orchestrations()
.surface_orchestration_questions(
orchestration_id,
second_task_id,
r#"[{"text":"Second child question"}]"#,
)
.await
.expect("failed to inspect occupied relay");
let claimed_orchestration = database
.orchestrations()
.load_orchestration_for_controller("controller")
.await
.expect("failed to load claimed orchestration")
.expect("orchestration should exist");
database
.orchestrations()
.clear_orchestration_questions(orchestration_id)
.await
.expect("failed to release first relay");
let second_surfaced = database
.orchestrations()
.surface_orchestration_questions(
orchestration_id,
second_task_id,
r#"[{"text":"Second child question"}]"#,
)
.await
.expect("failed to surface second child question");
let second_claimed_orchestration = database
.orchestrations()
.load_orchestration_for_controller("controller")
.await
.expect("failed to load second claimed orchestration")
.expect("orchestration should exist");
assert!(first_surfaced);
assert!(!second_blocked);
assert_eq!(
claimed_orchestration.relayed_question_task_id,
Some(first_task_id)
);
assert!(second_surfaced);
assert_eq!(
second_claimed_orchestration.relayed_question_task_id,
Some(second_task_id)
);
}
#[tokio::test]
async fn infrastructure_retries_are_bounded_and_completion_is_idempotent() {
let database = controller_fixture().await;
let orchestration_id = database
.orchestrations()
.insert_orchestration("controller", &OrchestrationStatus::Running.to_string(), 2)
.await
.expect("failed to insert orchestration");
let task_id = database
.orchestrations()
.upsert_orchestration_task(planned_task(orchestration_id, "alpha"))
.await
.expect("failed to insert task");
let first_status = database
.orchestrations()
.record_orchestration_spawn_failure(task_id, "provider unavailable", 2)
.await
.expect("failed to record first retry");
let second_status = database
.orchestrations()
.record_orchestration_spawn_failure(task_id, "provider unavailable", 2)
.await
.expect("failed to record second retry");
let final_status = database
.orchestrations()
.record_orchestration_spawn_failure(task_id, "provider unavailable", 2)
.await
.expect("failed to exhaust retries");
database
.orchestrations()
.update_orchestration_status(
orchestration_id,
&OrchestrationStatus::Integrating.to_string(),
)
.await
.expect("failed to begin integration completion");
let completed = database
.orchestrations()
.complete_orchestration_campaign(orchestration_id)
.await
.expect("failed to complete campaign");
let duplicate_completion = database
.orchestrations()
.complete_orchestration_campaign(orchestration_id)
.await
.expect("failed to inspect duplicate completion");
let orchestration = database
.orchestrations()
.load_orchestration_for_controller("controller")
.await
.expect("failed to load completed campaign")
.expect("completed campaign should exist");
let task = database
.orchestrations()
.load_orchestration_tasks(orchestration_id)
.await
.expect("failed to load exhausted task")
.remove(0);
let controller = database
.sessions()
.load_session("controller")
.await
.expect("failed to load completed controller")
.expect("controller should exist");
assert_eq!(first_status, OrchestrationTaskStatus::Planned.to_string());
assert_eq!(second_status, OrchestrationTaskStatus::Planned.to_string());
assert_eq!(final_status, OrchestrationTaskStatus::Failed.to_string());
assert!(completed);
assert!(!duplicate_completion);
assert_eq!(task.infrastructure_retry_count, 3);
assert_eq!(task.status, OrchestrationTaskStatus::Failed.to_string());
assert_eq!(orchestration.status, OrchestrationStatus::Done.to_string());
assert_eq!(controller.status, "Done");
}
#[tokio::test]
async fn test_missing_orchestration_returns_none() {
let database = controller_fixture().await;
let orchestration = database
.orchestrations()
.load_orchestration_for_controller("controller")
.await
.expect("failed to load orchestration");
assert_eq!(orchestration, None);
}
}