Skip to content
Merged
14 changes: 14 additions & 0 deletions CHANGES.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,20 @@ Version 2.5.0

To be released.

### @fedify/fedify

- Fixed log entries emitted during an individual outbound activity delivery
or inbound activity processing carrying the enclosing HTTP request's or
queue worker's `traceId`/`spanId` instead of the delivery or processing
operation's own. Warning and error logs now match the `traceId`/`spanId`
of the `TraceActivityRecord` that operation produces (as exposed by
`@fedify/fedify/otel`), including when the operation runs inside a
background queue worker, which makes it possible to correlate a failure
log with the specific activity it belongs to. [[#1030], [#1205] by u-zzn\]

[#1030]: https://github.com/fedify-dev/fedify/issues/1030
[#1205]: https://github.com/fedify-dev/fedify/pull/1205

### @fedify/cli

- Changed the `fedify nodeinfo` favicon selector to pick a usable bitmap
Expand Down
13 changes: 13 additions & 0 deletions changes.d/fedify/align-log-context-with-activity-spans.md
Comment thread
dahlia marked this conversation as resolved.
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
---
links:
'#1030': https://github.com/fedify-dev/fedify/issues/1030
'#1205': https://github.com/fedify-dev/fedify/pull/1205
---
- Fixed log entries emitted during an individual outbound activity delivery
or inbound activity processing carrying the enclosing HTTP request's or
queue worker's `traceId`/`spanId` instead of the delivery or processing
operation's own. Warning and error logs now match the `traceId`/`spanId`
of the `TraceActivityRecord` that operation produces (as exposed by
`@fedify/fedify/otel`), including when the operation runs inside a
background queue worker, which makes it possible to correlate a failure
log with the specific activity it belongs to. [[#1030], [#1205] by u-zzn]
330 changes: 330 additions & 0 deletions packages/fedify/src/federation/handler.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,8 +21,18 @@ import {
assertEquals,
assertGreaterOrEqual,
assertInstanceOf,
assertNotEquals,
assertRejects,
} from "@std/assert";
import {
configure,
getLogger,
type LogRecord,
reset,
withContext,
} from "@logtape/logtape";
import { AsyncLocalStorage } from "node:async_hooks";
import { FedifySpanExporter } from "../otel/exporter.ts";
import { parseAcceptSignature } from "../sig/accept.ts";
import { signRequest } from "../sig/http.ts";
import { generateCryptoKeyPair } from "../sig/key.ts";
Expand Down Expand Up @@ -4938,6 +4948,326 @@ test("handleInbox() records OpenTelemetry span events", async () => {
);
});

test("handleInbox() aligns the LogTape context with its own span", async (t) => {
async function buildSignedRequest(activityId: string, noteId: string) {
const activity = new Create({
id: new URL(activityId),
actor: new URL("https://example.com/users/someone"),
object: new Note({
id: new URL(noteId),
content: "Hello, world!",
}),
});
const request = new Request("https://example.com/users/someone/inbox", {
method: "POST",
headers: { "Content-Type": "application/activity+json" },
body: JSON.stringify(await activity.toJsonLd()),
});
return await signRequest(
request,
rsaPrivateKey3,
new URL("https://example.com/users/someone#main-key"),
);
}

const actorDispatcher: ActorDispatcher<void> = (ctx, identifier) => {
if (identifier !== "someone") return null;
return new Person({
id: ctx.getActorUri(identifier),
name: "Someone",
inbox: new URL("https://example.com/users/someone/inbox"),
publicKey: rsaPublicKey2,
});
};

async function callHandleInbox(
federation: ReturnType<typeof createFederation<void>>,
kv: MemoryKvStore,
tracerProvider: ReturnType<typeof createTestTracerProvider>[0],
listeners: ActivityListenerSet<InboxContext<void>>,
signed: Request,
) {
const context = createRequestContext<void>({
federation,
request: signed,
url: new URL(signed.url),
data: undefined,
documentLoader: mockDocumentLoader,
contextLoader: mockDocumentLoader,
getActorUri(identifier: string) {
return new URL(`https://example.com/users/${identifier}`);
},
});
return await handleInbox(signed, {
recipient: "someone",
context,
inboxContextFactory(_activity) {
return createInboxContext({ ...context, clone: undefined });
},
kv,
kvPrefixes: {
activityIdempotence: ["activityIdempotence"],
publicKey: ["publicKey"],
acceptSignatureNonce: ["acceptSignatureNonce"],
},
actorDispatcher,
inboxListeners: listeners,
inboxErrorHandler: undefined,
onNotFound: (_request) => new Response("Not found", { status: 404 }),
signatureTimeWindow: false,
skipSignatureVerification: true,
tracerProvider,
});
}

const records: LogRecord[] = [];
await reset();
try {
await configure({
sinks: { buffer: (record: LogRecord) => records.push(record) },
filters: {},
loggers: [
{ category: [], sinks: ["buffer"], lowestLevel: "debug" },
{ category: ["logtape", "meta"], sinks: [] },
],
contextLocalStorage: new AsyncLocalStorage(),
});

const outerContext = {
requestId: "outer-request",
traceId: "outer0000000000000000000000sentinel",
spanId: "outer00000sentinel",
};

await t.step(
"a listener failure log matches the inbox's own TraceActivityRecord, " +
"and the outer context is restored afterward",
async () => {
const [tracerProvider, exporter] = createTestTracerProvider();
try {
const kv = new MemoryKvStore();
const federation = createFederation<void>({ kv, tracerProvider });
const listeners = new ActivityListenerSet<InboxContext<void>>();
listeners.add(Create, () => {
throw new Error("listener boom");
});

const signed = await buildSignedRequest(
"https://example.com/activity/listener-failure",
"https://example.com/note/listener-failure",
);

await withContext(outerContext, async () => {
const response = await callHandleInbox(
federation,
kv,
tracerProvider,
listeners,
signed,
);
assertEquals(response.status, 500);
getLogger(["fedify", "federation", "handler.test"]).info(
"after handled failure",
);
});

const span = exporter.getSpan("activitypub.inbox");
assert(span != null);
const { traceId, spanId } = span.spanContext();

const fedifyExporter = new FedifySpanExporter(new MemoryKvStore());
await new Promise<void>((resolve) => {
fedifyExporter.export([span], () => resolve());
});
const [record] = await fedifyExporter.getActivitiesByTraceId(
traceId,
);
assert(record != null);
assertEquals(record.direction, "inbound");
assertEquals(record.spanId, spanId);

const failureLog = records.find((r) =>
String(r.rawMessage).startsWith(
"Failed to process the incoming activity",
)
);
assert(failureLog != null);
assertEquals(failureLog.properties.spanId, record.spanId);
assertEquals(failureLog.properties.traceId, record.traceId);
assertEquals(failureLog.properties.requestId, "outer-request");

const afterLog = records.find((r) =>
r.rawMessage === "after handled failure"
);
assert(afterLog != null);
assertEquals(afterLog.properties.requestId, outerContext.requestId);
assertEquals(afterLog.properties.traceId, outerContext.traceId);
assertEquals(afterLog.properties.spanId, outerContext.spanId);
} finally {
exporter.clear();
records.length = 0;
}
},
);

await t.step(
"outer context is restored after successful inbox processing",
async () => {
const [tracerProvider, exporter] = createTestTracerProvider();
try {
const kv = new MemoryKvStore();
const federation = createFederation<void>({ kv, tracerProvider });
const listeners = new ActivityListenerSet<InboxContext<void>>();
let received: Activity | null = null;
listeners.add(Create, (_ctx, activity) => {
received = activity;
});

const signed = await buildSignedRequest(
"https://example.com/activity/listener-success",
"https://example.com/note/listener-success",
);

await withContext(outerContext, async () => {
const response = await callHandleInbox(
federation,
kv,
tracerProvider,
listeners,
signed,
);
assertEquals(response.status, 202);
getLogger(["fedify", "federation", "handler.test"]).info(
"after success",
);
});
assert(received != null);

const span = exporter.getSpan("activitypub.inbox");
assert(span != null);
const { spanId } = span.spanContext();

const afterLog = records.find((r) =>
r.rawMessage === "after success"
);
assert(afterLog != null);
assertNotEquals(afterLog.properties.spanId, spanId);
assertEquals(afterLog.properties.requestId, outerContext.requestId);
assertEquals(afterLog.properties.traceId, outerContext.traceId);
assertEquals(afterLog.properties.spanId, outerContext.spanId);
} finally {
exporter.clear();
records.length = 0;
}
},
);

await t.step(
"concurrent inbox requests do not leak context into each other",
async () => {
const [tracerProvider, exporter] = createTestTracerProvider();
try {
const kv = new MemoryKvStore();
const federation = createFederation<void>({ kv, tracerProvider });

let arrived = 0;
let release: () => void = () => {};
const gate = new Promise<void>((resolve) => {
release = resolve;
});
const listeners = new ActivityListenerSet<InboxContext<void>>();
listeners.add(Create, async () => {
arrived++;
if (arrived >= 2) release();
await gate;
throw new Error("listener boom");
});

const [signedA, signedB] = await Promise.all([
buildSignedRequest(
"https://example.com/activity/concurrent-a",
"https://example.com/note/concurrent-a",
),
buildSignedRequest(
"https://example.com/activity/concurrent-b",
"https://example.com/note/concurrent-b",
),
]);

const [resultA, resultB] = await Promise.allSettled([
callHandleInbox(
federation,
kv,
tracerProvider,
listeners,
signedA,
).finally(release),
callHandleInbox(
federation,
kv,
tracerProvider,
listeners,
signedB,
).finally(release),
]);
assert(resultA.status === "fulfilled");
assert(resultB.status === "fulfilled");
assertEquals(resultA.value.status, 500);
assertEquals(resultB.value.status, 500);

const spans = exporter.getSpans("activitypub.inbox");
assertEquals(spans.length, 2);
const spanA = spans.find((s) =>
s.attributes["activitypub.activity.id"] ===
"https://example.com/activity/concurrent-a"
);
const spanB = spans.find((s) =>
s.attributes["activitypub.activity.id"] ===
"https://example.com/activity/concurrent-b"
);
assert(spanA != null && spanB != null);
assertNotEquals(
spanA.spanContext().spanId,
spanB.spanContext().spanId,
);

const logA = records.find((r) =>
r.properties.activityId ===
"https://example.com/activity/concurrent-a" &&
String(r.rawMessage).startsWith(
"Failed to process the incoming activity",
)
);
const logB = records.find((r) =>
r.properties.activityId ===
"https://example.com/activity/concurrent-b" &&
String(r.rawMessage).startsWith(
"Failed to process the incoming activity",
)
);
assert(logA != null && logB != null);
assertEquals(logA.properties.spanId, spanA.spanContext().spanId);
assertEquals(logB.properties.spanId, spanB.spanContext().spanId);
assertEquals(
logA.properties.traceId,
spanA.spanContext().traceId,
);
assertEquals(
logB.properties.traceId,
spanB.spanContext().traceId,
);
assertNotEquals(logA.properties.spanId, logB.properties.spanId);
} finally {
exporter.clear();
records.length = 0;
}
},
);
} finally {
await reset();
}
});

test("handleInbox() records fedify.queue.task.enqueued when queued", async () => {
const [meterProvider, recorder] = createTestMeterProvider();
const kv = new MemoryKvStore();
Expand Down
Loading
Loading