diff --git a/CHANGES.md b/CHANGES.md index ab79070..f060970 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -8,6 +8,14 @@ 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. 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 @@ -45,10 +53,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..1da1e9a --- /dev/null +++ b/changes.d/botkit/account-move.md @@ -0,0 +1,13 @@ +--- +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. 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/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..7078bfd 100644 --- a/docs/concepts/session.md +++ b/docs/concepts/session.md @@ -193,6 +193,104 @@ 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. + +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 + +### 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(); +} +~~~~ + +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..dff7296 100644 --- a/packages/botkit/src/follow-impl.ts +++ b/packages/botkit/src/follow-impl.ts @@ -47,9 +47,21 @@ 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(signal) != 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..0ccd474 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 = { @@ -672,6 +674,39 @@ export class InstanceImpl return []; } + readonly #sharing = new Map>(); + + /** + * 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. + * @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 3586b50..4c33085 100644 --- a/packages/botkit/src/message-impl.ts +++ b/packages/botkit/src/message-impl.ts @@ -163,6 +163,17 @@ export class MessageImpl async share( options: MessageShareOptions = {}, ): Promise> { + 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/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..e37ffef --- /dev/null +++ b/packages/botkit/src/move.test.ts @@ -0,0 +1,826 @@ +// 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 { + Accept, + type Activity, + Announce, + Article, + Create, + Follow, + isActor, + Move, + Note, + Person, + PUBLIC_COLLECTION, + QuoteRequest, + 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 { mention, text } from "./text.ts"; +import { createMessage } from "./message-impl.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); +}); + +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); +}); + +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/pages.test.ts b/packages/botkit/src/pages.test.ts index fdb871f..802df0f 100644 --- a/packages/botkit/src/pages.test.ts +++ b/packages/botkit/src/pages.test.ts @@ -390,3 +390,117 @@ 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* { + 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 9cffef4..f289c60 100644 --- a/packages/botkit/src/pages.tsx +++ b/packages/botkit/src/pages.tsx @@ -90,6 +90,24 @@ export const app = new Hono(); 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) + .catch(() => null); + signal?.throwIfAborted(); + return target == null + ? successor + : bot.instance.getBotWebUrl(target, ctx.origin); +} + async function profilePage( c: PageContext, bot: BotImpl, @@ -111,6 +129,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 +214,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 +600,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..9fabb64 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,198 @@ 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(); + 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 { + 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 +459,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,87 +614,94 @@ export class SessionImpl implements Session { object: msg, published: published.toTemporalInstant(), }); - 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, 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"); },