diff --git a/packages/relay/src/base.ts b/packages/relay/src/base.ts index d7a18e075..8161b0ddc 100644 --- a/packages/relay/src/base.ts +++ b/packages/relay/src/base.ts @@ -1,13 +1,40 @@ -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 { + type Actor, + Announce, + Create, + Delete, + Follow, + isActor, + Move, + Object as APObject, + Undo, + Update, +} from "@fedify/vocab"; +import type { Logger } from "@logtape/logtape"; +import { + handleUndoFollow, + sendFollowResponse, + validateFollowActivity, +} from "./follow.ts"; import { isRelayFollowerData, 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 +46,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, @@ -126,11 +156,97 @@ export abstract class BaseRelay implements Relay { return await this.parseFollowerData(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 sender = await activity.getActor(ctx); + const excludeBaseUris = sender?.id == null ? [] : [new URL(sender.id)]; + 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/litepub.test.ts b/packages/relay/src/litepub.test.ts index 899bf6d8b..a8955e1e9 100644 --- a/packages/relay/src/litepub.test.ts +++ b/packages/relay/src/litepub.test.ts @@ -583,64 +583,64 @@ describe("LitePubRelay", () => { strictEqual(followerData.state, "accepted"); }); - test("handles Undo Follow activity", async () => { - const kv = new MemoryKvStore(); + for (const state of ["pending", "accepted"] as const) { + test(`handles Undo Follow activity for ${state} follower`, async () => { + const kv = new MemoryKvStore(); - // Pre-populate with an accepted follower - 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 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: "accepted" }, - ); + 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 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 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, - }); + 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 }), - ), - }); + 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, - ); + request = await signRequest( + request, + rsaKeyPair.privateKey, + rsaPublicKey.id, + ); - await relay.fetch(request); + await relay.fetch(request); - // Verify follower was removed - const followerData = await kv.get(["follower", followerId]); - strictEqual(followerData, undefined); - }); + 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..653d7c08b 100644 --- a/packages/relay/src/litepub.ts +++ b/packages/relay/src/litepub.ts @@ -1,23 +1,14 @@ -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 { RELAY_SERVER_ACTOR, type RelayFollowerData, @@ -34,13 +25,43 @@ 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, + follower: Actor, + ): Promise { + if (follower.id == null) return true; + const existingFollow = await ctx.data.kv.get([ + "follower", + follower.id.href, + ]); + return existingFollow?.state === "pending"; + } + + protected override async afterFollowApproved( ctx: InboxContext, - activity: Create | Delete | Move | Update | Announce, + 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 +81,36 @@ 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(); + 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 (followerData == 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 = { ...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..3c1d27b0a 100644 --- a/packages/relay/src/mastodon.test.ts +++ b/packages/relay/src/mastodon.test.ts @@ -2,6 +2,7 @@ import { MemoryKvStore, signRequest } from "@fedify/fedify"; import { createRelay, type RelayOptions } from "@fedify/relay"; import { + Announce, Create, Delete, Follow, @@ -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,45 @@ describe("MastodonRelay", () => { ok(response.status === 200 || response.status === 202); }); + test("handles Announce activity forwarding", async () => { + const kv = new MemoryKvStore(); + + 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"), + }); + + let request = new Request("https://relay.example.com/inbox", { + method: "POST", + headers: { + "Content-Type": "application/activity+json", + }, + body: JSON.stringify( + await announceActivity.toJsonLd({ contextLoader: mockDocumentLoader }), + ), + }); + + request = await signRequest( + request, + rsaKeyPair.privateKey, + rsaPublicKey.id, + ); + + const response = await relay.fetch(request); + + // Verify the request was accepted + ok(response.status === 200 || response.status === 202); + }); + 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..da1671dcb 100644 --- a/packages/relay/src/types.ts +++ b/packages/relay/src/types.ts @@ -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; } /**