diff --git a/CHANGES.md b/CHANGES.md index 4d3a4f422..d1b02bc83 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -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 diff --git a/changes.d/fedify/align-log-context-with-activity-spans.md b/changes.d/fedify/align-log-context-with-activity-spans.md new file mode 100644 index 000000000..4217839dc --- /dev/null +++ b/changes.d/fedify/align-log-context-with-activity-spans.md @@ -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] diff --git a/packages/fedify/src/federation/handler.test.ts b/packages/fedify/src/federation/handler.test.ts index 4d33cac6f..b15df577c 100644 --- a/packages/fedify/src/federation/handler.test.ts +++ b/packages/fedify/src/federation/handler.test.ts @@ -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"; @@ -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 = (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>, + kv: MemoryKvStore, + tracerProvider: ReturnType[0], + listeners: ActivityListenerSet>, + signed: Request, + ) { + const context = createRequestContext({ + 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({ kv, tracerProvider }); + const listeners = new ActivityListenerSet>(); + 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((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({ kv, tracerProvider }); + const listeners = new ActivityListenerSet>(); + 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({ kv, tracerProvider }); + + let arrived = 0; + let release: () => void = () => {}; + const gate = new Promise((resolve) => { + release = resolve; + }); + const listeners = new ActivityListenerSet>(); + 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(); diff --git a/packages/fedify/src/federation/handler.ts b/packages/fedify/src/federation/handler.ts index 14ecec4ee..a5e4c3b8b 100644 --- a/packages/fedify/src/federation/handler.ts +++ b/packages/fedify/src/federation/handler.ts @@ -32,7 +32,7 @@ import { readBoundedText, } from "../utils/body.ts"; import jsonld from "@fedify/vocab-runtime/jsonld"; -import { getLogger } from "@logtape/logtape"; +import { getLogger, withContext } from "@logtape/logtape"; import type { MeterProvider, Span, @@ -1825,8 +1825,12 @@ export async function handleInbox( if (options.recipient != null) { span.setAttribute("fedify.inbox.recipient", options.recipient); } + const spanContext = span.spanContext(); try { - return await handleInboxInternal(request, options, span); + return await withContext( + { traceId: spanContext.traceId, spanId: spanContext.spanId }, + () => handleInboxInternal(request, options, span), + ); } catch (e) { options.observation!.hasException = true; span.setStatus({ code: SpanStatusCode.ERROR, message: String(e) }); diff --git a/packages/fedify/src/federation/middleware.test.ts b/packages/fedify/src/federation/middleware.test.ts index 4ff0de460..be8ec6777 100644 --- a/packages/fedify/src/federation/middleware.test.ts +++ b/packages/fedify/src/federation/middleware.test.ts @@ -15,7 +15,12 @@ import { Person, } from "@fedify/vocab"; import { FetchError, getDocumentLoader, UrlError } from "@fedify/vocab-runtime"; -import { configure, type LogRecord, reset } from "@logtape/logtape"; +import { + configure, + type LogRecord, + reset, + withContext, +} from "@logtape/logtape"; import { metrics, SpanStatusCode } from "@opentelemetry/api"; import { DataPointType, @@ -42,6 +47,7 @@ import { strictEqual, throws, } from "node:assert/strict"; +import { AsyncLocalStorage } from "node:async_hooks"; import dns from "node:dns/promises"; import createFixture from "../../../fixture/src/fixtures/example.com/create.json" with { type: "json", @@ -8427,6 +8433,115 @@ test("FederationImpl.processQueuedTask() permanent failure", async (t) => { fetchMock.hardReset(); }); +test( + "FederationImpl.processQueuedTask() aligns the LogTape context with " + + "the delivery span inside the outbox worker", + async () => { + await withLogtapeLock(async () => { + fetchMock.spyGlobal(); + fetchMock.post("https://example.com/inbox-queue-failing", { + status: 500, + body: "Internal Server Error", + }); + + 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 [tracerProvider, exporter] = createTestTracerProvider(); + const queuedMessages: Message[] = []; + const queue: MessageQueue = { + enqueue(message, _options) { + queuedMessages.push(message); + return Promise.resolve(); + }, + listen(_handler, _options) { + return Promise.resolve(); + }, + }; + const federation = new FederationImpl({ + kv: new MemoryKvStore(), + queue, + allowPrivateAddress: true, + tracerProvider, + }); + + const message = { + type: "outbox", + id: crypto.randomUUID(), + baseUrl: "https://example.com", + keys: [], + activity: { + "@context": "https://www.w3.org/ns/activitystreams", + type: "Create", + id: "https://example.com/activity/queue-failing", + actor: "https://example.com/users/alice", + object: { type: "Note", content: "test" }, + }, + activityType: "https://www.w3.org/ns/activitystreams#Create", + inbox: "https://example.com/inbox-queue-failing", + sharedInbox: false, + started: new Date().toISOString(), + attempt: 0, + headers: {}, + traceContext: {}, + } satisfies OutboxMessage; + + await withContext({ batchId: "outer-batch" }, async () => { + await federation.processQueuedTask(undefined, message); + }); + + assertEquals(queuedMessages.length, 1); + + const deliverySpan = exporter.getSpan("activitypub.send_activity"); + assert(deliverySpan != null); + const deliveryCtx = deliverySpan.spanContext(); + + const workerSpan = exporter.getSpan("activitypub.outbox"); + assert(workerSpan != null); + const workerCtx = workerSpan.spanContext(); + + assertNotEquals(deliveryCtx.spanId, workerCtx.spanId); + + const deliveryLog = records.find((r) => + String(r.rawMessage).startsWith( + "Failed to send activity {activityId} to {inbox} ({status}", + ) + ); + assert(deliveryLog != null); + assertEquals(deliveryLog.properties.spanId, deliveryCtx.spanId); + assertEquals(deliveryLog.properties.traceId, deliveryCtx.traceId); + assertEquals(deliveryLog.properties.messageId, message.id); + assertEquals(deliveryLog.properties.batchId, "outer-batch"); + + const retryLog = records.find((r) => + String(r.rawMessage).startsWith( + "Failed to send activity {activityId} to {inbox} (attempt", + ) + ); + assert(retryLog != null); + assertEquals(retryLog.properties.spanId, workerCtx.spanId); + assertEquals(retryLog.properties.traceId, workerCtx.traceId); + assertNotEquals(retryLog.properties.spanId, deliveryCtx.spanId); + assertEquals(retryLog.properties.messageId, message.id); + assertEquals(retryLog.properties.batchId, "outer-batch"); + } finally { + await reset(); + fetchMock.hardReset(); + } + }); + }, +); + test("FederationImpl.processQueuedTask() circuit breaker", async (t) => { fetchMock.spyGlobal(); diff --git a/packages/fedify/src/federation/send.test.ts b/packages/fedify/src/federation/send.test.ts index 1889754e1..0473c1ee0 100644 --- a/packages/fedify/src/federation/send.test.ts +++ b/packages/fedify/src/federation/send.test.ts @@ -23,6 +23,13 @@ import { assertNotEquals, assertRejects, } from "@std/assert"; +import { + configure, + getLogger, + type LogRecord, + reset, + withContext, +} from "@logtape/logtape"; import { AggregationTemporality, InMemoryMetricExporter, @@ -30,7 +37,9 @@ import { PeriodicExportingMetricReader, } from "@opentelemetry/sdk-metrics"; import fetchMock from "fetch-mock"; +import { AsyncLocalStorage } from "node:async_hooks"; import dns from "node:dns/promises"; +import { FedifySpanExporter } from "../otel/exporter.ts"; import { verifyRequest } from "../sig/http.ts"; import { doesActorOwnKey } from "../sig/owner.ts"; import { @@ -39,6 +48,7 @@ import { rsaPrivateKey2, rsaPublicKey2, } from "../testing/keys.ts"; +import { MemoryKvStore } from "./kv.ts"; import { extractInboxes, sendActivity, SendActivityError } from "./send.ts"; @@ -564,6 +574,305 @@ test("sendActivity() records OpenTelemetry span events", async (t) => { }); }); +function createActivity(id: string) { + return { + "@context": "https://www.w3.org/ns/activitystreams", + type: "Create", + id, + actor: "https://example.com/person", + }; +} + +function deliveryParams( + activityId: string, + inbox: string, + tracerProvider: ReturnType[0], +) { + return { + activity: createActivity(activityId), + activityId, + activityType: "https://www.w3.org/ns/activitystreams#Create", + keys: [{ + keyId: new URL("https://example.com/person#key"), + privateKey: rsaPrivateKey2, + }], + inbox: new URL(inbox), + tracerProvider, + }; +} + +test("sendActivity() aligns the LogTape context with its own span", async (t) => { + const [tracerProvider, exporter] = createTestTracerProvider(); + fetchMock.spyGlobal(); + + 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 warning logged during a successful send matches the resulting " + + "TraceActivityRecord", + async () => { + try { + fetchMock.post("https://example.com/inbox-warn", { status: 202 }); + + await withContext({ requestId: "outer-request-warn" }, async () => { + await sendActivity({ + ...deliveryParams( + "https://example.com/activity/warn", + "https://example.com/inbox-warn", + tracerProvider, + ), + keys: [{ + keyId: ed25519Multikey.id!, + privateKey: ed25519PrivateKey, + }], + }); + }); + + const span = exporter.getSpan("activitypub.send_activity"); + assert(span != null); + const { traceId, spanId } = span.spanContext(); + + const kv = new MemoryKvStore(); + const fedifyExporter = new FedifySpanExporter(kv); + await new Promise((resolve) => { + fedifyExporter.export([span], () => resolve()); + }); + const [record] = await fedifyExporter.getActivitiesByTraceId( + traceId, + ); + assert(record != null); + assertEquals(record.direction, "outbound"); + assertEquals(record.spanId, spanId); + + const warnLog = records.find((r) => + String(r.rawMessage).startsWith("No supported key found to sign") + ); + assert(warnLog != null); + assertEquals(warnLog.properties.spanId, record.spanId); + assertEquals(warnLog.properties.traceId, record.traceId); + assertEquals(warnLog.properties.requestId, "outer-request-warn"); + } finally { + exporter.clear(); + records.length = 0; + fetchMock.hardReset(); + fetchMock.spyGlobal(); + } + }, + ); + + await t.step( + "a failed delivery's error log matches its own span, and the outer " + + "context is restored once the rejection propagates", + async () => { + try { + fetchMock.post("https://example.com/inbox-failing", { + status: 500, + body: "Internal Server Error", + }); + + await withContext(outerContext, async () => { + await assertRejects( + () => + sendActivity( + deliveryParams( + "https://example.com/activity/failing", + "https://example.com/inbox-failing", + tracerProvider, + ), + ), + SendActivityError, + ); + getLogger(["fedify", "federation", "send.test"]).info( + "after failure", + ); + }); + + const span = exporter.getSpan("activitypub.send_activity"); + assert(span != null); + const { traceId, spanId } = span.spanContext(); + + const kv = new MemoryKvStore(); + const fedifyExporter = new FedifySpanExporter(kv); + await new Promise((resolve) => { + fedifyExporter.export([span], () => resolve()); + }); + assertEquals( + await fedifyExporter.getActivitiesByTraceId(traceId), + [], + ); + + const failureLog = records.find((r) => + String(r.rawMessage).startsWith( + "Failed to send activity {activityId} to {inbox} ({status}", + ) + ); + assert(failureLog != null); + assertEquals(failureLog.properties.spanId, spanId); + assertEquals(failureLog.properties.traceId, traceId); + assertEquals( + failureLog.properties.requestId, + outerContext.requestId, + ); + + const afterLog = records.find((r) => + r.rawMessage === "after 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; + fetchMock.hardReset(); + fetchMock.spyGlobal(); + } + }, + ); + + await t.step( + "outer context is restored after a successful send", + async () => { + try { + fetchMock.post("https://example.com/inbox-ok", { status: 202 }); + + await withContext(outerContext, async () => { + await sendActivity( + deliveryParams( + "https://example.com/activity/ok", + "https://example.com/inbox-ok", + tracerProvider, + ), + ); + getLogger(["fedify", "federation", "send.test"]).info( + "after success", + ); + }); + + const span = exporter.getSpan("activitypub.send_activity"); + 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; + fetchMock.hardReset(); + fetchMock.spyGlobal(); + } + }, + ); + + await t.step( + "concurrent deliveries do not leak context into each other", + async () => { + try { + let arrived = 0; + let release: () => void = () => {}; + const gate = new Promise((resolve) => { + release = resolve; + }); + const handler = () => async () => { + arrived++; + if (arrived >= 2) release(); + await gate; + return { status: 500, body: "Internal Server Error" }; + }; + fetchMock.post("https://example.com/inbox-concurrent-a", handler()); + fetchMock.post("https://example.com/inbox-concurrent-b", handler()); + + const [resultA, resultB] = await Promise.allSettled([ + sendActivity( + deliveryParams( + "https://example.com/activity/a", + "https://example.com/inbox-concurrent-a", + tracerProvider, + ), + ).finally(release), + sendActivity( + deliveryParams( + "https://example.com/activity/b", + "https://example.com/inbox-concurrent-b", + tracerProvider, + ), + ).finally(release), + ]); + assertEquals(resultA.status, "rejected"); + assertEquals(resultB.status, "rejected"); + + const spans = exporter.getSpans("activitypub.send_activity"); + assertEquals(spans.length, 2); + const spanA = spans.find((s) => + s.attributes["activitypub.activity.id"] === + "https://example.com/activity/a" + ); + const spanB = spans.find((s) => + s.attributes["activitypub.activity.id"] === + "https://example.com/activity/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/a" && + String(r.rawMessage).startsWith("Failed to send activity") + ); + const logB = records.find((r) => + r.properties.activityId === "https://example.com/activity/b" && + String(r.rawMessage).startsWith("Failed to send 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; + fetchMock.hardReset(); + } + }, + ); + } finally { + await reset(); + fetchMock.hardReset(); + } +}); + test("sendActivity() records OpenTelemetry delivery metrics", async (t) => { const [meterProvider, recorder] = createTestMeterProvider(); fetchMock.spyGlobal(); diff --git a/packages/fedify/src/federation/send.ts b/packages/fedify/src/federation/send.ts index 656b50014..ab89d327a 100644 --- a/packages/fedify/src/federation/send.ts +++ b/packages/fedify/src/federation/send.ts @@ -1,6 +1,6 @@ import type { Recipient } from "@fedify/vocab"; import { FetchError, UrlError, validatePublicUrl } from "@fedify/vocab-runtime"; -import { getLogger } from "@logtape/logtape"; +import { getLogger, withContext } from "@logtape/logtape"; import { type Attributes, type MeterProvider, @@ -294,8 +294,12 @@ export function sendActivity( if (options.activityType != null) { span.setAttribute("activitypub.activity.type", options.activityType); } + const spanContext = span.spanContext(); try { - await sendActivityInternal({ ...options, tracerProvider }, span); + await withContext( + { traceId: spanContext.traceId, spanId: spanContext.spanId }, + () => sendActivityInternal({ ...options, tracerProvider }, span), + ); } catch (e) { span.setStatus({ code: SpanStatusCode.ERROR, message: String(e) }); throw e;