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
36 changes: 36 additions & 0 deletions apps/server/src/provider/Layers/EventNdjsonLogger.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,13 @@ import * as NodePath from "node:path";
import { ThreadId } from "@t3tools/contracts";
import { assert, describe, it } from "@effect/vitest";
import * as Effect from "effect/Effect";
import * as Logger from "effect/Logger";
import * as Schema from "effect/Schema";

import { makeEventNdjsonLogger } from "./EventNdjsonLogger.ts";

const encodeUnknownJson = Schema.encodeUnknownSync(Schema.UnknownFromJsonString);

function parseLogLine(line: string) {
const match = /^\[([^\]]+)\] ([A-Z]+): (.+)$/.exec(line);
assert.notEqual(match, null);
Expand All @@ -29,6 +33,38 @@ function parseLogLine(line: string) {
}

describe("EventNdjsonLogger", () => {
it.effect("logs bounded diagnostics when an event cannot be serialized", () => {
const messages: Array<unknown> = [];
const logCapture = Logger.make<unknown, void>(({ message }) => {
if (Array.isArray(message)) {
messages.push(...message);
} else {
messages.push(message);
}
});
const secret = "secret-circular-event-value";

return Effect.gen(function* () {
const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-"));
const basePath = NodePath.join(tempDir, "provider-native.ndjson");
const circular: Record<string, unknown> = { secret };
circular.self = circular;

try {
const logger = yield* makeEventNdjsonLogger(basePath, { stream: "native" });
assert.exists(logger);
if (!logger) return;
yield* logger.write(circular, ThreadId.make("thread-1"));

const serialized = encodeUnknownJson(messages);
assert.notInclude(serialized, secret);
assert.include(serialized, '"errorTag":"SchemaError"');
} finally {
NodeFS.rmSync(tempDir, { recursive: true, force: true });
}
}).pipe(Effect.provide(Logger.layer([logCapture], { mergeWithExisting: false })));
});

it.effect("writes effect-style lines to thread-scoped files", () =>
Effect.gen(function* () {
const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-log-"));
Expand Down
17 changes: 9 additions & 8 deletions apps/server/src/provider/Layers/EventNdjsonLogger.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import * as NodePath from "node:path";

import type { ThreadId } from "@t3tools/contracts";
import { RotatingFileSink } from "@t3tools/shared/logging";
import { errorTag } from "@t3tools/shared/observability";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Logger from "effect/Logger";
Expand All @@ -31,8 +32,8 @@ export type EventNdjsonStream = "native" | "canonical" | "orchestration";

export interface EventNdjsonLogger {
readonly filePath: string;
write: (event: unknown, threadId: ThreadId | null) => Effect.Effect<void>;
close: () => Effect.Effect<void>;
write: (event: unknown, threadId: ThreadId | null) => Effect.Effect<void, never, never>;
close: () => Effect.Effect<void, never, never>;
}

export interface EventNdjsonLoggerOptions {
Expand Down Expand Up @@ -91,9 +92,9 @@ const toLogMessage = Effect.fn("toLogMessage")(function* (
): Effect.fn.Return<string | undefined> {
return yield* encodeUnknownJsonString(event).pipe(
Effect.catch((error) =>
logWarning("failed to serialize provider event log record", { error }).pipe(
Effect.as(undefined),
),
logWarning("failed to serialize provider event log record", {
errorTag: errorTag(error),
}).pipe(Effect.as(undefined)),
),
);
});
Expand Down Expand Up @@ -124,7 +125,7 @@ const makeThreadWriter = Effect.fn("makeThreadWriter")(function* (input: {
if (!sinkResult.ok) {
yield* logWarning("failed to initialize provider thread log file", {
filePath: input.filePath,
error: sinkResult.error,
errorTag: errorTag(sinkResult.error),
});
return undefined;
}
Expand All @@ -149,7 +150,7 @@ const makeThreadWriter = Effect.fn("makeThreadWriter")(function* (input: {
if (!flushResult.ok) {
yield* logWarning("provider event log batch flush failed", {
filePath: input.filePath,
error: flushResult.error,
errorTag: errorTag(flushResult.error),
});
}
}),
Expand Down Expand Up @@ -187,7 +188,7 @@ export const makeEventNdjsonLogger = Effect.fn("makeEventNdjsonLogger")(function
if (directoryReady !== true) {
yield* logWarning("failed to create provider event log directory", {
filePath,
error: directoryReady.error,
errorTag: errorTag(directoryReady.error),
});
return undefined;
}
Expand Down
137 changes: 137 additions & 0 deletions apps/server/src/provider/acp/AcpNativeLogging.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,137 @@
import * as NodeServices from "@effect/platform-node/NodeServices";
import { ProviderDriverKind, ThreadId } from "@t3tools/contracts";
import { assert, it } from "@effect/vitest";
import * as Cause from "effect/Cause";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Logger from "effect/Logger";
import * as Schema from "effect/Schema";
import * as AcpErrors from "effect-acp/errors";

import type { EventNdjsonLogger } from "../Layers/EventNdjsonLogger.ts";
import { makeAcpNativeLoggerFactory } from "./AcpNativeLogging.ts";

const nodeServicesIt = it.layer(NodeServices.layer);
const encodeUnknownJson = Schema.encodeUnknownSync(Schema.UnknownFromJsonString);

nodeServicesIt("ACP native logging", (it) => {
it.effect("records bounded request and protocol diagnostics without raw payloads", () =>
Effect.gen(function* () {
const records: Array<unknown> = [];
const nativeEventLogger: EventNdjsonLogger = {
filePath: "/tmp/provider-native.ndjson",
write: (event) => Effect.sync(() => void records.push(event)),
close: () => Effect.void,
};
const makeLogger = yield* makeAcpNativeLoggerFactory();
const logger = makeLogger({
nativeEventLogger,
provider: ProviderDriverKind.make("cursor"),
threadId: ThreadId.make("thread-1"),
});
const secret = "secret-token-value";
const requestLogger = logger.requestLogger;
const protocolLogger = logger.protocolLogging?.logger;
assert.exists(requestLogger);
assert.exists(protocolLogger);
if (!requestLogger || !protocolLogger) return;

yield* requestLogger({
method: "session/prompt",
payload: { prompt: secret, sessionId: secret },
status: "failed",
cause: Cause.fail(AcpErrors.AcpRequestError.internalError(secret, { token: secret })),
});
yield* protocolLogger({
direction: "incoming",
stage: "raw",
payload: `{"token":"${secret}"}`,
});
yield* protocolLogger({
direction: "outgoing",
stage: "decoded",
payload: {
_tag: "Request",
tag: "session/prompt",
payload: { prompt: secret },
},
});

const serialized = encodeUnknownJson(records);
assert.notInclude(serialized, secret);
assert.include(serialized, '"method":"session/prompt"');
assert.include(serialized, '"errorTag":"AcpRequestError"');
assert.include(serialized, '"reasonCount":1');
assert.include(serialized, '"valueType":"string"');
assert.include(serialized, '"messageTag":"Request"');
}),
);

it.effect("logs a structural tag when the native writer defects", () => {
const messages: Array<unknown> = [];
const logCapture = Logger.make<unknown, void>(({ message }) => {
if (Array.isArray(message)) {
messages.push(...message);
} else {
messages.push(message);
}
});
const secret = "secret-writer-failure";

return Effect.gen(function* () {
const makeLogger = yield* makeAcpNativeLoggerFactory();
const logger = makeLogger({
nativeEventLogger: {
filePath: "/tmp/provider-native.ndjson",
write: () => Effect.die(new Error(secret)),
close: () => Effect.void,
},
provider: ProviderDriverKind.make("cursor"),
threadId: ThreadId.make("thread-1"),
});
const requestLogger = logger.requestLogger;
assert.exists(requestLogger);
if (!requestLogger) return;

yield* requestLogger({
method: "session/prompt",
payload: {},
status: "started",
});

const serialized = encodeUnknownJson(messages);
assert.notInclude(serialized, secret);
assert.include(serialized, '"errorTag":"Die"');
assert.include(serialized, '"reasonCount":1');
}).pipe(Effect.provide(Logger.layer([logCapture], { mergeWithExisting: false })));
});

it.effect("preserves native writer interruption", () =>
Effect.gen(function* () {
const makeLogger = yield* makeAcpNativeLoggerFactory();
const logger = makeLogger({
nativeEventLogger: {
filePath: "/tmp/provider-native.ndjson",
write: () => Effect.interrupt,
close: () => Effect.void,
},
provider: ProviderDriverKind.make("cursor"),
threadId: ThreadId.make("thread-1"),
});
const requestLogger = logger.requestLogger;
assert.exists(requestLogger);
if (!requestLogger) return;

const exit = yield* requestLogger({
method: "session/prompt",
payload: {},
status: "started",
}).pipe(Effect.exit);

assert.isTrue(Exit.isFailure(exit));
if (Exit.isFailure(exit)) {
assert.isTrue(Cause.hasInterruptsOnly(exit.cause));
}
}),
);
});
71 changes: 60 additions & 11 deletions apps/server/src/provider/acp/AcpNativeLogging.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import type { ProviderDriverKind, ThreadId } from "@t3tools/contracts";
import { causeErrorTag, errorTag } from "@t3tools/shared/observability";
import * as Cause from "effect/Cause";
import * as Crypto from "effect/Crypto";
import * as DateTime from "effect/DateTime";
Expand All @@ -8,13 +9,58 @@ import type * as EffectAcpProtocol from "effect-acp/protocol";
import type { EventNdjsonLogger } from "../Layers/EventNdjsonLogger.ts";
import type * as AcpSessionRuntime from "./AcpSessionRuntime.ts";

function structuralMethod(value: string): string {
return value.length <= 128 && /^[A-Za-z][A-Za-z0-9._:/-]*$/.test(value) ? value : "unknown";
}

function summarizePayload(payload: unknown): Readonly<Record<string, unknown>> {
if (payload === null) return { valueType: "null" };
if (typeof payload === "string") {
return { valueType: "string", byteLength: new TextEncoder().encode(payload).byteLength };
}
if (payload instanceof Uint8Array) {
return { valueType: "bytes", byteLength: payload.byteLength };
}
if (Array.isArray(payload)) {
return { valueType: "array", itemCount: payload.length };
}
if (typeof payload !== "object") {
return { valueType: typeof payload };
}

try {
const record = payload as Record<string, unknown>;
return {
valueType: "object",
fieldCount: Object.keys(record).length,
...(typeof record._tag === "string" ? { messageTag: errorTag(record) } : {}),
...(typeof record.tag === "string" ? { method: structuralMethod(record.tag) } : {}),
};
} catch {
return { valueType: "object" };
}
}

function formatRequestLogPayload(event: AcpSessionRuntime.AcpSessionRequestLogEvent) {
return {
method: event.method,
method: structuralMethod(event.method),
status: event.status,
request: event.payload,
...(event.result !== undefined ? { result: event.result } : {}),
...(event.cause !== undefined ? { cause: Cause.pretty(event.cause) } : {}),
request: summarizePayload(event.payload),
...(event.result !== undefined ? { result: summarizePayload(event.result) } : {}),
...(event.cause !== undefined
? {
errorTag: causeErrorTag(event.cause),
reasonCount: event.cause.reasons.length,
}
: {}),
};
}

function formatProtocolLogPayload(event: EffectAcpProtocol.AcpProtocolLogEvent) {
return {
direction: event.direction,
stage: event.stage,
payload: summarizePayload(event.payload),
};
}

Expand Down Expand Up @@ -47,12 +93,15 @@ export const makeAcpNativeLoggerFactory = Effect.fn("makeAcpNativeLoggerFactory"
input.threadId,
);
}).pipe(
Effect.catch((cause) =>
Effect.logWarning("Failed to write native ACP event log.", {
cause,
provider: input.provider,
threadId: input.threadId,
}),
Effect.catchCause((cause) =>
Cause.hasInterrupts(cause)
? Effect.interrupt
: Effect.logWarning("Failed to write native ACP event log.", {
errorTag: causeErrorTag(cause),
reasonCount: cause.reasons.length,
provider: input.provider,
threadId: input.threadId,
}),
),
);

Expand All @@ -70,7 +119,7 @@ export const makeAcpNativeLoggerFactory = Effect.fn("makeAcpNativeLoggerFactory"
logger: (event: EffectAcpProtocol.AcpProtocolLogEvent) =>
writeNativeAcpLog({
kind: "protocol",
payload: event,
payload: formatProtocolLogPayload(event),
}),
} satisfies NonNullable<AcpSessionRuntime.AcpSessionRuntimeOptions["protocolLogging"]>,
}
Expand Down
9 changes: 9 additions & 0 deletions packages/shared/src/observability.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,11 +15,20 @@ import * as Tracer from "effect/Tracer";
import {
causeErrorTag,
compactTraceAttributes,
errorTag,
makeLocalFileTracer,
makeTraceSink,
type TraceRecord,
} from "./observability.ts";

describe("errorTag", () => {
it("reports structural tags without retaining arbitrary values", () => {
assert.equal(errorTag({ _tag: "AcpRequestError" }), "AcpRequestError");
assert.equal(errorTag(new TypeError("secret-token-value")), "TypeError");
assert.equal(errorTag({ _tag: "secret token value" }), "TaggedError");
});
});

describe("causeErrorTag", () => {
it("reports the tagged failure value instead of the Cause reason wrapper", () => {
assert.equal(
Expand Down
Loading
Loading