diff --git a/apps/server/src/mcp/OrchestratorMcpService.test.ts b/apps/server/src/mcp/OrchestratorMcpService.test.ts index eaa776173af3..6f753fae31cf 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.test.ts @@ -1355,6 +1355,108 @@ describe("OrchestratorMcpService provider resolution", () => { }), ); + it.effect("re-probes a driver-only target's instances once before refusing it", () => + Effect.gen(function* () { + const claudeInstanceId = ProviderInstanceId.make("claudeAgent"); + const claudeDriver = ProviderDriverKind.make("claudeAgent"); + const task = { + id: taskId, + threadId: parentThreadId, + runId: parentRunId, + parentNodeId, + origin: "app_owned", + createdBy: "agent", + driver: claudeDriver, + providerInstanceId: claudeInstanceId, + providerThreadId: null, + childThreadId, + nativeTaskRef: null, + prompt: "Review the diff.", + title: null, + model: "claude-opus-5-5", + status: "running", + result: null, + startedAt: null, + completedAt: null, + }; + const healthy = providerSnapshot({ + instanceId: claudeInstanceId, + driver: claudeDriver, + model: "claude-opus-5-5", + }); + const timedOut: ServerProvider = { + ...healthy, + status: "error", + message: + "Claude Agent CLI is installed but failed to run. Timed out while running command.", + }; + const codex = providerSnapshot({ + instanceId: codexInstanceId, + driver: ProviderDriverKind.make("codex"), + model: "gpt-5.4", + }); + // The cache still holds a startup probe timeout; a re-probe reports `probeSucceeds`. + let probeSucceeds = false; + const probed = yield* Ref.make>([]); + const dispatched = yield* Ref.make(0); + const layerDependencies = Layer.mergeAll( + NodeServices.layer, + Layer.mock(ThreadManagementService.ThreadManagementService)({ + getThreadRecords: (threadId) => + Effect.succeed( + threadId === parentThreadId ? parentProjection([task]) : childProjection, + ), + dispatch: () => + Ref.update(dispatched, (count) => count + 1).pipe( + Effect.as({ + sequence: 1, + storedEvents: [ + { + sequence: 1, + commandId: null, + event: { type: "subagent.updated", payload: task }, + }, + ], + } as never), + ), + }), + Layer.mock(ProviderRegistry.ProviderRegistry)({ + getProviders: Effect.succeed([codex, timedOut]), + refreshInstance: (instanceId) => + Ref.update(probed, (ids) => [...ids, instanceId]).pipe( + Effect.as([codex, probeSucceeds ? healthy : timedOut]), + ), + }), + adapterRegistryLayer([codexInstanceId, claudeInstanceId]), + Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), + Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), + ); + + yield* Effect.gen(function* () { + const service = yield* OrchestratorMcpService.OrchestratorMcpService; + const delegate = (clientRequestId: string) => + service.delegateTask(scope, { + task: "Review the diff.", + target: { driverKind: claudeDriver }, + mode: "async", + clientRequestId, + }); + + const error = yield* delegate("delegate-driver-recheck-1").pipe(Effect.flip); + assert.equal(error.code, "provider_unavailable"); + assert.deepStrictEqual(yield* Ref.get(probed), [claudeInstanceId]); + assert.equal(yield* Ref.get(dispatched), 0); + + probeSucceeds = true; + const result = yield* delegate("delegate-driver-recheck-2"); + assert.equal(result.providerInstanceId, claudeInstanceId); + assert.deepStrictEqual(yield* Ref.get(probed), [claudeInstanceId, claudeInstanceId]); + assert.equal(yield* Ref.get(dispatched), 1); + }).pipe(Effect.provide(OrchestratorMcpService.layer.pipe(Layer.provide(layerDependencies)))); + }), + ); + it.effect( "inherits an available parent instance for driver-only targets and otherwise selects a healthy peer", () => diff --git a/apps/server/src/mcp/OrchestratorMcpService.ts b/apps/server/src/mcp/OrchestratorMcpService.ts index b218a76d1ab5..f459ef6af1aa 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.ts @@ -1065,23 +1065,35 @@ const make = Effect.gen(function* () { /** * Provider snapshots only re-probe while a client is in the foreground, so an * unattended agent can see a provider as unavailable after it was fixed. - * Re-probe the requested instance once before refusing it. + * Re-probe the requested instance (or, for a driver-only target, each enabled + * and installed instance of that driver) once before refusing it. */ const resolveTargetRechecking = (input: Parameters[0]) => { - const instanceId = + const requestedDriver = input.target?.driverKind; + const requestedInstanceId = input.target?.providerInstanceId ?? - (input.target?.driverKind === undefined - ? input.parent.thread.modelSelection.instanceId - : undefined); + (requestedDriver === undefined ? input.parent.thread.modelSelection.instanceId : undefined); + const instanceIds = + requestedInstanceId === undefined + ? input.providers + .filter( + (provider) => + provider.driver === requestedDriver && provider.enabled && provider.installed, + ) + .map((provider) => provider.instanceId) + : [requestedInstanceId]; const resolved = resolveTarget(input); - if (instanceId === undefined) return resolved; + if (instanceIds.length === 0) return resolved; return resolved.pipe( Effect.catchIf( (error) => error.code === "provider_unavailable", () => - providerRegistry - .refreshInstance(instanceId) - .pipe(Effect.flatMap((providers) => resolveTarget({ ...input, providers }))), + // Each refresh returns the whole registry, so the last one carries every re-probe. + Effect.forEach(instanceIds, providerRegistry.refreshInstance).pipe( + Effect.flatMap((snapshots) => + resolveTarget({ ...input, providers: snapshots.at(-1) ?? input.providers }), + ), + ), ), ); };