Skip to content
Closed
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
6 changes: 6 additions & 0 deletions docs/observability/op-catalog.generated.json
Original file line number Diff line number Diff line change
Expand Up @@ -626,6 +626,12 @@
"site": "packages/keiko-server/src/observability/server-logger.ts:180",
"package": "keiko-server"
},
{
"op": "chat.turn.started",
"category": "gateway",
"site": "packages/keiko-server/src/chat-handlers.ts:1687",
"package": "keiko-server"
},
{
"op": "embedding.memory.failed",
"category": "embedding",
Expand Down
20 changes: 10 additions & 10 deletions docs/qa/package-coverage-baseline.json
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
"keiko-cli": {
"files": 44,
"uncoveredFiles": 0,
"uncoveredLines": 387,
"uncoveredLines": 385,
"totalLines": 4880,
"coverage": {
"lines": 92.11,
Expand Down Expand Up @@ -234,15 +234,15 @@
}
},
"keiko-server": {
"files": 580,
"files": 581,
"uncoveredFiles": 0,
"uncoveredLines": 4516,
"totalLines": 56260,
"uncoveredLines": 4503,
"totalLines": 56309,
"coverage": {
"lines": 91.98,
"statements": 89.27,
"branches": 81.89,
"functions": 94.86
"lines": 92,
"statements": 89.3,
"branches": 81.91,
"functions": 94.91
}
},
"keiko-tools": {
Expand All @@ -261,10 +261,10 @@
"files": 419,
"uncoveredFiles": 3,
"uncoveredLines": 2967,
"totalLines": 39452,
"totalLines": 39455,
"coverage": {
"lines": 92.48,
"statements": 89.52,
"statements": 89.53,
"branches": 81.8,
"functions": 91.28
}
Expand Down
34 changes: 34 additions & 0 deletions packages/keiko-server/src/atlassian/syncRoutes.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -685,3 +685,37 @@ describe("Confluence sync — request validation fail-closed", () => {
);
});
});

describe("Confluence sync — governed (agent-initiated) start correlation", () => {
it("threads the request's own correlation id into an authority-denied governed start instead of minting one", async () => {
// ADR-0173 D5 / g12: ctx.correlationId is minted at request entry (server.ts) and is already
// in scope in handleStartAtlassianConnectorSync — the governed denial record must reuse it,
// not a disconnected randomUUID(). An authority referencing a runId the server-side registry
// never issued denies fast (authority-invalid) without needing a real envelope setup.
const port = createInMemoryConfluenceFixture({ baseUrl: BASE_URL, spaces: [] });
const { deps, credential } = depsFor(port);
const ctx = {
...ctxFor(
"POST",
{ authRef: credential.authRef },
{
spaceKeys: ["ENG"],
authority: {
runId: "unregistered-agent-run",
envelopeDigest: "0".repeat(64),
workspaceRoot: "/nonexistent/workspace",
},
},
),
correlationId: "req-governed-thread-01",
};

const result = await handleStartAtlassianConnectorSync(ctx, deps);

expect(result.status).toBe(200);
expect(result.body).toMatchObject({
disposition: "denied",
correlationId: "req-governed-thread-01",
});
});
});
13 changes: 11 additions & 2 deletions packages/keiko-server/src/atlassian/syncRoutes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -384,10 +384,10 @@ async function startSyncGoverned(
credential: AtlassianCredentialMetadata,
body: StartSyncBody,
authority: AtlassianActionAuthorityContext,
correlationId: string,
): Promise<RouteResult> {
const actionType = SYNC_ACTION_TYPE_FOR_PROVIDER[credential.provider];
const connectorId = connectorIdForAuthRef(credential.authRef);
const correlationId = randomUUID();
const targetRef = syncScopeTargetRef(body);
const outcome = decideGovernedAtlassianAction(actionType, authority, deps);
const denied = (reasonCode: AtlassianConnectorActivityReasonCode): RouteResult =>
Expand Down Expand Up @@ -428,7 +428,16 @@ export function handleStartAtlassianConnectorSync(
const credential = requireAtlassianCredential(ctx, guard);
const body = validateStartSyncBody(await readJsonObject(ctx.req), credential.provider);
if (body.authority !== undefined) {
return startSyncGoverned(deps, guard, credential, body, body.authority);
// Threads the request's own correlation id (ADR-0173 D5 / g12) into the governed-start
// denial/pending-approval/allowed records instead of a disconnected mint.
return startSyncGoverned(
deps,
guard,
credential,
body,
body.authority,
ctx.correlationId ?? randomUUID(),
);
}
// Direct human-triggered start: human-approved by construction (ADR-0129; ADR-0128 D5) —
// recorded as `allowed` + `human-initiated` on the run's activity record.
Expand Down
27 changes: 27 additions & 0 deletions packages/keiko-server/src/atlassian/writeActionRoutes.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1185,3 +1185,30 @@ describe("write-action route — clear-field validation (KEIKO-0319)", () => {
expect(result.status).toBe(400);
});
});

describe("write-action route — governed action correlation", () => {
it("threads the request's own correlation id into a policy-denied response instead of minting one", async () => {
// ADR-0173 D5 / g12: ctx.correlationId is minted at request entry (server.ts) and is already
// in scope in handleExecuteAtlassianConnectorAction — the governed-action denial record must
// reuse it, not a disconnected randomUUID(). An envelope with no write scope denies fast
// (policy-denied) without needing a provider round-trip.
const guard = guardWith({ count: 0, requests: [] });
const deniedAuthority = registerEnvelope("autonomous-delivery", []);
const result = (await handleExecuteAtlassianConnectorAction(
{
...ctx(
{ action: ACTION_REQUESTS["transition-issue"], authority: deniedAuthority },
{ authRef: JIRA_AUTH_REF },
),
correlationId: "req-write-thread-01",
},
deps(guard, "autonomous-delivery"),
)) as { status: number; body: Record<string, unknown> };

expect(result.status).toBe(200);
expect(result.body).toMatchObject({
disposition: "denied",
correlationId: "req-write-thread-01",
});
});
});
5 changes: 4 additions & 1 deletion packages/keiko-server/src/atlassian/writeActionRoutes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -690,9 +690,9 @@ function governedActionResult(
credential: AtlassianCredentialMetadata,
authority: AtlassianActionAuthorityContext,
plan: GovernedActionPlan,
correlationId: string,
): Promise<RouteResult> | RouteResult {
const connectorId = connectorIdForAuthRef(credential.authRef);
const correlationId = randomUUID();
const denied = (reasonCode: AtlassianConnectorActivityReasonCode): RouteResult =>
deniedAtlassianActionResult({
connectorId,
Expand Down Expand Up @@ -785,11 +785,14 @@ export function handleExecuteAtlassianConnectorAction(
throw invalid("authority must carry runId, envelopeDigest, and workspaceRoot");
}
const input = validateGovernedActionInput(body.action, credential.provider);
// Threads the request's own correlation id (ADR-0173 D5 / g12) into the governed-action
// denial/pending-approval/allowed records instead of a disconnected mint.
return governedActionResult(
deps,
credential,
authority,
actionPlanFor(guard, credential, input),
ctx.correlationId ?? randomUUID(),
);
});
}
Expand Down
49 changes: 48 additions & 1 deletion packages/keiko-server/src/browser-routes.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,10 @@ import { createRunRegistry } from "./runs.js";
import { createUiServer, UI_HOST } from "./server.js";
import { EventEmitter } from "node:events";
import type { ServerResponse } from "node:http";
import { openBrowserSseStream } from "./browser.js";
import { handleBrowserEvents, openBrowserSseStream } from "./browser.js";
import type { SseBackpressureSignal } from "./sse-write.js";
import type { RouteContext } from "./routes.js";
import type { ServerDiagnosticRecord, ServerDiagnosticSink } from "./diagnostics-log.js";
import {
BrowserToolError,
type BrowserEventEmitter,
Expand Down Expand Up @@ -821,3 +823,48 @@ describe("openBrowserSseStream backpressure (KEIKO-0142)", () => {
expect(fake.writes).toHaveLength(writesAfterClose);
});
});

describe("handleBrowserEvents backpressure correlation (ADR-0173 D5 / g12)", () => {
it("threads the request's own correlation id into the backpressure diagnostic instead of minting one", () => {
const fake = makeFakeSseRes();
fake.writeReturns = false; // rejects the ready frame -> immediate backpressure kill.
const manager = new FakeBrowserSessionManager();
manager.opened.push("session-thread");
const records: ServerDiagnosticRecord[] = [];
const diagnostics: ServerDiagnosticSink = {
record: (entry) => {
records.push(entry);
},
};
const baseDeps: UiHandlerDeps = {
config: undefined,
configPresent: false,
evidenceStore: {
put: (): string => "",
list: (): readonly string[] => [],
get: (): undefined => undefined,
delete: (): undefined => undefined,
},
env: process.env,
redactor: buildRedactor({}),
registry: createRunRegistry(),
modelPortFactory: (): undefined => undefined,
store: createInMemoryUiStore(),
browser: manager,
diagnostics,
};
const ctx: RouteContext = {
req: { on: (): void => undefined } as unknown as RouteContext["req"],
res: fake.res,
params: { sessionId: "session-thread" },
url: new URL("http://127.0.0.1/api/browser/sessions/session-thread/events"),
correlationId: "req-browser-thread-01",
};

handleBrowserEvents(ctx, baseDeps);

expect(records).toHaveLength(1);
expect(records[0]?.source).toBe("sse.browser.backpressure");
expect(records[0]?.correlationId).toBe("req-browser-thread-01");
});
});
4 changes: 3 additions & 1 deletion packages/keiko-server/src/browser.ts
Original file line number Diff line number Diff line change
Expand Up @@ -260,12 +260,14 @@ export function handleBrowserEvents(ctx: RouteContext, deps: UiHandlerDeps): Han
if (!guard.hasSession(sessionId)) {
return { status: 404, body: errorBody("SESSION_NOT_FOUND", "Browser session not found.") };
}
// Threads the request's own correlation id (ADR-0173 D5 / g12) so a later backpressure kill
// joins back to the request that opened this stream instead of a disconnected mint.
openBrowserSseStream(
ctx.res,
guard,
sessionId,
deps.redactor,
sseBackpressureReporter(deps, "browser"),
sseBackpressureReporter(deps, "browser", ctx.correlationId),
);
ctx.req.on("close", () => {
ctx.res.end();
Expand Down
7 changes: 7 additions & 0 deletions packages/keiko-server/src/chat-compaction-evidence.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -425,6 +425,13 @@ describe("chat compaction evidence wiring (ADR-0057 D3)", () => {
expect(records[0]?.message).toBe("Audit or evidence persistence failed.");
expect(records[0]?.errorClass).toMatch(/^[A-Z][A-Za-z0-9]*$/);
expect(records[0]?.correlationId).toMatch(/^[A-Za-z0-9._-]{8,128}$/);
// ADR-0173 D5 / g12: the failure's correlationId is THIS attempt's own runId (same
// derivation the successful-persist tests above pin), not a disconnected `randomUUID()` —
// an operator can join the failure back to the compaction attempt it belongs to. Before the
// fix this was a random UUID (with dashes) and never matched the runId shape below.
expect(records[0]?.correlationId).toBe(
`chat-${sha256Hex("chat-compaction-diagnostic").slice(0, 16)}-t4`,
);
expect(JSON.stringify(records)).not.toContain(SECRET);
expect(consoleWarn).not.toHaveBeenCalled();
} finally {
Expand Down
11 changes: 9 additions & 2 deletions packages/keiko-server/src/chat-compaction-evidence.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import { resolveCostClass } from "@oscharko-dev/keiko-model-gateway";
import { sha256Hex } from "@oscharko-dev/keiko-security";
import type { ContextCompactionRecord } from "@oscharko-dev/keiko-contracts";
import { randomUUID } from "node:crypto";
import { isValidCorrelationId } from "./correlation.js";
import type { UiHandlerDeps } from "./deps.js";
import { currentAuditRedactString, currentRedactionSecrets } from "./deps.js";
import {
Expand Down Expand Up @@ -49,11 +50,16 @@ export function persistChatCompactionEvidence(
return;
}
const record = input.compaction;
// Computed before the try so a persistence failure can still report under it (ADR-0173 D5 /
// g12): this run's own runId already ties the diagnostic back to the SAME compaction evidence
// attempt an operator would otherwise have to guess at from a disconnected mint.
let runId: string | undefined;
try {
const chatIdHash = sha256Hex(input.chatId);
runId = compactionRunId(chatIdHash, input.messageCount);
persistCompactionEvidence(
{
runId: compactionRunId(chatIdHash, input.messageCount),
runId,
modelId: input.modelId,
records: [record],
startedAt: input.startedAt,
Expand All @@ -75,10 +81,11 @@ export function persistChatCompactionEvidence(
// Best-effort stays best-effort — the send is unaffected — but the failure is no longer a
// `console.warn` carrying the raw error object on a channel production never overrode. It goes to
// the server's single redacted diagnostic sink so a compaction-evidence gap is observable.
const correlationId = runId !== undefined && isValidCorrelationId(runId) ? runId : randomUUID();
emitServerDiagnostic(
deps.diagnostics,
serverDiagnosticFromError({
correlationId: randomUUID(),
correlationId,
operation: "chat.compaction.evidence",
source: "chat-compaction-evidence",
error,
Expand Down
Loading
Loading