diff --git a/.agentworkforce/trajectories/completed/2026-08/traj_u0oum79ywq6o.trace.json b/.agentworkforce/trajectories/completed/2026-08/traj_u0oum79ywq6o.trace.json new file mode 100644 index 000000000..d9b9adef6 --- /dev/null +++ b/.agentworkforce/trajectories/completed/2026-08/traj_u0oum79ywq6o.trace.json @@ -0,0 +1,1651 @@ +{ + "version": "1.0.0", + "id": "e05255c0-d937-4a0e-8091-89e19e999ff2", + "timestamp": "2026-08-17T21:38:26.114Z", + "trajectory": "traj_u0oum79ywq6o", + "files": [ + { + "path": ".agentworkforce/trajectories/active/traj_u0oum79ywq6o/trajectory.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 50, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": ".agentworkforce/trajectories/completed/2026-08/traj_h0wnqnadvrr8/summary.md", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 33, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": ".agentworkforce/trajectories/completed/2026-08/traj_h0wnqnadvrr8/trajectory.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 70, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": ".agentworkforce/trajectories/completed/2026-08/traj_hx6d4m1gz0wt/summary.md", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 33, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": ".agentworkforce/trajectories/completed/2026-08/traj_hx6d4m1gz0wt/trajectory.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 66, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": ".agentworkforce/trajectories/completed/2026-08/traj_vud9qvfm0vwg/summary.md", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 33, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": ".agentworkforce/trajectories/completed/2026-08/traj_vud9qvfm0vwg/trajectory.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 82, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": ".agentworkforce/trajectories/completed/2026-08/traj_zo0vdyhitsbv/summary.md", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 33, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": ".agentworkforce/trajectories/completed/2026-08/traj_zo0vdyhitsbv/trajectory.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 78, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "CHANGELOG.md", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 11, + "end_line": 40, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "crates/broker/src/listen_api.rs", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 80, + "end_line": 93, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 453, + "end_line": 462, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1184, + "end_line": 1207, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1753, + "end_line": 1765, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1773, + "end_line": 1794, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1893, + "end_line": 1907, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 2142, + "end_line": 2159, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 3923, + "end_line": 3998, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 5395, + "end_line": 5439, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "crates/broker/src/node_control.rs", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 598, + "end_line": 608, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 761, + "end_line": 809, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 975, + "end_line": 1020, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "crates/broker/src/relaycast/ws.rs", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 44, + "end_line": 130, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 207, + "end_line": 220, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 232, + "end_line": 395, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 420, + "end_line": 442, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 460, + "end_line": 512, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1254, + "end_line": 1267, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1297, + "end_line": 1319, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1347, + "end_line": 1354, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1665, + "end_line": 2100, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "crates/broker/src/runtime/api.rs", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 5, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 11, + "end_line": 30, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 927, + "end_line": 941, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1031, + "end_line": 1049, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1352, + "end_line": 1368, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "crates/broker/src/runtime/event_loop.rs", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 231, + "end_line": 240, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "crates/broker/src/runtime/fleet.rs", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 9, + "end_line": 17, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 22, + "end_line": 39, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 103, + "end_line": 150, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1388, + "end_line": 1401, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1552, + "end_line": 1628, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1922, + "end_line": 2071, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 2488, + "end_line": 2506, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 3458, + "end_line": 3876, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "crates/broker/src/runtime/init.rs", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 697, + "end_line": 703, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "crates/broker/src/runtime/maintenance.rs", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 14, + "end_line": 20, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 152, + "end_line": 162, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 222, + "end_line": 232, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 706, + "end_line": 727, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "crates/broker/src/runtime/relaycast_events.rs", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 64, + "end_line": 112, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 290, + "end_line": 300, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 311, + "end_line": 317, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 381, + "end_line": 402, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 406, + "end_line": 413, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 461, + "end_line": 470, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 554, + "end_line": 560, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 879, + "end_line": 1092, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1223, + "end_line": 1270, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "crates/broker/src/runtime/tests.rs", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1390, + "end_line": 1515, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "crates/broker/src/runtime/worker_events.rs", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 717, + "end_line": 743, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "crates/broker/src/terminal_control.rs", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 59, + "end_line": 102, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 331, + "end_line": 337, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 346, + "end_line": 352, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 366, + "end_line": 405, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 460, + "end_line": 467, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 510, + "end_line": 523, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 526, + "end_line": 532, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 566, + "end_line": 671, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "crates/broker/src/worker.rs", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 238, + "end_line": 252, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 511, + "end_line": 544, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 589, + "end_line": 622, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 2519, + "end_line": 2677, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/brand/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/broker-darwin-arm64/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/broker-darwin-x64/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/broker-linux-arm64/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/broker-linux-x64/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/broker-win32-x64/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cli/README.md", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 39, + "end_line": 45, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 134, + "end_line": 140, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 165, + "end_line": 173, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cli/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 43, + "end_line": 55, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/agent-relay-mcp.startup.test.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 578, + "end_line": 584, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 593, + "end_line": 599, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 608, + "end_line": 615, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 627, + "end_line": 634, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 636, + "end_line": 648, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/agent-relay-mcp.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 591, + "end_line": 600, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 614, + "end_line": 620, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 624, + "end_line": 637, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 640, + "end_line": 647, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 654, + "end_line": 666, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1080, + "end_line": 1096, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1117, + "end_line": 1124, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 1138, + "end_line": 1145, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/bootstrap.test.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 40, + "end_line": 46, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/commands/fleet-agent.test.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 537, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/commands/fleet-agent.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [] + } + ] + }, + { + "path": "packages/cli/src/cli/commands/fleet.test.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 30, + "end_line": 41, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 157, + "end_line": 208, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 476, + "end_line": 483, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 515, + "end_line": 521, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 544, + "end_line": 550, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 701, + "end_line": 708, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 732, + "end_line": 738, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/commands/fleet.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 16, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 117, + "end_line": 138, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 145, + "end_line": 151, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 171, + "end_line": 177, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 212, + "end_line": 218, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 232, + "end_line": 245, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 380, + "end_line": 583, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/commands/local-agent.test.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 11, + "end_line": 19, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 499, + "end_line": 654, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/commands/local-agent.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 7, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 338, + "end_line": 364, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 375, + "end_line": 472, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 592, + "end_line": 610, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/lib/attach-fleet-node.test.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 7, + "end_line": 13, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 98, + "end_line": 114, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 339, + "end_line": 409, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/lib/attach-input-recovery.test.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 3, + "end_line": 9, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 100, + "end_line": 120, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 144, + "end_line": 213, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cli/src/cli/lib/attach-input-recovery.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 41, + "end_line": 49, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 51, + "end_line": 87, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 150, + "end_line": 162, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 195, + "end_line": 202, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 245, + "end_line": 251, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 264, + "end_line": 286, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/cloud/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 62, + "end_line": 68, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/config/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/evals/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 71, + "end_line": 78, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/fleet/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 26, + "end_line": 33, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/fleet/src/index.test.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 62, + "end_line": 68, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 78, + "end_line": 84, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/fleet/src/index.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 146, + "end_line": 152, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 257, + "end_line": 266, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/fleet/src/serve-node.test.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 234, + "end_line": 240, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 251, + "end_line": 257, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/harness-driver/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 56, + "end_line": 71, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/harness-driver/src/client.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 45, + "end_line": 51, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 708, + "end_line": 746, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/harness-driver/src/list-fleet-inventory.test.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 94, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/harness-driver/src/protocol.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 251, + "end_line": 270, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/harness-driver/src/pty-input-stream.test.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 174, + "end_line": 223, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/harness-driver/src/transport.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 325, + "end_line": 345, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/harness-driver/src/types.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 130, + "end_line": 136, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/harnesses/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 26, + "end_line": 33, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/integration-prompts/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/policy/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 25, + "end_line": 31, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/sdk-py/pyproject.toml", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 4, + "end_line": 10, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/sdk/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/sdk/src/messaging/thin-client.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 56, + "end_line": 63, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/session/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "packages/utils/package.json", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 1, + "end_line": 6, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 109, + "end_line": 115, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + }, + { + "path": "tests/e2e/fleet/fleet-e2e.test.ts", + "conversations": [ + { + "contributor": { + "type": "ai" + }, + "ranges": [ + { + "start_line": 699, + "end_line": 709, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + }, + { + "start_line": 714, + "end_line": 727, + "revision": "0098957c40f900fb79bc6056e883da35a2181272" + } + ] + } + ] + } + ] +} \ No newline at end of file diff --git a/.agentworkforce/trajectories/completed/2026-08/traj_u0oum79ywq6o/summary.md b/.agentworkforce/trajectories/completed/2026-08/traj_u0oum79ywq6o/summary.md new file mode 100644 index 000000000..9cfb9f1a7 --- /dev/null +++ b/.agentworkforce/trajectories/completed/2026-08/traj_u0oum79ywq6o/summary.md @@ -0,0 +1,40 @@ +# Trajectory: Fix restored fleet delivery acknowledgement ordering + +> **Status:** ✅ Completed +> **Task:** relay#1543 +> **Confidence:** 92% +> **Started:** August 17, 2026 at 10:56 PM +> **Completed:** August 17, 2026 at 11:38 PM + +--- + +## Summary + +Implemented restart-safe per-agent ordering for echo-confirmed fleet delivery acknowledgements, including a durable ACK floor and explicit out-of-order, in-order, and second-restart regression coverage. + +**Approach:** Stored worker-confirmed sequences in each fleet cursor, restored pending siblings in sequence order, retained higher confirmations until the lower contiguous prefix lands, persisted the first required sequence across snapshots, and verified the broker suite plus formatting and lint checks. + +--- + +## Key Decisions + +### Restore pending fleet cursors from sequenced siblings and hold confirmed gaps per agent +- **Chose:** Restore pending fleet cursors from sequenced siblings and hold confirmed gaps per agent +- **Reasoning:** Pre-seeding alone drops the higher delivery's eventual ack. The delivery book now retains out-of-order confirmations, while the pending entry remains persisted and maintenance excludes only actively held confirmations; the contiguous prefix releases as one cumulative ack. + +--- + +## Chapters + +### 1. Work +*Agent: default* + +- Restore pending fleet cursors from sequenced siblings and hold confirmed gaps per agent: Restore pending fleet cursors from sequenced siblings and hold confirmed gaps per agent +- The initial in-memory per-agent hold/release closes the immediate out-of-order acknowledgement, but a higher confirmed entry can outlive the lower entry and another restart. Persisting each agent's first still-required sequence alongside every withheld ACK preserves the gap across that second restart; replaying pending siblings in explicit sequence order then releases only a contiguous confirmed prefix. + +--- + +## Artifacts + +**Commits:** 0098957c4, d570c6ee2, 79a69fd91, 5fd61e9fe, 28337030d, 0e4744407, 9f3b24e44, e3217d290, 58198b1ad, 006fd5108, e369f0e03, 6fb4c2f8c +**Files changed:** 67 diff --git a/.agentworkforce/trajectories/completed/2026-08/traj_u0oum79ywq6o/trajectory.json b/.agentworkforce/trajectories/completed/2026-08/traj_u0oum79ywq6o/trajectory.json new file mode 100644 index 000000000..7b5eec1c5 --- /dev/null +++ b/.agentworkforce/trajectories/completed/2026-08/traj_u0oum79ywq6o/trajectory.json @@ -0,0 +1,162 @@ +{ + "id": "traj_u0oum79ywq6o", + "version": 1, + "task": { + "title": "Fix restored fleet delivery acknowledgement ordering", + "source": { + "system": "plain", + "id": "relay#1543" + } + }, + "status": "completed", + "startedAt": "2026-08-17T20:56:37.912Z", + "completedAt": "2026-08-17T21:38:25.571Z", + "agents": [ + { + "name": "default", + "role": "lead", + "joinedAt": "2026-08-17T21:08:44.349Z" + } + ], + "chapters": [ + { + "id": "chap_06fnqnxotn6e", + "title": "Work", + "agentName": "default", + "startedAt": "2026-08-17T21:08:44.349Z", + "endedAt": "2026-08-17T21:38:25.571Z", + "events": [ + { + "ts": 1787000924351, + "type": "decision", + "content": "Restore pending fleet cursors from sequenced siblings and hold confirmed gaps per agent: Restore pending fleet cursors from sequenced siblings and hold confirmed gaps per agent", + "raw": { + "question": "Restore pending fleet cursors from sequenced siblings and hold confirmed gaps per agent", + "chosen": "Restore pending fleet cursors from sequenced siblings and hold confirmed gaps per agent", + "alternatives": [], + "reasoning": "Pre-seeding alone drops the higher delivery's eventual ack. The delivery book now retains out-of-order confirmations, while the pending entry remains persisted and maintenance excludes only actively held confirmations; the contiguous prefix releases as one cumulative ack." + }, + "significance": "high" + }, + { + "ts": 1787002697468, + "type": "reflection", + "content": "The initial in-memory per-agent hold/release closes the immediate out-of-order acknowledgement, but a higher confirmed entry can outlive the lower entry and another restart. Persisting each agent's first still-required sequence alongside every withheld ACK preserves the gap across that second restart; replaying pending siblings in explicit sequence order then releases only a contiguous confirmed prefix.", + "raw": { + "focalPoints": [ + "restart ordering", + "terminal lower failure", + "durable ACK floor", + "explicit boundary tests" + ], + "adjustments": "Added a persisted per-agent ACK floor, normalized it on snapshot load, propagated the live cursor floor at insertion, and added a second-restart regression test.", + "confidence": 0.92 + }, + "significance": "high", + "tags": [ + "focal:restart ordering", + "focal:terminal lower failure", + "focal:durable ACK floor", + "focal:explicit boundary tests", + "confidence:0.92" + ] + } + ] + } + ], + "retrospective": { + "summary": "Implemented restart-safe per-agent ordering for echo-confirmed fleet delivery acknowledgements, including a durable ACK floor and explicit out-of-order, in-order, and second-restart regression coverage.", + "approach": "Stored worker-confirmed sequences in each fleet cursor, restored pending siblings in sequence order, retained higher confirmations until the lower contiguous prefix lands, persisted the first required sequence across snapshots, and verified the broker suite plus formatting and lint checks.", + "confidence": 0.92 + }, + "commits": [ + "0098957c4", + "d570c6ee2", + "79a69fd91", + "5fd61e9fe", + "28337030d", + "0e4744407", + "9f3b24e44", + "e3217d290", + "58198b1ad", + "006fd5108", + "e369f0e03", + "6fb4c2f8c" + ], + "filesChanged": [ + ".agentworkforce/trajectories/active/traj_u0oum79ywq6o/trajectory.json", + ".agentworkforce/trajectories/completed/2026-08/traj_h0wnqnadvrr8/summary.md", + ".agentworkforce/trajectories/completed/2026-08/traj_h0wnqnadvrr8/trajectory.json", + ".agentworkforce/trajectories/completed/2026-08/traj_hx6d4m1gz0wt/summary.md", + ".agentworkforce/trajectories/completed/2026-08/traj_hx6d4m1gz0wt/trajectory.json", + ".agentworkforce/trajectories/completed/2026-08/traj_vud9qvfm0vwg/summary.md", + ".agentworkforce/trajectories/completed/2026-08/traj_vud9qvfm0vwg/trajectory.json", + ".agentworkforce/trajectories/completed/2026-08/traj_zo0vdyhitsbv/summary.md", + ".agentworkforce/trajectories/completed/2026-08/traj_zo0vdyhitsbv/trajectory.json", + "CHANGELOG.md", + "crates/broker/src/listen_api.rs", + "crates/broker/src/node_control.rs", + "crates/broker/src/relaycast/ws.rs", + "crates/broker/src/runtime/api.rs", + "crates/broker/src/runtime/event_loop.rs", + "crates/broker/src/runtime/fleet.rs", + "crates/broker/src/runtime/init.rs", + "crates/broker/src/runtime/maintenance.rs", + "crates/broker/src/runtime/relaycast_events.rs", + "crates/broker/src/runtime/tests.rs", + "crates/broker/src/runtime/worker_events.rs", + "crates/broker/src/terminal_control.rs", + "crates/broker/src/worker.rs", + "package.json", + "packages/brand/package.json", + "packages/broker-darwin-arm64/package.json", + "packages/broker-darwin-x64/package.json", + "packages/broker-linux-arm64/package.json", + "packages/broker-linux-x64/package.json", + "packages/broker-win32-x64/package.json", + "packages/cli/README.md", + "packages/cli/package.json", + "packages/cli/src/cli/agent-relay-mcp.startup.test.ts", + "packages/cli/src/cli/agent-relay-mcp.ts", + "packages/cli/src/cli/bootstrap.test.ts", + "packages/cli/src/cli/commands/fleet-agent.test.ts", + "packages/cli/src/cli/commands/fleet-agent.ts", + "packages/cli/src/cli/commands/fleet.test.ts", + "packages/cli/src/cli/commands/fleet.ts", + "packages/cli/src/cli/commands/local-agent.test.ts", + "packages/cli/src/cli/commands/local-agent.ts", + "packages/cli/src/cli/lib/attach-fleet-node.test.ts", + "packages/cli/src/cli/lib/attach-input-recovery.test.ts", + "packages/cli/src/cli/lib/attach-input-recovery.ts", + "packages/cloud/package.json", + "packages/config/package.json", + "packages/evals/package.json", + "packages/fleet/package.json", + "packages/fleet/src/index.test.ts", + "packages/fleet/src/index.ts", + "packages/fleet/src/serve-node.test.ts", + "packages/harness-driver/package.json", + "packages/harness-driver/src/client.ts", + "packages/harness-driver/src/list-fleet-inventory.test.ts", + "packages/harness-driver/src/protocol.ts", + "packages/harness-driver/src/pty-input-stream.test.ts", + "packages/harness-driver/src/transport.ts", + "packages/harness-driver/src/types.ts", + "packages/harnesses/package.json", + "packages/integration-prompts/package.json", + "packages/policy/package.json", + "packages/sdk-py/pyproject.toml", + "packages/sdk/package.json", + "packages/sdk/src/messaging/thin-client.ts", + "packages/session/package.json", + "packages/utils/package.json", + "tests/e2e/fleet/fleet-e2e.test.ts" + ], + "projectId": "AgentWorkforce/relay", + "tags": [], + "_trace": { + "startRef": "f8ee6e7da1859f64599ae45fab8166e4d1809e9c", + "endRef": "0098957c40f900fb79bc6056e883da35a2181272", + "traceId": "e05255c0-d937-4a0e-8091-89e19e999ff2" + } +} \ No newline at end of file diff --git a/CHANGELOG.md b/CHANGELOG.md index 895fc61f0..2697497cc 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,7 +5,11 @@ All notable changes to Agent Relay will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). -## [Unreleased] +## [Unreleased - Patch] + +### Fixed + +- Fleet node delivery acknowledgements now wait for worker receipt and preserve per-agent sequence order across broker restarts, preventing cumulative acknowledgements from covering lower undelivered messages. ## [11.7.0] - 2026-08-17 diff --git a/crates/broker/src/node_control.rs b/crates/broker/src/node_control.rs index 7f776d8de..604dedab3 100644 --- a/crates/broker/src/node_control.rs +++ b/crates/broker/src/node_control.rs @@ -598,6 +598,11 @@ struct AgentDeliveryCursor { acked_up_to_seq: u64, received_up_to_seq: u64, seen_msg_ids: SeenMsgIds, + /// Worker-confirmed sequenced deliveries that cannot advance the + /// cumulative ACK yet because a lower sequence is still unconfirmed. + /// Kept on the per-agent cursor so confirmations can be released in order + /// without letting one agent block another. + confirmed_delivery_seqs: BTreeMap, /// Whether a sequenced (`seq >= 1`) position has been established for this /// identity, either by a resume handshake seeding the cursor or by adopting /// the first sequenced delivery. Tracked separately from the cursor's @@ -756,11 +761,50 @@ impl FleetDeliveryBook { acked_up_to_seq: up_to_seq, received_up_to_seq: up_to_seq, seen_msg_ids: SeenMsgIds::default(), + confirmed_delivery_seqs: BTreeMap::new(), has_sequenced_position: true, }, ); } + /// Reconstruct the received cursor for a post-restart agent from the + /// withheld fleet deliveries that survived in the pending snapshot. + /// + /// The persisted ACK floor establishes the first still-required sequence, + /// even when that lower delivery has already left the pending map after a + /// terminal failure. Replaying the remaining snapshot in sequence order + /// then restores as much of the received frontier as is contiguous. A live + /// cursor learned from normal delivery or a resume handshake is never + /// rewound; pending siblings can only extend its received frontier. + pub(crate) fn restore_pending_agent(&mut self, deliveries: &[Deliver], ack_floor: Option) { + let mut sequenced = deliveries + .iter() + .filter(|deliver| deliver.seq > 0) + .collect::>(); + sequenced.sort_by_key(|deliver| deliver.seq); + let Some(first) = sequenced.first().copied() else { + return; + }; + let has_sequenced_position = self + .agents + .get(first.agent_id.as_str()) + .is_some_and(|cursor| cursor.has_sequenced_position); + if !has_sequenced_position { + if !self.bind_identity(&first.agent, &first.agent_id, false) { + return; + } + let first_required_seq = ack_floor.unwrap_or(first.seq).min(first.seq); + self.seed_cursor( + first.agent.clone(), + first.agent_id.clone(), + first_required_seq.saturating_sub(1), + ); + } + for deliver in sequenced { + self.commit_received(deliver); + } + } + fn active_up_to_seq(&self, agent: &str) -> u64 { self.active_agent_bindings_by_name .get(agent) @@ -932,9 +976,49 @@ impl FleetDeliveryBook { return None; } cursor.acked_up_to_seq = receipt.seq; + cursor.confirmed_delivery_seqs.remove(&receipt.seq); Some(cursor.acked_up_to_seq) } + /// Record worker confirmation of a surfaced fleet delivery, advancing the + /// cumulative ACK only across a contiguous confirmed prefix. An + /// out-of-order confirmation stays held on this agent's cursor until every + /// lower received sequence has also confirmed. + pub(crate) fn commit_confirmed_delivery(&mut self, deliver: &Deliver) -> Option { + self.commit_received(deliver); + let cursor = self.agents.get_mut(deliver.agent_id.as_str())?; + if deliver.seq == 0 { + cursor.seen_msg_ids.insert(&deliver.msg_id); + return Some(cursor.acked_up_to_seq); + } + if deliver.seq <= cursor.acked_up_to_seq { + return None; + } + cursor.confirmed_delivery_seqs.insert(deliver.seq, ()); + if deliver.seq > cursor.received_up_to_seq { + return None; + } + let before = cursor.acked_up_to_seq; + loop { + let next = cursor.acked_up_to_seq.saturating_add(1); + if next > cursor.received_up_to_seq + || cursor.confirmed_delivery_seqs.remove(&next).is_none() + { + break; + } + cursor.acked_up_to_seq = next; + } + (cursor.acked_up_to_seq > before).then_some(cursor.acked_up_to_seq) + } + + pub(crate) fn is_delivery_confirmation_held(&self, deliver: &Deliver) -> bool { + deliver.seq > 0 + && self + .agents + .get(deliver.agent_id.as_str()) + .is_some_and(|cursor| cursor.confirmed_delivery_seqs.contains_key(&deliver.seq)) + } + pub(crate) fn can_ack_receipt(&self, receipt: &RelaycastDeliveryReceipt) -> bool { let Some(cursor) = self.agents.get(receipt.agent_id.as_str()) else { return false; @@ -963,6 +1047,13 @@ impl FleetDeliveryBook { .map_or(0, |cursor| cursor.acked_up_to_seq) } + pub(crate) fn next_ack_seq(&self, agent_id: &str) -> Option { + self.agents + .get(agent_id) + .filter(|cursor| cursor.has_sequenced_position) + .map(|cursor| cursor.acked_up_to_seq.saturating_add(1)) + } + #[cfg(test)] pub(crate) fn received_up_to_seq(&self, agent_id: &str) -> u64 { self.agents diff --git a/crates/broker/src/runtime/dead_letter.rs b/crates/broker/src/runtime/dead_letter.rs index 8866590d0..113c4adc1 100644 --- a/crates/broker/src/runtime/dead_letter.rs +++ b/crates/broker/src/runtime/dead_letter.rs @@ -207,6 +207,11 @@ pub(crate) fn requeue_dead_letter( next_retry_at: Instant::now(), queued_at_ms: entry.queued_at_ms, last_error: None, + // The dead-lettered entry's withheld fleet ack (if any) was already + // dropped when it was dead-lettered — see relay#1310. A requeue is a + // fresh redelivery attempt, not a continuation of that withheld ack. + withheld_fleet_ack: None, + withheld_fleet_ack_floor: None, }; pending_deliveries.insert(pending.delivery.delivery_id.clone(), pending.clone()); Some(pending) diff --git a/crates/broker/src/runtime/delivery.rs b/crates/broker/src/runtime/delivery.rs index b15e190d9..d7ef928b3 100644 --- a/crates/broker/src/runtime/delivery.rs +++ b/crates/broker/src/runtime/delivery.rs @@ -12,6 +12,24 @@ pub(crate) struct PendingDelivery { pub(super) next_retry_at: Instant, pub(super) queued_at_ms: u64, pub(super) last_error: Option, + /// Fleet (engine-facing) `delivery_ack` withheld until the worker confirms + /// this specific PTY injection landed — echo-verified, or its bounded + /// timeout fallback — rather than acked the instant the write is merely + /// handed to the worker. See relay#1310. + /// + /// Lives on the `PendingDelivery` itself, not a second map keyed by + /// `DeliveryId`, so it cannot outlive the delivery it belongs to: every + /// path that disposes of a `PendingDelivery` (echo confirmation, dead + /// letter, worker teardown) disposes of its withheld ack with it, by + /// construction, instead of needing a matching removal remembered at + /// every one of those call sites. See relay#1543. + pub(super) withheld_fleet_ack: Option, + /// Lowest sequenced delivery that must be confirmed before this withheld + /// cumulative ACK may advance. Persisted independently of the lower + /// delivery entry so a terminal failure followed by another broker + /// restart cannot make a higher pending sequence look like a safe new + /// baseline. + pub(super) withheld_fleet_ack_floor: Option, } /// Serializable snapshot of pending deliveries for crash recovery. @@ -26,6 +44,15 @@ pub(crate) struct PersistedPendingDelivery { pub(super) queued_at_ms: u64, #[serde(default)] pub(super) last_error: Option, + /// See `PendingDelivery::withheld_fleet_ack`. `#[serde(default)]` so a + /// snapshot written before this field existed (or by a broker version + /// that predates it) deserializes as `None` instead of failing to load + /// — the same "nothing withheld" state that field already gets from a + /// fresh delivery. See relay#1543's restart-persistence follow-up. + #[serde(default)] + pub(super) withheld_fleet_ack: Option, + #[serde(default)] + pub(super) withheld_fleet_ack_floor: Option, } #[derive(Debug, Clone, PartialEq)] @@ -139,6 +166,8 @@ pub(crate) fn save_pending_deliveries( failed_attempts: pd.failed_attempts, queued_at_ms: pd.queued_at_ms, last_error: pd.last_error.clone(), + withheld_fleet_ack: pd.withheld_fleet_ack.clone(), + withheld_fleet_ack_floor: pd.withheld_fleet_ack_floor, }) .collect(); crate::util::fs::write_json_atomic(path, &persisted) @@ -153,7 +182,7 @@ pub(crate) fn load_pending_deliveries(path: &Path) -> HashMap v, Err(_) => return HashMap::new(), }; - persisted + let mut loaded = persisted .into_iter() .map(|p| { let id = p.delivery.delivery_id.clone(); @@ -171,10 +200,55 @@ pub(crate) fn load_pending_deliveries(path: &Path) -> HashMap) { + let mut floors = HashMap::::new(); + for pending in deliveries.values() { + let Some(deliver) = pending + .withheld_fleet_ack + .as_ref() + .filter(|deliver| deliver.seq > 0) + else { + continue; + }; + let candidate = pending + .withheld_fleet_ack_floor + .unwrap_or(deliver.seq) + .min(deliver.seq); + floors + .entry(deliver.agent_id.clone()) + .and_modify(|floor| *floor = (*floor).min(candidate)) + .or_insert(candidate); + } + for pending in deliveries.values_mut() { + let Some(deliver) = pending + .withheld_fleet_ack + .as_ref() + .filter(|deliver| deliver.seq > 0) + else { + pending.withheld_fleet_ack_floor = None; + continue; + }; + pending.withheld_fleet_ack_floor = floors.get(&deliver.agent_id).copied(); + } } // These payload structs were used by the stdio protocol handler (handle_sdk_frame). @@ -585,7 +659,17 @@ pub(crate) async fn try_inject_pending_relay_message( worker_name: &str, msg: &PendingRelayMessage, retry_interval: Duration, -) -> Result<()> { + // Fleet-originated deliveries pass the withheld engine ack through so it + // is embedded into the `PendingDelivery` at the moment of insertion — + // synchronously, before the handoff attempt below can time out. Deliver + // it any later (e.g. as a follow-up step keyed off this function's + // return value) and a handoff that outlives `retry_interval` loses the + // race: the timeout below fires, the `DeliveryId` never reaches the + // caller, and the ack is never registered even though the delivery is + // still very much alive and retryable. See relay#1310 / relay#1543. + withheld_fleet_ack: Option, + withheld_fleet_ack_floor: Option, +) -> Result { let event_id = msg .event_id .clone() @@ -610,6 +694,8 @@ pub(crate) async fn try_inject_pending_relay_message( msg.priority, msg.mode.clone(), retry_interval, + withheld_fleet_ack, + withheld_fleet_ack_floor, ), ) .await @@ -679,7 +765,9 @@ pub(crate) async fn queue_and_try_delivery_raw( priority: u8, injection_mode: MessageInjectionMode, retry_interval: Duration, -) -> Result<()> { + withheld_fleet_ack: Option, + withheld_fleet_ack_floor: Option, +) -> Result { let delivery = RelayDelivery { delivery_id: DeliveryId::new(format!("del_{}", Uuid::new_v4().simple())), event_id: EventId::new(event_id), @@ -692,7 +780,59 @@ pub(crate) async fn queue_and_try_delivery_raw( priority: Some(priority), injection_mode, }; + insert_and_attempt_delivery( + workers, + pending_deliveries, + worker_name, + delivery, + retry_interval, + withheld_fleet_ack, + withheld_fleet_ack_floor, + ) + .await +} + +/// Register a delivery and make its first handoff attempt, in one atomic +/// step: the `PendingDelivery` — including any withheld fleet ack — is +/// inserted into `pending_deliveries` before the handoff attempt starts, so +/// a slow or cancelled attempt can never separate "this delivery exists and +/// is retryable" from "its withheld ack is registered". Shared by the +/// broker-generated-id path (`queue_and_try_delivery_raw`) and any caller +/// that already has a fully-built [`RelayDelivery`] (the fleet +/// `WorkerMissing` injection path, which must keep the engine's own +/// `delivery_id`). +pub(crate) async fn insert_and_attempt_delivery( + workers: &mut WorkerRegistry, + pending_deliveries: &mut HashMap, + worker_name: &str, + delivery: RelayDelivery, + retry_interval: Duration, + withheld_fleet_ack: Option, + explicit_withheld_fleet_ack_floor: Option, +) -> Result { let delivery_id = delivery.delivery_id.clone(); + let withheld_fleet_ack_floor = withheld_fleet_ack + .as_ref() + .filter(|deliver| deliver.seq > 0) + .map(|deliver| { + pending_deliveries + .values() + .filter_map(|pending| { + let sibling = pending.withheld_fleet_ack.as_ref()?; + (sibling.agent_id == deliver.agent_id && sibling.seq > 0).then_some( + pending + .withheld_fleet_ack_floor + .unwrap_or(sibling.seq) + .min(sibling.seq), + ) + }) + .fold( + explicit_withheld_fleet_ack_floor + .unwrap_or(deliver.seq) + .min(deliver.seq), + u64::min, + ) + }); pending_deliveries.insert( delivery_id.clone(), PendingDelivery { @@ -703,6 +843,8 @@ pub(crate) async fn queue_and_try_delivery_raw( next_retry_at: Instant::now(), queued_at_ms: unix_timestamp_millis(), last_error: None, + withheld_fleet_ack, + withheld_fleet_ack_floor, }, ); @@ -717,7 +859,7 @@ pub(crate) async fn queue_and_try_delivery_raw( pending_deliveries.insert(pending.delivery.delivery_id.clone(), *pending); anyhow::bail!(last_error); } - Ok(()) + Ok(delivery_id) } pub(crate) async fn retry_pending_delivery( @@ -857,6 +999,23 @@ pub(crate) async fn emit_delivery_attempt_outcome( }, ) .await; + // A dead-lettered delivery never actually landed, so any fleet + // (engine-facing) ack withheld pending its confirmation must be + // dropped rather than sent — the engine keeps its own record of + // this delivery as un-acked and will redeliver it. See relay#1310: + // the whole point of withholding the ack is that "enqueued for + // injection" must not be reported the same as "delivered". The + // withheld ack lives on `pending` itself, so it is dropped here + // simply by `pending` going out of scope — nothing to remember to + // clean up separately. See relay#1543. + if pending.withheld_fleet_ack.is_some() { + tracing::info!( + target = "relay_broker::fleet", + worker = %pending.worker_name, + delivery_id = %pending.delivery.delivery_id, + "dropping withheld fleet delivery_ack for dead-lettered delivery" + ); + } dead_letter_pending_delivery(sdk_out_tx, dead_letters, &pending, &last_error).await; } DeliveryAttemptOutcome::Noop => {} @@ -888,6 +1047,11 @@ pub(crate) fn take_pending_for_worker( .collect() } +/// Choke point for every worker-exit / teardown disposition (agent release, +/// permanent worker death, unsupervised exit): whatever removed these +/// `PendingDelivery`s from `pending_deliveries` (via [`take_pending_for_worker`]) +/// already carried their withheld fleet acks along as a struct field, so this +/// is also the single place that drops them. See relay#1543. pub(crate) async fn emit_dropped_delivery_failures( sdk_out_tx: &mpsc::Sender>, dead_letters: &mut DeadLetterStore, @@ -895,6 +1059,15 @@ pub(crate) async fn emit_dropped_delivery_failures( reason: &str, ) -> Result<()> { for pending in dropped { + if pending.withheld_fleet_ack.is_some() { + tracing::info!( + target = "relay_broker::fleet", + worker = %pending.worker_name, + delivery_id = %pending.delivery.delivery_id, + reason = reason, + "dropping withheld fleet delivery_ack for a delivery dropped from the pending map" + ); + } // Notify best-effort: a send failure must not `?`-abort the loop and // strand the remaining dropped deliveries out of the dead-letter store. // The DLQ capture below runs regardless of the send's outcome. diff --git a/crates/broker/src/runtime/fleet.rs b/crates/broker/src/runtime/fleet.rs index fc5e2a8a3..b4af8a660 100644 --- a/crates/broker/src/runtime/fleet.rs +++ b/crates/broker/src/runtime/fleet.rs @@ -177,9 +177,17 @@ pub(super) fn verified_spawn_failed_result(invocation_id: String, error: &str) - } } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] +#[derive(Debug, Clone, PartialEq, Eq)] enum FleetDeliverySurfaceOutcome { + /// Ack the engine now: no PTY injection occurred (ambient receipt/reaction + /// or an unrecognized payload type), so there is nothing to verify. Acknowledge, + /// A PTY injection was handed to the worker. The engine ack is withheld + /// until the worker confirms it landed (echo-verified, or the bounded + /// timeout fallback) — see relay#1310. The withheld ack was already + /// registered on the corresponding `PendingDelivery` at insertion time + /// (see relay#1543), so there is nothing left to carry here. + AcknowledgeAfterEcho, HoldForManualFlush, } @@ -769,6 +777,25 @@ impl BrokerRuntime { Ok(FleetDeliverySurfaceOutcome::Acknowledge) => { self.fleet_delivery_book.commit_delivered(&deliver) } + Ok(FleetDeliverySurfaceOutcome::AcknowledgeAfterEcho) => { + // Advance the *received* cursor now — the broker has taken + // ownership of this message and the engine's next sequenced + // frame must not be treated as a gap while injection is + // in flight. The *acked* cursor (and the engine-facing + // delivery_ack below) is withheld until + // `handle_worker_event` observes the worker actually + // confirm this specific injection (echo or timeout + // fallback) — see relay#1310. + // + // The withheld ack itself was already registered on the + // `PendingDelivery` at insertion time, before this + // injection attempt even started (see + // `try_inject_pending_relay_message` / + // `insert_and_attempt_delivery`), so there is nothing left + // to record here — see relay#1543. + self.fleet_delivery_book.commit_received(&deliver); + return; + } Ok(FleetDeliverySurfaceOutcome::HoldForManualFlush) => { self.fleet_delivery_book.commit_received(&deliver); return; @@ -830,6 +857,9 @@ impl BrokerRuntime { // `ListenApiRequest::Send`, so if manual_flush isn't honored // here, it's never honored anywhere. FleetDeliverySurfacing::Inject => { + let withheld_fleet_ack_floor = self + .fleet_delivery_book + .next_ack_seq(deliver.agent_id.as_str()); let fields = fleet_delivery_fields(&deliver.payload, &deliver.agent); // Mirror the `relay_inbound` dashboard event that the HTTP @@ -962,44 +992,83 @@ impl BrokerRuntime { // engine to redeliver it); backlog injection failures // are logged and otherwise don't block the ack, since // their own delivery frames already governed their acks. - let mut current_result = Ok(()); + let mut current_result: Result<(), anyhow::Error> = Ok(()); + let mut current_injected = false; for queued in to_drain { let is_current = queued.event_id.as_deref() == Some(deliver.msg_id.as_str()); - if let Err(error) = try_inject_pending_relay_message( + // Only the current delivery's own frame carries the + // withheld engine ack forward — backlog messages + // drained alongside it are governed by their own + // delivery frames (see the comment above this loop). + match try_inject_pending_relay_message( &mut self.workers, &mut self.pending_deliveries, &deliver.agent, &queued, self.delivery_retry_interval, + is_current.then(|| deliver.clone()), + is_current.then_some(withheld_fleet_ack_floor).flatten(), ) .await { - if is_current { - current_result = Err(error); - } else { - tracing::warn!( - target = "relay_broker::fleet", - agent = %deliver.agent, - from = %queued.from, - error = %error, - "failed to inject drained backlog message" - ); + Ok(_delivery_id) => { + if is_current { + current_injected = true; + } + } + Err(error) => { + if is_current { + current_result = Err(error); + } else { + tracing::warn!( + target = "relay_broker::fleet", + agent = %deliver.agent, + from = %queued.from, + error = %error, + "failed to inject drained backlog message" + ); + } } } } - current_result.map(|()| FleetDeliverySurfaceOutcome::Acknowledge) + current_result.and_then(|()| { + if current_injected { + Ok(FleetDeliverySurfaceOutcome::AcknowledgeAfterEcho) + } else { + Err(anyhow::anyhow!( + "current delivery '{}' missing from drain batch", + deliver.msg_id + )) + } + }) } InboundQueueOutcome::RejectedFull => anyhow::bail!( "manual delivery queue is full for '{}'; retaining Relaycast ownership", deliver.agent ), InboundQueueOutcome::WorkerMissing => { + // Route through the same pending-delivery/retry path as + // `DrainNow` (`insert_and_attempt_delivery`) instead of + // a bare one-shot `workers.deliver`, so this injection + // also gets a `PendingDelivery` entry: its withheld ack + // is registered at insertion time and is guaranteed to + // be cleaned up by worker-teardown / retry-exhaustion + // paths if the worker disappears before echoing. A + // one-shot `workers.deliver` outside `pending_deliveries` + // had no such guarantee — see relay#1543. let relay_delivery = self.fleet_relay_delivery(deliver); - self.workers - .deliver(&deliver.agent, relay_delivery) - .await - .map(|()| FleetDeliverySurfaceOutcome::Acknowledge) + insert_and_attempt_delivery( + &mut self.workers, + &mut self.pending_deliveries, + &deliver.agent, + relay_delivery, + self.delivery_retry_interval, + Some(deliver.clone()), + withheld_fleet_ack_floor, + ) + .await + .map(|_delivery_id| FleetDeliverySurfaceOutcome::AcknowledgeAfterEcho) } } } @@ -1460,6 +1529,115 @@ fn fleet_spawn_outcome( } } +/// Resolve a fleet (engine-facing) `delivery_ack` withheld pending +/// confirmation of a specific PTY injection (relay#1310: the ack must not +/// fire before the worker confirms the write landed). Called with the +/// `delivery_id`/`event_id` from the worker's own internal `delivery_ack` +/// event (`worker_events.rs`), which pty_worker.rs sends only after echo +/// verification succeeds or its bounded timeout fallback fires — never at +/// write-enqueue time. +/// +/// Returns the `(agent, up_to_seq)` to send to the engine once resolved, or +/// `None` when there is nothing withheld for `delivery_id` (already +/// resolved, dead-lettered, or never a fleet-originated delivery) or when +/// `event_id` doesn't match the withheld delivery's `msg_id` — the same +/// stale/reused-id guard `clear_pending_delivery_if_event_matches` applies to +/// the broker's own pending-delivery map. +/// +/// The withheld ack itself lives on the `PendingDelivery` +/// (`withheld_fleet_ack`, see relay#1543), so the delivery_id/event_id +/// matching that used to happen against a second map here is already done by +/// `clear_pending_delivery_if_event_matches` before this is called — this +/// only extracts what that lookup found and commits it to the delivery book. +/// +/// Split out from `handle_fleet_deliver`/`handle_worker_event` so the ack +/// gating is testable without a whole `BrokerRuntime`. +pub(super) fn resolve_pending_fleet_ack( + pending: Option<&PendingDelivery>, + fleet_delivery_book: &mut FleetDeliveryBook, +) -> Option<(String, u64)> { + let deliver = pending?.withheld_fleet_ack.as_ref()?; + let up_to_seq = fleet_delivery_book.commit_confirmed_delivery(deliver)?; + Some((deliver.agent.clone(), up_to_seq)) +} + +/// Confirm one pending worker delivery and resolve only the contiguous prefix +/// of fleet ACKs that is now safe to report cumulatively. +/// +/// Before removing the matching pending entry, this restores a missing +/// post-restart cursor from every withheld sibling for the same immutable +/// agent identity. If the confirmed sequence is ahead of a lower unconfirmed +/// sibling, the entry is put back into the pending map and excluded from +/// retries until the lower confirmation releases both. Keeping the entry in +/// the persisted pending store also makes another broker restart safe: the +/// broker may retry an already-landed delivery, but it cannot falsely ACK an +/// undelivered lower one. +pub(super) fn confirm_pending_delivery_and_resolve_fleet_ack( + pending_deliveries: &mut HashMap, + delivery_id: &str, + event_id: Option<&str>, + worker_name: &str, + worker_signal: &str, + fleet_delivery_book: &mut FleetDeliveryBook, +) -> (Option, Option<(String, u64)>) { + if let Some(deliver) = pending_deliveries + .get(delivery_id) + .and_then(|pending| pending.withheld_fleet_ack.as_ref()) + { + let same_agent_pending = pending_deliveries + .values() + .filter_map(|pending| pending.withheld_fleet_ack.as_ref()) + .filter(|sibling| sibling.agent_id == deliver.agent_id) + .cloned() + .collect::>(); + let ack_floor = pending_deliveries + .values() + .filter(|pending| { + pending + .withheld_fleet_ack + .as_ref() + .is_some_and(|sibling| sibling.agent_id == deliver.agent_id) + }) + .filter_map(|pending| pending.withheld_fleet_ack_floor) + .min(); + fleet_delivery_book.restore_pending_agent(&same_agent_pending, ack_floor); + } + + let pending = clear_pending_delivery_if_event_matches( + pending_deliveries, + delivery_id, + event_id, + worker_name, + worker_signal, + ); + let Some(pending) = pending else { + return (None, None); + }; + let already_held = pending + .withheld_fleet_ack + .as_ref() + .is_some_and(|deliver| fleet_delivery_book.is_delivery_confirmation_held(deliver)); + let resolved = resolve_pending_fleet_ack(Some(&pending), fleet_delivery_book); + + match (&pending.withheld_fleet_ack, &resolved) { + (Some(deliver), Some((_, up_to_seq))) if deliver.seq > 0 => { + pending_deliveries.retain(|_, sibling| { + !sibling.withheld_fleet_ack.as_ref().is_some_and(|sibling| { + sibling.agent_id == deliver.agent_id + && sibling.seq > 0 + && sibling.seq <= *up_to_seq + }) + }); + } + (Some(_), None) => { + pending_deliveries.insert(pending.delivery.delivery_id.clone(), pending.clone()); + } + _ => {} + } + + ((!already_held).then_some(pending), resolved) +} + fn fleet_spawn_action_result( invocation_id: &str, name: &WorkerName, @@ -1490,6 +1668,20 @@ pub(super) struct FlushPendingRelayResult { /// Inject a worker's held queue in FIFO order. A failed item and every item /// behind it remain queued. Relaycast ACKs advance only after the corresponding /// PTY write succeeds, so the emitted cursor is always an injected prefix. +/// +/// relay#1310 design note: unlike `handle_fleet_deliver`'s `DrainNow` path, +/// this loop's ack timing is intentionally **not** deferred to echo +/// verification in this change. Each item's `can_ack_receipt` precondition +/// (`node_control.rs`) requires the *previous* item's ack to have already +/// committed, since `acked_up_to_seq` is a strictly ordered cumulative +/// cursor; committing here happens synchronously in this one loop so a +/// multi-item backlog drains in a single call. Deferring commit to async +/// echo confirmation would stall the loop after the first item — the second +/// item's `can_ack_receipt` check would never pass until the first item's +/// echo resolves — turning a batch drain into a call-per-item serialization. +/// A real fix needs an ack-cursor model that tolerates a pipeline of +/// outstanding (unconfirmed) commits, which is a separate, larger change; +/// this flush path still acks on write, not on echo, until that lands. pub(super) async fn flush_pending_relay_messages( delivery_states: &mut HashMap, workers: &mut WorkerRegistry, diff --git a/crates/broker/src/runtime/maintenance.rs b/crates/broker/src/runtime/maintenance.rs index aa1c7b517..9019130bd 100644 --- a/crates/broker/src/runtime/maintenance.rs +++ b/crates/broker/src/runtime/maintenance.rs @@ -152,7 +152,11 @@ impl BrokerRuntime { let due_ids: Vec = pending_deliveries .iter() .filter_map(|(delivery_id, pending)| { - if pending.next_retry_at <= now { + let confirmation_is_held = + pending.withheld_fleet_ack.as_ref().is_some_and(|deliver| { + fleet_delivery_book.is_delivery_confirmation_held(deliver) + }); + if pending.next_retry_at <= now && !confirmation_is_held { Some(delivery_id.clone()) } else { None diff --git a/crates/broker/src/runtime/tests.rs b/crates/broker/src/runtime/tests.rs index c0615466f..f5fb48985 100644 --- a/crates/broker/src/runtime/tests.rs +++ b/crates/broker/src/runtime/tests.rs @@ -50,8 +50,9 @@ use super::{ relaycast_ws_should_apply_local_spawn_echo_dedup, relaycast_ws_spawn_token, requeue_dead_letter, resolve_exit_after_task, resolve_workspace, retry_pending_delivery, save_dead_letters, seed_supplied_agent_token, send_broker_event, sender_is_dashboard_label, - should_clear_pending_delivery_for_event, synthetic_delivery_read_ack_reason, AgentRuntime, - DeadLetterEntry, DeadLetterStore, DeliveryAttemptOutcome, InboundContext, InboundQueueOutcome, + should_clear_pending_delivery_for_event, synthetic_delivery_read_ack_reason, + take_pending_for_worker, try_inject_pending_relay_message, AgentRuntime, DeadLetterEntry, + DeadLetterStore, DeliveryAttemptOutcome, InboundContext, InboundQueueOutcome, ObserverTokenMintError, ObserverTokenMintOutcome, PendingDelivery, PendingDeliveryStore, ProtocolHeadlessProvider, RelayWorkspace, TypedThreadMessage, MAX_DEAD_LETTERS, MAX_DELIVERY_RETRIES, @@ -126,6 +127,67 @@ async fn make_worker_registry_with_worker(name: &str) -> WorkerRegistry { registry } +/// A worker whose command channel accepts frames but never completes them — +/// no writer task ever drains `command_rx`, so `deliver()` hangs forever. +/// Models a handoff that outlives `retry_interval` deterministically (no +/// timing race): the receiver stays alive (a dropped one would fail the +/// send instead of hanging it), so the send always succeeds and the +/// subsequent completion wait never returns on its own. +async fn make_worker_registry_with_stalled_worker(name: &str) -> WorkerRegistry { + let (tx, _rx) = mpsc::channel::(16); + let mut registry = WorkerRegistry::new( + tx, + Vec::new(), + PathBuf::from("/tmp/agent-relay-broker-tests"), + Instant::now(), + ); + let child = tokio::process::Command::new("cat") + .stdin(Stdio::piped()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .expect("test worker process should spawn"); + let generation = Uuid::new_v4(); + let (command_tx, command_rx) = mpsc::channel(128); + // Deliberately leaked, not spawned as a writer: keeps the receiver alive + // (so sends succeed) without anything ever draining it. + std::mem::forget(command_rx); + registry.workers.insert( + WorkerName::from(name), + WorkerHandle { + generation, + spec: AgentSpec { + name: WorkerName::from(name), + runtime: AgentRuntime::Pty, + provider: None, + cli: Some("cat".to_string()), + session_id: None, + harness_config: None, + model: None, + cwd: None, + team: None, + shadow_of: None, + shadow_mode: None, + args: Vec::new(), + channels: Vec::new(), + restart_policy: None, + }, + parent: None, + workspace_id: Some(WorkspaceId::new("ws_demo")), + child, + command_tx, + harness_pid: None, + spawned_at: Instant::now(), + ready_at: Some(Instant::now()), + last_activity_at: Instant::now(), + context_budget_pct: None, + state: AgentWorkState::Working, + exit_reason: None, + }, + ); + registry +} + async fn cleanup_worker_registry(mut registry: WorkerRegistry) { for handle in registry.workers.values_mut() { let _ = handle.child.start_kill(); @@ -203,6 +265,8 @@ fn pending_delivery(worker_name: &str, delivery_id: &str, event_id: &str) -> Pen next_retry_at: Instant::now(), queued_at_ms: super::unix_timestamp_millis(), last_error: None, + withheld_fleet_ack: None, + withheld_fleet_ack_floor: None, } } @@ -495,6 +559,8 @@ fn make_pending_delivery(delivery_id: &str, worker: &str) -> PendingDelivery { next_retry_at: Instant::now(), queued_at_ms: super::unix_timestamp_millis(), last_error: None, + withheld_fleet_ack: None, + withheld_fleet_ack_floor: None, } } @@ -544,6 +610,89 @@ fn pending_delivery_load_defaults_legacy_failure_count() { assert_eq!(loaded["del_legacy"].failed_attempts, 0); } +// relay#1543 delivery.rs:190 MUST-FIRE (P1, blocker): a withheld fleet +// (engine-facing) ack must survive a broker restart along with the delivery +// it belongs to. Before this fix, `load_pending_deliveries` unconditionally +// reset `withheld_fleet_ack` to `None` on every reload — so a delivery that +// was persisted mid-flight (worker handed the injection but hadn't confirmed +// yet) came back after restart with its ack silently dropped. The retried +// delivery could still reach the worker and get echo-confirmed, but +// `resolve_pending_fleet_ack` would then have nothing to release: the engine +// stays unacknowledged and may redeliver a message the worker already has. +// This simulates a real restart end-to-end — shutdown persist, then startup +// load — rather than asserting on the intermediate `PersistedPendingDelivery` +// struct, so it catches a regression anywhere in that round trip. +#[test] +fn withheld_fleet_ack_survives_a_simulated_restart() { + let dir = tempfile::tempdir().expect("tempdir should create"); + let path = dir.path().join("pending-deliveries.json"); + let mut delivery = make_pending_delivery("del_inflight", "worker-a"); + delivery.withheld_fleet_ack = Some(withheld_ack_for("del_inflight")); + delivery.withheld_fleet_ack_floor = Some(1); + let deliveries = HashMap::from([(DeliveryId::new("del_inflight"), delivery)]); + + // Shutdown: persist whatever is still pending, exactly as the broker + // does before exiting. + persist_pending_on_shutdown(&path, true, &deliveries); + + // Startup: reload from disk, exactly as the broker does on the next boot. + let reloaded = load_pending_deliveries(&path); + let pending = reloaded + .get("del_inflight") + .expect("the in-flight delivery must survive the restart"); + assert_eq!( + pending + .withheld_fleet_ack + .as_ref() + .map(|d| d.msg_id.as_str()), + Some("evt_del_inflight"), + "the withheld fleet ack must survive the restart along with the delivery it belongs \ + to — otherwise a retried delivery that goes on to land has no ack left to release, \ + and the engine stays permanently unacknowledged for it" + ); + assert_eq!(pending.withheld_fleet_ack_floor, Some(1)); +} + +// relay#1543 delivery.rs:190 companion: a snapshot written by a broker +// version that predates `withheld_fleet_ack` (or a fresh delivery that never +// had one) must still load cleanly with the field defaulted to `None`, +// mirroring `pending_delivery_load_defaults_legacy_failure_count` above for +// `failed_attempts`. This is the other half of the P1 fix — `#[serde(default)]` +// must actually work, not just be present in the struct definition. +#[test] +fn legacy_pending_delivery_snapshot_without_withheld_ack_field_loads_as_none() { + let dir = tempfile::tempdir().expect("tempdir should create"); + let path = dir.path().join("pending-deliveries.json"); + let delivery = make_pending_delivery("del_legacy_ack", "worker-a"); + let deliveries = HashMap::from([(DeliveryId::new("del_legacy_ack"), delivery)]); + super::save_pending_deliveries(&path, &deliveries).expect("pending delivery should save"); + let mut json: Value = serde_json::from_slice( + &std::fs::read(&path).expect("pending delivery snapshot should read"), + ) + .expect("pending delivery snapshot should parse"); + json[0] + .as_object_mut() + .expect("pending delivery entry should be an object") + .remove("withheld_fleet_ack"); + json[0] + .as_object_mut() + .expect("pending delivery entry should be an object") + .remove("withheld_fleet_ack_floor"); + std::fs::write( + &path, + serde_json::to_vec(&json).expect("legacy snapshot encodes"), + ) + .expect("legacy pending snapshot should write"); + + let loaded = load_pending_deliveries(&path); + assert_eq!( + loaded["del_legacy_ack"].withheld_fleet_ack, None, + "a pre-relay#1543 snapshot has no withheld_fleet_ack field at all — it must load as \ + None (the same state that delivery actually had), not fail to deserialize" + ); + assert_eq!(loaded["del_legacy_ack"].withheld_fleet_ack_floor, None); +} + #[test] fn shutdown_removes_pending_file_only_when_empty() { let dir = tempfile::tempdir().expect("tempdir should create"); @@ -913,6 +1062,585 @@ async fn retry_exhaustion_dead_letters_instead_of_discarding() { assert_eq!(dead_frame.payload["reason"], "failed writing frame"); } +fn withheld_ack_for(delivery_id: &str) -> Deliver { + Deliver { + v: FLEET_WIRE_VERSION, + agent: "agent-a".to_string(), + agent_id: "agent-a-id".to_string(), + delivery_id: delivery_id.to_string(), + msg_id: format!("evt_{delivery_id}"), + seq: 1, + mode: DeliveryMode::Wait, + payload: json!({}), + } +} + +// relay#1543 delivery.rs:588 MUST-FIRE (P1, blocker): when the initial +// worker handoff outlives `retry_interval`, the withheld fleet ack must +// already be registered on the `PendingDelivery` — not dependent on the +// timed-out call's `Ok(DeliveryId)` ever reaching the caller. Before the +// fix, `try_inject_pending_relay_message` returned only a bare `Result` +// derived from the timed-out future, and the fleet caller registered the +// withheld ack as a *separate* follow-up step keyed off that return value; +// a handoff that timed out returned `Err`, so the ack was simply never +// registered even though the delivery itself remained alive and retryable — +// a later successful retry's echo would then have had nothing to resolve. +#[tokio::test] +async fn timed_out_initial_handoff_still_registers_its_withheld_fleet_ack() { + let mut registry = make_worker_registry_with_stalled_worker("worker-a").await; + let deliver = fleet_deliver(1); + let msg = held_fleet_message(&deliver); + let mut pending_deliveries = HashMap::new(); + + let outcome = tokio::time::timeout( + Duration::from_secs(5), + try_inject_pending_relay_message( + &mut registry, + &mut pending_deliveries, + "worker-a", + &msg, + Duration::from_millis(20), + Some(deliver.clone()), + Some(deliver.seq), + ), + ) + .await + .expect( + "the test's own generous bound must never fire — only the short \ + retry_interval passed to try_inject_pending_relay_message should", + ); + + assert!( + outcome.is_err(), + "a handoff that never completes must time out, not hang forever" + ); + assert_eq!( + pending_deliveries.len(), + 1, + "the delivery must remain registered and retryable even though the initial handoff timed out" + ); + let registered = pending_deliveries + .values() + .next() + .expect("checked len() == 1 above"); + assert_eq!( + registered + .withheld_fleet_ack + .as_ref() + .map(|d| d.msg_id.as_str()), + Some(deliver.msg_id.as_str()), + "the withheld fleet ack must be registered before the timeout can expire, so a later \ + successful retry can still resolve it" + ); + + cleanup_worker_registry(registry).await; +} + +// relay#1543 MUST-FIRE (parameterised over every terminal disposition): a +// `PendingDelivery`'s withheld fleet ack must never survive the delivery it +// belongs to. Before the structural fix, `pending_fleet_acks` was a second +// map that none of these dispositions — except one very manually-threaded +// retry-exhaustion call site — knew to clean up. The ack now lives on +// `PendingDelivery` itself, so every path that disposes of the delivery +// disposes of the ack with it, by construction. +#[tokio::test] +async fn every_terminal_disposition_drops_its_withheld_fleet_ack() { + // Disposition 1: retry-exhaustion dead-letter (delivery.rs:816's thread) + // — the `emit_delivery_attempt_outcome` `Failed` arm. + { + let (tx, _rx) = mpsc::channel::(16); + let mut workers = WorkerRegistry::new( + tx, + Vec::new(), + PathBuf::from("/tmp/agent-relay-broker-tests"), + Instant::now(), + ); + let mut exhausted = make_pending_delivery("del_exhausted_ack", "ghost"); + exhausted.attempts = MAX_DELIVERY_RETRIES; + exhausted.failed_attempts = MAX_DELIVERY_RETRIES; + exhausted.withheld_fleet_ack = Some(withheld_ack_for("del_exhausted_ack")); + let mut pending_deliveries = + HashMap::from([(DeliveryId::new("del_exhausted_ack"), exhausted)]); + + let outcome = retry_pending_delivery( + &DeliveryId::new("del_exhausted_ack"), + &mut workers, + &mut pending_deliveries, + Duration::from_millis(1), + ) + .await + .expect("exhausted retries should classify as terminal failure"); + match &outcome { + DeliveryAttemptOutcome::Failed { pending, .. } => assert!( + pending.withheld_fleet_ack.is_some(), + "fixture must carry a withheld ack for this case to be meaningful" + ), + other => panic!("expected terminal failure, got {other:?}"), + } + + let (sdk_out_tx, mut sdk_out_rx) = mpsc::channel(4); + let mut dead_letters = DeadLetterStore::default(); + emit_delivery_attempt_outcome( + &sdk_out_tx, + &mut dead_letters, + &DeliveryId::new("del_exhausted_ack"), + true, + outcome, + ) + .await + .expect("terminal outcome should emit"); + + assert!(!pending_deliveries.contains_key("del_exhausted_ack")); + let mut book = FleetDeliveryBook::default(); + assert_eq!( + super::fleet::resolve_pending_fleet_ack( + pending_deliveries.get("del_exhausted_ack"), + &mut book + ), + None, + "a retry-exhausted delivery must never resolve into an engine ack" + ); + let _ = tokio::time::timeout(Duration::from_secs(1), sdk_out_rx.recv()).await; + let _ = tokio::time::timeout(Duration::from_secs(1), sdk_out_rx.recv()).await; + } + + // Dispositions 2 & 3: worker-exit and `delivery_failed` both dispose of a + // `PendingDelivery` via `emit_dropped_delivery_failures` — the single + // choke point every worker-teardown path (`take_pending_for_worker`, + // maintenance.rs:26 / event_loop.rs:255's threads) and the + // `delivery_failed` worker-event path share. This is driven through a + // real `pending_deliveries` map (via `take_pending_for_worker`, the same + // removal every one of those call sites uses) so the final assertion + // observes state the code under test actually produced, mirroring + // dispositions 1 and 4 below — not a hardcoded `None` that would pass + // for any implementation. See relay#1543 tests.rs:1142's review thread. + for reason in ["worker_exited", "delivery_failed"] { + let mut pending = make_pending_delivery("del_dropped_ack", "ghost"); + pending.withheld_fleet_ack = Some(withheld_ack_for("del_dropped_ack")); + let mut pending_deliveries = HashMap::from([(DeliveryId::new("del_dropped_ack"), pending)]); + + let dropped = take_pending_for_worker(&mut pending_deliveries, "ghost"); + assert_eq!( + dropped.len(), + 1, + "fixture must carry exactly the one delivery being torn down for {reason}" + ); + + let (sdk_out_tx, mut sdk_out_rx) = mpsc::channel(4); + let mut dead_letters = DeadLetterStore::default(); + emit_dropped_delivery_failures(&sdk_out_tx, &mut dead_letters, &dropped, reason) + .await + .expect("dropped delivery outcome should emit"); + + let mut book = FleetDeliveryBook::default(); + assert_eq!( + super::fleet::resolve_pending_fleet_ack( + pending_deliveries.get("del_dropped_ack"), + &mut book + ), + None, + "a delivery dropped for {reason} must never resolve into an engine ack" + ); + let _ = tokio::time::timeout(Duration::from_secs(1), sdk_out_rx.recv()).await; + let _ = tokio::time::timeout(Duration::from_secs(1), sdk_out_rx.recv()).await; + } + + // Disposition 4: a `WorkerMissing` fleet injection whose recipient never + // existed (fleet.rs:741's thread) — before the fix this injected via a + // bare `workers.deliver` call outside `pending_deliveries`, so nothing + // ever tracked its withheld ack at all. Routed through + // `insert_and_attempt_delivery` like `DrainNow`, it is tracked from the + // first attempt and reaches the exact same terminal cleanup as every + // other disposition above. + { + let (tx, _rx) = mpsc::channel::(16); + let mut workers = WorkerRegistry::new( + tx, + Vec::new(), + PathBuf::from("/tmp/agent-relay-broker-tests"), + Instant::now(), + ); // no worker ever registered + let relay_delivery = RelayDelivery { + delivery_id: DeliveryId::new("del_worker_missing"), + event_id: EventId::new("evt_worker_missing"), + workspace_id: None, + workspace_alias: None, + from: "Alice".to_string(), + target: MessageTarget::new("ghost"), + body: "hello".to_string(), + thread_id: None, + priority: Some(2), + injection_mode: MessageInjectionMode::Wait, + }; + let mut pending_deliveries = HashMap::new(); + + let first_attempt = super::insert_and_attempt_delivery( + &mut workers, + &mut pending_deliveries, + "ghost", + relay_delivery, + Duration::from_millis(1), + Some(withheld_ack_for("del_worker_missing")), + Some(1), + ) + .await; + assert!( + first_attempt.is_err(), + "a missing recipient must fail the handoff" + ); + let tracked = pending_deliveries.get("del_worker_missing").expect( + "the delivery must remain tracked for the terminal-failure path to dead-letter it, \ + not vanish silently", + ); + assert!( + tracked.withheld_fleet_ack.is_some(), + "the withheld ack must have survived the failed first attempt" + ); + + // The next retry attempt (e.g. the maintenance sweep) observes the + // same missing recipient and reaches the terminal `Failed` outcome + // that `emit_delivery_attempt_outcome` dead-letters and drops the + // ack for — same as every other disposition in this test. + let outcome = retry_pending_delivery( + &DeliveryId::new("del_worker_missing"), + &mut workers, + &mut pending_deliveries, + Duration::from_millis(1), + ) + .await + .expect("a still-missing recipient should classify as terminal failure"); + + let (sdk_out_tx, mut sdk_out_rx) = mpsc::channel(4); + let mut dead_letters = DeadLetterStore::default(); + emit_delivery_attempt_outcome( + &sdk_out_tx, + &mut dead_letters, + &DeliveryId::new("del_worker_missing"), + true, + outcome, + ) + .await + .expect("terminal outcome should emit"); + + assert!(!pending_deliveries.contains_key("del_worker_missing")); + let mut book = FleetDeliveryBook::default(); + assert_eq!( + super::fleet::resolve_pending_fleet_ack( + pending_deliveries.get("del_worker_missing"), + &mut book + ), + None, + "a delivery to a permanently missing worker must never resolve into an engine ack" + ); + let _ = tokio::time::timeout(Duration::from_secs(1), sdk_out_rx.recv()).await; + let _ = tokio::time::timeout(Duration::from_secs(1), sdk_out_rx.recv()).await; + } +} + +// relay#1310 MUST-NOT-FIRE: once the worker confirms the injection landed +// (echo-verified, or its bounded timeout fallback — pty_worker.rs sends the +// same internal `delivery_ack` event either way), the engine ack must still +// fire, with the delivery's own (agent, up_to_seq) — i.e. the happy path is +// unchanged, just correctly gated on confirmation instead of write-enqueue. +// Exercises the full wiring: a real handoff through +// `try_inject_pending_relay_message`, then the same two-step resolution +// `handle_worker_event`'s `delivery_ack` arm performs +// (`clear_pending_delivery_if_event_matches` then `resolve_pending_fleet_ack`). +#[tokio::test] +async fn successful_injection_still_resolves_its_withheld_fleet_ack() { + let worker_name = "worker-a"; + let mut registry = make_worker_registry_with_worker(worker_name).await; + let deliver = fleet_deliver(1); + let msg = held_fleet_message(&deliver); + let mut pending_deliveries = HashMap::new(); + + let delivery_id = try_inject_pending_relay_message( + &mut registry, + &mut pending_deliveries, + worker_name, + &msg, + Duration::from_secs(2), + Some(deliver.clone()), + Some(deliver.seq), + ) + .await + .expect("a registered worker should accept the handoff"); + + assert!( + pending_deliveries + .get(&delivery_id) + .expect("the delivery must be tracked pending the worker's confirmation") + .withheld_fleet_ack + .is_some(), + "a successful handoff must still withhold the ack pending echo confirmation" + ); + + let resolved_pending = clear_pending_delivery_if_event_matches( + &mut pending_deliveries, + delivery_id.as_str(), + Some(deliver.msg_id.as_str()), + worker_name, + "delivery_ack", + ); + assert!( + resolved_pending.is_some(), + "a matching event_id must clear the pending delivery" + ); + + let mut fleet_delivery_book = FleetDeliveryBook::default(); + let resolved = super::fleet::resolve_pending_fleet_ack( + resolved_pending.as_ref(), + &mut fleet_delivery_book, + ); + assert_eq!( + resolved, + Some((deliver.agent.clone(), deliver.seq)), + "a genuinely landed delivery must still resolve its withheld engine ack" + ); + assert!(!pending_deliveries.contains_key(&delivery_id)); + + cleanup_worker_registry(registry).await; +} + +// relay#1543 restart-ordering MUST-NOT-FIRE / MUST-FIRE boundary. The +// confirmation order is explicit: HashMap iteration never decides which +// delivery lands first. With both restored sequences pending, confirming seq +// 42 first must not emit a cumulative ack that lies about seq 41. Once seq 41 +// confirms, the held seq 42 confirmation must be released by one cumulative +// ack through seq 42. +#[test] +fn restored_out_of_order_confirmation_waits_for_the_lower_sequence() { + // Use a mid-stream cursor to exercise the restart-only "adopt first + // position" branch rather than accidentally relying on a fresh seq-1 + // stream. + let first = fleet_deliver(41); + let second = fleet_deliver(42); + let mut first_pending = pending_delivery( + "worker-a", + first.delivery_id.as_str(), + first.msg_id.as_str(), + ); + first_pending.withheld_fleet_ack = Some(first.clone()); + let mut second_pending = pending_delivery( + "worker-a", + second.delivery_id.as_str(), + second.msg_id.as_str(), + ); + second_pending.withheld_fleet_ack = Some(second.clone()); + let mut pending_deliveries = HashMap::from([ + (DeliveryId::from(&first.delivery_id), first_pending), + (DeliveryId::from(&second.delivery_id), second_pending), + ]); + let mut fleet_delivery_book = FleetDeliveryBook::default(); + + let (confirmed_second, second_ack) = + super::fleet::confirm_pending_delivery_and_resolve_fleet_ack( + &mut pending_deliveries, + second.delivery_id.as_str(), + Some(second.msg_id.as_str()), + "worker-a", + "delivery_ack", + &mut fleet_delivery_book, + ); + assert!(confirmed_second.is_some()); + assert_eq!( + second_ack, None, + "seq 42 must remain withheld while restored seq 41 is still unconfirmed" + ); + assert!( + pending_deliveries.contains_key(second.delivery_id.as_str()), + "the held seq 42 confirmation must remain durable until seq 41 confirms" + ); + assert!( + fleet_delivery_book.is_delivery_confirmation_held(&second), + "maintenance must be able to distinguish the confirmed hold from a retryable delivery" + ); + assert_eq!( + fleet_delivery_book.acked_up_to_seq(first.agent_id.as_str()), + first.seq - 1, + "the restored cursor must remain immediately below the lowest pending sequence" + ); + + let (confirmed_first, first_ack) = super::fleet::confirm_pending_delivery_and_resolve_fleet_ack( + &mut pending_deliveries, + first.delivery_id.as_str(), + Some(first.msg_id.as_str()), + "worker-a", + "delivery_ack", + &mut fleet_delivery_book, + ); + assert!(confirmed_first.is_some()); + assert_eq!( + first_ack, + Some((first.agent.clone(), second.seq)), + "confirming seq 41 must release its already-confirmed seq 42 sibling" + ); + assert!( + pending_deliveries.is_empty(), + "the cumulative seq 42 ack must release both pending entries" + ); + assert_eq!( + fleet_delivery_book.acked_up_to_seq(first.agent_id.as_str()), + second.seq + ); + assert!(!fleet_delivery_book.is_delivery_confirmation_held(&second)); +} + +// The ordering gap itself must survive another restart after the lower entry +// has left the pending map. Otherwise the remaining higher sequence would be +// mistaken for a new baseline and could once again cumulatively ACK the +// dead-lettered lower delivery. +#[test] +fn restored_ack_floor_survives_lower_failure_and_a_second_restart() { + let dir = tempfile::tempdir().expect("temp dir"); + let path = dir.path().join("pending.json"); + let first = fleet_deliver(41); + let second = fleet_deliver(42); + let mut first_pending = pending_delivery( + "worker-a", + first.delivery_id.as_str(), + first.msg_id.as_str(), + ); + first_pending.withheld_fleet_ack = Some(first.clone()); + let mut second_pending = pending_delivery( + "worker-a", + second.delivery_id.as_str(), + second.msg_id.as_str(), + ); + second_pending.withheld_fleet_ack = Some(second.clone()); + let pending_deliveries = HashMap::from([ + (DeliveryId::from(&first.delivery_id), first_pending), + (DeliveryId::from(&second.delivery_id), second_pending), + ]); + super::save_pending_deliveries(&path, &pending_deliveries).expect("save first snapshot"); + + let mut after_first_restart = load_pending_deliveries(&path); + assert_eq!( + after_first_restart[second.delivery_id.as_str()].withheld_fleet_ack_floor, + Some(first.seq) + ); + after_first_restart.remove(first.delivery_id.as_str()); + super::save_pending_deliveries(&path, &after_first_restart) + .expect("save snapshot after lower terminal failure"); + + let mut after_second_restart = load_pending_deliveries(&path); + assert_eq!( + after_second_restart[second.delivery_id.as_str()].withheld_fleet_ack_floor, + Some(first.seq), + "the higher entry must retain the failed lower sequence as its ACK floor" + ); + let mut fleet_delivery_book = FleetDeliveryBook::default(); + let (_, higher_ack) = super::fleet::confirm_pending_delivery_and_resolve_fleet_ack( + &mut after_second_restart, + second.delivery_id.as_str(), + Some(second.msg_id.as_str()), + "worker-a", + "delivery_ack", + &mut fleet_delivery_book, + ); + assert_eq!( + higher_ack, None, + "seq 42 must still not ACK through the absent, failed seq 41" + ); + assert!(after_second_restart.contains_key(second.delivery_id.as_str())); + + let mut retried_first = pending_delivery( + "worker-a", + first.delivery_id.as_str(), + first.msg_id.as_str(), + ); + retried_first.withheld_fleet_ack = Some(first.clone()); + after_second_restart.insert(DeliveryId::from(&first.delivery_id), retried_first); + let (_, released_ack) = super::fleet::confirm_pending_delivery_and_resolve_fleet_ack( + &mut after_second_restart, + first.delivery_id.as_str(), + Some(first.msg_id.as_str()), + "worker-a", + "delivery_ack", + &mut fleet_delivery_book, + ); + assert_eq!(released_ack, Some((first.agent.clone(), second.seq))); + assert!(after_second_restart.is_empty()); +} + +// The paired happy-path boundary: adding an ordering hold must not turn +// ordinary in-order worker confirmations into acknowledgements that never +// fire. +#[test] +fn in_order_confirmation_still_acknowledges_each_sequence_immediately() { + let mut fleet_delivery_book = FleetDeliveryBook::default(); + let mut pending_deliveries = HashMap::new(); + + for expected_seq in [1, 2] { + let deliver = fleet_deliver(expected_seq); + let mut pending = pending_delivery( + "worker-a", + deliver.delivery_id.as_str(), + deliver.msg_id.as_str(), + ); + pending.withheld_fleet_ack = Some(deliver.clone()); + pending_deliveries.insert(DeliveryId::from(&deliver.delivery_id), pending); + + let (confirmed, resolved) = super::fleet::confirm_pending_delivery_and_resolve_fleet_ack( + &mut pending_deliveries, + deliver.delivery_id.as_str(), + Some(deliver.msg_id.as_str()), + "worker-a", + "delivery_ack", + &mut fleet_delivery_book, + ); + assert!(confirmed.is_some()); + assert_eq!( + resolved, + Some((deliver.agent.clone(), expected_seq)), + "an in-order confirmation must ack without waiting for another event" + ); + assert!(pending_deliveries.is_empty()); + } +} + +// A worker delivery_ack whose event_id doesn't match the withheld delivery's +// event_id (stale or reused delivery_id) must not resolve into an engine +// ack. The matching itself is `clear_pending_delivery_if_event_matches`'s +// job (see `clear_pending_delivery_returns_none_for_stale_event_id` below) +// — this exercises that guard against a real pending delivery that actually +// carries a withheld ack, then feeds its *return value* into +// `resolve_pending_fleet_ack` exactly as `handle_worker_event`'s +// `delivery_ack` arm does, so a regression that made the guard incorrectly +// clear on a mismatch would surface here as a resolved ack. See relay#1543 +// tests.rs:1309's review thread — the prior version passed a hardcoded +// `None` straight to `resolve_pending_fleet_ack`, which is `None` for every +// implementation and never exercised the guard at all. +#[tokio::test] +async fn mismatched_event_id_leaves_nothing_for_resolve_pending_fleet_ack() { + let mut pending = make_pending_delivery("del_reused", "worker-a"); + pending.withheld_fleet_ack = Some(withheld_ack_for("del_reused")); + let mut pending_deliveries = HashMap::from([(DeliveryId::new("del_reused"), pending)]); + + let cleared = clear_pending_delivery_if_event_matches( + &mut pending_deliveries, + "del_reused", + Some("evt_stale_reused_id"), + "worker-a", + "delivery_ack", + ); + assert!( + cleared.is_none(), + "a mismatched event_id must not clear the pending delivery" + ); + assert!( + pending_deliveries.contains_key("del_reused"), + "a mismatched event must not consume the withheld entry" + ); + + let mut fleet_delivery_book = FleetDeliveryBook::default(); + assert_eq!( + super::fleet::resolve_pending_fleet_ack(cleared.as_ref(), &mut fleet_delivery_book), + None, + "no pending delivery (because the event_id guard declined to clear one) means nothing to resolve" + ); +} + #[tokio::test] async fn delivery_retry_fails_promptly_when_recipient_is_gone() { let (tx, _rx) = mpsc::channel::(16); @@ -943,6 +1671,8 @@ async fn delivery_retry_fails_promptly_when_recipient_is_gone() { next_retry_at: Instant::now(), queued_at_ms: super::unix_timestamp_millis(), last_error: Some("failed writing frame".to_string()), + withheld_fleet_ack: None, + withheld_fleet_ack_floor: None, }, )]); @@ -1001,6 +1731,8 @@ async fn initial_delivery_failure_stays_owned_until_dead_lettered() { 2, MessageInjectionMode::Wait, Duration::from_millis(1), + None, + None, ) .await .expect_err("missing recipient should fail the initial handoff"); @@ -1060,6 +1792,8 @@ async fn delivery_retry_transient_blip_emits_failed_event_for_present_worker() { next_retry_at: Instant::now(), queued_at_ms: super::unix_timestamp_millis(), last_error: None, + withheld_fleet_ack: None, + withheld_fleet_ack_floor: None, }, )]); @@ -1200,6 +1934,8 @@ async fn delivery_retry_success_clears_stale_last_error() { next_retry_at: Instant::now(), queued_at_ms: super::unix_timestamp_millis(), last_error: Some("old transient failure".to_string()), + withheld_fleet_ack: None, + withheld_fleet_ack_floor: None, }, )]); @@ -2110,6 +2846,8 @@ fn drop_pending_for_worker_removes_only_matching_entries() { next_retry_at: Instant::now(), queued_at_ms: super::unix_timestamp_millis(), last_error: None, + withheld_fleet_ack: None, + withheld_fleet_ack_floor: None, }, ); pending.insert( @@ -2133,6 +2871,8 @@ fn drop_pending_for_worker_removes_only_matching_entries() { next_retry_at: Instant::now(), queued_at_ms: super::unix_timestamp_millis(), last_error: None, + withheld_fleet_ack: None, + withheld_fleet_ack_floor: None, }, ); @@ -2163,6 +2903,8 @@ async fn dropped_pending_deliveries_emit_terminal_message_failures() { next_retry_at: Instant::now(), queued_at_ms: super::unix_timestamp_millis(), last_error: Some("previous blip".to_string()), + withheld_fleet_ack: None, + withheld_fleet_ack_floor: None, }; let (sdk_out_tx, mut sdk_out_rx) = mpsc::channel(4); let mut dead_letters = DeadLetterStore::default(); @@ -2233,6 +2975,8 @@ fn should_clear_pending_delivery_when_event_id_matches() { next_retry_at: Instant::now(), queued_at_ms: super::unix_timestamp_millis(), last_error: None, + withheld_fleet_ack: None, + withheld_fleet_ack_floor: None, }; assert!(should_clear_pending_delivery_for_event( @@ -2268,6 +3012,8 @@ fn clear_pending_delivery_returns_none_for_stale_event_id() { next_retry_at: Instant::now(), queued_at_ms: super::unix_timestamp_millis(), last_error: None, + withheld_fleet_ack: None, + withheld_fleet_ack_floor: None, }, )]); @@ -2685,6 +3431,8 @@ fn should_clear_pending_delivery_without_event_id_for_compatibility() { next_retry_at: Instant::now(), queued_at_ms: super::unix_timestamp_millis(), last_error: None, + withheld_fleet_ack: None, + withheld_fleet_ack_floor: None, }; assert!(should_clear_pending_delivery_for_event( diff --git a/crates/broker/src/runtime/worker_events.rs b/crates/broker/src/runtime/worker_events.rs index 70073e796..0132b590b 100644 --- a/crates/broker/src/runtime/worker_events.rs +++ b/crates/broker/src/runtime/worker_events.rs @@ -1,8 +1,9 @@ use super::fleet::{ - fail_terminal_session, refresh_fleet_inventory_session_ref, try_send_terminal, - verified_spawn_ready_result, + confirm_pending_delivery_and_resolve_fleet_ack, fail_terminal_session, + refresh_fleet_inventory_session_ref, try_send_terminal, verified_spawn_ready_result, }; use super::*; +use crate::node_control::delivery_ack; use crate::terminal_control::{TerminalControlCommand, TerminalToCloud}; use crate::worker::AgentWorkState; @@ -619,6 +620,7 @@ impl BrokerRuntime { let pending_verified_spawns = &mut self.pending_verified_spawns; let delivery_retry_interval = self.delivery_retry_interval; let fleet_control_tx = &self.fleet_control_tx; + let fleet_delivery_book = &mut self.fleet_delivery_book; let fleet_inventory = &mut self.fleet_inventory; let delivery_states = &self.delivery_states; let terminal_control_tx = &self.terminal_control_tx; @@ -715,16 +717,34 @@ impl BrokerRuntime { let pending_for_confirmation = if let Ok(ack) = serde_json::from_value::(payload.clone()) { - let pending = clear_pending_delivery_if_event_matches( - pending_deliveries, - &ack.delivery_id, - Some(&ack.event_id), - &name, - "delivery_ack", - ); + let (pending, resolved_fleet_ack) = + confirm_pending_delivery_and_resolve_fleet_ack( + pending_deliveries, + &ack.delivery_id, + Some(&ack.event_id), + &name, + "delivery_ack", + fleet_delivery_book, + ); if pending.is_some() { terminal_failed_deliveries.remove(&ack.delivery_id); } + + // Resolve a fleet (engine-facing) ack withheld + // pending confirmation of this exact PTY + // injection (relay#1310). Restored sibling + // confirmations are held and released only as + // a contiguous per-agent prefix, so a higher + // sequence can never cumulatively ACK a lower + // delivery that has not landed (relay#1543). + if let Some((agent, up_to_seq)) = resolved_fleet_ack { + let _ = fleet_control_tx + .send(FleetControlCommand::Send(delivery_ack( + agent, up_to_seq, + ))) + .await; + } + pending } else { None @@ -1409,6 +1429,8 @@ impl BrokerRuntime { 2, MessageInjectionMode::Wait, delivery_retry_interval, + None, + None, ) .await { @@ -1740,6 +1762,8 @@ impl BrokerRuntime { 2, MessageInjectionMode::Wait, delivery_retry_interval, + None, + None, ) .await {