Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 46 additions & 0 deletions packages/backend/src/common/sentry-sampler.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
import { sampleRateFor } from './sentry-sampler';

const rates = { base: 0.05, mcp: 0.01 };

describe('sampleRateFor', () => {
it('samples MCP calls at the MCP rate', () => {
expect(
sampleRateFor(
{ name: 'POST', attributes: { 'http.method': 'POST', 'http.target': '/mcp/abc123?x=1' } },
rates,
),
).toBe(0.01);
expect(
sampleRateFor(
{ name: 'POST', normalizedRequest: { method: 'POST', url: 'https://cloud.example.com/mcp/abc' } },
rates,
),
).toBe(0.01);
expect(sampleRateFor({ name: 'POST /mcp/abc' }, rates)).toBe(0.01);
});

it('samples other requests at the base rate', () => {
expect(
sampleRateFor(
{ name: 'GET', attributes: { 'http.request.method': 'GET', 'url.path': '/api/connectors' } },
rates,
),
).toBe(0.05);
// Not the MCP route, only a lookalike prefix.
expect(
sampleRateFor({ name: 'GET', attributes: { 'http.method': 'GET', 'http.target': '/mcpx' } }, rates),
).toBe(0.05);
});

it('never samples health probes', () => {
expect(
sampleRateFor({ name: 'GET', attributes: { 'http.method': 'GET', 'http.target': '/health' } }, rates),
).toBe(0);
});

it('never samples root spans that are not requests', () => {
expect(sampleRateFor({ name: 'prisma:client:operation' }, rates)).toBe(0);
expect(sampleRateFor({ name: 'pg-pool.connect', attributes: { 'db.system': 'postgresql' } }, rates)).toBe(0);
expect(sampleRateFor({ name: 'SELECT "public"."kg_edges"."id" FROM "public"."kg_edges"' }, rates)).toBe(0);
});
});
74 changes: 74 additions & 0 deletions packages/backend/src/common/sentry-sampler.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
/**
* Which traces reach Sentry.
*
* A flat 5 % rate sent about 700,000 spans a day in September 2026, far past
* what a Sentry plan includes. Most of them were not requests at all:
*
* - root spans with no request behind them: Prisma and pg-pool spans from
* the knowledge-graph cron and other background work, each one sampled as
* a transaction of its own;
* - MCP tool calls, the bulk of the traffic, each carrying dozens of database
* spans.
*
* So: only incoming HTTP requests are traced; MCP calls at a lower rate than
* the rest (they are many and alike); health probes never. The decision of
* an incoming trace is not inherited: MCP clients send their own sampled
* traceparent, and that must not decide our volume.
*/

export interface SamplerRates {
/** Every HTTP request other than MCP and health probes. */
base: number;
/** POST /mcp/:serverId and the other MCP transport routes. */
mcp: number;
}

interface SamplingInput {
name: string;
attributes?: Record<string, unknown>;
normalizedRequest?: { url?: string; method?: string };
}

function str(v: unknown): string | undefined {
return typeof v === 'string' && v.length > 0 ? v : undefined;
}

/** The request path, from whichever of the places it can be at sampling time. */
function requestPath(ctx: SamplingInput): string | undefined {
const a = ctx.attributes ?? {};
const raw =
str(a['url.path']) ??
str(a['http.target']) ??
str(a['http.route']) ??
str(a['url.full']) ??
str(a['http.url']) ??
str(ctx.normalizedRequest?.url);
if (!raw) {
// "POST /mcp/abc" — the span name once the route is known.
const m = /^[A-Z]+ (\/\S*)/.exec(ctx.name);
return m?.[1];
}
try {
return raw.startsWith('/') ? raw.split('?')[0] : new URL(raw).pathname;
} catch {
return raw.split('?')[0];
}
}

function isHttpRequest(ctx: SamplingInput): boolean {
const a = ctx.attributes ?? {};
return Boolean(
str(a['http.request.method']) ??
str(a['http.method']) ??
str(ctx.normalizedRequest?.method) ??
(/^(GET|POST|PUT|PATCH|DELETE|HEAD|OPTIONS) /.test(ctx.name) ? 'yes' : undefined),
);
}

export function sampleRateFor(ctx: SamplingInput, rates: SamplerRates): number {
if (!isHttpRequest(ctx)) return 0;
const path = requestPath(ctx) ?? '';
if (/^\/health(\/|$)/.test(path)) return 0;
if (/^\/mcp(\/|$)/.test(path)) return rates.mcp;
return rates.base;
}
10 changes: 9 additions & 1 deletion packages/backend/src/instrument.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@
*/
import * as Sentry from '@sentry/nestjs';
import { scrubBreadcrumb, scrubEvent } from './common/sentry-scrub';
import { sampleRateFor } from './common/sentry-sampler';

const dsn = process.env.SENTRY_DSN;

Expand All @@ -30,7 +31,14 @@ if (dsn) {
release: process.env.SENTRY_RELEASE || process.env.npm_package_version,

// Tracing is opt-in on top of error reporting because it adds overhead.
tracesSampleRate: sample(process.env.SENTRY_TRACES_SAMPLE_RATE, 0.0),
// SENTRY_TRACES_SAMPLE_RATE applies to HTTP requests; MCP calls take
// SENTRY_MCP_TRACES_SAMPLE_RATE (default: a fifth of it); background work
// and health probes are not traced. See ./common/sentry-sampler.ts.
tracesSampler: (() => {
const base = sample(process.env.SENTRY_TRACES_SAMPLE_RATE, 0.0);
const rates = { base, mcp: sample(process.env.SENTRY_MCP_TRACES_SAMPLE_RATE, base / 5) };
return (ctx: Parameters<typeof sampleRateFor>[0]) => sampleRateFor(ctx, rates);
})(),
profilesSampleRate: sample(process.env.SENTRY_PROFILES_SAMPLE_RATE, 0.0),

sendDefaultPii: false,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -225,4 +225,40 @@ describe('KgObservationalService.ingestOrganization', () => {
const kinds = prisma.kgEdge.create.mock.calls.map((c: any[]) => c[0].data.kind).sort();
expect(kinds).toEqual(['produces_consumes', 'same_identity']);
});

it('writes each edge once per pass, however many values link the same nodes', async () => {
const occ = (hash: string) => [
{ valueHash: hash, connectorId: 'c2', entity: 'order', field: 'customer_id', direction: 'output' },
{ valueHash: hash, connectorId: 'c1', entity: 'customer', field: 'id', direction: 'input' },
];
const { svc, prisma } = make({
pages: [[row('r1', 'c2', 1, 'erp_get_order')]],
payloads: { r1: { input: {}, output: { order: { customer_id: 'CUST-12345' } } } },
valueSeen: ['h1', 'h2', 'h3', 'h4'].flatMap(occ),
});
// produces_consumes already seen 3 times; same_identity is new.
prisma.kgEdge.findUnique.mockImplementation(async (args: any) =>
args.where.organizationId_sourceNodeId_targetNodeId_kind.kind === 'produces_consumes'
? { id: 'e1', observations: 3, isManual: false, status: 'active' }
: null,
);
await svc.ingestOrganization(ORG);

expect(prisma.kgEdge.findUnique).toHaveBeenCalledTimes(2);
expect(prisma.kgEdge.update).toHaveBeenCalledTimes(1);
const upd = prisma.kgEdge.update.mock.calls[0][0].data;
// Four values, four observations, as four sequential bumps would leave it.
expect(upd.observations).toBe(7);
expect(upd.confidence).toBeCloseTo(0.55 + 0.05 * 6);
expect(upd.matchKey).toBe('id');

expect(prisma.kgEdge.create).toHaveBeenCalledTimes(1);
const created = prisma.kgEdge.create.mock.calls[0][0].data;
expect(created.kind).toBe('same_identity');
expect(created.observations).toBe(4);
expect(created.confidence).toBeCloseTo(0.2 + 0.05 * 3);
expect(created.status).toBe('suggested');

expect(prisma.kgEdge.updateMany).toHaveBeenCalledTimes(1);
});
});
93 changes: 69 additions & 24 deletions packages/backend/src/knowledge-graph/kg-observational.service.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { randomUUID } from 'node:crypto';
import { Injectable, Logger } from '@nestjs/common';
import * as Sentry from '@sentry/nestjs';
import { PrismaService } from '../common/prisma.service';
import { extractEntity } from './static/entity-extraction';
import { fkCandidate } from './static/fk-inference';
Expand Down Expand Up @@ -46,6 +47,18 @@ interface PageRow {
tool: { name: string } | null;
}

/** Observations of one edge within a correlate pass, written once at the end. */
interface EdgeBump {
src: string;
tgt: string;
kind: string;
times: number;
matchKey?: string;
base: number;
cap: number;
status: string;
}

interface ValueRow {
organizationId: string;
connectorId: string;
Expand Down Expand Up @@ -91,7 +104,10 @@ export class KgObservationalService {
}
this.inFlight.add(organizationId);
try {
return await this.ingest(organizationId);
// Not traced. The ingest is scheduled from inside an MCP tool call, so
// its hundreds of queries became spans of that call's trace: most of
// the cloud's Sentry volume in September 2026. Errors are still reported.
return await Sentry.suppressTracing(() => this.ingest(organizationId));
} finally {
this.inFlight.delete(organizationId);
}
Expand Down Expand Up @@ -280,7 +296,11 @@ export class KgObservationalService {
const src = await this.nodeId(refCache, organizationId, b.connectorId, b.from);
const tgt = await this.nodeId(refCache, organizationId, b.connectorId, b.to);
if (src && tgt && src !== tgt) {
await this.bumpEdge(organizationId, src, tgt, 'references', {
await this.bumpEdge(organizationId, {
src,
tgt,
kind: 'references',
times: 1,
matchKey: b.field,
base: 0.6,
cap: 0.9,
Expand All @@ -306,7 +326,11 @@ export class KgObservationalService {
const nb = await this.nodeId(coCache, organizationId, cb, eb);
if (!na || !nb || na === nb) continue;
const [src, tgt] = na < nb ? [na, nb] : [nb, na];
await this.bumpEdge(organizationId, src, tgt, 'related', {
await this.bumpEdge(organizationId, {
src,
tgt,
kind: 'related',
times: 1,
base: 0.3,
cap: 0.7,
status: 'suggested',
Expand Down Expand Up @@ -464,6 +488,23 @@ export class KgObservationalService {
}

const nodeCache = new Map<string, string>(); // `${connectorId}::${entity}` -> nodeId
// Many values link the same two nodes, so one pass used to write the same
// edge again for every one of them: a read and an update per value, most
// of the database work behind an MCP call. The bumps are counted here and
// each distinct edge is written once, with the same end state.
const bumps = new Map<string, EdgeBump>();
const confirmedRefs = new Map<string, [string, string]>();
const bump = (src: string, tgt: string, kind: string, b: Omit<EdgeBump, 'src' | 'tgt' | 'kind' | 'times'>) => {
const key = `${src}|${tgt}|${kind}`;
const prev = bumps.get(key);
if (prev) {
prev.times++;
// Sequential updates left the last matchKey that was set.
if (b.matchKey !== undefined) prev.matchKey = b.matchKey;
} else {
bumps.set(key, { src, tgt, kind, times: 1, ...b });
}
};
let edgeCount = 0;

for (const occ of byHash.values()) {
Expand All @@ -477,17 +518,14 @@ export class KgObservationalService {
const src = await this.nodeId(nodeCache, organizationId, p.connectorId, p.entity);
const tgt = await this.nodeId(nodeCache, organizationId, c.connectorId, c.entity);
if (!src || !tgt || src === tgt) continue;
await this.bumpEdge(organizationId, src, tgt, 'produces_consumes', {
bump(src, tgt, 'produces_consumes', {
matchKey: c.field,
base: 0.55,
cap: 0.95,
status: 'active',
});
// The data confirms any static FK guess between the same nodes.
await this.prisma.kgEdge.updateMany({
where: { organizationId, sourceNodeId: src, targetNodeId: tgt, kind: 'references', source: 'STATIC' },
data: { source: 'OBSERVED', confidence: 0.8, lastSeenAt: new Date() },
});
confirmedRefs.set(`${src}|${tgt}`, [src, tgt]);
edgeCount++;
}
}
Expand All @@ -505,7 +543,7 @@ export class KgObservationalService {
const nb = await this.nodeId(nodeCache, organizationId, b.connectorId, b.entity);
if (!na || !nb || na === nb) continue;
const [src, tgt] = na < nb ? [na, nb] : [nb, na];
await this.bumpEdge(organizationId, src, tgt, 'same_identity', {
bump(src, tgt, 'same_identity', {
matchKey: a.field === b.field ? a.field : undefined,
base: 0.2,
cap: 0.6,
Expand All @@ -517,6 +555,16 @@ export class KgObservationalService {
}
}
}

for (const b of bumps.values()) {
await this.bumpEdge(organizationId, b);
}
for (const [src, tgt] of confirmedRefs.values()) {
await this.prisma.kgEdge.updateMany({
where: { organizationId, sourceNodeId: src, targetNodeId: tgt, kind: 'references', source: 'STATIC' },
data: { source: 'OBSERVED', confidence: 0.8, lastSeenAt: new Date() },
});
}
return edgeCount;
}

Expand Down Expand Up @@ -546,13 +594,9 @@ export class KgObservationalService {
return node.id;
}

private async bumpEdge(
organizationId: string,
sourceNodeId: string,
targetNodeId: string,
kind: string,
opts: { matchKey?: string; base: number; cap: number; status: string },
): Promise<void> {
/** Add `times` observations to an edge, creating it if it is new. */
private async bumpEdge(organizationId: string, b: EdgeBump): Promise<void> {
const { src: sourceNodeId, tgt: targetNodeId, kind, times } = b;
const existing = await this.prisma.kgEdge.findUnique({
where: {
organizationId_sourceNodeId_targetNodeId_kind: {
Expand All @@ -564,35 +608,36 @@ export class KgObservationalService {
},
select: { id: true, observations: true, isManual: true, status: true },
});
const observations = (existing?.observations ?? 0) + times;
const confidence = Math.min(b.cap, b.base + 0.05 * (observations - 1));
if (!existing) {
await this.prisma.kgEdge.create({
data: {
organizationId,
sourceNodeId,
targetNodeId,
kind,
matchKey: opts.matchKey,
matchKey: b.matchKey,
source: 'OBSERVED',
confidence: Math.min(opts.cap, opts.base),
observations: 1,
status: opts.status,
confidence,
observations,
status: b.status,
},
});
return;
}
const observations = existing.observations + 1;
await this.prisma.kgEdge.update({
where: { id: existing.id },
data: {
observations,
confidence: Math.min(opts.cap, opts.base + 0.05 * (observations - 1)),
confidence,
source: 'OBSERVED',
matchKey: opts.matchKey,
matchKey: b.matchKey,
lastSeenAt: new Date(),
// Never silently re-open a link the user rejected, nor downgrade a confirmed one.
...(existing.isManual || existing.status === 'rejected'
? {}
: { status: opts.status }),
: { status: b.status }),
},
});
}
Expand Down
Loading