fix: restore V3 HTTP observability

Project one typed HttpExecutionObservation per logical V3 execution through a
closed composition-root projector: only registered diagnostic context keys and
bucketed values reach the sinks, and terminal non-abort failures now emit
exactly one api.request.failed telemetry event. Caller cancellation and scope
fencing record a diagnostic but never a failure event.

routeId becomes a required input at the installed operation-executor boundary
so the feature gateway's low-cardinality route identity survives to the sink.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
DongHyeonka
2026-08-13 22:42:51 +09:00
co-authored by Claude Opus 5
parent f7bec8274b
commit 67cc5b6d2c
9 changed files with 559 additions and 59 deletions
@@ -98,7 +98,7 @@ describe("HTTP operation execution contract", () => {
executor.execute(
nonKeyedCommand,
{ name: "created" },
{ scope, intent },
{ routeId: "TEST_ROUTE", scope, intent },
),
).resolves.toMatchObject({ kind: "SUCCESS" });
expect(observedKeys).toEqual([null]);
@@ -109,7 +109,7 @@ describe("HTTP operation execution contract", () => {
executor.execute(
nonKeyedCommand,
{ name: "created" },
{ scope, intent: { ...intent, idempotencyKey: "unexpected-key" } },
{ routeId: "TEST_ROUTE", scope, intent: { ...intent, idempotencyKey: "unexpected-key" } },
),
).resolves.toMatchObject({
kind: "CONTRACT_VIOLATION",
@@ -0,0 +1,319 @@
import { describe, expect, it, vi } from "vitest";
import { createContractHttpExecutor } from "../../src/adapters/http/http-execution-v3.ts";
import type {
HttpExecutionObservation,
} from "../../src/adapters/http/http-execution-v3.ts";
import { createHttpObservationProjector } from "../../src/bootstrap/runtime-adapters.ts";
import { createReferenceFeatureInstalledInput } from "../../src/features/reference-feature/adapters/create-reference-feature-input.ts";
import {
DIAGNOSTIC_CONTEXT_ALLOWLIST,
projectDiagnosticRecord,
type DiagnosticRecordInput,
} from "../../src/contracts/diagnostics.ts";
import {
projectTelemetryEvent,
type TelemetryEventName,
} from "../../src/contracts/telemetry.ts";
import {
TEST_CREATE_HTTP_CONTRACT,
TEST_LIST_HTTP_CONTRACT,
} from "../helpers/external-contract-fixture.ts";
const ROUTE_ID = "TEST_ROUTE";
function scopeSnapshot(isCurrent: () => boolean = () => true) {
return Object.freeze({
generation: 1,
fingerprint: "scope-1",
identities: Object.freeze({}) as never,
signal: new AbortController().signal,
isCurrent,
});
}
type RecordedDiagnostic = DiagnosticRecordInput;
type RecordedTelemetry = Readonly<{
eventName: TelemetryEventName;
attributes: Readonly<Record<string, unknown>>;
}>;
function recordingSinks() {
const diagnostics: RecordedDiagnostic[] = [];
const telemetry: RecordedTelemetry[] = [];
return {
diagnostics,
telemetry,
projector: createHttpObservationProjector({
diagnostics: {
record(input: DiagnosticRecordInput) {
diagnostics.push(input);
},
},
telemetry: {
emit(
eventName: TelemetryEventName,
attributes: Record<string, unknown>,
) {
telemetry.push(Object.freeze({ eventName, attributes }));
},
},
}),
};
}
describe("V3 HTTP observability projection", () => {
it("projects every V3 terminal outcome through the closed diagnostics allowlist", async () => {
const sinks = recordingSinks();
const cases: readonly Readonly<{
label: string;
fetcher: typeof fetch;
}>[] = [
{
label: "SUCCESS",
fetcher: (async () =>
Response.json([{ id: "a", name: "A" }])) as unknown as typeof fetch,
},
{
label: "PROBLEM",
fetcher: (async () =>
Response.json(
{ type: "about:blank", title: "nope", status: 400 },
{ status: 400 },
)) as unknown as typeof fetch,
},
{
label: "TRANSPORT_FAILURE",
fetcher: (async () => {
throw new TypeError("network down");
}) as unknown as typeof fetch,
},
{
label: "CONTRACT_VIOLATION",
fetcher: (async () =>
new Response("<html/>", {
status: 200,
headers: { "content-type": "text/html" },
})) as unknown as typeof fetch,
},
];
for (const testCase of cases) {
const executor = createContractHttpExecutor({
baseUrl: "https://api.example/",
maxRetryAttempts: 0,
attachCredentials: () => ({
kind: "READY" as const,
headers: {},
credentials: "omit" as const,
}),
fetcher: testCase.fetcher,
observe: sinks.projector,
});
await executor.execute(
TEST_LIST_HTTP_CONTRACT,
{ limit: 5 },
{ routeId: ROUTE_ID, scope: scopeSnapshot() },
);
}
expect(sinks.diagnostics).toHaveLength(cases.length);
for (const recorded of sinks.diagnostics) {
expect(recorded.eventId).toBe("http.request.completed");
const contextKeys = Object.keys(recorded.context ?? {});
expect(contextKeys).toEqual(
expect.arrayContaining([
"route_id",
"operation_id",
"operation",
"outcome",
"error_kind",
"http_status_group",
"attempt_count_bucket",
"duration_bucket",
]),
);
for (const key of contextKeys) {
expect(DIAGNOSTIC_CONTEXT_ALLOWLIST).toContain(key);
}
expect(contextKeys).not.toContain("attempts");
expect(contextKeys).not.toContain("certainty");
const projection = projectDiagnosticRecord(recorded);
expect(projection.success).toBe(true);
}
});
it("emits one failure telemetry event for a non-abort terminal failure", async () => {
const sinks = recordingSinks();
const executor = createContractHttpExecutor({
baseUrl: "https://api.example/",
maxRetryAttempts: 0,
attachCredentials: () => ({
kind: "READY" as const,
headers: {},
credentials: "omit" as const,
}),
fetcher: (async () => {
throw new TypeError("network down");
}) as unknown as typeof fetch,
observe: sinks.projector,
});
const outcome = await executor.execute(
TEST_LIST_HTTP_CONTRACT,
{ limit: 5 },
{ routeId: ROUTE_ID, scope: scopeSnapshot() },
);
expect(outcome.kind).toBe("TRANSPORT_FAILURE");
expect(sinks.telemetry).toHaveLength(1);
const emitted = sinks.telemetry[0];
expect(emitted?.eventName).toBe("api.request.failed");
expect(emitted?.attributes).toMatchObject({
route_id: ROUTE_ID,
operation_id: "TEST_LIST_ENTITIES",
error_kind: "NETWORK_FAILURE",
http_status_group: "none",
attempt_count_bucket: "1",
});
const projected = projectTelemetryEvent(
emitted?.eventName ?? "api.request.failed",
emitted?.attributes ?? {},
);
expect(projected.success).toBe(true);
expect(sinks.diagnostics).toHaveLength(1);
});
it("does not emit failure telemetry for caller cancellation or scope fencing", async () => {
const cancelled = recordingSinks();
const callerController = new AbortController();
callerController.abort();
const cancelledExecutor = createContractHttpExecutor({
baseUrl: "https://api.example/",
maxRetryAttempts: 0,
attachCredentials: () => ({
kind: "READY" as const,
headers: {},
credentials: "omit" as const,
}),
fetcher: (async () =>
Response.json([{ id: "a", name: "A" }])) as unknown as typeof fetch,
observe: cancelled.projector,
});
const cancelledOutcome = await cancelledExecutor.execute(
TEST_LIST_HTTP_CONTRACT,
{ limit: 5 },
{
routeId: ROUTE_ID,
scope: scopeSnapshot(),
signal: callerController.signal,
},
);
expect(cancelledOutcome.kind).toBe("CANCELLED");
expect(cancelled.diagnostics).toHaveLength(1);
expect(cancelled.telemetry).toHaveLength(0);
const fenced = recordingSinks();
const fencedExecutor = createContractHttpExecutor({
baseUrl: "https://api.example/",
maxRetryAttempts: 0,
attachCredentials: () => ({
kind: "READY" as const,
headers: {},
credentials: "omit" as const,
}),
fetcher: (async () =>
Response.json([{ id: "a", name: "A" }])) as unknown as typeof fetch,
observe: fenced.projector,
});
const fencedOutcome = await fencedExecutor.execute(
TEST_LIST_HTTP_CONTRACT,
{ limit: 5 },
{ routeId: ROUTE_ID, scope: scopeSnapshot(() => false) },
);
expect(fencedOutcome.kind).toBe("CONTRACT_VIOLATION");
expect(fenced.diagnostics).toHaveLength(1);
expect(fenced.telemetry).toHaveLength(0);
});
it("preserves the feature route id through the installed operation executor", async () => {
const seen: Array<Readonly<Record<string, unknown>>> = [];
const installed = createReferenceFeatureInstalledInput({
contractOperations: Object.freeze({
async execute(
_operationId: string,
_input: unknown,
context?: Readonly<Record<string, unknown>>,
) {
seen.push(Object.freeze({ ...(context ?? {}) }));
return Object.freeze({
kind: "SUCCESS" as const,
value: [],
metadata: Object.freeze({ status: 200 }),
effect: "NOT_APPLICABLE" as const,
});
},
}),
});
await installed.input.listResources({ limit: 20 });
await installed.input.getResource("resource-1");
expect(seen.map((context) => context.routeId)).toEqual([
"REFERENCE_RESOURCE_LIST",
"REFERENCE_RESOURCE_DETAIL",
]);
});
it("cannot change the HTTP result when diagnostics or telemetry throws", async () => {
const projector = createHttpObservationProjector({
diagnostics: {
record() {
throw new Error("diagnostics sink exploded");
},
},
telemetry: {
emit() {
throw new Error("telemetry sink exploded");
},
},
});
const observe = vi.fn((observation: HttpExecutionObservation) => {
projector(observation);
});
const executor = createContractHttpExecutor({
baseUrl: "https://api.example/",
maxRetryAttempts: 0,
attachCredentials: () => ({
kind: "READY" as const,
headers: {},
credentials: "omit" as const,
}),
fetcher: (async () =>
Response.json({ id: "created", name: "Created" }, {
status: 201,
})) as unknown as typeof fetch,
observe,
});
const outcome = await executor.execute(
TEST_CREATE_HTTP_CONTRACT,
{ name: "Created" },
{
routeId: ROUTE_ID,
scope: scopeSnapshot(),
intent: Object.freeze({
intentId: "intent-1",
operationId: "TEST_CREATE_ENTITY",
canonicalInputIdentity: "opaque-input-identity",
idempotencyKey: "key-1",
createdAtMonotonicMs: 1,
}),
},
);
expect(outcome.kind).toBe("SUCCESS");
expect(observe).toHaveBeenCalledTimes(1);
});
});
+17 -10
View File
@@ -4,7 +4,10 @@ import path from "node:path";
import { afterAll, beforeAll, describe, expect, it } from "vitest";
import { readBoundedBytes } from "../../src/adapters/http/bounded-body-reader.ts";
import { createContractHttpExecutor } from "../../src/adapters/http/http-execution-v3.ts";
import {
createContractHttpExecutor,
type HttpExecutionObservation,
} from "../../src/adapters/http/http-execution-v3.ts";
import {
type InstalledHttpContract,
} from "../../src/contracts/external-contract-runtime.ts";
@@ -30,6 +33,13 @@ import {
type HttpScenarioOperationId,
} from "../mocks/scenarios/catalog.ts";
/** The reference gateway owns these low-cardinality route identities. */
function routeIdFor(operationId: HttpScenarioOperationId): string {
return operationId === "GET_REFERENCE_RESOURCE"
? "REFERENCE_RESOURCE_DETAIL"
: "REFERENCE_RESOURCE_LIST";
}
const RECEIPT_PATH = path.resolve(
"artifacts/tests/http-scenario-executions.json",
);
@@ -191,11 +201,7 @@ async function executeScenario(
): Promise<HttpScenarioAssertionGroups> {
const physicalAttempts: AttemptTrace[] = [];
const sleeps: RetryReason[] = [];
const observations: Array<Readonly<{
outcome: string;
attempts: number;
certainty: string;
}>> = [];
const observations: HttpExecutionObservation[] = [];
const caller = new AbortController();
const scopeLifetime = new AbortController();
let scopeCurrent = true;
@@ -265,6 +271,7 @@ async function executeScenario(
try {
const execution = executor.execute(operation, inputFor(entry.operationId), {
routeId: routeIdFor(entry.operationId),
scope,
signal: caller.signal,
...(entry.operationId === "CREATE_REFERENCE_RESOURCE"
@@ -296,7 +303,7 @@ async function executeScenario(
const observation = observations[0]!;
const observedSignal = scopeLifetime.signal.aborted ? "ABORTED" : "ACTIVE";
const cancellationOwner =
observation.certainty === "TIMEOUT"
observation.terminalReason === "TIMEOUT"
? "DEADLINE"
: caller.signal.aborted
? "CALLER"
@@ -314,13 +321,13 @@ async function executeScenario(
}),
effect: Object.freeze({
outcome: String(outcome.effect),
observer: observation.certainty,
observer: observation.terminalReason,
}),
retry: Object.freeze({ count: sleeps.length, reasons: Object.freeze(sleeps) }),
fetch: Object.freeze({
count: physicalAttempts.length,
observerAttempts: observation.attempts,
agrees: physicalAttempts.length === observation.attempts,
observerAttempts: observation.attemptCount,
agrees: physicalAttempts.length === observation.attemptCount,
}),
media: Object.freeze({
attempts: Object.freeze(physicalAttempts.map((attempt) => attempt.media)),
+27 -15
View File
@@ -10,6 +10,8 @@ import {
const installed: InstalledHttpContract<unknown, unknown, unknown> =
TEST_LIST_HTTP_CONTRACT;
const ROUTE_ID = "TEST_ROUTE";
const scope = Object.freeze({
generation: 1,
fingerprint: "scope-1",
@@ -82,7 +84,7 @@ describe("descriptor-driven HTTP execution lifetime", () => {
});
await expect(
executor.execute(installed, { limit: 20 }, { scope }),
executor.execute(installed, { limit: 20 }, { routeId: ROUTE_ID, scope }),
).resolves.toMatchObject({
kind: "RATE_LIMITED",
effect: "NOT_APPLICABLE",
@@ -119,7 +121,7 @@ describe("descriptor-driven HTTP execution lifetime", () => {
executor.execute(
createInstalled,
{ name: "created" },
{ scope, ...(intent === undefined ? {} : { intent }) },
{ routeId: ROUTE_ID, scope, ...(intent === undefined ? {} : { intent }) },
),
).resolves.toMatchObject({
kind: "CONTRACT_VIOLATION",
@@ -154,7 +156,7 @@ describe("descriptor-driven HTTP execution lifetime", () => {
});
await expect(
executor.execute(installed, { limit: 20 }, { scope, intent }),
executor.execute(installed, { limit: 20 }, { routeId: ROUTE_ID, scope, intent }),
).resolves.toMatchObject({
kind: "CONTRACT_VIOLATION",
violation: {
@@ -173,14 +175,14 @@ describe("descriptor-driven HTTP execution lifetime", () => {
label: "query with the canonical reserved header",
operation: installed,
input: { limit: 20 },
context: { scope },
context: { routeId: ROUTE_ID, scope },
headerName: "Idempotency-Key",
},
{
label: "valid KEYED command with a case-variant reserved header",
operation: createInstalled,
input: { name: "created" },
context: { scope, intent: mutationIntent() },
context: { routeId: ROUTE_ID, scope, intent: mutationIntent() },
headerName: "iDeMpOtEnCy-KeY",
},
])(
@@ -252,7 +254,7 @@ describe("descriptor-driven HTTP execution lifetime", () => {
executor.execute(
retryingCreate,
{ name: "created" },
{ scope, intent: mutationIntent({ idempotencyKey: "logical-key" }) },
{ routeId: ROUTE_ID, scope, intent: mutationIntent({ idempotencyKey: "logical-key" }) },
),
).resolves.toMatchObject({ kind: "SUCCESS" });
expect(fetcher).toHaveBeenCalledTimes(2);
@@ -284,7 +286,7 @@ describe("descriptor-driven HTTP execution lifetime", () => {
executor.execute(
installed,
{ limit: 20 },
{ scope },
{ routeId: ROUTE_ID, scope },
),
).resolves.toMatchObject({ kind: "SUCCESS" });
expect(fetcher).toHaveBeenCalledOnce();
@@ -317,6 +319,7 @@ describe("descriptor-driven HTTP execution lifetime", () => {
createInstalled,
{ name: "created" },
{
routeId: ROUTE_ID,
scope,
intent: Object.freeze({
intentId: "private-intent-id",
@@ -353,7 +356,7 @@ describe("descriptor-driven HTTP execution lifetime", () => {
fetcher,
});
await expect(executor.execute(installed, input, { scope })).resolves.toMatchObject({
await expect(executor.execute(installed, input, { routeId: ROUTE_ID, scope })).resolves.toMatchObject({
kind: "SUCCESS",
});
expect(String((fetcher.mock.calls as unknown[][])[0]?.[0])).toContain(
@@ -389,11 +392,11 @@ describe("descriptor-driven HTTP execution lifetime", () => {
},
} as unknown as typeof installed;
await expect(executor.execute(throwing, { limit: 20 }, { scope })).resolves.toMatchObject({
await expect(executor.execute(throwing, { limit: 20 }, { routeId: ROUTE_ID, scope })).resolves.toMatchObject({
kind: "CONTRACT_VIOLATION",
effect: "NOT_APPLICABLE",
});
await expect(executor.execute(malformed, { limit: 20 }, { scope })).resolves.toMatchObject({
await expect(executor.execute(malformed, { limit: 20 }, { routeId: ROUTE_ID, scope })).resolves.toMatchObject({
kind: "CONTRACT_VIOLATION",
effect: "NOT_APPLICABLE",
});
@@ -416,6 +419,7 @@ describe("descriptor-driven HTTP execution lifetime", () => {
createInstalled,
{ name: "created" },
{
routeId: ROUTE_ID,
scope,
intent: mutationIntent(),
},
@@ -458,6 +462,7 @@ describe("descriptor-driven HTTP execution lifetime", () => {
createInstalled,
{ name: "created" },
{
routeId: ROUTE_ID,
scope: fencedScope,
intent: mutationIntent({ intentId: "intent-2", idempotencyKey: "key-2" }),
},
@@ -479,7 +484,7 @@ describe("descriptor-driven HTTP execution lifetime", () => {
});
const result = executor
.execute(operation({ deadlineMs: 5 }), { limit: 20 }, { scope })
.execute(operation({ deadlineMs: 5 }), { limit: 20 }, { routeId: ROUTE_ID, scope })
.then((outcome) => {
settled = true;
return outcome;
@@ -532,7 +537,7 @@ describe("descriptor-driven HTTP execution lifetime", () => {
const result = executor.execute(
operation({ deadlineMs: 5 }),
{ limit: 20 },
{ scope },
{ routeId: ROUTE_ID, scope },
);
await vi.advanceTimersByTimeAsync(5);
await flushMicrotasks();
@@ -548,7 +553,14 @@ describe("descriptor-driven HTTP execution lifetime", () => {
expect(observe).toHaveBeenCalledTimes(iterationCount);
for (const [observation] of observe.mock.calls) {
expect(observation).toEqual(
expect.objectContaining({ attempts: 1, certainty: "TIMEOUT" }),
expect.objectContaining({
attemptCount: 1,
errorKind: "TIMEOUT",
terminalReason: "TIMEOUT",
cancellationOwner: "DEADLINE",
routeId: ROUTE_ID,
operationId: "TEST_LIST_ENTITIES",
}),
);
}
} finally {
@@ -582,7 +594,7 @@ describe("descriptor-driven HTTP execution lifetime", () => {
});
const result = executor
.execute(installed, { limit: 20 }, { scope, signal: caller.signal })
.execute(installed, { limit: 20 }, { routeId: ROUTE_ID, scope, signal: caller.signal })
.then((outcome) => {
settled = true;
return outcome;
@@ -617,7 +629,7 @@ describe("descriptor-driven HTTP execution lifetime", () => {
executor.execute(
operation({ responseBody: "NONE" }),
{ limit: 20 },
{ scope },
{ routeId: ROUTE_ID, scope },
),
).resolves.toMatchObject({
kind: "TRANSPORT_FAILURE",