use super::*;
use std::time::Duration;
#[derive(Clone, Debug)]
pub struct SubAgentLimits {
pub max_sessions: usize,
pub max_parallel_requests: usize,
pub max_requests: Option<u64>,
pub max_tokens: Option<u64>,
pub deadline_ms: Option<u64>,
pub max_message_bytes: usize,
pub max_pending_messages: usize,
pub max_events: usize,
pub terminal_retention_secs: u64,
}
impl Default for SubAgentLimits {
fn default() -> Self {
Self {
max_sessions: 64,
max_parallel_requests: 8,
max_requests: None,
max_tokens: None,
deadline_ms: None,
max_message_bytes: 64 * 1024,
max_pending_messages: 42,
max_events: 128,
terminal_retention_secs: 3600,
}
}
}
#[derive(Clone, Debug, Deserialize, Serialize, PartialEq, Eq)]
pub struct ExecutionIdentity {
pub id: String,
pub root_id: String,
pub parent_id: Option<String>,
}
#[derive(Clone, Copy, Debug, Deserialize, Serialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum SubAgentEventKind {
Started,
TurnCompleted,
Interrupted,
Failed,
Closed,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct SubAgentEvent {
pub sequence: u64,
pub execution: ExecutionIdentity,
pub agent: String,
pub session: String,
pub turn: u64,
pub kind: SubAgentEventKind,
pub summary: Option<String>,
}
#[derive(Clone, Debug, Serialize)]
pub struct SubAgentEvents {
pub cursor: u64,
pub lagged: bool,
pub timed_out: bool,
pub events: Vec<SubAgentEvent>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum WaitMode {
Any,
All,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum MessageDelivery {
QueueOnly,
TriggerTurn,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct SubAgentMessage {
pub sender: String,
pub id: String,
pub content: String,
#[serde(default)]
pub resources: Vec<Resource>,
}
#[derive(Default)]
struct ScopeState {
owner: Option<Principal>,
activity: HashMap<String, bool>,
sessions: usize,
requests_admitted: u64,
usage: Usage,
sequence: u64,
events: VecDeque<SubAgentEvent>,
terminals: BTreeMap<(Principal, String, String), (u64, Json)>,
}
struct ScopeInner {
id: String,
limits: SubAgentLimits,
state: Mutex<ScopeState>,
changed: tokio::sync::watch::Sender<u64>,
requests: Arc<tokio::sync::Semaphore>,
}
#[derive(Clone)]
pub struct SubAgentScope(Arc<ScopeInner>);
impl Default for SubAgentScope {
fn default() -> Self {
Self::new(SubAgentLimits::default())
}
}
impl SubAgentScope {
pub fn new(limits: SubAgentLimits) -> Self {
Self::restore(format!("root_{:032x}", rand::random::<u128>()), limits)
}
pub fn restore(id: String, limits: SubAgentLimits) -> Self {
let (changed, _) = tokio::sync::watch::channel(0);
let requests = limits
.max_parallel_requests
.clamp(1, tokio::sync::Semaphore::MAX_PERMITS);
Self(Arc::new(ScopeInner {
id,
limits,
state: Mutex::new(ScopeState::default()),
changed,
requests: Arc::new(tokio::sync::Semaphore::new(requests)),
}))
}
pub fn id(&self) -> &str {
&self.0.id
}
pub fn limits(&self) -> &SubAgentLimits {
&self.0.limits
}
pub fn restore_usage(&self, usage: Usage, admitted_requests: u64) {
let mut state = self.0.state.lock();
state.usage.input_tokens = state.usage.input_tokens.max(usage.input_tokens);
state.usage.output_tokens = state.usage.output_tokens.max(usage.output_tokens);
state.usage.cached_tokens = state.usage.cached_tokens.max(usage.cached_tokens);
state.usage.requests = state.usage.requests.max(usage.requests);
state.requests_admitted = state.requests_admitted.max(admitted_requests);
}
pub fn usage(&self) -> Usage {
self.0.state.lock().usage.clone()
}
pub fn admitted_requests(&self) -> u64 {
self.0.state.lock().requests_admitted
}
pub(crate) fn bind_caller(&self, caller: Principal) -> Result<(), BoxError> {
let mut state = self.0.state.lock();
if state.owner.is_some_and(|owner| owner != caller) {
return Err("subagent scope belongs to another caller".into());
}
state.owner = Some(caller);
Ok(())
}
pub(super) fn set_activity(&self, id: &str, busy: bool) {
self.0.state.lock().activity.insert(id.into(), busy);
}
pub(super) fn activity(&self, id: &str) -> Option<bool> {
self.0.state.lock().activity.get(id).copied()
}
pub(super) fn forget_activity(&self, id: &str) {
self.0.state.lock().activity.remove(id);
}
pub(super) fn identity(&self, parent: Option<ExecutionIdentity>) -> ExecutionIdentity {
ExecutionIdentity {
id: format!("exec_{:032x}", rand::random::<u128>()),
root_id: self.id().into(),
parent_id: parent.filter(|p| p.root_id == self.id()).map(|p| p.id),
}
}
pub(crate) async fn deadline(&self) {
match self.limits().deadline_ms {
Some(deadline) => {
tokio::time::sleep(Duration::from_millis(deadline.saturating_sub(unix_ms()))).await
}
None => std::future::pending::<()>().await,
}
}
pub fn cursor(&self) -> u64 {
self.0.state.lock().sequence
}
pub(crate) fn check_deadline(&self) -> Result<(), BoxError> {
if self
.limits()
.deadline_ms
.is_some_and(|deadline| unix_ms() >= deadline)
{
return Err("subagent root deadline exceeded".into());
}
Ok(())
}
pub(super) fn reserve_session(&self) -> Result<ScopePermit, BoxError> {
self.check_deadline()?;
let mut state = self.0.state.lock();
if state.sessions >= self.limits().max_sessions {
return Err("subagent resident session limit reached".into());
}
state.sessions += 1;
Ok(ScopePermit {
scope: self.clone(),
session: true,
_request: None,
})
}
fn check_request_budget(&self, state: &ScopeState) -> Result<(), BoxError> {
self.check_deadline()?;
if self
.limits()
.max_requests
.is_some_and(|n| state.requests_admitted >= n)
|| self.limits().max_tokens.is_some_and(|n| {
state
.usage
.input_tokens
.saturating_add(state.usage.output_tokens)
>= n
})
{
return Err("subagent root model budget exhausted".into());
}
Ok(())
}
pub(crate) async fn admit_request(&self) -> Result<ScopePermit, BoxError> {
self.check_request_budget(&self.0.state.lock())?;
let request = self.0.requests.clone().acquire_owned().await?;
let mut state = self.0.state.lock();
self.check_request_budget(&state)?;
state.requests_admitted = state.requests_admitted.saturating_add(1);
Ok(ScopePermit {
scope: self.clone(),
session: false,
_request: Some(request),
})
}
pub(crate) fn record_usage(&self, usage: &Usage) {
self.0.state.lock().usage.accumulate(usage);
}
pub(super) fn publish(&self, mut event: SubAgentEvent) {
let mut state = self.0.state.lock();
state.sequence += 1;
event.sequence = state.sequence;
if let Some(summary) = &mut event.summary {
truncate_utf8_to_max_bytes(summary, 2000);
}
state.events.push_back(event);
while state.events.len() > self.limits().max_events.clamp(1, 4096) {
state.events.pop_front();
}
self.0.changed.send_replace(state.sequence);
}
pub fn events(&self, after: u64, targets: &[String]) -> SubAgentEvents {
let state = self.0.state.lock();
SubAgentEvents {
cursor: state.sequence,
lagged: after > state.sequence
|| state
.events
.front()
.is_some_and(|e| after.saturating_add(1) < e.sequence),
timed_out: false,
events: state
.events
.iter()
.filter(|e| {
e.sequence > after && (targets.is_empty() || targets.contains(&e.execution.id))
})
.cloned()
.collect(),
}
}
pub async fn wait(
&self,
after: u64,
targets: &[String],
mode: WaitMode,
timeout: Duration,
cancellation: anda_core::CancellationToken,
) -> Result<SubAgentEvents, BoxError> {
if mode == WaitMode::All && targets.is_empty() {
return Err("wait-all requires execution IDs".into());
}
let mut changes = self.0.changed.subscribe();
let deadline = tokio::time::Instant::now() + timeout.min(Duration::from_secs(60));
loop {
let result = self.events(after, targets);
let ready = match mode {
WaitMode::Any => !result.events.is_empty(),
WaitMode::All => targets.iter().all(|id| {
result
.events
.iter()
.any(|e| &e.execution.id == id && e.kind != SubAgentEventKind::Started)
}),
};
if ready || result.lagged {
return Ok(result);
}
tokio::select! {
biased;
_ = cancellation.cancelled() => return Err("subagent wait cancelled".into()),
_ = tokio::time::sleep_until(deadline) => return Ok(SubAgentEvents { timed_out: true, ..self.events(after, targets) }),
_ = changes.changed() => {}
}
}
}
pub(super) fn remember(&self, caller: Principal, agent: String, session: String, detail: Json) {
let mut state = self.0.state.lock();
state
.terminals
.insert((caller, agent, session), (unix_ms(), detail));
while state.terminals.len() > self.limits().max_events.clamp(1, 4096) {
let oldest = state
.terminals
.iter()
.min_by_key(|(_, (at, _))| *at)
.map(|(key, _)| key.clone())
.unwrap();
state.terminals.remove(&oldest);
}
}
pub(super) fn terminal(&self, caller: Principal, agent: &str, session: &str) -> Option<Json> {
let mut state = self.0.state.lock();
state.terminals.retain(|_, (at, _)| {
unix_ms().saturating_sub(*at)
<= self.limits().terminal_retention_secs.saturating_mul(1000)
});
state
.terminals
.get(&(caller, agent.into(), session.into()))
.map(|(_, detail)| detail.clone())
}
}
pub(crate) struct ScopePermit {
scope: SubAgentScope,
session: bool,
_request: Option<tokio::sync::OwnedSemaphorePermit>,
}
impl Drop for ScopePermit {
fn drop(&mut self) {
if self.session {
self.scope.0.state.lock().sessions -= 1;
}
}
}