Skip to content

Commit 654af39

Browse files
committed
feat(agent-core): add linked inbox routing
Authorize linked Member and Coordinator side quests from persisted turn context, make call receipts transactional and idempotent, keep formal claims isolated, and delete causal user-directed rows safely before Team sessions.\n\nVerification:\n- cargo test -p agent_core --lib — 3413 passed, 0 failed, 2 ignored\n- packaged BuildFast linked side quest, restart recovery, and Delete scenarios — passed Pre-commit hook ran. Total eslint: 5, total circular: 0
1 parent 17ad45e commit 654af39

17 files changed

Lines changed: 1610 additions & 67 deletions

File tree

src-tauri/crates/agent-core/src/core/coordination/agent_org_formal_triggers/claim.rs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,8 @@ pub(crate) fn claim_for_coordinator_turn(
3030
if context.org_run_id != org_run_id
3131
|| context.turn_kind
3232
!= crate::coordination::agent_org_turn_contexts::AgentOrgTurnKind::Coordinator
33+
|| context.source_kind
34+
!= crate::coordination::agent_org_turn_contexts::AgentOrgTurnSourceKind::RootTurn
3335
{
3436
return Err(
3537
"FormalTriggerReceipt claim requires exact Coordinator Turn authority".to_string(),
@@ -51,6 +53,7 @@ pub(crate) fn claim_for_coordinator_turn(
5153
WHERE attempt.session_id=?1 AND attempt.turn_intent_id=?2
5254
AND attempt.status IN ('queued','running')
5355
AND receipt.org_run_id=?3 AND receipt.status='materialized'
56+
AND inbox.delivery_class='formal_work'
5457
AND inbox.read_at IS NULL
5558
ORDER BY receipt.created_at,receipt.inbox_id,receipt.receipt_id",
5659
)
@@ -117,6 +120,7 @@ pub(crate) fn claim_for_coordinator_turn(
117120
AND attempt.status IN ('queued','running')
118121
WHERE receipt.org_run_id=?1 AND receipt.status IN ('pending','materialized')
119122
AND receipt.doorbell_status IN ('missing','delivered')
123+
AND inbox.delivery_class='formal_work'
120124
AND inbox.read_at IS NULL
121125
AND NOT EXISTS (
122126
SELECT 1 FROM agent_org_runtime_inbox_delivery_resolutions resolution

src-tauri/crates/agent-core/src/core/coordination/agent_org_runs/store/delete.rs

Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,63 @@ impl AgentOrgRunDeleteOutcome {
1616
}
1717
}
1818

19+
fn delete_user_directed_work_for_run_with_connection(
20+
conn: &Connection,
21+
run_id: &str,
22+
) -> Result<(), String> {
23+
conn.execute(
24+
"DELETE FROM agent_org_runtime_user_directed_coordinator_bindings
25+
WHERE org_run_id=?1",
26+
params![run_id],
27+
)
28+
.map_err(|err| {
29+
format!("failed to delete user-directed Coordinator bindings for {run_id}: {err}")
30+
})?;
31+
32+
let mut remaining = conn
33+
.query_row(
34+
"SELECT COUNT(*)
35+
FROM agent_org_runtime_user_directed_deliveries
36+
WHERE org_run_id=?1",
37+
params![run_id],
38+
|row| row.get::<_, usize>(0),
39+
)
40+
.map_err(|err| format!("failed to count user-directed deliveries for {run_id}: {err}"))?;
41+
while remaining > 0 {
42+
let deleted = conn
43+
.execute(
44+
"DELETE FROM agent_org_runtime_user_directed_deliveries
45+
WHERE delivery_id IN (
46+
SELECT parent.delivery_id
47+
FROM agent_org_runtime_user_directed_deliveries parent
48+
WHERE parent.org_run_id=?1
49+
AND NOT EXISTS (
50+
SELECT 1
51+
FROM agent_org_runtime_user_directed_deliveries child
52+
WHERE child.parent_delivery_id=parent.delivery_id
53+
)
54+
)",
55+
params![run_id],
56+
)
57+
.map_err(|err| {
58+
format!("failed to delete user-directed delivery leaves for {run_id}: {err}")
59+
})?;
60+
if deleted == 0 || deleted > remaining {
61+
return Err(format!(
62+
"user_directed_delete_causal_cycle: Run {run_id} has {remaining} delivery row(s) but no removable causal leaf"
63+
));
64+
}
65+
remaining -= deleted;
66+
}
67+
68+
conn.execute(
69+
"DELETE FROM agent_org_runtime_user_directed_roots WHERE org_run_id=?1",
70+
params![run_id],
71+
)
72+
.map_err(|err| format!("failed to delete user-directed roots for {run_id}: {err}"))?;
73+
Ok(())
74+
}
75+
1976
impl AgentOrgRunStore {
2077
pub fn delete_by_id(run_id: &str) -> Result<(), String> {
2178
let outcome = with_sessions_writer(|| -> Result<AgentOrgRunDeleteOutcome, String> {
@@ -58,6 +115,13 @@ impl AgentOrgRunStore {
58115
.map_err(|err| err.to_string())?
59116
};
60117

118+
// User-directed deliveries deliberately use RESTRICT for their
119+
// parent-Inbox and parent-delivery authority. Tear the causal graph
120+
// down from its leaves before deleting Turn intents or Inbox rows;
121+
// otherwise those cascades can try to remove a parent while a linked
122+
// child still proves that it was derived from it.
123+
delete_user_directed_work_for_run_with_connection(conn, run_id)?;
124+
61125
// Intent ownership is explicit. The hierarchy delete caller rejects
62126
// nested run roots before reaching this helper; standalone run cleanup
63127
// still deletes only rows owned by the requested run.

src-tauri/crates/agent-core/src/core/coordination/agent_org_runs/store/lifecycle.rs

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -158,8 +158,16 @@ impl AgentOrgRunStore {
158158
}
159159
let is_coordinator = context.participant_id == COORDINATOR_MEMBER_ID
160160
&& context.turn_kind
161-
== crate::coordination::agent_org_turn_contexts::AgentOrgTurnKind::Coordinator;
162-
let is_user_directed_writer = if context.turn_kind
161+
== crate::coordination::agent_org_turn_contexts::AgentOrgTurnKind::Coordinator
162+
&& context.source_kind
163+
== crate::coordination::agent_org_turn_contexts::AgentOrgTurnSourceKind::RootTurn;
164+
let is_user_directed_coordinator = context.participant_id == COORDINATOR_MEMBER_ID
165+
&& context.turn_kind
166+
== crate::coordination::agent_org_turn_contexts::AgentOrgTurnKind::Coordinator
167+
&& context.source_kind
168+
== crate::coordination::agent_org_turn_contexts::AgentOrgTurnSourceKind::MemberInbox
169+
&& context.activation_generation.is_none();
170+
let is_user_directed_member_writer = if context.turn_kind
163171
== crate::coordination::agent_org_turn_contexts::AgentOrgTurnKind::UserDirectedWork
164172
&& context.participant_id != COORDINATOR_MEMBER_ID
165173
&& context.dispatch_member_id.as_deref() == Some(context.participant_id.as_str())
@@ -182,7 +190,7 @@ impl AgentOrgRunStore {
182190
} else {
183191
false
184192
};
185-
if !is_coordinator && !is_user_directed_writer {
193+
if !is_coordinator && !is_user_directed_coordinator && !is_user_directed_member_writer {
186194
return Err(
187195
"task_graph_writer_idle_activation_requires_canonical_writer_turn".to_string(),
188196
);
@@ -212,6 +220,7 @@ impl AgentOrgRunStore {
212220
SET activation_generation=?4
213221
WHERE session_id=?1 AND turn_intent_id=?2 AND org_run_id=?3
214222
AND participant_id='coordinator' AND turn_kind='coordinator'
223+
AND source_kind='root_turn'
215224
AND activation_generation=?5",
216225
params![
217226
session_id,

src-tauri/crates/agent-core/src/core/session/turn/event_handler/mod.rs

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -514,6 +514,51 @@ impl UnifiedEventHandler {
514514
}
515515
}
516516

517+
/// Bind every Direct/Group/Linked assistant event to the exact durable
518+
/// UDW receipt. PR10 can project replies from this causal authority without
519+
/// guessing from timestamps, display names, or adjacent transcript rows.
520+
fn attach_agent_org_user_directed_reply(
521+
&self,
522+
session_id: &str,
523+
event: &mut SessionEvent,
524+
) -> bool {
525+
let Some(turn_intent_id) = self.config.agent_org_turn_intent_id.as_deref() else {
526+
return true;
527+
};
528+
match crate::coordination::agent_org_user_directed_work::causal_reply_for_turn(
529+
session_id,
530+
turn_intent_id,
531+
) {
532+
Ok(Some(authority)) => {
533+
let Some(result) = event.result.as_object_mut() else {
534+
self.record_assistant_persistence_error(
535+
"assistant UDW causal reply target is not an object".to_string(),
536+
);
537+
return false;
538+
};
539+
match serde_json::to_value(authority) {
540+
Ok(authority) => {
541+
result.insert("agent_org_user_directed_reply".to_string(), authority);
542+
true
543+
}
544+
Err(error) => {
545+
self.record_assistant_persistence_error(format!(
546+
"assistant UDW causal reply serialization failed: {error}"
547+
));
548+
false
549+
}
550+
}
551+
}
552+
Ok(None) => true,
553+
Err(error) => {
554+
self.record_assistant_persistence_error(format!(
555+
"assistant UDW causal reply lookup failed: {error}"
556+
));
557+
false
558+
}
559+
}
560+
}
561+
517562
/// Bind a backend-issued completion certificate to the exact final
518563
/// assistant event. The model's prose is deliberately not authoritative:
519564
/// consumers can project Delivered only from this typed metadata.
@@ -569,6 +614,7 @@ impl UnifiedEventHandler {
569614
event: &mut SessionEvent,
570615
) -> bool {
571616
self.attach_agent_org_direct_reply(session_id, event)
617+
&& self.attach_agent_org_user_directed_reply(session_id, event)
572618
&& self.attach_agent_org_completion_certificate(session_id, event)
573619
}
574620

src-tauri/crates/agent-core/src/core/session/turn/processor/inbox_drain/drain.rs

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -153,6 +153,11 @@ fn drain_and_render_deferred_impl(
153153
}
154154
Some(crate::coordination::agent_org_turn_contexts::AgentOrgTurnKind::Coordinator) => {
155155
let context = turn_context.expect("Coordinator arm requires persisted context");
156+
if context.source_kind
157+
== crate::coordination::agent_org_turn_contexts::AgentOrgTurnSourceKind::MemberInbox
158+
{
159+
return DrainGuard::empty(&org_context.run_id, recipient_member_id_value);
160+
}
156161
match crate::coordination::agent_org_final_summary::is_summary_turn(
157162
&context.session_id,
158163
&context.turn_intent_id,

src-tauri/crates/agent-core/src/core/session/turn/processor/inbox_drain/guard.rs

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -71,11 +71,13 @@ impl DrainGuard {
7171
mut self,
7272
turn_context: &crate::coordination::agent_org_turn_contexts::AgentOrgTurnContext,
7373
) -> Self {
74-
if matches!(
75-
turn_context.turn_kind,
76-
crate::coordination::agent_org_turn_contexts::AgentOrgTurnKind::Coordinator
77-
| crate::coordination::agent_org_turn_contexts::AgentOrgTurnKind::TaskExecution
78-
) {
74+
if turn_context.turn_kind
75+
== crate::coordination::agent_org_turn_contexts::AgentOrgTurnKind::TaskExecution
76+
|| (turn_context.turn_kind
77+
== crate::coordination::agent_org_turn_contexts::AgentOrgTurnKind::Coordinator
78+
&& turn_context.source_kind
79+
== crate::coordination::agent_org_turn_contexts::AgentOrgTurnSourceKind::RootTurn)
80+
{
7981
self.formal_turn_intent_id = Some(turn_context.turn_intent_id.clone());
8082
}
8183
self

src-tauri/crates/agent-core/src/core/session/turn/processor/prompt.rs

Lines changed: 33 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -233,15 +233,39 @@ impl UnifiedMessageProcessor {
233233
match tokio::task::spawn_blocking(move || {
234234
let (task_snapshot, presented_revision, completion_candidate) = if coordinator_prompt {
235235
match coordinator_turn_intent_id.as_deref() {
236-
Some(turn_intent_id) => match crate::coordination::agent_org_runs::AgentOrgRunStore::stage_coordinator_work_revision_and_load_tasks(
237-
&context_snapshot.run_id,
238-
&coordinator_session_id,
239-
turn_intent_id,
240-
&projected_inbox_ids,
241-
) {
242-
Ok((revision, tasks, candidate)) => (Ok(tasks), revision, Some(candidate)),
243-
Err(error) => (Err(error), None, None),
244-
},
236+
Some(turn_intent_id) => {
237+
let member_inbox_side_quest = crate::coordination::agent_org_turn_contexts::optional_context_for_session(
238+
&coordinator_session_id,
239+
turn_intent_id,
240+
)
241+
.map(|context| {
242+
context.is_some_and(|context| {
243+
context.turn_kind
244+
== crate::coordination::agent_org_turn_contexts::AgentOrgTurnKind::Coordinator
245+
&& context.source_kind
246+
== crate::coordination::agent_org_turn_contexts::AgentOrgTurnSourceKind::MemberInbox
247+
})
248+
});
249+
match member_inbox_side_quest {
250+
Ok(true) => (
251+
crate::coordination::agent_org_tasks::AgentOrgTaskStore::list_operational(
252+
&context_snapshot.run_id,
253+
),
254+
None,
255+
None,
256+
),
257+
Ok(false) => match crate::coordination::agent_org_runs::AgentOrgRunStore::stage_coordinator_work_revision_and_load_tasks(
258+
&context_snapshot.run_id,
259+
&coordinator_session_id,
260+
turn_intent_id,
261+
&projected_inbox_ids,
262+
) {
263+
Ok((revision, tasks, candidate)) => (Ok(tasks), revision, Some(candidate)),
264+
Err(error) => (Err(error), None, None),
265+
},
266+
Err(error) => (Err(error), None, None),
267+
}
268+
}
245269
None => (
246270
Err("Coordinator prompt requires an exact Turn intent id".to_string()),
247271
None,

src-tauri/crates/agent-core/src/core/tools/call_context.rs

Lines changed: 27 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -272,16 +272,33 @@ impl CallContext {
272272
.is_some_and(|context| {
273273
let persisted_profile = match context.turn_kind {
274274
crate::coordination::agent_org_turn_contexts::AgentOrgTurnKind::Coordinator => {
275-
let is_summary_turn = database::db::get_connection()
276-
.map_err(|error| error.to_string())
277-
.and_then(|conn| {
278-
crate::coordination::agent_org_final_summary::is_summary_turn_with_connection(
279-
&conn,
280-
&context.session_id,
281-
&context.turn_intent_id,
282-
)
283-
});
284-
coordinator_replay_profile(is_summary_turn)
275+
if context.source_kind
276+
== crate::coordination::agent_org_turn_contexts::AgentOrgTurnSourceKind::MemberInbox
277+
{
278+
match crate::coordination::agent_org_runs::AgentOrgRunStore::get_run_status(
279+
&context.org_run_id,
280+
) {
281+
Ok(Some(
282+
crate::coordination::agent_org_runs::AgentOrgRunStatus::Paused,
283+
)) => Some(AgentOrgTurnToolProfile::SummaryOnly),
284+
Ok(Some(
285+
crate::coordination::agent_org_runs::AgentOrgRunStatus::Running
286+
| crate::coordination::agent_org_runs::AgentOrgRunStatus::Idle,
287+
)) => Some(AgentOrgTurnToolProfile::CoordinatorOrchestration),
288+
_ => None,
289+
}
290+
} else {
291+
let is_summary_turn = database::db::get_connection()
292+
.map_err(|error| error.to_string())
293+
.and_then(|conn| {
294+
crate::coordination::agent_org_final_summary::is_summary_turn_with_connection(
295+
&conn,
296+
&context.session_id,
297+
&context.turn_intent_id,
298+
)
299+
});
300+
coordinator_replay_profile(is_summary_turn)
301+
}
285302
}
286303
crate::coordination::agent_org_turn_contexts::AgentOrgTurnKind::TaskExecution => {
287304
Some(AgentOrgTurnToolProfile::TaskExecution)

0 commit comments

Comments
 (0)