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
10 changes: 7 additions & 3 deletions docs/operations/observability.md
Original file line number Diff line number Diff line change
Expand Up @@ -50,8 +50,10 @@ Important fields common to both record types:
- `attributes`: structured context
- `events`: embedded logs and custom events

`effect-span` records also contain `exit` with `Success`, `Failure`, or `Interrupted`. `otlp-span`
records instead carry OTLP resource, scope, and optional status fields.
`effect-span` records also contain `exit` with `Success`, `Failure`, or `Interrupted`. Their
attribute strings and event names keep the first 500 characters (`db.query.text` keeps 200). A
failure's `cause` keeps the first 8,000 characters, which fits a normal stack and `[cause]` chain.
`otlp-span` records instead carry OTLP resource, scope, and optional status fields.

The `TraceRecord`, `EffectTraceRecord`, and `OtlpTraceRecord` schemas live in
`packages/shared/src/observability.ts`.
Expand Down Expand Up @@ -533,7 +535,9 @@ yield * Effect.logInfo("starting provider turn");
yield * Effect.logDebug("waiting for approval response");
```

Those messages show up as span events because `Logger.tracerLogger` is installed.
Those messages show up as span events because `Logger.tracerLogger` is installed. A span keeps
its newest 128 events and counts the older ones it drops in its `span.dropped_events_count`
attribute. Logs from a fiber that outlives its span are not written.

### Use The Pipeable Metrics API

Expand Down
101 changes: 100 additions & 1 deletion packages/shared/src/observability.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ const TraceRecordLine = Schema.Struct({
exit: Schema.optional(
Schema.Struct({
_tag: Schema.String,
cause: Schema.optional(Schema.String),
}),
),
});
Expand Down Expand Up @@ -97,7 +98,7 @@ const readTraceRecords = Effect.fn("readTraceRecords")(function* (tracePath: str
.map((line) => decodeTraceRecordLine(line));
});

const makeTestLayer = (tracePath: string) =>
const makeTestLayer = (tracePath: string, delegate?: Tracer.Tracer) =>
Layer.mergeAll(
Layer.effect(
Tracer.Tracer,
Expand All @@ -106,12 +107,27 @@ const makeTestLayer = (tracePath: string) =>
maxBytes: 1024 * 1024,
maxFiles: 2,
batchWindowMs: 10_000,
...(delegate ? { delegate } : {}),
}),
),
Logger.layer([Logger.tracerLogger], { mergeWithExisting: false }),
Layer.succeed(References.MinimumLogLevel, "Info"),
);

// A delegate tracer that keeps its spans, to see what the local tracer
// forwards to it.
const makeRecordingDelegate = () => {
const spans: Array<Tracer.NativeSpan> = [];
const delegate = Tracer.make({
span: (options) => {
const span = new Tracer.NativeSpan(options);
spans.push(span);
return span;
},
});
return { delegate, spans };
};

const nodeServicesIt = it.layer(NodeServices.layer);

describe("truncateTraceAttributes", () => {
Expand Down Expand Up @@ -506,6 +522,89 @@ describe("observability", () => {
}),
),
);

it.effect("clamps oversized failure causes and log event names", () =>
Effect.scoped(
Effect.gen(function* () {
const fileSystem = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const tempDir = yield* fileSystem.makeTempDirectoryScoped({ prefix: "t3-local-tracer-" });
const tracePath = path.join(tempDir, "shared.trace.ndjson");

yield* Effect.scoped(
Effect.exit(
Effect.logWarning("w".repeat(5_000)).pipe(
Effect.andThen(Effect.fail("f".repeat(20_000))),
Effect.withSpan("failing-span"),
Effect.provide(makeTestLayer(tracePath)),
),
),
);

const records = yield* readTraceRecords(tracePath);
assert.equal(records[0]?.events[0]?.name, `${"w".repeat(500)}鈥truncated]`);
const cause = records[0]?.exit?.cause ?? "";
assert.equal(cause.length, 8_000 + "鈥truncated]".length);
assert.isTrue(cause.endsWith("鈥truncated]"));
}),
),
);

it.effect("ignores events added after a span has ended", () =>
Effect.scoped(
Effect.gen(function* () {
const fileSystem = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const tempDir = yield* fileSystem.makeTempDirectoryScoped({ prefix: "t3-local-tracer-" });
const tracePath = path.join(tempDir, "shared.trace.ndjson");

const span = yield* Effect.scoped(
Effect.currentSpan.pipe(
Effect.tap((span) => Effect.sync(() => span.event("during", 1n))),
Effect.withSpan("ended-span"),
Effect.provide(makeTestLayer(tracePath)),
),
);
span.event("after end", 2n);

// A fiber that outlives its span would grow this list.
assert.deepPropertyVal(span, "events", [["during", 1n, {}]]);
}),
),
);

it.effect("keeps the newest events per span and records how many were dropped", () =>
Effect.scoped(
Effect.gen(function* () {
const fileSystem = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const tempDir = yield* fileSystem.makeTempDirectoryScoped({ prefix: "t3-local-tracer-" });
const tracePath = path.join(tempDir, "shared.trace.ndjson");
const { delegate, spans } = makeRecordingDelegate();

yield* Effect.scoped(
Effect.forEach(Arr.range(1, 200), (index) => Effect.logInfo(`event ${index}`), {
discard: true,
}).pipe(
Effect.withSpan("chatty-span"),
Effect.provide(makeTestLayer(tracePath, delegate)),
),
);

const records = yield* readTraceRecords(tracePath);
const names = records[0]?.events.map((event) => event.name);
assert.deepEqual(
names,
Arr.range(73, 200).map((index) => `event ${index}`),
);
assert.equal(records[0]?.attributes["span.dropped_events_count"], 72);
assert.deepEqual(
spans[0]?.events.map(([name]) => name),
names,
);
}),
),
);
});
});

Expand Down
50 changes: 39 additions & 11 deletions packages/shared/src/observability.ts
Original file line number Diff line number Diff line change
Expand Up @@ -259,19 +259,25 @@ export function compactTraceAttributes(
return Object.fromEntries(entries);
}

// A failure repeats its pretty cause on every failing ancestor span. Normal
// causes (message, stack, and [cause] chain) are well under the cap; it only
// stops a pathological one, such as an error that embeds a large payload.
const TRACE_CAUSE_MAX_LENGTH = 8_000;
Comment thread
coderabbitai[bot] marked this conversation as resolved.

function formatTraceExit(exit: Exit.Exit<unknown, unknown>): EffectTraceRecord["exit"] {
if (ExitRuntime.isSuccess(exit)) {
return { _tag: "Success" };
}
const cause = truncateTraceString(Cause.pretty(exit.cause), TRACE_CAUSE_MAX_LENGTH);
if (Cause.hasInterruptsOnly(exit.cause)) {
return {
_tag: "Interrupted",
cause: Cause.pretty(exit.cause),
cause,
};
}
return {
_tag: "Failure",
cause: Cause.pretty(exit.cause),
cause,
};
}

Expand All @@ -280,14 +286,18 @@ const TRACE_ATTRIBUTE_TRUNCATED_LENGTH = 200;
const TRACE_ATTRIBUTE_TRUNCATION_SUFFIX = "鈥truncated]";
const ALWAYS_TRUNCATED_TRACE_ATTRIBUTES: ReadonlySet<string> = new Set(["db.query.text"]);

function truncateTraceString(value: string, maxLength: number): string {
return value.length <= maxLength
? value
: `${value.slice(0, maxLength)}${TRACE_ATTRIBUTE_TRUNCATION_SUFFIX}`;
}

// Clamps strings nested inside already-normalized attribute values (arrays and
// plain objects from normalizeJsonValue, e.g. an Error's `stack`). Returns the
// input reference when nothing was clamped.
function truncateNestedValue(value: unknown): unknown {
if (typeof value === "string") {
return value.length <= TRACE_ATTRIBUTE_MAX_LENGTH
? value
: `${value.slice(0, TRACE_ATTRIBUTE_MAX_LENGTH)}${TRACE_ATTRIBUTE_TRUNCATION_SUFFIX}`;
return truncateTraceString(value, TRACE_ATTRIBUTE_MAX_LENGTH);
}
if (Array.isArray(value)) {
const truncated = value.map(truncateNestedValue);
Expand Down Expand Up @@ -318,8 +328,7 @@ export function truncateTraceAttributes(attributes: TraceAttributes): TraceAttri
if (typeof value === "string" && ALWAYS_TRUNCATED_TRACE_ATTRIBUTES.has(key)) {
if (value.length <= TRACE_ATTRIBUTE_TRUNCATED_LENGTH) continue;
truncated ??= { ...attributes };
truncated[key] =
`${value.slice(0, TRACE_ATTRIBUTE_TRUNCATED_LENGTH)}${TRACE_ATTRIBUTE_TRUNCATION_SUFFIX}`;
truncated[key] = truncateTraceString(value, TRACE_ATTRIBUTE_TRUNCATED_LENGTH);
continue;
}
const next = truncateNestedValue(value);
Expand Down Expand Up @@ -349,7 +358,8 @@ function spanToTraceRecord(span: SerializableSpan): EffectTraceRecord {
compactTraceAttributes(Object.fromEntries(span.attributes)),
),
events: span.events.map(([name, startTime, attributes]) => ({
name,
// A log event is named after its whole formatted message.
name: truncateTraceString(name, TRACE_ATTRIBUTE_MAX_LENGTH),
timeUnixNano: String(startTime),
attributes: truncateTraceAttributes(compactTraceAttributes(attributes)),
})),
Expand Down Expand Up @@ -457,6 +467,11 @@ export const makeTraceSink = Effect.fn("makeTraceSink")(function* (options: Trac
} satisfies TraceSink;
});

// Long-lived spans (a whole RPC stream, for example) would otherwise keep
// every log event in memory until they end. Like the OpenTelemetry SDK, a
// span keeps its newest 128 events and drops the oldest.
const TRACE_SPAN_MAX_EVENTS = 128;
Comment thread
coderabbitai[bot] marked this conversation as resolved.

class LocalFileSpan implements Tracer.Span {
readonly _tag = "Span";
readonly name: string;
Expand All @@ -471,6 +486,7 @@ class LocalFileSpan implements Tracer.Span {
status: Tracer.SpanStatus;
attributes: Map<string, unknown>;
events: Array<[name: string, startTime: bigint, attributes: Record<string, unknown>]>;
private droppedEventCount = 0;
private readonly delegate: Tracer.Span;
private readonly push: (record: EffectTraceRecord) => void;

Expand Down Expand Up @@ -498,12 +514,20 @@ class LocalFileSpan implements Tracer.Span {
}

end(endTime: bigint, exit: Exit.Exit<unknown, unknown>): void {
if (this.droppedEventCount > 0) {
this.attribute("span.dropped_events_count", this.droppedEventCount);
}
this.status = {
_tag: "Ended",
startTime: this.status.startTime,
endTime,
exit,
};
// The delegate gets the kept events only now, so it holds the same bounded
// set. The OTLP delegate reads a span's events only at end.
for (const [name, startTime, attributes] of this.events) {
this.delegate.event(name, startTime, attributes);
}
this.delegate.end(endTime, exit);

if (this.sampled) {
Expand All @@ -516,10 +540,14 @@ class LocalFileSpan implements Tracer.Span {
this.delegate.attribute(key, value);
}

// Events on an ended span are never written, so they are not kept either.
event(name: string, startTime: bigint, attributes?: Record<string, unknown>): void {
const nextAttributes = attributes ?? {};
this.events.push([name, startTime, nextAttributes]);
this.delegate.event(name, startTime, nextAttributes);
if (this.status._tag === "Ended") return;
if (this.events.length >= TRACE_SPAN_MAX_EVENTS) {
this.events.shift();
this.droppedEventCount += 1;
}
this.events.push([name, startTime, attributes ?? {}]);
}

addLinks(links: ReadonlyArray<Tracer.SpanLink>): void {
Expand Down
Loading