@@ -27,6 +27,18 @@ class AcpAgentSessionTest {
2727
2828 private static final Duration TIMEOUT = Duration .ofSeconds (5 );
2929
30+ private static final Duration PROMPT_RESPONSE_DELAY = Duration .ofMillis (250 );
31+
32+ private static final long AGENT_TRANSPORT_SUBSCRIPTION_DELAY_MILLIS = 100 ;
33+
34+ private static final long CLIENT_TRANSPORT_SUBSCRIPTION_DELAY_MILLIS = 50 ;
35+
36+ private static final int ACTIVE_PROMPT_ERROR_CODE = -32000 ;
37+
38+ private static final String SESSION_1 = "session-1" ;
39+
40+ private static final String SESSION_2 = "session-2" ;
41+
3042 @ Test
3143 void constructorValidatesArguments () {
3244 var transportPair = InMemoryTransportPair .create ();
@@ -60,8 +72,7 @@ void handlesIncomingRequest() throws Exception {
6072
6173 new AcpAgentSession (TIMEOUT , transportPair .agentTransport (), requestHandlers , Map .of ());
6274
63- // Allow transport to start
64- Thread .sleep (100 );
75+ allowAgentTransportSubscription ();
6576
6677 // Send a request from the client side
6778 CountDownLatch latch = new CountDownLatch (1 );
@@ -75,7 +86,7 @@ void handlesIncomingRequest() throws Exception {
7586 latch .countDown ();
7687 }).then (Mono .empty ())).subscribe ();
7788
78- Thread . sleep ( 50 );
89+ allowClientTransportSubscription ( );
7990 transportPair .clientTransport ().sendMessage (request ).block (TIMEOUT );
8091
8192 assertThat (latch .await (5 , TimeUnit .SECONDS )).isTrue ();
@@ -97,7 +108,7 @@ void handlesMethodNotFound() throws Exception {
97108 // Create session with no handlers
98109 new AcpAgentSession (TIMEOUT , transportPair .agentTransport (), Map .of (), Map .of ());
99110
100- Thread . sleep ( 100 );
111+ allowAgentTransportSubscription ( );
101112
102113 // Send a request for unknown method
103114 CountDownLatch latch = new CountDownLatch (1 );
@@ -111,7 +122,7 @@ void handlesMethodNotFound() throws Exception {
111122 latch .countDown ();
112123 }).then (Mono .empty ())).subscribe ();
113124
114- Thread . sleep ( 50 );
125+ allowClientTransportSubscription ( );
115126 transportPair .clientTransport ().sendMessage (request ).block (TIMEOUT );
116127
117128 assertThat (latch .await (5 , TimeUnit .SECONDS )).isTrue ();
@@ -141,14 +152,14 @@ void handlesNotification() throws Exception {
141152
142153 new AcpAgentSession (TIMEOUT , transportPair .agentTransport (), Map .of (), notificationHandlers );
143154
144- Thread . sleep ( 100 );
155+ allowAgentTransportSubscription ( );
145156
146157 // Send a notification from client
147158 AcpSchema .JSONRPCNotification notification = new AcpSchema .JSONRPCNotification (AcpSchema .JSONRPC_VERSION ,
148- AcpSchema .METHOD_SESSION_CANCEL , new AcpSchema .CancelNotification ("session-1" ));
159+ AcpSchema .METHOD_SESSION_CANCEL , new AcpSchema .CancelNotification (SESSION_1 ));
149160
150161 transportPair .clientTransport ().connect (mono -> mono .then (Mono .empty ())).subscribe ();
151- Thread . sleep ( 50 );
162+ allowClientTransportSubscription ( );
152163 transportPair .clientTransport ().sendMessage (notification ).block (TIMEOUT );
153164
154165 assertThat (notificationLatch .await (5 , TimeUnit .SECONDS )).isTrue ();
@@ -170,14 +181,14 @@ void singleTurnEnforcementRejectsConcurrentPromptsForSameSession() throws Except
170181 params -> Mono .defer (() -> {
171182 handlerInvocations .incrementAndGet ();
172183 handlerStarted .countDown ();
173- return Mono .delay (Duration . ofMillis ( 250 ) )
184+ return Mono .delay (PROMPT_RESPONSE_DELAY )
174185 .map (ignored -> new AcpSchema .PromptResponse (AcpSchema .StopReason .END_TURN ));
175186 }));
176187
177188 AcpAgentSession session = new AcpAgentSession (TIMEOUT , transportPair .agentTransport (), requestHandlers ,
178189 Map .of ());
179190
180- Thread . sleep ( 100 );
191+ allowAgentTransportSubscription ( );
181192
182193 CountDownLatch responseLatch = new CountDownLatch (2 );
183194 List <AcpSchema .JSONRPCResponse > responses = new CopyOnWriteArrayList <>();
@@ -189,19 +200,19 @@ void singleTurnEnforcementRejectsConcurrentPromptsForSameSession() throws Except
189200 responseLatch .countDown ();
190201 }).then (Mono .empty ())).subscribe ();
191202
192- Thread . sleep ( 50 );
203+ allowClientTransportSubscription ( );
193204
194- transportPair .clientTransport ().sendMessage (promptRequest ("1" , "session-1" , "first" )).block (TIMEOUT );
205+ transportPair .clientTransport ().sendMessage (promptRequest ("1" , SESSION_1 , "first" )).block (TIMEOUT );
195206 assertThat (handlerStarted .await (5 , TimeUnit .SECONDS )).isTrue ();
196- assertThat (session .hasActivePrompt ("session-1" )).isTrue ();
207+ assertThat (session .hasActivePrompt (SESSION_1 )).isTrue ();
197208
198- transportPair .clientTransport ().sendMessage (promptRequest ("2" , "session-1" , "second" )).block (TIMEOUT );
209+ transportPair .clientTransport ().sendMessage (promptRequest ("2" , SESSION_1 , "second" )).block (TIMEOUT );
199210
200211 assertThat (responseLatch .await (5 , TimeUnit .SECONDS )).isTrue ();
201212
202213 AcpSchema .JSONRPCResponse rejectedResponse = responseById (responses , "2" );
203214 assertThat (rejectedResponse .error ()).isNotNull ();
204- assertThat (rejectedResponse .error ().code ()).isEqualTo (- 32000 );
215+ assertThat (rejectedResponse .error ().code ()).isEqualTo (ACTIVE_PROMPT_ERROR_CODE );
205216 assertThat (rejectedResponse .error ().message ()).contains ("already an active prompt" );
206217 assertThat (handlerInvocations .get ()).isEqualTo (1 );
207218 assertThat (session .hasActivePrompt ()).isFalse ();
@@ -222,14 +233,14 @@ void singleTurnEnforcementAllowsConcurrentPromptsForDifferentSessions() throws E
222233 params -> Mono .defer (() -> {
223234 handlerInvocations .incrementAndGet ();
224235 handlersStarted .countDown ();
225- return Mono .delay (Duration . ofMillis ( 250 ) )
236+ return Mono .delay (PROMPT_RESPONSE_DELAY )
226237 .map (ignored -> new AcpSchema .PromptResponse (AcpSchema .StopReason .END_TURN ));
227238 }));
228239
229240 AcpAgentSession session = new AcpAgentSession (TIMEOUT , transportPair .agentTransport (), requestHandlers ,
230241 Map .of ());
231242
232- Thread . sleep ( 100 );
243+ allowAgentTransportSubscription ( );
233244
234245 CountDownLatch responseLatch = new CountDownLatch (2 );
235246 List <AcpSchema .JSONRPCResponse > responses = new CopyOnWriteArrayList <>();
@@ -241,15 +252,15 @@ void singleTurnEnforcementAllowsConcurrentPromptsForDifferentSessions() throws E
241252 responseLatch .countDown ();
242253 }).then (Mono .empty ())).subscribe ();
243254
244- Thread . sleep ( 50 );
255+ allowClientTransportSubscription ( );
245256
246- transportPair .clientTransport ().sendMessage (promptRequest ("1" , "session-1" , "first" )).block (TIMEOUT );
247- transportPair .clientTransport ().sendMessage (promptRequest ("2" , "session-2" , "second" )).block (TIMEOUT );
257+ transportPair .clientTransport ().sendMessage (promptRequest ("1" , SESSION_1 , "first" )).block (TIMEOUT );
258+ transportPair .clientTransport ().sendMessage (promptRequest ("2" , SESSION_2 , "second" )).block (TIMEOUT );
248259
249260 assertThat (handlersStarted .await (5 , TimeUnit .SECONDS )).isTrue ();
250- assertThat (session .hasActivePrompt ("session-1" )).isTrue ();
251- assertThat (session .hasActivePrompt ("session-2" )).isTrue ();
252- assertThat (session .getActivePromptSessionIds ()).containsExactlyInAnyOrder ("session-1" , "session-2" );
261+ assertThat (session .hasActivePrompt (SESSION_1 )).isTrue ();
262+ assertThat (session .hasActivePrompt (SESSION_2 )).isTrue ();
263+ assertThat (session .getActivePromptSessionIds ()).containsExactlyInAnyOrder (SESSION_1 , SESSION_2 );
253264
254265 assertThat (responseLatch .await (5 , TimeUnit .SECONDS )).isTrue ();
255266
@@ -273,37 +284,37 @@ void hasActivePromptReturnsCorrectState() throws Exception {
273284 Map <String , AcpAgentSession .RequestHandler <?>> requestHandlers = Map .of (AcpSchema .METHOD_SESSION_PROMPT ,
274285 params -> Mono .defer (() -> {
275286 handlerStarted .countDown ();
276- return Mono .delay (Duration . ofMillis ( 250 ) )
287+ return Mono .delay (PROMPT_RESPONSE_DELAY )
277288 .map (ignored -> new AcpSchema .PromptResponse (AcpSchema .StopReason .END_TURN ));
278289 }));
279290
280291 AcpAgentSession session = new AcpAgentSession (TIMEOUT , transportPair .agentTransport (), requestHandlers ,
281292 Map .of ());
282293
283- Thread . sleep ( 100 );
294+ allowAgentTransportSubscription ( );
284295
285296 assertThat (session .hasActivePrompt ()).isFalse ();
286- assertThat (session .hasActivePrompt ("session-1" )).isFalse ();
297+ assertThat (session .hasActivePrompt (SESSION_1 )).isFalse ();
287298 assertThat (session .getActivePromptSessionId ()).isNull ();
288299 assertThat (session .getActivePromptSessionIds ()).isEmpty ();
289300
290301 CountDownLatch responseLatch = new CountDownLatch (1 );
291302 transportPair .clientTransport ().connect (mono -> mono .doOnNext (msg -> responseLatch .countDown ())
292303 .then (Mono .empty ())).subscribe ();
293304
294- Thread . sleep ( 50 );
295- transportPair .clientTransport ().sendMessage (promptRequest ("1" , "session-1" , "hello" )).block (TIMEOUT );
305+ allowClientTransportSubscription ( );
306+ transportPair .clientTransport ().sendMessage (promptRequest ("1" , SESSION_1 , "hello" )).block (TIMEOUT );
296307
297308 assertThat (handlerStarted .await (5 , TimeUnit .SECONDS )).isTrue ();
298309 assertThat (session .hasActivePrompt ()).isTrue ();
299- assertThat (session .hasActivePrompt ("session-1" )).isTrue ();
300- assertThat (session .getActivePromptSessionIds ()).containsExactly ("session-1" );
301- assertThat (session .getActivePromptSessionId ()).isEqualTo ("session-1" );
310+ assertThat (session .hasActivePrompt (SESSION_1 )).isTrue ();
311+ assertThat (session .getActivePromptSessionIds ()).containsExactly (SESSION_1 );
312+ assertThat (session .getActivePromptSessionId ()).isEqualTo (SESSION_1 );
302313
303314 assertThat (responseLatch .await (5 , TimeUnit .SECONDS )).isTrue ();
304315
305316 assertThat (session .hasActivePrompt ()).isFalse ();
306- assertThat (session .hasActivePrompt ("session-1" )).isFalse ();
317+ assertThat (session .hasActivePrompt (SESSION_1 )).isFalse ();
307318 assertThat (session .getActivePromptSessionIds ()).isEmpty ();
308319 assertThat (session .getActivePromptSessionId ()).isNull ();
309320 }
@@ -318,7 +329,7 @@ void closeGracefullyCompletes() throws Exception {
318329
319330 AcpAgentSession session = new AcpAgentSession (TIMEOUT , transportPair .agentTransport (), Map .of (), Map .of ());
320331
321- Thread . sleep ( 100 );
332+ allowAgentTransportSubscription ( );
322333
323334 // Should complete without error
324335 session .closeGracefully ().block (TIMEOUT );
@@ -335,7 +346,7 @@ void handlerErrorReturnsJsonRpcError() throws Exception {
335346
336347 new AcpAgentSession (TIMEOUT , transportPair .agentTransport (), requestHandlers , Map .of ());
337348
338- Thread . sleep ( 100 );
349+ allowAgentTransportSubscription ( );
339350
340351 CountDownLatch latch = new CountDownLatch (1 );
341352 AtomicReference <AcpSchema .JSONRPCMessage > response = new AtomicReference <>();
@@ -348,7 +359,7 @@ void handlerErrorReturnsJsonRpcError() throws Exception {
348359 latch .countDown ();
349360 }).then (Mono .empty ())).subscribe ();
350361
351- Thread . sleep ( 50 );
362+ allowClientTransportSubscription ( );
352363 transportPair .clientTransport ().sendMessage (request ).block (TIMEOUT );
353364
354365 assertThat (latch .await (5 , TimeUnit .SECONDS )).isTrue ();
@@ -372,4 +383,17 @@ private static AcpSchema.JSONRPCResponse responseById(List<AcpSchema.JSONRPCResp
372383 return responses .stream ().filter (response -> id .equals (response .id ())).findFirst ().orElseThrow ();
373384 }
374385
386+ private static void allowAgentTransportSubscription () throws InterruptedException {
387+ // AcpAgentSession subscribes to the in-memory transport in its constructor.
388+ // subscribe() is asynchronous, so give the unicast sink subscriber a short
389+ // window to attach before the test sends client messages.
390+ Thread .sleep (AGENT_TRANSPORT_SUBSCRIPTION_DELAY_MILLIS );
391+ }
392+
393+ private static void allowClientTransportSubscription () throws InterruptedException {
394+ // clientTransport.connect(...).subscribe() also attaches asynchronously. Without
395+ // this small wait, an immediate agent response can race the test subscriber.
396+ Thread .sleep (CLIENT_TRANSPORT_SUBSCRIPTION_DELAY_MILLIS );
397+ }
398+
375399}
0 commit comments