Skip to content
Draft
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
201 changes: 149 additions & 52 deletions apps/server/src/provider/EventNdjsonLogger.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
// @effect-diagnostics nodeBuiltinImport:off
import * as NodeCrypto from "node:crypto";
import * as NodeFS from "node:fs";
import * as NodeOS from "node:os";
import * as NodePath from "node:path";
Expand All @@ -17,6 +18,7 @@ import {
makeEventNdjsonLogger,
makeEventNdjsonLogStore,
type PendingRecord,
providerEventLogThreadSegment,
writeBatchedMessages,
} from "./EventNdjsonLogger.ts";

Expand All @@ -30,6 +32,16 @@ function ownedLogPath(basePath: string, segment: string): string {
return NodePath.join(NodePath.dirname(basePath), `${stem}.${segment}.log`);
}

function threadLogPath(basePath: string, threadId: string): string {
return ownedLogPath(basePath, providerEventLogThreadSegment(threadId));
}

function logFileNames(directory: string): Array<string> {
return NodeFS.readdirSync(directory)
.filter((name) => name.endsWith(".log"))
.toSorted();
}

function parseLogLine(line: string) {
const match = /^\[([^\]]+)\] ([A-Z]+): (.+)$/.exec(line);
assert.notEqual(match, null);
Expand Down Expand Up @@ -77,7 +89,7 @@ describe("EventNdjsonLogger", () => {
const serialized = encodeUnknownJson(messages);
assert.notInclude(serialized, secret);
const line = parseLogLine(
NodeFS.readFileSync(ownedLogPath(basePath, "thread-1"), "utf8").trim(),
NodeFS.readFileSync(threadLogPath(basePath, "thread-1"), "utf8").trim(),
);
assert.equal(line.payload, '{"truncated":true}');
} finally {
Expand Down Expand Up @@ -108,8 +120,8 @@ describe("EventNdjsonLogger", () => {
);
yield* logger.close();

const threadOnePath = ownedLogPath(basePath, "thread-1");
const threadTwoPath = ownedLogPath(basePath, "thread-2");
const threadOnePath = threadLogPath(basePath, "thread-1");
const threadTwoPath = threadLogPath(basePath, "thread-2");
assert.equal(NodeFS.existsSync(threadOnePath), true);
assert.equal(NodeFS.existsSync(threadTwoPath), true);

Expand All @@ -132,41 +144,39 @@ describe("EventNdjsonLogger", () => {
}),
);

it.effect(
"falls back to a global segment when orchestration thread id is missing or invalid",
() =>
Effect.gen(function* () {
const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-"));
const basePath = NodePath.join(tempDir, "provider-canonical.ndjson");

try {
const logger = yield* makeEventNdjsonLogger(basePath, { stream: "orchestration" });
assert.notEqual(logger, undefined);
if (!logger) {
return;
}

yield* logger.write({ id: "evt-no-thread" }, null);
yield* logger.write({ id: "evt-invalid-thread" }, "!!!" as unknown as ThreadId);
yield* logger.close();

const globalPath = ownedLogPath(basePath, "_global");
assert.equal(NodeFS.existsSync(globalPath), true);
const lines = NodeFS.readFileSync(globalPath, "utf8")
.trim()
.split("\n")
.map((line) => parseLogLine(line));
assert.equal(lines.length, 2);
assert.equal(Number.isNaN(Date.parse(lines[0]?.observedAt ?? "")), false);
assert.equal(Number.isNaN(Date.parse(lines[1]?.observedAt ?? "")), false);
assert.equal(lines[0]?.stream, "ORCH");
assert.equal(lines[0]?.payload, '{"id":"evt-no-thread"}');
assert.equal(lines[1]?.stream, "ORCH");
assert.equal(lines[1]?.payload, '{"id":"evt-invalid-thread"}');
} finally {
NodeFS.rmSync(tempDir, { recursive: true, force: true });
it.effect("falls back to a global segment when the orchestration thread id is missing", () =>
Effect.gen(function* () {
const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-"));
const basePath = NodePath.join(tempDir, "provider-canonical.ndjson");

try {
const logger = yield* makeEventNdjsonLogger(basePath, { stream: "orchestration" });
assert.notEqual(logger, undefined);
if (!logger) {
return;
}
}),

yield* logger.write({ id: "evt-no-thread" }, null);
yield* logger.write({ id: "evt-empty-thread" }, "" as unknown as ThreadId);
yield* logger.close();

const globalPath = ownedLogPath(basePath, "_global");
assert.equal(NodeFS.existsSync(globalPath), true);
const lines = NodeFS.readFileSync(globalPath, "utf8")
.trim()
.split("\n")
.map((line) => parseLogLine(line));
assert.equal(lines.length, 2);
assert.equal(Number.isNaN(Date.parse(lines[0]?.observedAt ?? "")), false);
assert.equal(Number.isNaN(Date.parse(lines[1]?.observedAt ?? "")), false);
assert.equal(lines[0]?.stream, "ORCH");
assert.equal(lines[0]?.payload, '{"id":"evt-no-thread"}');
assert.equal(lines[1]?.stream, "ORCH");
assert.equal(lines[1]?.payload, '{"id":"evt-empty-thread"}');
} finally {
NodeFS.rmSync(tempDir, { recursive: true, force: true });
}
}),
);

it.effect("shares one thread writer across native and canonical streams", () =>
Expand All @@ -184,7 +194,7 @@ describe("EventNdjsonLogger", () => {
yield* canonical.write({ type: "item.completed", id: "canonical-event" }, threadId);
yield* store.close();

const lines = NodeFS.readFileSync(ownedLogPath(basePath, "thread-shared"), "utf8")
const lines = NodeFS.readFileSync(threadLogPath(basePath, "thread-shared"), "utf8")
.trim()
.split("\n")
.map(parseLogLine);
Expand Down Expand Up @@ -221,7 +231,7 @@ describe("EventNdjsonLogger", () => {
yield* canonical.write({ type: "item.completed", id: "after-close" }, threadId);
yield* store.close();

const lines = NodeFS.readFileSync(ownedLogPath(basePath, "thread-shared-close"), "utf8")
const lines = NodeFS.readFileSync(threadLogPath(basePath, "thread-shared-close"), "utf8")
.trim()
.split("\n")
.map(parseLogLine);
Expand All @@ -246,7 +256,7 @@ describe("EventNdjsonLogger", () => {
Effect.gen(function* () {
const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-"));
const basePath = NodePath.join(tempDir, "events.log");
const threadPath = ownedLogPath(basePath, "thread-batched");
const threadPath = threadLogPath(basePath, "thread-batched");

try {
const store = yield* makeEventNdjsonLogStore(basePath, { batchWindowMs: 1_000 });
Expand All @@ -268,7 +278,7 @@ describe("EventNdjsonLogger", () => {
Effect.gen(function* () {
const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-"));
const basePath = NodePath.join(tempDir, "events.log");
const threadPath = ownedLogPath(basePath, "thread-interrupted");
const threadPath = threadLogPath(basePath, "thread-interrupted");

try {
const store = yield* makeEventNdjsonLogStore(basePath, { batchWindowMs: 1_000 });
Expand Down Expand Up @@ -413,7 +423,7 @@ describe("EventNdjsonLogger", () => {
yield* native.write({ type: "turn.completed", id: "native-final" }, threadId);
yield* store.close();

const lines = NodeFS.readFileSync(ownedLogPath(basePath, "thread-filtered"), "utf8")
const lines = NodeFS.readFileSync(threadLogPath(basePath, "thread-filtered"), "utf8")
.trim()
.split("\n")
.map(parseLogLine);
Expand Down Expand Up @@ -457,7 +467,7 @@ describe("EventNdjsonLogger", () => {
);
yield* store.close();

const contents = NodeFS.readFileSync(ownedLogPath(basePath, "large-history"), "utf8");
const contents = NodeFS.readFileSync(threadLogPath(basePath, "large-history"), "utf8");
assert.isBelow(Buffer.byteLength(contents), 2_048);
const record = decodeUnknownJson(parseLogLine(contents.trim()).payload);
assert.nestedPropertyVal(record, "event.payload.id", 42);
Expand Down Expand Up @@ -491,7 +501,7 @@ describe("EventNdjsonLogger", () => {
yield* logger.write({ id: "escaped", output: "\u0000".repeat(20_000) }, threadId);
yield* store.close();

const contents = NodeFS.readFileSync(ownedLogPath(basePath, "large-error"), "utf8");
const contents = NodeFS.readFileSync(threadLogPath(basePath, "large-error"), "utf8");
const records = contents
.trim()
.split("\n")
Expand Down Expand Up @@ -532,7 +542,7 @@ describe("EventNdjsonLogger", () => {
threadId,
);
yield* store.close();
const contents = NodeFS.readFileSync(ownedLogPath(basePath, "large-diff"), "utf8");
const contents = NodeFS.readFileSync(threadLogPath(basePath, "large-diff"), "utf8");
assert.isBelow(Buffer.byteLength(contents), 2_048);
const record = decodeUnknownJson(parseLogLine(contents.trim()).payload);
assert.propertyVal(record, "type", "turn.diff.updated");
Expand Down Expand Up @@ -573,7 +583,7 @@ describe("EventNdjsonLogger", () => {
yield* store.close();

const payloads = NodeFS.readFileSync(
ownedLogPath(basePath, "thread-tool-lifecycle"),
threadLogPath(basePath, "thread-tool-lifecycle"),
"utf8",
)
.trim()
Expand Down Expand Up @@ -614,7 +624,7 @@ describe("EventNdjsonLogger", () => {
yield* logger.write(hostile, ThreadId.make("thread-hostile"));
yield* logger.close();

const contents = NodeFS.readFileSync(ownedLogPath(basePath, "thread-hostile"), "utf8");
const contents = NodeFS.readFileSync(threadLogPath(basePath, "thread-hostile"), "utf8");
assert.notInclude(contents, "blocked");
} finally {
NodeFS.rmSync(tempDir, { recursive: true, force: true });
Expand Down Expand Up @@ -692,7 +702,7 @@ describe("EventNdjsonLogger", () => {
}
yield* store.close();

const fileStem = NodePath.basename(ownedLogPath(basePath, "thread-rotate"));
const fileStem = NodePath.basename(threadLogPath(basePath, "thread-rotate"));
const matchingFiles = NodeFS.readdirSync(tempDir)
.filter((entry) => entry === fileStem || entry.startsWith(`${fileStem}.`))
.toSorted();
Expand All @@ -719,9 +729,9 @@ describe("EventNdjsonLogger", () => {
Effect.gen(function* () {
const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-"));
const basePath = NodePath.join(tempDir, "events.log");
const expiredPath = ownedLogPath(basePath, "expired");
const oldPath = ownedLogPath(basePath, "old");
const newPath = ownedLogPath(basePath, "new");
const expiredPath = threadLogPath(basePath, "expired");
const oldPath = threadLogPath(basePath, "old");
const newPath = threadLogPath(basePath, "new");
const unrelatedLogPath = NodePath.join(tempDir, "unrelated.log");
const legacyLogPath = NodePath.join(tempDir, "legacy-thread.log");
const ignoredPath = NodePath.join(tempDir, "ignored.txt");
Expand Down Expand Up @@ -763,7 +773,7 @@ describe("EventNdjsonLogger", () => {
Effect.gen(function* () {
const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-"));
const basePath = NodePath.join(tempDir, "events.log");
const activePath = ownedLogPath(basePath, "active");
const activePath = threadLogPath(basePath, "active");

try {
yield* TestClock.setTime(1_800_000_000_000);
Expand Down Expand Up @@ -822,6 +832,93 @@ describe("EventNdjsonLogger", () => {
assert.deepEqual(attributed, [records[0]]);
});

it.effect("keeps delegated threads that share a long id prefix in separate files", () =>
Effect.gen(function* () {
const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-"));
const basePath = NodePath.join(tempDir, "events.log");
const first =
"thread:delegated-task:command%3Amcp%3Ad41693b3-cbc4-439b-9911-f51c5a0ba5b3%3Adelegate-task%3Amimo-review-live-driver-effort-20260929";
const retry = `${first}-retry1`;

try {
const logger = yield* makeEventNdjsonLogger(basePath, { stream: "native" });
assert.notEqual(logger, undefined);
if (!logger) return;
yield* logger.write({ id: "first-run" }, ThreadId.make(first));
yield* logger.write({ id: "retry-run" }, ThreadId.make(retry));
yield* logger.write({ id: "first-run-late" }, ThreadId.make(first));
yield* logger.close();

const files = logFileNames(tempDir);
assert.lengthOf(files, 2, "each delegated thread owns its own log file");
const contents = files.map((name) =>
NodeFS.readFileSync(NodePath.join(tempDir, name), "utf8")
.trim()
.split("\n")
.map((line) => parseLogLine(line).payload),
);
assert.sameDeepMembers(contents, [
['{"id":"first-run"}', '{"id":"first-run-late"}'],
['{"id":"retry-run"}'],
]);
} finally {
NodeFS.rmSync(tempDir, { recursive: true, force: true });
}
}),
);

it.effect("keeps thread ids that normalize alike in separate files", () =>
Effect.gen(function* () {
const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-"));
const basePath = NodePath.join(tempDir, "events.log");
const ids = [
"Thread-Case",
"thread-case",
"thread.dot",
"thread-dot",
`thread-${"\u2603".repeat(200)}`,
`thread-${"\u2600".repeat(200)}`,
"!!!",
"???",
];

try {
const logger = yield* makeEventNdjsonLogger(basePath, { stream: "native" });
assert.notEqual(logger, undefined);
if (!logger) return;
for (const id of ids) {
yield* logger.write({ id }, id as ThreadId);
}
yield* logger.close();

const files = logFileNames(tempDir);
assert.lengthOf(files, ids.length);
assert.notInclude(files, "events._global.log", "a supplied id is never global");
for (const name of files) {
const payloads = NodeFS.readFileSync(NodePath.join(tempDir, name), "utf8")
.trim()
.split("\n");
assert.lengthOf(payloads, 1, `${name} holds one thread's records`);
}
} finally {
NodeFS.rmSync(tempDir, { recursive: true, force: true });
}
}),
);

it("names thread files deterministically with a bounded length", () => {
const hash = NodeCrypto.createHash("sha256").update("thread-1", "utf8").digest("hex");
assert.equal(providerEventLogThreadSegment("thread-1"), `thread-1-${hash.slice(0, 32)}`);
assert.equal(
providerEventLogThreadSegment("thread-1"),
providerEventLogThreadSegment("thread-1"),
);
const long = providerEventLogThreadSegment(`thread-${"x\u00e9".repeat(5_000)}`);
assert.isAtMost(long.length, 100);
assert.match(long, /^[a-z0-9_-]+$/);
assert.match(providerEventLogThreadSegment("!!!"), /^thread-[0-9a-f]{32}$/);
});

it.effect("reports logical provider log writes to resource attribution", () =>
Effect.gen(function* () {
const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-"));
Expand Down
29 changes: 25 additions & 4 deletions apps/server/src/provider/EventNdjsonLogger.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
* Native and canonical views share batching, rotation, and retention state so
* they cannot race while appending to the same thread-scoped file.
*/
import * as NodeCrypto from "node:crypto";
import * as NodeFS from "node:fs";
import * as NodePath from "node:path";

Expand Down Expand Up @@ -37,6 +38,8 @@ const MAX_RECORD_CHARACTERS = 64 * 1024;
const MAX_RECORD_FIELDS = 1_024;
const MAX_RECORD_DEPTH = 16;
const GLOBAL_THREAD_SEGMENT = "_global";
const THREAD_SEGMENT_PREFIX_MAX_CHARS = 64;
const THREAD_SEGMENT_HASH_HEX_CHARS = 32;
const LOG_SCOPE = "provider-observability";
const encodeUnknownJsonString = Schema.encodeUnknownEffect(Schema.fromJsonString(Schema.Unknown));

Expand Down Expand Up @@ -168,9 +171,22 @@ function logWarning(message: string, context: Record<string, unknown>): Effect.E
return Effect.logWarning(message, context).pipe(Effect.annotateLogs({ scope: LOG_SCOPE }));
}

function resolveThreadSegment(raw: string | null | undefined): string {
const normalized = typeof raw === "string" ? toSafeThreadAttachmentSegment(raw) : null;
return normalized ?? GLOBAL_THREAD_SEGMENT;
/**
* File-name segment for one thread's provider log. The readable prefix is
* lossy (case, punctuation and length), so a hash of the complete raw id keeps
* threads that normalize alike, such as retries of one delegated task, in
* separate files. Records without a thread share the global file.
*/
export function providerEventLogThreadSegment(threadId: string | null | undefined): string {
if (typeof threadId !== "string" || threadId.length === 0) return GLOBAL_THREAD_SEGMENT;
const readable = (toSafeThreadAttachmentSegment(threadId) ?? "thread")
.slice(0, THREAD_SEGMENT_PREFIX_MAX_CHARS)
.replace(/[-_]+$/g, "");
const hash = NodeCrypto.createHash("sha256")
.update(threadId, "utf8")
.digest("hex")
.slice(0, THREAD_SEGMENT_HASH_HEX_CHARS);
return `${readable.length > 0 ? readable : "thread"}-${hash}`;
}

function resolveStreamLabel(stream: EventNdjsonStream): string {
Expand Down Expand Up @@ -753,7 +769,12 @@ export const makeEventNdjsonLogStore = Effect.fnUntraced(function* (
return Effect.succeed([{ flush: false }, state] as const);
}
const pending = state.pending;
pending.push({ stream, threadSegment: resolveThreadSegment(threadId), line, bytes });
pending.push({
stream,
threadSegment: providerEventLogThreadSegment(threadId),
line,
bytes,
});
const pendingBytes = state.pendingBytes + bytes;
const flush =
resolved.batchWindowMs === 0 ||
Expand Down