use bamboo_subagent::{
AgentRef, ChildFrame, ChildLink, ChildOutcome, InboxKind, InboxMessage, MsgId, ParentFrame,
TransportError, TransportResult,
};
use chrono::Utc;
use crate::client::BrokerClient;
use crate::error::{BrokerError, BrokerResult};
pub struct BrokerChildLink {
client: BrokerClient,
child: String,
me: AgentRef,
run_id: Option<MsgId>,
done: bool,
}
impl BrokerChildLink {
pub async fn connect(
endpoint: &str,
parent: AgentRef,
token: &str,
child: impl Into<String>,
) -> BrokerResult<Self> {
let mut client = BrokerClient::connect(endpoint, parent.clone(), token).await?;
client.subscribe().await?;
Ok(Self {
client,
child: child.into(),
me: parent,
run_id: None,
done: false,
})
}
fn msg(
&self,
kind: InboxKind,
body: serde_json::Value,
correlation: Option<MsgId>,
) -> InboxMessage {
InboxMessage {
id: MsgId::new(),
from: self.me.clone(),
kind,
body,
created_at: Utc::now(),
correlation_id: correlation,
}
}
pub async fn send(&mut self, frame: ParentFrame) -> BrokerResult<()> {
match frame {
ParentFrame::Run(spec) => {
let body = serde_json::to_value(spec)
.map_err(|e| BrokerError::Transport(format!("encode RunSpec: {e}")))?;
let m = self.msg(InboxKind::Run, body, None);
self.run_id = Some(m.id.clone());
self.done = false;
self.client.deliver(&self.child, m).await?;
}
ParentFrame::Cancel => {
if let Some(rid) = self.run_id.clone() {
self.client.cancel(&self.child, &rid).await?;
}
}
ParentFrame::Message { text } => {
let m = self.msg(
InboxKind::Steer,
serde_json::json!({ "text": text }),
self.run_id.clone(),
);
self.client.deliver(&self.child, m).await?;
}
ParentFrame::SessionMessage { delivery } => {
let body = serde_json::to_value(delivery).map_err(|e| {
BrokerError::Transport(format!("encode SessionMessageDelivery: {e}"))
})?;
let m = self.msg(InboxKind::Steer, body, self.run_id.clone());
self.client.deliver(&self.child, m).await?;
}
ParentFrame::ApprovalReply { id, approved } => {
let m = self.msg(
InboxKind::ApprovalReply,
serde_json::json!({ "id": id, "approved": approved }),
self.run_id.clone(),
);
self.client.deliver(&self.child, m).await?;
}
}
Ok(())
}
pub async fn next_frame(&mut self) -> BrokerResult<Option<ChildFrame>> {
if self.done {
return Ok(None);
}
loop {
let Some(msg) = self.client.next_message().await else {
return Ok(None);
};
let id = msg.id.clone();
if self.run_id.is_some() && msg.correlation_id != self.run_id {
self.client.ack(id).await.ok();
continue;
}
let frame = match msg.kind {
InboxKind::Event => Some(ChildFrame::Event { event: msg.body }),
InboxKind::SessionMessageAdmitted => {
let confirmation = serde_json::from_value(msg.body).map_err(|e| {
BrokerError::Transport(format!(
"decode SessionMessageAdmissionConfirmation: {e}"
))
})?;
Some(ChildFrame::SessionMessageAdmitted { confirmation })
}
InboxKind::ApprovalRequest => {
let id = msg
.body
.get("id")
.and_then(|v| v.as_str())
.unwrap_or_default()
.to_string();
let body = msg
.body
.get("request")
.cloned()
.unwrap_or_else(|| serde_json::json!({}));
Some(ChildFrame::ApprovalRequest { id, body })
}
InboxKind::Outcome => {
let oc: ChildOutcome = serde_json::from_value(msg.body)
.map_err(|e| BrokerError::Transport(format!("decode ChildOutcome: {e}")))?;
self.done = true;
Some(ChildFrame::Terminal {
status: oc.status,
result: oc.result,
error: oc.error,
transcript: oc.transcript,
})
}
_ => None,
};
self.client.ack(id).await.ok();
if let Some(f) = frame {
return Ok(Some(f));
}
}
}
}
#[async_trait::async_trait]
impl ChildLink for BrokerChildLink {
async fn send(&mut self, frame: ParentFrame) -> TransportResult<()> {
BrokerChildLink::send(self, frame)
.await
.map_err(|e| TransportError::Protocol(format!("broker link send: {e}")))
}
async fn next_frame(&mut self) -> TransportResult<Option<ChildFrame>> {
BrokerChildLink::next_frame(self)
.await
.map_err(|e| TransportError::Protocol(format!("broker link recv: {e}")))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::core::BrokerCore;
use crate::serve::serve_executor;
use crate::server::BrokerServer;
use bamboo_subagent::{EchoExecutor, RunSpec};
use std::sync::Arc;
use std::time::Duration;
use tokio::net::TcpListener;
#[tokio::test]
async fn drives_a_child_run_over_the_bus() {
let dir = tempfile::tempdir().unwrap();
let core = Arc::new(BrokerCore::new(dir.path()));
let server = Arc::new(BrokerServer::new(core, "t"));
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move {
let _ = server.serve(listener).await;
});
let endpoint = format!("ws://{addr}");
let worker_ep = endpoint.clone();
tokio::spawn(async move {
let _ = serve_executor(
&worker_ep,
AgentRef {
session_id: "child".into(),
role: None,
},
"t",
Arc::new(EchoExecutor),
)
.await;
});
let mut link = BrokerChildLink::connect(
&endpoint,
AgentRef {
session_id: "parent".into(),
role: None,
},
"t",
"child",
)
.await
.unwrap();
link.send(ParentFrame::Run(RunSpec {
assignment: "hello world".into(),
logical_session: None,
project_id: None,
reasoning_effort: None,
permission_policy: None,
messages: vec![],
activation_run_id: None,
initial_session_messages: Vec::new(),
secrets: Default::default(),
}))
.await
.unwrap();
let mut events = 0usize;
let mut terminal = None;
loop {
match tokio::time::timeout(Duration::from_secs(5), link.next_frame())
.await
.expect("a frame arrives")
.expect("link ok")
{
Some(ChildFrame::Event { .. }) => events += 1,
Some(ChildFrame::Terminal { status, result, .. }) => {
terminal = Some((status, result));
break;
}
Some(_) => {}
None => break,
}
}
assert!(events >= 1, "expected streamed events, got {events}");
let (status, result) = terminal.expect("a terminal frame");
assert_eq!(status, bamboo_subagent::TerminalStatus::Completed);
assert_eq!(result.as_deref(), Some("echo: hello world"));
assert!(link.next_frame().await.unwrap().is_none());
}
async fn start_broker() -> (String, tempfile::TempDir) {
let dir = tempfile::tempdir().unwrap();
let core = Arc::new(BrokerCore::new(dir.path()));
let server = Arc::new(BrokerServer::new(core, "t"));
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
tokio::spawn(async move {
let _ = server.serve(listener).await;
});
(format!("ws://{addr}"), dir)
}
fn spawn_worker(endpoint: &str, exec: Arc<dyn bamboo_subagent::ChildExecutor>) {
let ep = endpoint.to_string();
tokio::spawn(async move {
let _ = serve_executor(
&ep,
AgentRef {
session_id: "child".into(),
role: None,
},
"t",
exec,
)
.await;
});
}
async fn connect_parent(endpoint: &str) -> BrokerChildLink {
BrokerChildLink::connect(
endpoint,
AgentRef {
session_id: "parent".into(),
role: None,
},
"t",
"child",
)
.await
.expect("connect link")
}
struct SteerEcho;
#[async_trait::async_trait]
impl bamboo_subagent::ChildExecutor for SteerEcho {
async fn run(
&self,
_spec: RunSpec,
events: bamboo_subagent::EventSink,
mut steer: bamboo_subagent::SteerInbox,
_cancel: tokio_util::sync::CancellationToken,
) -> bamboo_subagent::ChildOutcome {
events.emit(serde_json::json!({ "type": "ready" }));
let s = steer.recv().await.unwrap_or_default();
bamboo_subagent::ChildOutcome::completed(format!("steered: {s}"))
}
}
struct TypedSteer;
#[async_trait::async_trait]
impl bamboo_subagent::ChildExecutor for TypedSteer {
async fn run(
&self,
spec: RunSpec,
events: bamboo_subagent::EventSink,
mut steer: bamboo_subagent::SteerInbox,
_cancel: tokio_util::sync::CancellationToken,
) -> bamboo_subagent::ChildOutcome {
assert_eq!(
spec.logical_session,
Some(bamboo_subagent::LogicalSessionIdentity {
session_id: "logical-child".to_string(),
parent_session_id: Some("logical-parent".to_string()),
root_session_id: "logical-root".to_string(),
})
);
events.emit(serde_json::json!({ "type": "ready" }));
let bamboo_subagent::SteerMessage::SessionMessage(delivery) =
steer.recv_message().await.expect("typed delivery")
else {
panic!("expected typed delivery");
};
events.confirm_session_message(bamboo_subagent::SessionMessageAdmissionConfirmation {
target_session_id: delivery.target_session_id,
envelope_id: delivery.envelope.id.as_str().to_string(),
canonical_claim_generation: delivery.canonical_claim_generation,
activation_run_id: delivery.activation_run_id,
});
bamboo_subagent::ChildOutcome::completed("confirmed")
}
}
struct AskApproval;
#[async_trait::async_trait]
impl bamboo_subagent::ChildExecutor for AskApproval {
async fn run(
&self,
_spec: RunSpec,
events: bamboo_subagent::EventSink,
_steer: bamboo_subagent::SteerInbox,
_cancel: tokio_util::sync::CancellationToken,
) -> bamboo_subagent::ChildOutcome {
let approved = match events.host().cloned() {
Some(host) => host
.approval_call(
serde_json::json!({ "tool_name": "Bash", "resource": "rm -rf /" }),
)
.await
.ok()
.and_then(|v| v.get("approved").and_then(|b| b.as_bool()))
.unwrap_or(false),
None => false,
};
bamboo_subagent::ChildOutcome::completed(
if approved { "approved" } else { "denied" }.to_string(),
)
}
}
#[tokio::test]
async fn carries_an_in_band_steer_over_the_bus() {
let (endpoint, _dir) = start_broker().await;
spawn_worker(&endpoint, Arc::new(SteerEcho));
let mut link = connect_parent(&endpoint).await;
link.send(ParentFrame::Run(RunSpec {
assignment: "go".into(),
logical_session: Some(bamboo_subagent::LogicalSessionIdentity {
session_id: "logical-child".to_string(),
parent_session_id: Some("logical-parent".to_string()),
root_session_id: "logical-root".to_string(),
}),
project_id: None,
reasoning_effort: None,
permission_policy: None,
messages: vec![],
activation_run_id: None,
initial_session_messages: Vec::new(),
secrets: Default::default(),
}))
.await
.unwrap();
match tokio::time::timeout(Duration::from_secs(5), link.next_frame())
.await
.expect("a frame")
.expect("ok")
{
Some(ChildFrame::Event { .. }) => {}
other => panic!("expected ready event first, got {other:?}"),
}
link.send(ParentFrame::Message {
text: "turn-left".into(),
})
.await
.unwrap();
let result = loop {
match tokio::time::timeout(Duration::from_secs(5), link.next_frame())
.await
.expect("a frame")
.expect("ok")
{
Some(ChildFrame::Terminal { result, .. }) => break result,
_ => continue,
}
};
assert_eq!(result.as_deref(), Some("steered: turn-left"));
}
#[tokio::test]
async fn carries_typed_session_message_confirmation_over_the_bus() {
let (endpoint, _dir) = start_broker().await;
spawn_worker(&endpoint, Arc::new(TypedSteer));
let mut link = connect_parent(&endpoint).await;
link.send(ParentFrame::Run(RunSpec {
assignment: "go".into(),
logical_session: Some(bamboo_subagent::LogicalSessionIdentity {
session_id: "logical-child".to_string(),
parent_session_id: Some("logical-parent".to_string()),
root_session_id: "logical-root".to_string(),
}),
project_id: None,
reasoning_effort: None,
permission_policy: None,
messages: vec![],
activation_run_id: None,
initial_session_messages: Vec::new(),
secrets: Default::default(),
}))
.await
.unwrap();
assert!(matches!(
tokio::time::timeout(Duration::from_secs(5), link.next_frame())
.await
.unwrap()
.unwrap(),
Some(ChildFrame::Event { .. })
));
let mut envelope =
bamboo_domain::SessionMessageEnvelope::user_input("logical-child", "typed");
envelope.id = bamboo_domain::SessionMessageId::parse("broker-typed-id").unwrap();
link.send(ParentFrame::SessionMessage {
delivery: bamboo_subagent::SessionMessageDelivery {
target_session_id: "logical-child".to_string(),
envelope,
canonical_claim_generation: 9,
activation_run_id: "run-broker-9".to_string(),
activation_policy: bamboo_domain::SessionActivationPolicy::InterruptSpecificWait,
},
})
.await
.unwrap();
match tokio::time::timeout(Duration::from_secs(5), link.next_frame())
.await
.unwrap()
.unwrap()
{
Some(ChildFrame::SessionMessageAdmitted { confirmation }) => {
assert_eq!(confirmation.target_session_id, "logical-child");
assert_eq!(confirmation.envelope_id, "broker-typed-id");
assert_eq!(confirmation.canonical_claim_generation, 9);
assert_eq!(confirmation.activation_run_id, "run-broker-9");
}
other => panic!("expected typed confirmation, got {other:?}"),
}
}
#[tokio::test]
async fn carries_an_approval_round_trip_over_the_bus() {
let (endpoint, _dir) = start_broker().await;
spawn_worker(&endpoint, Arc::new(AskApproval));
let mut link = connect_parent(&endpoint).await;
link.send(ParentFrame::Run(RunSpec {
assignment: "do the dangerous thing".into(),
logical_session: None,
project_id: None,
reasoning_effort: None,
permission_policy: None,
messages: vec![],
activation_run_id: None,
initial_session_messages: Vec::new(),
secrets: Default::default(),
}))
.await
.unwrap();
let mut saw_request = false;
let result = loop {
match tokio::time::timeout(Duration::from_secs(5), link.next_frame())
.await
.expect("a frame")
.expect("ok")
{
Some(ChildFrame::ApprovalRequest { id, body }) => {
saw_request = true;
assert_eq!(body.get("tool_name").and_then(|v| v.as_str()), Some("Bash"));
link.send(ParentFrame::ApprovalReply { id, approved: true })
.await
.unwrap();
}
Some(ChildFrame::Terminal { result, .. }) => break result,
_ => continue,
}
};
assert!(saw_request, "the worker must proxy an approval request up");
assert_eq!(result.as_deref(), Some("approved"));
}
}