diff --git a/apps/server/src/mcp/McpHttpServer.ts b/apps/server/src/mcp/McpHttpServer.ts index 5c8f146bcadf..f494ed46304d 100644 --- a/apps/server/src/mcp/McpHttpServer.ts +++ b/apps/server/src/mcp/McpHttpServer.ts @@ -11,13 +11,23 @@ import * as Schema from "effect/Schema"; import * as Sink from "effect/Sink"; import * as Stream from "effect/Stream"; import type * as Types from "effect/Types"; -import { McpProtocol, McpSchema, McpServer, Tool } from "effect/unstable/ai"; +import { McpProtocol, McpSchema, McpServer, Tool, type Toolkit } from "effect/unstable/ai"; import { HttpRouter, HttpServerRequest, HttpServerResponse } from "effect/unstable/http"; -import { PreviewAutomationError } from "@t3tools/contracts"; +import { PreviewAutomationError, ThreadId } from "@t3tools/contracts"; import packageJson from "../../package.json" with { type: "json" }; import * as ServerConfig from "../config.ts"; import * as DeviceService from "../device/DeviceService.ts"; +import * as ProviderRegistry from "../provider/Services/ProviderRegistry.ts"; +import * as AnalyticsService from "../telemetry/AnalyticsService.ts"; +import { + type AgentThread, + agentToolProperties, + handoffSettings, + MAX_REPORTED_DEPTH, + toolOutcome, +} from "../telemetry/ProviderDimensions.ts"; +import * as ThreadManagement from "../orchestration-v2/ThreadManagementService.ts"; import * as McpInvocationContext from "./McpInvocationContext.ts"; import * as OrchestratorMcpService from "./OrchestratorMcpService.ts"; import { PreviewControlsToolkit } from "./toolkits/previewControls/tools.ts"; @@ -647,61 +657,225 @@ const registerDeviceScreenshot = Effect.fn("McpHttpServer.registerDeviceScreensh ); }); -const PreviewStandardToolkitRegistrationLive = McpServer.toolkit(PreviewStandardToolkit).pipe( +/** + * Tools that hand work to another thread's agent: they create it, message it, + * steer it, or move context into it. Their events name the receiving agent too. + */ +const HANDOFF_TOOLS: ReadonlySet = new Set([ + "delegate_task", + "create_threads", + "t3_thread_launch", + "t3_thread_send", + "t3_thread_send_attachments", + "t3_thread_fork", + "t3_thread_merge_back", + "t3_queue_edit", + "t3_queue_promote_to_steer", + "t3_pending_request_respond", + "schedule_task", + "run_scheduled_task_now", +]); + +const resultThreadKeys = ["childThreadId", "targetThreadId", "boundThreadId", "threadId"] as const; + +/** Thread ids a handoff tool result names: the thread that received the work. */ +const resultThreadIds = (content: unknown): ReadonlyArray => { + if (typeof content !== "object" || content === null) return []; + const record = content as Readonly>; + if (Array.isArray(record.threads)) return record.threads.flatMap(resultThreadIds); + const key = resultThreadKeys.find((candidate) => typeof record[candidate] === "string"); + return key === undefined ? [] : [record[key] as string]; +}; + +/** + * Runs a tool registration against a server that records one anonymous + * `mcp.tool.invoked` event per call: the tool, the calling agent's provider, + * the outcome with its failure code, and the duration. A handoff tool also + * reports the caller's model, origin, and delegation depth, the settings the + * agent chose, and the provider, model, and modes of each thread that received + * the work. Ids, prompts, and free text are never recorded. + */ +const withToolAnalytics = (registration: Effect.Effect) => + Effect.gen(function* () { + const server = yield* McpServer.McpServer; + const analytics = yield* Effect.serviceOption(AnalyticsService.AnalyticsService); + if (Option.isNone(analytics)) return yield* registration; + const registry = yield* Effect.serviceOption(ProviderRegistry.ProviderRegistry); + const threads = yield* Effect.serviceOption(ThreadManagement.ThreadManagementService); + const shellOf = (threadId: string) => + Option.isNone(threads) + ? Effect.succeed(null) + : threads.value + .getThreadShell(ThreadId.make(threadId)) + .pipe(Effect.orElseSucceed(() => null)); + // Delegation depth: parents walked from the caller, capped. + const depthOf = (thread: AgentThread) => + Effect.gen(function* () { + let depth = 0; + let parentId = + thread.lineage.relationshipToParent === "subagent" ? thread.lineage.parentThreadId : null; + while (parentId !== null && depth < MAX_REPORTED_DEPTH) { + depth += 1; + const parent = yield* shellOf(parentId); + parentId = + parent?.lineage.relationshipToParent === "subagent" + ? parent.lineage.parentThreadId + : null; + } + return depth; + }); + // Whether the caller's active run was started by a scheduled task. + const scheduledRunOf = (threadId: string, runId: string | null) => + Option.isNone(threads) || runId === null + ? Effect.succeed(false) + : threads.value + .getThreadRecords(ThreadId.make(threadId), ["runs", "messages"], { + runIds: [runId as never], + messageRoles: ["user"], + }) + .pipe( + Effect.map((records) => { + const run = records.runs.find((candidate) => candidate.id === runId); + return records.messages.some( + (message) => + message.id === run?.userMessageId && message.scheduledTaskId !== undefined, + ); + }), + Effect.orElseSucceed(() => false), + ); + const record = ( + tool: string, + args: unknown, + result: McpSchema.CallToolResult | undefined, + durationMs: number, + ) => + Effect.gen(function* () { + const invocation = Option.getOrUndefined( + yield* Effect.serviceOption(McpInvocationContext.McpInvocationContext), + ); + const providers = Option.isSome(registry) ? yield* registry.value.getProviders : []; + // Only handoff tools pay for thread lookups; the rest report the + // caller's provider from the credential. + const handoff = HANDOFF_TOOLS.has(tool); + const caller = + handoff && invocation !== undefined + ? ((yield* shellOf(invocation.threadId)) ?? undefined) + : undefined; + const targetIds = handoff + ? [...new Set(resultThreadIds(result?.structuredContent))].filter( + (threadId) => threadId !== invocation?.threadId, + ) + : []; + const targets = yield* Effect.forEach(targetIds, shellOf); + for (const properties of agentToolProperties({ + tool, + providers, + caller, + ...(invocation === undefined + ? {} + : { callerProviderInstanceId: invocation.providerInstanceId }), + ...(caller === undefined || invocation === undefined + ? {} + : { + callerDepth: yield* depthOf(caller), + callerScheduledRun: yield* scheduledRunOf(invocation.threadId, caller.activeRunId), + }), + targets: targets.filter((thread) => thread !== null), + outcome: toolOutcome(result), + ...(handoff ? { settings: handoffSettings(tool, args, result?.structuredContent) } : {}), + durationMs, + })) { + yield* analytics.value.record("mcp.tool.invoked", properties); + } + }).pipe(Effect.ignoreCause); + const recordingServer = McpServer.McpServer.of({ + ...server, + addTool: (options) => + server.addTool({ + ...options, + handle: (payload) => + Effect.gen(function* () { + const startedAt = yield* Clock.currentTimeMillis; + return yield* options + .handle(payload) + .pipe( + Effect.onExit((exit) => + Clock.currentTimeMillis.pipe( + Effect.flatMap((endedAt) => + record( + options.tool.name, + payload, + exit._tag === "Success" ? exit.value : undefined, + endedAt - startedAt, + ), + ), + ), + ), + ); + }), + }), + }); + return yield* registration.pipe(Effect.provideService(McpServer.McpServer, recordingServer)); + }); + +const toolkit = >(tools: Toolkit.Toolkit) => + Layer.effectDiscard(withToolAnalytics(McpServer.registerToolkit(tools))).pipe( + Layer.provide(McpServer.McpServer.layer), + ); + +const PreviewStandardToolkitRegistrationLive = toolkit(PreviewStandardToolkit).pipe( Layer.provide(PreviewStandardToolkitHandlersLive), ); -const PreviewSnapshotRegistrationLive = Layer.effectDiscard(registerPreviewSnapshot()).pipe( - Layer.provide(PreviewSnapshotToolkitHandlersLive), -); +const PreviewSnapshotRegistrationLive = Layer.effectDiscard( + withToolAnalytics(registerPreviewSnapshot()), +).pipe(Layer.provide(PreviewSnapshotToolkitHandlersLive)); export const PreviewToolkitRegistrationLive = Layer.mergeAll( PreviewStandardToolkitRegistrationLive, PreviewSnapshotRegistrationLive, ); -export const OrchestratorToolkitRegistrationLive = McpServer.toolkit(OrchestratorToolkit).pipe( +export const OrchestratorToolkitRegistrationLive = toolkit(OrchestratorToolkit).pipe( Layer.provide(OrchestratorToolkitHandlersLive), Layer.provide(OrchestratorMcpService.layer), Layer.provide(ThreadMetadataMcpService.layer), ); -export const ThreadToolkitRegistrationLive = McpServer.toolkit(ThreadToolkit).pipe( +export const ThreadToolkitRegistrationLive = toolkit(ThreadToolkit).pipe( Layer.provide(ThreadToolkitHandlersLive), ); -const WorktreeToolkitRegistrationLive = McpServer.toolkit(WorktreeToolkit).pipe( +const WorktreeToolkitRegistrationLive = toolkit(WorktreeToolkit).pipe( Layer.provide(WorktreeToolkitHandlersLive), Layer.provide(WorktreeMcpService.layer), ); -const PreviewControlsRegistrationLive = McpServer.toolkit(PreviewControlsToolkit).pipe( +const PreviewControlsRegistrationLive = toolkit(PreviewControlsToolkit).pipe( Layer.provide(PreviewControlsHandlersLive), ); -const EnvironmentRegistrationLive = McpServer.toolkit(EnvironmentToolkit).pipe( +const EnvironmentRegistrationLive = toolkit(EnvironmentToolkit).pipe( Layer.provide(EnvironmentHandlersLive), ); -const ProjectRegistrationLive = McpServer.toolkit(ProjectToolkit).pipe( - Layer.provide(ProjectHandlersLive), -); +const ProjectRegistrationLive = toolkit(ProjectToolkit).pipe(Layer.provide(ProjectHandlersLive)); -const AttachmentRegistrationLive = McpServer.toolkit(AttachmentToolkit).pipe( +const AttachmentRegistrationLive = toolkit(AttachmentToolkit).pipe( Layer.provide(AttachmentHandlersLive), ); -export const PullRequestsToolkitRegistrationLive = McpServer.toolkit(PullRequestsToolkit).pipe( +export const PullRequestsToolkitRegistrationLive = toolkit(PullRequestsToolkit).pipe( Layer.provide(PullRequestsToolkitHandlersLive), ); -const DeviceStandardToolkitRegistrationLive = McpServer.toolkit(DeviceStandardToolkit).pipe( +const DeviceStandardToolkitRegistrationLive = toolkit(DeviceStandardToolkit).pipe( Layer.provide(DeviceStandardToolkitHandlersLive), ); -const DeviceScreenshotRegistrationLive = Layer.effectDiscard(registerDeviceScreenshot()).pipe( - Layer.provide(DeviceScreenshotToolkitHandlersLive), -); +const DeviceScreenshotRegistrationLive = Layer.effectDiscard( + withToolAnalytics(registerDeviceScreenshot()), +).pipe(Layer.provide(DeviceScreenshotToolkitHandlersLive)); export const DeviceToolkitRegistrationLive = Layer.mergeAll( DeviceStandardToolkitRegistrationLive, diff --git a/apps/server/src/mcp/toolkits/core.test.ts b/apps/server/src/mcp/toolkits/core.test.ts index 339e5eb9b2a8..58a4ea76add1 100644 --- a/apps/server/src/mcp/toolkits/core.test.ts +++ b/apps/server/src/mcp/toolkits/core.test.ts @@ -17,6 +17,8 @@ import { OrchestratorProjectionError } from "../../orchestration-v2/Orchestrator import * as ThreadManagement from "../../orchestration-v2/ThreadManagementService.ts"; import * as McpHttpServer from "../McpHttpServer.ts"; import * as McpInvocationContext from "../McpInvocationContext.ts"; +import * as ProviderRegistry from "../../provider/Services/ProviderRegistry.ts"; +import * as AnalyticsService from "../../telemetry/AnalyticsService.ts"; import { OrchestratorToolkit } from "./orchestrator/tools.ts"; import { PreviewToolkit } from "./preview/tools.ts"; import { PreviewControlsToolkit } from "./previewControls/tools.ts"; @@ -157,6 +159,105 @@ it.effect("returns a bounded public failure without serializing storage causes", ), ); +it.effect("records the calling agent and the agent a tool acted on", () => { + const recorded: Array<{ event: string; properties: unknown }> = []; + const shell = (id: ThreadId, instanceId: string, model: string) => + ({ + id, + projectId: "mcp-core-project", + providerInstanceId: ProviderInstanceId.make(instanceId), + modelSelection: { instanceId: ProviderInstanceId.make(instanceId), model }, + runtimeMode: "full-access", + interactionMode: "default", + archivedAt: null, + deletedAt: null, + activeRunId: "mcp-core-run", + createdBy: "agent", + creationSource: "mcp", + lineage: { + parentThreadId: id === threadId ? ThreadId.make("mcp-core-parent") : null, + relationshipToParent: id === threadId ? "subagent" : null, + rootThreadId: ThreadId.make("mcp-core-parent"), + }, + }) as never; + return Effect.gen(function* () { + const server = yield* McpServer.McpServer; + yield* server + .callTool({ name: "t3_thread_fork", arguments: { sourcePoint: { type: "latest_stable" } } }) + .pipe( + Effect.provideService(McpInvocationContext.McpInvocationContext, scope), + Effect.provideService(McpSchema.McpServerClient, client), + ); + expect(recorded).toHaveLength(1); + const [{ event, properties }] = recorded as [{ event: string; properties: object }]; + expect(event).toBe("mcp.tool.invoked"); + expect(properties).toMatchObject({ + tool: "t3_thread_fork", + outcome: "ok", + callerProvider: "codex", + callerModel: "gpt-5.5", + callerOrigin: "agent", + callerDepth: 1, + targetProvider: "claudeAgent", + targetModel: "claude-opus-5-5", + targetRuntimeMode: "full-access", + targetInteractionMode: "default", + crossProvider: true, + }); + expect(properties).toHaveProperty("durationMs"); + expect(Object.values(properties).some((value) => String(value).includes("mcp-core"))).toBe( + false, + ); + }).pipe( + Effect.provide( + McpHttpServer.ThreadToolkitRegistrationLive.pipe( + Layer.provideMerge(McpServer.McpServer.layer), + Layer.provide(NodeCrypto.layer), + Layer.provide( + Layer.mock(ThreadManagement.ThreadManagementService)({ + getThreadShell: (id) => + Effect.succeed( + id === threadId + ? shell(threadId, "codex", "gpt-5.5") + : shell(id, "claudeAgent", "claude-opus-5-5"), + ), + getThreadRecords: () => Effect.succeed({ runs: [], messages: [] } as never), + getProjectThreadRecords: () => + Effect.succeed({ thread: shell(threadId, "codex", "gpt-5.5") } as never), + dispatch: () => Effect.succeed({ sequence: 1 } as never), + }), + ), + Layer.provide( + Layer.mock(ProviderRegistry.ProviderRegistry)({ + getProviders: Effect.succeed([ + { + instanceId: ProviderInstanceId.make("codex"), + driver: "codex", + models: [{ slug: "gpt-5.5", isCustom: false }], + } as never, + { + instanceId: ProviderInstanceId.make("claudeAgent"), + driver: "claudeAgent", + models: [{ slug: "claude-opus-5-5", isCustom: false }], + } as never, + ]), + }), + ), + Layer.provide( + Layer.succeed( + AnalyticsService.AnalyticsService, + AnalyticsService.AnalyticsService.of({ + record: (event, properties) => + Effect.sync(() => void recorded.push({ event, properties })), + flush: Effect.void, + }), + ), + ), + ), + ), + ); +}); + it("keeps MCP preference output allowlisted and Unicode-bounded", () => { const settings = { ...DEFAULT_SERVER_SETTINGS, diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 96639a4d41aa..3412b025d08d 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -46,6 +46,7 @@ import * as SqlitePersistence from "./persistence/Layers/Sqlite.ts"; import * as PullRequestFilesViewed from "./persistence/PullRequestFilesViewed.ts"; import * as ServerLifecycleEvents from "./serverLifecycleEvents.ts"; import * as AnalyticsService from "./telemetry/AnalyticsService.ts"; +import * as DelegatedTaskAnalytics from "./telemetry/DelegatedTaskAnalytics.ts"; import * as ProviderEventIngestor from "./orchestration-v2/ProviderEventIngestor.ts"; import * as ModelManifest from "./provider/ModelManifest.ts"; import * as ResetCreditCoordinator from "./provider/Layers/resetCreditCoordinator.ts"; @@ -470,6 +471,10 @@ const ThreadSettlementWorkerLive = Layer.effectDiscard( ThreadSettlementService.make.pipe(Effect.flatMap((service) => service.start())), ).pipe(Layer.provide(PullRequestServiceLive), Layer.provide(ProjectionStoreV2.layer)); +const DelegatedTaskAnalyticsLive = Layer.effectDiscard( + DelegatedTaskAnalytics.make.pipe(Effect.flatMap((service) => service.start())), +); + const ThreadPullRequestWorkerLive = Layer.effectDiscard( ThreadPullRequestService.make.pipe(Effect.flatMap((service) => service.start())), ).pipe(Layer.provide(PullRequestServiceLive)); @@ -513,6 +518,7 @@ const RuntimeCoreDependenciesBaseLive = Layer.mergeAll( Layer.provide(ProjectionStoreV2.layer), ), ThreadPullRequestWorkerLive, + DelegatedTaskAnalyticsLive, Layer.effectDiscard( Effect.gen(function* () { const service = yield* PullRequestSyncReactor.PullRequestSyncReactor; diff --git a/apps/server/src/telemetry/DelegatedTaskAnalytics.test.ts b/apps/server/src/telemetry/DelegatedTaskAnalytics.test.ts new file mode 100644 index 000000000000..8c7a7f6f6108 --- /dev/null +++ b/apps/server/src/telemetry/DelegatedTaskAnalytics.test.ts @@ -0,0 +1,88 @@ +import { ProviderInstanceId, ThreadId } from "@t3tools/contracts"; +import { expect, it } from "@effect/vitest"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; + +import * as Orchestrator from "../orchestration-v2/Orchestrator.ts"; +import * as ProviderRegistry from "../provider/Services/ProviderRegistry.ts"; +import * as AnalyticsService from "./AnalyticsService.ts"; +import * as DelegatedTaskAnalytics from "./DelegatedTaskAnalytics.ts"; + +const parentThreadId = ThreadId.make("parent-private"); + +const subagentEvent = (status: string, origin = "app_owned") => + ({ + type: "subagent.updated", + payload: { + id: "task-private", + threadId: parentThreadId, + origin, + status, + providerInstanceId: ProviderInstanceId.make("claudeAgent"), + model: "claude-opus-5-5", + prompt: "secret task", + completionWake: "always", + startedAt: DateTime.makeUnsafe(0), + completedAt: DateTime.makeUnsafe(90_000), + }, + }) as never; + +it.effect("records each delegated task's final outcome once, without ids or prompts", () => { + const recorded: Array<{ event: string; properties: unknown }> = []; + return Effect.gen(function* () { + const analytics = yield* DelegatedTaskAnalytics.make; + yield* analytics.onSubagent(subagentEvent("running")); + yield* analytics.onSubagent(subagentEvent("completed", "provider_native")); + yield* analytics.onSubagent(subagentEvent("completed")); + yield* analytics.onSubagent(subagentEvent("completed")); + expect(recorded).toEqual([ + { + event: "mcp.delegated_task.finished", + properties: { + status: "completed", + callerProvider: "codex", + callerModel: "gpt-5.5", + targetProvider: "claudeAgent", + targetModel: "claude-opus-5-5", + crossProvider: true, + completionWake: "always", + durationSeconds: 90, + }, + }, + ]); + }).pipe( + Effect.provide( + Layer.mergeAll( + Layer.mock(Orchestrator.OrchestratorV2)({ + getThreadShell: () => + Effect.succeed({ + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5.5" }, + } as never), + }), + Layer.mock(ProviderRegistry.ProviderRegistry)({ + getProviders: Effect.succeed([ + { + instanceId: ProviderInstanceId.make("codex"), + driver: "codex", + models: [{ slug: "gpt-5.5", isCustom: false }], + } as never, + { + instanceId: ProviderInstanceId.make("claudeAgent"), + driver: "claudeAgent", + models: [{ slug: "claude-opus-5-5", isCustom: false }], + } as never, + ]), + }), + Layer.succeed( + AnalyticsService.AnalyticsService, + AnalyticsService.AnalyticsService.of({ + record: (event, properties) => + Effect.sync(() => void recorded.push({ event, properties })), + flush: Effect.void, + }), + ), + ), + ), + ); +}); diff --git a/apps/server/src/telemetry/DelegatedTaskAnalytics.ts b/apps/server/src/telemetry/DelegatedTaskAnalytics.ts new file mode 100644 index 000000000000..3ec884962de6 --- /dev/null +++ b/apps/server/src/telemetry/DelegatedTaskAnalytics.ts @@ -0,0 +1,84 @@ +/** + * Records one anonymous `mcp.delegated_task.finished` event when an app-owned + * delegated task reaches a final status, so we can see whether delegation + * between providers actually completes. The event carries the parent and child + * provider and model, the final status, and the child's runtime in seconds. + * + * @module DelegatedTaskAnalytics + */ +import type { OrchestrationV2DomainEvent } from "@t3tools/contracts"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Stream from "effect/Stream"; + +import * as Orchestrator from "../orchestration-v2/Orchestrator.ts"; +import * as ProviderRegistry from "../provider/Services/ProviderRegistry.ts"; +import { forkParked } from "../serverActivation.ts"; +import * as AnalyticsService from "./AnalyticsService.ts"; +import { providerDimensions } from "./ProviderDimensions.ts"; + +type SubagentEvent = Extract; + +const FINAL_STATUSES: ReadonlySet = new Set([ + "completed", + "failed", + "cancelled", + "interrupted", +]); + +export const make = Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const analytics = yield* AnalyticsService.AnalyticsService; + const registry = yield* ProviderRegistry.ProviderRegistry; + // A task can be re-reported with the same final status; count it once per process. + const reported = new Set(); + + const onSubagent = (event: SubagentEvent) => + Effect.gen(function* () { + const task = event.payload; + if (task.origin !== "app_owned" || !FINAL_STATUSES.has(task.status)) return; + if (reported.has(task.id)) return; + reported.add(task.id); + const parent = yield* orchestrator.getThreadShell(task.threadId); + const providers = yield* registry.getProviders; + const parentDimensions = + parent === null + ? { provider: "unknown" } + : providerDimensions(providers, parent.modelSelection); + const child = providerDimensions(providers, { + instanceId: task.providerInstanceId, + model: task.model ?? "", + }); + const durationSeconds = + task.startedAt !== null && task.completedAt !== null + ? Math.max( + 0, + Math.round( + (DateTime.toEpochMillis(task.completedAt) - + DateTime.toEpochMillis(task.startedAt)) / + 1000, + ), + ) + : undefined; + yield* analytics.record("mcp.delegated_task.finished", { + status: task.status, + callerProvider: parentDimensions.provider, + ...(parentDimensions.model === undefined ? {} : { callerModel: parentDimensions.model }), + targetProvider: child.provider, + ...(child.model === undefined ? {} : { targetModel: child.model }), + crossProvider: parentDimensions.provider !== child.provider, + ...(task.completionWake === undefined ? {} : { completionWake: task.completionWake }), + ...(durationSeconds === undefined ? {} : { durationSeconds }), + }); + }).pipe(Effect.ignoreCause); + + const start = Effect.fn("DelegatedTaskAnalytics.start")(function* () { + yield* forkParked( + Stream.runForEach(orchestrator.streamDomainEvents, (event) => + event.type === "subagent.updated" ? onSubagent(event) : Effect.void, + ), + ); + }); + + return { start, onSubagent }; +}); diff --git a/apps/server/src/telemetry/ProviderDimensions.test.ts b/apps/server/src/telemetry/ProviderDimensions.test.ts new file mode 100644 index 000000000000..7ea340d3e8dc --- /dev/null +++ b/apps/server/src/telemetry/ProviderDimensions.test.ts @@ -0,0 +1,204 @@ +import { ProviderInstanceId, ThreadId, type ServerProvider } from "@t3tools/contracts"; +import { expect, it } from "@effect/vitest"; + +import { + type AgentThread, + agentToolProperties, + handoffSettings, + callerOrigin, + toolOutcome, +} from "./ProviderDimensions.ts"; + +const provider = ( + instanceId: string, + driver: string, + models: ReadonlyArray<{ slug: string; isCustom: boolean }>, +) => + ({ + instanceId: ProviderInstanceId.make(instanceId), + driver, + models, + }) as unknown as ServerProvider; + +const providers = [ + provider("work-codex", "codex", [{ slug: "gpt-5.5", isCustom: false }]), + provider("claudeAgent", "claudeAgent", [ + { slug: "claude-opus-5-5", isCustom: false }, + { slug: "my-private-finetune", isCustom: true }, + ]), + provider("opencode", "opencode", [{ slug: "ollama/acme-internal", isCustom: false }]), +]; + +const thread = ( + instanceId: string, + model: string, + overrides: Partial = {}, +): AgentThread => ({ + modelSelection: { instanceId: ProviderInstanceId.make(instanceId), model }, + runtimeMode: "full-access", + interactionMode: "default", + createdBy: "user", + creationSource: "web", + lineage: { parentThreadId: null, relationshipToParent: null, rootThreadId: ThreadId.make("t") }, + ...overrides, +}); + +it("relates the calling agent to each agent that received work", () => { + expect( + agentToolProperties({ + tool: "create_threads", + providers, + caller: thread("work-codex", "gpt-5.5", { createdBy: "agent", creationSource: "mcp" }), + callerDepth: 1, + targets: [ + thread("claudeAgent", "claude-opus-5-5", { interactionMode: "plan" }), + thread("work-codex", "gpt-5.5"), + ], + outcome: { outcome: "ok" }, + settings: { batchSize: 2, targetChosen: true }, + durationMs: 42, + }), + ).toEqual([ + { + tool: "create_threads", + callerProvider: "codex", + callerModel: "gpt-5.5", + callerOrigin: "agent", + callerDepth: 1, + durationMs: 42, + outcome: "ok", + batchSize: 2, + targetChosen: true, + targetProvider: "claudeAgent", + targetModel: "claude-opus-5-5", + targetRuntimeMode: "full-access", + targetInteractionMode: "plan", + crossProvider: true, + }, + { + tool: "create_threads", + callerProvider: "codex", + callerModel: "gpt-5.5", + callerOrigin: "agent", + callerDepth: 1, + durationMs: 42, + outcome: "ok", + batchSize: 2, + targetChosen: true, + targetProvider: "codex", + targetModel: "gpt-5.5", + targetRuntimeMode: "full-access", + targetInteractionMode: "default", + crossProvider: false, + }, + ]); +}); + +it("reports a non-handoff tool from the credential's provider alone", () => { + expect( + agentToolProperties({ + tool: "preview_click", + providers, + caller: undefined, + callerProviderInstanceId: ProviderInstanceId.make("claudeAgent"), + targets: [], + outcome: { outcome: "error", errorCode: "capability_denied" }, + durationMs: 3, + }), + ).toEqual([ + { + tool: "preview_click", + callerProvider: "claudeAgent", + durationMs: 3, + outcome: "error", + errorCode: "capability_denied", + }, + ]); +}); + +it("omits custom, user-configured, and unknown models and instance names", () => { + expect( + agentToolProperties({ + tool: "t3_thread_send", + providers, + caller: thread("claudeAgent", "my-private-finetune"), + targets: [thread("opencode", "ollama/acme-internal"), thread("removed-instance", "gpt-5.5")], + outcome: { outcome: "ok" }, + }).map(({ callerProvider, callerModel, targetProvider, targetModel }) => ({ + callerProvider, + callerModel, + targetProvider, + targetModel, + })), + ).toEqual([ + { + callerProvider: "claudeAgent", + callerModel: undefined, + targetProvider: "opencode", + targetModel: undefined, + }, + { + callerProvider: "claudeAgent", + callerModel: undefined, + targetProvider: "unknown", + targetModel: undefined, + }, + ]); +}); + +it("reads handoff settings from enums and flags, never from prompts or ids", () => { + expect( + handoffSettings( + "delegate_task", + { task: "secret", mode: "wait", target: { driverKind: "claudeAgent" } }, + { taskId: "t", waitTimedOut: true }, + ), + ).toEqual({ targetChosen: true, mode: "wait", waitTimedOut: true }); + expect(handoffSettings("delegate_task", { task: "secret" }, { waitTimedOut: false })).toEqual({ + targetChosen: false, + mode: "async", + }); + expect( + handoffSettings( + "t3_thread_launch", + { message: "secret", workspaceStrategy: { type: "worktree", branch: "private" } }, + {}, + ), + ).toEqual({ targetChosen: false, workspace: "worktree" }); + expect(handoffSettings("t3_thread_launch", { scratch: true }, {})).toMatchObject({ + workspace: "scratch", + }); + expect( + handoffSettings( + "create_threads", + { threads: [{ prompt: "a" }, { prompt: "b", target: {} }] }, + {}, + ), + ).toEqual({ batchSize: 2, targetChosen: true }); + expect(handoffSettings("t3_thread_send", { message: "secret" }, { delivery: "steered" })).toEqual( + { delivery: "steered" }, + ); +}); + +it("reads failure codes from declared tool errors even without isError", () => { + expect( + toolOutcome({ + structuredContent: { + _tag: "OrchestratorMcpFailure", + code: "capability_denied", + message: "private detail", + }, + }), + ).toEqual({ outcome: "error", errorCode: "capability_denied" }); + expect(toolOutcome({ isError: false, structuredContent: { threadId: "t" } })).toEqual({ + outcome: "ok", + }); + expect(toolOutcome(undefined)).toEqual({ outcome: "error", errorCode: "exception" }); +}); + +it("tells user, agent, system, and scheduled work apart", () => { + expect(callerOrigin({ thread: { createdBy: "user" }, scheduledRun: false })).toBe("user"); + expect(callerOrigin({ thread: { createdBy: "agent" }, scheduledRun: false })).toBe("agent"); + expect(callerOrigin({ thread: { createdBy: "system" }, scheduledRun: false })).toBe("system"); + expect(callerOrigin({ thread: { createdBy: "user" }, scheduledRun: true })).toBe("scheduler"); +}); diff --git a/apps/server/src/telemetry/ProviderDimensions.ts b/apps/server/src/telemetry/ProviderDimensions.ts new file mode 100644 index 000000000000..71f7a9a7118f --- /dev/null +++ b/apps/server/src/telemetry/ProviderDimensions.ts @@ -0,0 +1,194 @@ +import type { + ModelSelection, + OrchestrationV2ThreadShell, + ServerProvider, +} from "@t3tools/contracts"; + +/** + * Drivers whose model catalogs are published by the vendor. Other drivers + * (OpenCode, Pi, ACP agents) list models from the user's own configuration, + * such as local Ollama tags, so their model names stay out of analytics. + */ +const VENDOR_CATALOG_DRIVERS: ReadonlySet = new Set([ + "codex", + "claudeAgent", + "cursor", + "grok", + "antigravity", +]); + +export interface ProviderDimensions { + readonly provider: string; + readonly model?: string; +} + +/** + * Anonymous provider and model for a model selection. Instance ids are + * user-defined, so only the driver kind is reported. The model is reported + * only when it comes from a vendor catalog and is not a user-added custom model. + */ +export function providerDimensions( + providers: ReadonlyArray, + selection: Pick, +): ProviderDimensions { + const provider = providers.find((candidate) => candidate.instanceId === selection.instanceId); + if (provider === undefined) return { provider: "unknown" }; + const model = provider.models.find((candidate) => candidate.slug === selection.model); + return VENDOR_CATALOG_DRIVERS.has(provider.driver) && model !== undefined && !model.isCustom + ? { provider: provider.driver, model: model.slug } + : { provider: provider.driver }; +} + +/** The thread fields agent analytics reads. Ids are used for lookups, never reported. */ +export type AgentThread = Pick< + OrchestrationV2ThreadShell, + "modelSelection" | "runtimeMode" | "interactionMode" | "createdBy" | "creationSource" | "lineage" +>; + +type Fields = Readonly>; + +const asFields = (value: unknown): Fields => + typeof value === "object" && value !== null ? (value as Fields) : {}; + +const stringField = (fields: Fields, key: string) => + typeof fields[key] === "string" ? (fields[key] as string) : undefined; + +const defined = (entries: Readonly>) => + Object.fromEntries(Object.entries(entries).filter(([, value]) => value !== undefined)); + +/** + * How deep a thread sits in a delegation tree: 0 for a top-level thread, 1 for + * its subagent, and so on. Forks are top-level threads of their own. Depth is + * counted by walking parents, capped so a corrupt lineage cannot loop. + */ +export const MAX_REPORTED_DEPTH = 5; + +/** + * Who started the work that made this call. Scheduled runs keep the task + * creator's provenance, so the run's message marks them instead. + */ +export function callerOrigin(input: { + readonly thread: Pick; + readonly scheduledRun: boolean; +}) { + if (input.scheduledRun) return "scheduler"; + if (input.thread.createdBy === "agent") return "agent"; + return input.thread.createdBy === "system" ? "system" : "user"; +} + +/** + * Settings the agent chose for a handoff, read only from closed enums and + * booleans in the tool's arguments and result. Prompts, titles, and ids are + * never read. + */ +export function handoffSettings(tool: string, args: unknown, result: unknown): Fields { + const input = asFields(args); + const output = asFields(result); + switch (tool) { + case "delegate_task": + return defined({ + targetChosen: input.target !== undefined, + mode: stringField(input, "mode") ?? "async", + ...(output.waitTimedOut === true ? { waitTimedOut: true } : {}), + }); + case "create_threads": { + const threads = Array.isArray(input.threads) ? input.threads : []; + return { + batchSize: threads.length, + targetChosen: threads.some((thread) => asFields(thread).target !== undefined), + }; + } + case "t3_thread_launch": + return defined({ + targetChosen: input.modelSelection !== undefined, + workspace: + input.scratch === true + ? "scratch" + : (stringField(asFields(input.workspaceStrategy), "type") ?? "root"), + }); + case "t3_thread_send": + return defined({ delivery: stringField(output, "delivery") }); + case "t3_thread_send_attachments": + return { delivery: "auto" }; + case "schedule_task": + return { bindToCurrentThread: input.bindToCurrentThread !== false }; + default: + return {}; + } +} + +/** Outcome of a tool call, with the stable failure code when the call failed. */ +export function toolOutcome( + result: + | { readonly isError?: boolean | undefined; readonly structuredContent?: unknown } + | undefined, +) { + if (result === undefined) return { outcome: "error", errorCode: "exception" }; + const content = asFields(result.structuredContent); + // T3 tool failures are declared errors whose structured content carries a + // closed `code`; the MCP layer does not always set isError for them. + const code = stringField(content, "code") ?? stringField(asFields(content.error), "_tag"); + if (result.isError === true || (content._tag !== undefined && code !== undefined)) { + return { outcome: "error", errorCode: code ?? "unknown" }; + } + return { outcome: "ok" }; +} + +/** + * Event properties for an agent tool call: the calling agent, and when the call + * handed work to other threads, the agents running them. Each target counts + * once, so a create_threads batch on two providers reports both. + */ +export function agentToolProperties(input: { + readonly tool: string; + readonly providers: ReadonlyArray; + readonly caller: AgentThread | undefined; + readonly callerProviderInstanceId?: ModelSelection["instanceId"]; + readonly callerDepth?: number; + readonly callerScheduledRun?: boolean; + readonly targets: ReadonlyArray; + readonly outcome: Fields; + readonly settings?: Fields; + readonly durationMs?: number; +}): ReadonlyArray>> { + const caller = + input.caller !== undefined + ? providerDimensions(input.providers, input.caller.modelSelection) + : input.callerProviderInstanceId !== undefined + ? { + provider: providerDimensions(input.providers, { + instanceId: input.callerProviderInstanceId, + model: "", + }).provider, + } + : { provider: "unknown" }; + const base = defined({ + tool: input.tool, + callerProvider: caller.provider, + callerModel: caller.model, + callerOrigin: + input.caller === undefined + ? undefined + : callerOrigin({ + thread: input.caller, + scheduledRun: input.callerScheduledRun === true, + }), + callerDepth: + input.callerDepth === undefined ? undefined : Math.min(input.callerDepth, MAX_REPORTED_DEPTH), + durationMs: input.durationMs, + ...input.outcome, + ...input.settings, + }); + if (input.targets.length === 0) return [base]; + return input.targets.map((thread) => { + const target = providerDimensions(input.providers, thread.modelSelection); + return defined({ + ...base, + targetProvider: target.provider, + targetModel: target.model, + targetRuntimeMode: thread.runtimeMode, + targetInteractionMode: thread.interactionMode, + crossProvider: caller.provider !== target.provider, + }); + }); +} diff --git a/docs/internals/product-analytics.md b/docs/internals/product-analytics.md index e9123640a541..eb476bfbfdc3 100644 --- a/docs/internals/product-analytics.md +++ b/docs/internals/product-analytics.md @@ -41,6 +41,16 @@ Unknown counts stay absent. Partial usage contains valid observed counts but cannot establish a whole-turn total. Keep these distinctions when changing token normalization or building reports. +Agent tool events report driver kinds, never instance ids: users name +instances, so an instance id can identify an account or employer. +[Provider dimensions](../../apps/server/src/telemetry/ProviderDimensions.ts) +report a model only when it is a non-custom entry in a vendor catalog. OpenCode, +Pi, and ACP agents list models from user configuration, such as local Ollama +tags, so their models stay out. Handoff settings come from closed enums and +booleans in tool arguments; prompts, titles, and ids are never read. Thread ids +are used only to look up the threads on each side and never leave the server, +which is also why tool events carry no per-thread sequence. + ## Delivery A send can fail after PostHog has stored the batch, so every retry is a copy.