From c886ccaaa7795086520444dec239aca38bb65c7e Mon Sep 17 00:00:00 2001 From: Hong Minhee Date: Sat, 3 Oct 2026 17:33:42 +0900 Subject: [PATCH 1/4] Let bots move their followers to a new actor Operators need to retire a bot without losing its followers. Validate migration targets against their declared aliases, persist the successor before submitting Update and Move, and retain that state if notifications fail so they can be republished safely. Keep moved bots from publishing or accepting new follows, expose their successor in actor documents and public pages, and preserve existing messages. Support the migration state in every built-in repository with an atomic first-write-wins contract for custom repositories. Closes https://github.com/fedify-dev/botkit/issues/50 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 | 26 ++ changes.d/botkit-postgres/account-move.md | 8 + changes.d/botkit-redis/account-move.md | 8 + changes.d/botkit-sqlite/account-move.md | 8 + changes.d/botkit/account-move.md | 11 + docs/concepts/repository.md | 29 ++ docs/concepts/session.md | 92 ++++ packages/botkit-postgres/src/mod.test.ts | 39 ++ packages/botkit-postgres/src/mod.ts | 48 ++ packages/botkit-redis/src/mod.test.ts | 38 ++ packages/botkit-redis/src/mod.ts | 28 ++ packages/botkit-sqlite/src/mod.test.ts | 37 ++ packages/botkit-sqlite/src/mod.ts | 32 ++ packages/botkit/src/bot-impl.ts | 31 +- packages/botkit/src/follow-impl.ts | 5 + packages/botkit/src/follow-local.test.ts | 74 +++ packages/botkit/src/follow.ts | 2 +- packages/botkit/src/instance-impl.ts | 2 + packages/botkit/src/message-impl.ts | 1 + packages/botkit/src/message.ts | 2 +- packages/botkit/src/mod.ts | 1 + packages/botkit/src/move.test.ts | 535 ++++++++++++++++++++++ packages/botkit/src/pages.test.ts | 84 ++++ packages/botkit/src/pages.tsx | 57 ++- packages/botkit/src/repository.test.ts | 77 ++++ packages/botkit/src/repository.ts | 128 ++++++ packages/botkit/src/session-impl.ts | 197 +++++++- packages/botkit/src/session.ts | 46 ++ packages/botkit/src/successor.ts | 33 ++ packages/botkit/src/text.test.ts | 6 + 30 files changed, 1680 insertions(+), 5 deletions(-) create mode 100644 changes.d/botkit-postgres/account-move.md create mode 100644 changes.d/botkit-redis/account-move.md create mode 100644 changes.d/botkit-sqlite/account-move.md create mode 100644 changes.d/botkit/account-move.md create mode 100644 packages/botkit/src/move.test.ts create mode 100644 packages/botkit/src/successor.ts diff --git a/CHANGES.md b/CHANGES.md index ab79070..f840417 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -8,6 +8,12 @@ To be released. ### @fedify/botkit + - Added `Session.move()` to move a bot's followers to a linked actor, + persist its `movedTo` redirect, and show its new home on the profile. + Moved bots reject new follows and cannot publish, reply, or share new + messages. `Session.republishMove()` resends failed migration notifications. + Custom repositories must now implement the required `getSuccessor()` and + atomic, write-once `setSuccessor()` methods. [[#50], [#57]] - Added an `aliases` option to `CreateBotOptions` and `BotProfile` so existing accounts can move their followers to a BotKit bot. Actor URIs listed in the option are published as `alsoKnownAs`, and are available @@ -45,10 +51,30 @@ To be released. [#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 +[#50]: https://github.com/fedify-dev/botkit/issues/50 [#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 +[#57]: https://github.com/fedify-dev/botkit/pull/57 + +### @fedify/botkit-postgres + + - Added persistent account migration state through `getSuccessor()` and + atomic, write-once `setSuccessor()`, keeping moved bots inactive across + restarts. [[#50], [#57]] + +### @fedify/botkit-redis + + - Added persistent account migration state through `getSuccessor()` and + atomic, write-once `setSuccessor()`, keeping moved bots inactive across + restarts. [[#50], [#57]] + +### @fedify/botkit-sqlite + + - Added persistent account migration state through `getSuccessor()` and + atomic, write-once `setSuccessor()`, keeping moved bots inactive across + restarts. [[#50], [#57]] Version 0.5.6 diff --git a/changes.d/botkit-postgres/account-move.md b/changes.d/botkit-postgres/account-move.md new file mode 100644 index 0000000..0b1ba01 --- /dev/null +++ b/changes.d/botkit-postgres/account-move.md @@ -0,0 +1,8 @@ +--- +links: + '#50': https://github.com/fedify-dev/botkit/issues/50 + '#57': https://github.com/fedify-dev/botkit/pull/57 +--- + - Added persistent account migration state through `getSuccessor()` and + atomic, write-once `setSuccessor()`, keeping moved bots inactive across + restarts. [[#50], [#57]] diff --git a/changes.d/botkit-redis/account-move.md b/changes.d/botkit-redis/account-move.md new file mode 100644 index 0000000..0b1ba01 --- /dev/null +++ b/changes.d/botkit-redis/account-move.md @@ -0,0 +1,8 @@ +--- +links: + '#50': https://github.com/fedify-dev/botkit/issues/50 + '#57': https://github.com/fedify-dev/botkit/pull/57 +--- + - Added persistent account migration state through `getSuccessor()` and + atomic, write-once `setSuccessor()`, keeping moved bots inactive across + restarts. [[#50], [#57]] diff --git a/changes.d/botkit-sqlite/account-move.md b/changes.d/botkit-sqlite/account-move.md new file mode 100644 index 0000000..0b1ba01 --- /dev/null +++ b/changes.d/botkit-sqlite/account-move.md @@ -0,0 +1,8 @@ +--- +links: + '#50': https://github.com/fedify-dev/botkit/issues/50 + '#57': https://github.com/fedify-dev/botkit/pull/57 +--- + - Added persistent account migration state through `getSuccessor()` and + atomic, write-once `setSuccessor()`, keeping moved bots inactive across + restarts. [[#50], [#57]] diff --git a/changes.d/botkit/account-move.md b/changes.d/botkit/account-move.md new file mode 100644 index 0000000..81465cf --- /dev/null +++ b/changes.d/botkit/account-move.md @@ -0,0 +1,11 @@ +--- +links: + '#50': https://github.com/fedify-dev/botkit/issues/50 + '#57': https://github.com/fedify-dev/botkit/pull/57 +--- + - Added `Session.move()` to move a bot's followers to a linked actor, + persist its `movedTo` redirect, and show its new home on the profile. + Moved bots reject new follows and cannot publish, reply, or share new + messages. `Session.republishMove()` resends failed migration notifications. + Custom repositories must now implement the required `getSuccessor()` and + atomic, write-once `setSuccessor()` methods. [[#50], [#57]] diff --git a/docs/concepts/repository.md b/docs/concepts/repository.md index e463af3..bda2ed8 100644 --- a/docs/concepts/repository.md +++ b/docs/concepts/repository.md @@ -43,6 +43,35 @@ the quote's authorization state. [FEP-044f]: https://w3id.org/fep/044f +Account migration storage +------------------------- + +Since BotKit 0.6.0, custom repositories must implement two additional methods: + + - `~Repository.getSuccessor(identifier, signal?)` returns the successor actor + `URL`, or `undefined` if the bot has not moved. Return a fresh `URL` so + callers cannot mutate the stored value. + - `~Repository.setSuccessor(identifier, successorId, signal?)` atomically + records the first successor and returns `true`. If one already exists, + return `false`, including when the URI is identical. Never overwrite it. + +These are required even if your application never calls `Session.move()`: +actor dispatch and publishing read the moved state. BotKit rejects repositories +missing either method during construction, including repositories passed +through a cache or the single-bot compatibility wrapper. Honor an aborted +signal before a read or write, and never report cancellation after a write has +committed. See [moving a bot](./session.md#moving-the-bot-to-another-actor) for +the lifecycle and recovery behavior. + +All built-in repositories implement this contract. SQL repositories create an +additional table automatically, and KV/Redis repositories store a bot-scoped +successor key. `MemoryCachedRepository` reads successors directly from its +backing repository, so another process's move takes effect immediately. A +`KvRepository` without CAS can serialize successor writes only within the +same repository instance; use a CAS-capable store or a dedicated SQL/Redis +repository when several instances may move the same bot concurrently. + + `KvRepository` -------------- diff --git a/docs/concepts/session.md b/docs/concepts/session.md index 5ef7396..676d7af 100644 --- a/docs/concepts/session.md +++ b/docs/concepts/session.md @@ -193,6 +193,98 @@ followers. Call it after your application updates the bot profile and you want the change to propagate without waiting for the next post. +Moving the bot to another actor +------------------------------- + +*This API is available since BotKit 0.6.0.* + +Call `~Session.move()` to move the bot's followers to another BotKit bot or +an account on a different server. First configure the destination to list +this bot's *actor URI* in its `alsoKnownAs` aliases. For a BotKit destination, +use [`aliases`](./bot.md#createbotoptions-aliases) and deploy it before starting +the move. The URI is `session.actorId`, rather than the bot's profile page URL. + +~~~~ typescript twoslash +import type { Session } from "@fedify/botkit"; +declare const session: Session; +// ---cut-before--- +await session.move("@mybot@new.example"); +~~~~ + +The target can also be an actor `URL`, a URI string, or an `Actor` object. +BotKit fetches its current actor document and verifies that its aliases contain +this bot's actor URI. It rejects the bot itself, targets without an inbox, +and targets that have already moved. + +BotKit stores the successor, publishes an actor `Update` carrying `movedTo`, +then sends a push-mode [FEP-7628] `Move` to the old followers. This includes +followers on the same instance. Only followers move: posts, followed accounts, +and other data stay on the old server. Keep that server running while remote +servers process the migration. A successful call means the notifications were +submitted, rather than that every follower has already moved. + +The moved state survives a restart when the repository is persistent. The +old actor's `successorId` points to the destination, its profile shows a link, +and it rejects incoming follow requests without invoking `onFollow`, regardless +of `followerPolicy`. `Session.publish()`, `Message.reply()`, and +`Message.share()` throw `TypeError` on a moved bot. Existing posts remain +accessible, and their editing and deletion remain available. + +Other event handlers still run. Deploy a moved-state check in handlers that +publish or reply *before* starting the move: + +~~~~ typescript twoslash +import { type Bot, text } from "@fedify/botkit"; +declare const bot: Bot; +// ---cut-before--- +bot.onMention = async (session, message) => { + if ((await session.getActor()).successorId != null) return; + await message.reply(text`Thanks for mentioning me!`); +}; +~~~~ + +Without this check, a reply attempt throws from the handler and fails the +incoming activity's processing; a configured queue may retry it. Pending +follow requests retained before the move also cannot be accepted afterwards, +but can still be rejected. There is no API to undo a move or change its +stored destination. + +[FEP-7628]: https://w3id.org/fep/7628 + +### Recovering notification failures + +A storage or delivery failure can occur after the successor has been stored. +BotKit keeps the moved state because some servers may already have processed +the migration. It attempts the `Move` even if submitting the `Update` fails. +Post-commit notification failures throw `AggregateError`; a storage error can +also leave the write's outcome uncertain. On any error, check +`(await session.getActor()).successorId` before choosing how to retry. + +Calling `move()` on a moved bot throws `TypeError`. Use +`~Session.republishMove()` to revalidate the stored successor's alias and +resend both notifications to the remaining followers: + +~~~~ typescript twoslash +import type { Session } from "@fedify/botkit"; +declare const session: Session; +// ---cut-before--- +if ((await session.getActor()).successorId != null) { + await session.republishMove(); +} +~~~~ + +This also reaches followers whose acceptance was already in progress when the +bot moved. It does not change the successor. The destination must still list +the old actor as an alias, even if it has since moved again. With a configured +queue, Fedify retries delivery of notifications it has accepted; without a +queue, delivery happens during the call and can partially fail. + +Both methods accept `{ signal: AbortSignal }`. `move()` honours cancellation +until successor storage commits, then continues both notifications. +`republishMove()` honours cancellation during preparation, including the +follower snapshot, then continues both notifications once submission starts. + + Publishing a message -------------------- diff --git a/packages/botkit-postgres/src/mod.test.ts b/packages/botkit-postgres/src/mod.test.ts index 2f268f7..bd6a884 100644 --- a/packages/botkit-postgres/src/mod.test.ts +++ b/packages/botkit-postgres/src/mod.test.ts @@ -113,6 +113,44 @@ if (postgresUrl == null) { test("PostgresRepository integration tests", { skip: true }, () => {}); } else { describe("PostgresRepository", () => { + test("successor is atomic, scoped and persistent", async (t) => { + const harness = createHarness(); + const repo = harness.repository; + const target = new URL("https://new.example/actor"); + try { + assert.strictEqual(await repo.getSuccessor("old", t.signal), undefined); + const results = await Promise.all([ + repo.setSuccessor("old", target, t.signal), + repo.setSuccessor( + "old", + new URL("https://other.example/actor"), + t.signal, + ), + ]); + assert.strictEqual(results.filter(Boolean).length, 1); + const successor = await repo.getSuccessor("old", t.signal); + assert.ok(successor); + assert.ok(!await repo.setSuccessor("old", successor)); + assert.strictEqual(await repo.getSuccessor("sibling"), undefined); + await assert.rejects( + repo.setSuccessor("sibling", target, AbortSignal.abort()), + { name: "AbortError" }, + ); + const second = new PostgresRepository({ + url: postgresUrl, + schema: harness.schema, + }); + try { + assert.deepStrictEqual(await second.getSuccessor("old"), successor); + assert.ok(!await second.setSuccessor("old", target)); + } finally { + await second.close(); + } + } finally { + await harness.cleanup(); + } + }); + test("initializes schema explicitly", async () => { const sql = createSql(postgresUrl); const schema = createSchemaName(); @@ -129,6 +167,7 @@ if (postgresUrl == null) { assert.deepStrictEqual( tables.map((row) => row.table_name), [ + "bot_successors", "botkit_metadata", "follow_requests", "followees", diff --git a/packages/botkit-postgres/src/mod.ts b/packages/botkit-postgres/src/mod.ts index 17da99a..92bbfab 100644 --- a/packages/botkit-postgres/src/mod.ts +++ b/packages/botkit-postgres/src/mod.ts @@ -179,6 +179,15 @@ async function initializePostgresRepositorySchemaInTransaction( [], prepare, ); + await execute( + sql, + `CREATE TABLE IF NOT EXISTS "${validatedSchema}"."bot_successors" ( + bot_id TEXT PRIMARY KEY, + successor_id TEXT NOT NULL + )`, + [], + prepare, + ); await execute( sql, `CREATE TABLE IF NOT EXISTS "${validatedSchema}"."key_pairs" ( @@ -595,6 +604,45 @@ export class PostgresRepository implements Repository, AsyncDisposable { } } + /** {@inheritDoc Repository.getSuccessor} */ + async getSuccessor( + identifier: string, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + await this.ensureReady(); + signal?.throwIfAborted(); + const rows = await this.query<{ readonly successor_id: string }>( + this.sql, + `SELECT successor_id FROM ${ + this.table("bot_successors") + } WHERE bot_id = $1`, + [identifier], + ); + return rows[0] === undefined ? undefined : new URL(rows[0].successor_id); + } + + /** {@inheritDoc Repository.setSuccessor} */ + async setSuccessor( + identifier: string, + successorId: URL, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + const href = successorId.href; + await this.ensureReady(); + signal?.throwIfAborted(); + const rows = await this.query<{ readonly bot_id: string }>( + this.sql, + `INSERT INTO ${ + this.table("bot_successors") + } (bot_id, successor_id) VALUES ($1, $2) + ON CONFLICT (bot_id) DO NOTHING RETURNING bot_id`, + [identifier, href], + ); + return rows.length > 0; + } + async setKeyPairs( identifier: string, keyPairs: CryptoKeyPair[], diff --git a/packages/botkit-redis/src/mod.test.ts b/packages/botkit-redis/src/mod.test.ts index f3aca3e..78414b2 100644 --- a/packages/botkit-redis/src/mod.test.ts +++ b/packages/botkit-redis/src/mod.test.ts @@ -300,6 +300,44 @@ if (redisUrl == null) { ); }); + test("successor is atomic, scoped and persistent", async (t) => { + const harness = createHarness(); + const repo = harness.repository; + const target = new URL("https://new.example/actor"); + try { + assert.strictEqual(await repo.getSuccessor("old", t.signal), undefined); + const results = await Promise.all([ + repo.setSuccessor("old", target, t.signal), + repo.setSuccessor( + "old", + new URL("https://other.example/actor"), + t.signal, + ), + ]); + assert.strictEqual(results.filter(Boolean).length, 1); + const successor = await repo.getSuccessor("old", t.signal); + assert.ok(successor); + assert.ok(!await repo.setSuccessor("old", successor)); + assert.strictEqual(await repo.getSuccessor("sibling"), undefined); + await assert.rejects( + repo.setSuccessor("sibling", target, AbortSignal.abort()), + { name: "AbortError" }, + ); + const second = new RedisRepository({ + url: redisUrl, + prefix: harness.prefix, + }); + try { + assert.deepStrictEqual(await second.getSuccessor("old"), successor); + assert.ok(!await second.setSuccessor("old", target)); + } finally { + await second.close(); + } + } finally { + await harness.cleanup(); + } + }); + test("key pairs", async () => { const { repository, cleanup } = createHarness(); try { diff --git a/packages/botkit-redis/src/mod.ts b/packages/botkit-redis/src/mod.ts index a0ff8e0..cf3884a 100644 --- a/packages/botkit-redis/src/mod.ts +++ b/packages/botkit-redis/src/mod.ts @@ -440,6 +440,34 @@ export class RedisRepository implements Repository, AsyncDisposable { } } + /** {@inheritDoc Repository.getSuccessor} */ + async getSuccessor( + identifier: string, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + const value = await this.get(this.botKey(identifier, "successor")); + return value === undefined ? undefined : new URL(value); + } + + /** {@inheritDoc Repository.setSuccessor} */ + async setSuccessor( + identifier: string, + successorId: URL, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + const href = successorId.href; + await this.ensureReady(); + signal?.throwIfAborted(); + return await this.command([ + "SET", + this.botKey(identifier, "successor"), + href, + "NX", + ]) === "OK"; + } + async setKeyPairs( identifier: string, keyPairs: CryptoKeyPair[], diff --git a/packages/botkit-sqlite/src/mod.test.ts b/packages/botkit-sqlite/src/mod.test.ts index 54bb032..f7f09e8 100644 --- a/packages/botkit-sqlite/src/mod.test.ts +++ b/packages/botkit-sqlite/src/mod.test.ts @@ -914,3 +914,40 @@ describe("SqliteRepository.migrate() with empty-string identifiers", () => { } }); }); + +test("SQLite successor persists across reopen and cannot be overwritten", async (t) => { + const directory = await mkdtemp(join(tmpdir(), "botkit-move-")); + const path = join(directory, "data.sqlite"); + const target = new URL("https://new.example/actor"); + try { + const repo = new SqliteRepository({ path }); + try { + assert.strictEqual(await repo.getSuccessor("old", t.signal), undefined); + const results = await Promise.all([ + repo.setSuccessor("old", target, t.signal), + repo.setSuccessor( + "old", + new URL("https://other.example/actor"), + t.signal, + ), + ]); + assert.deepStrictEqual(results, [true, false]); + assert.strictEqual(await repo.getSuccessor("sibling"), undefined); + await assert.rejects( + repo.setSuccessor("sibling", target, AbortSignal.abort()), + { name: "AbortError" }, + ); + } finally { + repo.close(); + } + const reopened = new SqliteRepository({ path }); + try { + assert.deepStrictEqual(await reopened.getSuccessor("old"), target); + assert.ok(!await reopened.setSuccessor("old", target)); + } finally { + reopened.close(); + } + } finally { + await rm(directory, { recursive: true, force: true }); + } +}); diff --git a/packages/botkit-sqlite/src/mod.ts b/packages/botkit-sqlite/src/mod.ts index 23a37ce..a94c7ac 100644 --- a/packages/botkit-sqlite/src/mod.ts +++ b/packages/botkit-sqlite/src/mod.ts @@ -351,6 +351,10 @@ export class SqliteRepository implements Repository, Disposable { private initializeTables(): void { this.rebuildLegacyTables(); + this.db.exec(`CREATE TABLE IF NOT EXISTS bot_successors ( + bot_id TEXT PRIMARY KEY NOT NULL, + successor_id TEXT NOT NULL + )`); // Key pairs table this.db.exec(` @@ -485,6 +489,34 @@ export class SqliteRepository implements Repository, Disposable { `); } + /** {@inheritDoc Repository.getSuccessor} */ + async getSuccessor( + identifier: string, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + const row = this.db.prepare( + "SELECT successor_id FROM bot_successors WHERE bot_id = ?", + ).get(identifier); + return await Promise.resolve( + row === undefined ? undefined : new URL(String(row.successor_id)), + ); + } + + /** {@inheritDoc Repository.setSuccessor} */ + async setSuccessor( + identifier: string, + successorId: URL, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + const href = successorId.href; + const result = this.db.prepare( + "INSERT INTO bot_successors (bot_id, successor_id) VALUES (?, ?) ON CONFLICT (bot_id) DO NOTHING", + ).run(identifier, href); + return await Promise.resolve(result.changes > 0); + } + async setKeyPairs( identifier: string, keyPairs: CryptoKeyPair[], diff --git a/packages/botkit/src/bot-impl.ts b/packages/botkit/src/bot-impl.ts index d44d77c..726b59c 100644 --- a/packages/botkit/src/bot-impl.ts +++ b/packages/botkit/src/bot-impl.ts @@ -131,6 +131,8 @@ import { SessionImpl } from "./session-impl.ts"; import type { Session } from "./session.ts"; import type { Text } from "./text.ts"; +import { assertSuccessorRepository } from "./successor.ts"; + const logger = getLogger(["botkit", "bot"]); const uuidPattern = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i; @@ -352,6 +354,7 @@ export class BotImpl implements Bot { preferredUsername: this.username, // Fedify may mutate its array during lazy alias resolution. aliases: [...this.aliases], + successor: await this.repository.getSuccessor(), name: this.name, summary: summary == null ? null : summary.text, attachments: pairs, @@ -659,9 +662,15 @@ export class BotImpl implements Bot { follow, follower, ); + if (await this.repository.getSuccessor() != null) { + await followRequest.reject(); + return; + } await this.onFollow?.(session, followRequest); if (followRequest.state === "pending") { - if (this.followerPolicy === "accept") await followRequest.accept(); + if (await this.repository.getSuccessor() != null) { + await followRequest.reject(); + } else if (this.followerPolicy === "accept") await followRequest.accept(); else if (this.followerPolicy === "reject") await followRequest.reject(); } } @@ -2349,11 +2358,31 @@ export function wrapBotImpl( * @internal */ export class MigrationGatedRepository implements Repository { + /** {@inheritDoc Repository.getSuccessor} */ + async getSuccessor( + identifier: string, + signal?: AbortSignal, + ): Promise { + await this.#migration; + return await this.#repository.getSuccessor(identifier, signal); + } + + /** {@inheritDoc Repository.setSuccessor} */ + async setSuccessor( + identifier: string, + successorId: URL, + signal?: AbortSignal, + ): Promise { + await this.#migration; + return await this.#repository.setSuccessor(identifier, successorId, signal); + } + readonly #repository: Repository; readonly #migration: Promise; readonly #missingQuoteAuthorizationReferenceMethods = new Set(); constructor(repository: Repository, identifier: string) { + assertSuccessorRepository(repository); this.#repository = repository; this.#migration = repository.migrate?.(identifier) ?? Promise.resolve(); // The rejection is re-thrown by the first awaiting operation; this diff --git a/packages/botkit/src/follow-impl.ts b/packages/botkit/src/follow-impl.ts index 69b9002..44b348a 100644 --- a/packages/botkit/src/follow-impl.ts +++ b/packages/botkit/src/follow-impl.ts @@ -50,6 +50,11 @@ export class FollowRequestImpl implements FollowRequest { if (this.#state !== "pending") { throw new TypeError("The follow request is not pending."); } + if (await this.session.bot.repository.getSuccessor() != null) { + throw new TypeError( + "The bot has moved and cannot accept follow requests.", + ); + } await this.session.context.sendActivity( this.session.bot, this.follower, diff --git a/packages/botkit/src/follow-local.test.ts b/packages/botkit/src/follow-local.test.ts index 7065516..8e9bbc3 100644 --- a/packages/botkit/src/follow-local.test.ts +++ b/packages/botkit/src/follow-local.test.ts @@ -237,3 +237,77 @@ test("a sibling's rejection reaches the sender through real HTTP inboxes", async await local.close(); } }); + +test("Session.move migrates local and remote followers through real HTTP", async (t) => { + const local = await createTestServer(t.signal); + const remote = await createTestServer(t.signal); + try { + const old = local.instance.createBot("old", { username: "old" }); + const session = old.getSession(local.origin); + const target = local.instance.createBot("target", { + username: "target", + aliases: [session.actorId], + }); + const targetId = target.getSession(local.origin).actorId; + const localFollower = local.instance.createBot("local", { + username: "local", + }); + const remoteFollower = remote.instance.createBot("remote", { + username: "remote", + }); + const oldActor = await session.getActor(); + const migrations: string[] = []; + for ( + const [follower, server] of [[localFollower, local], [ + remoteFollower, + remote, + ]] as const + ) { + follower.onFolloweeMove = (_session, origin, destination) => { + assert.deepStrictEqual(origin.id, session.actorId); + assert.deepStrictEqual(destination.id, targetId); + migrations.push(follower.identifier); + }; + await follower.getSession(server.origin).follow(oldActor); + assert.ok( + await server.repository.getFollowee( + follower.identifier, + session.actorId, + ), + ); + } + assert.strictEqual(await local.repository.countFollowers("old"), 2); + await session.move(targetId, { signal: t.signal }); + for ( + const [follower, server] of [[localFollower, local], [ + remoteFollower, + remote, + ]] as const + ) { + assert.ok( + await server.repository.getFollowee(follower.identifier, targetId), + ); + assert.strictEqual( + await server.repository.getFollowee( + follower.identifier, + session.actorId, + ), + undefined, + ); + assert.ok( + await local.repository.hasFollower( + "target", + follower.getSession(server.origin).actorId, + ), + ); + } + assert.deepStrictEqual(migrations.sort(), ["local", "remote"]); + assert.deepStrictEqual((await session.getActor()).successorId, targetId); + assert.strictEqual(await local.repository.countFollowers("old"), 0); + assert.deepStrictEqual(local.errors, []); + assert.deepStrictEqual(remote.errors, []); + } finally { + await local.close(); + await remote.close(); + } +}); diff --git a/packages/botkit/src/follow.ts b/packages/botkit/src/follow.ts index 9231866..366e1c6 100644 --- a/packages/botkit/src/follow.ts +++ b/packages/botkit/src/follow.ts @@ -45,7 +45,7 @@ export interface FollowRequest { /** * Accepts the follow request. - * @throws {TypeError} The follow request is not pending. + * @throws {TypeError} If the follow request is not pending or the bot has moved. */ accept(): Promise; diff --git a/packages/botkit/src/instance-impl.ts b/packages/botkit/src/instance-impl.ts index 1cf4c9d..9126909 100644 --- a/packages/botkit/src/instance-impl.ts +++ b/packages/botkit/src/instance-impl.ts @@ -81,6 +81,7 @@ import { app, multiApp } from "./pages.tsx"; import { KvRepository, type Repository } from "./repository.ts"; import type { Session } from "./session.ts"; import { parseLocalUri, rewriteLegacyObjectPath } from "./uri.ts"; +import { assertSuccessorRepository } from "./successor.ts"; interface FolloweeMoveResult { readonly oldActor: Actor; @@ -181,6 +182,7 @@ export class InstanceImpl this.kv = options.kv; this.queue = options.queue; this.repository = options.repository ?? new KvRepository(options.kv); + assertSuccessorRepository(this.repository); this.software = options.software; this.behindProxy = options.behindProxy ?? false; this.pages = { diff --git a/packages/botkit/src/message-impl.ts b/packages/botkit/src/message-impl.ts index 3586b50..97489da 100644 --- a/packages/botkit/src/message-impl.ts +++ b/packages/botkit/src/message-impl.ts @@ -163,6 +163,7 @@ export class MessageImpl async share( options: MessageShareOptions = {}, ): Promise> { + await this.session.ensureActive(); const published = new Date(); const id = uuidv7({ msecs: +published }) as Uuid; const visibility = options.visibility ?? this.visibility; diff --git a/packages/botkit/src/message.ts b/packages/botkit/src/message.ts index b7f198c..bc39c52 100644 --- a/packages/botkit/src/message.ts +++ b/packages/botkit/src/message.ts @@ -205,7 +205,7 @@ export interface Message { * @param options The options for sharing the message. * @returns The shared message. * @throws {TypeError} If the visibility of the message is not `"public"` or - * `"unlisted"`. + * `"unlisted"`, or if the bot has moved. */ share( options?: MessageShareOptions, diff --git a/packages/botkit/src/mod.ts b/packages/botkit/src/mod.ts index 833f14e..d11186f 100644 --- a/packages/botkit/src/mod.ts +++ b/packages/botkit/src/mod.ts @@ -109,6 +109,7 @@ export { export type { Session, SessionGetOutboxOptions, + SessionMoveOptions, SessionPublishOptions, SessionPublishOptionsWithClass, } from "./session.ts"; diff --git a/packages/botkit/src/move.test.ts b/packages/botkit/src/move.test.ts new file mode 100644 index 0000000..5618d68 --- /dev/null +++ b/packages/botkit/src/move.test.ts @@ -0,0 +1,535 @@ +// 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 Context, + type InboxContext, + MemoryKvStore, + type SendActivityOptions, +} from "@fedify/fedify/federation"; +import type { DocumentLoader } from "@fedify/vocab-runtime"; +import { + type Activity, + Article, + Follow, + isActor, + Move, + Note, + Person, + PUBLIC_COLLECTION, + Reject, + Update, +} from "@fedify/vocab"; +import assert from "node:assert/strict"; +import { test } from "node:test"; +import { BotImpl, MigrationGatedRepository } from "./bot-impl.ts"; +import { FollowRequestImpl } from "./follow-impl.ts"; +import { hideRepositoryMethods } from "./helpers.ts"; +import { createInstance } from "./instance.ts"; +import { MemoryCachedRepository, MemoryRepository } from "./repository.ts"; +import { SessionImpl } from "./session-impl.ts"; +import { text } from "./text.ts"; + +function mockLoader( + loader: DocumentLoader, +): Context["getDocumentLoader"] { + function get( + identity: { identifier: string } | { username: string }, + ): Promise; + function get(identity: { keyId: URL; privateKey: CryptoKey }): DocumentLoader; + function get( + identity: { identifier: string } | { username: string } | { + keyId: URL; + privateKey: CryptoKey; + }, + ): Promise | DocumentLoader { + return "keyId" in identity ? loader : Promise.resolve(loader); + } + return get; +} + +function fixture(repository = new MemoryRepository()) { + const bot = new BotImpl({ + kv: new MemoryKvStore(), + repository, + username: "old", + }); + const context = bot.federation.createContext( + new URL("https://old.example"), + undefined, + ); + const session = new SessionImpl(bot, context); + const sent: { + readonly activity: Activity; + readonly options?: SendActivityOptions; + }[] = []; + context.sendActivity = (_sender, _recipients, activity, options) => { + sent.push({ activity, options }); + return Promise.resolve(); + }; + const target = new Person({ + id: new URL("https://new.example/actor"), + inbox: new URL("https://new.example/inbox"), + aliases: [session.actorId], + }); + const fetched: string[] = []; + const loader: DocumentLoader = async (url, options) => { + options?.signal?.throwIfAborted(); + fetched.push(url); + return { + documentUrl: url, + contextUrl: null, + document: await target.toJsonLd({ format: "compact" }), + }; + }; + context.getDocumentLoader = mockLoader(loader); + return { repository, bot, context, session, target, sent, fetched }; +} + +async function seedFollower( + repository: MemoryRepository, + signal?: AbortSignal, +) { + signal?.throwIfAborted(); + await repository.addFollower( + "bot", + new URL("https://follower.example/follow/1"), + new Person({ + id: new URL("https://follower.example/actor"), + inbox: new URL("https://follower.example/inbox"), + }), + ); +} + +test("move persists successor and sends Update then push-mode Move", async (t) => { + const { repository, session, context, target, sent, fetched } = fixture(); + await seedFollower(repository, t.signal); + const previous = await session.publish(text`Keep this post.`); + sent.length = 0; + await session.move(target, { signal: t.signal }); + assert.deepStrictEqual(fetched, [target.id!.href]); + assert.deepStrictEqual(await repository.getSuccessor("bot"), target.id); + assert.deepStrictEqual((await session.getActor()).successorId, target.id); + assert.strictEqual(sent.length, 2); + assert.ok(sent[0].activity instanceof Update); + const updatedActor = await sent[0].activity.getObject(context); + assert.ok(isActor(updatedActor)); + assert.deepStrictEqual(updatedActor.successorId, target.id); + const move = sent[1].activity; + assert.ok(move instanceof Move); + assert.ok(move.id); + assert.notStrictEqual(move.id.href, sent[0].activity.id?.href); + assert.deepStrictEqual(move.actorId, session.actorId); + assert.deepStrictEqual(move.objectId, session.actorId); + assert.deepStrictEqual(move.targetId, target.id); + assert.deepStrictEqual(move.toIds, [PUBLIC_COLLECTION]); + assert.deepStrictEqual(move.ccIds, [context.getFollowersUri("bot")]); + for (const entry of sent) { + assert.deepStrictEqual(entry.options?.excludeBaseUris, []); + assert.strictEqual(entry.options?.orderingKey, session.actorId.href); + } + await assert.rejects(session.move(target), TypeError); + await assert.rejects(session.publish(text`Inactive.`), TypeError); + await assert.rejects( + session.publish(text`Inactive.`, { class: Article }), + TypeError, + ); + await assert.rejects(previous.reply(text`Inactive reply.`), TypeError); + await assert.rejects(previous.share(), TypeError); + assert.strictEqual(await repository.countMessages("bot"), 1); + assert.strictEqual(await repository.countFollowers("bot"), 1); + assert.strictEqual(sent.length, 2); + const fresh = fixture(repository); + assert.deepStrictEqual( + (await fresh.session.getActor()).successorId, + target.id, + ); + await assert.rejects(fresh.session.publish(text`Still inactive.`), TypeError); +}); + +test("move validates authoritative aliases, not the passed actor", async () => { + const { session, context, target, repository, sent } = fixture(); + const actual = new Person({ id: target.id, inbox: target.inboxId }); + context.getDocumentLoader = mockLoader(async (url) => ({ + documentUrl: url, + contextUrl: null, + document: await actual.toJsonLd({ format: "compact" }), + })); + await assert.rejects(session.move(target), TypeError); + assert.strictEqual(await repository.getSuccessor("bot"), undefined); + assert.strictEqual(sent.length, 0); +}); + +test("move rejects invalid targets without changing state", async (t) => { + for ( + const kind of [ + "self", + "idless", + "no-inbox", + "moved", + "mismatched-id", + "cross-origin", + "non-actor", + "scheme", + ] as const + ) { + await t.test(kind, async () => { + const { session, context, target, repository } = fixture(); + let object = target; + if (kind === "no-inbox") { + object = new Person({ id: target.id, aliases: [session.actorId] }); + } + if (kind === "moved") { + object = target.clone({ + successor: new URL("https://next.example/actor"), + }); + } + if (kind === "mismatched-id") { + object = target.clone({ id: new URL("https://new.example/other") }); + } + context.getDocumentLoader = mockLoader(async (url) => ({ + documentUrl: kind === "cross-origin" + ? "https://evil.example/actor" + : url, + contextUrl: null, + document: + await (kind === "non-actor" ? new Note({ id: target.id }) : object) + .toJsonLd({ format: "compact" }), + })); + const input = kind === "self" + ? session.actorId + : kind === "idless" + ? new Person({}) + : kind === "scheme" + ? new URL("javascript:alert(1)") + : target; + await assert.rejects(session.move(input), TypeError); + assert.strictEqual(await repository.getSuccessor("bot"), undefined); + }); + } +}); + +test("move resolves handles and URI strings, preserving fetch failures", async (t) => { + for (const form of ["uri", "handle"] as const) { + await t.test(form, async () => { + const { session, context, target, fetched } = fixture(); + context.lookupObject = () => Promise.resolve(target); + await session.move(form === "uri" ? target.id!.href : "@new@new.example"); + assert.deepStrictEqual(fetched, [target.id!.href]); + }); + } + const { session, context, target, repository } = fixture(); + const failure = new TypeError("Network unavailable."); + context.getDocumentLoader = mockLoader(() => Promise.reject(failure)); + await assert.rejects(session.move(target.id!), (error) => error === failure); + assert.strictEqual(await repository.getSuccessor("bot"), undefined); +}); + +test("concurrent move calls choose one successor and one notification pair", async (t) => { + const f = fixture(); + await seedFollower(f.repository, t.signal); + const second = new SessionImpl(f.bot, f.context); + const results = await Promise.allSettled([ + f.session.move(f.target), + second.move(f.target), + ]); + assert.strictEqual(results.filter((r) => r.status === "fulfilled").length, 1); + assert.strictEqual(f.sent.length, 2); +}); + +test("move without followers still persists inactive state", async () => { + const { session, target, sent } = fixture(); + await session.move(target); + assert.deepStrictEqual((await session.getActor()).successorId, target.id); + assert.strictEqual(sent.length, 0); +}); + +test("move attempts both notifications on failure and republishMove recovers", async (t) => { + for (const failing of [Update, Move]) { + await t.test(failing.name, async () => { + const f = fixture(); + await seedFollower(f.repository, t.signal); + f.context.sendActivity = (_sender, _recipients, activity) => { + f.sent.push({ activity }); + return activity instanceof failing + ? Promise.reject(new TypeError("Delivery failed.")) + : Promise.resolve(); + }; + await assert.rejects(f.session.move(f.target), AggregateError); + assert.strictEqual(f.sent.length, 2); + assert.deepStrictEqual( + (await f.session.getActor()).successorId, + f.target.id, + ); + f.context.sendActivity = (_sender, _recipients, activity) => { + f.sent.push({ activity }); + return Promise.resolve(); + }; + await f.session.republishMove(); + assert.strictEqual(f.sent.length, 4); + assert.notStrictEqual( + f.sent[1].activity.id?.href, + f.sent[3].activity.id?.href, + ); + }); + } + const f = fixture(); + await assert.rejects(f.session.republishMove(), TypeError); +}); + +test("move honours pre-commit abort but completes after committing", async (t) => { + const f = fixture(); + await assert.rejects( + f.session.move(f.target, { signal: AbortSignal.abort() }), + { name: "AbortError" }, + ); + assert.strictEqual(await f.repository.getSuccessor("bot"), undefined); + const controller = new AbortController(); + const set = f.repository.setSuccessor.bind(f.repository); + f.repository.setSuccessor = async (id, target, signal) => { + const result = await set(id, target, signal); + controller.abort(); + return result; + }; + await seedFollower(f.repository, t.signal); + await f.session.move(f.target, { signal: controller.signal }); + assert.strictEqual(f.sent.length, 2); + await assert.rejects( + f.session.republishMove({ signal: AbortSignal.abort() }), + { name: "AbortError" }, + ); + const resend = new AbortController(); + f.context.sendActivity = (_sender, _recipients, activity) => { + f.sent.push({ activity }); + resend.abort(); + return Promise.resolve(); + }; + await f.session.republishMove({ signal: resend.signal }); + assert.strictEqual(f.sent.length, 4); +}); + +test("local target URLs, Actors and handles need no self HTTP", async (t) => { + for (const form of ["url", "actor", "handle"] as const) { + await t.test(form, async () => { + const instance = createInstance({ + kv: new MemoryKvStore(), + repository: new MemoryRepository(), + }); + const old = instance.createBot("old", { username: "old" }); + const session = old.getSession("https://local.example"); + const target = instance.createBot("target", { + username: "new", + aliases: [session.actorId], + }); + const targetSession = target.getSession("https://local.example"); + session.context.lookupObject = () => + Promise.reject(new TypeError("Self HTTP is forbidden.")); + session.context.getDocumentLoader = mockLoader(() => + Promise.reject(new TypeError("Self HTTP is forbidden.")) + ); + await session.move( + form === "url" + ? targetSession.actorId + : form === "actor" + ? await targetSession.getActor() + : "@new@local.example", + ); + assert.deepStrictEqual( + (await session.getActor()).successorId, + targetSession.actorId, + ); + }); + } +}); + +test("moved bots reject follows under every policy, without onFollow", async (t) => { + for (const policy of ["accept", "manual", "reject"] as const) { + await t.test(policy, async () => { + const repository = new MemoryRepository(); + const bot = new BotImpl({ + kv: new MemoryKvStore(), + repository, + username: "old", + followerPolicy: policy, + }); + const context: InboxContext = Object.assign( + bot.federation.createContext(new URL("https://old.example"), undefined), + { + recipient: "bot", + forwardActivity: () => Promise.resolve(), + clone: (_data: void): InboxContext => context, + }, + ); + const sent: Activity[] = []; + context.sendActivity = (_sender, _recipients, activity) => { + sent.push(activity); + return Promise.resolve(); + }; + bot.onFollow = () => assert.fail("Moved onFollow must not fire."); + const follower = new Person({ + id: new URL("https://follower.example/actor"), + inbox: new URL("https://follower.example/inbox"), + }); + const follow = new Follow({ + id: new URL("https://follower.example/follow/1"), + actor: follower, + object: context.getActorUri("bot"), + }); + const request = new FollowRequestImpl( + new SessionImpl(bot, context), + follow, + follower, + ); + await repository.setSuccessor( + "bot", + new URL("https://new.example/actor"), + ); + await assert.rejects(request.accept(), TypeError); + assert.strictEqual(request.state, "pending"); + await bot.onFollowed(context, follow); + assert.strictEqual(sent.length, 1); + assert.ok(sent[0] instanceof Reject); + assert.strictEqual(await repository.countFollowers("bot"), 0); + await request.reject(); + }); + } +}); + +test("custom repositories fail before missing successor methods are obscured", () => { + for (const method of ["getSuccessor", "setSuccessor"]) { + const repository = hideRepositoryMethods(new MemoryRepository(), [method]); + assert.throws( + () => createInstance({ kv: new MemoryKvStore(), repository }), + TypeError, + ); + assert.throws( + () => new MigrationGatedRepository(repository, "bot"), + TypeError, + ); + assert.throws(() => new MemoryCachedRepository(repository), TypeError); + } +}); + +test("publishing cannot cross a move while rendering content", async () => { + const f = fixture(); + const original = text`Slow content.`; + const content = { + ...original, + type: "block" as const, + async *getHtml() { + await f.session.move(f.target); + yield "Slow content."; + }, + getTags: original.getTags.bind(original), + getCachedObjects: original.getCachedObjects.bind(original), + }; + await assert.rejects(f.session.publish(content), TypeError); + assert.strictEqual(await f.repository.countMessages("bot"), 0); +}); + +test("migration notifications keep their snapshot when followers change", async (t) => { + const f = fixture(); + await seedFollower(f.repository, t.signal); + const audiences: string[][] = []; + f.context.sendActivity = async (_sender, recipients, activity) => { + assert.ok(Array.isArray(recipients)); + audiences.push(recipients.map((actor) => actor.id!.href)); + if (activity instanceof Update) { + await f.repository.removeFollower( + "bot", + new URL("https://follower.example/follow/1"), + new URL("https://follower.example/actor"), + ); + } + }; + await f.session.move(f.target); + assert.deepStrictEqual(audiences, [["https://follower.example/actor"], [ + "https://follower.example/actor", + ]]); +}); + +test("republishMove revalidates aliases, allows moved successors and aborts preparation", async (t) => { + const f = fixture(); + await seedFollower(f.repository, t.signal); + await f.session.move(f.target); + f.sent.length = 0; + const movedTarget = f.target.clone({ + successor: new URL("https://next.example/actor"), + }); + f.context.getDocumentLoader = mockLoader(async (url) => ({ + documentUrl: url, + contextUrl: null, + document: await movedTarget.toJsonLd({ format: "compact" }), + })); + await f.session.republishMove(); + assert.strictEqual(f.sent.length, 2); + f.sent.length = 0; + const controller = new AbortController(); + const getFollowers = f.repository.getFollowers.bind(f.repository); + f.repository.getFollowers = (identifier, options) => { + controller.abort(); + return getFollowers(identifier, options); + }; + await assert.rejects(f.session.republishMove({ signal: controller.signal }), { + name: "AbortError", + }); + assert.strictEqual(f.sent.length, 0); + const unlinked = new Person({ id: f.target.id, inbox: f.target.inboxId }); + f.context.getDocumentLoader = mockLoader(async (url) => ({ + documentUrl: url, + contextUrl: null, + document: await unlinked.toJsonLd({ format: "compact" }), + })); + await assert.rejects(f.session.republishMove(), TypeError); + assert.strictEqual(f.sent.length, 0); +}); + +test("move storage errors leave their actual commit state inspectable", async (t) => { + const f = fixture(); + await seedFollower(f.repository, t.signal); + const set = f.repository.setSuccessor.bind(f.repository); + const failure = new TypeError("Storage response lost."); + f.repository.setSuccessor = async (identifier, target, signal) => { + await set(identifier, target, signal); + throw failure; + }; + await assert.rejects(f.session.move(f.target), (error) => error === failure); + assert.deepStrictEqual((await f.session.getActor()).successorId, f.target.id); + assert.strictEqual(f.sent.length, 0); + await f.session.republishMove(); + assert.strictEqual(f.sent.length, 2); +}); + +test("dynamic targets resolve without HTTP and remain separate from the source", async () => { + const repository = new MemoryRepository(); + const instance = createInstance({ + kv: new MemoryKvStore(), + repository, + }); + const old = instance.createBot("old", { username: "old" }); + const session = old.getSession("https://local.example"); + instance.createBot((_context, identifier) => + identifier === "new" + ? { username: "new", aliases: [session.actorId] } + : null + ); + session.context.getDocumentLoader = mockLoader(() => + Promise.reject(new TypeError("Self fetch is forbidden.")) + ); + const targetId = session.context.getActorUri("new"); + await session.move(targetId); + assert.deepStrictEqual((await session.getActor()).successorId, targetId); + assert.strictEqual(await repository.getSuccessor("new"), undefined); +}); diff --git a/packages/botkit/src/pages.test.ts b/packages/botkit/src/pages.test.ts index fdb871f..0a60568 100644 --- a/packages/botkit/src/pages.test.ts +++ b/packages/botkit/src/pages.test.ts @@ -390,3 +390,87 @@ describe("message rendering", () => { assert.ok(!html.includes("javascript:alert")); }); }); + +test("moved profiles show their successor and reject direct follow forms", async () => { + const repository = new MemoryRepository(); + const instance = new InstanceImpl({ + kv: new MemoryKvStore(), + repository, + }); + instance.createBot("old", { username: "old" }); + instance.createBot("active", { username: "active" }); + const successor = new URL("https://new.example/actor?name=Alice&from=old"); + await seedMessage(repository, "old"); + await repository.setSuccessor("old", successor); + const profile = await instance.fetch( + new Request("https://example.com/@old"), + undefined, + ); + const html = await profile.text(); + assert.ok(html.includes("This bot has moved")); + assert.ok(html.includes("https://new.example/actor?name=Alice&from=old")); + assert.ok(!html.includes('action="/@old/follow"')); + assert.ok(html.includes("Hello, world!")); + const active = await instance.fetch( + new Request("https://example.com/@active"), + undefined, + ); + assert.ok((await active.text()).includes('action="/@active/follow"')); + const follow = await instance.fetch( + new Request("https://example.com/@old/follow", { method: "POST" }), + undefined, + ); + assert.strictEqual(follow.status, 409); + assert.ok((await follow.text()).includes("This bot has moved")); +}); + +test("single-bot moved profile escapes unsafe successor schemes", async () => { + const repository = new MemoryRepository(); + const bot = new BotImpl({ + kv: new MemoryKvStore(), + repository, + username: "old", + }); + await repository.setSuccessor("bot", new URL("javascript:alert(1)")); + const profile = await bot.fetch( + new Request("https://example.com/"), + undefined, + ); + const html = await profile.text(); + assert.ok(html.includes("This bot has moved")); + assert.ok(!html.includes('href="javascript:')); + assert.ok(!html.includes('action="/follow"')); +}); + +test("moved notices link local successors to their public profile", async () => { + const repository = new MemoryRepository(); + const instance = new InstanceImpl({ + kv: new MemoryKvStore(), + repository, + }); + instance.createBot("old", { username: "old" }); + instance.createBot("target-id", { username: "new-name" }); + const target = new URL("https://example.com/ap/actor/target-id"); + await repository.setSuccessor("old", target); + for ( + const request of [ + new Request("https://example.com/@old"), + new Request("https://example.com/@old/follow", { method: "POST" }), + ] + ) { + const response = await instance.fetch(request, undefined); + const html = await response.text(); + const link = + /(?:This bot has moved to|Its new actor is)\s*(); app.get("/", (c) => profilePage(c, c.env.bot, c.env.contextData, "")); +// Actor URIs identify the successor; local bots have separate HTML profiles. +async function successorWebUrl( + bot: BotImpl, + ctx: Context, + successor: URL, + signal?: AbortSignal, +): Promise { + signal?.throwIfAborted(); + const parsed = ctx.parseUri(successor); + if (parsed?.type !== "actor") return successor; + const target = await bot.instance.resolveBot(ctx, parsed.identifier); + signal?.throwIfAborted(); + return target == null + ? successor + : bot.instance.getBotWebUrl(target, ctx.origin); +} + async function profilePage( c: PageContext, bot: BotImpl, @@ -111,6 +128,10 @@ async function profilePage( : bot.image; const imageWidth = bot.image instanceof Image ? bot.image.width : null; const imageHeight = bot.image instanceof Image ? bot.image.height : null; + const successor = await bot.repository.getSuccessor(); + const successorUrl = successor == null + ? undefined + : await successorWebUrl(bot, ctx, successor); const followersCount = await bot.repository.countFollowers(); const summaryChunks = bot.summary?.getHtml(session); const postsCount = await bot.repository.countMessages(); @@ -192,6 +213,14 @@ async function profilePage( + {successor != null && ( +

+ This bot has moved to {successor.protocol === "http:" || + successor.protocol === "https:" + ? {successor.href} + : successor.href}. +

+ )} {summary && (
- + {successor == null && ( + + )} @@ -568,6 +599,30 @@ async function followPage( const ctx = bot.federation.createContext(c.req.raw, contextData); const url = new URL(c.req.url); + const successor = await bot.repository.getSuccessor(); + const successorUrl = successor == null + ? undefined + : await successorWebUrl(bot, ctx, successor); + if (successor != null) { + return c.html( + +
+

This bot has moved

+

+ Its new actor is{" "} + {successor.protocol === "http:" || successor.protocol === "https:" + ? {successor.href} + : successor.href}. +

+

+ Go back +

+
+
, + 409, + ); + } + const formData = await c.req.formData(); let followerHandle = formData.get("handle")?.toString(); diff --git a/packages/botkit/src/repository.test.ts b/packages/botkit/src/repository.test.ts index 71ebdbe..0a320fd 100644 --- a/packages/botkit/src/repository.test.ts +++ b/packages/botkit/src/repository.test.ts @@ -2430,3 +2430,80 @@ describe("KvRepository.migrate()", () => { ); }); }); + +for (const [name, factory] of Object.entries(factories)) { + test(`${name} successor is write-once and bot-scoped`, async (t) => { + const repository = factory(); + const target = new URL("https://new.example/actor"); + assert.strictEqual( + await repository.getSuccessor("old", t.signal), + undefined, + ); + const results = await Promise.all([ + repository.setSuccessor("old", target, t.signal), + repository.setSuccessor( + "old", + new URL("https://other.example/actor"), + t.signal, + ), + ]); + assert.strictEqual(results.filter(Boolean).length, 1); + const successor = await repository.getSuccessor("old", t.signal); + assert.ok(successor); + assert.ok(!await repository.setSuccessor("old", successor, t.signal)); + assert.strictEqual(await repository.getSuccessor("sibling"), undefined); + const href = successor.href; + successor.pathname = "/mutated"; + target.pathname = "/mutated"; + assert.strictEqual((await repository.getSuccessor("old"))?.href, href); + const signal = AbortSignal.abort(); + await assert.rejects(repository.setSuccessor("sibling", target, signal), { + name: "AbortError", + }); + assert.strictEqual(await repository.getSuccessor("sibling"), undefined); + }); +} + +test("MemoryCachedRepository observes successor changes through other views", async () => { + const underlying = new MemoryRepository(); + const first = new MemoryCachedRepository(underlying); + const second = new MemoryCachedRepository(underlying); + assert.strictEqual(await first.getSuccessor("old"), undefined); + const target = new URL("https://new.example/actor"); + assert.ok(await second.setSuccessor("old", target)); + assert.deepStrictEqual(await first.getSuccessor("old"), target); + const scoped = first.forIdentifier("old"); + assert.deepStrictEqual(await scoped.getSuccessor(), target); + assert.ok(!await scoped.setSuccessor(target)); +}); + +test("KvRepository successor CAS coordinates independent views", async () => { + const kv = new MemoryKvStore(); + const first = new KvRepository(kv); + const second = new KvRepository(kv); + const results = await Promise.all([ + first.setSuccessor("bot", new URL("https://new.example/first")), + second.setSuccessor("bot", new URL("https://new.example/second")), + ]); + assert.strictEqual(results.filter(Boolean).length, 1); + assert.deepStrictEqual( + await first.getSuccessor("bot"), + await second.getSuccessor("bot"), + ); +}); + +test("KvRepository successor fallback serializes within one non-CAS view", async () => { + const kv = new Proxy(new MemoryKvStore(), { + get(target, property, receiver) { + if (property === "cas") return undefined; + const value = Reflect.get(target, property, receiver); + return typeof value === "function" ? value.bind(target) : value; + }, + }); + const repository = new KvRepository(kv); + const results = await Promise.all([ + repository.setSuccessor("bot", new URL("https://new.example/first")), + repository.setSuccessor("bot", new URL("https://new.example/second")), + ]); + assert.deepStrictEqual(results, [true, false]); +}); diff --git a/packages/botkit/src/repository.ts b/packages/botkit/src/repository.ts index 179b9e1..46639fd 100644 --- a/packages/botkit/src/repository.ts +++ b/packages/botkit/src/repository.ts @@ -30,6 +30,8 @@ import { getLogger } from "@logtape/logtape"; export type { KvKey, KvStore } from "@fedify/fedify/federation"; export { Announce, Create } from "@fedify/vocab"; +import { assertSuccessorRepository } from "./successor.ts"; + const logger = getLogger(["botkit", "repository"]); const kvLockPollIntervalMs = 100; const uuidPattern = @@ -92,6 +94,34 @@ interface QuoteAuthorizationReferenceData { * @since 0.3.0 */ export interface Repository { + /** + * Gets the actor URI to which a bot has moved. + * @param identifier The owning bot actor's identifier. + * @param signal The signal to cancel before reading. + * @returns The successor URI, or `undefined` if the bot has not moved. + * @since 0.6.0 + */ + getSuccessor( + identifier: string, + signal?: AbortSignal, + ): Promise; + + /** + * Records a bot's successor only if it has not already moved. + * This must be an atomic first-write-wins operation. Existing values, + * including the same URI, must never be overwritten. + * @param identifier The owning bot actor's identifier. + * @param successorId The successor's actor URI. + * @param signal The signal to cancel before committing the write. + * @returns `true` if stored, `false` if a successor already exists. + * @since 0.6.0 + */ + setSuccessor( + identifier: string, + successorId: URL, + signal?: AbortSignal, + ): Promise; + /** * Sets the key pairs of a bot actor. * @param identifier The identifier of the bot actor that owns the key pairs. @@ -510,6 +540,27 @@ export interface Repository { * @since 0.5.0 */ export class ActorScopedRepository { + /** + * Gets this bot's successor actor URI. + * @param signal The signal to cancel before reading. + * @returns The successor URI, or `undefined` if the bot has not moved. + * @since 0.6.0 + */ + getSuccessor(signal?: AbortSignal): Promise { + return this.repository.getSuccessor(this.identifier, signal); + } + + /** + * Records this bot's successor only if it has not already moved. + * @param successorId The successor's actor URI. + * @param signal The signal to cancel before committing the write. + * @returns Whether the successor was stored. + * @since 0.6.0 + */ + setSuccessor(successorId: URL, signal?: AbortSignal): Promise { + return this.repository.setSuccessor(this.identifier, successorId, signal); + } + /** * The underlying repository. */ @@ -939,6 +990,37 @@ export interface KvRepositoryOptions { * A repository for storing bot data using a key-value store. */ export class KvRepository implements Repository { + /** {@inheritDoc Repository.getSuccessor} */ + async getSuccessor( + identifier: string, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + const value = await this.kv.get(this.#key(identifier, "successor")); + return value === undefined ? undefined : new URL(value); + } + + /** + * {@inheritDoc Repository.setSuccessor} + * Without CAS, writes are serialized only within this repository instance. + */ + async setSuccessor( + identifier: string, + successorId: URL, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + const key = this.#key(identifier, "successor"); + const href = successorId.href; + if (this.kv.cas != null) return await this.kv.cas(key, undefined, href); + return await this.#withNonCasKvLock(key, async () => { + if (await this.kv.get(key) !== undefined) return false; + signal?.throwIfAborted(); + await this.kv.set(key, href); + return true; + }); + } + readonly kv: KvStore; /** @@ -2132,6 +2214,7 @@ function extractTimestamp(uuid: string): number { } interface MemoryActorData { + successor?: string; keyPairs?: CryptoKeyPair[]; messages: Map; followers: Map; @@ -2150,6 +2233,31 @@ interface MemoryActorData { * persistent and is only suitable for testing or development. */ export class MemoryRepository implements Repository { + /** {@inheritDoc Repository.getSuccessor} */ + async getSuccessor( + identifier: string, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + const value = this.#data.get(identifier)?.successor; + return await Promise.resolve( + value === undefined ? undefined : new URL(value), + ); + } + + /** {@inheritDoc Repository.setSuccessor} */ + async setSuccessor( + identifier: string, + successorId: URL, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + const data = this.#bucket(identifier); + if (data.successor !== undefined) return await Promise.resolve(false); + data.successor = successorId.href; + return await Promise.resolve(true); + } + #data: Map = new Map(); #bucket(identifier: string): MemoryActorData { @@ -2567,6 +2675,25 @@ export class MemoryRepository implements Repository { * @since 0.3.0 */ export class MemoryCachedRepository implements Repository { + // Successor reads intentionally bypass the cache: another repository view + // may have moved the bot since a previous read. + /** {@inheritDoc Repository.getSuccessor} */ + getSuccessor( + identifier: string, + signal?: AbortSignal, + ): Promise { + return this.underlying.getSuccessor(identifier, signal); + } + + /** {@inheritDoc Repository.setSuccessor} */ + setSuccessor( + identifier: string, + successorId: URL, + signal?: AbortSignal, + ): Promise { + return this.underlying.setSuccessor(identifier, successorId, signal); + } + private underlying: Repository; private cache: MemoryRepository; @@ -2577,6 +2704,7 @@ export class MemoryCachedRepository implements Repository { * If not provided, a new one will be created internally. */ constructor(underlying: Repository, cache?: MemoryRepository) { + assertSuccessorRepository(underlying); this.underlying = underlying; this.cache = cache ?? new MemoryRepository(); } diff --git a/packages/botkit/src/session-impl.ts b/packages/botkit/src/session-impl.ts index 9faa5b9..1a2d725 100644 --- a/packages/botkit/src/session-impl.ts +++ b/packages/botkit/src/session-impl.ts @@ -16,7 +16,7 @@ import "./temporal.ts"; import type { Context } from "@fedify/fedify/federation"; import { quoteInteraction } from "@fedify/interaction-controls"; -import { LanguageString } from "@fedify/vocab-runtime"; +import { type DocumentLoader, LanguageString } from "@fedify/vocab-runtime"; import { type Actor, Collection, @@ -25,8 +25,10 @@ import { isActor, Link, Mention, + Move, Note, type Object, + Object as APObject, PUBLIC_COLLECTION, QuoteRequest, Undo, @@ -52,6 +54,7 @@ import { serializeQuotePolicy } from "./quote.ts"; import type { Session, SessionGetOutboxOptions, + SessionMoveOptions, SessionPublishOptions, SessionPublishOptionsWithClass, SessionPublishOptionsWithQuestion, @@ -226,6 +229,194 @@ export class SessionImpl implements Session { return follow != null; } + async move( + target: Actor | URL | string, + options: SessionMoveOptions = {}, + ): Promise { + const signal = options.signal; + signal?.throwIfAborted(); + if (await this.bot.repository.getSuccessor(signal) != null) { + throw new TypeError("The bot has already moved."); + } + const actor = await this.#resolveMoveTarget(target, false, signal); + const successorId = actor.id!; + signal?.throwIfAborted(); + if (!await this.bot.repository.setSuccessor(successorId, signal)) { + throw new TypeError("The bot has already moved."); + } + // From here cancellation must not interrupt the notification pair. + // Never roll back a move which another server may already have processed. + try { + await this.#notifyMove(successorId); + } catch (error) { + if (error instanceof AggregateError) throw error; + throw new AggregateError( + [error], + "The bot moved, but its migration notifications failed.", + ); + } + } + + async republishMove(options: SessionMoveOptions = {}): Promise { + const signal = options.signal; + signal?.throwIfAborted(); + const successorId = await this.bot.repository.getSuccessor(signal); + if (successorId == null) throw new TypeError("The bot has not moved."); + await this.#resolveMoveTarget(successorId, true, signal); + await this.#notifyMove(successorId, signal); + } + + async #resolveMoveTarget( + target: Actor | URL | string, + allowMoved: boolean, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + let id: URL | null; + if (isActor(target)) { + id = target.id; + } else if (target instanceof URL) { + id = target; + } else { + const handle = target.replace(/^acct:/, "").match( + /^@?([^@/:]+)@([^@/]+)$/, + ); + if ( + handle != null && + handle[2].toLowerCase() === this.context.host.toLowerCase() + ) { + const bot = await this.bot.instance.resolveBotByUsername( + this.context, + handle[1], + ); + if (bot == null) { + throw new TypeError("The migration target could not be resolved."); + } + id = this.context.getActorUri(bot.identifier); + } else if (handle != null) { + const documentLoader = await this.context.getDocumentLoader(this.bot); + const actor = await this.context.lookupObject(target, { + documentLoader, + signal, + }); + signal?.throwIfAborted(); + if (!isActor(actor)) { + throw new TypeError("The migration target could not be resolved."); + } + id = actor.id; + } else { + id = new URL(target); + } + } + if (id == null) { + throw new TypeError("The migration target does not have an ID."); + } + if (id.protocol !== "http:" && id.protocol !== "https:") { + throw new TypeError( + "The migration target must have an HTTP or HTTPS actor URI.", + ); + } + if (id.href === this.actorId.href) { + throw new TypeError("The bot cannot move to itself."); + } + let actor: APObject | null; + const local = this.context.parseUri(id); + if (local?.type === "actor") { + const bot = await this.bot.instance.resolveBot( + this.context, + local.identifier, + ); + actor = await bot?.dispatchActor(this.context, local.identifier) ?? null; + } else { + const loader = await this.context.getDocumentLoader(this.bot); + const documentLoader: DocumentLoader = (url, options) => + loader(url, { ...options, signal: signal ?? options?.signal }); + const document = await documentLoader(id.href); + const documentUrl = new URL(document.documentUrl); + if (documentUrl.origin !== id.origin) { + throw new TypeError( + "The migration target document has a different origin.", + ); + } + actor = await APObject.fromJsonLd(document.document, { + baseUrl: documentUrl, + documentLoader, + contextLoader: (url, options) => + this.context.contextLoader(url, { + ...options, + signal: signal ?? options?.signal, + }), + }); + } + signal?.throwIfAborted(); + if ( + !isActor(actor) || actor.id?.href !== id.href || actor.inboxId == null + ) { + throw new TypeError( + "The migration target is not a valid actor with an inbox.", + ); + } + if (!allowMoved && actor.successorId != null) { + throw new TypeError("The migration target has already moved."); + } + if (!actor.aliasIds.some((alias) => alias.href === this.actorId.href)) { + throw new TypeError( + "The migration target does not list the bot as an alias.", + ); + } + return actor; + } + + async #notifyMove(successorId: URL, signal?: AbortSignal): Promise { + signal?.throwIfAborted(); + // Freeze the audience for both sends: a local Move can cause immediate + // Undo(Follow) deliveries which mutate the live followers collection. + const followers = await Array.fromAsync(this.bot.repository.getFollowers()); + const actor = await this.getActor(); + const followersId = this.context.getFollowersUri(this.bot.identifier); + const update = new Update({ + id: new URL(`#update-profile/${crypto.randomUUID()}`, this.actorId), + actor: this.actorId, + to: followersId, + object: actor, + }); + const move = new Move({ + id: new URL(`#move/${crypto.randomUUID()}`, this.actorId), + actor: this.actorId, + object: this.actorId, + target: successorId, + to: PUBLIC_COLLECTION, + cc: followersId, + }); + signal?.throwIfAborted(); + if (followers.length === 0) return; + const errors: unknown[] = []; + for (const activity of [update, move]) { + try { + await this.context.sendActivity(this.bot, followers, activity, { + preferSharedInbox: true, + excludeBaseUris: [], + orderingKey: this.actorId.href, + }); + } catch (error) { + errors.push(error); + } + } + if (errors.length > 0) { + throw new AggregateError( + errors, + "The migration notifications could not all be submitted.", + ); + } + } + + /** Checks whether this bot may publish new messages. @internal */ + async ensureActive(signal?: AbortSignal): Promise { + if (await this.bot.repository.getSuccessor(signal) != null) { + throw new TypeError("The bot has moved and cannot publish new messages."); + } + } + async republishProfile(): Promise { const actor = await this.getActor(); const update = new Update({ @@ -264,6 +455,7 @@ export class SessionImpl implements Session { | SessionImplPublishOptionsWithClass | SessionImplPublishOptionsWithQuestion = {}, ): Promise> { + await this.ensureActive(); const published = new Date(); const id = uuidv7({ msecs: +published }) as Uuid; const cls = "class" in options ? options.class : Note; @@ -418,6 +610,9 @@ export class SessionImpl implements Session { object: msg, published: published.toTemporalInstant(), }); + // Rendering text can await remote lookups or user code. Recheck before + // storing anything in case the bot moved during that preparation. + await this.ensureActive(); await this.bot.repository.addMessage(id, activity); const preferSharedInbox = visibility === "public" || visibility === "unlisted" || visibility === "followers"; diff --git a/packages/botkit/src/session.ts b/packages/botkit/src/session.ts index 8476ab7..0177066 100644 --- a/packages/botkit/src/session.ts +++ b/packages/botkit/src/session.ts @@ -110,6 +110,40 @@ export interface Session { */ follows(actor: Actor | URL | string): Promise; + /** + * Moves the bot's followers to a linked actor and makes the bot inactive. + * The destination must advertise this bot's actor URI in `alsoKnownAs`. + * Cancellation is honoured until the successor is committed; notifications + * continue after that point. On failure, check {@link Session.getActor}'s + * `successorId` before retrying or calling {@link Session.republishMove}. + * @param target The destination actor, actor URI, or fediverse handle. + * @param options Options for cancelling preparation. + * @returns A promise that resolves after submitting the notifications. + * @throws {TypeError} If the bot has moved, or the target is invalid, + * inactive, the bot itself, or does not list its alias. + * @throws {AggregateError} If notification preparation or submission fails + * after the successor has been committed. + * @since 0.6.0 + */ + move( + target: Actor | URL | string, + options?: SessionMoveOptions, + ): Promise; + + /** + * Resends the stored account migration to the remaining followers. + * Revalidates the destination's alias without changing the successor. + * Cancellation is honoured through preparation, until the first activity + * is submitted; both notifications continue once submission starts. + * @param options Options for cancelling preparation. + * @returns A promise that resolves after submitting the notifications. + * @throws {TypeError} If the bot has not moved or its successor is invalid + * or no longer lists the bot as an alias. + * @throws {AggregateError} If notification submission fails. + * @since 0.6.0 + */ + republishMove(options?: SessionMoveOptions): Promise; + /** * Republishes the bot profile to its followers. * @@ -127,6 +161,7 @@ export interface Session { * @param text The content of the note. * @param options The options for publishing the message. * @returns The published message. + * @throws {TypeError} If the bot has moved. */ publish( content: Text<"block", TContextData>, @@ -139,6 +174,7 @@ export interface Session { * @param text The content of the note. * @param options The options for publishing the message. * @returns The published message. + * @throws {TypeError} If the bot has moved. */ publish( content: Text<"block", TContextData>, @@ -150,6 +186,7 @@ export interface Session { * @param content The content of the question. * @param options The options for publishing the question. * @returns The published question. + * @throws {TypeError} If the bot has moved. * @since 0.3.0 */ publish( @@ -167,6 +204,15 @@ export interface Session { ): AsyncIterable>; } +/** + * Options for preparing an account migration or resending its notifications. + * @since 0.6.0 + */ +export interface SessionMoveOptions { + /** The signal to cancel preparation before committing or submitting. */ + readonly signal?: AbortSignal; +} + /** * Options for publishing a message. * @typeParam TContextData The type of the context data. diff --git a/packages/botkit/src/successor.ts b/packages/botkit/src/successor.ts new file mode 100644 index 0000000..72d2b4a --- /dev/null +++ b/packages/botkit/src/successor.ts @@ -0,0 +1,33 @@ +// 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 { Repository } from "./repository.ts"; + +/** + * Checks the account migration storage contract before wrapping a repository. + * @param repository The repository to check. + * @throws {TypeError} If either required successor method is missing. + * @internal + */ +export function assertSuccessorRepository(repository: Repository): void { + if ( + typeof repository.getSuccessor !== "function" || + typeof repository.setSuccessor !== "function" + ) { + throw new TypeError( + "Repository must implement getSuccessor() and setSuccessor().", + ); + } +} diff --git a/packages/botkit/src/text.test.ts b/packages/botkit/src/text.test.ts index 3b72732..298f815 100644 --- a/packages/botkit/src/text.test.ts +++ b/packages/botkit/src/text.test.ts @@ -156,6 +156,12 @@ const bot: BotWithVoidContextData = { publish() { throw new Error("Not implemented"); }, + move() { + throw new Error("Not implemented"); + }, + republishMove() { + throw new Error("Not implemented"); + }, republishProfile() { throw new Error("Not implemented"); }, From 159187d95fdf6f45f92851e04de889d5b3921a6f Mon Sep 17 00:00:00 2001 From: Hong Minhee Date: Sat, 3 Oct 2026 21:43:20 +0900 Subject: [PATCH 2/4] Keep moved pages usable after target lookup fails A dynamic successor dispatcher can fail independently of the retired bot. Fall back to the stored actor URI so its profile still renders and direct follow requests still return the persisted-move 409 response. https://github.com/fedify-dev/botkit/pull/57#discussion_r4173149759 Assisted-by: Codex:gpt-6.1-sol --- packages/botkit/src/pages.test.ts | 30 ++++++++++++++++++++++++++++++ packages/botkit/src/pages.tsx | 3 ++- 2 files changed, 32 insertions(+), 1 deletion(-) diff --git a/packages/botkit/src/pages.test.ts b/packages/botkit/src/pages.test.ts index 0a60568..802df0f 100644 --- a/packages/botkit/src/pages.test.ts +++ b/packages/botkit/src/pages.test.ts @@ -474,3 +474,33 @@ test("moved notices link local successors to their public profile", async () => assert.ok((await destination.text()).includes("@new-name")); } }); + +test("moved pages survive a failing dynamic successor dispatcher", async () => { + const repository = new MemoryRepository(); + const instance = new InstanceImpl({ + kv: new MemoryKvStore(), + repository, + }); + instance.createBot("old", { username: "old" }); + instance.createBot(async () => { + await Promise.resolve(); + throw new TypeError("The target profile is unavailable."); + }); + const successor = new URL("https://example.com/ap/actor/target"); + await repository.setSuccessor("old", successor); + for ( + const [path, method, status] of [ + ["/@old", "GET", 200], + ["/@old/follow", "POST", 409], + ] as const + ) { + const response = await instance.fetch( + new Request(`https://example.com${path}`, { method }), + undefined, + ); + assert.strictEqual(response.status, status); + const html = await response.text(); + assert.ok(html.includes("This bot has moved")); + assert.ok(html.includes(`href="${successor.href}"`)); + } +}); diff --git a/packages/botkit/src/pages.tsx b/packages/botkit/src/pages.tsx index 5805de6..f289c60 100644 --- a/packages/botkit/src/pages.tsx +++ b/packages/botkit/src/pages.tsx @@ -100,7 +100,8 @@ async function successorWebUrl( signal?.throwIfAborted(); const parsed = ctx.parseUri(successor); if (parsed?.type !== "actor") return successor; - const target = await bot.instance.resolveBot(ctx, parsed.identifier); + const target = await bot.instance.resolveBot(ctx, parsed.identifier) + .catch(() => null); signal?.throwIfAborted(); return target == null ? successor From 45ad85c99641a6e8925a374bbedff0b3d43a2226 Mon Sep 17 00:00:00 2001 From: Hong Minhee Date: Sat, 3 Oct 2026 21:44:43 +0900 Subject: [PATCH 3/4] Finish sharing before committing a bot move An initial active-state check lets a move commit while a share is still being stored or submitted. Serialize successor writes with the entire share operation, including both submissions, for each bot on an instance. Keep target validation and migration notifications outside this boundary. Release the boundary on failures and check cancellation before queued successor writes. Document that separate instances/processes still need application-level coordination. https://github.com/fedify-dev/botkit/pull/57#discussion_r4173149751 Assisted-by: Codex:gpt-6.1-sol --- CHANGES.md | 8 +-- changes.d/botkit/account-move.md | 8 +-- docs/concepts/session.md | 5 ++ packages/botkit/src/instance-impl.ts | 33 ++++++++++++ packages/botkit/src/message-impl.ts | 12 ++++- packages/botkit/src/move.test.ts | 75 ++++++++++++++++++++++++++++ packages/botkit/src/session-impl.ts | 10 ++-- 7 files changed, 141 insertions(+), 10 deletions(-) diff --git a/CHANGES.md b/CHANGES.md index f840417..be7166b 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -11,9 +11,11 @@ To be released. - Added `Session.move()` to move a bot's followers to a linked actor, persist its `movedTo` redirect, and show its new home on the profile. Moved bots reject new follows and cannot publish, reply, or share new - messages. `Session.republishMove()` resends failed migration notifications. - Custom repositories must now implement the required `getSuccessor()` and - atomic, write-once `setSuccessor()` methods. [[#50], [#57]] + messages. In-progress shares finish before a move commits within the + same instance. `Session.republishMove()` resends failed migration + notifications. Custom repositories must now implement the required + `getSuccessor()` and atomic, write-once `setSuccessor()` methods. + [[#50], [#57]] - Added an `aliases` option to `CreateBotOptions` and `BotProfile` so existing accounts can move their followers to a BotKit bot. Actor URIs listed in the option are published as `alsoKnownAs`, and are available diff --git a/changes.d/botkit/account-move.md b/changes.d/botkit/account-move.md index 81465cf..3832c08 100644 --- a/changes.d/botkit/account-move.md +++ b/changes.d/botkit/account-move.md @@ -6,6 +6,8 @@ links: - Added `Session.move()` to move a bot's followers to a linked actor, persist its `movedTo` redirect, and show its new home on the profile. Moved bots reject new follows and cannot publish, reply, or share new - messages. `Session.republishMove()` resends failed migration notifications. - Custom repositories must now implement the required `getSuccessor()` and - atomic, write-once `setSuccessor()` methods. [[#50], [#57]] + messages. In-progress shares finish before a move commits within the + same instance. `Session.republishMove()` resends failed migration + notifications. Custom repositories must now implement the required + `getSuccessor()` and atomic, write-once `setSuccessor()` methods. + [[#50], [#57]] diff --git a/docs/concepts/session.md b/docs/concepts/session.md index 676d7af..9c391e7 100644 --- a/docs/concepts/session.md +++ b/docs/concepts/session.md @@ -249,6 +249,11 @@ follow requests retained before the move also cannot be accepted afterwards, but can still be rejected. There is no API to undo a move or change its stored destination. +Within one instance, an in-progress share finishes storing and submitting its +notifications before `move()` can commit. Shares starting after that commit +are rejected. Applications serving the same bot from several processes must +coordinate sharing and migration between those processes. + [FEP-7628]: https://w3id.org/fep/7628 ### Recovering notification failures diff --git a/packages/botkit/src/instance-impl.ts b/packages/botkit/src/instance-impl.ts index 9126909..2a8297a 100644 --- a/packages/botkit/src/instance-impl.ts +++ b/packages/botkit/src/instance-impl.ts @@ -674,6 +674,39 @@ export class InstanceImpl return []; } + readonly #sharing = new Map>(); + + /** + * Serializes sharing and successor writes for one bot on this instance. + * @param identifier The bot whose state is protected. + * @param operation The operation to run after preceding work finishes. + * @param signal The signal for cancelling before the operation starts. + * @returns The operation's result. + * @throws If the operation fails or its signal is aborted. + * @internal + */ + async withSharingLock( + identifier: string, + operation: (signal?: AbortSignal) => Promise, + signal?: AbortSignal, + ): Promise { + signal?.throwIfAborted(); + const previous = this.#sharing.get(identifier) ?? Promise.resolve(); + const run = previous.then(() => { + signal?.throwIfAborted(); + return operation(signal); + }); + const tail = run.then(() => {}, () => {}); + this.#sharing.set(identifier, tail); + try { + return await run; + } finally { + if (this.#sharing.get(identifier) === tail) { + this.#sharing.delete(identifier); + } + } + } + // Serializes copies delivered through shared and personal inboxes. This // is process-local; repository relationships gate subsequent deliveries. readonly #moves = new Map>(); diff --git a/packages/botkit/src/message-impl.ts b/packages/botkit/src/message-impl.ts index 97489da..4c33085 100644 --- a/packages/botkit/src/message-impl.ts +++ b/packages/botkit/src/message-impl.ts @@ -163,7 +163,17 @@ export class MessageImpl async share( options: MessageShareOptions = {}, ): Promise> { - await this.session.ensureActive(); + return await this.session.bot.instance.withSharingLock( + this.session.bot.identifier, + (signal) => this.#share(options, signal), + ); + } + + async #share( + options: MessageShareOptions, + signal?: AbortSignal, + ): Promise> { + await this.session.ensureActive(signal); const published = new Date(); const id = uuidv7({ msecs: +published }) as Uuid; const visibility = options.visibility ?? this.visibility; diff --git a/packages/botkit/src/move.test.ts b/packages/botkit/src/move.test.ts index 5618d68..ac323e2 100644 --- a/packages/botkit/src/move.test.ts +++ b/packages/botkit/src/move.test.ts @@ -22,6 +22,7 @@ import { import type { DocumentLoader } from "@fedify/vocab-runtime"; import { type Activity, + Announce, Article, Follow, isActor, @@ -533,3 +534,77 @@ test("dynamic targets resolve without HTTP and remain separate from the source", assert.deepStrictEqual((await session.getActor()).successorId, targetId); assert.strictEqual(await repository.getSuccessor("new"), undefined); }); + +test("move waits for sharing storage and both submissions across sessions", async (t) => { + for (const pause of ["storage", "followers", "author"] as const) { + await t.test(pause, async () => { + const f = fixture(); + const message = await f.session.publish(text`Original.`); + await seedFollower(f.repository, t.signal); + const started = Promise.withResolvers(); + const release = Promise.withResolvers(); + const fetched = Promise.withResolvers(); + const getLoader = f.context.getDocumentLoader; + const loader = await getLoader(f.bot); + f.context.getDocumentLoader = mockLoader(async (url, options) => { + const result = await loader(url, options); + fetched.resolve(); + return result; + }); + const add = f.repository.addMessage.bind(f.repository); + f.repository.addMessage = async (identifier, id, activity) => { + if (pause === "storage" && activity instanceof Announce) { + started.resolve(); + await release.promise; + } + return await add(identifier, id, activity); + }; + let submissions = 0; + f.context.sendActivity = async (_sender, _recipients, activity) => { + if (!(activity instanceof Announce)) return; + submissions++; + if ( + (pause === "followers" && submissions === 1) || + (pause === "author" && submissions === 2) + ) { + started.resolve(); + await release.promise; + } + assert.strictEqual(await f.repository.getSuccessor("bot"), undefined); + }; + const sharing = message.share(); + await started.promise; + const moving = new SessionImpl(f.bot, f.context).move(f.target); + await fetched.promise; + // Let the authoritative actor parse finish while the share is suspended. + await new Promise((resolve) => setTimeout(resolve, 20)); + try { + assert.strictEqual(await f.repository.getSuccessor("bot"), undefined); + } finally { + release.resolve(); + await Promise.allSettled([sharing, moving]); + } + await sharing; + await moving; + assert.strictEqual(submissions, 2); + assert.deepStrictEqual( + await f.repository.getSuccessor("bot"), + f.target.id, + ); + await assert.rejects(message.share(), TypeError); + assert.strictEqual(await f.repository.countMessages("bot"), 2); + }); + } +}); + +test("failed sharing releases the move serialization boundary", async () => { + const f = fixture(); + const message = await f.session.publish(text`Original.`); + f.context.sendActivity = async (_sender, _recipients, activity) => { + await Promise.resolve(); + if (activity instanceof Announce) throw new TypeError("Submission failed."); + }; + await assert.rejects(message.share(), TypeError); + await f.session.move(f.target); + assert.deepStrictEqual(await f.repository.getSuccessor("bot"), f.target.id); +}); diff --git a/packages/botkit/src/session-impl.ts b/packages/botkit/src/session-impl.ts index 1a2d725..740f774 100644 --- a/packages/botkit/src/session-impl.ts +++ b/packages/botkit/src/session-impl.ts @@ -241,9 +241,13 @@ export class SessionImpl implements Session { const actor = await this.#resolveMoveTarget(target, false, signal); const successorId = actor.id!; signal?.throwIfAborted(); - if (!await this.bot.repository.setSuccessor(successorId, signal)) { - throw new TypeError("The bot has already moved."); - } + const committed = await this.bot.instance.withSharingLock( + this.bot.identifier, + async (signal) => + await this.bot.repository.setSuccessor(successorId, signal), + signal, + ); + if (!committed) throw new TypeError("The bot has already moved."); // From here cancellation must not interrupt the notification pair. // Never roll back a move which another server may already have processed. try { From 12b92307791d5f775318e5bde5cc8e9b9a332a43 Mon Sep 17 00:00:00 2001 From: Hong Minhee Date: Sat, 3 Oct 2026 22:57:32 +0900 Subject: [PATCH 4/4] Keep actor writes from racing with bot moves A final successor check alone can still let a move commit while a post is being stored or submitted. Hold the existing per-bot instance lock through the final active check, persistence, and every publication submission, while keeping text rendering outside the lock. Use the same boundary for follow acceptance so the move's audience snapshot includes accepted followers. Check both pending and moved state inside the lock, including requests queued behind a move. Document these guarantees and the separate-process coordination limit. https://github.com/fedify-dev/botkit/pull/57#discussion_r4173246347 https://github.com/fedify-dev/botkit/pull/57#pullrequestreview-5400870019 Assisted-by: Codex:gpt-6.1-sol --- CHANGES.md | 10 +- changes.d/botkit/account-move.md | 10 +- docs/concepts/session.md | 13 +- packages/botkit/src/follow-impl.ts | 9 +- packages/botkit/src/instance-impl.ts | 2 +- packages/botkit/src/move.test.ts | 218 ++++++++++++++++++++++++++- packages/botkit/src/session-impl.ts | 170 +++++++++++---------- 7 files changed, 330 insertions(+), 102 deletions(-) diff --git a/CHANGES.md b/CHANGES.md index be7166b..f060970 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -11,11 +11,11 @@ To be released. - Added `Session.move()` to move a bot's followers to a linked actor, persist its `movedTo` redirect, and show its new home on the profile. Moved bots reject new follows and cannot publish, reply, or share new - messages. In-progress shares finish before a move commits within the - same instance. `Session.republishMove()` resends failed migration - notifications. Custom repositories must now implement the required - `getSuccessor()` and atomic, write-once `setSuccessor()` methods. - [[#50], [#57]] + messages. Publishing, sharing, and follow acceptance are serialized + with moves within the same instance. `Session.republishMove()` + resends failed migration notifications. Custom repositories must now + implement the required `getSuccessor()` and atomic, write-once + `setSuccessor()` methods. [[#50], [#57]] - Added an `aliases` option to `CreateBotOptions` and `BotProfile` so existing accounts can move their followers to a BotKit bot. Actor URIs listed in the option are published as `alsoKnownAs`, and are available diff --git a/changes.d/botkit/account-move.md b/changes.d/botkit/account-move.md index 3832c08..1da1e9a 100644 --- a/changes.d/botkit/account-move.md +++ b/changes.d/botkit/account-move.md @@ -6,8 +6,8 @@ links: - Added `Session.move()` to move a bot's followers to a linked actor, persist its `movedTo` redirect, and show its new home on the profile. Moved bots reject new follows and cannot publish, reply, or share new - messages. In-progress shares finish before a move commits within the - same instance. `Session.republishMove()` resends failed migration - notifications. Custom repositories must now implement the required - `getSuccessor()` and atomic, write-once `setSuccessor()` methods. - [[#50], [#57]] + messages. Publishing, sharing, and follow acceptance are serialized + with moves within the same instance. `Session.republishMove()` + resends failed migration notifications. Custom repositories must now + implement the required `getSuccessor()` and atomic, write-once + `setSuccessor()` methods. [[#50], [#57]] diff --git a/docs/concepts/session.md b/docs/concepts/session.md index 9c391e7..7078bfd 100644 --- a/docs/concepts/session.md +++ b/docs/concepts/session.md @@ -249,10 +249,12 @@ follow requests retained before the move also cannot be accepted afterwards, but can still be rejected. There is no API to undo a move or change its stored destination. -Within one instance, an in-progress share finishes storing and submitting its -notifications before `move()` can commit. Shares starting after that commit -are rejected. Applications serving the same bot from several processes must -coordinate sharing and migration between those processes. +Within one instance, once publishing, sharing, or follow acceptance passes its +final state check, `move()` waits for its storage and activity submissions to +finish. Text rendering happens before that check, so a move during rendering +rejects publication before it stores anything. Applications serving the same +bot from several processes must coordinate these operations and migration +between those processes. [FEP-7628]: https://w3id.org/fep/7628 @@ -278,8 +280,7 @@ if ((await session.getActor()).successorId != null) { } ~~~~ -This also reaches followers whose acceptance was already in progress when the -bot moved. It does not change the successor. The destination must still list +It does not change the successor. The destination must still list the old actor as an alias, even if it has since moved again. With a configured queue, Fedify retries delivery of notifications it has accepted; without a queue, delivery happens during the call and can partially fail. diff --git a/packages/botkit/src/follow-impl.ts b/packages/botkit/src/follow-impl.ts index 44b348a..dff7296 100644 --- a/packages/botkit/src/follow-impl.ts +++ b/packages/botkit/src/follow-impl.ts @@ -47,10 +47,17 @@ export class FollowRequestImpl implements FollowRequest { } async accept(): Promise { + await this.session.bot.instance.withSharingLock( + this.session.bot.identifier, + (signal) => this.#accept(signal), + ); + } + + async #accept(signal?: AbortSignal): Promise { if (this.#state !== "pending") { throw new TypeError("The follow request is not pending."); } - if (await this.session.bot.repository.getSuccessor() != null) { + if (await this.session.bot.repository.getSuccessor(signal) != null) { throw new TypeError( "The bot has moved and cannot accept follow requests.", ); diff --git a/packages/botkit/src/instance-impl.ts b/packages/botkit/src/instance-impl.ts index 2a8297a..0ccd474 100644 --- a/packages/botkit/src/instance-impl.ts +++ b/packages/botkit/src/instance-impl.ts @@ -677,7 +677,7 @@ export class InstanceImpl readonly #sharing = new Map>(); /** - * Serializes sharing and successor writes for one bot on this instance. + * Serializes actor writes with migration for one bot on this instance. * @param identifier The bot whose state is protected. * @param operation The operation to run after preceding work finishes. * @param signal The signal for cancelling before the operation starts. diff --git a/packages/botkit/src/move.test.ts b/packages/botkit/src/move.test.ts index ac323e2..e37ffef 100644 --- a/packages/botkit/src/move.test.ts +++ b/packages/botkit/src/move.test.ts @@ -21,15 +21,18 @@ import { } from "@fedify/fedify/federation"; import type { DocumentLoader } from "@fedify/vocab-runtime"; import { + Accept, type Activity, Announce, Article, + Create, Follow, isActor, Move, Note, Person, PUBLIC_COLLECTION, + QuoteRequest, Reject, Update, } from "@fedify/vocab"; @@ -41,7 +44,8 @@ import { hideRepositoryMethods } from "./helpers.ts"; import { createInstance } from "./instance.ts"; import { MemoryCachedRepository, MemoryRepository } from "./repository.ts"; import { SessionImpl } from "./session-impl.ts"; -import { text } from "./text.ts"; +import { mention, text } from "./text.ts"; +import { createMessage } from "./message-impl.ts"; function mockLoader( loader: DocumentLoader, @@ -608,3 +612,215 @@ test("failed sharing releases the move serialization boundary", async () => { await f.session.move(f.target); assert.deepStrictEqual(await f.repository.getSuccessor("bot"), f.target.id); }); + +test("move waits for publishing storage and every activity submission", async (t) => { + for (const pause of [0, 1, 2, 3, 4, 5]) { + await t.test(pause === 0 ? "storage" : `submission ${pause}`, async () => { + const f = fixture(); + const remote = new Person({ + id: new URL("https://remote.example/actor"), + inbox: new URL("https://remote.example/inbox"), + }); + const original = await createMessage( + new Note({ + id: new URL("https://remote.example/note/1"), + attribution: remote.id, + content: "Original.", + to: PUBLIC_COLLECTION, + }), + f.session, + { [remote.id!.href]: remote }, + ); + const started = Promise.withResolvers(); + const release = Promise.withResolvers(); + const fetched = Promise.withResolvers(); + const loader = await f.context.getDocumentLoader(f.bot); + f.context.getDocumentLoader = mockLoader(async (url, options) => { + const result = await loader(url, options); + fetched.resolve(); + return result; + }); + const add = f.repository.addMessage.bind(f.repository); + f.repository.addMessage = async (identifier, id, activity) => { + if (pause === 0) { + started.resolve(); + await release.promise; + } + return await add(identifier, id, activity); + }; + let submissions = 0; + f.context.sendActivity = async (_sender, _recipients, activity) => { + if (!(activity instanceof Create || activity instanceof QuoteRequest)) { + return; + } + submissions++; + if (pause === submissions) { + started.resolve(); + await release.promise; + } + assert.strictEqual(await f.repository.getSuccessor("bot"), undefined); + }; + const publishing = f.session.publish( + text`Hello ${mention("remote", remote)}.`, + { + replyTarget: original, + quoteTarget: original, + }, + ); + await started.promise; + const moving = new SessionImpl(f.bot, f.context).move(f.target); + await fetched.promise; + await new Promise((resolve) => setTimeout(resolve, 20)); + try { + assert.strictEqual(await f.repository.getSuccessor("bot"), undefined); + } finally { + release.resolve(); + await Promise.allSettled([publishing, moving]); + } + await publishing; + await moving; + assert.strictEqual(submissions, 5); + assert.strictEqual(await f.repository.countMessages("bot"), 1); + }); + } +}); + +test("move snapshots followers after pending acceptance finishes", async (t) => { + for (const pause of ["delivery", "storage"] as const) { + await t.test(pause, async () => { + const f = fixture(); + const follower = new Person({ + id: new URL("https://follower.example/actor"), + inbox: new URL("https://follower.example/inbox"), + }); + const request = new FollowRequestImpl( + f.session, + new Follow({ + id: new URL("https://follower.example/follow/1"), + actor: follower, + object: f.session.actorId, + }), + follower, + ); + const started = Promise.withResolvers(); + const release = Promise.withResolvers(); + const fetched = Promise.withResolvers(); + const loader = await f.context.getDocumentLoader(f.bot); + f.context.getDocumentLoader = mockLoader(async (url, options) => { + const result = await loader(url, options); + fetched.resolve(); + return result; + }); + const audiences: string[][] = []; + f.context.sendActivity = async (_sender, recipients, activity) => { + if (activity instanceof Accept && pause === "delivery") { + started.resolve(); + await release.promise; + } + if (activity instanceof Update || activity instanceof Move) { + assert.ok(Array.isArray(recipients)); + audiences.push(recipients.map((actor) => actor.id!.href)); + } + }; + const add = f.repository.addFollower.bind(f.repository); + f.repository.addFollower = async (identifier, id, actor) => { + if (pause === "storage") { + started.resolve(); + await release.promise; + } + return await add(identifier, id, actor); + }; + const accepting = request.accept(); + await started.promise; + const moving = new SessionImpl(f.bot, f.context).move(f.target); + await fetched.promise; + await new Promise((resolve) => setTimeout(resolve, 20)); + try { + assert.strictEqual(await f.repository.getSuccessor("bot"), undefined); + } finally { + release.resolve(); + await Promise.allSettled([accepting, moving]); + } + await accepting; + await moving; + assert.strictEqual(request.state, "accepted"); + assert.deepStrictEqual(audiences, [[follower.id!.href], [ + follower.id!.href, + ]]); + }); + } +}); + +test("queued follow acceptance rechecks moved state inside the lock", async () => { + const f = fixture(); + const follower = new Person({ + id: new URL("https://follower.example/actor"), + inbox: new URL("https://follower.example/inbox"), + }); + const request = new FollowRequestImpl( + f.session, + new Follow({ + id: new URL("https://follower.example/follow/1"), + actor: follower, + object: f.session.actorId, + }), + follower, + ); + const started = Promise.withResolvers(); + const release = Promise.withResolvers(); + const set = f.repository.setSuccessor.bind(f.repository); + f.repository.setSuccessor = async (identifier, successor, signal) => { + started.resolve(); + await release.promise; + return await set(identifier, successor, signal); + }; + const moving = f.session.move(f.target); + await started.promise; + const accepting = request.accept(); + await new Promise((resolve) => setTimeout(resolve, 20)); + release.resolve(); + const results = await Promise.allSettled([moving, accepting]); + assert.strictEqual(results[0].status, "fulfilled"); + assert.strictEqual(results[1].status, "rejected"); + if (results[1].status === "rejected") { + assert.ok(results[1].reason instanceof TypeError); + } + assert.strictEqual(request.state, "pending"); + assert.strictEqual(await f.repository.countFollowers("bot"), 0); + assert.ok(!f.sent.some(({ activity }) => activity instanceof Accept)); +}); + +test("concurrent acceptance checks pending state after preceding acceptance", async () => { + const f = fixture(); + const follower = new Person({ + id: new URL("https://follower.example/actor"), + inbox: new URL("https://follower.example/inbox"), + }); + const request = new FollowRequestImpl( + f.session, + new Follow({ + id: new URL("https://follower.example/follow/1"), + actor: follower, + object: f.session.actorId, + }), + follower, + ); + const results = await Promise.allSettled([ + request.accept(), + request.accept(), + ]); + assert.strictEqual( + results.filter((result) => result.status === "fulfilled").length, + 1, + ); + assert.strictEqual( + results.filter((result) => result.status === "rejected").length, + 1, + ); + assert.strictEqual(request.state, "accepted"); + assert.strictEqual( + f.sent.filter(({ activity }) => activity instanceof Accept).length, + 1, + ); + assert.strictEqual(await f.repository.countFollowers("bot"), 1); +}); diff --git a/packages/botkit/src/session-impl.ts b/packages/botkit/src/session-impl.ts index 740f774..9fabb64 100644 --- a/packages/botkit/src/session-impl.ts +++ b/packages/botkit/src/session-impl.ts @@ -614,90 +614,94 @@ export class SessionImpl implements Session { object: msg, published: published.toTemporalInstant(), }); - // Rendering text can await remote lookups or user code. Recheck before - // storing anything in case the bot moved during that preparation. - await this.ensureActive(); - await this.bot.repository.addMessage(id, activity); - const preferSharedInbox = visibility === "public" || - visibility === "unlisted" || visibility === "followers"; - const excludeBaseUris = [new URL(this.context.origin)]; - if (preferSharedInbox) { - await this.context.sendActivity( - this.bot, - "followers", - activity, - { preferSharedInbox, excludeBaseUris }, - ); - } const cachedObjects: Record = {}; - const textObjects = [ - ...content.getCachedObjects(), - ...(summary?.getCachedObjects() ?? []), - ]; - for (const cachedObject of textObjects) { - if (cachedObject.id == null) continue; - cachedObjects[cachedObject.id.href] = cachedObject; - } - if (mentionedActorIds.length > 0) { - const documentLoader = await this.context.getDocumentLoader(this.bot); - const promises: Promise[] = []; - for (const mentionedActorId of mentionedActorIds) { - const cachedObject = cachedObjects[mentionedActorId.href]; - const promise = cachedObject == null - ? this.context.lookupObject( - mentionedActorId, - { documentLoader }, - ) - : Promise.resolve(cachedObject); - promises.push(promise); - } - const objects = await Promise.all(promises); - const mentionedActors = objects.filter(isActor); - await this.context.sendActivity( - this.bot, - mentionedActors, - activity, - { preferSharedInbox, excludeBaseUris }, - ); - } - if (options.replyTarget != null) { - await this.context.sendActivity( - this.bot, - options.replyTarget.actor, - activity, - { preferSharedInbox, excludeBaseUris, fanout: "skip" }, - ); - } - if (options.quoteTarget != null) { - await this.context.sendActivity( - this.bot, - options.quoteTarget.actor, - activity, - { preferSharedInbox, excludeBaseUris, fanout: "skip" }, - ); - if ( - options.quoteTarget.actor.id != null && - options.quoteTarget.actor.id.href !== - this.context.getActorUri(this.bot.identifier).href - ) { - const request = quoteInteraction.createRequest({ - id: this.context.getObjectUri(QuoteRequest, { - identifier: this.bot.identifier, - id, - }), - actor: this.context.getActorUri(this.bot.identifier), - object: options.quoteTarget.id, - instrument: msgId, - to: options.quoteTarget.actor.id, - }); - await this.context.sendActivity( - this.bot, - options.quoteTarget.actor, - request, - { preferSharedInbox, excludeBaseUris, fanout: "skip" }, - ); - } - } + await this.bot.instance.withSharingLock( + this.bot.identifier, + async (signal) => { + // Recheck after rendering, under the same lock as successor commits. + await this.ensureActive(signal); + await this.bot.repository.addMessage(id, activity); + const preferSharedInbox = visibility === "public" || + visibility === "unlisted" || visibility === "followers"; + const excludeBaseUris = [new URL(this.context.origin)]; + if (preferSharedInbox) { + await this.context.sendActivity( + this.bot, + "followers", + activity, + { preferSharedInbox, excludeBaseUris }, + ); + } + const textObjects = [ + ...content.getCachedObjects(), + ...(summary?.getCachedObjects() ?? []), + ]; + for (const cachedObject of textObjects) { + if (cachedObject.id == null) continue; + cachedObjects[cachedObject.id.href] = cachedObject; + } + if (mentionedActorIds.length > 0) { + const documentLoader = await this.context.getDocumentLoader(this.bot); + const promises: Promise[] = []; + for (const mentionedActorId of mentionedActorIds) { + const cachedObject = cachedObjects[mentionedActorId.href]; + const promise = cachedObject == null + ? this.context.lookupObject( + mentionedActorId, + { documentLoader }, + ) + : Promise.resolve(cachedObject); + promises.push(promise); + } + const objects = await Promise.all(promises); + const mentionedActors = objects.filter(isActor); + await this.context.sendActivity( + this.bot, + mentionedActors, + activity, + { preferSharedInbox, excludeBaseUris }, + ); + } + if (options.replyTarget != null) { + await this.context.sendActivity( + this.bot, + options.replyTarget.actor, + activity, + { preferSharedInbox, excludeBaseUris, fanout: "skip" }, + ); + } + if (options.quoteTarget != null) { + await this.context.sendActivity( + this.bot, + options.quoteTarget.actor, + activity, + { preferSharedInbox, excludeBaseUris, fanout: "skip" }, + ); + if ( + options.quoteTarget.actor.id != null && + options.quoteTarget.actor.id.href !== + this.context.getActorUri(this.bot.identifier).href + ) { + const request = quoteInteraction.createRequest({ + id: this.context.getObjectUri(QuoteRequest, { + identifier: this.bot.identifier, + id, + }), + actor: this.context.getActorUri(this.bot.identifier), + object: options.quoteTarget.id, + instrument: msgId, + to: options.quoteTarget.actor.id, + }); + await this.context.sendActivity( + this.bot, + options.quoteTarget.actor, + request, + { preferSharedInbox, excludeBaseUris, fanout: "skip" }, + ); + } + } + }, + ); return await createMessage( msg, this,