Skip to content

Commit 74ad371

Browse files
committed
fix(harness): correlate subagent lifecycle events
1 parent a425660 commit 74ad371

2 files changed

Lines changed: 21 additions & 2 deletions

File tree

agentscope-harness/src/main/java/io/agentscope/harness/agent/tool/AgentSpawnTool.java

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -780,19 +780,20 @@ private Mono<Msg> execLocalSync(
780780
if (emitterOpt.isPresent()) {
781781
AgentEventEmitter parentEmitter = emitterOpt.get();
782782
String sourcePath = buildSourcePath(spawned, parentCtx);
783+
String replyId = UUID.randomUUID().toString().replace("-", "");
783784
AgentEventEmitter taggedEmitter =
784785
event -> parentEmitter.emit(event.withSource(sourcePath));
785786

786787
parentEmitter.emit(
787-
new AgentStartEvent(spawned.sessionId(), null, spawned.agentId())
788+
new AgentStartEvent(spawned.sessionId(), replyId, spawned.agentId())
788789
.withSource(sourcePath));
789790

790791
AtomicBoolean endEmitted = new AtomicBoolean();
791792
Runnable emitEnd =
792793
() -> {
793794
if (endEmitted.compareAndSet(false, true)) {
794795
parentEmitter.emit(
795-
new AgentEndEvent(null).withSource(sourcePath));
796+
new AgentEndEvent(replyId).withSource(sourcePath));
796797
}
797798
};
798799

agentscope-harness/src/test/java/io/agentscope/harness/agent/tool/AgentSpawnToolCancelEndEventTest.java

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
package io.agentscope.harness.agent.tool;
1717

1818
import static org.junit.jupiter.api.Assertions.assertEquals;
19+
import static org.junit.jupiter.api.Assertions.assertNotNull;
1920
import static org.junit.jupiter.api.Assertions.assertTrue;
2021
import static org.mockito.ArgumentMatchers.any;
2122

@@ -74,6 +75,14 @@ public void emit(AgentEvent event) {
7475
long count(Class<? extends AgentEvent> type) {
7576
return events.stream().filter(type::isInstance).count();
7677
}
78+
79+
<T extends AgentEvent> T first(Class<T> type) {
80+
return events.stream()
81+
.filter(type::isInstance)
82+
.map(type::cast)
83+
.findFirst()
84+
.orElseThrow();
85+
}
7786
}
7887

7988
@Test
@@ -125,6 +134,7 @@ void parentCancel_emitsAgentEndEvent() throws Exception {
125134
"parent cancel must still close the child's event stream — an AgentStartEvent"
126135
+ " without a matching AgentEndEvent leaves consumers rendering the"
127136
+ " subagent as running forever (doOnTerminate does not fire on cancel)");
137+
assertReplyIdPair(emitter);
128138
}
129139

130140
@Test
@@ -153,6 +163,14 @@ void normalCompletion_emitsPairedEvents() throws Exception {
153163

154164
assertEquals(1, emitter.count(AgentStartEvent.class), "expected one start event");
155165
assertEquals(1, emitter.count(AgentEndEvent.class), "expected one end event");
166+
assertReplyIdPair(emitter);
167+
}
168+
169+
private static void assertReplyIdPair(RecordingEmitter emitter) {
170+
AgentStartEvent start = emitter.first(AgentStartEvent.class);
171+
AgentEndEvent end = emitter.first(AgentEndEvent.class);
172+
assertNotNull(start.getReplyId(), "subagent start event should have a replyId");
173+
assertEquals(start.getReplyId(), end.getReplyId());
156174
}
157175

158176
private static final class NoopTaskRepository implements TaskRepository {

0 commit comments

Comments
 (0)