feat(telemetry): propose connector telemetry events behind proposal markers (#13582)

Adds connection.created / connection.updated / connection.invoked as proposed-telemetry wrappers (unregistered; the TelemetryClient cannot queue or send them) plus the server-side emitter wiring and tests. No telemetry is transmitted by this change. Proposal: #13578. Coordinated release plan and review package: PAP-24 (Step A.1).
This commit is contained in:
Michael Nguyen authored and GitHub committed 2026-09-23 22:32:40 +00:00
1 parent 6681c71b40
commit 30182c704c
11 files changed
+1881 -2

No files matched your search

+11 -1
View File
@@ -116,11 +116,21 @@ export class TelemetryClient {
* backend event schema.
*/
track<K extends TelemetryEventName>(eventName: K, ...args: TrackArgs<K>): 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
+69
View File
@@ -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 },
+3
View File
@@ -18,6 +18,9 @@ export {
trackErrorHandlerCrash,
trackInteractionCreated,
trackInteractionResolved,
trackConnectionCreated,
trackConnectionUpdated,
trackConnectionInvoked,
} from "./events.js";
export type {
TelemetryConfig,
@@ -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);
});
});
@@ -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<typeof createDb>;
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<ReturnType<typeof startEmbeddedPostgresTestDatabase>> | 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<typeof toolAccessService>[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);
});
@@ -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<typeof import("../services/activity-log.js")>();
return {
...actual,
logActivity: vi.fn(
async (...args: Parameters<typeof actual.logActivity>) => {
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<typeof createDb>;
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<ReturnType<typeof startEmbeddedPostgresTestDatabase>> | 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<typeof createToolGatewayService>[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);
});
@@ -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<string, unknown> = {}): 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<string, unknown> = {}): 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);
});
});
+226
View File
@@ -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<string> = 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<ToolConnectionRow, "config">,
): 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<ToolConnectionRow, "connectionPurpose">): 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<ToolConnectionRow, "status" | "enabled">,
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<void> {
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",
);
}
}
+73 -1
View File
@@ -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);
},
@@ -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;
}
+36
View File
@@ -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,