diff --git a/packages/shared/src/telemetry/client.ts b/packages/shared/src/telemetry/client.ts index 4a131d2029..03a10a5550 100644 --- a/packages/shared/src/telemetry/client.ts +++ b/packages/shared/src/telemetry/client.ts @@ -116,11 +116,21 @@ export class TelemetryClient { * backend event schema. */ track(eventName: K, ...args: TrackArgs): void { - if (!Object.hasOwn(PAPERCLIP_EVENTS, eventName)) return; + if (!this.isRegisteredEventName(eventName)) return; const [dimensions] = args; this.enqueue(eventName, dimensions); } + /** + * Whether `track` would enqueue this first-party event name rather than drop + * it as unregistered. Lets a caller with an expensive payload (e.g. one that + * needs database reads) skip building dimensions for a proposed event the + * client would discard anyway. + */ + isRegisteredEventName(eventName: string): boolean { + return Object.hasOwn(PAPERCLIP_EVENTS, eventName); + } + /** * Tracks plugin telemetry bridge events whose names are built dynamically * from third-party plugin input. The backend accepts only explicitly diff --git a/packages/shared/src/telemetry/events.ts b/packages/shared/src/telemetry/events.ts index 5b38b75d78..71f78bb2f3 100644 --- a/packages/shared/src/telemetry/events.ts +++ b/packages/shared/src/telemetry/events.ts @@ -153,6 +153,75 @@ export function trackAgentTaskRun( }); } +export function trackConnectionCreated( + client: TelemetryClient, + dims: { + connector_key: RawDimension<"custom">; + transport: RawDimension<"mcp_remote" | "rest_api" | "local_stdio">; + auth_kind: RawDimension<"oauth" | "api_key" | "none">; + setup_flow: RawDimension<"gallery" | "api" | "example">; + status: RawDimension<"draft" | "active" | "disabled" | "archived">; + enabled: boolean; + }, +): void { + client.track( + // @ts-expect-error -- proposed-telemetry(https://github.com/paperclipai/paperclip/issues/13578): measure which catalog connectors installations create connections for + "connection.created", + dims, + ); +} + +export function trackConnectionUpdated( + client: TelemetryClient, + dims: { + connector_key: RawDimension<"custom">; + transport: RawDimension<"mcp_remote" | "rest_api" | "local_stdio">; + auth_kind: RawDimension<"oauth" | "api_key" | "none">; + change_source: RawDimension< + | "api" + | "gallery" + | "oauth_callback" + | "credential_refresh" + | "archive" + | "example" + >; + previous_status: RawDimension<"draft" | "active" | "disabled" | "archived">; + status: RawDimension<"draft" | "active" | "disabled" | "archived">; + previous_enabled: boolean; + enabled: boolean; + }, +): void { + client.track( + // @ts-expect-error -- proposed-telemetry(https://github.com/paperclipai/paperclip/issues/13578): measure connector lifecycle transitions (configured, paused, archived) after creation + "connection.updated", + dims, + ); +} + +export function trackConnectionInvoked( + client: TelemetryClient, + dims: { + connector_key: RawDimension<"custom">; + transport: RawDimension<"mcp_remote" | "rest_api" | "local_stdio">; + status: RawDimension< + | "succeeded" + | "failed" + | "denied" + | "cancelled" + | "timed_out" + | "rate_limited" + >; + origin: RawDimension<"setup_test" | "agent" | "user" | "system" | "plugin">; + duration_seconds?: number; + }, +): void { + client.track( + // @ts-expect-error -- proposed-telemetry(https://github.com/paperclipai/paperclip/issues/13578): measure whether connected connectors are successfully used and where invocations fail; records completed invocation attempts and their terminal status, never invocation starts + "connection.invoked", + dims, + ); +} + export function trackErrorHandlerCrash( client: TelemetryClient, dims: { errorCode: string }, diff --git a/packages/shared/src/telemetry/index.ts b/packages/shared/src/telemetry/index.ts index 7628cac592..61f945ad83 100644 --- a/packages/shared/src/telemetry/index.ts +++ b/packages/shared/src/telemetry/index.ts @@ -18,6 +18,9 @@ export { trackErrorHandlerCrash, trackInteractionCreated, trackInteractionResolved, + trackConnectionCreated, + trackConnectionUpdated, + trackConnectionInvoked, } from "./events.js"; export type { TelemetryConfig, diff --git a/packages/shared/src/telemetry/proposed-connector-events.test.ts b/packages/shared/src/telemetry/proposed-connector-events.test.ts new file mode 100644 index 0000000000..a8a6826e9d --- /dev/null +++ b/packages/shared/src/telemetry/proposed-connector-events.test.ts @@ -0,0 +1,124 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import { TelemetryClient } from "./client.js"; +import { resolveTelemetryConfig } from "./config.js"; +import { + trackConnectionCreated, + trackConnectionUpdated, + trackConnectionInvoked, +} from "./events.js"; +import type { TelemetryState } from "./types.js"; + +const TEST_STATE: TelemetryState = { + installId: "test-install", + salt: "test-salt", + createdAt: "2026-01-01T00:00:00Z", + firstSeenVersion: "0.0.0", +}; + +function makeClient(config?: { enabled?: boolean }) { + const stateFactory = vi.fn(() => TEST_STATE); + return { + client: new TelemetryClient( + { enabled: config?.enabled ?? true, endpoint: "http://localhost:9999/ingest" }, + stateFactory, + "0.0.0-test", + () => 0.5, + ), + stateFactory, + }; +} + +function trackAllProposedConnectorEvents(client: TelemetryClient) { + trackConnectionCreated(client, { + connector_key: "github", + transport: "mcp_remote", + auth_kind: "oauth", + setup_flow: "gallery", + status: "active", + enabled: true, + }); + trackConnectionUpdated(client, { + connector_key: "github", + transport: "mcp_remote", + auth_kind: "oauth", + change_source: "api", + previous_status: "draft", + status: "active", + previous_enabled: false, + enabled: true, + }); + trackConnectionInvoked(client, { + connector_key: "github", + transport: "mcp_remote", + status: "succeeded", + origin: "agent", + duration_seconds: 3, + }); +} + +describe("proposed connector events against the real TelemetryClient", () => { + afterEach(() => { + vi.restoreAllMocks(); + vi.unstubAllEnvs(); + }); + + it("cannot queue or send: state stays untouched and nothing hits the wire", async () => { + vi.stubGlobal("fetch", vi.fn().mockResolvedValue({ ok: true })); + const { client, stateFactory } = makeClient(); + + trackAllProposedConnectorEvents(client); + await client.flush(); + + expect(stateFactory).not.toHaveBeenCalled(); + expect(fetch).not.toHaveBeenCalled(); + }); + + it("differential control: a registered event through the same client does send", async () => { + vi.stubGlobal("fetch", vi.fn().mockResolvedValue({ ok: true })); + const { client } = makeClient(); + + trackAllProposedConnectorEvents(client); + client.track("project.created", {}); + await client.flush(); + + expect(fetch).toHaveBeenCalledTimes(1); + const body = JSON.parse( + String((vi.mocked(fetch).mock.calls[0]?.[1] as RequestInit).body), + ); + expect(body.events).toEqual([ + expect.objectContaining({ name: "project.created" }), + ]); + }); + + it("reports the three connector event names as unregistered", () => { + const { client } = makeClient(); + expect(client.isRegisteredEventName("connection.created")).toBe(false); + expect(client.isRegisteredEventName("connection.updated")).toBe(false); + expect(client.isRegisteredEventName("connection.invoked")).toBe(false); + expect(client.isRegisteredEventName("project.created")).toBe(true); + }); + + it("a disabled client drops even registered events", async () => { + vi.stubGlobal("fetch", vi.fn().mockResolvedValue({ ok: true })); + const { client, stateFactory } = makeClient({ enabled: false }); + + client.track("project.created", {}); + await client.flush(); + + expect(stateFactory).not.toHaveBeenCalled(); + expect(fetch).not.toHaveBeenCalled(); + }); + + it("environment suppression resolves telemetry off under CI and opt-out flags", () => { + vi.stubEnv("PAPERCLIP_TELEMETRY_DISABLED", "1"); + expect(resolveTelemetryConfig().enabled).toBe(false); + vi.unstubAllEnvs(); + + vi.stubEnv("DO_NOT_TRACK", "1"); + expect(resolveTelemetryConfig().enabled).toBe(false); + vi.unstubAllEnvs(); + + vi.stubEnv("GITHUB_ACTIONS", "true"); + expect(resolveTelemetryConfig().enabled).toBe(false); + }); +}); diff --git a/server/src/__tests__/tool-access-connector-telemetry.test.ts b/server/src/__tests__/tool-access-connector-telemetry.test.ts new file mode 100644 index 0000000000..fd41786598 --- /dev/null +++ b/server/src/__tests__/tool-access-connector-telemetry.test.ts @@ -0,0 +1,493 @@ +import { randomUUID } from "node:crypto"; +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { eq } from "drizzle-orm"; +import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "vitest"; +import { + activityLog, + companies, + companyMemberships, + companySecretBindings, + companySecretVersions, + companySecrets, + connectionGrants, + createDb, + secretAccessEvents, + toolApplications, + toolCatalogEntries, + toolConnectionInstalls, + toolConnections, + toolProfileBindings, + toolProfileEntries, + toolProfiles, +} from "@paperclipai/db"; +import { + getEmbeddedPostgresTestSupport, + startEmbeddedPostgresTestDatabase, +} from "./helpers/embedded-postgres.js"; + +const track = vi.fn(); +const isRegisteredEventName = vi.fn(() => true); +vi.mock("../telemetry.js", () => ({ + getTelemetryClient: () => ({ track, isRegisteredEventName }), +})); + +const { toolAccessService } = await import("../services/tool-access.js"); +const { secretService } = await import("../services/secrets.js"); +const { instanceSettingsService } = await import("../services/instance-settings.js"); + +const embeddedPostgresSupport = await getEmbeddedPostgresTestSupport(); +const describeEmbeddedPostgres = embeddedPostgresSupport.supported + ? describe + : describe.skip; + +type Db = ReturnType; + +function createdEvents() { + return track.mock.calls.filter(([name]) => name === "connection.created"); +} + +function updatedEvents() { + return track.mock.calls.filter(([name]) => name === "connection.updated"); +} + +async function createCompany(db: Db) { + const company = await db + .insert(companies) + .values({ + name: `Lifecycle telemetry ${randomUUID()}`, + issuePrefix: `LT${randomUUID().slice(0, 6).toUpperCase()}`, + }) + .returning() + .then((rows) => rows[0]!); + await db.insert(companyMemberships).values({ + companyId: company.id, + principalType: "user", + principalId: actor.actorId, + status: "active", + membershipRole: "admin", + }); + return company; +} + +const actor = { + actorType: "user" as const, + actorId: "lifecycle-telemetry-user", + actorSource: "local_implicit" as const, +}; + +const mcpTool = (name: string) => ({ + name, + description: name, + inputSchema: { type: "object", properties: {} }, + annotations: { readOnlyHint: true }, +}); + +/** + * Minimal remote MCP server: answers initialize, the initialized notification, + * and tools/list. Mirrors the fixture the remote-connector lifecycle tests use + * so the direct Composio MCP path here is the real gallery path, not a stub. + */ +function remoteMcpFixture() { + const requests: { url: string; headers: Headers }[] = []; + const remoteHttpRequest = async (url: string, init: RequestInit) => { + const body = JSON.parse(String(init.body)) as { id: number; method: string }; + requests.push({ url, headers: new Headers(init.headers) }); + if (body.method === "initialize") { + return Response.json( + { + id: body.id, + jsonrpc: "2.0", + result: { protocolVersion: "2025-06-18", capabilities: {} }, + }, + { headers: { "Mcp-Session-Id": "session" } }, + ); + } + if (body.method === "notifications/initialized") + return new Response(null, { status: 202 }); + return Response.json({ + id: body.id, + jsonrpc: "2.0", + result: { tools: [mcpTool("read"), mcpTool("ask")] }, + }); + }; + return { requests, remoteHttpRequest }; +} + +/** + * Retained legacy broker records (#13758): a `rest_api` parent keyed to the + * catalog slug and a `provider: composio` child. Master keeps such rows only + * for inspection and explicit removal; they cannot execute or reconnect. + */ +async function insertRetiredComposioRecords(db: Db, companyId: string) { + const [application] = await db + .insert(toolApplications) + .values({ companyId, name: "Composio", type: "rest_api", status: "active" }) + .returning(); + const [parent] = await db + .insert(toolConnections) + .values({ + companyId, + applicationId: application!.id, + name: "Composio (legacy broker)", + uid: `composio/${randomUUID()}`, + transport: "rest_api", + authKind: "api_key", + status: "active", + enabled: true, + config: { sourceTemplateKey: "composio" }, + transportConfig: { sourceTemplateKey: "composio" }, + }) + .returning(); + const [child] = await db + .insert(toolConnections) + .values({ + companyId, + applicationId: application!.id, + name: "GitHub (via Composio)", + uid: `composio/github/${randomUUID()}`, + transport: "mcp_remote", + authKind: "none", + status: "active", + enabled: true, + config: { + provider: "composio", + parentConnectionId: parent!.id, + toolkitSlug: "github", + connectedAccountId: "account-github", + }, + transportConfig: {}, + }) + .returning(); + return { parent: parent!, child: child! }; +} + +describeEmbeddedPostgres("connector lifecycle telemetry (tool-access)", () => { + let db!: Db; + let tempDb: Awaited> | null = + null; + let keyDir: string | null = null; + + beforeAll(async () => { + keyDir = await mkdtemp(join(tmpdir(), "lifecycle-telemetry-secrets-")); + vi.stubEnv("PAPERCLIP_SECRETS_MASTER_KEY_FILE", join(keyDir, "key")); + tempDb = await startEmbeddedPostgresTestDatabase( + "paperclip-lifecycle-telemetry-", + ); + db = createDb(tempDb.connectionString); + await instanceSettingsService(db).updateExperimental({ + enableMcpAggregators: true, + }); + }, 30_000); + + afterEach(async () => { + track.mockClear(); + await db.delete(activityLog); + await db.delete(toolProfileEntries); + await db.delete(toolProfileBindings); + await db.delete(toolProfiles); + await db.delete(toolCatalogEntries); + await db.delete(connectionGrants); + await db.delete(toolConnectionInstalls); + await db.delete(toolConnections); + await db.delete(toolApplications); + await db.delete(secretAccessEvents); + await db.delete(companySecretBindings); + await db.delete(companySecretVersions); + await db.delete(companySecrets); + await db.delete(companyMemberships); + await db.delete(companies); + }); + + afterAll(async () => { + await tempDb?.cleanup(); + vi.unstubAllEnvs(); + if (keyDir) await rm(keyDir, { recursive: true, force: true }); + }); + + function service(options: Parameters[1] = {}) { + return toolAccessService(db, { + remoteHttpEndpointLookup: async () => [{ address: "8.8.8.8", family: 4 }], + ...options, + }); + } + + it("createConnection emits one created event with the catalog key from the committed row", async () => { + const company = await createCompany(db); + const svc = service(); + await svc.createConnection(company.id, { + name: "GitHub fixture", + transport: "mcp_remote", + config: { url: "https://fixture.example/mcp", sourceTemplateKey: "github" }, + enabled: true, + status: "active", + }); + const events = createdEvents(); + expect(events).toHaveLength(1); + expect(events[0]?.[1]).toMatchObject({ + connector_key: "github", + transport: "mcp_remote", + setup_flow: "api", + status: "active", + enabled: true, + }); + }, 30_000); + + it("updateConnection emits committed transitions only; metadata saves stay silent", async () => { + const company = await createCompany(db); + const svc = service(); + const connection = await svc.createConnection(company.id, { + name: "Custom fixture", + transport: "mcp_remote", + config: { url: "https://fixture.example/mcp" }, + enabled: true, + status: "active", + }); + track.mockClear(); + + await svc.updateConnection(connection.id, { enabled: false }); + let events = updatedEvents(); + expect(events).toHaveLength(1); + expect(events[0]?.[1]).toMatchObject({ + connector_key: "custom", + change_source: "api", + previous_enabled: true, + enabled: false, + previous_status: "active", + status: "active", + }); + + track.mockClear(); + await svc.updateConnection(connection.id, { name: "Renamed fixture" }); + expect(updatedEvents()).toHaveLength(0); + }, 30_000); + + it("direct Composio MCP setup from the gallery emits catalog identity only, never the session URL", async () => { + const company = await createCompany(db); + const remote = remoteMcpFixture(); + const svc = service({ + deploymentMode: "local_trusted", + deploymentExposure: "private", + remoteHttpRequest: remote.remoteHttpRequest, + }); + // The same server path the Apps gallery and inline task cards call + // (POST /companies/:companyId/tools/apps/connect): Composio Connect with an + // externally configured MCP session URL. + const connected = await svc.connectGalleryApp( + company.id, + { + galleryKey: "composio", + connectionMethodKey: "mcp", + link: "https://mcp.composio.dev/session/fixture?token=fixture-secret", + authMode: "none", + }, + actor, + ); + expect(remote.requests.length).toBeGreaterThan(0); + const created = createdEvents(); + expect(created).toHaveLength(1); + expect(created[0]?.[1]).toMatchObject({ + connector_key: "composio", + transport: "mcp_remote", + auth_kind: "none", + setup_flow: "gallery", + status: connected.connection.status, + enabled: connected.connection.enabled, + }); + expect(JSON.stringify(track.mock.calls)).not.toContain("fixture-secret"); + expect(JSON.stringify(track.mock.calls)).not.toContain("mcp.composio.dev"); + + track.mockClear(); + const finished = await svc.finishGalleryAppConnection( + company.id, + connected.connectionId, + { + enabledCatalogEntryIds: connected.catalog.map((entry) => entry.id), + askFirstCatalogEntryIds: [], + access: "all_agents", + }, + actor, + ); + // Connect commits a draft; the finish step is the committed draft -> active + // transition, so it is the one `gallery` update this setup reports. + expect(connected.connection).toMatchObject({ status: "draft", enabled: false }); + expect(finished.connection).toMatchObject({ status: "active", enabled: true }); + const updated = updatedEvents(); + expect(updated).toHaveLength(1); + expect(updated[0]?.[1]).toMatchObject({ + connector_key: "composio", + transport: "mcp_remote", + change_source: "gallery", + previous_status: "draft", + status: "active", + previous_enabled: false, + enabled: true, + }); + expect(createdEvents()).toHaveLength(0); + }, 30_000); + + it("pausing and resuming a direct Composio MCP connection reports api transitions with the catalog key", async () => { + const company = await createCompany(db); + const remote = remoteMcpFixture(); + const svc = service({ + deploymentMode: "local_trusted", + deploymentExposure: "private", + remoteHttpRequest: remote.remoteHttpRequest, + }); + const connected = await svc.connectGalleryApp( + company.id, + { + galleryKey: "composio", + connectionMethodKey: "mcp", + link: "https://mcp.composio.dev/session/fixture?token=fixture-secret", + authMode: "none", + }, + actor, + ); + await svc.finishGalleryAppConnection( + company.id, + connected.connectionId, + { + enabledCatalogEntryIds: connected.catalog.map((entry) => entry.id), + askFirstCatalogEntryIds: [], + access: "all_agents", + }, + actor, + ); + track.mockClear(); + + await svc.updateConnection(connected.connectionId, { enabled: false }); + await svc.updateConnection(connected.connectionId, { enabled: true }); + const updated = updatedEvents(); + expect(updated).toHaveLength(2); + expect(updated.map(([, dims]) => dims)).toEqual([ + expect.objectContaining({ + connector_key: "composio", + transport: "mcp_remote", + change_source: "api", + previous_enabled: true, + enabled: false, + }), + expect.objectContaining({ + connector_key: "composio", + change_source: "api", + previous_enabled: false, + enabled: true, + }), + ]); + // No per-app child rows exist for direct MCP, so nothing else cascades. + expect(createdEvents()).toHaveLength(0); + }, 30_000); + + it("finalizing OAuth access reports the committed draft -> active transition as a gallery step", async () => { + // Direct MCP aggregator OAuth (Composio Connect via DCR) from the Apps + // gallery or an inline task card: the callback stores a personal grant, + // then finalizeOAuthAccess — POST .../apps/:id/finalize-oauth-access, also + // called by connection-intent completion — activates the draft row itself + // before handing off to the finish step. + const company = await createCompany(db); + const svc = service(); + const [application] = await db + .insert(toolApplications) + .values({ companyId: company.id, name: "Composio", type: "mcp_remote", status: "active" }) + .returning(); + const [connection] = await db + .insert(toolConnections) + .values({ + companyId: company.id, + applicationId: application!.id, + name: "Composio", + uid: `composio/${randomUUID()}`, + transport: "mcp_remote", + authKind: "oauth", + credentialPolicy: "per_user", + status: "draft", + enabled: false, + config: { + sourceTemplateKey: "composio", + connectionMethodKey: "mcp", + url: "https://connect.composio.dev/mcp", + }, + transportConfig: { url: "https://connect.composio.dev/mcp" }, + }) + .returning(); + const accessToken = await secretService(db).create(company.id, { + name: `OAuth access token ${randomUUID().slice(0, 8)}`, + key: `tool_app.${randomUUID()}.oauth_access_token`, + provider: "local_encrypted", + value: "fixture-access-token", + }); + await db.insert(connectionGrants).values({ + companyId: company.id, + connectionId: connection!.id, + kind: "user", + subjectUserId: actor.actorId, + status: "active", + credentialSecretRefs: [ + { + secretId: accessToken.id, + versionSelector: "latest", + configPath: "oauth.access_token", + required: true, + label: "OAuth access token", + }, + ], + }); + track.mockClear(); + + const finished = await svc.finalizeOAuthAccess( + company.id, + connection!.id, + { grantKind: "user" }, + actor, + ); + expect(finished.connection).toMatchObject({ status: "active", enabled: true }); + const updated = updatedEvents(); + // Exactly one transition: finalize's own write. The finish step it calls + // afterwards re-reads an already-active row and stays silent. + expect(updated).toHaveLength(1); + expect(updated[0]?.[1]).toMatchObject({ + connector_key: "composio", + transport: "mcp_remote", + auth_kind: "oauth", + change_source: "gallery", + previous_status: "draft", + status: "active", + previous_enabled: false, + enabled: true, + }); + expect(createdEvents()).toHaveLength(0); + expect(JSON.stringify(track.mock.calls)).not.toContain("fixture-access-token"); + }, 30_000); + + it("retained legacy broker records emit only their explicit removal, distinguishable by transport", async () => { + const company = await createCompany(db); + const { parent, child } = await insertRetiredComposioRecords(db, company.id); + const svc = service(); + track.mockClear(); + + await svc.archiveConnection(parent.id); + await svc.archiveConnection(child.id); + const updated = updatedEvents(); + expect(updated).toHaveLength(2); + expect(updated[0]?.[1]).toMatchObject({ + // Legacy parent: catalog slug survives, but `rest_api` marks it as the + // retired broker rather than a direct MCP connection. + connector_key: "composio", + transport: "rest_api", + change_source: "archive", + status: "archived", + }); + expect(updated[1]?.[1]).toMatchObject({ + // Legacy child: no catalog key; the raw toolkit slug never leaves. + connector_key: "custom", + transport: "mcp_remote", + change_source: "archive", + status: "archived", + }); + expect(JSON.stringify(track.mock.calls)).not.toContain("github"); + expect(JSON.stringify(track.mock.calls)).not.toContain("account-github"); + expect(createdEvents()).toHaveLength(0); + }, 30_000); +}); diff --git a/server/src/__tests__/tool-gateway-connector-telemetry.test.ts b/server/src/__tests__/tool-gateway-connector-telemetry.test.ts new file mode 100644 index 0000000000..1ead5d320e --- /dev/null +++ b/server/src/__tests__/tool-gateway-connector-telemetry.test.ts @@ -0,0 +1,571 @@ +import { randomUUID } from "node:crypto"; +import { eq } from "drizzle-orm"; +import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "vitest"; +import { + activityLog, + agents, + companies, + companyMemberships, + connectionGrants, + createDb, + heartbeatRuns, + issues, + projects, + toolAccessAuditEvents, + toolActionRequests, + toolApplications, + toolCallEvents, + toolCatalogEntries, + toolConnections, + toolGatewayRateLimitCounters, + toolGatewaySessions, + toolInvocations, + toolPolicies, + toolProfileBindings, + toolProfiles, + toolRuntimeSlots, +} from "@paperclipai/db"; +import { + getEmbeddedPostgresTestSupport, + startEmbeddedPostgresTestDatabase, +} from "./helpers/embedded-postgres.js"; + +// Proposal wrappers drop unregistered names inside the real client, so the +// integration assertions use a client double that reports the proposed event +// as registered. `track` observes exactly what the emitter would hand the +// real client after schema adoption. +const track = vi.fn(); +const isRegisteredEventName = vi.fn(() => true); +vi.mock("../telemetry.js", () => ({ + getTelemetryClient: () => ({ track, isRegisteredEventName }), +})); + +// Deterministic reproduction of "post-execution bookkeeping failure": the +// gateway's success path awaits writeAudit -> logActivity after the succeeded +// save; failing that call lands in the catch that overwrites the row to +// failed. +const failLogActivityForAction = { current: null as string | null }; +vi.mock("../services/activity-log.js", async (importOriginal) => { + const actual = + await importOriginal(); + return { + ...actual, + logActivity: vi.fn( + async (...args: Parameters) => { + const input = args[1] as { action?: string }; + if ( + failLogActivityForAction.current && + input?.action === failLogActivityForAction.current + ) { + throw new Error(`forced activity write failure: ${input.action}`); + } + return actual.logActivity(...args); + }, + ), + }; +}); + +const { createToolGatewayService } = await import( + "../services/tool-gateway.js" +); + +const embeddedPostgresSupport = await getEmbeddedPostgresTestSupport(); +const describeEmbeddedPostgres = embeddedPostgresSupport.supported + ? describe + : describe.skip; + +type Db = ReturnType; + +const SIGNING_SECRET = "connector-telemetry-signing-secret"; + +function mcpSuccessResponse(init: RequestInit): Response { + const body = JSON.parse(String(init.body ?? "{}")); + return new Response( + JSON.stringify({ + jsonrpc: "2.0", + id: body?.id ?? "fixture", + result: { content: [{ type: "text", text: "ok" }] }, + }), + { status: 200, headers: { "content-type": "application/json" } }, + ); +} + +async function createCompany(db: Db) { + return db + .insert(companies) + .values({ + name: `Connector telemetry ${randomUUID()}`, + issuePrefix: `CT${randomUUID().slice(0, 6).toUpperCase()}`, + }) + .returning() + .then((rows) => rows[0]!); +} + +async function createAgent(db: Db, companyId: string) { + return db + .insert(agents) + .values({ + companyId, + name: `Agent ${randomUUID()}`, + role: "engineer", + adapterType: "process", + adapterConfig: {}, + runtimeConfig: {}, + permissions: {}, + }) + .returning() + .then((rows) => rows[0]!); +} + +async function createIssueAndRun(db: Db, companyId: string, agentId: string) { + const project = await db + .insert(projects) + .values({ companyId, name: `Project ${randomUUID()}` }) + .returning() + .then((rows) => rows[0]!); + const issue = await db + .insert(issues) + .values({ + companyId, + projectId: project.id, + title: `Telemetry issue ${randomUUID()}`, + status: "in_progress", + assigneeAgentId: agentId, + }) + .returning() + .then((rows) => rows[0]!); + const run = await db + .insert(heartbeatRuns) + .values({ + companyId, + agentId, + invocationSource: "assignment", + status: "running", + contextSnapshot: { issueId: issue.id, projectId: project.id }, + }) + .returning() + .then((rows) => rows[0]!); + return { issue, run }; +} + +async function allowAllToolsForAgent(db: Db, companyId: string, agentId: string) { + const profile = await db + .insert(toolProfiles) + .values({ + companyId, + profileKey: `telemetry-all-${randomUUID()}`, + name: `Telemetry profile ${randomUUID()}`, + defaultAction: "allow", + }) + .returning() + .then((rows) => rows[0]!); + await db.insert(toolProfileBindings).values({ + companyId, + profileId: profile.id, + targetType: "agent", + targetId: agentId, + }); +} + +async function createRemoteMcpTool(db: Db, companyId: string) { + const [application] = await db + .insert(toolApplications) + .values({ + companyId, + applicationKey: `telemetry-app-${randomUUID().slice(0, 8)}`, + name: `Remote app ${randomUUID()}`, + type: "mcp_http", + status: "active", + }) + .returning(); + const config = { + url: "https://mcp.telemetry.test/mcp", + sourceTemplateKey: "github", + }; + const [connection] = await db + .insert(toolConnections) + .values({ + companyId, + applicationId: application!.id, + name: `Remote connection ${randomUUID()}`, + uid: `test/${randomUUID()}`, + transport: "mcp_remote", + status: "active", + enabled: true, + healthStatus: "ok", + config, + transportConfig: config, + credentialRefs: [], + credentialSecretRefs: [], + }) + .returning(); + await db.insert(connectionGrants).values({ + companyId, + connectionId: connection!.id, + kind: "organization", + credentialSecretRefs: [], + status: "active", + isDefault: true, + }); + const [catalogEntry] = await db + .insert(toolCatalogEntries) + .values({ + companyId, + applicationId: application!.id, + connectionId: connection!.id, + entryKind: "tool", + name: `kv_set-${randomUUID()}`, + toolName: "kv_set", + title: "KV Set", + description: "Set a value", + inputSchema: { + type: "object", + properties: { key: { type: "string" }, value: { type: "string" } }, + required: ["key", "value"], + additionalProperties: false, + }, + annotations: { readOnlyHint: false }, + riskLevel: "write", + isReadOnly: false, + isWrite: true, + isDestructive: false, + status: "active", + versionHash: randomUUID(), + }) + .returning(); + return { application: application!, connection: connection!, catalogEntry: catalogEntry! }; +} + +async function settledEvents(expectedCount: number) { + await vi.waitFor(() => { + expect( + track.mock.calls.filter( + ([name]) => name === "connection.invoked", + ).length, + ).toBeGreaterThanOrEqual(expectedCount); + }); + // The emitter runs as detached background work; give any surplus emission a + // chance to land before asserting the exact count. + await new Promise((resolve) => setTimeout(resolve, 150)); + return track.mock.calls.filter( + ([name]) => name === "connection.invoked", + ); +} + +describeEmbeddedPostgres("gateway connector invocation telemetry", () => { + let db!: Db; + let tempDb: Awaited> | null = + null; + + beforeAll(async () => { + tempDb = await startEmbeddedPostgresTestDatabase( + "paperclip-connector-telemetry-", + ); + db = createDb(tempDb.connectionString); + }, 30_000); + + afterEach(async () => { + track.mockClear(); + failLogActivityForAction.current = null; + await db.delete(activityLog); + await db.delete(toolCallEvents); + await db.delete(toolRuntimeSlots); + await db.delete(toolGatewaySessions); + await db.delete(toolGatewayRateLimitCounters); + await db.delete(toolActionRequests); + await db.delete(toolInvocations); + await db.delete(toolAccessAuditEvents); + await db.delete(toolPolicies); + await db.delete(toolProfileBindings); + await db.delete(toolProfiles); + await db.delete(toolCatalogEntries); + await db.delete(connectionGrants); + await db.delete(toolConnections); + await db.delete(toolApplications); + await db.delete(heartbeatRuns); + await db.delete(issues); + await db.delete(projects); + await db.delete(agents); + await db.delete(companyMemberships); + await db.delete(companies); + }); + + afterAll(async () => { + await tempDb?.cleanup(); + }); + + async function fixture(options: Parameters[1] = {}) { + const company = await createCompany(db); + const agent = await createAgent(db, company.id); + const { run } = await createIssueAndRun(db, company.id, agent.id); + const remoteTool = await createRemoteMcpTool(db, company.id); + await allowAllToolsForAgent(db, company.id, agent.id); + const gateway = createToolGatewayService(db, { + toolActionSigningSecret: SIGNING_SECRET, + remoteHttpRequest: async (_url, init) => mcpSuccessResponse(init), + ...options, + }); + const session = await gateway.createSession({ + companyId: company.id, + agentId: agent.id, + runId: run.id, + }); + const tool = (await gateway.listToolsForSession(session.token)).find( + (candidate) => candidate.connectionId === remoteTool.connection.id, + ); + expect(tool).toBeDefined(); + return { company, agent, run, remoteTool, gateway, session, tool: tool! }; + } + + it("emits exactly one succeeded event for a successful agent execution", async () => { + const f = await fixture(); + await expect( + f.gateway.executeTool({ + sessionToken: f.session.token, + tool: f.tool.name, + parameters: { key: "a", value: "1" }, + }), + ).resolves.toMatchObject({ status: "completed" }); + + const events = await settledEvents(1); + expect(events).toHaveLength(1); + expect(events[0]?.[1]).toMatchObject({ + connector_key: "github", + transport: "mcp_remote", + status: "succeeded", + origin: "agent", + }); + }, 30_000); + + it("emits a single failed event when post-success audit work throws (no contradictory pair)", async () => { + const f = await fixture(); + failLogActivityForAction.current = "tool_gateway.call_completed"; + + await expect( + f.gateway.executeTool({ + sessionToken: f.session.token, + tool: f.tool.name, + parameters: { key: "a", value: "1" }, + }), + ).rejects.toThrow(/forced activity write failure/); + + const [invocation] = await db.select().from(toolInvocations); + expect(invocation?.status).toBe("failed"); + + const events = await settledEvents(1); + expect(events).toHaveLength(1); + expect(events[0]?.[1]).toMatchObject({ status: "failed", origin: "agent" }); + expect( + events.some(([, dims]) => (dims as { status: string }).status === "succeeded"), + ).toBe(false); + }, 30_000); + + it("emits one failed event when the remote execution itself fails", async () => { + const f = await fixture({ + remoteHttpRequest: async () => + new Response("upstream exploded", { status: 500 }), + }); + await expect( + f.gateway.executeTool({ + sessionToken: f.session.token, + tool: f.tool.name, + parameters: { key: "a", value: "1" }, + }), + ).rejects.toThrow(); + + const events = await settledEvents(1); + expect(events).toHaveLength(1); + expect(events[0]?.[1]).toMatchObject({ status: "failed" }); + }, 30_000); + + it("emits one denied event for a policy block and one rate_limited event past the limit", async () => { + const f = await fixture(); + await db.insert(toolPolicies).values({ + companyId: f.company.id, + name: "Block telemetry tool", + policyType: "block", + selectors: { connectionId: f.remoteTool.connection.id }, + priority: 1, + }); + await expect( + f.gateway.executeTool({ + sessionToken: f.session.token, + tool: f.tool.name, + parameters: { key: "a", value: "1" }, + }), + ).rejects.toMatchObject({ reasonCode: expect.any(String) }); + + let events = await settledEvents(1); + expect(events).toHaveLength(1); + expect(events[0]?.[1]).toMatchObject({ status: "denied", origin: "agent" }); + + track.mockClear(); + await db.delete(toolPolicies); + await db.insert(toolPolicies).values({ + companyId: f.company.id, + name: "One call only", + policyType: "rate_limit", + selectors: { connectionId: f.remoteTool.connection.id }, + config: { limit: 1, windowSeconds: 60 }, + priority: 1, + }); + await expect( + f.gateway.executeTool({ + sessionToken: f.session.token, + tool: f.tool.name, + parameters: { key: "rate", value: "first" }, + }), + ).resolves.toMatchObject({ status: "completed" }); + await expect( + f.gateway.executeTool({ + sessionToken: f.session.token, + tool: f.tool.name, + parameters: { key: "rate", value: "second" }, + }), + ).rejects.toMatchObject({ status: 429 }); + + events = await settledEvents(2); + expect(events).toHaveLength(2); + expect(events.map(([, dims]) => (dims as { status: string }).status).sort()).toEqual( + ["rate_limited", "succeeded"], + ); + }, 30_000); + + it("approval flow: pending emits nothing, decline emits denied, approve emits succeeded", async () => { + const f = await fixture(); + await db.insert(toolPolicies).values({ + companyId: f.company.id, + name: "Review telemetry tool", + policyType: "require_approval", + selectors: { connectionId: f.remoteTool.connection.id }, + priority: 1, + }); + + await f.gateway + .executeTool({ + sessionToken: f.session.token, + tool: f.tool.name, + parameters: { key: "a", value: "declined" }, + }) + .then( + () => { + throw new Error("expected approval_required"); + }, + (error) => expect(error).toMatchObject({ reasonCode: "approval_required" }), + ); + // Awaiting approval is in-flight: no completion event may exist yet. The + // approval card itself emits the registered interaction.created event, so + // filter to the proposed connector event. + await new Promise((resolve) => setTimeout(resolve, 150)); + expect( + track.mock.calls.filter(([name]) => name === "connection.invoked"), + ).toHaveLength(0); + + const [pendingRequest] = await db.select().from(toolActionRequests); + await f.gateway.declineActionRequest({ + companyId: f.company.id, + actionRequestId: pendingRequest!.id, + actor: { userId: "board" }, + }); + let events = await settledEvents(1); + expect(events).toHaveLength(1); + expect(events[0]?.[1]).toMatchObject({ status: "denied", origin: "agent" }); + + track.mockClear(); + await f.gateway + .executeTool({ + sessionToken: f.session.token, + tool: f.tool.name, + parameters: { key: "a", value: "approved" }, + }) + .then( + () => { + throw new Error("expected approval_required"); + }, + (error) => expect(error).toMatchObject({ reasonCode: "approval_required" }), + ); + const approvable = await db + .select() + .from(toolActionRequests) + .where(eq(toolActionRequests.status, "pending")) + .then((rows) => rows[0]!); + await f.gateway.approveActionRequest({ + companyId: f.company.id, + actionRequestId: approvable.id, + actor: { userId: "board" }, + }); + events = await settledEvents(1); + expect(events).toHaveLength(1); + expect(events[0]?.[1]).toMatchObject({ status: "succeeded", origin: "agent" }); + }, 40_000); + + it("labels connection-test executions as setup_test", async () => { + const f = await fixture(); + await db.insert(companyMemberships).values({ + companyId: f.company.id, + principalType: "user", + principalId: "board-user", + status: "active", + membershipRole: "member", + }); + const testCall = await f.gateway.executeTestCall({ + companyId: f.company.id, + connectionId: f.remoteTool.connection.id, + agentId: f.agent.id, + userId: "board-user", + toolName: f.tool.name, + parameters: { key: "t", value: "1" }, + }); + expect(testCall).toMatchObject({ decision: "allowed" }); + expect(testCall).not.toHaveProperty("error"); + + const events = await settledEvents(1); + expect(events).toHaveLength(1); + expect(events[0]?.[1]).toMatchObject({ + status: "succeeded", + origin: "setup_test", + }); + }, 30_000); + + it("an idempotent replay returns the stored result without a second event", async () => { + const f = await fixture(); + const idempotencyKey = `telemetry-replay-${randomUUID()}`; + await expect( + f.gateway.executeTool({ + sessionToken: f.session.token, + tool: f.tool.name, + parameters: { key: "a", value: "1" }, + idempotencyKey, + }), + ).resolves.toMatchObject({ status: "completed" }); + await settledEvents(1); + await expect( + f.gateway.executeTool({ + sessionToken: f.session.token, + tool: f.tool.name, + parameters: { key: "a", value: "1" }, + idempotencyKey, + }), + ).resolves.toMatchObject({ status: "replayed" }); + + const events = await settledEvents(1); + expect(events).toHaveLength(1); + }, 30_000); + + it("telemetry sink failures never fail the tool call", async () => { + const f = await fixture(); + track.mockImplementation(() => { + throw new Error("sink offline"); + }); + await expect( + f.gateway.executeTool({ + sessionToken: f.session.token, + tool: f.tool.name, + parameters: { key: "a", value: "1" }, + }), + ).resolves.toMatchObject({ status: "completed" }); + const [invocation] = await db.select().from(toolInvocations); + expect(invocation?.status).toBe("succeeded"); + track.mockReset(); + }, 30_000); +}); diff --git a/server/src/services/connector-telemetry.test.ts b/server/src/services/connector-telemetry.test.ts new file mode 100644 index 0000000000..2fef1f55cc --- /dev/null +++ b/server/src/services/connector-telemetry.test.ts @@ -0,0 +1,269 @@ +import { toolConnections, toolInvocations, type Db } from "@paperclipai/db"; +import { beforeEach, describe, expect, it, vi } from "vitest"; + +const track = vi.fn(); +const isRegisteredEventName = vi.fn(() => true); +let telemetryClient: { + track: typeof track; + isRegisteredEventName: typeof isRegisteredEventName; +} | null = { track, isRegisteredEventName }; + +vi.mock("../telemetry.js", () => ({ + getTelemetryClient: () => telemetryClient, +})); + +const { + connectorKeyForConnection, + emitConnectionCreated, + emitConnectionUpdated, + emitConnectionInvoked, + invocationOrigin, +} = await import("./connector-telemetry.js"); + +type ToolConnectionRow = typeof toolConnections.$inferSelect; +type ToolInvocationRow = typeof toolInvocations.$inferSelect; + +function connectionRow(overrides: Record = {}): ToolConnectionRow { + return { + id: "conn-1", + connectionPurpose: "tool", + transport: "mcp_remote", + authKind: "oauth", + status: "active", + enabled: true, + config: { sourceTemplateKey: "github" }, + ...overrides, + } as unknown as ToolConnectionRow; +} + +function invocationRow(overrides: Record = {}): ToolInvocationRow { + return { + id: "inv-1", + connectionId: "conn-1", + status: "succeeded", + actorType: "agent", + runId: "run-1", + issueId: "issue-1", + gatewayId: "gw-1", + startedAt: new Date("2026-09-17T00:00:00.000Z"), + completedAt: new Date("2026-09-17T00:00:07.400Z"), + ...overrides, + } as unknown as ToolInvocationRow; +} + +function fakeDb( + invocation: ToolInvocationRow | null, + connection: ToolConnectionRow | null, +): Db & { selectCalls: () => number } { + let selectCalls = 0; + return { + selectCalls: () => selectCalls, + select: () => { + selectCalls += 1; + return { + from: () => ({ + innerJoin: () => ({ + where: () => + Promise.resolve( + // Mirror the emitter's single joined read: an invocation with + // no (or a deleted) connection joins to zero rows. + invocation && connection && invocation.connectionId + ? [{ invocation, connection }] + : [], + ), + }), + }), + }; + }, + } as unknown as Db & { selectCalls: () => number }; +} + +beforeEach(() => { + track.mockReset(); + isRegisteredEventName.mockReset(); + isRegisteredEventName.mockReturnValue(true); + telemetryClient = { track, isRegisteredEventName }; +}); + +describe("connectorKeyForConnection", () => { + it("keeps a sourceTemplateKey that resolves in the shared catalog", () => { + expect(connectorKeyForConnection(connectionRow())).toBe("github"); + }); + + it("reports custom for keys outside the catalog", () => { + expect( + connectorKeyForConnection(connectionRow({ config: { sourceTemplateKey: "my-internal-server" } })), + ).toBe("custom"); + }); + + it("reports custom when the key is missing or not a string", () => { + expect(connectorKeyForConnection(connectionRow({ config: {} }))).toBe("custom"); + expect( + connectorKeyForConnection(connectionRow({ config: { sourceTemplateKey: 42 } })), + ).toBe("custom"); + }); +}); + +describe("invocationOrigin", () => { + it("labels user-driven connection tests without run context as setup_test", () => { + expect( + invocationOrigin( + invocationRow({ actorType: "user", runId: null, issueId: null, gatewayId: null }), + ), + ).toBe("setup_test"); + }); + + it("keeps the actor type for run-attached invocations", () => { + expect(invocationOrigin(invocationRow())).toBe("agent"); + expect(invocationOrigin(invocationRow({ actorType: "user" }))).toBe("user"); + }); +}); + +describe("emitConnectionCreated", () => { + it("emits catalog identity and lifecycle state for tool connections", () => { + emitConnectionCreated(connectionRow(), "gallery"); + expect(track).toHaveBeenCalledWith("connection.created", { + connector_key: "github", + transport: "mcp_remote", + auth_kind: "oauth", + setup_flow: "gallery", + status: "active", + enabled: true, + }); + }); + + it("never emits for channel or AI connections", () => { + emitConnectionCreated(connectionRow({ connectionPurpose: "channel" }), "api"); + expect(track).not.toHaveBeenCalled(); + }); + + it("is a no-op without an initialized telemetry client", () => { + telemetryClient = null; + expect(() => emitConnectionCreated(connectionRow(), "gallery")).not.toThrow(); + }); + + it("swallows client failures instead of failing the setup path", () => { + track.mockImplementation(() => { + throw new Error("sink offline"); + }); + expect(() => emitConnectionCreated(connectionRow(), "gallery")).not.toThrow(); + }); +}); + +describe("emitConnectionUpdated", () => { + it("stays silent when the committed write left status and enabled unchanged", () => { + emitConnectionUpdated( + connectionRow(), + { status: "active", enabled: true }, + "api", + ); + expect(track).not.toHaveBeenCalled(); + }); + + it("emits the transition when lifecycle status changes", () => { + emitConnectionUpdated( + connectionRow(), + { status: "draft", enabled: true }, + "oauth_callback", + ); + expect(track).toHaveBeenCalledWith("connection.updated", { + connector_key: "github", + transport: "mcp_remote", + auth_kind: "oauth", + change_source: "oauth_callback", + previous_status: "draft", + status: "active", + previous_enabled: true, + enabled: true, + }); + }); + + it("emits when only the enabled flag flips", () => { + emitConnectionUpdated( + connectionRow({ enabled: false }), + { status: "active", enabled: true }, + "api", + ); + expect(track).toHaveBeenCalledTimes(1); + expect(track.mock.calls[0]?.[1]).toMatchObject({ + previous_enabled: true, + enabled: false, + }); + }); + + it("never emits for non-tool connections even on lifecycle changes", () => { + emitConnectionUpdated( + connectionRow({ connectionPurpose: "channel" }), + { status: "draft", enabled: true }, + "api", + ); + expect(track).not.toHaveBeenCalled(); + }); +}); + +describe("emitConnectionInvoked", () => { + it("emits terminal outcome, origin, and duration from the committed row", async () => { + await emitConnectionInvoked(fakeDb(invocationRow(), connectionRow()), "inv-1"); + expect(track).toHaveBeenCalledWith("connection.invoked", { + connector_key: "github", + transport: "mcp_remote", + status: "succeeded", + origin: "agent", + duration_seconds: 7, + }); + }); + + it("omits duration when the invocation never recorded a start time", async () => { + await emitConnectionInvoked( + fakeDb(invocationRow({ startedAt: null, status: "denied" }), connectionRow()), + "inv-1", + ); + expect(track.mock.calls[0]?.[1]).not.toHaveProperty("duration_seconds"); + expect(track.mock.calls[0]?.[1]).toMatchObject({ status: "denied" }); + }); + + it("ignores in-flight statuses so approval waits never count as failures", async () => { + await emitConnectionInvoked( + fakeDb(invocationRow({ status: "awaiting_approval" }), connectionRow()), + "inv-1", + ); + expect(track).not.toHaveBeenCalled(); + }); + + it("ignores invocations without a connection or with a non-tool connection", async () => { + await emitConnectionInvoked( + fakeDb(invocationRow({ connectionId: null }), connectionRow()), + "inv-1", + ); + await emitConnectionInvoked( + fakeDb(invocationRow(), connectionRow({ connectionPurpose: "channel" })), + "inv-1", + ); + expect(track).not.toHaveBeenCalled(); + }); + + it("resolves without emitting when the invocation row is missing", async () => { + await expect( + emitConnectionInvoked(fakeDb(null, null), "inv-missing"), + ).resolves.toBeUndefined(); + expect(track).not.toHaveBeenCalled(); + }); + + it("reads nothing from the database while the event name is unregistered", async () => { + isRegisteredEventName.mockReturnValue(false); + const db = fakeDb(invocationRow(), connectionRow()); + await emitConnectionInvoked(db, "inv-1"); + expect(db.selectCalls()).toBe(0); + expect(track).not.toHaveBeenCalled(); + expect(isRegisteredEventName).toHaveBeenCalledWith( + "connection.invoked", + ); + }); + + it("loads the invocation and connection with a single read once registered", async () => { + const db = fakeDb(invocationRow(), connectionRow()); + await emitConnectionInvoked(db, "inv-1"); + expect(db.selectCalls()).toBe(1); + expect(track).toHaveBeenCalledTimes(1); + }); +}); diff --git a/server/src/services/connector-telemetry.ts b/server/src/services/connector-telemetry.ts new file mode 100644 index 0000000000..66305e42bd --- /dev/null +++ b/server/src/services/connector-telemetry.ts @@ -0,0 +1,226 @@ +import { eq } from "drizzle-orm"; +import { toolConnections, toolInvocations, type Db } from "@paperclipai/db"; +import { getConnectableAppDefinition } from "@paperclipai/shared"; +import { + trackConnectionCreated, + trackConnectionUpdated, + trackConnectionInvoked, +} from "@paperclipai/shared/telemetry"; +import { logger } from "../middleware/logger.js"; +import { getTelemetryClient } from "../telemetry.js"; + +type ToolConnectionRow = typeof toolConnections.$inferSelect; +type ToolInvocationRow = typeof toolInvocations.$inferSelect; + +export type ConnectorSetupFlow = "gallery" | "api" | "example"; +export type ConnectorChangeSource = + | "api" + | "gallery" + | "oauth_callback" + | "credential_refresh" + | "archive" + | "example"; + +/** + * Terminal `tool_invocations.status` values. `pending`, `authorized`, + * `awaiting_approval`, and `executing` are in-flight and never emit: pending + * approval is not a failure. `cancelled` is declared in the schema enum but no + * current code path writes it to an invocation (signing-path cancellations + * land on the action request and fail the invocation); it stays here so a + * future writer is counted without a telemetry change. + */ +const TERMINAL_INVOCATION_STATUSES: ReadonlySet = new Set([ + "succeeded", + "failed", + "denied", + "cancelled", + "timed_out", + "rate_limited", +]); + +/** + * Connector identity for telemetry is the reviewed first-party catalog slug + * only: `config.sourceTemplateKey` validated against the shared app-definitions + * catalog. Anything else — custom MCP servers, user-named connections, keys + * that no longer resolve to a catalog entry — reports the literal `custom` so + * no user- or provider-controlled string leaves the process. + */ +export function connectorKeyForConnection( + connection: Pick, +): string { + const raw = connection.config?.sourceTemplateKey; + const key = typeof raw === "string" ? raw : null; + return key && getConnectableAppDefinition(key) ? key : "custom"; +} + +/** + * Setup-test invocations (Apps → Test tab) are separated from agent usage on + * durable invocation columns, mirroring the gateway's own test-origin + * predicate, so the split survives reloads and approval-driven completion. + */ +export function invocationOrigin( + invocation: Pick< + ToolInvocationRow, + "actorType" | "runId" | "issueId" | "gatewayId" | "connectionId" + >, +): string { + const isTestOrigin = + invocation.actorType === "user" && + invocation.runId === null && + invocation.issueId === null && + invocation.gatewayId === null && + invocation.connectionId !== null; + return isTestOrigin ? "setup_test" : invocation.actorType; +} + +function isToolPurpose(connection: Pick): boolean { + return connection.connectionPurpose === "tool"; +} + +/** + * Emits one proposed `connection.created` event for a tool-purpose + * connection row that a caller already committed. Channel and AI connections + * never emit. This function never throws; a telemetry failure must never fail + * a connection setup path. + */ +export function emitConnectionCreated( + connection: ToolConnectionRow, + setupFlow: ConnectorSetupFlow, +): void { + try { + const client = getTelemetryClient(); + if (!client) return; + if (!isToolPurpose(connection)) return; + trackConnectionCreated(client, { + connector_key: connectorKeyForConnection(connection), + transport: connection.transport, + auth_kind: connection.authKind, + setup_flow: setupFlow, + status: connection.status, + enabled: connection.enabled, + }); + } catch (err) { + logger.warn( + { err, connectionId: connection.id }, + "failed to emit connection.created telemetry", + ); + } +} + +/** + * Emits one proposed `connection.updated` event when a committed + * write changed a tool-purpose connection's persisted lifecycle state + * (`status` or `enabled`). Metadata-only saves, health polls, credential + * rotation, and catalog refreshes do not change either field and therefore + * never emit. This function never throws. + */ +export function emitConnectionUpdated( + connection: ToolConnectionRow, + previous: Pick, + changeSource: ConnectorChangeSource, +): void { + try { + const client = getTelemetryClient(); + if (!client) return; + if (!isToolPurpose(connection)) return; + if ( + previous.status === connection.status && + previous.enabled === connection.enabled + ) { + return; + } + trackConnectionUpdated(client, { + connector_key: connectorKeyForConnection(connection), + transport: connection.transport, + auth_kind: connection.authKind, + change_source: changeSource, + previous_status: previous.status, + status: connection.status, + previous_enabled: previous.enabled, + enabled: connection.enabled, + }); + } catch (err) { + logger.warn( + { err, connectionId: connection.id }, + "failed to emit connection.updated telemetry", + ); + } +} + +/** + * Emits one proposed `connection.invoked` event for an invocation + * a caller already wrote to a terminal status. Despite the name, this is a + * completion event: it records finished invocation attempts and their + * terminal status, never invocation starts. Never await it: like + * `agent.task_run`, this is best-effort background work and must not delay + * the caller's own response or lifecycle writes. + * + * Call placement inside the gateway is deliberate and asymmetric. Failure + * paths call it right after their terminal save. Success paths call it only + * after the post-save bookkeeping (action-request settlement, tool-call + * event, audit) has completed, because a bookkeeping failure lands in a catch + * that overwrites the row to `failed` and emits there — emitting the success + * beforehand would let one execution report both outcomes. + * + * Self-guards, so callers do not have to re-check anything: non-terminal + * statuses, invocations without a tool-purpose connection, and an + * unregistered (still-proposed) event name all return without emitting. The + * registration gate keeps the completion path free of telemetry reads until + * schema adoption: the runtime client drops unregistered names in `track`, + * so there is no point loading rows for an event that cannot be queued. + * + * Delivery is best-effort per execution, not exactly-once. Known residual + * duplications, accepted for the proposal stage: an outer safety-net catch + * can re-save `failed` and emit again when a failure path's own bookkeeping + * throws after its emit (test-invocation approval flow), and a crash-side + * replay that re-writes a terminal status emits again. Deduplication state + * is out of scope for this proposal. + */ +export async function emitConnectionInvoked( + db: Db, + invocationId: string, +): Promise { + try { + const client = getTelemetryClient(); + if (!client) return; + if (!client.isRegisteredEventName("connection.invoked")) return; + + const joined = await db + .select({ invocation: toolInvocations, connection: toolConnections }) + .from(toolInvocations) + .innerJoin( + toolConnections, + eq(toolInvocations.connectionId, toolConnections.id), + ) + .where(eq(toolInvocations.id, invocationId)) + .then((rows) => rows[0] ?? null); + if (!joined) return; + const { invocation, connection } = joined; + if (!TERMINAL_INVOCATION_STATUSES.has(invocation.status)) return; + if (!isToolPurpose(connection)) return; + + const startedAtMs = invocation.startedAt + ? new Date(invocation.startedAt).getTime() + : null; + const completedAtMs = invocation.completedAt + ? new Date(invocation.completedAt).getTime() + : null; + const durationSeconds = + startedAtMs !== null && completedAtMs !== null + ? Math.max(0, Math.round((completedAtMs - startedAtMs) / 1000)) + : undefined; + + trackConnectionInvoked(client, { + connector_key: connectorKeyForConnection(connection), + transport: connection.transport, + status: invocation.status, + origin: invocationOrigin(invocation), + ...(durationSeconds === undefined ? {} : { duration_seconds: durationSeconds }), + }); + } catch (err) { + logger.warn( + { err, invocationId }, + "failed to emit connection.invoked telemetry", + ); + } +} diff --git a/server/src/services/tool-access.ts b/server/src/services/tool-access.ts index a8fc6224be..e4dc96b87b 100644 --- a/server/src/services/tool-access.ts +++ b/server/src/services/tool-access.ts @@ -2,6 +2,10 @@ import { isRemoteMcpConnectorMethod, connectionPurposeTransportSchema } from "@p import { instanceSettingsService } from "./instance-settings.js"; import { githubBotRequest } from "./chat-github-client.js"; import { syncConnectionCredentialBindings } from "./connection-credential-bindings.js"; +import { + emitConnectionCreated, + emitConnectionUpdated, +} from "./connector-telemetry.js"; import { canBrowseProjectRepositoryGrant, mergeProjectRepository } from "./project-repositories.js"; import { captureRunIdentity } from "./run-identity.js"; import { createHash, randomBytes, randomUUID } from "node:crypto"; @@ -6415,6 +6419,7 @@ export function toolAccessService( return { connection: updatedConnection, applicationArchived }; }); + emitConnectionUpdated(archived.connection, connection, "archive"); // Only now, with every access path closed, revoke the credentials. Each // `secrets.remove` marks the row deleted before it calls the provider, so a @@ -8017,6 +8022,7 @@ export function toolAccessService( await ensureDefaultOrganizationGrant(updated); await syncCredentialBindings(updated); await ensureRuntimeSlot(updated); + emitConnectionUpdated(updated, existing, "example"); return { row: updated, created: false }; } const connectionId = randomUUID(); @@ -8045,6 +8051,7 @@ export function toolAccessService( await ensureDefaultOrganizationGrant(created); await syncCredentialBindings(created); await ensureRuntimeSlot(created); + emitConnectionCreated(created, "example"); return { row: created, created: true }; } @@ -10387,7 +10394,10 @@ export function toolAccessService( ), ) .returning(); - if (updated) await syncCredentialBindings(updated); + if (updated) { + await syncCredentialBindings(updated); + emitConnectionUpdated(updated, connection, "credential_refresh"); + } return updated ?? null; } @@ -11275,6 +11285,11 @@ export function toolAccessService( ) .returning(); await syncCredentialBindings(reauthorizationRequired); + emitConnectionUpdated( + reauthorizationRequired, + latestConnection, + "credential_refresh", + ); } else { const activeGrantRefs = await db .select({ @@ -12724,6 +12739,15 @@ export function toolAccessService( }) .returning(); } + if (revivedConnectionPrevious) { + emitConnectionUpdated( + connectionRow, + revivedConnectionPrevious, + "gallery", + ); + } else { + emitConnectionCreated(connectionRow, "gallery"); + } if (personalIdentityUserId) { // "Just me" (PAP-17835 seam #4). The credential is committed straight to // the caller's own grant; the connection row keeps only the header @@ -13749,6 +13773,11 @@ export function toolAccessService( return { profileId, profileBindings, policies, updatedConnection }; }); + emitConnectionUpdated( + transactionResult.updatedConnection, + connection, + "gallery", + ); const details = await profileDetails( transactionResult.profileId, @@ -14859,6 +14888,10 @@ export function toolAccessService( stateRow.connectionId, stateRow.companyId, ); + const preCloudCallbackLifecycle = { + status: connection.status, + enabled: connection.enabled, + }; // The connection lifecycle, not the incidental presence of its app profile, // distinguishes setup from reauthorization. New connections and connections // revived after removal are drafts until this callback completes. A profile @@ -15220,6 +15253,10 @@ export function toolAccessService( .where(eq(issueThreadInteractions.id, stateRow.interactionId)); } }); + // The transaction above committed the connection's lifecycle write; a + // Paperclip Cloud connector callback is the managed variant of an OAuth + // callback completion. + emitConnectionUpdated(connection, preCloudCallbackLifecycle, "oauth_callback"); if (githubMetadata) { const [githubGrant] = await db .select({ id: connectionGrants.id }) @@ -15452,6 +15489,10 @@ export function toolAccessService( : null, }); } + const preActivationLifecycle = { + status: connection.status, + enabled: connection.enabled, + }; [connection] = await db .update(toolConnections) .set({ @@ -15468,6 +15509,11 @@ export function toolAccessService( ), ) .returning(); + emitConnectionUpdated( + connection, + preActivationLifecycle, + "oauth_callback", + ); await db .update(toolApplications) .set({ status: "active", updatedAt: now() }) @@ -15566,6 +15612,10 @@ export function toolAccessService( // Reauthorization refreshes credentials and catalog without rebuilding the // operator's action profile, policy rules, or access bindings. const shouldFinalizeDefaults = connection.status === "draft"; + const preCallbackLifecycle = { + status: connection.status, + enabled: connection.enabled, + }; const sourceTemplateKey = typeof connection.config.sourceTemplateKey === "string" ? connection.config.sourceTemplateKey @@ -15841,6 +15891,11 @@ export function toolAccessService( tx, ); }); + emitConnectionUpdated( + connection, + preCallbackLifecycle, + "oauth_callback", + ); // Personal OAuth used to return immediately after saving the grant. That // left the connection draft/paused and its catalog empty, so the person @@ -16059,6 +16114,11 @@ export function toolAccessService( await ensureDefaultOrganizationGrant(connection, tx); await syncCredentialBindings(connection, [], tx); }); + emitConnectionUpdated( + connection, + preCallbackLifecycle, + "oauth_callback", + ); await checkConnectionHealth(connection.id, input.actor); const refresh = await refreshCatalog(connection.id, input.actor, { @@ -16115,6 +16175,10 @@ export function toolAccessService( if (requestingAgentId) await assertAgentsInCompany(companyId, [requestingAgentId]); let connection = await getConnectionRow(connectionId, companyId); + const preFinalizeLifecycle = { + status: connection.status, + enabled: connection.enabled, + }; if (connection.authKind !== "oauth") throw badRequest("This connection does not use browser sign-in"); if (connection.status === "archived") @@ -16424,6 +16488,12 @@ export function toolAccessService( for (const secretId of personalSecretIds) await secrets.remove(secretId); } + // Both identity branches above commit their own lifecycle write (draft -> + // active) before `finishGalleryAppConnection` re-reads the row, so that + // step alone would never observe this transition. This is the OAuth + // access-finalization step of the same gallery setup flow (Apps and inline + // task cards both reach it), hence `gallery`. + emitConnectionUpdated(connection, preFinalizeLifecycle, "gallery"); const catalog = await db .select() .from(toolCatalogEntries) @@ -17349,6 +17419,7 @@ export function toolAccessService( await ensureDefaultOrganizationGrant(row); await syncCredentialBindings(row); await ensureRuntimeSlot(row); + emitConnectionCreated(row, "api"); return toConnection(row); }, @@ -18239,6 +18310,7 @@ export function toolAccessService( .returning(); await syncCredentialBindings(row); await ensureRuntimeSlot(row); + emitConnectionUpdated(row, existing, "api"); return toConnection(row); }, diff --git a/server/src/services/tool-action-review.ts b/server/src/services/tool-action-review.ts index 01450174b2..b429b35ce6 100644 --- a/server/src/services/tool-action-review.ts +++ b/server/src/services/tool-action-review.ts @@ -9,6 +9,7 @@ import { } from "@paperclipai/db"; import { conflict, forbidden, notFound } from "../errors.js"; import { assertIssueThreadInteractionResolverAudience } from "./issue-thread-interaction-resolution.js"; +import { emitConnectionInvoked } from "./connector-telemetry.js"; import { toolAccessPolicyService } from "./tool-access-policy.js"; import { logActivity, @@ -237,6 +238,11 @@ export async function commitToolActionReview( } return updated; }); + // A rejection writes the invocation's terminal `denied` status inside the + // transaction above; emit only after that commit. Approvals stay in-flight + // and reach their terminal status in the gateway execution paths. + if (input.decision === "rejected") + void emitConnectionInvoked(db, source.invocationId); for (const publication of publications) publishActivity(publication); return result; } diff --git a/server/src/services/tool-gateway.ts b/server/src/services/tool-gateway.ts index 56ab25c088..c815717f19 100644 --- a/server/src/services/tool-gateway.ts +++ b/server/src/services/tool-gateway.ts @@ -8,6 +8,7 @@ import { githubGuestBotConnectionForSession, githubBotToolsForSession } from "./ import { githubChatReviewService } from "./chat-github-reviews.js"; import { runIdentityContexts } from "@paperclipai/db"; import { captureRunIdentity } from "./run-identity.js"; +import { emitConnectionInvoked } from "./connector-telemetry.js"; import { resolveManagedGitHubIdentitySelection } from "./git-credentials.js"; import { extractRemoteMcpPending } from "./remote-mcp-pending.js"; import { logger } from "../middleware/logger.js"; @@ -2223,6 +2224,7 @@ export function createToolGatewayService( updatedAt: new Date(), }) .where(eq(toolInvocations.id, input.invocation.id)); + void emitConnectionInvoked(db, input.invocation.id); await writeToolCallEvent({ invocationId: input.invocation.id, actionRequestId: input.actionRequest?.id ?? null, @@ -2260,6 +2262,7 @@ export function createToolGatewayService( updatedAt: new Date(), }) .where(eq(toolInvocations.id, input.invocation.id)); + void emitConnectionInvoked(db, input.invocation.id); throw new ToolGatewayHttpError( 500, "Approval request was not created", @@ -2309,6 +2312,7 @@ export function createToolGatewayService( updatedAt: new Date(), }) .where(eq(toolInvocations.id, input.invocation.id)); + void emitConnectionInvoked(db, input.invocation.id); throw new ToolGatewayHttpError( 500, error.message, @@ -2475,6 +2479,7 @@ export function createToolGatewayService( updatedAt: new Date(), }) .where(eq(toolInvocations.id, input.invocation.id)); + void emitConnectionInvoked(db, input.invocation.id); throw new ToolGatewayHttpError( 409, "The approval request was resolved before it could be signed", @@ -7177,6 +7182,10 @@ export function createToolGatewayService( execution: connectedMcpExecution.execution, }, }); + // After the bookkeeping above: a throw there is caught below, overwrites + // the row to failed, and emits — an earlier success emit would make one + // execution report both outcomes. + void emitConnectionInvoked(db, args.invocationId); return { decision: "allowed" as const, invocationId: args.invocationId, @@ -7206,6 +7215,7 @@ export function createToolGatewayService( updatedAt: new Date(), }) .where(eq(toolInvocations.id, args.invocationId)); + void emitConnectionInvoked(db, args.invocationId); await writeToolCallEvent({ invocationId: args.invocationId, session: args.session, @@ -7320,6 +7330,7 @@ export function createToolGatewayService( updatedAt: new Date(), }) .where(eq(toolInvocations.id, invocation.id)); + void emitConnectionInvoked(db, invocation.id); await reflectToolActionInteractionLifecycle({ actionRequestId, status: "failed", @@ -7363,6 +7374,7 @@ export function createToolGatewayService( updatedAt: new Date(), }) .where(eq(toolInvocations.id, invocation.id)); + void emitConnectionInvoked(db, invocation.id); await reflectToolActionInteractionLifecycle({ actionRequestId, status: "failed", @@ -7528,6 +7540,7 @@ export function createToolGatewayService( return true; }); if (!settled) return { reasonCode, message, settled: false }; + void emitConnectionInvoked(db, input.invocationId); await reflectToolActionInteractionLifecycle({ actionRequestId: input.actionRequestId, status: "failed", @@ -7580,6 +7593,7 @@ export function createToolGatewayService( updatedAt: now, }) .where(eq(toolInvocations.id, input.invocationId)); + void emitConnectionInvoked(db, input.invocationId); await reflectToolActionInteractionLifecycle({ actionRequestId: expired.id, status: "expired", @@ -7976,6 +7990,7 @@ export function createToolGatewayService( updatedAt: now, }) .where(eq(toolInvocations.id, invocation.id)); + void emitConnectionInvoked(db, invocation.id); await db .update(toolActionRequests) .set({ status: "executed", resolvedAt: now, updatedAt: now }) @@ -9045,6 +9060,7 @@ export function createToolGatewayService( updatedAt: new Date(), }) .where(eq(toolInvocations.id, invocationId)); + void emitConnectionInvoked(db, invocationId); throw new ToolGatewayHttpError( 500, error.message, @@ -9117,6 +9133,9 @@ export function createToolGatewayService( } if (!accessDecision.allowed) { + // recordInvocation inserted this row already terminal (denied or + // rate_limited), so this is its only completion boundary. + void emitConnectionInvoked(db, invocationId); await writeAudit({ session, companyId: input.companyId, @@ -9316,6 +9335,7 @@ export function createToolGatewayService( updatedAt: now, }) .where(eq(toolInvocations.id, row.invocationId)); + void emitConnectionInvoked(db, row.invocationId); await reflectToolActionInteractionLifecycle({ actionRequestId: row.id, status, @@ -10281,6 +10301,9 @@ export function createToolGatewayService( }); } if (!accessDecision.allowed) { + // recordInvocation inserted this row already terminal (denied or + // rate_limited), so this is its only completion boundary. + void emitConnectionInvoked(db, invocationId); await writeAudit({ session, companyId: session.companyId, @@ -10485,6 +10508,10 @@ export function createToolGatewayService( execution: connectedMcpExecution?.execution ?? undefined, }, }); + // After the bookkeeping above: a throw there is caught below, overwrites + // the row to failed, and emits — an earlier success emit would make one + // execution report both outcomes. + void emitConnectionInvoked(db, invocationId); return { invocationId, status: "completed" as const, @@ -10541,6 +10568,7 @@ export function createToolGatewayService( updatedAt: completedAt, }) .where(eq(toolInvocations.id, invocationId)); + void emitConnectionInvoked(db, invocationId); if (input.approvedActionRequestId) { const [failedRequest] = await db .update(toolActionRequests) @@ -10722,6 +10750,9 @@ export function createToolGatewayService( } if (!accessDecision.allowed) { + // recordInvocation inserted this row already terminal (denied or + // rate_limited), so this is its only completion boundary. + void emitConnectionInvoked(db, invocationId); await writeAudit({ session: sessionLike, companyId: input.runContext.companyId, @@ -10836,6 +10867,10 @@ export function createToolGatewayService( resultSummary: resultValidation.summary, }, }); + // After the bookkeeping above: a throw there is caught below, overwrites + // the row to failed, and emits — an earlier success emit would make one + // execution report both outcomes. + void emitConnectionInvoked(db, invocationId); return resultValidation.value as typeof result; } catch (err) { const status = err instanceof ToolGatewayHttpError ? err.status : 502; @@ -10856,6 +10891,7 @@ export function createToolGatewayService( updatedAt: new Date(), }) .where(eq(toolInvocations.id, invocationId)); + void emitConnectionInvoked(db, invocationId); await writeToolCallEvent({ invocationId, session: sessionLike,