Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -1492,9 +1492,7 @@ mod tests {

#[test]
fn superseded_delivery_can_follow_a_real_replacement_chain_but_not_cycle() {
use crate::definitions::orgs::{
HierarchyMode, OrgDefinition, OrgMember, PlanApprovalPolicy,
};
use crate::definitions::orgs::{FlatOrgMember, OrgDefinition, PlanApprovalPolicy};

let _sandbox = sandbox_with_inbox_schema();
let run_id = "run-delivery-chain";
Expand All @@ -1505,39 +1503,43 @@ mod tests {
role: "Coordinator".into(),
agent_id: "coordinator-agent".into(),
description: None,
hierarchy_mode: HierarchyMode::Soft,
plan_approval_policy: PlanApprovalPolicy::Coordinator,
children: vec![
OrgMember {
id: "member-a".into(),
members: vec![
FlatOrgMember {
member_id: "member-a".into(),
name: "Member A".into(),
role: "worker".into(),
agent_id: "agent-a".into(),
runtime_config: None,
children: Vec::new(),
},
OrgMember {
id: "member-b".into(),
FlatOrgMember {
member_id: "member-b".into(),
name: "Member B".into(),
role: "worker".into(),
agent_id: "agent-b".into(),
runtime_config: None,
children: Vec::new(),
},
OrgMember {
id: "member-c".into(),
FlatOrgMember {
member_id: "member-c".into(),
name: "Member C".into(),
role: "worker".into(),
agent_id: "agent-c".into(),
runtime_config: None,
children: Vec::new(),
},
],
additional_task_graph_writer_member_ids: Vec::new(),
member_communication_links: Vec::new(),
};
let conn = get_connection().expect("open sandbox database");
conn.execute(
"UPDATE agent_org_runs SET org_snapshot_json=?1 WHERE id=?2",
params![serde_json::to_string(&org).unwrap(), run_id],
params![
serde_json::to_string(&crate::definitions::orgs::AgentOrgLaunchSnapshot::from(
&org
))
.unwrap(),
run_id
],
)
.expect("seed roster snapshot");
let now = chrono::Utc::now().to_rfc3339();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ use crate::coordination::agent_org_runs::{
use crate::coordination::agent_org_tasks::{
AgentOrgTaskStore, CreateTaskParams, TaskStatus, TASK_METADATA_EXECUTION_MODE,
};
use crate::definitions::orgs::{HierarchyMode, OrgDefinition, OrgMember};
use crate::definitions::orgs::{FlatOrgMember, OrgDefinition};

fn setup(policy: PlanApprovalPolicy) -> (test_helpers::test_env::SandboxGuard, AgentOrgRunContext) {
let sandbox = test_helpers::test_env::sandbox();
Expand Down Expand Up @@ -48,32 +48,31 @@ fn setup(policy: PlanApprovalPolicy) -> (test_helpers::test_env::SandboxGuard, A
role: "lead".into(),
agent_id: "coord-agent".into(),
description: None,
hierarchy_mode: HierarchyMode::Soft,
plan_approval_policy: policy,
children: vec![
OrgMember {
id: "planner".into(),
members: vec![
FlatOrgMember {
member_id: "planner".into(),
name: "Planner".into(),
role: "plan".into(),
agent_id: "planner-agent".into(),
runtime_config: None,
children: Vec::new(),
},
OrgMember {
id: "builder".into(),
FlatOrgMember {
member_id: "builder".into(),
name: "Builder".into(),
role: "build".into(),
agent_id: "builder-agent".into(),
runtime_config: None,
children: Vec::new(),
},
],
additional_task_graph_writer_member_ids: Vec::new(),
member_communication_links: Vec::new(),
};
let run = AgentOrgRunStore::create(CreateAgentOrgRunParams {
org_id: org.id.clone(),
coordinator_agent_id: org.agent_id.clone(),
root_session_id: Some("root-plan-approval".into()),
org_snapshot: org,
org_snapshot: (&org).into(),
entry_mode: AgentOrgRunEntryMode::StandaloneSession,
status: AgentOrgRunStatus::Running,
work_item_id: None,
Expand All @@ -95,18 +94,16 @@ fn setup(policy: PlanApprovalPolicy) -> (test_helpers::test_env::SandboxGuard, A
name: "Planner".into(),
role: "plan".into(),
agent_id: "planner-agent".into(),
parent_member_id: None,
},
AgentOrgContextMember {
member_id: "builder".into(),
name: "Builder".into(),
role: "build".into(),
agent_id: "builder-agent".into(),
parent_member_id: None,
},
],
hierarchy_mode: HierarchyMode::Soft,
plan_approval_policy: policy,
capability_index: Default::default(),
root_session_id: Some("root-plan-approval".into()),
};
(sandbox, context)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ use crate::coordination::agent_org_runs::COORDINATOR_MEMBER_ID;
use crate::coordination::agent_org_tasks::{
AgentOrgTaskStore, TaskExecutionMode, TaskOutput, TaskStatus, TASK_METADATA_EXECUTION_MODE,
};
use crate::definitions::orgs::{OrgDefinition, OrgMember};
use crate::definitions::orgs::AgentOrgLaunchSnapshot;

use super::artifact::validate_owned_plan_path_with_connection;
use super::persistence::insert_record;
Expand Down Expand Up @@ -299,24 +299,19 @@ fn participant_agent_ids_in_tx(
.map_err(|err| format!("agent_org_run_not_mutable: {run_id}: {err}"))?;
let mut participants = HashMap::new();
if let Some(snapshot_json) = snapshot_json {
let snapshot: OrgDefinition = serde_json::from_str(&snapshot_json).map_err(|err| {
format!("failed to parse Agent Org launch snapshot for run {run_id}: {err}")
})?;
collect_participant_agent_ids(&snapshot.children, &mut participants);
let snapshot: AgentOrgLaunchSnapshot =
serde_json::from_str(&snapshot_json).map_err(|err| {
format!("failed to parse Agent Org launch snapshot for run {run_id}: {err}")
})?;
crate::definitions::orgs::validate_launch_snapshot(&snapshot)
.map_err(|err| format!("invalid Agent Org launch snapshot for run {run_id}: {err}"))?;
for member in snapshot.members {
participants.insert(member.member_id, member.agent_id);
}
}
Ok((coordinator_agent_id, participants))
}

fn collect_participant_agent_ids(
members: &[OrgMember],
participants: &mut HashMap<String, String>,
) {
for member in members {
participants.insert(member.id.clone(), member.agent_id.clone());
collect_participant_agent_ids(&member.children, participants);
}
}

pub(super) fn plan_approval_request_message(approval: &AgentOrgPlanApproval) -> AgentMessage {
let plan_char_count = approval.plan_content.chars().count();
let mut inline_plan_content =
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
use rusqlite::{params, Connection, OptionalExtension, Result as SqliteResult};

use crate::definitions::orgs::{AgentOrgsStore, OrgDefinition, OrgMember};
use crate::definitions::orgs::{
validate_launch_snapshot, AgentOrgCapabilityIndex, AgentOrgLaunchSnapshot,
};
use database::db::get_connection;

use super::{
Expand Down Expand Up @@ -146,64 +148,58 @@ pub(super) fn row_to_run(row: &rusqlite::Row<'_>) -> SqliteResult<AgentOrgRunRec

pub(super) fn context_for_run_record(
run: &AgentOrgRunRecord,
org_store: &AgentOrgsStore,
) -> Result<AgentOrgRunContext, String> {
if let Some(snapshot_json) = run.org_snapshot_json.as_deref() {
let snapshot: OrgDefinition = serde_json::from_str(snapshot_json).map_err(|err| {
format!(
"failed to parse Agent Org launch snapshot for run {}: {}",
run.id, err
)
})?;
return Ok(context_from_run_and_org(run, &snapshot));
}

let org = org_store.get(&run.org_id)?;
Ok(context_from_run_and_org(run, &org))
let snapshot_json = run
.org_snapshot_json
.as_deref()
.ok_or_else(|| format!("Agent Org run {} has no immutable launch snapshot", run.id))?;
let snapshot: AgentOrgLaunchSnapshot = serde_json::from_str(snapshot_json).map_err(|err| {
format!(
"failed to parse Agent Org launch snapshot for run {}: {}",
run.id, err
)
})?;
validate_launch_snapshot(&snapshot).map_err(|err| {
format!(
"Agent Org run {} has invalid launch snapshot: {err}",
run.id
)
})?;
Ok(context_from_run_and_snapshot(run, &snapshot))
}

pub(super) fn context_from_run_and_org(
pub(super) fn context_from_run_and_snapshot(
run: &AgentOrgRunRecord,
org: &OrgDefinition,
snapshot: &AgentOrgLaunchSnapshot,
) -> AgentOrgRunContext {
AgentOrgRunContext {
run_id: run.id.clone(),
org_id: org.id.clone(),
org_name: org.name.clone(),
org_role: org.role.clone(),
org_id: snapshot.org_id.clone(),
org_name: snapshot.org_name.clone(),
org_role: snapshot.coordinator_role.clone(),
coordinator_agent_id: run.coordinator_agent_id.clone(),
coordinator_name: DEFAULT_COORDINATOR_DISPLAY_NAME.to_string(),
coordinator_role: org.role.clone(),
members: flatten_members(&org.children, None),
hierarchy_mode: org.hierarchy_mode,
plan_approval_policy: org.plan_approval_policy,
coordinator_role: snapshot.coordinator_role.clone(),
members: flatten_members(&snapshot.members),
plan_approval_policy: snapshot.plan_approval_policy,
capability_index: AgentOrgCapabilityIndex::from_snapshot(snapshot),
root_session_id: run.root_session_id.clone(),
}
}

/// Flatten the `OrgMember` tree into a `Vec<AgentOrgContextMember>`,
/// preserving each member's parent id (the immediate parent in
/// `OrgDefinition.children`). A `None` parent means the member is a
/// direct report of the coordinator.
///
/// In `HierarchyMode::Flat` the parent ids are still emitted but the
/// system prompt and routing layer ignore them.
/// Project the immutable flat snapshot roster into runtime context rows.
pub(super) fn flatten_members(
members: &[OrgMember],
parent_id: Option<&str>,
members: &[crate::definitions::orgs::FlatOrgMember],
) -> Vec<AgentOrgContextMember> {
let mut flattened = Vec::new();
for member in members {
flattened.push(AgentOrgContextMember {
member_id: member.id.clone(),
members
.iter()
.map(|member| AgentOrgContextMember {
member_id: member.member_id.clone(),
name: member.name.clone(),
role: member.role.clone(),
agent_id: member.agent_id.clone(),
parent_member_id: parent_id.map(|id| id.to_string()),
});
flattened.extend(flatten_members(&member.children, Some(&member.id)));
}
flattened
})
.collect()
}

pub(super) fn insert_run(conn: &Connection, run: &AgentOrgRunRecord) -> SqliteResult<()> {
Expand Down
Loading
Loading