diff --git a/reference/jsonrpc/src/main/java/org/a2aproject/sdk/server/apps/quarkus/A2AServerRoutes.java b/reference/jsonrpc/src/main/java/org/a2aproject/sdk/server/apps/quarkus/A2AServerRoutes.java index 44e267cf7..af7ab2ec6 100644 --- a/reference/jsonrpc/src/main/java/org/a2aproject/sdk/server/apps/quarkus/A2AServerRoutes.java +++ b/reference/jsonrpc/src/main/java/org/a2aproject/sdk/server/apps/quarkus/A2AServerRoutes.java @@ -171,7 +171,10 @@ public class A2AServerRoutes { @Inject - JSONRPCHandler jsonRpcHandler; + Instance jsonRpcHandler; + + @Inject + Instance multiAgentRegistry; @Inject AgentCardCacheMetadata cacheMetadata; @@ -201,16 +204,31 @@ public class A2AServerRoutes { * @param router the Vert.x Web Router instance to configure */ void setupRoutes(@Observes Router router) { - // Main JSON-RPC endpoint: POST / - // BodyHandler is per-route (not global) to avoid interfering with gRPC routes - // ordered=false: delegation via Vert.x WebClient can share the same event loop context as the outer request; ordered=true would serialize them, causing a 30s deadlock. - router.post("/") + if (!multiAgentRegistry.isUnsatisfied()) { + // Multi-agent mode + Map agents = multiAgentRegistry.get().getAgents(); + for (Map.Entry entry : agents.entrySet()) { + String agentId = entry.getKey(); + String pathPrefix = "/" + agentId; + registerAgentRoutes(router, pathPrefix, entry.getValue()); + } + } else if (!jsonRpcHandler.isUnsatisfied()) { + // Single-agent mode (default) + registerAgentRoutes(router, "", jsonRpcHandler.get()); + } + } + + private void registerAgentRoutes(Router router, String pathPrefix, JSONRPCHandler handler) { + String rpcPath = pathPrefix.isEmpty() ? "/" : pathPrefix; + String cardPath = pathPrefix + "/.well-known/agent-card.json"; + + router.post(rpcPath) .consumes(APPLICATION_JSON) .handler(BodyHandler.create()) .blockingHandler(ctx -> { try { vertxSecurityHelper.runInRequestContextDeferred(ctx, () -> { - invokeJSONRPCHandler(ctx.body().asString(), ctx); + invokeJSONRPCHandler(ctx.body().asString(), ctx, handler); }); } catch (UnauthorizedException | ForbiddenException e) { vertxSecurityHelper.handleAuthError(ctx, e); @@ -219,12 +237,11 @@ void setupRoutes(@Observes Router router) { } }, false); - // Agent card endpoint: GET /.well-known/agent-card.json - router.get("/.well-known/agent-card.json") + router.get(cardPath) .produces(APPLICATION_JSON) .handler(ctx -> { try { - String agentCard = getAgentCard(ctx); + String agentCard = getAgentCard(ctx, handler); ctx.response() .setStatusCode(200) .putHeader(CONTENT_TYPE, APPLICATION_JSON) @@ -308,7 +325,7 @@ void setupRoutes(@Observes Router router) { * @throws A2AError if request processing fails */ @Authenticated - public void invokeJSONRPCHandler(String body, RoutingContext rc) { + public void invokeJSONRPCHandler(String body, RoutingContext rc, JSONRPCHandler handler) { boolean streaming = false; ServerCallContext context = createCallContext(rc); A2AResponse nonStreamingResponse = null; @@ -318,10 +335,10 @@ public void invokeJSONRPCHandler(String body, RoutingContext rc) { A2ARequest request = JSONRPCUtils.parseRequestBody(body, extractTenant(rc)); context.getState().put(METHOD_NAME_KEY, request.getMethod()); if (request instanceof NonStreamingJSONRPCRequest nonStreamingRequest) { - nonStreamingResponse = processNonStreamingRequest(nonStreamingRequest, context); + nonStreamingResponse = processNonStreamingRequest(nonStreamingRequest, context, handler); } else { streaming = true; - streamingResponse = processStreamingRequest(request, context); + streamingResponse = processStreamingRequest(request, context, handler); } } catch (A2AError e) { error = new A2AErrorResponse(e); @@ -406,10 +423,10 @@ public void invokeJSONRPCHandler(String body, RoutingContext rc) { * @throws JsonProcessingException if serialization fails * @see JSONRPCHandler#getAgentCard() */ - public String getAgentCard(RoutingContext rc) throws JsonProcessingException { + public String getAgentCard(RoutingContext rc, JSONRPCHandler handler) throws JsonProcessingException { // Add caching headers per A2A specification section 8.6 cacheMetadata.getHttpHeadersMap().forEach((k, v) -> rc.response().putHeader(k, v)); - return JsonUtil.toJson(jsonRpcHandler.getAgentCard()); + return JsonUtil.toJson(handler.getAgentCard()); } /** @@ -435,33 +452,33 @@ public String getAgentCard(RoutingContext rc) throws JsonProcessingException { * @param context the server call context * @return the JSON-RPC response */ - private A2AResponse processNonStreamingRequest(NonStreamingJSONRPCRequest request, ServerCallContext context) { + private A2AResponse processNonStreamingRequest(NonStreamingJSONRPCRequest request, ServerCallContext context, JSONRPCHandler handler) { if (request instanceof GetTaskRequest req) { - return jsonRpcHandler.onGetTask(req, context); + return handler.onGetTask(req, context); } if (request instanceof CancelTaskRequest req) { - return jsonRpcHandler.onCancelTask(req, context); + return handler.onCancelTask(req, context); } if (request instanceof ListTasksRequest req) { - return jsonRpcHandler.onListTasks(req, context); + return handler.onListTasks(req, context); } if (request instanceof CreateTaskPushNotificationConfigRequest req) { - return jsonRpcHandler.setPushNotificationConfig(req, context); + return handler.setPushNotificationConfig(req, context); } if (request instanceof GetTaskPushNotificationConfigRequest req) { - return jsonRpcHandler.getPushNotificationConfig(req, context); + return handler.getPushNotificationConfig(req, context); } if (request instanceof SendMessageRequest req) { - return jsonRpcHandler.onMessageSend(req, context); + return handler.onMessageSend(req, context); } if (request instanceof ListTaskPushNotificationConfigsRequest req) { - return jsonRpcHandler.listPushNotificationConfigs(req, context); + return handler.listPushNotificationConfigs(req, context); } if (request instanceof DeleteTaskPushNotificationConfigRequest req) { - return jsonRpcHandler.deletePushNotificationConfig(req, context); + return handler.deletePushNotificationConfig(req, context); } if (request instanceof GetExtendedAgentCardRequest req) { - return jsonRpcHandler.onGetExtendedCardRequest(req, context); + return handler.onGetExtendedCardRequest(req, context); } return generateErrorResponse(request, new UnsupportedOperationError()); } @@ -483,20 +500,20 @@ private A2AResponse processNonStreamingRequest(NonStreamingJSONRPCRequest * @return a Multi stream of JSON-RPC responses */ private Multi> processStreamingRequest( - A2ARequest request, ServerCallContext context) throws A2AError { + A2ARequest request, ServerCallContext context, JSONRPCHandler handler) throws A2AError { if (request instanceof SendStreamingMessageRequest req) { - jsonRpcHandler.authorizeTaskAccess(req.getParams().message().taskId(), context, + handler.authorizeTaskAccess(req.getParams().message().taskId(), context, TaskOperation.MESSAGE_SEND_STREAM); } else if (request instanceof SubscribeToTaskRequest req) { - jsonRpcHandler.authorizeTaskAccess(req.getParams().id(), context, + handler.authorizeTaskAccess(req.getParams().id(), context, TaskOperation.SUBSCRIBE_TO_TASK); } try { Flow.Publisher> publisher; if (request instanceof SendStreamingMessageRequest req) { - publisher = jsonRpcHandler.onMessageSendStream(req, context); + publisher = handler.onMessageSendStream(req, context); } else if (request instanceof SubscribeToTaskRequest req) { - publisher = jsonRpcHandler.onSubscribeToTask(req, context); + publisher = handler.onSubscribeToTask(req, context); } else { return Multi.createFrom().item(generateErrorResponse(request, new UnsupportedOperationError())); } diff --git a/reference/jsonrpc/src/main/java/org/a2aproject/sdk/server/apps/quarkus/registry/MultiAgentRegistry.java b/reference/jsonrpc/src/main/java/org/a2aproject/sdk/server/apps/quarkus/registry/MultiAgentRegistry.java new file mode 100644 index 000000000..c86d6c7c8 --- /dev/null +++ b/reference/jsonrpc/src/main/java/org/a2aproject/sdk/server/apps/quarkus/registry/MultiAgentRegistry.java @@ -0,0 +1,16 @@ +package org.a2aproject.sdk.server.apps.quarkus.registry; + +import java.util.Map; +import org.a2aproject.sdk.transport.jsonrpc.handler.JSONRPCHandler; + +/** + * Registry for supporting multiple agents in a single Quarkus application. + * If a CDI bean implements this interface, the server will register routes for each + * agent in the registry under // and //.well-known/agent-card.json. + */ +public interface MultiAgentRegistry { + /** + * @return a map of agent ID (path segment) to their JSONRPCHandler + */ + Map getAgents(); +} diff --git a/reference/jsonrpc/src/test/java/org/a2aproject/sdk/server/apps/quarkus/A2AServerRoutesTest.java b/reference/jsonrpc/src/test/java/org/a2aproject/sdk/server/apps/quarkus/A2AServerRoutesTest.java index f2a746b2c..955c166eb 100644 --- a/reference/jsonrpc/src/test/java/org/a2aproject/sdk/server/apps/quarkus/A2AServerRoutesTest.java +++ b/reference/jsonrpc/src/test/java/org/a2aproject/sdk/server/apps/quarkus/A2AServerRoutesTest.java @@ -75,7 +75,8 @@ public class A2AServerRoutesTest { private A2AServerRoutes routes; - private JSONRPCHandler mockJsonRpcHandler; + private Instance mockJsonRpcHandler; + private JSONRPCHandler mockHandlerInstance; private Executor mockExecutor; private Instance mockCallContextFactory; private RoutingContext mockRoutingContext; @@ -87,7 +88,14 @@ public class A2AServerRoutesTest { @BeforeEach public void setUp() { routes = new A2AServerRoutes(); - mockJsonRpcHandler = mock(JSONRPCHandler.class); + mockJsonRpcHandler = mock(Instance.class); + mockHandlerInstance = mock(JSONRPCHandler.class); + when(mockHandlerInstance.isUnsatisfied()).thenReturn(false); + when(mockHandlerInstance.get()).thenReturn(mockHandlerInstance); + + Instance multiAgentRegistry = mock(Instance.class); + when(multiAgentRegistry.isUnsatisfied()).thenReturn(true); + setField(routes, "multiAgentRegistry", multiAgentRegistry); mockExecutor = mock(Executor.class); mockCallContextFactory = mock(Instance.class); mockRoutingContext = mock(RoutingContext.class); @@ -154,16 +162,16 @@ public void testSendMessage_MethodNameSetInContext() { .status(new TaskStatus(TaskState.TASK_STATE_SUBMITTED)) .build(); SendMessageResponse realResponse = new SendMessageResponse("1", responseTask); - when(mockJsonRpcHandler.onMessageSend(any(SendMessageRequest.class), any(ServerCallContext.class))) + when(mockHandlerInstance.onMessageSend(any(SendMessageRequest.class), any(ServerCallContext.class))) .thenReturn(realResponse); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).onMessageSend(any(SendMessageRequest.class), contextCaptor.capture()); + verify(mockHandlerInstance).onMessageSend(any(SendMessageRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); assertEquals(SEND_MESSAGE_METHOD, capturedContext.getState().get(METHOD_NAME_KEY)); @@ -201,16 +209,16 @@ public void testSendStreamingMessage_MethodNameSetInContext() { @SuppressWarnings("unchecked") Flow.Publisher mockPublisher = mock(Flow.Publisher.class); - when(mockJsonRpcHandler.onMessageSendStream(any(SendStreamingMessageRequest.class), + when(mockHandlerInstance.onMessageSendStream(any(SendStreamingMessageRequest.class), any(ServerCallContext.class))).thenReturn(mockPublisher); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).onMessageSendStream(any(SendStreamingMessageRequest.class), + verify(mockHandlerInstance).onMessageSendStream(any(SendStreamingMessageRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); @@ -239,16 +247,16 @@ public void testGetTask_MethodNameSetInContext() { .status(new TaskStatus(TaskState.TASK_STATE_SUBMITTED)) .build(); GetTaskResponse realResponse = new GetTaskResponse("1", responseTask); - when(mockJsonRpcHandler.onGetTask(any(GetTaskRequest.class), any(ServerCallContext.class))) + when(mockHandlerInstance.onGetTask(any(GetTaskRequest.class), any(ServerCallContext.class))) .thenReturn(realResponse); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).onGetTask(any(GetTaskRequest.class), contextCaptor.capture()); + verify(mockHandlerInstance).onGetTask(any(GetTaskRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); assertEquals(GET_TASK_METHOD, capturedContext.getState().get(METHOD_NAME_KEY)); @@ -276,16 +284,16 @@ public void testCancelTask_MethodNameSetInContext() { .status(new TaskStatus(TaskState.TASK_STATE_CANCELED)) .build(); CancelTaskResponse realResponse = new CancelTaskResponse("1", responseTask); - when(mockJsonRpcHandler.onCancelTask(any(CancelTaskRequest.class), any(ServerCallContext.class))) + when(mockHandlerInstance.onCancelTask(any(CancelTaskRequest.class), any(ServerCallContext.class))) .thenReturn(realResponse); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).onCancelTask(any(CancelTaskRequest.class), contextCaptor.capture()); + verify(mockHandlerInstance).onCancelTask(any(CancelTaskRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); assertEquals(CANCEL_TASK_METHOD, capturedContext.getState().get(METHOD_NAME_KEY)); @@ -308,16 +316,16 @@ public void testTaskResubscription_MethodNameSetInContext() { @SuppressWarnings("unchecked") Flow.Publisher mockPublisher = mock(Flow.Publisher.class); - when(mockJsonRpcHandler.onSubscribeToTask(any(SubscribeToTaskRequest.class), + when(mockHandlerInstance.onSubscribeToTask(any(SubscribeToTaskRequest.class), any(ServerCallContext.class))).thenReturn(mockPublisher); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).onSubscribeToTask(any(SubscribeToTaskRequest.class), + verify(mockHandlerInstance).onSubscribeToTask(any(SubscribeToTaskRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); @@ -353,16 +361,16 @@ public void testCreateTaskPushNotificationConfig_MethodNameSetInContext() { .build(); CreateTaskPushNotificationConfigResponse realResponse = new CreateTaskPushNotificationConfigResponse("1", responseConfig); - when(mockJsonRpcHandler.setPushNotificationConfig(any(CreateTaskPushNotificationConfigRequest.class), + when(mockHandlerInstance.setPushNotificationConfig(any(CreateTaskPushNotificationConfigRequest.class), any(ServerCallContext.class))).thenReturn(realResponse); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).setPushNotificationConfig(any(CreateTaskPushNotificationConfigRequest.class), + verify(mockHandlerInstance).setPushNotificationConfig(any(CreateTaskPushNotificationConfigRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); @@ -392,16 +400,16 @@ public void testGetTaskPushNotificationConfig_MethodNameSetInContext() { .url("https://example.com/callback") .build(); GetTaskPushNotificationConfigResponse realResponse = new GetTaskPushNotificationConfigResponse("1", responseConfig); - when(mockJsonRpcHandler.getPushNotificationConfig(any(GetTaskPushNotificationConfigRequest.class), + when(mockHandlerInstance.getPushNotificationConfig(any(GetTaskPushNotificationConfigRequest.class), any(ServerCallContext.class))).thenReturn(realResponse); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).getPushNotificationConfig(any(GetTaskPushNotificationConfigRequest.class), + verify(mockHandlerInstance).getPushNotificationConfig(any(GetTaskPushNotificationConfigRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); @@ -432,16 +440,16 @@ public void testListTaskPushNotificationConfigs_MethodNameSetInContext() { .url("https://example.com/callback") .build(); ListTaskPushNotificationConfigsResponse realResponse = new ListTaskPushNotificationConfigsResponse("1", new ListTaskPushNotificationConfigsResult(singletonList(config))); - when(mockJsonRpcHandler.listPushNotificationConfigs(any(ListTaskPushNotificationConfigsRequest.class), + when(mockHandlerInstance.listPushNotificationConfigs(any(ListTaskPushNotificationConfigsRequest.class), any(ServerCallContext.class))).thenReturn(realResponse); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).listPushNotificationConfigs(any(ListTaskPushNotificationConfigsRequest.class), + verify(mockHandlerInstance).listPushNotificationConfigs(any(ListTaskPushNotificationConfigsRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); @@ -466,16 +474,16 @@ public void testDeleteTaskPushNotificationConfig_MethodNameSetInContext() { // Create a real response with id DeleteTaskPushNotificationConfigResponse realResponse = new DeleteTaskPushNotificationConfigResponse("1"); - when(mockJsonRpcHandler.deletePushNotificationConfig(any(DeleteTaskPushNotificationConfigRequest.class), + when(mockHandlerInstance.deletePushNotificationConfig(any(DeleteTaskPushNotificationConfigRequest.class), any(ServerCallContext.class))).thenReturn(realResponse); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).deletePushNotificationConfig(any(DeleteTaskPushNotificationConfigRequest.class), + verify(mockHandlerInstance).deletePushNotificationConfig(any(DeleteTaskPushNotificationConfigRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); @@ -502,17 +510,17 @@ public void testGetExtendedCard_MethodNameSetInContext() { .supportedInterfaces(Collections.singletonList(new AgentInterface("jsonrpc", "http://localhost:9999"))) .build(); GetExtendedAgentCardResponse realResponse = new GetExtendedAgentCardResponse(1, agentCard); - when(mockJsonRpcHandler.onGetExtendedCardRequest( + when(mockHandlerInstance.onGetExtendedCardRequest( any(GetExtendedAgentCardRequest.class), any(ServerCallContext.class))) .thenReturn(realResponse); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).onGetExtendedCardRequest( + verify(mockHandlerInstance).onGetExtendedCardRequest( any(GetExtendedAgentCardRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); @@ -542,16 +550,16 @@ public void testTenantExtraction_MultiSegmentPath() { .status(new TaskStatus(TaskState.TASK_STATE_SUBMITTED)) .build(); GetTaskResponse realResponse = new GetTaskResponse("1", responseTask); - when(mockJsonRpcHandler.onGetTask(any(GetTaskRequest.class), any(ServerCallContext.class))) + when(mockHandlerInstance.onGetTask(any(GetTaskRequest.class), any(ServerCallContext.class))) .thenReturn(realResponse); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).onGetTask(any(GetTaskRequest.class), contextCaptor.capture()); + verify(mockHandlerInstance).onGetTask(any(GetTaskRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); assertEquals("test/titi", capturedContext.getState().get(TENANT_KEY)); @@ -579,16 +587,16 @@ public void testTenantExtraction_RootPath() { .status(new TaskStatus(TaskState.TASK_STATE_SUBMITTED)) .build(); GetTaskResponse realResponse = new GetTaskResponse("1", responseTask); - when(mockJsonRpcHandler.onGetTask(any(GetTaskRequest.class), any(ServerCallContext.class))) + when(mockHandlerInstance.onGetTask(any(GetTaskRequest.class), any(ServerCallContext.class))) .thenReturn(realResponse); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).onGetTask(any(GetTaskRequest.class), contextCaptor.capture()); + verify(mockHandlerInstance).onGetTask(any(GetTaskRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); assertEquals("", capturedContext.getState().get(TENANT_KEY)); @@ -616,16 +624,16 @@ public void testTenantExtraction_SingleSegmentPath() { .status(new TaskStatus(TaskState.TASK_STATE_SUBMITTED)) .build(); GetTaskResponse realResponse = new GetTaskResponse("1", responseTask); - when(mockJsonRpcHandler.onGetTask(any(GetTaskRequest.class), any(ServerCallContext.class))) + when(mockHandlerInstance.onGetTask(any(GetTaskRequest.class), any(ServerCallContext.class))) .thenReturn(realResponse); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).onGetTask(any(GetTaskRequest.class), contextCaptor.capture()); + verify(mockHandlerInstance).onGetTask(any(GetTaskRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); assertEquals("tenant1", capturedContext.getState().get(TENANT_KEY)); @@ -653,16 +661,16 @@ public void testTenantExtraction_ThreeSegmentPath() { .status(new TaskStatus(TaskState.TASK_STATE_SUBMITTED)) .build(); GetTaskResponse realResponse = new GetTaskResponse("1", responseTask); - when(mockJsonRpcHandler.onGetTask(any(GetTaskRequest.class), any(ServerCallContext.class))) + when(mockHandlerInstance.onGetTask(any(GetTaskRequest.class), any(ServerCallContext.class))) .thenReturn(realResponse); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).onGetTask(any(GetTaskRequest.class), contextCaptor.capture()); + verify(mockHandlerInstance).onGetTask(any(GetTaskRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); assertEquals("tenant1/api/v1", capturedContext.getState().get(TENANT_KEY)); @@ -700,16 +708,16 @@ public void testTenantExtraction_StreamingRequest() { @SuppressWarnings("unchecked") Flow.Publisher mockPublisher = mock(Flow.Publisher.class); - when(mockJsonRpcHandler.onMessageSendStream(any(SendStreamingMessageRequest.class), + when(mockHandlerInstance.onMessageSendStream(any(SendStreamingMessageRequest.class), any(ServerCallContext.class))).thenReturn(mockPublisher); ArgumentCaptor contextCaptor = ArgumentCaptor.forClass(ServerCallContext.class); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert - verify(mockJsonRpcHandler).onMessageSendStream(any(SendStreamingMessageRequest.class), + verify(mockHandlerInstance).onMessageSendStream(any(SendStreamingMessageRequest.class), contextCaptor.capture()); ServerCallContext capturedContext = contextCaptor.getValue(); assertNotNull(capturedContext); @@ -723,7 +731,7 @@ public void testJsonParseError_ContentTypeIsApplicationJson() { when(mockRequestBody.asString()).thenReturn(invalidJson); // Act - routes.invokeJSONRPCHandler(invalidJson, mockRoutingContext); + routes.invokeJSONRPCHandler(invalidJson, mockRoutingContext, mockHandlerInstance); // Assert verify(mockHttpResponse).putHeader(CONTENT_TYPE, APPLICATION_JSON); @@ -742,7 +750,7 @@ public void testMethodNotFound_ContentTypeIsApplicationJson() { when(mockRequestBody.asString()).thenReturn(jsonRpcRequest); // Act - routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext); + routes.invokeJSONRPCHandler(jsonRpcRequest, mockRoutingContext, mockHandlerInstance); // Assert verify(mockHttpResponse).putHeader(CONTENT_TYPE, APPLICATION_JSON);