From fe651b2761674eefc0f25084986cb985d9c81ba5 Mon Sep 17 00:00:00 2001 From: Hong Minhee Date: Fri, 2 Oct 2026 23:09:22 +0900 Subject: [PATCH] Preserve follows across account migrations Follow an alias-verified destination when a followed account sends a push-mode Move, then unfollow the origin and notify onFolloweeMove handlers. Route migrations to all following bots, including dynamic groups, and retry transient document failures without trusting embeds. Allow Follow, Accept, Reject, and Undo delivery between sibling bots so migrations to a local actor complete through real HTTP inboxes. Cover validation, retries, duplicates, callbacks, and delivery in both runtimes, and document the event's request and acceptance semantics. Design and implementation were AI-assisted, with independent design and code reviews. Validated with mise run test and mise run docs:build. Fixes https://github.com/fedify-dev/botkit/issues/49 Assisted-by: OpenCode:deepseek-flash Assisted-by: Codex:gpt-6.1-sol Assisted-by: Claude Code:claude-fable-5-1 Assisted-by: Claude Code:claude-opus-5-5 --- CHANGES.md | 8 + changes.d/botkit/followee-move.md | 8 + changes.d/botkit/local-follows.md | 8 + docs/concepts/events.md | 60 +- docs/concepts/session.md | 5 + packages/botkit/src/bot-impl.ts | 11 + packages/botkit/src/bot.ts | 16 + packages/botkit/src/events.ts | 16 + packages/botkit/src/follow-impl.ts | 5 +- packages/botkit/src/follow-local.test.ts | 239 ++++++++ packages/botkit/src/follow-move.test.ts | 681 +++++++++++++++++++++++ packages/botkit/src/instance-impl.ts | 266 +++++++++ packages/botkit/src/session-impl.ts | 5 +- packages/botkit/src/uri.ts | 19 + 14 files changed, 1342 insertions(+), 5 deletions(-) create mode 100644 changes.d/botkit/followee-move.md create mode 100644 changes.d/botkit/local-follows.md create mode 100644 packages/botkit/src/follow-local.test.ts create mode 100644 packages/botkit/src/follow-move.test.ts diff --git a/CHANGES.md b/CHANGES.md index 2c4d3c2..ab79070 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -12,6 +12,9 @@ To be released. existing accounts can move their followers to a BotKit bot. Actor URIs listed in the option are published as `alsoKnownAs`, and are available through `Bot.aliases` and `Session.bot.aliases`. [[#48], [#55]] + - Added automatic re-following when an account a bot follows moves to a + verified new account, with an `onFolloweeMove` event after submitting the + new follow request and unfollowing the old account. [[#49], [#56]] - Added names, inline summaries, and custom URLs when publishing or updating messages. Updated messages can remove these fields by setting them to `null`, and the new `inline` template composes paragraphless rich text for @@ -27,6 +30,9 @@ To be released. - Fixed a bug where a remote server could approve a quote with a quote authorization stamp other than the one named in its `Accept` activity, as long as the substituted stamp was on the same origin. [[#52], [#53]] + - Fixed delivery of follow requests, acceptances, rejections, and unfollows + between bots hosted on the same instance, including follow requests to + account migration targets on that instance. [[#49], [#56]] - Upgraded Fedify to 2.4.0, which adds support for [FEP-ef61] portable objects and hardens HTTP Signature verification and document loading. @@ -38,9 +44,11 @@ To be released. [#42]: https://github.com/fedify-dev/botkit/pull/42 [#43]: https://github.com/fedify-dev/botkit/pull/43 [#48]: https://github.com/fedify-dev/botkit/issues/48 +[#49]: https://github.com/fedify-dev/botkit/issues/49 [#52]: https://github.com/fedify-dev/botkit/issues/52 [#53]: https://github.com/fedify-dev/botkit/pull/53 [#55]: https://github.com/fedify-dev/botkit/pull/55 +[#56]: https://github.com/fedify-dev/botkit/pull/56 Version 0.5.6 diff --git a/changes.d/botkit/followee-move.md b/changes.d/botkit/followee-move.md new file mode 100644 index 0000000..9e953b0 --- /dev/null +++ b/changes.d/botkit/followee-move.md @@ -0,0 +1,8 @@ +--- +links: + '#49': https://github.com/fedify-dev/botkit/issues/49 + '#56': https://github.com/fedify-dev/botkit/pull/56 +--- + - Added automatic re-following when an account a bot follows moves to a + verified new account, with an `onFolloweeMove` event after submitting the + new follow request and unfollowing the old account. [[#49], [#56]] diff --git a/changes.d/botkit/local-follows.md b/changes.d/botkit/local-follows.md new file mode 100644 index 0000000..1411590 --- /dev/null +++ b/changes.d/botkit/local-follows.md @@ -0,0 +1,8 @@ +--- +links: + '#49': https://github.com/fedify-dev/botkit/issues/49 + '#56': https://github.com/fedify-dev/botkit/pull/56 +--- + - Fixed delivery of follow requests, acceptances, rejections, and unfollows + between bots hosted on the same instance, including follow requests to + account migration targets on that instance. [[#49], [#56]] diff --git a/docs/concepts/events.md b/docs/concepts/events.md index 8311811..80e51ea 100644 --- a/docs/concepts/events.md +++ b/docs/concepts/events.md @@ -23,7 +23,7 @@ bot.onMention = async (session, message) => { ~~~~ Every event handler receives a [session](./session.md) object as the first -argument, and the event-specific object as the second argument. +argument, followed by the event-specific objects. BotKit invokes these event handlers only for activities whose signatures were verified by Fedify. As of Fedify 2.1.0, BotKit also acknowledges certain @@ -152,6 +152,64 @@ bot.onRejectFollow = async (session, rejecter) => { ~~~~ +Followee move +------------- + +*This event is available since BotKit 0.6.0.* + +When an account your bot follows moves, BotKit automatically submits a follow +request to its new account and unfollows the old one. It accepts push-mode +`Move` activities sent by the old account only, and verifies that the new +account lists the old actor URI in its `alsoKnownAs`. A target embedded in +an activity is checked against the target account's own actor document. + +The `~Bot.onFolloweeMove` handler receives the bot's session, the old `Actor`, +and the new `Actor`, in that order. It runs after the new request is submitted +and the old account is unfollowed. The new request may still await acceptance; +`~Bot.onAcceptFollow` or `~Bot.onRejectFollow` reports the eventual response. +With an outgoing queue, submitting the request means enqueueing it, rather +than completing delivery. + +~~~~ typescript twoslash +import type { Bot } from "@fedify/botkit"; +declare const bot: Bot; +// ---cut-before--- +bot.onFolloweeMove = (session, oldActor, newActor) => { + console.info( + session.bot.identifier, + "followed account moved", + oldActor.id?.href, + newActor.id?.href, + ); +}; +~~~~ + +You can also supply this handler as the `onFolloweeMove` option to +`createBot()`, or assign it to a +[dynamic bot group](./instance.md#dynamic-bots). The `FolloweeMoveEventHandler` +type is exported by *@fedify/botkit*. + +Only accepted follows of the old account are migrated. If the bot already +follows the target, it keeps that follow and only unfollows the old account. +A repeat delivery does nothing once the old follow has been removed. A pending +request to the target may receive another follow request, since it is not yet +an accepted follow. + +There is no migration policy option. A handler can unfollow an already accepted +target with `~Session.unfollow()`; to decline a target whose request is still +pending, unfollow it from `~Bot.onAcceptFollow`. `~Session.unfollow()` does not +cancel pending requests. A rejected target leaves the bot following neither +account. + +> [!NOTE] +> These changes are not a transaction across servers. Failure to submit the +> new request preserves the old follow, but an outgoing queue's later delivery +> failure cannot restore it. The old follow is removed before its `Undo` is +> submitted; if that submission fails, the event does not run, and a repeated +> `Move` does not retry the `Undo`. Event handlers are likewise not replayed +> after the old follow has been removed. + + Mention ------- diff --git a/docs/concepts/session.md b/docs/concepts/session.md index 7ab5e69..5ef7396 100644 --- a/docs/concepts/session.md +++ b/docs/concepts/session.md @@ -255,6 +255,11 @@ bot.onFollow = async (session, followRequest) => { > If you try to follow an actor that is already followed, the method will just > do nothing. +When an account the bot follows moves to another account, BotKit automatically +submits a follow request to the verified target and unfollows the old account. +See [the followee move event](./events.md#followee-move) for validation rules, +request timing, and the `onFolloweeMove` callback. + Unfollowing an actor -------------------- diff --git a/packages/botkit/src/bot-impl.ts b/packages/botkit/src/bot-impl.ts index 468a335..d44d77c 100644 --- a/packages/botkit/src/bot-impl.ts +++ b/packages/botkit/src/bot-impl.ts @@ -77,6 +77,7 @@ import { } from "./emoji.ts"; import type { AcceptEventHandler, + FolloweeMoveEventHandler, FollowEventHandler, LikeEventHandler, MentionEventHandler, @@ -162,6 +163,7 @@ export interface BotImplOptions */ export const botEventHandlerNames = [ "onFollow", + "onFolloweeMove", "onUnfollow", "onAcceptFollow", "onRejectFollow", @@ -239,6 +241,7 @@ export class BotImpl implements Bot { } onFollow?: FollowEventHandler; + onFolloweeMove?: FolloweeMoveEventHandler; onUnfollow?: UnfollowEventHandler; onAcceptFollow?: AcceptEventHandler; onRejectFollow?: RejectEventHandler; @@ -258,6 +261,7 @@ export class BotImpl implements Bot { onVote?: VoteEventHandler; constructor(options: BotImplOptions) { + this.onFolloweeMove = options.onFolloweeMove; this.identifier = options.identifier ?? "bot"; this.class = options.class ?? Service; this.username = options.username; @@ -2225,6 +2229,12 @@ export function wrapBotImpl( set onFollow(value) { bot.onFollow = value; }, + get onFolloweeMove() { + return bot.onFolloweeMove; + }, + set onFolloweeMove(value) { + bot.onFolloweeMove = value; + }, get onUnfollow() { return bot.onUnfollow; }, @@ -2679,6 +2689,7 @@ export class BotGroupImpl implements BotGroup { ) => string | null | Promise; onFollow?: FollowEventHandler; + onFolloweeMove?: FolloweeMoveEventHandler; onUnfollow?: UnfollowEventHandler; onAcceptFollow?: AcceptEventHandler; onRejectFollow?: RejectEventHandler; diff --git a/packages/botkit/src/bot.ts b/packages/botkit/src/bot.ts index c922913..71cf64f 100644 --- a/packages/botkit/src/bot.ts +++ b/packages/botkit/src/bot.ts @@ -25,6 +25,7 @@ import { BotImpl, wrapBotImpl } from "./bot-impl.ts"; import type { CustomEmoji, DeferredCustomEmoji } from "./emoji.ts"; import type { AcceptEventHandler, + FolloweeMoveEventHandler, FollowEventHandler, LikeEventHandler, MentionEventHandler, @@ -63,6 +64,14 @@ export interface BotEventHandlers { */ onFollow?: FollowEventHandler; + /** + * Invoked after submitting a follow request to a followed actor's verified + * migration target and unfollowing the old actor. The new request may + * still await acceptance. + * @since 0.6.0 + */ + onFolloweeMove?: FolloweeMoveEventHandler; + /** * An event handler for an unfollow event from the bot. */ @@ -339,6 +348,13 @@ export interface BotWithVoidContextData extends Bot { * Options for creating a bot. */ export interface CreateBotOptions { + /** + * The handler invoked after a followed actor moves to a verified target. + * It can also be assigned through {@link Bot.onFolloweeMove} afterwards. + * @since 0.6.0 + */ + readonly onFolloweeMove?: FolloweeMoveEventHandler; + /** * The internal identifier of the bot. Since it is used for the actor URI, * it *should not* be changed after the bot is federated. diff --git a/packages/botkit/src/events.ts b/packages/botkit/src/events.ts index 1535eaa..be6c966 100644 --- a/packages/botkit/src/events.ts +++ b/packages/botkit/src/events.ts @@ -37,6 +37,22 @@ export type FollowEventHandler = ( followRequest: FollowRequest, ) => void | Promise; +/** + * An event handler invoked after a followed actor moves to another account. + * The new follow request has been submitted, but may still await acceptance. + * @typeParam TContextData The type of the context data. + * @param session The session of the bot. + * @param oldActor The actor the bot previously followed. + * @param newActor The actor to which the account moved. + * @returns Nothing, or a promise that resolves when handling completes. + * @since 0.6.0 + */ +export type FolloweeMoveEventHandler = ( + session: Session, + oldActor: Actor, + newActor: Actor, +) => void | Promise; + /** * An event handler for an unfollow event from the bot. * @typeParam TContextData The type of the context data. diff --git a/packages/botkit/src/follow-impl.ts b/packages/botkit/src/follow-impl.ts index 7f33709..69b9002 100644 --- a/packages/botkit/src/follow-impl.ts +++ b/packages/botkit/src/follow-impl.ts @@ -16,6 +16,7 @@ import { Accept, type Actor, type Follow, Reject } from "@fedify/vocab"; import type { FollowRequest } from "./follow.ts"; import type { SessionImpl } from "./session-impl.ts"; +import { getFollowDeliveryOptions } from "./uri.ts"; export class FollowRequestImpl implements FollowRequest { readonly session: SessionImpl; @@ -58,7 +59,7 @@ export class FollowRequestImpl implements FollowRequest { to: this.follower.id, object: this.raw, }), - { excludeBaseUris: [new URL(this.session.context.origin)] }, + getFollowDeliveryOptions(this.session.context, this.follower.id), ); await this.session.bot.repository.addFollower(this.id, this.follower); this.#state = "accepted"; @@ -77,7 +78,7 @@ export class FollowRequestImpl implements FollowRequest { to: this.follower.id, object: this.raw, }), - { excludeBaseUris: [new URL(this.session.context.origin)] }, + getFollowDeliveryOptions(this.session.context, this.follower.id), ); this.#state = "rejected"; } diff --git a/packages/botkit/src/follow-local.test.ts b/packages/botkit/src/follow-local.test.ts new file mode 100644 index 0000000..7065516 --- /dev/null +++ b/packages/botkit/src/follow-local.test.ts @@ -0,0 +1,239 @@ +// BotKit by Fedify: A framework for creating ActivityPub bots +// Copyright (C) 2025–2026 Hong Minhee +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as +// published by the Free Software Foundation, either version 3 of the +// License, or (at your option) any later version. +// +// This program is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . +import { MemoryKvStore } from "@fedify/fedify/federation"; +import { Follow, Move, Person } from "@fedify/vocab"; +import assert from "node:assert/strict"; +import { + createServer, + type IncomingMessage, + type ServerResponse, +} from "node:http"; +import { test } from "node:test"; +import { createInstance, type Instance } from "./instance.ts"; +import { MemoryRepository } from "./repository.ts"; + +async function respond( + request: IncomingMessage, + response: ServerResponse, + instance: Instance, + origin: URL, + signal?: AbortSignal, +): Promise { + signal?.throwIfAborted(); + const headers = new Headers(); + for (const [key, value] of Object.entries(request.headers)) { + if (value != null) { + headers.set(key, Array.isArray(value) ? value.join(", ") : value); + } + } + const chunks: Uint8Array[] = []; + for await (const chunk of request) { + if (!(chunk instanceof Uint8Array)) { + throw new TypeError("Expected request bytes."); + } + chunks.push(chunk); + } + const body = new Uint8Array( + chunks.reduce((length, chunk) => length + chunk.length, 0), + ); + let offset = 0; + for (const chunk of chunks) { + body.set(chunk, offset); + offset += chunk.length; + } + const result = await instance.fetch( + new Request(new URL(request.url ?? "/", origin), { + method: request.method, + headers, + body: body.length > 0 ? body : undefined, + signal, + }), + ); + response.writeHead(result.status, Object.fromEntries(result.headers)); + response.end(new Uint8Array(await result.arrayBuffer())); +} + +async function createTestServer(signal?: AbortSignal) { + const repository = new MemoryRepository(); + const instance = createInstance({ + kv: new MemoryKvStore(), + repository, + federationOptions: { allowPrivateAddress: true }, + }); + let origin = new URL("http://127.0.0.1"); + const errors: unknown[] = []; + const server = createServer((request, response) => { + void respond(request, response, instance, origin, signal).catch( + (error: unknown) => { + errors.push(error); + response.writeHead(500); + response.end(); + }, + ); + }); + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(0, "127.0.0.1", () => { + server.off("error", reject); + resolve(); + }); + }); + const address = server.address(); + assert.ok(address != null && typeof address !== "string"); + origin = new URL(`http://127.0.0.1:${address.port}`); + return { + instance, + repository, + origin, + errors, + close: () => + new Promise((resolve, reject) => { + server.close((error) => error == null ? resolve() : reject(error)); + server.closeAllConnections(); + }), + }; +} + +test("signed Move migrates a follow to a sibling bot through real HTTP inboxes", async (t) => { + const remote = await createTestServer(t.signal); + const local = await createTestServer(t.signal); + try { + const oldBot = remote.instance.createBot("old", { username: "old" }); + const oldSession = oldBot.getSession(remote.origin); + const oldActor = await oldSession.getActor(); + const target = local.instance.createBot("target", { + username: "target", + aliases: [oldSession.actorId], + }); + const follower = local.instance.createBot("follower", { + username: "follower", + }); + const followerSession = follower.getSession(local.origin); + const followerActor = await followerSession.getActor(); + const followId = followerSession.context.getObjectUri(Follow, { + identifier: follower.identifier, + id: "018f6db5-27d2-7000-8000-000000000003", + }); + await local.repository.addFollowee( + follower.identifier, + oldSession.actorId, + new Follow({ + id: followId, + actor: followerSession.actorId, + object: oldSession.actorId, + }), + ); + await remote.repository.addFollower( + oldBot.identifier, + followId, + followerActor, + ); + const moved: string[] = []; + follower.onFolloweeMove = (session, oldActor, newActor) => { + assert.deepStrictEqual(oldActor.id, oldSession.actorId); + assert.deepStrictEqual( + newActor.id, + target.getSession(local.origin).actorId, + ); + moved.push(session.bot.identifier); + }; + await oldSession.context.sendActivity( + { identifier: oldBot.identifier }, + followerActor, + new Move({ + id: new URL("#move", oldSession.actorId), + actor: new Person({ id: oldActor.id }), + object: oldSession.actorId, + target: target.getSession(local.origin).actorId, + to: followerSession.actorId, + }), + ); + const targetId = target.getSession(local.origin).actorId; + assert.ok( + await local.repository.getFollowee(follower.identifier, targetId), + ); + assert.ok( + await local.repository.hasFollower( + target.identifier, + followerSession.actorId, + ), + ); + assert.ok( + !await local.repository.getFollowee( + follower.identifier, + oldSession.actorId, + ), + ); + assert.ok( + !await remote.repository.hasFollower( + oldBot.identifier, + followerSession.actorId, + ), + ); + assert.deepStrictEqual(moved, ["follower"]); + assert.deepStrictEqual(local.errors, []); + assert.deepStrictEqual(remote.errors, []); + + // Undo to a sibling must also reach its actual inbox. + await followerSession.unfollow( + await target.getSession(local.origin).getActor(), + ); + assert.ok( + !await local.repository.hasFollower( + target.identifier, + followerSession.actorId, + ), + ); + } finally { + await Promise.all([remote.close(), local.close()]); + } +}); + +test("a sibling's rejection reaches the sender through real HTTP inboxes", async (t) => { + const local = await createTestServer(t.signal); + try { + const target = local.instance.createBot("target", { + username: "target", + followerPolicy: "reject", + }); + const follower = local.instance.createBot("follower", { + username: "follower", + }); + let rejected = 0; + follower.onRejectFollow = () => { + rejected++; + }; + await follower.getSession(local.origin).follow( + await target.getSession(local.origin).getActor(), + ); + assert.strictEqual(rejected, 1); + assert.ok( + !await local.repository.hasFollower( + target.identifier, + follower.getSession(local.origin).actorId, + ), + ); + assert.ok( + !await local.repository.getFollowee( + follower.identifier, + target.getSession(local.origin).actorId, + ), + ); + assert.deepStrictEqual(local.errors, []); + } finally { + await local.close(); + } +}); diff --git a/packages/botkit/src/follow-move.test.ts b/packages/botkit/src/follow-move.test.ts new file mode 100644 index 0000000..06616a3 --- /dev/null +++ b/packages/botkit/src/follow-move.test.ts @@ -0,0 +1,681 @@ +// BotKit by Fedify: A framework for creating ActivityPub bots +// Copyright (C) 2025–2026 Hong Minhee +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as +// published by the Free Software Foundation, either version 3 of the +// License, or (at your option) any later version. +// +// This program is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . +import { type InboxContext, MemoryKvStore } from "@fedify/fedify/federation"; +import { + Accept, + type Activity, + Follow, + Move, + Note, + Person, + Undo, +} from "@fedify/vocab"; +import { + type DocumentLoader, + FetchError, + UrlError, +} from "@fedify/vocab-runtime"; +import assert from "node:assert/strict"; +import { test } from "node:test"; +import { BotImpl } from "./bot-impl.ts"; +import { createBot } from "./bot.ts"; +import { InstanceImpl } from "./instance-impl.ts"; +import { MemoryRepository } from "./repository.ts"; + +const oldId = new URL("https://old.example/users/alice"); +const newId = new URL("https://new.example/users/alice"); +const oldActor = new Person({ id: oldId, inbox: new URL("inbox", oldId) }); +const newActor = new Person({ + id: newId, + inbox: new URL("inbox", newId), + aliases: [oldId], +}); +const move = new Move({ + id: new URL("#move", oldId), + actor: oldActor, + object: oldId, + target: newId, +}); + +async function harness(signal?: AbortSignal) { + signal?.throwIfAborted(); + const repository = new MemoryRepository(); + const instance = new InstanceImpl({ + kv: new MemoryKvStore(), + repository, + }); + const alpha = instance.createBot("alpha", { + username: "alpha", + aliases: [oldId], + }); + const beta = instance.createBot("beta", { username: "beta" }); + const ctx = instance.federation.createContext( + new URL("https://example.com"), + undefined, + ) as InboxContext; + Object.defineProperty(ctx, "recipient", { value: null, configurable: true }); + const sent: Activity[] = []; + ctx.sendActivity = (_sender, _recipient, activity) => { + sent.push(activity); + return Promise.resolve(); + }; + const loads: string[] = []; + const documents = new Map([ + [newId.href, await newActor.toJsonLd()], + [oldId.href, await oldActor.toJsonLd()], + ]); + const loader: DocumentLoader = (url) => { + loads.push(url); + return Promise.resolve({ + contextUrl: null, + documentUrl: url, + document: documents.get(url), + }); + }; + Object.defineProperty(ctx, "getDocumentLoader", { + value: () => Promise.resolve(loader), + configurable: true, + }); + for (const bot of [alpha, beta]) { + await repository.addFollowee( + bot.identifier, + oldId, + new Follow({ + id: ctx.getObjectUri(Follow, { + identifier: bot.identifier, + id: "018f6db5-27d2-7000-8000-000000000001", + }), + actor: ctx.getActorUri(bot.identifier), + object: oldId, + to: oldId, + }), + ); + } + return { instance, repository, alpha, beta, ctx, sent, loads, documents }; +} + +test("a verified push Move migrates all followed bots and emits events", async (t) => { + const { instance, repository, alpha, beta, ctx, sent, loads } = await harness( + t.signal, + ); + const moved: string[] = []; + for (const bot of [alpha, beta]) { + bot.onFolloweeMove = async (session, oldActor, newActor) => { + assert.deepStrictEqual(oldActor.id, oldId); + assert.deepStrictEqual(newActor.id, newId); + assert.ok(!await repository.getFollowee(bot.identifier, oldId)); + assert.ok(!await session.follows(newActor)); // Still awaiting Accept. + moved.push(session.bot.identifier); + }; + } + await instance.onMoved(ctx, move, t.signal); + assert.deepStrictEqual(moved, ["alpha", "beta"]); + assert.deepStrictEqual(sent.map((a) => a.constructor), [ + Follow, + Undo, + Follow, + Undo, + ]); + assert.deepStrictEqual(sent[0].objectId, newId); + assert.deepStrictEqual(loads.filter((url) => url === newId.href), [ + newId.href, + ]); + await instance.onMoved(ctx, move, t.signal); + assert.strictEqual(sent.length, 4); + assert.strictEqual(moved.length, 2); +}); + +for (const recipient of ["alpha", "unrelated"]) { + test(`personal inbox ${recipient} migrates every bot following the origin`, async (t) => { + const h = await harness(t.signal); + Object.defineProperty(h.ctx, "recipient", { value: recipient }); + await h.instance.onMoved(h.ctx, move, t.signal); + assert.strictEqual(h.sent.length, 4); + assert.ok(!await h.repository.getFollowee("alpha", oldId)); + assert.ok(!await h.repository.getFollowee("beta", oldId)); + }); +} + +test("unrelated bots and missing dynamic bots cause no target fetch", async (t) => { + const h = await harness(t.signal); + await h.repository.removeFollowee("alpha", oldId); + await h.repository.removeFollowee("beta", oldId); + await h.repository.addFollowee( + "missing", + oldId, + new Follow({ object: oldId }), + ); + await h.instance.onMoved(h.ctx, move, t.signal); + assert.deepStrictEqual(h.sent, []); + assert.deepStrictEqual(h.loads, []); +}); + +const invalidMoves = [ + ["missing id", new Move({ actor: oldActor, object: oldId, target: newId })], + ["missing actor", new Move({ id: move.id, object: oldId, target: newId })], + ["missing object", new Move({ id: move.id, actor: oldActor, target: newId })], + ["missing target", new Move({ id: move.id, actor: oldActor, object: oldId })], + ["pull mode", move.clone({ actor: newActor })], + ["different object", move.clone({ object: newId })], + ["same target", move.clone({ target: oldId })], + ["non-HTTP target", move.clone({ target: new URL("urn:alice") })], + ["multiple actors", move.clone({ actors: [oldActor, newActor] })], + ["multiple objects", move.clone({ objects: [oldId, newId] })], + ["multiple targets", move.clone({ targets: [newId, oldId] })], +] as const; +for (const [name, invalid] of invalidMoves) { + test(`ignores Move with ${name} before lookup`, async (t) => { + const h = await harness(t.signal); + await h.instance.onMoved(h.ctx, invalid, t.signal); + assert.deepStrictEqual(h.sent, []); + assert.deepStrictEqual(h.loads, []); + assert.ok(await h.repository.getFollowee("alpha", oldId)); + }); +} + +const invalidTargets = [ + ["no aliases", newActor.clone({ aliases: [] })], + [ + "wrong alias", + newActor.clone({ aliases: [new URL("https://old.example/@alice")] }), + ], + [ + "wrong id", + newActor.clone({ id: new URL("https://new.example/users/bob") }), + ], + ["missing id", new Person({ inbox: newActor.inboxId, aliases: [oldId] })], + ["missing inbox", new Person({ id: newId, aliases: [oldId] })], + ["non-actor", new Note({ id: newId })], +] as const; +for (const [name, actor] of invalidTargets) { + test(`ignores target with ${name}, even if embedded target claims valid aliases`, async (t) => { + const h = await harness(t.signal); + h.documents.set(newId.href, await actor.toJsonLd()); + await h.instance.onMoved(h.ctx, move.clone({ target: newActor }), t.signal); + assert.deepStrictEqual(h.sent, []); + assert.ok(await h.repository.getFollowee("alpha", oldId)); + assert.deepStrictEqual(h.loads, [newId.href]); + }); +} + +for (const doc of [null, { "@context": { "@vocab": 42 }, type: "Person" }]) { + test(`ignores malformed target JSON-LD (${JSON.stringify(doc)})`, async (t) => { + const h = await harness(t.signal); + h.documents.set(newId.href, doc); + await h.instance.onMoved(h.ctx, move, t.signal); + assert.deepStrictEqual(h.sent, []); + assert.ok(await h.repository.getFollowee("alpha", oldId)); + }); +} + +test("ignores target document redirected to a different origin", async (t) => { + const h = await harness(t.signal); + Object.defineProperty(h.ctx, "getDocumentLoader", { + value: () => + Promise.resolve(() => + Promise.resolve({ + documentUrl: "https://untrusted.example/actor", + contextUrl: null, + document: h.documents.get(newId.href), + }) + ), + }); + await h.instance.onMoved(h.ctx, move, t.signal); + assert.deepStrictEqual(h.sent, []); +}); + +const fetchErrors = [ + ["non-JSON response", new SyntaxError("Invalid JSON."), true], + ["disallowed URL", new UrlError("Disallowed URL."), true], + ["DNS", new UrlError("DNS failed.", { reason: "dns" }), false], + ["network", new TypeError("Network failed."), false], + ["timeout", new FetchError(newId, "Timeout."), false], + ...[200, 301, 304, 403, 404, 410, 408, 429, 500].map((status) => + [ + `HTTP ${status}`, + new FetchError(newId, "Fetch failed.", new Response(null, { status })), + status < 400 || status < 500 && status !== 408 && status !== 429, + ] as const + ), +] as const; +for (const [name, error, permanent] of fetchErrors) { + test(`${permanent ? "ignores" : "retries"} target fetch failure: ${name}`, async (t) => { + const h = await harness(t.signal); + Object.defineProperty(h.ctx, "getDocumentLoader", { + value: () => Promise.resolve(() => Promise.reject(error)), + configurable: true, + }); + if (permanent) await h.instance.onMoved(h.ctx, move, t.signal); + else await assert.rejects(h.instance.onMoved(h.ctx, move, t.signal), error); + assert.deepStrictEqual(h.sent, []); + assert.ok(await h.repository.getFollowee("alpha", oldId)); + }); +} + +test("context loader failure is retried rather than mistaken for malformed JSON-LD", async (t) => { + const h = await harness(t.signal); + const error = new TypeError("Context fetch failed."); + h.documents.set(newId.href, { + "@context": "https://context.example/context", + id: newId.href, + }); + Object.defineProperty(h.ctx, "contextLoader", { + value: () => Promise.reject(error), + }); + await assert.rejects(h.instance.onMoved(h.ctx, move, t.signal), error); + assert.deepStrictEqual(h.sent, []); +}); + +test("already accepted destination is retained without a second Follow", async (t) => { + const h = await harness(t.signal); + for (const bot of [h.alpha, h.beta]) { + await h.repository.addFollowee( + bot.identifier, + newId, + new Follow({ + actor: h.ctx.getActorUri(bot.identifier), + object: newId, + }), + ); + } + await h.instance.onMoved(h.ctx, move, t.signal); + assert.deepStrictEqual(h.sent.map((a) => a.constructor), [Undo, Undo]); + assert.ok(await h.repository.getFollowee("alpha", newId)); +}); + +test("the destination's later Accept completes the new follow", async (t) => { + const h = await harness(t.signal); + await h.instance.onMoved(h.ctx, move, t.signal); + const follow = h.sent[0]; + assert.ok(follow instanceof Follow); + assert.ok(!await h.repository.getFollowee("alpha", newId)); + Object.defineProperty(h.ctx, "documentLoader", { + value: (url: string) => + Promise.resolve({ + contextUrl: null, + documentUrl: url, + document: h.documents.get(url), + }), + }); + await h.instance.onFollowAccepted( + h.ctx, + new Accept({ + id: new URL("#accept", newId), + actor: newActor, + object: follow, + }), + ); + assert.ok(await h.repository.getFollowee("alpha", newId)); +}); + +test("concurrent copies migrate each bot only once", async (t) => { + const h = await harness(t.signal); + const events: string[] = []; + for (const bot of [h.alpha, h.beta]) { + bot.onFolloweeMove = (s) => { + events.push(s.bot.identifier); + }; + } + await Promise.all([ + h.instance.onMoved(h.ctx, move, t.signal), + h.instance.onMoved(h.ctx, move, t.signal), + h.instance.onMoved( + h.ctx, + move.clone({ id: new URL("#copy", oldId) }), + t.signal, + ), + ]); + assert.strictEqual(h.sent.length, 4); + assert.deepStrictEqual(events, ["alpha", "beta"]); +}); + +test("a failed Follow preserves that bot's old follow and does not poison retries", async (t) => { + const h = await harness(t.signal); + const error = new TypeError("Follow submission failed."); + const send = h.ctx.sendActivity; + h.ctx.sendActivity = (sender, recipient, activity, options) => { + if ( + activity instanceof Follow && + activity.actorId?.href === h.ctx.getActorUri("alpha").href + ) return Promise.reject(error); + if (recipient === "followers") { + throw new TypeError("Unexpected collection delivery."); + } + return send(sender, recipient, activity, options); + }; + const events: string[] = []; + for (const bot of [h.alpha, h.beta]) { + bot.onFolloweeMove = (s) => { + events.push(s.bot.identifier); + }; + } + await assert.rejects(h.instance.onMoved(h.ctx, move, t.signal), error); + assert.ok(await h.repository.getFollowee("alpha", oldId)); + assert.ok(!await h.repository.getFollowee("beta", oldId)); + assert.deepStrictEqual(events, ["beta"]); + h.ctx.sendActivity = send; + await h.instance.onMoved(h.ctx, move, t.signal); + assert.deepStrictEqual(events, ["beta", "alpha"]); + assert.strictEqual(h.sent.length, 4); +}); + +test("callback failure does not prevent other bots from migrating or firing callbacks", async (t) => { + const h = await harness(t.signal); + const error = new TypeError("Callback failed."); + const events: string[] = []; + h.alpha.onFolloweeMove = () => { + throw error; + }; + h.beta.onFolloweeMove = (s) => { + events.push(s.bot.identifier); + }; + await assert.rejects(h.instance.onMoved(h.ctx, move, t.signal), error); + assert.deepStrictEqual(events, ["beta"]); + assert.strictEqual(h.sent.length, 4); + await h.instance.onMoved(h.ctx, move, t.signal); + assert.deepStrictEqual(events, ["beta"]); +}); + +test("a slow callback does not hold the migration lock", async (t) => { + const h = await harness(t.signal); + const started = Promise.withResolvers(); + const finish = Promise.withResolvers(); + h.alpha.onFolloweeMove = () => { + started.resolve(); + return finish.promise; + }; + const first = h.instance.onMoved(h.ctx, move, t.signal); + await started.promise; + try { + await h.instance.onMoved(h.ctx, move, t.signal); + assert.strictEqual(h.sent.length, 4); + } finally { + finish.resolve(); + await first; + } +}); + +test("Undo submission failure removes the old follow but does not emit or replay the event", async (t) => { + const h = await harness(t.signal); + const error = new TypeError("Undo failed."); + const send = h.ctx.sendActivity; + h.ctx.sendActivity = (sender, recipient, activity, options) => { + if (recipient === "followers") { + throw new TypeError("Unexpected collection delivery."); + } + return activity instanceof Undo + ? Promise.reject(error) + : send(sender, recipient, activity, options); + }; + let events = 0; + h.alpha.onFolloweeMove = h.beta.onFolloweeMove = () => { + events++; + }; + await assert.rejects(h.instance.onMoved(h.ctx, move, t.signal), error); + assert.ok(!await h.repository.getFollowee("alpha", oldId)); + assert.strictEqual(events, 0); + await h.instance.onMoved(h.ctx, move, t.signal); + assert.strictEqual(h.sent.length, 2); +}); + +test("partial embedded origin is fetched to recover its delivery inbox", async (t) => { + const h = await harness(t.signal); + await h.instance.onMoved( + h.ctx, + move.clone({ actor: new Person({ id: oldId }) }), + t.signal, + ); + assert.strictEqual(h.sent.length, 4); + assert.ok(h.loads.includes(oldId.href)); +}); + +test("an aborted Move never mutates follow relationships", async (t) => { + const h = await harness(t.signal); + const controller = new AbortController(); + controller.abort(); + await assert.rejects(h.instance.onMoved(h.ctx, move, controller.signal), { + name: "AbortError", + }); + assert.deepStrictEqual(h.sent, []); + assert.ok(await h.repository.getFollowee("alpha", oldId)); +}); + +test("CreateBotOptions and wrapped setter expose onFolloweeMove", () => { + const initial = () => {}; + const bot = createBot({ + kv: new MemoryKvStore(), + username: "bot", + onFolloweeMove: initial, + }); + assert.strictEqual(bot.onFolloweeMove, initial); + const next = () => {}; + bot.onFolloweeMove = next; + assert.strictEqual(bot.onFolloweeMove, next); + bot.onFolloweeMove = undefined; + assert.strictEqual(bot.onFolloweeMove, undefined); +}); + +test("dynamic groups forward the handler live for every resolved bot", async (t) => { + const h = await harness(t.signal); + await h.repository.removeFollowee("alpha", oldId); + await h.repository.removeFollowee("beta", oldId); + const group = h.instance.createBot((_ctx, id) => + id.startsWith("dynamic") ? { username: id } : null + ); + await group.getSession(h.ctx.origin, "dynamic-a"); + const moved: string[] = []; + group.onFolloweeMove = (s) => { + moved.push(s.bot.identifier); + }; + for (const id of ["dynamic-a", "dynamic-b"]) { + await h.repository.addFollowee( + id, + oldId, + new Follow({ + id: h.ctx.getObjectUri(Follow, { + identifier: id, + id: "018f6db5-27d2-7000-8000-000000000002", + }), + actor: h.ctx.getActorUri(id), + object: oldId, + }), + ); + } + await h.instance.onMoved(h.ctx, move, t.signal); + assert.deepStrictEqual(moved, ["dynamic-a", "dynamic-b"]); +}); + +test("migration to one following bot skips itself and still migrates its siblings", async (t) => { + const h = await harness(t.signal); + const alphaId = h.ctx.getActorUri("alpha"); + await h.instance.onMoved(h.ctx, move.clone({ target: alphaId }), t.signal); + assert.strictEqual(h.sent.length, 2); + assert.deepStrictEqual(h.sent[0].actorId, h.ctx.getActorUri("beta")); + assert.deepStrictEqual(h.sent[0].objectId, alphaId); + assert.ok(await h.repository.getFollowee("alpha", oldId)); + assert.ok(!await h.repository.getFollowee("beta", oldId)); + assert.deepStrictEqual(h.loads, []); // Local target is resolved authoritatively. +}); + +test("single-bot compatibility path migrates and invokes its configured handler", async (t) => { + const h = await harness(t.signal); + let events = 0; + const bot = new BotImpl({ + kv: new MemoryKvStore(), + username: "bot", + repository: h.repository, + onFolloweeMove: () => { + events++; + }, + }); + const ctx = bot.federation.createContext( + new URL(h.ctx.origin), + undefined, + ) as InboxContext; + Object.defineProperty(ctx, "recipient", { value: "bot" }); + Object.defineProperty(ctx, "getDocumentLoader", { + value: () => h.ctx.getDocumentLoader({ identifier: "alpha" }), + }); + ctx.sendActivity = h.ctx.sendActivity; + await h.repository.addFollowee( + "bot", + oldId, + new Follow({ + id: ctx.getObjectUri(Follow, { + identifier: "bot", + id: "018f6db5-27d2-7000-8000-000000000004", + }), + actor: ctx.getActorUri("bot"), + object: oldId, + }), + ); + await bot.instance.onMoved(ctx, move, t.signal); + assert.strictEqual(events, 1); + assert.ok(!await h.repository.getFollowee("bot", oldId)); + assert.strictEqual(h.sent.length, 2); +}); + +test("synchronous context-loader failure remains retriable", async (t) => { + const h = await harness(t.signal); + const error = new TypeError("Context fetch failed synchronously."); + h.documents.set(newId.href, { + "@context": "https://context.example/context", + id: newId.href, + }); + Object.defineProperty(h.ctx, "contextLoader", { + value: () => { + throw error; + }, + }); + await assert.rejects(h.instance.onMoved(h.ctx, move, t.signal), error); + assert.deepStrictEqual(h.sent, []); +}); + +test("origin lookup failure preserves all old follows", async (t) => { + const h = await harness(t.signal); + const error = new TypeError("Origin fetch failed."); + Object.defineProperty(h.ctx, "getDocumentLoader", { + value: () => + Promise.resolve((url: string) => + url === oldId.href ? Promise.reject(error) : Promise.resolve({ + documentUrl: url, + contextUrl: null, + document: h.documents.get(url), + }) + ), + }); + await assert.rejects( + h.instance.onMoved(h.ctx, move.clone({ actor: oldId }), t.signal), + error, + ); + assert.deepStrictEqual(h.sent, []); + assert.ok(await h.repository.getFollowee("alpha", oldId)); +}); + +test("cancellation is forwarded to target and origin context loaders", async (t) => { + const h = await harness(t.signal); + const loadContext = h.ctx.contextLoader; + let loaded = 0; + Object.defineProperty(h.ctx, "contextLoader", { + value: (url: string, options?: Parameters[1]) => { + loaded++; + assert.strictEqual(options?.signal, t.signal); + return loadContext(url, options); + }, + }); + await h.instance.onMoved(h.ctx, move.clone({ actor: oldId }), t.signal); + assert.ok(loaded > 0); + assert.strictEqual(h.sent.length, 4); +}); + +for (const status of [401, 403]) { + test(`target authorization denial (${status}) for one bot does not block its siblings`, async (t) => { + const h = await harness(t.signal); + const identities: string[] = []; + Object.defineProperty(h.ctx, "getDocumentLoader", { + value: (identity: { identifier: string }) => + Promise.resolve((url: string) => { + identities.push(identity.identifier); + if (identity.identifier === "alpha") { + return Promise.reject( + new FetchError(newId, "Denied.", new Response(null, { status })), + ); + } + return Promise.resolve({ + documentUrl: url, + contextUrl: null, + document: h.documents.get(url), + }); + }), + }); + await h.instance.onMoved(h.ctx, move, t.signal); + assert.deepStrictEqual(identities, ["alpha", "beta"]); + assert.strictEqual(h.sent.length, 4); + assert.ok(!await h.repository.getFollowee("beta", oldId)); + }); +} + +for (const status of [401, 403]) { + for (const embedded of [false, true]) { + test(`origin authorization denial (${status}, embedded=${embedded}) retries a sibling identity`, async (t) => { + const h = await harness(t.signal); + const origins: string[] = []; + Object.defineProperty(h.ctx, "getDocumentLoader", { + value: (identity: { identifier: string }) => + Promise.resolve((url: string) => { + if (url === oldId.href) { + origins.push(identity.identifier); + if (identity.identifier === "alpha") { + return Promise.reject( + new FetchError( + oldId, + "Denied.", + new Response(null, { status }), + ), + ); + } + } + return Promise.resolve({ + documentUrl: url, + contextUrl: null, + document: h.documents.get(url), + }); + }), + }); + await h.instance.onMoved( + h.ctx, + move.clone({ actor: embedded ? new Person({ id: oldId }) : oldId }), + t.signal, + ); + assert.deepStrictEqual(origins, ["alpha", "beta"]); + assert.strictEqual(h.sent.length, 4); + assert.ok(!await h.repository.getFollowee("beta", oldId)); + }); + } +} + +test("local target without the origin alias is ignored", async (t) => { + const h = await harness(t.signal); + await h.instance.onMoved( + h.ctx, + move.clone({ target: h.ctx.getActorUri("beta") }), + t.signal, + ); + assert.deepStrictEqual(h.sent, []); + assert.deepStrictEqual(h.loads, []); + assert.ok(await h.repository.getFollowee("alpha", oldId)); + assert.ok(await h.repository.getFollowee("beta", oldId)); +}); diff --git a/packages/botkit/src/instance-impl.ts b/packages/botkit/src/instance-impl.ts index 2b6e596..1cf4c9d 100644 --- a/packages/botkit/src/instance-impl.ts +++ b/packages/botkit/src/instance-impl.ts @@ -39,10 +39,13 @@ import { Endpoints, Follow, Image, + isActor, Like as RawLike, Link, Mention, + Move, Note, + Object as APObject, Question, QuoteAuthorization, QuoteRequest, @@ -50,6 +53,11 @@ import { Undo, Update, } from "@fedify/vocab"; +import { + type DocumentLoader, + FetchError, + UrlError, +} from "@fedify/vocab-runtime"; import { getLogger } from "@logtape/logtape"; import mimeDb from "mime-db"; import fs from "node:fs/promises"; @@ -74,6 +82,22 @@ import { KvRepository, type Repository } from "./repository.ts"; import type { Session } from "./session.ts"; import { parseLocalUri, rewriteLegacyObjectPath } from "./uri.ts"; +interface FolloweeMoveResult { + readonly oldActor: Actor; + readonly newActor: Actor; + readonly bots: readonly BotImpl[]; + readonly errors: readonly unknown[]; +} + +function isPermanentMoveFetchError(error: unknown): boolean { + if (error instanceof SyntaxError) return true; + if (error instanceof UrlError) return error.reason !== "dns"; + if (!(error instanceof FetchError) || error.response == null) return false; + const status = error.response.status; + return status < 400 || + status >= 400 && status < 500 && status !== 408 && status !== 429; +} + /** * The default identifier of the instance actor: an internal `Application` * actor that an {@link Instance} uses for signing shared-inbox related @@ -335,6 +359,7 @@ export class InstanceImpl this.onUnverifiedActivity(ctx, activity, reason) ) .on(Follow, (ctx, follow) => this.onFollowed(ctx, follow)) + .on(Move, (ctx, move) => this.onMoved(ctx, move)) .on(QuoteRequest, (ctx, request) => this.onQuoteRequested(ctx, request)) .on(Undo, (ctx, undo) => this.onUndone(ctx, undo)) .on(Accept, (ctx, accept) => this.onFollowAccepted(ctx, accept)) @@ -647,6 +672,247 @@ export class InstanceImpl return []; } + // Serializes copies delivered through shared and personal inboxes. This + // is process-local; repository relationships gate subsequent deliveries. + readonly #moves = new Map>(); + + /** + * Handles a verified push-mode account migration for all following bots. + * @param ctx The inbox context. + * @param move The verified incoming activity. + * @param signal The signal for cancelling processing. + * @returns A promise that resolves after transitions and callbacks complete. + * @throws If a transient lookup, transition, or event handler fails. + * @internal + */ + async onMoved( + ctx: InboxContext, + move: Move, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + const oldId = move.actorId; + const targetId = move.targetId; + if ( + move.id == null || oldId == null || targetId == null || + move.actorIds.length !== 1 || move.objectIds.length !== 1 || + move.targetIds.length !== 1 || move.objectId?.href !== oldId.href || + targetId.href === oldId.href || + (targetId.protocol !== "https:" && targetId.protocol !== "http:") + ) return; + + const key = JSON.stringify([ctx.origin, oldId.href]); + const previous = this.#moves.get(key) ?? Promise.resolve(); + const run = previous.then(() => + this.#processMove(ctx, move, oldId, targetId, signal) + ); + const tail = run.then(() => {}, () => {}); + this.#moves.set(key, tail); + let result: FolloweeMoveResult | null; + try { + result = await run; + } finally { + if (this.#moves.get(key) === tail) this.#moves.delete(key); + } + if (result == null) return; + const errors = [...result.errors]; + // Handlers run after releasing the lock and are read live from groups. + for (const bot of result.bots) { + try { + await bot.onFolloweeMove?.( + bot.getSession(ctx), + result.oldActor, + result.newActor, + ); + } catch (error) { + errors.push(error); + } + } + if (errors.length > 0) throw errors[0]; + } + + async #fetchMoveActor( + ctx: InboxContext, + id: URL, + documentLoader: DocumentLoader, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + const logger = getLogger(["botkit", "instance", "inbox"]); + const document = await documentLoader(id.href); + let loaderFailure: { readonly error: unknown } | undefined; + const observe = + (load: DocumentLoader): DocumentLoader => async (url, options) => { + try { + return await load(url, options); + } catch (error) { + loaderFailure = { error }; + throw error; + } + }; + let actor: APObject; + try { + const documentUrl = new URL(document.documentUrl); + if (documentUrl.origin !== id.origin) return null; + actor = await APObject.fromJsonLd(document.document, { + documentLoader: observe(documentLoader), + contextLoader: observe((url, options) => + ctx.contextLoader(url, { + ...options, + signal: signal ?? options?.signal, + }) + ), + baseUrl: documentUrl, + }); + } catch (error) { + signal?.throwIfAborted(); + if (loaderFailure != null) throw loaderFailure.error; + logger.debug( + "Ignoring Move with an invalid target document: {error}", + { error }, + ); + return null; + } + return isActor(actor) && actor.id?.href === id.href ? actor : null; + } + + async #processMove( + ctx: InboxContext, + move: Move, + oldId: URL, + targetId: URL, + signal?: AbortSignal, + ): Promise | null> { + signal?.throwIfAborted(); + // Always fan out, including personal deliveries: Fedify's inbox queue + // deduplicates an activity across recipients on the receiving origin. + // Snapshot first, since unfollowing mutates the repository's reverse index. + const identifiers = new Set( + await Array.fromAsync(this.repository.findFollowedBots(oldId)), + ); + const bots: BotImpl[] = []; + for (const identifier of identifiers) { + const bot = await this.resolveBot(ctx, identifier); + if ( + bot != null && ctx.getActorUri(identifier).href !== targetId.href && + await bot.repository.getFollowee(oldId) != null + ) bots.push(bot); + } + if (bots.length === 0) return null; + + const logger = getLogger(["botkit", "instance", "inbox"]); + const loader = await ctx.getDocumentLoader(bots[0]); + let documentLoader: DocumentLoader = (url, options) => + loader(url, { ...options, signal: signal ?? options?.signal }); + let newActor: Actor; + try { + const local = ctx.parseUri(targetId); + if (local?.type === "actor") { + const targetBot = await this.resolveBot(ctx, local.identifier); + const actor = await targetBot?.dispatchActor(ctx, local.identifier); + if (actor == null) return null; + newActor = actor; + } else { + // Do not use getTarget() (which trusts embeds) or lookupObject() + // (which suppresses transport errors and can prevent queue retries). + let actor: Actor | null = null; + for (let index = 0; index < bots.length; index++) { + try { + actor = await this.#fetchMoveActor( + ctx, + targetId, + documentLoader, + signal, + ); + break; + } catch (error) { + // A signature-specific denial must not prevent other following + // bots from verifying the destination with their own identity. + if ( + error instanceof FetchError && + (error.response?.status === 401 || + error.response?.status === 403) && + index + 1 < bots.length + ) { + const alternate = await ctx.getDocumentLoader(bots[index + 1]); + documentLoader = (url, options) => + alternate(url, { + ...options, + signal: signal ?? options?.signal, + }); + continue; + } + throw error; + } + } + if (actor == null) return null; + newActor = actor; + } + } catch (error) { + signal?.throwIfAborted(); + if (!isPermanentMoveFetchError(error)) throw error; + logger.debug("Ignoring Move with an unavailable target: {error}", { + error, + }); + return null; + } + if ( + newActor.id?.href !== targetId.href || newActor.inboxId == null || + !newActor.aliasIds.some((alias) => alias.href === oldId.href) + ) return null; + + let oldActor: Actor | null = null; + for (let index = 0; index < bots.length; index++) { + const originLoader = await ctx.getDocumentLoader(bots[index]); + const loadOrigin: DocumentLoader = (url, options) => + originLoader(url, { ...options, signal: signal ?? options?.signal }); + try { + oldActor = await move.getActor({ + documentLoader: loadOrigin, + contextLoader: (url, options) => + ctx.contextLoader(url, { + ...options, + signal: signal ?? options?.signal, + }), + }); + if (oldActor?.id?.href !== oldId.href) return null; + if (oldActor.inboxId == null) { + oldActor = await this.#fetchMoveActor(ctx, oldId, loadOrigin, signal); + if (oldActor == null || oldActor.inboxId == null) return null; + } + break; + } catch (error) { + signal?.throwIfAborted(); + if ( + error instanceof FetchError && + (error.response?.status === 401 || error.response?.status === 403) && + index + 1 < bots.length + ) continue; + if (!isPermanentMoveFetchError(error)) throw error; + return null; + } + } + if (oldActor == null) return null; + const migrated: BotImpl[] = []; + const errors: unknown[] = []; + for (const bot of bots) { + try { + signal?.throwIfAborted(); + if (await bot.repository.getFollowee(oldId) == null) continue; + const session = bot.getSession(ctx); + if (await bot.repository.getFollowee(targetId) == null) { + await session.follow(newActor); + } + signal?.throwIfAborted(); + await session.unfollow(oldActor); + migrated.push(bot); + } catch (error) { + errors.push(error); + } + } + return { oldActor, newActor, bots: migrated, errors }; + } + async onFollowed( ctx: InboxContext, follow: Follow, diff --git a/packages/botkit/src/session-impl.ts b/packages/botkit/src/session-impl.ts index 07bc146..9faa5b9 100644 --- a/packages/botkit/src/session-impl.ts +++ b/packages/botkit/src/session-impl.ts @@ -57,6 +57,7 @@ import type { SessionPublishOptionsWithQuestion, } from "./session.ts"; import { plainText, type Text } from "./text.ts"; +import { getFollowDeliveryOptions } from "./uri.ts"; const logger = getLogger(["botkit", "session"]); @@ -144,7 +145,7 @@ export class SessionImpl implements Session { this.bot, actor, follow, - { excludeBaseUris: [new URL(this.context.origin)] }, + getFollowDeliveryOptions(this.context, actor.id), ); } @@ -188,7 +189,7 @@ export class SessionImpl implements Session { object: follow, to: actor.id, }), - { excludeBaseUris: [new URL(this.context.origin)] }, + getFollowDeliveryOptions(this.context, actor.id), ); } } diff --git a/packages/botkit/src/uri.ts b/packages/botkit/src/uri.ts index 54c7f8b..ae24e4c 100644 --- a/packages/botkit/src/uri.ts +++ b/packages/botkit/src/uri.ts @@ -80,3 +80,22 @@ export function parseLocalUri( rewritten.pathname = rewrittenPath; return ctx.parseUri(rewritten); } + +/** + * Keeps follow activities deliverable to actors hosted on this instance. + * Remote recipients retain the usual exclusion of the instance's own inboxes. + * @param ctx The federation context. + * @param recipientId The recipient's actor URI. + * @returns The delivery options for the follow activity. + * @internal + */ +export function getFollowDeliveryOptions( + ctx: Context, + recipientId: URL | null, +): { excludeBaseUris: URL[] } { + return { + excludeBaseUris: ctx.parseUri(recipientId)?.type === "actor" + ? [] + : [new URL(ctx.origin)], + }; +}