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
9 changes: 5 additions & 4 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 2 additions & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ agentkit-plugins = "=0.10.7"
agentkit-provider-openai = "=0.10.8"
agentkit-provider-openrouter = "=0.10.8"
agentkit-task-manager = "=0.10.7"
agentkit-tool-compose = { version = "=0.10.10", default-features = false, features = ["runlet"] }
agentkit-tool-compose = { version = "=0.10.11", default-features = false, features = ["runlet"] }
agentkit-tool-skills = "=0.10.8"
agentkit-tools-core = "=0.10.5"
async-trait = "=0.1.92"
Expand Down Expand Up @@ -56,6 +56,7 @@ rmcp = { version = "=3.2.0", default-features = false, features = ["auth", "clie
serde = { version = "=1.0.229", features = ["derive"] }
serde_json = "=1.0.151"
shlex = "=2.0.1"
runlet = "=0.6.0"
sha2 = "=0.11.0"
subtle = "=2.6.1"
tar = { version = "=0.4.46", default-features = false }
Expand Down
176 changes: 161 additions & 15 deletions src/acp_child.rs
Original file line number Diff line number Diff line change
Expand Up @@ -840,6 +840,8 @@ fn harness_diagnostic(label: &str, line: &str) -> Option<String> {
Some(
crate::events::RuntimeEvent::ChildStarted { .. }
| crate::events::RuntimeEvent::ChildFinished { .. }
| crate::events::RuntimeEvent::RunletProgress { .. }
| crate::events::RuntimeEvent::RunletTransport { .. }
)
) {
return None;
Expand Down Expand Up @@ -898,8 +900,15 @@ async fn run(
&label,
ancestor_id.as_deref(),
|output| match output {
ForwardedStderr::RuntimeLine(line) | ForwardedStderr::Diagnostic(line) => {
eprintln!("{line}");
ForwardedStderr::RuntimeLine(line) => {
if let Some(transport) = crate::runlet_progress::transport::global() {
transport.publish_runtime_line(&line);
}
}
ForwardedStderr::Diagnostic(line) => {
if let Some(transport) = crate::runlet_progress::transport::global() {
transport.publish_line(&line);
}
}
ForwardedStderr::Cleanup(event) => crate::events::emit(&event),
},
Expand Down Expand Up @@ -1256,20 +1265,71 @@ async fn forward_stderr(
.map(str::to_owned)
.collect::<BTreeSet<_>>();
let mut lines = BufReader::new(stderr).lines();
while let Ok(Some(line)) = lines.next_line().await {
if let Some(event) = crate::events::parse(&line)
&& event.forward_from_child()
{
if let crate::events::RuntimeEvent::SubagentStateChanged {
parent_id: Some(parent_id),
..
} = event
{
ancestors.insert(parent_id);
let mut deadline = None;
let mut unavailable = false;
loop {
// A nested Kit transport has its own lease. The parent's healthy
// heartbeat cannot certify a stalled or failed descendant publisher.
let next = if let Some(expires) = deadline {
if tokio::time::Instant::now() >= expires {
Err(())
} else {
tokio::time::timeout_at(expires, lines.next_line())
.await
.map_err(|_| ())
}
} else {
Ok(lines.next_line().await)
};
let line = match next {
Ok(Ok(Some(line))) => line,
Ok(_) => break,
Err(()) => {
unavailable = true;
deadline = None;
output(ForwardedStderr::Cleanup(
crate::events::RuntimeEvent::RunletTransport { available: false },
));
continue;
}
};
if let Some(event) = crate::events::parse(&line) {
if unavailable {
continue;
}
match event {
crate::events::RuntimeEvent::RunletTransport { available: true } => {
deadline = Some(
tokio::time::Instant::now() + crate::runlet_progress::transport::LEASE,
);
continue;
}
crate::events::RuntimeEvent::RunletTransport { available: false } => {
unavailable = true;
deadline = None;
output(ForwardedStderr::RuntimeLine(line));
continue;
}
_ => {}
}
if deadline.is_some() {
deadline =
Some(tokio::time::Instant::now() + crate::runlet_progress::transport::LEASE);
}
if event.forward_from_child() {
if let crate::events::RuntimeEvent::SubagentStateChanged {
parent_id: Some(parent_id),
..
} = event
{
ancestors.insert(parent_id);
}
// Preserve recursively forwarded private events byte-for-byte.
output(ForwardedStderr::RuntimeLine(line));
continue;
}
// Preserve recursively forwarded private runtime events byte-for-byte.
output(ForwardedStderr::RuntimeLine(line));
} else if let Some(line) = harness_diagnostic(label, &line) {
}
if let Some(line) = harness_diagnostic(label, &line) {
output(ForwardedStderr::Diagnostic(line));
}
}
Expand Down Expand Up @@ -1408,6 +1468,27 @@ mod tests {
)
}

#[tokio::test]
async fn close_with_stalled_stderr() {
if !crate::events::test_support::with_stalled_stderr(
"acp_child::tests::close_with_stalled_stderr",
) {
return;
}
let (mut session, mut requests) = admission_test_session();
session.descendant_parent = Some("ancestor".into());
session.capabilities.session_capabilities.close = Some(Default::default());
let actor = async {
let Some(Request::Close(close)) = requests.recv().await else {
panic!("expected actual child close request");
};
assert_eq!(close.session_id.to_string(), "test");
close.reply.send(Ok(())).unwrap();
};
let (result, ()) = tokio::join!(session.close(), actor);
result.unwrap();
}

#[tokio::test]
async fn actor_events_share_ready_request_fatal_and_completion_backlogs() {
for enabled in [0b011u8, 0b101, 0b110, 0b111] {
Expand Down Expand Up @@ -2671,6 +2752,71 @@ mod tests {
mod forwards_subagent_events {
use super::*;

#[tokio::test(start_paused = true)]
async fn nested_transport_loss_preserves_diagnostics_not_stale_lifecycle() {
use crate::events::{EVENT_MARKER, RuntimeEvent};
use tokio::io::AsyncWriteExt;
for explicit in [false, true] {
let (mut writer, reader) = tokio::io::duplex(4096);
let (tx, mut rx) = mpsc::unbounded_channel();
let forward = tokio::spawn(async move {
forward_stderr(reader, "acp.kit", None, |item| {
tx.send(item).unwrap();
})
.await;
});
let heartbeat = format!(
"{EVENT_MARKER}{}\n",
serde_json::to_string(&RuntimeEvent::RunletTransport { available: true })
.unwrap()
);
let started = RuntimeEvent::ChildStarted {
call: "parent:compose:0".into(),
tool: "shell".into(),
summary: "working".into(),
at: 1,
};
let start = format!(
"{EVENT_MARKER}{}\n",
serde_json::to_string(&started).unwrap()
);
writer
.write_all(format!("{heartbeat}{start}").as_bytes())
.await
.unwrap();
assert!(matches!(
rx.recv().await.unwrap(),
ForwardedStderr::RuntimeLine(_)
));
if explicit {
let reset = format!(
"{EVENT_MARKER}{}\n",
serde_json::to_string(&RuntimeEvent::RunletTransport { available: false })
.unwrap()
);
writer.write_all(reset.as_bytes()).await.unwrap();
} else {
tokio::time::advance(crate::runlet_progress::transport::LEASE).await;
}
let reset = match rx.recv().await.unwrap() {
ForwardedStderr::RuntimeLine(line) => crate::events::parse(&line).unwrap(),
ForwardedStderr::Cleanup(event) => event,
_ => panic!("expected nested invalidation"),
};
assert_eq!(reset, RuntimeEvent::RunletTransport { available: false });
writer
.write_all(format!("{heartbeat}{start}later child error\n").as_bytes())
.await
.unwrap();
assert!(
matches!(rx.recv().await.unwrap(), ForwardedStderr::Diagnostic(line) if line.contains("later child error"))
);
drop(writer);
forward.await.unwrap();
assert!(rx.try_recv().is_err());
}
}

#[tokio::test]
async fn preserves_nested_roster_event_lines_exactly() {
let event = crate::events::RuntimeEvent::SubagentStateChanged {
Expand Down
42 changes: 30 additions & 12 deletions src/events.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,6 @@
//! ACP hosts never see the extra chatter.

use std::{
io::Write,
sync::OnceLock,
time::{SystemTime, UNIX_EPOCH},
};
Expand All @@ -36,6 +35,12 @@ pub const EVENTS_ENV: &str = "KIT_RUNTIME_EVENTS";
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
#[serde(tag = "event", rename_all = "snake_case")]
pub enum RuntimeEvent {
/// Process-wide progress transport lease/reset, not a source execution.
RunletTransport { available: bool },
/// Authoritative, value-free observations owned by an exact compose call.
RunletProgress {
progress: crate::runlet_progress::Progress,
},
/// Process-wide durability state, independent of the active ACP session.
StorageStatus { pending: bool, exhausted: bool },
/// A persisted ACP session was opened by the child runtime.
Expand Down Expand Up @@ -116,8 +121,10 @@ impl RuntimeEvent {
#[must_use]
pub fn parent_call(&self) -> Option<&str> {
let call = match self {
Self::RunletProgress { progress } => return Some(&progress.owner),
Self::ChildStarted { call, .. } | Self::ChildFinished { call, .. } => call,
Self::StorageStatus { .. }
Self::RunletTransport { .. }
| Self::StorageStatus { .. }
| Self::SessionStarted { .. }
| Self::CompactionStarted { .. }
| Self::CompactionFinished { .. }
Expand All @@ -135,25 +142,33 @@ pub fn enabled() -> bool {
*ENABLED.get_or_init(|| std::env::var_os(EVENTS_ENV).is_some())
}

/// Writes one event to stderr when emission is enabled.
/// Enqueues one event without waiting for stderr. Loss disables the transport;
/// its explicit reset (or the client lease on a stalled sink) hides source state.
pub fn emit(event: &RuntimeEvent) {
if !enabled() {
return;
}
let mut stderr = std::io::stderr().lock();
write_event(&mut stderr, event);
}

fn write_event(writer: &mut impl Write, event: &RuntimeEvent) {
if let Ok(line) = serde_json::to_string(event) {
let _ = writeln!(writer, "{EVENT_MARKER}{line}");
if let Some(transport) = crate::runlet_progress::transport::global() {
transport.publish_event(event);
}
}

/// Parses one stderr line, returning an event when the line carries one.
#[must_use]
pub fn parse(line: &str) -> Option<RuntimeEvent> {
serde_json::from_str(line.strip_prefix(EVENT_MARKER)?).ok()
let body = line.strip_prefix(EVENT_MARKER)?;
if body.len() > 64 * 1024 {
return None;
}
// Existing diagnostic events retain their historical parser shape. The new
// bounded payload is checked before it can reach retained UI state.
let event: RuntimeEvent = serde_json::from_str(body).ok()?;
if let RuntimeEvent::RunletProgress { progress } = &event
&& (body.len() > 4096 || !progress.bounded())
{
return None;
}
Some(event)
}

/// Milliseconds since the Unix epoch, saturating at zero on a broken clock.
Expand Down Expand Up @@ -236,7 +251,7 @@ mod tests {

use super::{
EVENT_MARKER, GenerationOutcome, RuntimeEvent, SubagentStatus, parse, summarize_input,
summarize_output, write_event,
summarize_output, test_support::write_event,
};

#[test]
Expand Down Expand Up @@ -394,3 +409,6 @@ mod tests {
);
}
}

#[cfg(test)]
pub(crate) mod test_support;
Loading
Loading