|
| 1 | +@file:OptIn(UnstableApi::class) |
| 2 | + |
| 3 | +package com.agentclientprotocol.agent |
| 4 | + |
| 5 | +import com.agentclientprotocol.annotations.UnstableApi |
| 6 | +import com.agentclientprotocol.client.ClientInfo |
| 7 | +import com.agentclientprotocol.common.Event |
| 8 | +import com.agentclientprotocol.common.SessionCreationParameters |
| 9 | +import com.agentclientprotocol.model.AcpMethod |
| 10 | +import com.agentclientprotocol.model.AcpNotification |
| 11 | +import com.agentclientprotocol.model.AcpRequest |
| 12 | +import com.agentclientprotocol.model.AcpResponse |
| 13 | +import com.agentclientprotocol.model.CancelNotification |
| 14 | +import com.agentclientprotocol.model.ContentBlock |
| 15 | +import com.agentclientprotocol.model.InitializeRequest |
| 16 | +import com.agentclientprotocol.model.LATEST_PROTOCOL_VERSION |
| 17 | +import com.agentclientprotocol.model.McpServer |
| 18 | +import com.agentclientprotocol.model.NewSessionRequest |
| 19 | +import com.agentclientprotocol.model.PromptRequest |
| 20 | +import com.agentclientprotocol.model.PromptResponse |
| 21 | +import com.agentclientprotocol.model.SessionId |
| 22 | +import com.agentclientprotocol.model.SessionUpdate |
| 23 | +import com.agentclientprotocol.model.StopReason |
| 24 | +import com.agentclientprotocol.protocol.Protocol |
| 25 | +import com.agentclientprotocol.rpc.ACPJson |
| 26 | +import com.agentclientprotocol.rpc.JsonRpcNotification |
| 27 | +import com.agentclientprotocol.rpc.JsonRpcResponse |
| 28 | +import kotlinx.atomicfu.atomic |
| 29 | +import kotlinx.coroutines.CoroutineScope |
| 30 | +import kotlinx.coroutines.delay |
| 31 | +import kotlinx.coroutines.flow.Flow |
| 32 | +import kotlinx.coroutines.flow.FlowCollector |
| 33 | +import kotlinx.coroutines.flow.flow |
| 34 | +import kotlinx.coroutines.runBlocking |
| 35 | +import kotlinx.serialization.json.JsonElement |
| 36 | +import kotlin.time.Duration |
| 37 | +import kotlin.time.Duration.Companion.milliseconds |
| 38 | +import kotlin.time.Duration.Companion.seconds |
| 39 | + |
| 40 | +class TestAgent(val agent: Agent, val agentSupport: TestAgentSupport, val transport: TestTransport) { |
| 41 | + suspend fun <TRequest : AcpRequest, TResponse : AcpResponse> testRequest( |
| 42 | + method: AcpMethod.AcpRequestResponseMethod<TRequest, TResponse>, |
| 43 | + request: TRequest |
| 44 | + ): Pair<TResponse?, List<JsonRpcNotification>> { |
| 45 | + val received = transport.fireTestRequest( |
| 46 | + methodName = method.methodName, |
| 47 | + params = ACPJson.encodeToJsonElement(method.requestSerializer, request) |
| 48 | + ) |
| 49 | + val response = (received.lastOrNull() as? JsonRpcResponse)?.result?.let { |
| 50 | + ACPJson.decodeFromJsonElement(method.responseSerializer, it) |
| 51 | + } |
| 52 | + val notifications = received.filterIsInstance<JsonRpcNotification>() |
| 53 | + return response to notifications |
| 54 | + } |
| 55 | + |
| 56 | + fun <TNotification : AcpNotification> testNotification( |
| 57 | + method: AcpMethod.AcpNotificationMethod<TNotification>, |
| 58 | + notification: TNotification |
| 59 | + ) { |
| 60 | + transport.fireTestNotification(method.methodName, ACPJson.encodeToJsonElement(method.serializer, notification)) |
| 61 | + } |
| 62 | + |
| 63 | + fun close() { |
| 64 | + agent.protocol.close() |
| 65 | + } |
| 66 | + |
| 67 | + suspend fun testInitialize(request: InitializeRequest) = testRequest(AcpMethod.AgentMethods.Initialize, request) |
| 68 | + suspend fun testNewSession(request: NewSessionRequest) = testRequest(AcpMethod.AgentMethods.SessionNew, request) |
| 69 | + suspend fun testPrompt(request: PromptRequest) = testRequest(AcpMethod.AgentMethods.SessionPrompt, request) |
| 70 | + |
| 71 | + fun testCancel(notification: CancelNotification) = testNotification(AcpMethod.AgentMethods.SessionCancel, notification) |
| 72 | +} |
| 73 | + |
| 74 | +suspend fun TestAgent.simplePrompt(prompt: String): Pair<PromptResponse, List<SessionUpdate>> { |
| 75 | + val session = agentSupport.createdSessions.values.single() |
| 76 | + val (resp, notifications) = testPrompt(PromptRequest(session.sessionId, listOf(ContentBlock.Text(prompt)))) |
| 77 | + checkNotNull(resp) |
| 78 | + |
| 79 | + return resp to notifications |
| 80 | + .filter { it.method == AcpMethod.ClientMethods.SessionUpdate.methodName } |
| 81 | + .mapNotNull { it.params } |
| 82 | + .map { ACPJson.decodeFromJsonElement(AcpMethod.ClientMethods.SessionUpdate.serializer, it).update } |
| 83 | +} |
| 84 | + |
| 85 | +class TestAgentSupport(val promptHandler: PromptHandler) : AgentSupport { |
| 86 | + var isInitialized = false |
| 87 | + val createdSessions = mutableMapOf<SessionId, TestAgentSession>() |
| 88 | + |
| 89 | + override suspend fun initialize(clientInfo: ClientInfo): AgentInfo { |
| 90 | + isInitialized = true |
| 91 | + return AgentInfo() |
| 92 | + } |
| 93 | + |
| 94 | + override suspend fun createSession(sessionParameters: SessionCreationParameters): AgentSession { |
| 95 | + val sessionId = SessionId("test-agent-session-${sessionId.incrementAndGet()}") |
| 96 | + val session = TestAgentSession(sessionId, promptHandler) |
| 97 | + createdSessions[sessionId] = session |
| 98 | + return session |
| 99 | + } |
| 100 | + |
| 101 | + companion object { |
| 102 | + private val sessionId = atomic(0) |
| 103 | + } |
| 104 | +} |
| 105 | + |
| 106 | +typealias PromptHandler = suspend FlowCollector<Event>.(List<ContentBlock>) -> Unit |
| 107 | + |
| 108 | +class TestAgentSession( |
| 109 | + override val sessionId: SessionId, |
| 110 | + val promptHandler: PromptHandler |
| 111 | +) : AgentSession { |
| 112 | + override suspend fun prompt(content: List<ContentBlock>, _meta: JsonElement?): Flow<Event> = flow { |
| 113 | + promptHandler(content) |
| 114 | + } |
| 115 | +} |
| 116 | + |
| 117 | +fun withTestAgent( |
| 118 | + timeout: Duration = 5.seconds, |
| 119 | + promptHandler: PromptHandler = echoPromptHandler, |
| 120 | + block: suspend CoroutineScope.(TestAgent) -> Unit |
| 121 | +) = runBlocking { |
| 122 | + val transport = TestTransport(timeout) |
| 123 | + val protocol = Protocol(this, transport) |
| 124 | + val agentSupport = TestAgentSupport(promptHandler) |
| 125 | + val agent = Agent(protocol, agentSupport) |
| 126 | + protocol.start() |
| 127 | + |
| 128 | + // wait a little after protocol start, if messages get sent right away they can get lost |
| 129 | + delay(100.milliseconds) |
| 130 | + |
| 131 | + val testAgent = TestAgent(agent, agentSupport, transport) |
| 132 | + block(testAgent) |
| 133 | + testAgent.close() |
| 134 | +} |
| 135 | + |
| 136 | +fun withInitializedTestAgent( |
| 137 | + timeout: Duration = 5.seconds, |
| 138 | + promptHandler: PromptHandler = echoPromptHandler, |
| 139 | + block: suspend CoroutineScope.(TestAgent) -> Unit |
| 140 | +) = withTestAgent( |
| 141 | + timeout = timeout, |
| 142 | + promptHandler = promptHandler, |
| 143 | +) { testAgent -> |
| 144 | + testAgent.testInitialize(InitializeRequest(LATEST_PROTOCOL_VERSION)) |
| 145 | + check(testAgent.agentSupport.isInitialized) |
| 146 | + block(testAgent) |
| 147 | +} |
| 148 | + |
| 149 | +fun withTestAgentSession( |
| 150 | + timeout: Duration = 5.seconds, |
| 151 | + promptHandler: PromptHandler = echoPromptHandler, |
| 152 | + cwd: String = ".", |
| 153 | + mcpServers: List<McpServer> = emptyList(), |
| 154 | + block: suspend CoroutineScope.(TestAgent, TestAgentSession) -> Unit |
| 155 | +) = withInitializedTestAgent( |
| 156 | + timeout = timeout, |
| 157 | + promptHandler = promptHandler, |
| 158 | +) { testAgent -> |
| 159 | + val (newSessionResponse) = testAgent.testNewSession(NewSessionRequest(cwd, mcpServers)) |
| 160 | + checkNotNull(newSessionResponse) |
| 161 | + val session = testAgent.agentSupport.createdSessions[newSessionResponse.sessionId] |
| 162 | + checkNotNull(session) |
| 163 | + block(testAgent, session) |
| 164 | +} |
| 165 | + |
| 166 | +val echoPromptHandler: PromptHandler = { prompt -> |
| 167 | + prompt.filterIsInstance<ContentBlock.Text>().forEach { |
| 168 | + emit(Event.SessionUpdateEvent(SessionUpdate.AgentMessageChunk(it))) |
| 169 | + } |
| 170 | + emit(Event.PromptResponseEvent(PromptResponse(StopReason.END_TURN))) |
| 171 | +} |
| 172 | + |
| 173 | +fun delayEchoPromptHandler(delay: Duration): PromptHandler = { prompt -> |
| 174 | + delay(delay) |
| 175 | + prompt.filterIsInstance<ContentBlock.Text>().forEach { |
| 176 | + emit(Event.SessionUpdateEvent(SessionUpdate.AgentMessageChunk(it))) |
| 177 | + } |
| 178 | + emit(Event.PromptResponseEvent(PromptResponse(StopReason.END_TURN))) |
| 179 | +} |
0 commit comments