Skip to content
Merged
Show file tree
Hide file tree
Changes from 4 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
21 changes: 21 additions & 0 deletions CHANGES.md
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,25 @@ To be released.
deliberately exclude raw URLs, query strings, and identifier values to
keep cardinality bounded. [[#316], [#736], [#757]]

- Added OpenTelemetry metrics for ActivityPub collection requests handled
by `Federation.fetch()` and custom collection handlers:

- `activitypub.collection.request` (counter)
- `activitypub.collection.dispatch.duration` (histogram)
- `activitypub.collection.page.items` (histogram)
- `activitypub.collection.total_items` (histogram)

The metrics expose bounded collection dimensions:
`activitypub.collection.kind`, `activitypub.collection.page`,
`activitypub.collection.result`, `fedify.collection.dispatcher`, and
optional `http.response.status_code`. Built-in collections are classified
Comment thread
dahlia marked this conversation as resolved.
as `inbox`, `outbox`, `following`, `followers`, `liked`, `featured`, or
`featured_tags`; application-defined collection routes are collapsed into
`custom`. Collection IDs, cursors, custom route names, actor identifiers,
and full URLs are deliberately excluded so dashboards can aggregate
collection rate, latency, item counts, and `totalItems` values without
attacker-controlled cardinality. [[#316], [#741], [#777]]

- Added OpenTelemetry queue task metrics covering Fedify's enqueue and
worker boundaries for inbox, outbox, and fanout work:

Expand Down Expand Up @@ -208,6 +227,7 @@ To be released.
[#738]: https://github.com/fedify-dev/fedify/issues/738
[#739]: https://github.com/fedify-dev/fedify/issues/739
[#740]: https://github.com/fedify-dev/fedify/issues/740
[#741]: https://github.com/fedify-dev/fedify/issues/741
[#742]: https://github.com/fedify-dev/fedify/issues/742
[#748]: https://github.com/fedify-dev/fedify/pull/748
[#752]: https://github.com/fedify-dev/fedify/issues/752
Expand All @@ -220,6 +240,7 @@ To be released.
[#770]: https://github.com/fedify-dev/fedify/pull/770
[#771]: https://github.com/fedify-dev/fedify/pull/771
[#772]: https://github.com/fedify-dev/fedify/pull/772
[#777]: https://github.com/fedify-dev/fedify/pull/777

### @fedify/fixture

Expand Down
196 changes: 126 additions & 70 deletions docs/manual/opentelemetry.md

Large diffs are not rendered by default.

311 changes: 311 additions & 0 deletions packages/fedify/src/federation/handler.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1143,6 +1143,161 @@ test("handleCollection()", async () => {
assertEquals(onUnauthorizedCalled, null);
});

test("handleCollection() records OpenTelemetry collection metrics", async () => {
const [meterProvider, recorder] = createTestMeterProvider();
const federation = createFederation<void>({
kv: new MemoryKvStore(),
meterProvider,
});
const context = createRequestContext<void>({
federation,
data: undefined,
url: new URL("https://example.com/users/someone/followers"),
request: new Request("https://example.com/users/someone/followers", {
headers: { Accept: "application/activity+json" },
}),
getActorUri(identifier: string) {
return new URL(`https://example.com/users/${identifier}`);
},
});
const dispatcher: CollectionDispatcher<
Activity,
RequestContext<void>,
void,
void
> = (_ctx, identifier) =>
identifier === "someone"
? {
items: [
new Create({ id: new URL("https://example.com/activities/1") }),
new Create({ id: new URL("https://example.com/activities/2") }),
new Create({ id: new URL("https://example.com/activities/3") }),
],
}
: null;
const counter: CollectionCounter<void, void> = (_ctx, identifier) =>
identifier === "someone" ? 3 : null;

const response = await handleCollection(context.request, {
context,
name: "followers",
identifier: "someone",
uriGetter(identifier) {
return new URL(`https://example.com/users/${identifier}/followers`);
},
collectionCallbacks: { dispatcher, counter },
meterProvider,
onNotFound: () => new Response("Not found", { status: 404 }),
onUnauthorized: () => new Response("Unauthorized", { status: 401 }),
});
assertEquals(response.status, 200);

const requests = recorder.getMeasurements("activitypub.collection.request");
assertEquals(requests.length, 1);
assertEquals(requests[0].type, "counter");
assertEquals(requests[0].value, 1);
assertEquals(
requests[0].attributes["activitypub.collection.kind"],
"followers",
);
assertEquals(requests[0].attributes["activitypub.collection.page"], false);
assertEquals(
requests[0].attributes["fedify.collection.dispatcher"],
"built_in",
);
assertEquals(
requests[0].attributes["activitypub.collection.result"],
"served",
);
assertEquals(requests[0].attributes["http.response.status_code"], 200);
assertEquals(
"activitypub.collection.id" in requests[0].attributes,
false,
);

const durations = recorder.getMeasurements(
"activitypub.collection.dispatch.duration",
);
assertEquals(durations.length, 1);
assertEquals(durations[0].type, "histogram");
assert(durations[0].value >= 0);
assertEquals(
durations[0].attributes["activitypub.collection.result"],
"served",
);

const items = recorder.getMeasurements("activitypub.collection.page.items");
assertEquals(items.length, 1);
assertEquals(items[0].type, "histogram");
assertEquals(items[0].value, 3);
assertEquals(items[0].attributes["activitypub.collection.page"], false);

const totalItems = recorder.getMeasurements(
"activitypub.collection.total_items",
);
assertEquals(totalItems.length, 1);
assertEquals(totalItems[0].type, "histogram");
assertEquals(totalItems[0].value, 3);
});

test("handleCollection() records not_found collection metrics", async () => {
const [meterProvider, recorder] = createTestMeterProvider();
const federation = createFederation<void>({
kv: new MemoryKvStore(),
meterProvider,
});
const context = createRequestContext<void>({
federation,
data: undefined,
url: new URL("https://example.com/users/nobody/outbox"),
request: new Request("https://example.com/users/nobody/outbox", {
headers: { Accept: "application/activity+json" },
}),
});
const dispatcher: CollectionDispatcher<
Activity,
RequestContext<void>,
void,
void
> = () => null;

const response = await handleCollection(context.request, {
context,
name: "outbox",
identifier: "nobody",
uriGetter(identifier) {
return new URL(`https://example.com/users/${identifier}/outbox`);
},
collectionCallbacks: { dispatcher },
meterProvider,
onNotFound: () => new Response("Not found", { status: 404 }),
onUnauthorized: () => new Response("Unauthorized", { status: 401 }),
});
assertEquals(response.status, 404);

const requests = recorder.getMeasurements("activitypub.collection.request");
assertEquals(requests.length, 1);
assertEquals(requests[0].attributes["activitypub.collection.kind"], "outbox");
assertEquals(
requests[0].attributes["activitypub.collection.result"],
"not_found",
);
assertEquals(requests[0].attributes["http.response.status_code"], 404);

const durations = recorder.getMeasurements(
"activitypub.collection.dispatch.duration",
);
assertEquals(durations.length, 1);
assertEquals(
durations[0].attributes["activitypub.collection.result"],
"not_found",
);
assertEquals(
recorder.getMeasurements("activitypub.collection.page.items").length,
0,
);
});

test("handleInbox()", async () => {
const activity = new Create({
id: new URL("https://example.com/activities/1"),
Expand Down Expand Up @@ -4019,6 +4174,162 @@ test("handleCustomCollection()", async () => {
assertEquals(onUnauthorizedCalled, null);
});

test("handleCustomCollection() records OpenTelemetry collection metrics", async () => {
const [meterProvider, recorder] = createTestMeterProvider();
const federation = createFederation<void>({
kv: new MemoryKvStore(),
meterProvider,
});
const context = createRequestContext<void>({
federation,
data: undefined,
url: new URL("https://example.com/users/someone/custom"),
request: new Request("https://example.com/users/someone/custom", {
headers: { Accept: "application/activity+json" },
}),
});
const dispatcher: CustomCollectionDispatcher<
Create,
string,
RequestContext<void>,
void
> = (_ctx, values) =>
values.identifier === "someone"
? {
items: [
new Create({ id: new URL("https://example.com/activities/1") }),
new Create({ id: new URL("https://example.com/activities/2") }),
],
}
: null;
const counter: CustomCollectionCounter<string, void> = (_ctx, values) =>
values.identifier === "someone" ? 2 : null;

const response = await handleCustomCollection(context.request, {
context,
name: "custom collection",
values: { identifier: "someone" },
collectionCallbacks: { dispatcher, counter },
meterProvider,
onNotFound: () => new Response("Not found", { status: 404 }),
onUnauthorized: () => new Response("Unauthorized", { status: 401 }),
});
assertEquals(response.status, 200);

const requests = recorder.getMeasurements("activitypub.collection.request");
assertEquals(requests.length, 1);
assertEquals(requests[0].attributes["activitypub.collection.kind"], "custom");
assertEquals(requests[0].attributes["activitypub.collection.page"], false);
assertEquals(
requests[0].attributes["fedify.collection.dispatcher"],
"custom",
);
assertEquals(
requests[0].attributes["activitypub.collection.result"],
"served",
);
assertEquals(requests[0].attributes["http.response.status_code"], 200);

const durations = recorder.getMeasurements(
"activitypub.collection.dispatch.duration",
);
assertEquals(durations.length, 1);
assertEquals(
durations[0].attributes["fedify.collection.dispatcher"],
"custom",
);

const items = recorder.getMeasurements("activitypub.collection.page.items");
assertEquals(items.length, 1);
assertEquals(items[0].value, 2);

const totalItems = recorder.getMeasurements(
"activitypub.collection.total_items",
);
assertEquals(totalItems.length, 1);
assertEquals(totalItems[0].value, 2);
});

test("handleCustomCollection() classifies deferred collection metrics as error", async () => {
const [meterProvider, recorder] = createTestMeterProvider();
const federation = createFederation<void>({
kv: new MemoryKvStore(),
meterProvider,
});
const context = createRequestContext<void>({
federation,
data: undefined,
url: new URL("https://example.com/users/someone/custom"),
request: new Request("https://example.com/users/someone/custom", {
headers: { Accept: "application/activity+json" },
}),
});
const brokenActivity = new Create({
id: new URL("https://example.com/activities/1"),
});
globalThis.Object.defineProperty(brokenActivity, "toJsonLd", {
value: () => {
throw new Error("serialization failed");
},
});
const dispatcher: CustomCollectionDispatcher<
Create,
string,
RequestContext<void>,
void
> = () => ({ items: [brokenActivity] });
const counter: CustomCollectionCounter<string, void> = () => 1;

await assertRejects(
() =>
handleCustomCollection(context.request, {
context,
name: "custom collection",
values: { identifier: "someone" },
collectionCallbacks: { dispatcher, counter },
meterProvider,
onNotFound: () => new Response("Not found", { status: 404 }),
onUnauthorized: () => new Response("Unauthorized", { status: 401 }),
}),
Error,
"serialization failed",
);

const requests = recorder.getMeasurements("activitypub.collection.request");
assertEquals(requests.length, 1);
assertEquals(
requests[0].attributes["activitypub.collection.result"],
"error",
);

const durations = recorder.getMeasurements(
"activitypub.collection.dispatch.duration",
);
assertEquals(durations.length, 1);
assertEquals(
durations[0].attributes["activitypub.collection.result"],
"error",
);

const items = recorder.getMeasurements("activitypub.collection.page.items");
assertEquals(items.length, 1);
assertEquals(items[0].value, 1);
assertEquals(
items[0].attributes["activitypub.collection.result"],
"error",
);

const totalItems = recorder.getMeasurements(
"activitypub.collection.total_items",
);
assertEquals(totalItems.length, 1);
assertEquals(totalItems[0].value, 1);
assertEquals(
totalItems[0].attributes["activitypub.collection.result"],
"error",
);
});

test("handleInbox() records OpenTelemetry span events", async () => {
const [tracerProvider, exporter] = createTestTracerProvider();
const [meterProvider, recorder] = createTestMeterProvider();
Expand Down
Loading