diff --git a/packages/relay/src/base.ts b/packages/relay/src/base.ts index d7a18e075..e70228cb4 100644 --- a/packages/relay/src/base.ts +++ b/packages/relay/src/base.ts @@ -1,13 +1,38 @@ -import type { Context, Federation, FederationBuilder } from "@fedify/fedify"; -import { isActor, Object as APObject } from "@fedify/vocab"; +import type { + Context, + Federation, + FederationBuilder, + InboxContext, + InboxListenerSetters, +} from "@fedify/fedify"; import { - isRelayFollowerData, + type Actor, + Announce, + Create, + Delete, + Follow, + Move, + Undo, + Update, +} from "@fedify/vocab"; +import type { Logger } from "@logtape/logtape"; +import { + handleUndoFollow, + sendFollowResponse, + validateFollowActivity, +} from "./follow.ts"; +import { + parseRelayFollowerData, type Relay, RELAY_SERVER_ACTOR, type RelayFollower, + type RelayFollowerState, type RelayOptions, } from "./types.ts"; +/** @internal */ +export type RelayableActivity = Create | Delete | Move | Update | Announce; + /** * Abstract base class for relay implementations. * Provides common infrastructure for both Mastodon and LitePub relays. @@ -19,6 +44,9 @@ export abstract class BaseRelay implements Relay { protected options: RelayOptions; protected federation?: Federation; + protected abstract readonly initialFollowerState: RelayFollowerState; + protected abstract readonly logger: Logger; + constructor( options: RelayOptions, relayBuilder: FederationBuilder, @@ -33,31 +61,6 @@ export abstract class BaseRelay implements Relay { }); } - /** - * Helper method to parse and validate follower data from storage. - * Deserializes JSON-LD actor data and validates it. - * - * @param actorId The actor ID of the follower - * @param data Raw data from KV store - * @returns RelayFollower object if valid, null otherwise - * @internal - */ - private async parseFollowerData( - actorId: string, - data: unknown, - ): Promise { - if (!isRelayFollowerData(data)) return null; - - const actor = await APObject.fromJsonLd(data.actor); - if (!isActor(actor)) return null; - - return { - actorId, - actor, - state: data.state, - }; - } - /** * Lists all followers of the relay. * @@ -88,7 +91,7 @@ export abstract class BaseRelay implements Relay { const actorId = entry.key[1]; if (typeof actorId !== "string") continue; - const follower = await this.parseFollowerData(actorId, entry.value); + const follower = await parseRelayFollowerData(actorId, entry.value); if (follower) yield follower; } } @@ -123,14 +126,100 @@ export abstract class BaseRelay implements Relay { */ async getFollower(actorId: string): Promise { const followerData = await this.options.kv.get(["follower", actorId]); - return await this.parseFollowerData(actorId, followerData); + return await parseRelayFollowerData(actorId, followerData); } - /** - * Set up inbox listeners for handling ActivityPub activities. - * Each relay type implements this method with protocol-specific logic. - */ - protected abstract setupInboxListeners(): void; + protected shouldSkipFollow( + _ctx: InboxContext, + _follower: Actor, + ): Promise { + return Promise.resolve(false); + } + + protected afterFollowApproved( + _ctx: InboxContext, + _follower: Actor, + ): Promise { + return Promise.resolve(); + } + + protected abstract deliverActivity( + ctx: InboxContext, + activity: RelayableActivity, + excludeBaseUris: URL[], + ): Promise; + + async #handleFollow( + ctx: InboxContext, + follow: Follow, + ): Promise { + const follower = await validateFollowActivity(ctx, follow); + if (follower?.id == null || await this.shouldSkipFollow(ctx, follower)) { + return; + } + + const approved = await this.options.subscriptionHandler(ctx, follower); + if (approved) { + await ctx.data.kv.set( + ["follower", follower.id.href], + { + actor: await follower.toJsonLd(), + state: this.initialFollowerState, + }, + ); + } + + await sendFollowResponse(ctx, follow, follower, approved); + if (approved) await this.afterFollowApproved(ctx, follower); + } + + async #relayActivity( + ctx: InboxContext, + activity: RelayableActivity, + ): Promise { + const senderId = activity.actorId; + const excludeBaseUris = senderId == null ? [] : [senderId]; + await this.deliverActivity(ctx, activity, excludeBaseUris); + } + + protected setupInboxListeners(): InboxListenerSetters { + if (this.federation == null) { + throw new Error("Federation must be initialized before inbox listeners"); + } + + const listeners = this.federation.setInboxListeners( + "/users/{identifier}/inbox", + "/inbox", + ); + listeners + .on(Follow, async (ctx, follow) => await this.#handleFollow(ctx, follow)) + .on( + Undo, + async (ctx, undo) => await handleUndoFollow(ctx, undo, this.logger), + ) + .on( + Create, + async (ctx, create) => await this.#relayActivity(ctx, create), + ) + .on( + Delete, + async (ctx, deleteActivity) => + await this.#relayActivity(ctx, deleteActivity), + ) + .on( + Move, + async (ctx, move) => await this.#relayActivity(ctx, move), + ) + .on( + Update, + async (ctx, update) => await this.#relayActivity(ctx, update), + ) + .on( + Announce, + async (ctx, announce) => await this.#relayActivity(ctx, announce), + ); + return listeners; + } async #getFederation(): Promise> { if (this.federation == null) { diff --git a/packages/relay/src/builder.ts b/packages/relay/src/builder.ts index 7bd40aa3f..deca78026 100644 --- a/packages/relay/src/builder.ts +++ b/packages/relay/src/builder.ts @@ -7,9 +7,9 @@ import { importJwk, } from "@fedify/fedify"; import type { Actor } from "@fedify/vocab"; -import { Application, isActor, Object } from "@fedify/vocab"; +import { Application } from "@fedify/vocab"; import { - isRelayFollowerData, + parseRelayFollowerData, RELAY_SERVER_ACTOR, type RelayOptions, } from "./types.ts"; @@ -78,12 +78,12 @@ async function getFollowerActors( ): Promise { const actors: Actor[] = []; - for await (const { value } of ctx.data.kv.list(["follower"])) { - if (!isRelayFollowerData(value)) continue; - if (value.state !== "accepted") continue; - const actor = await Object.fromJsonLd(value.actor); - if (!isActor(actor)) continue; - actors.push(actor); + for await (const { key, value } of ctx.data.kv.list(["follower"])) { + const actorId = key[1]; + if (typeof actorId !== "string") continue; + const follower = await parseRelayFollowerData(actorId, value); + if (follower?.state !== "accepted") continue; + actors.push(follower.actor); } return actors; diff --git a/packages/relay/src/litepub.test.ts b/packages/relay/src/litepub.test.ts index 899bf6d8b..bc6a92cb5 100644 --- a/packages/relay/src/litepub.test.ts +++ b/packages/relay/src/litepub.test.ts @@ -18,7 +18,7 @@ import { getDocumentLoader, type RemoteDocument, } from "@fedify/vocab-runtime"; -import { ok, strictEqual } from "node:assert"; +import { deepStrictEqual, ok, strictEqual } from "node:assert"; import test, { describe } from "node:test"; import { isRelayFollowerData } from "./types.ts"; @@ -466,60 +466,121 @@ describe("LitePubRelay", () => { strictEqual(followerData, undefined); }); - test("ignores duplicate Follow activity from pending follower", async () => { - const kv = new MemoryKvStore(); - let handlerCallCount = 0; - - const relay = createRelay("litepub", { - kv, - origin: "https://relay.example.com", - documentLoaderFactory: () => mockDocumentLoader, - authenticatedDocumentLoaderFactory: () => mockDocumentLoader, - subscriptionHandler: async (_ctx, _actor) => { - handlerCallCount++; - return await Promise.resolve(true); - }, + test("replaces malformed follower data on Follow", async () => { + const followerId = "https://remote.example.com/users/alice"; + const mismatchedActor = new Person({ + id: new URL("https://remote.example.com/users/bob"), }); + const malformedRows = [ + { state: "pending" }, + { actor: null, state: "pending" }, + { actor: {}, state: "pending" }, + { actor: await mismatchedActor.toJsonLd(), state: "pending" }, + ]; - const follower = new Person({ - id: new URL("https://remote.example.com/users/alice"), - preferredUsername: "alice", - inbox: new URL("https://remote.example.com/users/alice/inbox"), - }); + for (const malformedRow of malformedRows) { + const kv = new MemoryKvStore(); + await kv.set(["follower", followerId], malformedRow); + let handlerCallCount = 0; + + const relay = createRelay("litepub", { + kv, + origin: "https://relay.example.com", + documentLoaderFactory: () => mockDocumentLoader, + authenticatedDocumentLoaderFactory: () => mockDocumentLoader, + subscriptionHandler: () => { + handlerCallCount++; + return Promise.resolve(true); + }, + }); + strictEqual(await relay.getFollower(followerId), null); - // Pre-populate with pending follower - await kv.set( - ["follower", "https://remote.example.com/users/alice"], - { actor: await follower.toJsonLd(), state: "pending" }, - ); + const followActivity = new Follow({ + id: new URL("https://remote.example.com/activities/follow/1"), + actor: new URL(followerId), + object: new URL("https://relay.example.com/users/relay"), + }); + let request = new Request("https://relay.example.com/inbox", { + method: "POST", + headers: { "Content-Type": "application/activity+json" }, + body: JSON.stringify( + await followActivity.toJsonLd({ contextLoader: mockDocumentLoader }), + ), + }); + request = await signRequest( + request, + rsaKeyPair.privateKey, + rsaPublicKey.id, + ); - const followActivity = new Follow({ - id: new URL("https://remote.example.com/activities/follow/1"), - actor: follower.id, - object: new URL("https://relay.example.com/users/relay"), - }); + await relay.fetch(request); - let request = new Request("https://relay.example.com/inbox", { - method: "POST", - headers: { - "Content-Type": "application/activity+json", - }, - body: JSON.stringify( - await followActivity.toJsonLd({ contextLoader: mockDocumentLoader }), - ), - }); + strictEqual(handlerCallCount, 1); + const follower = await relay.getFollower(followerId); + ok(follower); + strictEqual(follower.state, "pending"); + strictEqual(follower.actor.id?.href, followerId); + } + }); - request = await signRequest( - request, - rsaKeyPair.privateKey, - rsaPublicKey.id, - ); + for (const state of ["pending", "accepted"] as const) { + test(`ignores duplicate Follow activity from ${state} follower`, async () => { + const kv = new MemoryKvStore(); + let handlerCallCount = 0; + + const relay = createRelay("litepub", { + kv, + origin: "https://relay.example.com", + documentLoaderFactory: () => mockDocumentLoader, + authenticatedDocumentLoaderFactory: () => mockDocumentLoader, + subscriptionHandler: async (_ctx, _actor) => { + handlerCallCount++; + return await Promise.resolve(true); + }, + }); - await relay.fetch(request); + const follower = new Person({ + id: new URL("https://remote.example.com/users/alice"), + preferredUsername: "alice", + inbox: new URL("https://remote.example.com/users/alice/inbox"), + }); + await kv.set( + ["follower", "https://remote.example.com/users/alice"], + { actor: await follower.toJsonLd(), state }, + ); - // Verify handler was NOT called (duplicate follow ignored) - strictEqual(handlerCallCount, 0); - }); + const followActivity = new Follow({ + id: new URL("https://remote.example.com/activities/follow/1"), + actor: follower.id, + object: new URL("https://relay.example.com/users/relay"), + }); + + let request = new Request("https://relay.example.com/inbox", { + method: "POST", + headers: { + "Content-Type": "application/activity+json", + }, + body: JSON.stringify( + await followActivity.toJsonLd({ contextLoader: mockDocumentLoader }), + ), + }); + request = await signRequest( + request, + rsaKeyPair.privateKey, + rsaPublicKey.id, + ); + + await relay.fetch(request); + + strictEqual(handlerCallCount, 0); + const followerData = await kv.get([ + "follower", + "https://remote.example.com/users/alice", + ]); + ok(isRelayFollowerData(followerData)); + strictEqual(followerData.state, state); + }); + } test("handles Accept activity completing reciprocal follow", async () => { const kv = new MemoryKvStore(); @@ -583,64 +644,120 @@ describe("LitePubRelay", () => { strictEqual(followerData.state, "accepted"); }); - test("handles Undo Follow activity", async () => { - const kv = new MemoryKvStore(); - - // Pre-populate with an accepted follower + test("ignores Accept activity for invalid follower data", async () => { const followerId = "https://remote.example.com/users/alice"; - const follower = new Person({ - id: new URL(followerId), - preferredUsername: "alice", - inbox: new URL("https://remote.example.com/users/alice/inbox"), + const mismatchedActor = new Person({ + id: new URL("https://remote.example.com/users/bob"), }); + const invalidRows = [ + { state: "pending" }, + { actor: null, state: "pending" }, + { actor: {}, state: "pending" }, + { actor: await mismatchedActor.toJsonLd(), state: "pending" }, + ]; - await kv.set( - ["follower", followerId], - { actor: await follower.toJsonLd(), state: "accepted" }, - ); + for (const invalidRow of invalidRows) { + const kv = new MemoryKvStore(); + await kv.set(["follower", followerId], invalidRow); - const relay = createRelay("litepub", { - kv, - origin: "https://relay.example.com", - documentLoaderFactory: () => mockDocumentLoader, - authenticatedDocumentLoaderFactory: () => mockDocumentLoader, - subscriptionHandler: () => Promise.resolve(true), - }); + const relay = createRelay("litepub", { + kv, + origin: "https://relay.example.com", + documentLoaderFactory: () => mockDocumentLoader, + authenticatedDocumentLoaderFactory: () => mockDocumentLoader, + subscriptionHandler: () => Promise.resolve(true), + }); - const originalFollow = new Follow({ - id: new URL("https://remote.example.com/activities/follow/1"), - actor: new URL(followerId), - object: new URL("https://relay.example.com/users/relay"), - }); + const relayFollow = new Follow({ + id: new URL("https://relay.example.com/activities/follow/1"), + actor: new URL("https://relay.example.com/users/relay"), + object: new URL(followerId), + }); + const acceptActivity = new Accept({ + id: new URL("https://remote.example.com/activities/accept/1"), + actor: new URL(followerId), + object: relayFollow, + }); - const undoActivity = new Undo({ - id: new URL("https://remote.example.com/activities/undo/1"), - actor: new URL(followerId), - object: originalFollow, - }); + let request = new Request("https://relay.example.com/inbox", { + method: "POST", + headers: { + "Content-Type": "application/activity+json", + }, + body: JSON.stringify( + await acceptActivity.toJsonLd({ contextLoader: mockDocumentLoader }), + ), + }); + request = await signRequest( + request, + rsaKeyPair.privateKey, + rsaPublicKey.id, + ); - let request = new Request("https://relay.example.com/inbox", { - method: "POST", - headers: { - "Content-Type": "application/activity+json", - }, - body: JSON.stringify( - await undoActivity.toJsonLd({ contextLoader: mockDocumentLoader }), - ), - }); + await relay.fetch(request); - request = await signRequest( - request, - rsaKeyPair.privateKey, - rsaPublicKey.id, - ); + deepStrictEqual(await kv.get(["follower", followerId]), invalidRow); + } + }); - await relay.fetch(request); + for (const state of ["pending", "accepted"] as const) { + test(`handles Undo Follow activity for ${state} follower`, async () => { + const kv = new MemoryKvStore(); - // Verify follower was removed - const followerData = await kv.get(["follower", followerId]); - strictEqual(followerData, undefined); - }); + const followerId = "https://remote.example.com/users/alice"; + const follower = new Person({ + id: new URL(followerId), + preferredUsername: "alice", + inbox: new URL("https://remote.example.com/users/alice/inbox"), + }); + + await kv.set( + ["follower", followerId], + { actor: await follower.toJsonLd(), state }, + ); + + const relay = createRelay("litepub", { + kv, + origin: "https://relay.example.com", + documentLoaderFactory: () => mockDocumentLoader, + authenticatedDocumentLoaderFactory: () => mockDocumentLoader, + subscriptionHandler: () => Promise.resolve(true), + }); + + const originalFollow = new Follow({ + id: new URL("https://remote.example.com/activities/follow/1"), + actor: new URL(followerId), + object: new URL("https://relay.example.com/users/relay"), + }); + + const undoActivity = new Undo({ + id: new URL("https://remote.example.com/activities/undo/1"), + actor: new URL(followerId), + object: originalFollow, + }); + + let request = new Request("https://relay.example.com/inbox", { + method: "POST", + headers: { + "Content-Type": "application/activity+json", + }, + body: JSON.stringify( + await undoActivity.toJsonLd({ contextLoader: mockDocumentLoader }), + ), + }); + + request = await signRequest( + request, + rsaKeyPair.privateKey, + rsaPublicKey.id, + ); + + await relay.fetch(request); + + const followerData = await kv.get(["follower", followerId]); + strictEqual(followerData, undefined); + }); + } test("handles Create activity with Announce forwarding", async () => { const kv = new MemoryKvStore(); diff --git a/packages/relay/src/litepub.ts b/packages/relay/src/litepub.ts index 8114e0f50..52b2bde79 100644 --- a/packages/relay/src/litepub.ts +++ b/packages/relay/src/litepub.ts @@ -1,24 +1,17 @@ -import type { InboxContext } from "@fedify/fedify"; +import type { InboxContext, InboxListenerSetters } from "@fedify/fedify"; import { Accept, + type Actor, Announce, - Create, - Delete, Follow, isActor, - Move, PUBLIC_COLLECTION, - Undo, - Update, } from "@fedify/vocab"; import { getLogger } from "@logtape/logtape"; -import { BaseRelay } from "./base.ts"; -import { - handleUndoFollow, - sendFollowResponse, - validateFollowActivity, -} from "./follow.ts"; +import { BaseRelay, type RelayableActivity } from "./base.ts"; import { + isRelayFollowerData, + parseRelayFollowerData, RELAY_SERVER_ACTOR, type RelayFollowerData, type RelayOptions, @@ -34,13 +27,47 @@ const logger = getLogger(["fedify", "relay", "litepub"]); * @since 2.0.0 */ export class LitePubRelay extends BaseRelay { - async #announceToFollowers( + protected readonly initialFollowerState = "pending"; + protected readonly logger = logger; + + protected override async shouldSkipFollow( ctx: InboxContext, - activity: Create | Delete | Move | Update | Announce, + follower: Actor, + ): Promise { + if (follower.id == null) return true; + const existingFollow = await ctx.data.kv.get([ + "follower", + follower.id.href, + ]); + const storedFollower = await parseRelayFollowerData( + follower.id.href, + existingFollow, + ); + return storedFollower != null; + } + + protected override async afterFollowApproved( + ctx: InboxContext, + follower: Actor, ): Promise { - const sender = await activity.getActor(ctx); - const excludeBaseUris = sender?.id ? [new URL(sender.id)] : []; + if (follower.id == null) return; + const relayActorUri = ctx.getActorUri(RELAY_SERVER_ACTOR); + await ctx.sendActivity( + { identifier: RELAY_SERVER_ACTOR }, + follower, + new Follow({ + actor: relayActorUri, + object: follower.id, + to: follower.id, + }), + ); + } + protected async deliverActivity( + ctx: InboxContext, + activity: RelayableActivity, + excludeBaseUris: URL[], + ): Promise { const announce = new Announce({ id: new URL(`/announce#${crypto.randomUUID()}`, ctx.origin), actor: ctx.getActorUri(RELAY_SERVER_ACTOR), @@ -60,105 +87,44 @@ export class LitePubRelay extends BaseRelay { ); } - protected setupInboxListeners(): void { - if (this.federation != null) { - this.federation.setInboxListeners("/users/{identifier}/inbox", "/inbox") - .on(Follow, async (ctx, follow) => { - const follower = await validateFollowActivity(ctx, follow); - if (!follower || !follower.id) return; - - // Litepub-specific: check if already in pending state - const existingFollow = await ctx.data.kv.get([ - "follower", - follower.id.href, - ]); - if (existingFollow?.state === "pending") return; - - const approved = await this.options.subscriptionHandler( - ctx, - follower, - ); + protected override setupInboxListeners(): InboxListenerSetters { + return super.setupInboxListeners().on(Accept, async (ctx, accept) => { + // Validate follow activity from accept activity + const follow = await accept.getObject({ + crossOrigin: "trust", + ...ctx, + }); + if (!(follow instanceof Follow)) return; + const relayActorId = follow.actorId; + if (relayActorId == null) return; - if (approved) { - // Litepub-specific: save with "pending" state - await ctx.data.kv.set( - ["follower", follower.id.href], - { actor: await follower.toJsonLd(), state: "pending" }, - ); + // Validate follower actor - accept activity sender + const followerActor = await accept.getActor(ctx); + if (!isActor(followerActor) || !followerActor.id) return; + const parsed = ctx.parseUri(relayActorId); + if (parsed == null || parsed.type !== "actor") return; - await sendFollowResponse(ctx, follow, follower, approved); + // Get follower from kv store + const followerData = await ctx.data.kv.get([ + "follower", + followerActor.id.href, + ]); + if (!isRelayFollowerData(followerData)) return; + const storedFollower = await parseRelayFollowerData( + followerActor.id.href, + followerData, + ); + if (storedFollower == null) return; - // Litepub-specific: send reciprocal follow - const relayActorUri = ctx.getActorUri(RELAY_SERVER_ACTOR); - await ctx.sendActivity( - { identifier: RELAY_SERVER_ACTOR }, - follower, - new Follow({ - actor: relayActorUri, - object: follower.id, - to: follower.id, - }), - ); - } else { - await sendFollowResponse(ctx, follow, follower, approved); - } - }) - .on(Accept, async (ctx, accept) => { - // Validate follow activity from accept activity - const follow = await accept.getObject({ - crossOrigin: "trust", - ...ctx, - }); - if (!(follow instanceof Follow)) return; - const relayActorId = follow.actorId; - if (relayActorId == null) return; - - // Validate follower actor - accept activity sender - const followerActor = await accept.getActor(); - if (!isActor(followerActor) || !followerActor.id) return; - const parsed = ctx.parseUri(relayActorId); - if (parsed == null || parsed.type !== "actor") return; - - // Get follower from kv store - const followerData = await ctx.data.kv.get([ - "follower", - followerActor.id.href, - ]); - if (followerData == null) return; - - // Update follower state to accepted - const updatedFollowerData = { ...followerData, state: "accepted" }; - await ctx.data.kv.set( - ["follower", followerActor.id.href], - updatedFollowerData, - ); - }) - .on( - Undo, - async (ctx, undo) => await handleUndoFollow(ctx, undo, logger), - ) - .on( - Create, - async (ctx, create) => await this.#announceToFollowers(ctx, create), - ) - .on( - Update, - async (ctx, update) => await this.#announceToFollowers(ctx, update), - ) - .on( - Move, - async (ctx, move) => await this.#announceToFollowers(ctx, move), - ) - .on( - Delete, - async (ctx, deleteActivity) => - await this.#announceToFollowers(ctx, deleteActivity), - ) - .on( - Announce, - async (ctx, announce) => - await this.#announceToFollowers(ctx, announce), - ); - } + // Update follower state to accepted + const updatedFollowerData: RelayFollowerData = { + ...followerData, + state: "accepted", + }; + await ctx.data.kv.set( + ["follower", followerActor.id.href], + updatedFollowerData, + ); + }); } } diff --git a/packages/relay/src/mastodon.test.ts b/packages/relay/src/mastodon.test.ts index d5a73b503..df36d631a 100644 --- a/packages/relay/src/mastodon.test.ts +++ b/packages/relay/src/mastodon.test.ts @@ -1,7 +1,8 @@ // deno-lint-ignore-file no-explicit-any -import { MemoryKvStore, signRequest } from "@fedify/fedify"; +import { MemoryKvStore, signJsonLd, signRequest } from "@fedify/fedify"; import { createRelay, type RelayOptions } from "@fedify/relay"; import { + Announce, Create, Delete, Follow, @@ -16,7 +17,7 @@ import { getDocumentLoader, type RemoteDocument, } from "@fedify/vocab-runtime"; -import { ok, strictEqual } from "node:assert"; +import { deepStrictEqual, ok, strictEqual } from "node:assert"; import test, { describe } from "node:test"; import { isRelayFollowerData } from "./types.ts"; @@ -389,12 +390,13 @@ describe("MastodonRelay", () => { strictEqual(handlerCalled, true); ok(handlerActor); - // Verify follower was stored + // Verify follower was immediately accepted const followerData = await kv.get([ "follower", "https://remote.example.com/users/alice", ]); - ok(followerData); + ok(isRelayFollowerData(followerData)); + strictEqual(followerData.state, "accepted"); }); test("handles Follow activity with subscription rejection", async () => { @@ -688,6 +690,88 @@ describe("MastodonRelay", () => { ok(response.status === 200 || response.status === 202); }); + test("handles Announce activity forwarding", async () => { + const kv = new MemoryKvStore(); + const follower = new Person({ + id: new URL("https://follower.example.com/users/bob"), + preferredUsername: "bob", + inbox: new URL("https://follower.example.com/users/bob/inbox"), + }); + await kv.set( + ["follower", follower.id!.href], + { actor: await follower.toJsonLd(), state: "accepted" }, + ); + + const relay = createRelay("mastodon", { + kv, + origin: "https://relay.example.com", + documentLoaderFactory: () => mockDocumentLoader, + authenticatedDocumentLoaderFactory: () => mockDocumentLoader, + subscriptionHandler: () => Promise.resolve(true), + }); + + const announceActivity = new Announce({ + id: new URL("https://remote.example.com/activities/announce/1"), + actor: new URL("https://remote.example.com/users/alice"), + object: new URL("https://remote.example.com/notes/1"), + }); + const signedAnnounce = await signJsonLd( + await announceActivity.toJsonLd({ contextLoader: mockDocumentLoader }), + rsaKeyPair.privateKey, + rsaPublicKey.id, + { contextLoader: mockDocumentLoader }, + ); + + let request = new Request("https://relay.example.com/inbox", { + method: "POST", + headers: { + "Content-Type": "application/activity+json", + }, + body: JSON.stringify(signedAnnounce), + }); + + request = await signRequest( + request, + rsaKeyPair.privateKey, + rsaPublicKey.id, + ); + + const originalFetch = globalThis.fetch; + let deliveryMethod: string | undefined; + let deliveredActivity: unknown; + globalThis.fetch = (async ( + input: URL | RequestInfo, + init?: RequestInit, + ) => { + const outboundRequest = input instanceof Request + ? input + : new Request(input, init); + if ( + outboundRequest.url === + "https://follower.example.com/users/bob/inbox" + ) { + deliveryMethod = outboundRequest.method; + deliveredActivity = await outboundRequest.json(); + return new Response(null, { status: 202 }); + } + return originalFetch(input, init); + }) as typeof fetch; + + try { + const response = await relay.fetch(request); + ok( + response.status === 200 || response.status === 202, + `Unexpected inbox response status: ${response.status}`, + ); + } finally { + globalThis.fetch = originalFetch; + } + + ok(deliveredActivity, "Expected Announce delivery to the follower inbox"); + strictEqual(deliveryMethod, "POST"); + deepStrictEqual(deliveredActivity, signedAnnounce); + }); + test("ignores Follow activity without required fields", async () => { const kv = new MemoryKvStore(); diff --git a/packages/relay/src/mastodon.ts b/packages/relay/src/mastodon.ts index a1136be6a..fc5b99000 100644 --- a/packages/relay/src/mastodon.ts +++ b/packages/relay/src/mastodon.ts @@ -1,20 +1,6 @@ import type { InboxContext } from "@fedify/fedify"; -import { - Announce, - Create, - Delete, - Follow, - Move, - Undo, - Update, -} from "@fedify/vocab"; import { getLogger } from "@logtape/logtape"; -import { BaseRelay } from "./base.ts"; -import { - handleUndoFollow, - sendFollowResponse, - validateFollowActivity, -} from "./follow.ts"; +import { BaseRelay, type RelayableActivity } from "./base.ts"; import { RELAY_SERVER_ACTOR, type RelayOptions } from "./types.ts"; const logger = getLogger(["fedify", "relay", "mastodon"]); @@ -27,13 +13,14 @@ const logger = getLogger(["fedify", "relay", "mastodon"]); * @since 2.0.0 */ export class MastodonRelay extends BaseRelay { - async #forwardToFollowers( + protected readonly initialFollowerState = "accepted"; + protected readonly logger = logger; + + protected async deliverActivity( ctx: InboxContext, - activity: Create | Delete | Move | Update | Announce, + _activity: RelayableActivity, + excludeBaseUris: URL[], ): Promise { - const sender = await activity.getActor(ctx); - const excludeBaseUris = sender?.id ? [new URL(sender.id)] : []; - await ctx.forwardActivity( { identifier: RELAY_SERVER_ACTOR }, "followers", @@ -44,55 +31,4 @@ export class MastodonRelay extends BaseRelay { }, ); } - - protected setupInboxListeners(): void { - if (this.federation != null) { - this.federation.setInboxListeners("/users/{identifier}/inbox", "/inbox") - .on(Follow, async (ctx, follow) => { - const follower = await validateFollowActivity(ctx, follow); - if (!follower || !follower.id) return; - - const approved = await this.options.subscriptionHandler( - ctx, - follower, - ); - - if (approved) { - // Mastodon-specific: immediately add to followers list with accepted state - await ctx.data.kv.set( - ["follower", follower.id.href], - { actor: await follower.toJsonLd(), state: "accepted" }, - ); - } - - await sendFollowResponse(ctx, follow, follower, approved); - }) - .on( - Undo, - async (ctx, undo) => await handleUndoFollow(ctx, undo, logger), - ) - .on( - Create, - async (ctx, create) => await this.#forwardToFollowers(ctx, create), - ) - .on( - Delete, - async (ctx, deleteActivity) => - await this.#forwardToFollowers(ctx, deleteActivity), - ) - .on( - Move, - async (ctx, move) => await this.#forwardToFollowers(ctx, move), - ) - .on( - Update, - async (ctx, update) => await this.#forwardToFollowers(ctx, update), - ) - .on( - Announce, - async (ctx, announce) => - await this.#forwardToFollowers(ctx, announce), - ); - } - } } diff --git a/packages/relay/src/types.ts b/packages/relay/src/types.ts index fd906bed7..a2dd63468 100644 --- a/packages/relay/src/types.ts +++ b/packages/relay/src/types.ts @@ -1,5 +1,5 @@ import type { Context, KvStore, MessageQueue } from "@fedify/fedify"; -import type { Actor } from "@fedify/vocab"; +import { type Actor, isActor, Object as APObject } from "@fedify/vocab"; import type { AuthenticatedDocumentLoaderFactory, DocumentLoaderFactory, @@ -12,6 +12,13 @@ export const RELAY_SERVER_ACTOR = "relay"; */ export type RelayType = "mastodon" | "litepub"; +/** + * A follower's subscription state. + * + * @internal + */ +export type RelayFollowerState = "pending" | "accepted"; + /** * Handler for subscription requests (Follow/Undo activities). */ @@ -79,7 +86,7 @@ export interface RelayFollowerData { /** The actor's JSON-LD representation (serialized for storage). */ readonly actor: unknown; /** The follower's state. */ - readonly state: "pending" | "accepted"; + readonly state: RelayFollowerState; } /** @@ -94,7 +101,7 @@ export interface RelayFollower { /** The validated Actor object. */ readonly actor: Actor; /** The follower's state. */ - readonly state: "pending" | "accepted"; + readonly state: RelayFollowerState; } /** @@ -162,3 +169,26 @@ export function isRelayFollowerData( (obj.state === "pending" || obj.state === "accepted") ); } + +/** + * Parses and semantically validates follower data from storage. + * + * @param actorId The actor ID used as the follower's storage key. + * @param value The stored follower data. + * @returns The parsed follower, or `null` if the row is invalid. + * @internal + */ +export async function parseRelayFollowerData( + actorId: string, + value: unknown, +): Promise { + if (!isRelayFollowerData(value)) return null; + + try { + const actor = await APObject.fromJsonLd(value.actor); + if (!isActor(actor) || actor.id?.href !== actorId) return null; + return { actorId, actor, state: value.state }; + } catch { + return null; + } +}