Skip to content
Merged
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
3 changes: 2 additions & 1 deletion crates/tinyagents-orchestration/src/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,8 @@ on orchestration policy.
- `subagent::SubAgentTool` starts a child asynchronously through typed
parent-context tool dispatch and immediately returns a stable job id.
- `subagent::SubAgentJobRegistry` records queued, running, completed, failed,
and cancelled jobs. `SubAgentJobsTool` queries them and
and cancelled jobs on top of `tinyagents_tasks::DetachedTaskRegistry`.
`SubAgentJobsTool` queries them and
`SubAgentMessageTool` sends messages to live children.
- `subagent::SubAgentSession` reuses one child and its transcript across turns.
- `subagent::SpawnPolicy` / `SpawnAdmission` bound child fan-out (per-parent
Expand Down
397 changes: 230 additions & 167 deletions crates/tinyagents-orchestration/src/subagent/invocation/jobs.rs

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -1082,3 +1082,7 @@ mod admission_test;
#[cfg(test)]
#[path = "mod_policy_tests.rs"]
mod policy_test;

#[cfg(test)]
Comment thread
senamakel marked this conversation as resolved.
#[path = "mod_registry_tests.rs"]
Comment thread
senamakel marked this conversation as resolved.
mod registry_tests;
Original file line number Diff line number Diff line change
Expand Up @@ -263,16 +263,16 @@ async fn terminal_jobs_never_return_to_running() {
async fn settled_jobs_release_their_cancellation_token() {
let jobs = SubAgentJobRegistry::new();
let (job_id, _steering) = jobs.create("worker", 1);
assert!(jobs.inner.read().unwrap()[&job_id].cancellation.is_some());
assert!(jobs.holds_live_cancellation(&job_id));
jobs.cancel_owned(job_id.as_str(), 1).expect("cancel");
assert!(jobs.inner.read().unwrap()[&job_id].cancellation.is_none());
assert!(!jobs.inner.read().unwrap()[&job_id].job.status.is_terminal());
assert!(!jobs.holds_live_cancellation(&job_id));
assert!(!jobs.get(job_id.as_str()).unwrap().status.is_terminal());
jobs.mark_result(&job_id, Err(TinyAgentsError::Cancelled));
assert!(jobs.inner.read().unwrap()[&job_id].job.status.is_terminal());
assert!(jobs.get(job_id.as_str()).unwrap().status.is_terminal());

let (aborted, _steering) = jobs.create("worker", 1);
jobs.mark_aborted(&aborted, true);
assert!(jobs.inner.read().unwrap()[&aborted].cancellation.is_none());
assert!(!jobs.holds_live_cancellation(&aborted));
}

#[tokio::test]
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
//! The job registry is an adapter over `tinyagents_tasks::DetachedTaskRegistry`:
//! these tests pin the shared-registry view and the lifecycle hand-offs.

use super::*;
use tinyagents_harness::error::TinyAgentsError;
use tinyagents_harness::ids::TaskId;

#[tokio::test]
async fn jobs_are_registered_in_the_shared_detached_registry() {
let jobs = SubAgentJobRegistry::new();
let (job_id, _steering) = jobs.create("worker", 7);
let task_id = TaskId::new(job_id.as_str());

let tasks = jobs.tasks();
assert_eq!(tasks.len().unwrap(), 1);
let snapshot = tasks.snapshot(&task_id, "7").unwrap();
assert_eq!(snapshot.status.status, SubAgentJobStatus::Queued);
assert_eq!(snapshot.status.agent, "worker");
// Ownership is the detached registry's, so a foreign owner is refused.
assert!(tasks.snapshot(&task_id, "8").is_err());
}

#[tokio::test]
async fn steering_is_reachable_while_live_and_released_when_settled() {
let jobs = SubAgentJobRegistry::new();
let (job_id, _steering) = jobs.create("worker", 1);
let task_id = TaskId::new(job_id.as_str());
assert!(jobs.tasks().steering_handle_trusted(&task_id).is_ok());

jobs.mark_result(&job_id, Err(TinyAgentsError::Cancelled));

assert!(jobs.tasks().steering_handle_trusted(&task_id).is_err());
// A settled job stays queryable; the registry never prunes it.
assert_eq!(
jobs.get(job_id.as_str()).unwrap().status,
SubAgentJobStatus::Cancelled
);
assert_eq!(jobs.tasks().len().unwrap(), 1);
}

#[tokio::test]
async fn terminal_jobs_survive_many_registrations() {
// The detached registry sweeps terminal entries at its soft cap; the job
// registry must never lose a settled job that way.
let jobs = SubAgentJobRegistry::new();
let mut ids = Vec::new();
for _ in 0..2_000 {
let (id, _s) = jobs.create("worker", 1);
jobs.mark_result(&id, Err(TinyAgentsError::Cancelled));
ids.push(id);
}
assert!(ids.iter().all(|id| jobs.get(id.as_str()).is_some()));
assert_eq!(jobs.list().len(), 2_000);
}

#[tokio::test]
async fn control_errors_keep_their_documented_order() {
let jobs = SubAgentJobRegistry::new();
let (live, _s) = jobs.create("worker", 1);
let (done, _s2) = jobs.create("worker", 1);
jobs.mark_result(&done, Err(TinyAgentsError::Cancelled));
let not_found = |id: &str| SubAgentJobError::NotFound(id.to_owned());

// Unknown id and foreign owner are both NotFound, even for a settled job.
assert_eq!(jobs.cancel_owned("nope", 1).unwrap_err(), not_found("nope"));
assert_eq!(
jobs.cancel_owned(live.as_str(), 2).unwrap_err(),
not_found(live.as_str())
);
assert_eq!(
jobs.cancel_owned(done.as_str(), 2).unwrap_err(),
not_found(done.as_str())
);
assert_eq!(
jobs.send_message_with_request_id(done.as_str(), 2, "m", None)
.unwrap_err(),
not_found(done.as_str())
);
// Own settled job is Terminal.
let terminal = SubAgentJobError::Terminal {
job_id: done.as_str().to_owned(),
status: SubAgentJobStatus::Cancelled,
};
assert_eq!(
jobs.cancel_owned(done.as_str(), 1).unwrap_err(),
terminal.clone()
);
assert_eq!(
jobs.send_message_with_request_id(done.as_str(), 1, "m", None)
.unwrap_err(),
terminal
);
// An oversized id on a live job is RequestIdTooLong...
let long = "x".repeat(10_000);
assert_eq!(
jobs.send_message_with_request_id(live.as_str(), 1, "m", Some(&long))
.unwrap_err(),
SubAgentJobError::RequestIdTooLong
);
// ...but a cancelling job reports Cancelling before the id is looked at.
jobs.cancel_owned(live.as_str(), 1).unwrap();
let cancelling = SubAgentJobError::Cancelling(live.as_str().to_owned());
assert_eq!(
jobs.send_message_with_request_id(live.as_str(), 1, "m", Some(&long))
.unwrap_err(),
cancelling
);
// Repeated cancels of a cancelling job stay Ok.
assert!(jobs.cancel_owned(live.as_str(), 1).is_ok());
}

#[tokio::test]
async fn listings_are_sorted_by_job_id() {
let jobs = SubAgentJobRegistry::new();
for _ in 0..40 {
jobs.create("worker", 1);
}
let all: Vec<String> = jobs.list().into_iter().map(|j| j.id.0).collect();
let mut sorted = all.clone();
sorted.sort();
assert_eq!(all, sorted);
assert_eq!(all.len(), 40);
}
62 changes: 45 additions & 17 deletions crates/tinyagents-orchestration/src/subagent/invocation/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,14 +12,15 @@
//! sibling `mod.rs` and its `*_tests.rs` files.

use std::collections::HashMap;
use std::sync::{Arc, RwLock};
use std::sync::{Arc, Mutex};

use serde_json::Value;

use tinyagents_harness::cancel::CancellationToken;
use tinyagents_harness::events::EventSink;
use tinyagents_harness::runtime::AgentHarness;
use tinyagents_harness::steering::{RecentRequestIds, SteeringHandle};
use tokio::sync::watch;

use tinyagents_tasks::{DetachedTaskRegistry, SteeringRegistry};
use tinyinference_llm::message::Message;

/// The argument key a [`SubAgentTool`] reads the child input from.
Expand Down Expand Up @@ -285,23 +286,50 @@ pub(crate) struct JobLink {
}

/// Shared registry behind asynchronous subagent spawning and host controls.
#[derive(Clone, Default)]
///
/// This is a thin adapter over [`tinyagents_tasks::DetachedTaskRegistry`], the
/// one live detached-task registry implementation. The registry holds each
/// job's [`SubAgentJob`] snapshot as its watched status (so ownership,
/// snapshots, steering lookup and request-id dedupe are the detached
/// registry's), registered cooperatively because a job is stopped through its
/// own cancellation token rather than hard-aborted. The adapter adds only what
/// is specific to subagent jobs: the lifecycle transitions, the
/// cancel-then-settle protocol, and retention of settled jobs.
#[derive(Clone)]
pub struct SubAgentJobRegistry {
pub(crate) inner: Arc<RwLock<HashMap<SubAgentJobId, SubAgentJobEntry>>>,
// Private to `jobs.rs`. Invariant: only cooperative calls (register,
// cancel_cooperative, release_cancellation, snapshots) may touch it, under
// `controls`. Its removing APIs (`cancel`, `wait`, `sweep_terminal`) would
// drop settled jobs, which must stay queryable.
pub(super) tasks: DetachedTaskRegistry<JobMeta, SubAgentJob>,
pub(crate) steering: SteeringRegistry,
/// Status senders and cancel flags. Also the transition gate: every
/// mutation holds this lock, so a settle and a cancel never interleave.
pub(crate) controls: Arc<Mutex<HashMap<SubAgentJobId, JobControl>>>,
}

impl Default for SubAgentJobRegistry {
fn default() -> Self {
let steering = SteeringRegistry::new();
Self {
// Settled jobs must stay queryable, so the soft cap that makes the
// detached registry sweep terminal entries is never reached.
tasks: DetachedTaskRegistry::new(steering.clone(), usize::MAX, |job: &SubAgentJob| {
job.status.is_terminal()
}),
steering,
controls: Arc::default(),
}
}
}

pub(crate) struct SubAgentJobEntry {
pub(crate) job: SubAgentJob,
/// Identity of the parent run that created this capability.
pub(crate) owner: u64,
pub(crate) steering: SteeringHandle,
/// The child run's own cancellation token (a linked child of the parent's),
/// so one job can be cancelled without touching the parent or siblings.
///
/// `None` once the job is terminal: a settled job holds no live token.
pub(crate) cancellation: Option<CancellationToken>,
/// Message `request_id`s already applied, so a retried message is queued once.
pub(crate) message_requests: RecentRequestIds,
/// Application metadata kept with each job in the detached registry.
#[derive(Clone, Debug, Default)]
pub(crate) struct JobMeta;

/// The adapter-owned half of a job: the status publisher and cancel flag.
pub(crate) struct JobControl {
pub(crate) status: watch::Sender<SubAgentJob>,
/// Whether cancellation was requested while the child was still running.
/// The job remains non-terminal until the child reports its result.
pub(crate) cancellation_requested: bool,
Expand Down
68 changes: 68 additions & 0 deletions crates/tinyagents-tasks/src/lib_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -309,6 +309,74 @@ async fn detached_registry_reports_a_poisoned_lock() {
join.abort();
}

#[tokio::test]
async fn cooperative_registration_cancels_without_removing_or_aborting() {
let registry = runtime_registry();
let task_id = TaskId::new("coop-cancel");
let (tx, rx, _cancellation, join) = detached_handles();
let cancellation = CancellationToken::new();
registry
.register_cooperative(
task_id.clone(),
"parent",
"meta".to_string(),
rx,
cancellation.clone(),
)
.unwrap();
assert!(registry.holds_cancellation(&task_id).unwrap());

assert_eq!(
registry.cancel_cooperative(&task_id, "other").unwrap_err(),
DetachedTaskRegistryError::NotOwned
);
assert!(!cancellation.is_cancelled());

let snapshot = registry.cancel_cooperative(&task_id, "parent").unwrap();
assert_eq!(snapshot.status, RuntimeStatus::Running);
assert!(cancellation.is_cancelled());
// The entry stays registered and the token is released.
assert_eq!(registry.len().unwrap(), 1);
assert!(!registry.holds_cancellation(&task_id).unwrap());

tx.send(RuntimeStatus::Completed("x".into())).unwrap();
assert_eq!(
registry.cancel_cooperative(&task_id, "parent").unwrap_err(),
DetachedTaskRegistryError::AlreadyDone
);
join.abort();
}

#[tokio::test]
async fn release_cancellation_drops_the_token_without_cancelling_it() {
let registry = runtime_registry();
let task_id = TaskId::new("coop-release");
let (_tx, rx, _c, join) = detached_handles();
let cancellation = CancellationToken::new();
registry
.register_cooperative(
task_id.clone(),
"p",
"m".to_string(),
rx,
cancellation.clone(),
)
.unwrap();
assert_eq!(
registry
.release_cancellation_trusted(&TaskId::new("missing"))
.unwrap_err(),
DetachedTaskRegistryError::Unknown
);
registry.release_cancellation_trusted(&task_id).unwrap();
assert!(!registry.holds_cancellation(&task_id).unwrap());
assert!(!cancellation.is_cancelled());
// Hard cancel still works on an entry that holds no abort handle.
let cancelled = registry.cancel_trusted(&task_id).unwrap();
assert_eq!(cancelled.metadata, "m");
join.abort();
}

fn unique_log_path(tag: &str) -> std::path::PathBuf {
// Deterministic-per-test path in the system temp dir (no clock/random ids).
std::env::temp_dir().join(format!("tinyagents-taskstore-{tag}.jsonl"))
Expand Down
Loading
Loading