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)], + }; +}