diff --git a/docs/channels/discord.md b/docs/channels/discord.md index ae1cb5c098e1..ab12a63bdcbe 100644 --- a/docs/channels/discord.md +++ b/docs/channels/discord.md @@ -815,25 +815,6 @@ See [Slash commands](/tools/slash-commands) for the command catalog and behavior - - Set `channels.discord.subagentProgress: true` to show background child activity on the Discord message that started the parent run. - -```json5 -{ - channels: { - discord: { - subagentProgress: true, - }, - }, -} -``` - - While child runs are active, OpenClaw keeps Discord typing active for up to one hour and replaces one count reaction (`1️⃣` through `🔟`) as the concurrent count changes; `🔟` also represents 10 or more. The count reaction is removed after the final child ends. A failed, timed-out, or killed child leaves a `🔴` reaction. - - This is opt-in and uses fixed internal timing and emoji defaults. The bot needs **Add Reactions** permission for reaction feedback. Account-level `channels.discord.accounts..subagentProgress` overrides the top-level value. - - - For stable "always-on" ACP workspaces, configure top-level typed ACP bindings targeting Discord conversations. diff --git a/extensions/discord/src/subagent-hooks.test.ts b/extensions/discord/src/subagent-hooks.test.ts index 68c2c942d969..1bea26f42d7b 100644 --- a/extensions/discord/src/subagent-hooks.test.ts +++ b/extensions/discord/src/subagent-hooks.test.ts @@ -61,7 +61,6 @@ const hookMocks = vi.hoisted(() => { listThreadBindingsBySessionKey: vi.fn((_params?: unknown): ThreadBindingRecord[] => []), unbindThreadBindingsBySessionKey: vi.fn(() => []), progressModuleFactory: vi.fn(), - handleDiscordSubagentProgress: vi.fn(), recoverDiscordSubagentProgress: vi.fn(), }; }); @@ -78,10 +77,7 @@ vi.mock("./monitor/thread-bindings.js", () => ({ })); vi.mock("./subagent-progress.js", () => { hookMocks.progressModuleFactory(); - return { - handleDiscordSubagentProgress: hookMocks.handleDiscordSubagentProgress, - recoverDiscordSubagentProgress: hookMocks.recoverDiscordSubagentProgress, - }; + return { recoverDiscordSubagentProgress: hookMocks.recoverDiscordSubagentProgress }; }); function registerHandlersForTest( @@ -221,11 +217,10 @@ describe("discord subagent hook handlers", () => { hookMocks.listThreadBindingsBySessionKey.mockClear(); hookMocks.unbindThreadBindingsBySessionKey.mockClear(); hookMocks.progressModuleFactory.mockClear(); - hookMocks.handleDiscordSubagentProgress.mockClear(); hookMocks.recoverDiscordSubagentProgress.mockClear(); }); - it("keeps progress runtime lazy for unrelated subagent hooks", async () => { + it("keeps progress cleanup lazy for unrelated subagent hooks", async () => { const handlers = registerHandlersForTest(); const handler = getRequiredHookHandler(handlers, "subagent_delivery_target"); @@ -244,7 +239,7 @@ describe("discord subagent hook handlers", () => { expect(hookMocks.progressModuleFactory).not.toHaveBeenCalled(); }); - it("loads progress runtime for gateway recovery", async () => { + it("loads retired progress cleanup on gateway startup", async () => { const handlers = registerHandlersForTest(); const handler = getRequiredHookHandler(handlers, "gateway_start"); diff --git a/extensions/discord/src/subagent-progress-config.ts b/extensions/discord/src/subagent-progress-config.ts deleted file mode 100644 index 9ef55d90d118..000000000000 --- a/extensions/discord/src/subagent-progress-config.ts +++ /dev/null @@ -1,63 +0,0 @@ -import { DEFAULT_EMOJIS } from "openclaw/plugin-sdk/channel-feedback"; -import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; -import { - normalizeOptionalLowercaseString, - normalizeOptionalStringifiedId, -} from "openclaw/plugin-sdk/string-coerce-runtime"; - -export const RUNNING_EMOJIS = ["1️⃣", "2️⃣", "3️⃣", "4️⃣", "5️⃣", "6️⃣", "7️⃣", "8️⃣", "9️⃣", "🔟"]; -export const FAILURE_EMOJI = "🔴"; - -export type DiscordProgressRequester = { - channel?: string; - accountId?: string; - to?: string; - threadId?: string | number; - channelId?: string | number; - messageId?: string | number; -}; - -function channelIdFromTarget(target?: string): string | undefined { - const trimmed = target?.trim(); - if (!trimmed) { - return undefined; - } - if (trimmed.startsWith("channel:")) { - return trimmed.slice("channel:".length).trim() || undefined; - } - return /^\d+$/u.test(trimmed) ? trimmed : undefined; -} - -export function resolveDiscordProgressTarget(requester?: DiscordProgressRequester) { - if (normalizeOptionalLowercaseString(requester?.channel) !== "discord") { - return undefined; - } - const channelId = - normalizeOptionalStringifiedId(requester?.channelId) ?? channelIdFromTarget(requester?.to); - const messageId = normalizeOptionalStringifiedId(requester?.messageId); - if (!channelId || !messageId) { - return undefined; - } - return { channelId, messageId }; -} - -export function reservedReactionEmojis(config: OpenClawConfig, ackReaction?: string): Set { - const reserved = new Set(Object.values(DEFAULT_EMOJIS)); - for (const emoji of [config.messages?.ackReaction, ackReaction]) { - if (emoji?.trim()) { - reserved.add(emoji.trim()); - } - } - for (const agent of config.agents?.list ?? []) { - const emoji = agent.identity?.emoji?.trim(); - if (emoji) { - reserved.add(emoji); - } - } - return reserved; -} - -export function reactionsAreAvailable(config: OpenClawConfig, ackReaction?: string): boolean { - const reserved = reservedReactionEmojis(config, ackReaction); - return !RUNNING_EMOJIS.some((emoji) => reserved.has(emoji)) && !reserved.has(FAILURE_EMOJI); -} diff --git a/extensions/discord/src/subagent-progress-state.ts b/extensions/discord/src/subagent-progress-state.ts deleted file mode 100644 index 8a3c32197505..000000000000 --- a/extensions/discord/src/subagent-progress-state.ts +++ /dev/null @@ -1,264 +0,0 @@ -import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; -import type { PluginStateKeyedStore } from "openclaw/plugin-sdk/plugin-state-runtime"; - -export const PROGRESS_STORE_TTL_MS = 7 * 24 * 60 * 60_000; -export const MAX_TRACKED_RUNS = 4_096; -const TERMINAL_TOMBSTONE_TTL_MS = 60 * 60_000; - -export type SubagentProgressOutcome = "ok" | "error" | "timeout" | "killed" | "unknown"; - -export type ProgressApi = { - config: OpenClawConfig; - logger: { debug?: (message: string) => void }; - runtime?: { - state: { - openKeyedStore(options: { - namespace: string; - maxEntries: number; - overflowPolicy: "reject-new"; - defaultTtlMs: number; - }): PluginStateKeyedStore; - }; - }; -}; - -type PersistedProgressRunBase = { - /** Plugin SQLite ownership lets an ended run clean reactions after gateway restart. */ - key: string; - accountId: string; - channelId: string; - messageId: string; - runningEmoji?: string; -}; - -type PersistedActiveProgressRun = PersistedProgressRunBase & { status: "active" }; -type PersistedCleanupProgressRun = PersistedProgressRunBase & { - status: "cleanup"; - outcome: SubagentProgressOutcome; -}; -export type PersistedProgressRun = PersistedActiveProgressRun | PersistedCleanupProgressRun; - -export function persistedTerminalOutcome( - persisted?: PersistedProgressRun, -): SubagentProgressOutcome | undefined { - return persisted?.status === "cleanup" ? persisted.outcome : undefined; -} - -export type ProgressTracker = { - accountId: string; - channelId: string; - messageId: string; - activeRunIds: Set; - persistedRunIds: Set; - runningEmoji?: string; - runningEmojiConfirmed: boolean; - reactionsEnabled: boolean; - typingTimer?: ReturnType; - typingExpiresAt: number; -}; - -type ProgressRunLookupResult = - | { status: "found"; value: PersistedProgressRun } - | { status: "missing" } - | { status: "error" }; - -export type PersistProgressResult = "persisted" | "terminal" | "conflict" | "error"; - -type ProgressStateForKeyResult = - | { - ok: true; - activeRunIds: string[]; - cleanupRuns: Array<{ runId: string; value: PersistedCleanupProgressRun }>; - ownedEmojis: string[]; - } - | { ok: false }; - -let progressStores = new WeakMap | null>(); -const trackerQueues = new Map>(); -const terminalRuns = new Map(); - -export function markRunTerminal(runId: string, outcome: SubagentProgressOutcome) { - const now = Date.now(); - for (const [trackedRunId, terminal] of terminalRuns) { - if (terminal.expiresAt <= now) { - terminalRuns.delete(trackedRunId); - } - } - terminalRuns.set(runId, { expiresAt: now + TERMINAL_TOMBSTONE_TTL_MS, outcome }); - if (terminalRuns.size > MAX_TRACKED_RUNS) { - const oldest = terminalRuns.keys().next().value; - if (oldest) { - terminalRuns.delete(oldest); - } - } -} - -export function terminalOutcome(runId: string): SubagentProgressOutcome | undefined { - const terminal = terminalRuns.get(runId); - if (!terminal) { - return undefined; - } - if (terminal.expiresAt <= Date.now()) { - terminalRuns.delete(runId); - return undefined; - } - return terminal.outcome; -} - -export function logFailure(api: ProgressApi, action: string, error: unknown) { - const message = error instanceof Error ? error.message : String(error); - api.logger.debug?.(`discord subagent progress ${action} failed: ${message}`); -} - -export function getProgressStore( - api: ProgressApi, -): PluginStateKeyedStore | undefined { - const cached = progressStores.get(api); - if (cached !== undefined) { - return cached ?? undefined; - } - if (!api.runtime?.state) { - progressStores.set(api, null); - return undefined; - } - try { - const store = api.runtime.state.openKeyedStore({ - namespace: "subagent-progress", - maxEntries: MAX_TRACKED_RUNS, - overflowPolicy: "reject-new", - defaultTtlMs: PROGRESS_STORE_TTL_MS, - }); - progressStores.set(api, store); - return store; - } catch (error) { - logFailure(api, "state store open", error); - return undefined; - } -} - -export function persistedProgressRunFromTracker(tracker: ProgressTracker): PersistedProgressRun { - return { - key: `${tracker.accountId}:${tracker.channelId}:${tracker.messageId}`, - accountId: tracker.accountId, - channelId: tracker.channelId, - messageId: tracker.messageId, - status: "active", - ...(tracker.runningEmoji ? { runningEmoji: tracker.runningEmoji } : {}), - }; -} - -export async function persistProgressRun( - api: ProgressApi, - runId: string, - tracker: ProgressTracker, -): Promise { - const store = getProgressStore(api); - if (!store) { - return "error"; - } - const value = persistedProgressRunFromTracker(tracker); - try { - if (await store.registerIfAbsent(runId, value)) { - return "persisted"; - } - const existing = await store.lookup(runId); - if (!existing) { - return "error"; - } - if (existing.status === "cleanup") { - return "terminal"; - } - return existing.key === value.key ? "persisted" : "conflict"; - } catch (error) { - logFailure(api, "state store write", error); - return "error"; - } -} - -export async function markProgressRunForCleanup( - api: ProgressApi, - runId: string, - persisted: PersistedProgressRun, - outcome: SubagentProgressOutcome, -) { - try { - await getProgressStore(api)?.register(runId, { ...persisted, status: "cleanup", outcome }); - return true; - } catch (error) { - logFailure(api, "state store cleanup mark", error); - return false; - } -} - -export async function lookupProgressRun( - api: ProgressApi, - runId: string, -): Promise { - const store = getProgressStore(api); - if (!store) { - return { status: "error" }; - } - try { - const value = await store.lookup(runId); - return value ? { status: "found", value } : { status: "missing" }; - } catch (error) { - logFailure(api, "state store read", error); - return { status: "error" }; - } -} - -export async function consumeProgressRun(api: ProgressApi, runId: string) { - try { - return await getProgressStore(api)?.consume(runId); - } catch (error) { - logFailure(api, "state store consume", error); - return undefined; - } -} - -export async function listProgressStateForKey( - api: ProgressApi, - key: string, -): Promise { - const store = getProgressStore(api); - if (!store) { - return { ok: false }; - } - try { - const entries = await store.entries(); - const matching = entries.filter((entry) => entry.value.key === key); - const cleanupRuns = matching.flatMap((entry) => - entry.value.status === "cleanup" ? [{ runId: entry.key, value: entry.value }] : [], - ); - return { - ok: true, - activeRunIds: matching - .filter((entry) => entry.value.status === "active") - .map((entry) => entry.key), - cleanupRuns, - ownedEmojis: Array.from(new Set(matching.flatMap((entry) => entry.value.runningEmoji ?? []))), - }; - } catch (error) { - logFailure(api, "state store list", error); - return { ok: false }; - } -} - -export async function runQueued(key: string, task: () => Promise) { - const previous = trackerQueues.get(key) ?? Promise.resolve(); - const current = previous.catch(() => undefined).then(task); - trackerQueues.set(key, current); - try { - await current; - } finally { - if (trackerQueues.get(key) === current) { - trackerQueues.delete(key); - } - } -} - -export function resetDiscordSubagentProgressStateForTest() { - trackerQueues.clear(); - terminalRuns.clear(); - progressStores = new WeakMap(); -} diff --git a/extensions/discord/src/subagent-progress.test.ts b/extensions/discord/src/subagent-progress.test.ts index 0df0f9bc08fa..f336c0fdcb33 100644 --- a/extensions/discord/src/subagent-progress.test.ts +++ b/extensions/discord/src/subagent-progress.test.ts @@ -1,959 +1,217 @@ -// Discord tests cover portable subagent progress presentation. +// Discord tests cover cleanup left by the retired subagent progress feature. import type { PluginStateKeyedStore } from "openclaw/plugin-sdk/plugin-state-runtime"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; -import { - handleDiscordSubagentProgress, - recoverDiscordSubagentProgress, -} from "./subagent-progress.js"; +import { recoverDiscordSubagentProgress } from "./subagent-progress.js"; const sendMocks = vi.hoisted(() => ({ - react: vi.fn(async (_channelId: string, _messageId: string, _emoji: string, _opts: unknown) => ({ - ok: true, - })), remove: vi.fn(async (_channelId: string, _messageId: string, _emoji: string, _opts: unknown) => ({ ok: true, })), - typing: vi.fn(async (_channelId: string, _opts: unknown) => ({ - ok: true, - channelId: "123", - })), })); -const resetDiscordSubagentProgressForTest = handleDiscordSubagentProgress.resetForTest; - vi.mock("./send.reactions.js", () => ({ - reactMessageDiscord: sendMocks.react, removeReactionDiscord: sendMocks.remove, })); -vi.mock("./send.typing.js", () => ({ sendTypingDiscord: sendMocks.typing })); -type StoredProgressRunBase = { +type StoredProgressRun = { key: string; accountId: string; channelId: string; messageId: string; runningEmoji?: string; + status: "active" | "cleanup"; + outcome?: "ok" | "error" | "timeout" | "killed" | "unknown"; }; -type StoredProgressRun = StoredProgressRunBase & - ( - | { status: "active" } - | { status: "cleanup"; outcome: "ok" | "error" | "timeout" | "killed" | "unknown" } - ); - function createProgressStore() { const values = new Map(); - return { - values, - store: { - register: vi.fn(async (key: string, value: StoredProgressRun) => { - values.set(key, value); - }), - registerIfAbsent: vi.fn(async (key: string, value: StoredProgressRun) => { - if (values.has(key)) { - return false; - } - values.set(key, value); - return true; - }), - lookup: vi.fn(async (key: string) => values.get(key)), - consume: vi.fn(async (key: string) => { - const value = values.get(key); - values.delete(key); - return value; - }), - delete: vi.fn(async (key: string) => values.delete(key)), - entries: vi.fn(async () => - Array.from(values, ([key, value]) => ({ key, value, createdAt: Date.now() })), - ), - clear: vi.fn(async () => { - values.clear(); - }), - } satisfies PluginStateKeyedStore, - }; + const store = { + register: vi.fn(async (key: string, value: StoredProgressRun) => { + values.set(key, value); + }), + registerIfAbsent: vi.fn(async () => false), + lookup: vi.fn(async (key: string) => values.get(key)), + consume: vi.fn(async (key: string) => { + const value = values.get(key); + values.delete(key); + return value; + }), + delete: vi.fn(async (key: string) => values.delete(key)), + entries: vi.fn(async () => + Array.from(values, ([key, value]) => ({ key, value, createdAt: Date.now() })), + ), + clear: vi.fn(async () => values.clear()), + } satisfies PluginStateKeyedStore; + return { values, store }; } -function createApi(config?: Record) { +function createApi(discord: Record = { token: "test-token" }) { const progressStore = createProgressStore(); const openKeyedStore = vi.fn(() => progressStore.store as unknown as PluginStateKeyedStore); const api = { - config: config ?? { - channels: { discord: { token: "test-token", subagentProgress: true } }, - }, + config: { channels: { discord } }, logger: { debug: vi.fn() }, - runtime: { - state: { - openKeyedStore, - }, - }, - } as unknown as Parameters[0]; + runtime: { state: { openKeyedStore } }, + } as unknown as Parameters[0]; + return { api, progressStore, openKeyedStore }; +} + +function storedRun(overrides: Partial = {}): StoredProgressRun { return { - api, - progressStore, - openKeyedStore, + key: "default:123:456", + accountId: "default", + channelId: "123", + messageId: "456", + runningEmoji: "1️⃣", + status: "active", + ...overrides, }; } -let api: ReturnType["api"]; -let progressStore: ReturnType; - -function started(runId: string) { - return { - phase: "started" as const, - runId, - requester: { - channel: "discord", - accountId: "default", - to: "channel:123", - messageId: "456", - }, - }; -} - -describe("Discord subagent progress", () => { +describe("retired Discord subagent progress cleanup", () => { beforeEach(() => { vi.useFakeTimers(); - ({ api, progressStore } = createApi()); - sendMocks.react.mockReset().mockResolvedValue({ ok: true }); sendMocks.remove.mockReset().mockResolvedValue({ ok: true }); - sendMocks.typing.mockReset().mockResolvedValue({ ok: true, channelId: "123" }); }); afterEach(() => { - resetDiscordSubagentProgressForTest(); + recoverDiscordSubagentProgress.resetForTest(); + vi.unstubAllEnvs(); vi.useRealTimers(); }); - it("updates one source-message count reaction for concurrent runs", async () => { - await handleDiscordSubagentProgress(api, started("run-1")); - await handleDiscordSubagentProgress(api, started("run-2")); - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-1", - outcome: "ok", - }); - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-2", - outcome: "ok", - }); - - expect(sendMocks.react.mock.calls.map((call) => call[2])).toEqual(["1️⃣", "2️⃣", "1️⃣"]); - expect(sendMocks.remove.mock.calls.map((call) => call[2])).toEqual(["1️⃣", "2️⃣", "1️⃣"]); - }); - - it("refreshes reaction ownership inside the source queue for concurrent endings", async () => { - await handleDiscordSubagentProgress(api, started("run-one")); - await handleDiscordSubagentProgress(api, started("run-two")); - sendMocks.remove.mockImplementation(async (_channelId, _messageId, emoji) => ({ - ok: emoji !== "1️⃣" || sendMocks.remove.mock.calls.length !== 3, - })); - - await Promise.all([ - handleDiscordSubagentProgress(api, { phase: "ended", runId: "run-one", outcome: "ok" }), - handleDiscordSubagentProgress(api, { phase: "ended", runId: "run-two", outcome: "ok" }), - ]); - await vi.advanceTimersByTimeAsync(1_000); - - expect(sendMocks.remove).toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(progressStore.values.size).toBe(0); - }); - - it("adds the internal failure marker after clearing the running count", async () => { - await handleDiscordSubagentProgress(api, started("run-error")); - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-error", - outcome: "error", - }); - - expect(sendMocks.remove).toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(sendMocks.react).toHaveBeenLastCalledWith("123", "456", "🔴", expect.any(Object)); - }); - - it("clears a count reaction whose successful add response was lost", async () => { - sendMocks.react.mockResolvedValueOnce({ ok: false }); - - await handleDiscordSubagentProgress(api, started("run-ambiguous-add")); - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-ambiguous-add", - outcome: "ok", - }); - - expect(sendMocks.remove).toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(progressStore.values.has("run-ambiguous-add")).toBe(false); - }); - - it("retries an unconfirmed replacement count reaction", async () => { - await handleDiscordSubagentProgress(api, started("run-one")); - await handleDiscordSubagentProgress(api, started("run-two")); - sendMocks.react.mockResolvedValueOnce({ ok: false }); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-two", - outcome: "ok", - }); - await vi.advanceTimersByTimeAsync(1_000); - - expect(sendMocks.react.mock.calls.filter((call) => call[2] === "1️⃣")).toHaveLength(3); - expect(progressStore.values.has("run-two")).toBe(false); - }); - - it("keeps typing alive only while the source has active runs", async () => { - await handleDiscordSubagentProgress(api, started("run-typing")); - expect(sendMocks.typing).toHaveBeenCalledTimes(1); - - await vi.advanceTimersByTimeAsync(8_500); - expect(sendMocks.typing).toHaveBeenCalledTimes(2); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-typing", - outcome: "ok", - }); - await vi.advanceTimersByTimeAsync(8_500); - expect(sendMocks.typing).toHaveBeenCalledTimes(2); - }); - - it("does nothing unless the single enable toggle is true", async () => { - await handleDiscordSubagentProgress( - createApi({ channels: { discord: { token: "test-token" } } }).api, - started("run-disabled"), + it("opens the shipped namespace and cleans active and terminal beta3 rows", async () => { + vi.stubEnv("VITEST", ""); + const { api, progressStore, openKeyedStore } = createApi(); + progressStore.values.set("active-run", storedRun()); + progressStore.values.set( + "cleanup-run", + storedRun({ status: "cleanup", outcome: "timeout", runningEmoji: "2️⃣" }), ); - expect(sendMocks.react).not.toHaveBeenCalled(); - expect(sendMocks.typing).not.toHaveBeenCalled(); - }); + await recoverDiscordSubagentProgress(api); - it("ignores a started event that arrives after its terminal event", async () => { - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-raced", - outcome: "ok", + expect(openKeyedStore).toHaveBeenCalledWith({ + namespace: "subagent-progress", + maxEntries: 4_096, + overflowPolicy: "reject-new", + defaultTtlMs: 7 * 24 * 60 * 60_000, }); - await handleDiscordSubagentProgress(api, started("run-raced")); - - expect(sendMocks.react).not.toHaveBeenCalled(); - expect(sendMocks.typing).not.toHaveBeenCalled(); - }); - - it("restores durable ownership to clear a running reaction after restart", async () => { - await handleDiscordSubagentProgress(api, started("run-restart")); - resetDiscordSubagentProgressForTest(); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-restart", - outcome: "ok", - }); - + expect(sendMocks.remove.mock.calls.map((call) => call[2])).toEqual(["1️⃣", "2️⃣"]); expect(progressStore.values.size).toBe(0); - expect(sendMocks.remove).toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); }); - it("restores typing for sibling runs when one ends after restart", async () => { - await handleDiscordSubagentProgress(api, started("run-sibling")); - await handleDiscordSubagentProgress(api, started("run-ending")); - resetDiscordSubagentProgressForTest(); - sendMocks.react.mockClear(); - sendMocks.typing.mockClear(); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-ending", - outcome: "ok", + it("removes the external reaction before consuming persisted ownership", async () => { + const { api, progressStore } = createApi(); + progressStore.values.set("owned-run", storedRun()); + sendMocks.remove.mockImplementationOnce(async () => { + expect(progressStore.values.has("owned-run")).toBe(true); + return { ok: true }; }); - expect(sendMocks.react).toHaveBeenLastCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(sendMocks.typing).toHaveBeenCalledTimes(1); + await recoverDiscordSubagentProgress(api); + + expect(progressStore.store.consume).toHaveBeenCalledWith("owned-run"); }); - it("does not restore sibling typing after the feature is disabled", async () => { - await handleDiscordSubagentProgress(api, started("run-sibling")); - await handleDiscordSubagentProgress(api, started("run-ending")); - resetDiscordSubagentProgressForTest(); - sendMocks.typing.mockClear(); - api.config = { channels: { discord: { token: "test-token", subagentProgress: false } } }; + it("does not remove glyphs reserved by current reaction configuration", async () => { + const { api, progressStore } = createApi({ token: "test-token", ackReaction: "1️⃣" }); + progressStore.values.set("reserved-run", storedRun()); - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-ending", - outcome: "ok", - }); + await recoverDiscordSubagentProgress(api); - expect(sendMocks.typing).not.toHaveBeenCalled(); + expect(sendMocks.remove).not.toHaveBeenCalled(); + expect(progressStore.values.size).toBe(0); }); - it("restores active counts before presenting a new run after restart", async () => { - await handleDiscordSubagentProgress(api, started("run-before-restart")); - resetDiscordSubagentProgressForTest(); + it("does not use persisted state to remove a non-historical reaction", async () => { + const { api, progressStore } = createApi(); + progressStore.values.set("unexpected-run", storedRun({ runningEmoji: "🔥" })); - await handleDiscordSubagentProgress(api, started("run-after-restart")); + await recoverDiscordSubagentProgress(api); - expect(sendMocks.remove).toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(sendMocks.react).toHaveBeenLastCalledWith("123", "456", "2️⃣", expect.any(Object)); + expect(sendMocks.remove).not.toHaveBeenCalled(); + expect(progressStore.values.size).toBe(0); }); - it("rebuilds a restored count when durable state outran Discord presentation", async () => { - await handleDiscordSubagentProgress(api, started("run-visible")); - progressStore.values.set("run-not-visible", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "active", - }); - resetDiscordSubagentProgressForTest(); - sendMocks.remove.mockClear(); + it.each([ + ["disabled", { enabled: false, token: "test-token" }], + ["missing credentials", {}], + [ + "configured-unavailable credentials", + { token: { source: "env", provider: "default", id: "MISSING_DISCORD_TOKEN" } }, + ], + ])("retains ownership while the account is %s", async (_label, discord) => { + const { api, progressStore } = createApi(discord); + progressStore.values.set("pending-run", storedRun()); - await handleDiscordSubagentProgress(api, started("run-after-restart")); + await recoverDiscordSubagentProgress(api); - expect(sendMocks.remove).toHaveBeenCalledTimes(1); - expect(sendMocks.remove).toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(sendMocks.remove).not.toHaveBeenCalledWith("123", "456", "2️⃣", expect.any(Object)); - expect(sendMocks.react).toHaveBeenLastCalledWith("123", "456", "3️⃣", expect.any(Object)); + expect(sendMocks.remove).not.toHaveBeenCalled(); + expect(progressStore.values.has("pending-run")).toBe(true); }); - it("keeps failed terminal cleanup durable until the next source event", async () => { - await handleDiscordSubagentProgress(api, started("run-cleanup")); - sendMocks.remove.mockRejectedValueOnce(new Error("transient Discord failure")); + it("does not clean a removed named account with the current root bot", async () => { + const { api, progressStore } = createApi({ token: "current-root-token" }); + progressStore.values.set("retired-run", storedRun({ accountId: "retired" })); - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-cleanup", - outcome: "ok", - }); + await recoverDiscordSubagentProgress(api); - expect(progressStore.values.get("run-cleanup")?.status).toBe("cleanup"); - resetDiscordSubagentProgressForTest(); - - await handleDiscordSubagentProgress(api, started("run-after-cleanup")); - - expect(progressStore.values.has("run-cleanup")).toBe(false); - expect(progressStore.values.get("run-after-cleanup")?.status).toBe("active"); - expect(sendMocks.react).toHaveBeenLastCalledWith("123", "456", "1️⃣", expect.any(Object)); + expect(sendMocks.remove).not.toHaveBeenCalled(); + expect(progressStore.store.consume).not.toHaveBeenCalled(); + expect(progressStore.values.has("retired-run")).toBe(true); }); - it("retries failed terminal reaction cleanup without another source event", async () => { - await handleDiscordSubagentProgress(api, started("run-retry-cleanup")); - sendMocks.remove.mockRejectedValueOnce(new Error("transient Discord failure")); + it("retains ownership and retries a failed reaction cleanup", async () => { + const { api, progressStore } = createApi(); + progressStore.values.set("retry-run", storedRun()); + sendMocks.remove.mockRejectedValueOnce(new Error("Discord unavailable")); - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-retry-cleanup", - outcome: "ok", - }); - expect(progressStore.values.get("run-retry-cleanup")?.status).toBe("cleanup"); + await recoverDiscordSubagentProgress(api); + expect(progressStore.values.has("retry-run")).toBe(true); await vi.advanceTimersByTimeAsync(1_000); - expect(progressStore.values.has("run-retry-cleanup")).toBe(false); - expect(sendMocks.remove).toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); + expect(sendMocks.remove).toHaveBeenCalledTimes(2); + expect(progressStore.values.has("retry-run")).toBe(false); }); - it("retries a resolved false terminal reaction removal", async () => { - await handleDiscordSubagentProgress(api, started("run-false-cleanup")); - sendMocks.remove.mockResolvedValueOnce({ ok: false }); + it("retains ownership when the persistence cleanup fails", async () => { + const { api, progressStore } = createApi(); + progressStore.values.set("retry-run", storedRun()); + progressStore.store.consume.mockRejectedValueOnce(new Error("state unavailable")); - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-false-cleanup", - outcome: "ok", - }); - expect(progressStore.values.get("run-false-cleanup")?.status).toBe("cleanup"); + await recoverDiscordSubagentProgress(api); + expect(progressStore.values.has("retry-run")).toBe(true); await vi.advanceTimersByTimeAsync(1_000); - expect(progressStore.values.has("run-false-cleanup")).toBe(false); + expect(sendMocks.remove).toHaveBeenCalledTimes(2); + expect(progressStore.values.has("retry-run")).toBe(false); }); - it("retries a failed terminal outcome marker before consuming ownership", async () => { - await handleDiscordSubagentProgress(api, started("run-failure-marker")); - sendMocks.react.mockResolvedValueOnce({ ok: false }).mockResolvedValue({ ok: true }); + it("bounds retries while retaining failed cleanup ownership", async () => { + const { api, progressStore } = createApi(); + progressStore.values.set("retry-run", storedRun()); + sendMocks.remove.mockRejectedValue(new Error("Discord unavailable")); - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-failure-marker", - outcome: "error", - }); - expect(progressStore.values.get("run-failure-marker")).toMatchObject({ - status: "cleanup", - outcome: "error", - }); - - await vi.advanceTimersByTimeAsync(1_000); - - expect(progressStore.values.has("run-failure-marker")).toBe(false); - expect(sendMocks.react.mock.calls.filter((call) => call[2] === "🔴")).toHaveLength(2); - }); - - it("preserves cleanup ownership when a retry refresh lookup fails", async () => { - await handleDiscordSubagentProgress(api, started("run-refresh-error")); - sendMocks.remove.mockResolvedValueOnce({ ok: false }); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-refresh-error", - outcome: "ok", - }); - progressStore.store.lookup.mockImplementationOnce(async (key) => progressStore.values.get(key)); - progressStore.store.lookup.mockRejectedValueOnce(new Error("state store unavailable")); - await vi.advanceTimersByTimeAsync(1_000); - - expect(sendMocks.remove).toHaveBeenLastCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(progressStore.values.has("run-refresh-error")).toBe(false); - }); - - it("bounds permanent terminal reaction retries while retaining ownership", async () => { - await handleDiscordSubagentProgress(api, started("run-permanent-failure")); - sendMocks.remove.mockResolvedValue({ ok: false }); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-permanent-failure", - outcome: "ok", - }); + await recoverDiscordSubagentProgress(api); await vi.advanceTimersByTimeAsync(2 * 60 * 60_000); expect(sendMocks.remove).toHaveBeenCalledTimes(13); - expect(progressStore.values.get("run-permanent-failure")?.status).toBe("cleanup"); + expect(progressStore.values.has("retry-run")).toBe(true); + expect(vi.getTimerCount()).toBe(0); }); - it("recovers durable cleanup rows on gateway startup", async () => { - progressStore.values.set("run-restart-cleanup", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "cleanup", - outcome: "ok", - runningEmoji: "1️⃣", - }); - - await recoverDiscordSubagentProgress(api); - - expect(sendMocks.remove).toHaveBeenCalledTimes(1); - expect(progressStore.values.has("run-restart-cleanup")).toBe(false); - }); - - it("restores a persisted terminal outcome marker on gateway startup", async () => { - progressStore.values.set("run-restart-timeout", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "cleanup", - outcome: "timeout", - runningEmoji: "1️⃣", - }); - - await recoverDiscordSubagentProgress(api); - - expect(sendMocks.remove).toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(sendMocks.react).toHaveBeenCalledWith("123", "456", "🔴", expect.any(Object)); - expect(progressStore.values.has("run-restart-timeout")).toBe(false); - }); - - it("cleans interrupted active rows on gateway startup", async () => { - progressStore.values.set("run-interrupted", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "active", - runningEmoji: "1️⃣", - }); - - await recoverDiscordSubagentProgress(api); - - expect(sendMocks.remove).toHaveBeenCalledTimes(1); - expect(progressStore.values.has("run-interrupted")).toBe(false); - }); - - it("does not remove a progress glyph newly reserved by configuration", async () => { - const collision = createApi({ - messages: { ackReaction: "1️⃣" }, - channels: { discord: { token: "test-token", subagentProgress: true } }, - }); - collision.progressStore.values.set("run-collision-cleanup", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "cleanup", - outcome: "ok", - runningEmoji: "1️⃣", - }); - - await recoverDiscordSubagentProgress(collision.api); - - expect(sendMocks.remove).not.toHaveBeenCalled(); - expect(sendMocks.remove).not.toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(collision.progressStore.values.has("run-collision-cleanup")).toBe(false); - }); - - it("removes owned progress reactions after the feature is disabled", async () => { - const disabled = createApi({ channels: { discord: { token: "test-token" } } }); - disabled.progressStore.values.set("run-disabled-cleanup", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "cleanup", - outcome: "ok", - runningEmoji: "1️⃣", - }); - - await recoverDiscordSubagentProgress(disabled.api); - - expect(sendMocks.remove).toHaveBeenCalledTimes(1); - expect(disabled.progressStore.values.has("run-disabled-cleanup")).toBe(false); - }); - - it("stops live presentation when the feature is disabled", async () => { - await handleDiscordSubagentProgress(api, started("run-one")); - await handleDiscordSubagentProgress(api, started("run-two")); - api.config.channels!.discord!.subagentProgress = false; - sendMocks.react.mockClear(); - sendMocks.remove.mockClear(); - const typingCalls = sendMocks.typing.mock.calls.length; - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-two", - outcome: "error", - }); - await vi.advanceTimersByTimeAsync(8_500); - - expect(sendMocks.remove).toHaveBeenCalledTimes(1); - expect(sendMocks.react).not.toHaveBeenCalled(); - expect(sendMocks.typing).toHaveBeenCalledTimes(typingCalls); - }); - - it("tears down live presentation when a disabled start arrives", async () => { - await handleDiscordSubagentProgress(api, started("run-active")); - api.config.channels!.discord!.subagentProgress = false; - sendMocks.remove.mockClear(); - const typingCalls = sendMocks.typing.mock.calls.length; - - await handleDiscordSubagentProgress(api, started("run-ignored")); - await vi.advanceTimersByTimeAsync(8_500); - - expect(sendMocks.remove).toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(sendMocks.typing).toHaveBeenCalledTimes(typingCalls); - expect(progressStore.values.has("run-ignored")).toBe(false); - }); - - it("retains cleanup ownership while the Discord account is disabled", async () => { - const disabled = createApi({ - channels: { - discord: { enabled: false, token: "test-token", subagentProgress: true }, - }, - }); - disabled.progressStore.values.set("run-account-disabled", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "cleanup", - outcome: "ok", - runningEmoji: "1️⃣", - }); - - await recoverDiscordSubagentProgress(disabled.api); - - expect(sendMocks.remove).not.toHaveBeenCalled(); - expect(disabled.progressStore.values.get("run-account-disabled")?.status).toBe("cleanup"); - }); - - it("retries startup recovery after a transient state listing failure", async () => { - progressStore.values.set("run-list-retry", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "cleanup", - outcome: "ok", - runningEmoji: "1️⃣", - }); - progressStore.store.entries.mockRejectedValueOnce(new Error("state store unavailable")); + it("retries a transient state listing failure", async () => { + const { api, progressStore } = createApi(); + progressStore.values.set("retry-run", storedRun()); + progressStore.store.entries.mockRejectedValueOnce(new Error("state unavailable")); await recoverDiscordSubagentProgress(api); await vi.advanceTimersByTimeAsync(1_000); - expect(progressStore.store.entries).toHaveBeenCalledTimes(3); - expect(progressStore.values.has("run-list-retry")).toBe(false); - }); - - it("uses listed ownership when startup lookup fails transiently", async () => { - progressStore.values.set("run-lookup-recovery", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "cleanup", - outcome: "ok", - runningEmoji: "1️⃣", - }); - progressStore.store.lookup.mockRejectedValueOnce(new Error("state store unavailable")); - - await recoverDiscordSubagentProgress(api); - - expect(sendMocks.remove).toHaveBeenCalledTimes(1); - expect(progressStore.values.has("run-lookup-recovery")).toBe(false); - }); - - it("rebuilds presentation for a duplicate start replayed after restart", async () => { - await handleDiscordSubagentProgress(api, started("run-replayed")); - resetDiscordSubagentProgressForTest(); - sendMocks.react.mockClear(); - sendMocks.remove.mockClear(); - sendMocks.typing.mockClear(); - - await handleDiscordSubagentProgress(api, started("run-replayed")); - - expect(progressStore.values.size).toBe(1); - expect(sendMocks.remove).toHaveBeenCalledTimes(1); - expect(sendMocks.react).toHaveBeenLastCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(sendMocks.typing).toHaveBeenCalledTimes(1); - }); - - it("retains cleanup when restoring the remaining count fails", async () => { - await handleDiscordSubagentProgress(api, started("run-remaining")); - await handleDiscordSubagentProgress(api, started("run-ending")); - resetDiscordSubagentProgressForTest(); - sendMocks.react.mockRejectedValueOnce(new Error("transient Discord failure")); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-ending", - outcome: "ok", - }); - - expect(progressStore.values.get("run-remaining")?.status).toBe("active"); - expect(progressStore.values.get("run-ending")?.status).toBe("cleanup"); - }); - - it("persists a replacement count before adding it during restart recovery", async () => { - await handleDiscordSubagentProgress(api, started("run-remaining")); - await handleDiscordSubagentProgress(api, started("run-ending")); - resetDiscordSubagentProgressForTest(); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-ending", - outcome: "ok", - }); - - expect(progressStore.values.get("run-remaining")?.runningEmoji).toBe("1️⃣"); - resetDiscordSubagentProgressForTest(); - sendMocks.remove.mockClear(); - await recoverDiscordSubagentProgress(api); - expect(sendMocks.remove).toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); - }); - - it("does not resurrect a durable cleanup row on a replayed start", async () => { - progressStore.values.set("run-terminal", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "cleanup", - outcome: "ok", - }); - - await handleDiscordSubagentProgress(api, started("run-terminal")); - - expect(progressStore.values.has("run-terminal")).toBe(false); - expect(sendMocks.react).not.toHaveBeenCalled(); - expect(sendMocks.typing).not.toHaveBeenCalled(); - }); - - it("restores other active runs when a terminal start is replayed", async () => { - progressStore.values.set("run-active", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "active", - }); - progressStore.values.set("run-terminal", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "cleanup", - outcome: "ok", - }); - - await handleDiscordSubagentProgress(api, started("run-terminal")); - - expect(progressStore.values.has("run-terminal")).toBe(false); - expect(sendMocks.react).toHaveBeenLastCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(sendMocks.typing).toHaveBeenCalledTimes(1); - }); - - it("still clears in-memory presentation when cleanup marking fails", async () => { - await handleDiscordSubagentProgress(api, started("run-store-error")); - progressStore.store.register.mockRejectedValueOnce(new Error("state store unavailable")); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-store-error", - outcome: "ok", - }); - await vi.advanceTimersByTimeAsync(8_500); - - expect(sendMocks.remove).toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(sendMocks.typing).toHaveBeenCalledTimes(1); - expect(progressStore.values.has("run-store-error")).toBe(false); - }); - - it("excludes the ending row when cleanup marking fails after restart", async () => { - await handleDiscordSubagentProgress(api, started("run-store-error")); - resetDiscordSubagentProgressForTest(); - sendMocks.react.mockClear(); - sendMocks.remove.mockClear(); - progressStore.store.register.mockRejectedValueOnce(new Error("state store unavailable")); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-store-error", - outcome: "ok", - }); - - expect(sendMocks.remove).toHaveBeenCalledTimes(1); - expect(sendMocks.react).not.toHaveBeenCalled(); - expect(progressStore.values.has("run-store-error")).toBe(false); - }); - - it("retains cleanup without mutating reactions when state listing fails", async () => { - await handleDiscordSubagentProgress(api, started("run-list-error")); - resetDiscordSubagentProgressForTest(); - sendMocks.remove.mockClear(); - progressStore.store.entries.mockRejectedValueOnce(new Error("state store unavailable")); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-list-error", - outcome: "ok", - }); - - expect(sendMocks.remove).not.toHaveBeenCalled(); - expect(progressStore.values.get("run-list-error")?.status).toBe("cleanup"); - }); - - it("repairs an in-memory run after a durable lookup failure", async () => { - await handleDiscordSubagentProgress(api, started("run-lookup-error")); - progressStore.store.lookup.mockRejectedValueOnce(new Error("state store unavailable")); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-lookup-error", - outcome: "ok", - }); - - expect(sendMocks.remove).toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(progressStore.values.has("run-lookup-error")).toBe(false); - }); - - it("retries a terminal lookup after restart", async () => { - await handleDiscordSubagentProgress(api, started("run-lookup-retry")); - resetDiscordSubagentProgressForTest(); - sendMocks.remove.mockClear(); - progressStore.store.lookup.mockRejectedValueOnce(new Error("state store unavailable")); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-lookup-retry", - outcome: "ok", - requester: started("ignored").requester, - }); - expect(sendMocks.remove).not.toHaveBeenCalled(); - - await vi.advanceTimersByTimeAsync(1_000); - - expect(sendMocks.remove).toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(progressStore.values.has("run-lookup-retry")).toBe(false); - }); - - it("reopens the state store after a transient open failure", async () => { - const reopened = createApi(); - api = reopened.api; - progressStore = reopened.progressStore; - progressStore.values.set("run-open-retry", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "active", - }); - reopened.openKeyedStore.mockImplementationOnce(() => { - throw new Error("state store unavailable"); - }); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-open-retry", - outcome: "ok", - requester: started("ignored").requester, - }); - await vi.advanceTimersByTimeAsync(1_000); - - expect(reopened.openKeyedStore).toHaveBeenCalledTimes(2); - expect(progressStore.values.has("run-open-retry")).toBe(false); - }); - - it("does not retry an unowned non-Discord terminal event", async () => { - progressStore.store.lookup.mockRejectedValue(new Error("state store unavailable")); - - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-slack", - outcome: "ok", - requester: { channel: "slack", channelId: "123", messageId: "456" }, - }); - await vi.advanceTimersByTimeAsync(60_000); - - expect(progressStore.store.lookup).toHaveBeenCalledTimes(1); - }); - - it("never reconciles reactions for a typing-only tracker", async () => { - const typingOnly = createApi(); - typingOnly.progressStore.store.registerIfAbsent.mockRejectedValueOnce( - new Error("state store unavailable"), - ); - - await handleDiscordSubagentProgress(typingOnly.api, started("run-typing-only")); - await handleDiscordSubagentProgress(typingOnly.api, { - phase: "ended", - runId: "run-typing-only", - outcome: "ok", - }); - - expect(sendMocks.remove).not.toHaveBeenCalled(); - expect(sendMocks.typing).toHaveBeenCalledTimes(1); - }); - - it("defers a new start when restored state cannot be listed", async () => { - progressStore.values.set("run-existing", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "active", - }); - progressStore.store.entries.mockRejectedValueOnce(new Error("state store unavailable")); - - await handleDiscordSubagentProgress(api, started("run-deferred")); - - expect(progressStore.values.has("run-deferred")).toBe(false); - expect(sendMocks.react).not.toHaveBeenCalled(); - expect(sendMocks.typing).not.toHaveBeenCalled(); - }); - - it("keeps consumed cleanup rows as process terminal tombstones", async () => { - progressStore.values.set("run-terminal", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "cleanup", - outcome: "ok", - }); - - await handleDiscordSubagentProgress(api, started("run-other")); - await handleDiscordSubagentProgress(api, started("run-terminal")); - - expect(progressStore.values.has("run-terminal")).toBe(false); - expect(progressStore.values.get("run-other")?.status).toBe("active"); - expect(sendMocks.react).toHaveBeenLastCalledWith("123", "456", "1️⃣", expect.any(Object)); - }); - - it("does not overwrite a cleanup tombstone while a sibling tracker is active", async () => { - await handleDiscordSubagentProgress(api, started("run-active")); - progressStore.values.set("run-terminal", { - key: "default:123:456", - accountId: "default", - channelId: "123", - messageId: "456", - status: "cleanup", - outcome: "ok", - }); - sendMocks.react.mockClear(); - - await handleDiscordSubagentProgress(api, started("run-terminal")); - - expect(progressStore.values.get("run-terminal")?.status).toBe("cleanup"); - expect(sendMocks.react).not.toHaveBeenCalled(); - }); - - it("does not treat a terminal-event locator as reaction ownership", async () => { - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-terminal-locator", - outcome: "ok", - requester: started("ignored").requester, - }); - - expect(sendMocks.remove).not.toHaveBeenCalled(); - }); - - it("keeps reactions off when an agent identity owns a progress emoji", async () => { - const collisionApi = createApi({ - agents: { list: [{ id: "main", identity: { emoji: "1️⃣" } }] }, - channels: { discord: { token: "test-token", subagentProgress: true } }, - }).api; - - await handleDiscordSubagentProgress(collisionApi, started("run-collision")); - - expect(sendMocks.react).not.toHaveBeenCalled(); - expect(sendMocks.typing).toHaveBeenCalledTimes(1); - }); - - it("preserves a live progress glyph newly reserved by configuration", async () => { - await handleDiscordSubagentProgress(api, started("run-before-reload")); - api.config.messages = { ackReaction: "1️⃣" }; - sendMocks.react.mockClear(); - sendMocks.remove.mockClear(); - - await handleDiscordSubagentProgress(api, started("run-after-reload")); - await handleDiscordSubagentProgress(api, { - phase: "ended", - runId: "run-before-reload", - outcome: "ok", - }); - - expect(sendMocks.remove).not.toHaveBeenCalledWith("123", "456", "1️⃣", expect.any(Object)); - expect(sendMocks.react).not.toHaveBeenCalled(); - }); - - it("retains ownership when collision cleanup fails", async () => { - await handleDiscordSubagentProgress(api, started("run-before-collision")); - api.config.messages = { ackReaction: "2️⃣" }; - sendMocks.remove.mockResolvedValueOnce({ ok: false }); - - await handleDiscordSubagentProgress(api, started("run-after-collision")); - - expect(progressStore.values.get("run-before-collision")?.runningEmoji).toBe("1️⃣"); - expect(progressStore.values.has("run-after-collision")).toBe(false); - }); - - it("uses the native source channel for Discord DM routes", async () => { - await handleDiscordSubagentProgress(api, { - ...started("run-dm"), - requester: { - channel: "discord", - accountId: "default", - to: "user:789", - channelId: "987", - messageId: "456", - }, - }); - - expect(sendMocks.react).toHaveBeenCalledWith("987", "456", "1️⃣", expect.any(Object)); - expect(sendMocks.typing).toHaveBeenCalledWith("987", expect.any(Object)); + expect(progressStore.store.entries).toHaveBeenCalledTimes(2); + expect(progressStore.values.has("retry-run")).toBe(false); }); }); diff --git a/extensions/discord/src/subagent-progress.ts b/extensions/discord/src/subagent-progress.ts index 73a7149e2e7b..8932eb8c1449 100644 --- a/extensions/discord/src/subagent-progress.ts +++ b/extensions/discord/src/subagent-progress.ts @@ -1,735 +1,188 @@ -// Discord plugin module maps portable subagent progress onto source-message feedback. -import { resolveDiscordAccount } from "./accounts.js"; -import { reactMessageDiscord, removeReactionDiscord } from "./send.reactions.js"; -import { sendTypingDiscord } from "./send.typing.js"; -import { - FAILURE_EMOJI, - RUNNING_EMOJIS, - reactionsAreAvailable, - reservedReactionEmojis, - resolveDiscordProgressTarget, - type DiscordProgressRequester, -} from "./subagent-progress-config.js"; -import { - MAX_TRACKED_RUNS, - PROGRESS_STORE_TTL_MS, - consumeProgressRun, - getProgressStore, - listProgressStateForKey, - logFailure, - lookupProgressRun, - markRunTerminal, - markProgressRunForCleanup, - persistedProgressRunFromTracker, - persistedTerminalOutcome, - persistProgressRun, - resetDiscordSubagentProgressStateForTest, - runQueued, - terminalOutcome, - type PersistedProgressRun, - type PersistProgressResult, - type ProgressApi, - type ProgressTracker, - type SubagentProgressOutcome, -} from "./subagent-progress-state.js"; +// Discord plugin module cleans up reactions left by the retired subagent progress feature. +import { DEFAULT_EMOJIS } from "openclaw/plugin-sdk/channel-feedback"; +import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; +import type { + PluginStateEntry, + PluginStateKeyedStore, +} from "openclaw/plugin-sdk/plugin-state-runtime"; +import { DEFAULT_ACCOUNT_ID } from "openclaw/plugin-sdk/routing"; +import { resolveDiscordAccount, resolveDiscordAccountConfig } from "./accounts.js"; +import { removeReactionDiscord } from "./send.reactions.js"; -const TYPING_INTERVAL_MS = 8_500; -const TYPING_TTL_MS = 60 * 60_000; -const TERMINAL_LOOKUP_RETRY_MS = 1_000; -const TERMINAL_RETRY_MAX_DELAY_MS = 60 * 60_000; -const TERMINAL_RETRY_MAX_ATTEMPTS = 12; -const STARTUP_RETRY_MAX_ATTEMPTS = 12; +// Keep beta3's namespace lifecycle so its reaction ownership remains discoverable after upgrade. +const PROGRESS_STORE_TTL_MS = 7 * 24 * 60 * 60_000; +const MAX_TRACKED_RUNS = 4_096; +const RETRY_BASE_DELAY_MS = 1_000; +const RETRY_MAX_DELAY_MS = 60 * 60_000; +const RETRY_MAX_ATTEMPTS = 12; +const HISTORICAL_RUNNING_EMOJIS = new Set([ + "1️⃣", + "2️⃣", + "3️⃣", + "4️⃣", + "5️⃣", + "6️⃣", + "7️⃣", + "8️⃣", + "9️⃣", + "🔟", +]); -type SubagentProgressEvent = - | { - phase: "started"; - runId: string; - requester?: DiscordProgressRequester; - } - | { - phase: "ended"; - runId: string; - outcome: SubagentProgressOutcome; - requester?: Extract["requester"]; +type PersistedProgressRun = { + accountId: string; + channelId: string; + messageId: string; + runningEmoji?: string; +}; + +type ProgressCleanupApi = { + config: OpenClawConfig; + logger: { debug?: (message: string) => void }; + runtime: { + state: { + openKeyedStore(options: { + namespace: string; + maxEntries: number; + overflowPolicy: "reject-new"; + defaultTtlMs: number; + }): PluginStateKeyedStore; }; - -type PersistedReconciliationResult = - | { ok: false } - | { - ok: true; - activeRunIds: string[]; - reactionsEnabled: boolean; - typingEnabled: boolean; - runningEmoji?: string; - }; -type StartupRecovery = { attempts: number; timer?: ReturnType }; - -const trackers = new Map(); - -const trackerKeyByRunId = new Map(); -const terminalRetryTimers = new Map>(); -const terminalRetryExpiresAt = new Map(); -const terminalRetryAttempts = new Map(); -const startupRecoveryRetries = new Map(); - -function isSubagentProgressEnabled(account: ReturnType): boolean { - return Boolean(process.env.VITEST) && account.enabled && account.config.subagentProgress === true; -} - -function clearTerminalRetry(runId: string) { - const timer = terminalRetryTimers.get(runId); - if (timer) { - clearTimeout(timer); - } - terminalRetryTimers.delete(runId); - terminalRetryExpiresAt.delete(runId); - terminalRetryAttempts.delete(runId); -} - -function cancelTerminalRetryTimer(runId: string) { - const timer = terminalRetryTimers.get(runId); - if (timer) { - clearTimeout(timer); - } - terminalRetryTimers.delete(runId); -} - -async function setReaction(api: ProgressApi, tracker: ProgressTracker, emoji: string) { - try { - const result = await reactMessageDiscord(tracker.channelId, tracker.messageId, emoji, { - cfg: api.config, - accountId: tracker.accountId, - }); - return result.ok; - } catch (error) { - logFailure(api, "reaction add", error); - return false; - } -} - -async function clearReaction(api: ProgressApi, tracker: ProgressTracker, emoji?: string) { - if (!emoji) { - return true; - } - try { - const result = await removeReactionDiscord(tracker.channelId, tracker.messageId, emoji, { - cfg: api.config, - accountId: tracker.accountId, - }); - return result.ok; - } catch (error) { - logFailure(api, "reaction remove", error); - return false; - } -} - -async function clearRunningReactions( - api: ProgressApi, - tracker: ProgressTracker, - emojis: readonly string[], -) { - const results = await Promise.all(emojis.map((emoji) => clearReaction(api, tracker, emoji))); - return results.every(Boolean); -} - -async function persistTrackerRunningEmoji(api: ProgressApi, tracker: ProgressTracker) { - const store = getProgressStore(api); - if (!store) { - return false; - } - try { - await Promise.all( - Array.from(tracker.persistedRunIds, (runId) => - store.register(runId, persistedProgressRunFromTracker(tracker)), - ), - ); - return true; - } catch (error) { - logFailure(api, "reaction ownership write", error); - return false; - } -} - -async function updateRunningReaction(api: ProgressApi, tracker: ProgressTracker) { - if (!tracker.reactionsEnabled) { - return true; - } - const nextEmoji = - tracker.activeRunIds.size > 0 - ? RUNNING_EMOJIS[Math.min(tracker.activeRunIds.size, RUNNING_EMOJIS.length) - 1] - : undefined; - if (nextEmoji === tracker.runningEmoji) { - if (!nextEmoji || tracker.runningEmojiConfirmed) { - return true; - } - if (!(await persistTrackerRunningEmoji(api, tracker))) { - return false; - } - tracker.runningEmojiConfirmed = await setReaction(api, tracker, nextEmoji); - return tracker.runningEmojiConfirmed; - } - if (!(await clearReaction(api, tracker, tracker.runningEmoji))) { - return false; - } - tracker.runningEmoji = undefined; - tracker.runningEmojiConfirmed = false; - if (nextEmoji) { - // A lost add response may still have applied; keep attempted ownership for cleanup. - tracker.runningEmoji = nextEmoji; - if (!(await persistTrackerRunningEmoji(api, tracker))) { - tracker.runningEmoji = undefined; - return false; - } - tracker.runningEmojiConfirmed = await setReaction(api, tracker, nextEmoji); - return tracker.runningEmojiConfirmed; - } - await persistTrackerRunningEmoji(api, tracker); - return true; -} - -async function disableTrackerReactionsOnCollision( - api: ProgressApi, - tracker: ProgressTracker, - ackReaction?: string, -) { - if (!tracker.reactionsEnabled || reactionsAreAvailable(api.config, ackReaction)) { - return true; - } - const reserved = reservedReactionEmojis(api.config, ackReaction); - if (tracker.runningEmoji && !reserved.has(tracker.runningEmoji)) { - if (!(await clearReaction(api, tracker, tracker.runningEmoji))) { - tracker.reactionsEnabled = false; - return false; - } - } - tracker.runningEmoji = undefined; - tracker.runningEmojiConfirmed = false; - tracker.reactionsEnabled = false; - await persistTrackerRunningEmoji(api, tracker); - return true; -} - -async function sendTyping(api: ProgressApi, tracker: ProgressTracker) { - try { - await sendTypingDiscord(tracker.channelId, { - cfg: api.config, - accountId: tracker.accountId, - }); - } catch (error) { - logFailure(api, "typing", error); - } -} - -function startTyping(api: ProgressApi, tracker: ProgressTracker) { - tracker.typingExpiresAt = Date.now() + TYPING_TTL_MS; - void sendTyping(api, tracker); - if (tracker.typingTimer) { - return; - } - tracker.typingTimer = setInterval(() => { - if (tracker.activeRunIds.size === 0 || Date.now() >= tracker.typingExpiresAt) { - stopTyping(tracker); - return; - } - void sendTyping(api, tracker); - }, TYPING_INTERVAL_MS); - tracker.typingTimer.unref?.(); -} - -function stopTyping(tracker: ProgressTracker) { - if (tracker.typingTimer) { - clearInterval(tracker.typingTimer); - tracker.typingTimer = undefined; - } -} - -async function handleStarted( - api: ProgressApi, - event: Extract, -) { - const runId = event.runId.trim(); - const target = resolveDiscordProgressTarget(event.requester); - if (!runId || !target || terminalOutcome(runId)) { - return; - } - const account = resolveDiscordAccount({ cfg: api.config, accountId: event.requester?.accountId }); - const key = `${account.accountId}:${target.channelId}:${target.messageId}`; - if (!isSubagentProgressEnabled(account)) { - await runQueued(key, async () => { - const tracker = trackers.get(key); - if (!tracker) { - return; - } - stopTyping(tracker); - tracker.reactionsEnabled = false; - if (!account.enabled || !tracker.runningEmoji) { - return; - } - const reserved = reservedReactionEmojis(api.config, account.config.ackReaction); - if ( - !reserved.has(tracker.runningEmoji) && - (await clearReaction(api, tracker, tracker.runningEmoji)) - ) { - tracker.runningEmoji = undefined; - tracker.runningEmojiConfirmed = false; - await persistTrackerRunningEmoji(api, tracker); - } - }); - return; - } - await runQueued(key, async () => { - let tracker = trackers.get(key); - let restoredCurrentRunWasTerminal = false; - if (!tracker) { - const reactionsEnabled = reactionsAreAvailable(api.config, account.config.ackReaction); - const restored = await listProgressStateForKey(api, key); - if (!restored.ok) { - return; - } - restoredCurrentRunWasTerminal = restored.cleanupRuns.some( - (cleanup) => cleanup.runId === runId, - ); - tracker = { - accountId: account.accountId, - channelId: target.channelId, - messageId: target.messageId, - activeRunIds: new Set(restored.activeRunIds), - persistedRunIds: new Set(restored.activeRunIds), - runningEmojiConfirmed: false, - reactionsEnabled, - typingExpiresAt: 0, - }; - // Recover post-registration crashes from bot-owned glyphs instead of guessing. - if (restored.activeRunIds.length > 0 || restored.cleanupRuns.length > 0) { - const reserved = reservedReactionEmojis(api.config, account.config.ackReaction); - const cleanupEmojis = restored.ownedEmojis.filter((emoji) => !reserved.has(emoji)); - const countsCleared = await clearRunningReactions(api, tracker, cleanupEmojis); - const failedCleanupRuns = restored.cleanupRuns.filter( - (cleanup) => cleanup.value.outcome !== "ok", - ); - const failurePresented = - failedCleanupRuns.length === 0 || - !reactionsEnabled || - (countsCleared && (await setReaction(api, tracker, FAILURE_EMOJI))); - if (countsCleared && failurePresented) { - for (const cleanup of restored.cleanupRuns) { - markRunTerminal(cleanup.runId, cleanup.value.outcome); - } - await Promise.all( - restored.cleanupRuns.map((cleanup) => consumeProgressRun(api, cleanup.runId)), - ); - } else { - tracker.reactionsEnabled = false; - for (const cleanup of restored.cleanupRuns) { - scheduleTerminalLookupRetry( - api, - { - phase: "ended", - runId: cleanup.runId, - outcome: cleanup.value.outcome, - }, - cleanup.value, - ); - } - } - } - trackers.set(key, tracker); - } - if (!(await disableTrackerReactionsOnCollision(api, tracker, account.config.ackReaction))) { - return; - } - if (restoredCurrentRunWasTerminal) { - if (tracker.activeRunIds.size > 0) { - await updateRunningReaction(api, tracker); - startTyping(api, tracker); - } else { - trackers.delete(key); - } - return; - } - if (tracker.activeRunIds.has(runId)) { - trackerKeyByRunId.set(runId, key); - await updateRunningReaction(api, tracker); - startTyping(api, tracker); - return; - } - let persistResult: PersistProgressResult = "error"; - if (tracker.reactionsEnabled) { - persistResult = await persistProgressRun(api, runId, tracker); - if (persistResult === "terminal") { - markRunTerminal(runId, "unknown"); - return; - } - if (persistResult === "conflict") { - api.logger.debug?.(`discord subagent progress ignored conflicting run id: ${runId}`); - return; - } - if (persistResult === "error") { - await clearReaction(api, tracker, tracker.runningEmoji); - tracker.runningEmoji = undefined; - tracker.runningEmojiConfirmed = false; - tracker.reactionsEnabled = false; - } - } - tracker.activeRunIds.add(runId); - trackerKeyByRunId.set(runId, key); - if (persistResult === "persisted") { - tracker.persistedRunIds.add(runId); - } - // Tombstone a fast child that ends before presentation so a late start cannot stick. - const endedOutcome = terminalOutcome(runId); - if (endedOutcome) { - const owned = persistedProgressRunFromTracker(tracker); - await markProgressRunForCleanup(api, runId, owned, endedOutcome); - const failurePresented = - endedOutcome === "ok" || - !tracker.reactionsEnabled || - (await setReaction(api, tracker, FAILURE_EMOJI)); - if (failurePresented) { - await consumeProgressRun(api, runId); - } else { - scheduleTerminalLookupRetry( - api, - { phase: "ended", runId, outcome: endedOutcome, requester: event.requester }, - { ...owned, status: "cleanup", outcome: endedOutcome }, - ); - } - tracker.activeRunIds.delete(runId); - tracker.persistedRunIds.delete(runId); - trackerKeyByRunId.delete(runId); - await updateRunningReaction(api, tracker); - if (tracker.activeRunIds.size === 0) { - stopTyping(tracker); - trackers.delete(key); - } - return; - } - await updateRunningReaction(api, tracker); - startTyping(api, tracker); - }); -} - -async function reconcilePersistedTracker( - api: ProgressApi, - persisted: PersistedProgressRun, - outcome: Extract["outcome"], - endingRunId: string, -): Promise { - const store = getProgressStore(api); - let activeRunIds: string[] = []; - if (store) { - try { - const entries = await store.entries(); - activeRunIds = entries - .filter( - (entry) => - entry.key !== endingRunId && - entry.value.key === persisted.key && - entry.value.status === "active", - ) - .map((entry) => entry.key); - } catch (error) { - logFailure(api, "state store list", error); - return { ok: false }; - } - } - const tracker: ProgressTracker = { - accountId: persisted.accountId, - channelId: persisted.channelId, - messageId: persisted.messageId, - activeRunIds: new Set(activeRunIds), - persistedRunIds: new Set(activeRunIds), - runningEmojiConfirmed: false, - reactionsEnabled: true, - typingExpiresAt: 0, - }; - const account = resolveDiscordAccount({ cfg: api.config, accountId: persisted.accountId }); - const typingEnabled = isSubagentProgressEnabled(account); - const reserved = reservedReactionEmojis(api.config, account.config.ackReaction); - const cleanupEmojis = - persisted.runningEmoji && !reserved.has(persisted.runningEmoji) ? [persisted.runningEmoji] : []; - const reactionsEnabled = - typingEnabled && reactionsAreAvailable(api.config, account.config.ackReaction); - // Preserve newly reserved keycaps, but remove every unreserved glyph that - // this feature could have left behind under the previous configuration. - const reactionsCleared = - account.enabled && (await clearRunningReactions(api, tracker, cleanupEmojis)); - const nextEmoji = RUNNING_EMOJIS[Math.min(activeRunIds.length, RUNNING_EMOJIS.length) - 1]; - let countPresented = true; - if (reactionsEnabled && reactionsCleared && nextEmoji) { - tracker.runningEmoji = nextEmoji; - countPresented = - (await persistTrackerRunningEmoji(api, tracker)) && - (await setReaction(api, tracker, nextEmoji)); - tracker.runningEmojiConfirmed = countPresented; - } - const outcomePresented = - outcome === "ok" || - !reactionsEnabled || - (reactionsCleared && countPresented && (await setReaction(api, tracker, FAILURE_EMOJI))); - if (!reactionsCleared || !countPresented || !outcomePresented) { - return { ok: false }; - } - return { - ok: true, - activeRunIds, - reactionsEnabled, - typingEnabled, - ...(reactionsEnabled && nextEmoji ? { runningEmoji: nextEmoji } : {}), }; +}; + +type RecoveryRetry = { attempts: number; timer?: ReturnType }; + +const recoveryRetries = new Map(); + +function logFailure(api: ProgressCleanupApi, action: string, error: unknown) { + const message = error instanceof Error ? error.message : String(error); + api.logger.debug?.(`discord retired subagent progress ${action} failed: ${message}`); } -function scheduleTerminalLookupRetry( - api: ProgressApi, - event: Extract, - owned?: PersistedProgressRun, -) { - const runId = event.runId.trim(); - if (!runId || terminalRetryTimers.has(runId)) { - return; - } - if (!owned) { - const target = resolveDiscordProgressTarget(event.requester); - const account = resolveDiscordAccount({ - cfg: api.config, - accountId: event.requester?.accountId, - }); - if (!target || !isSubagentProgressEnabled(account)) { - return; +function reservedReactionEmojis(api: ProgressCleanupApi, accountAckReaction?: string): Set { + const reserved = new Set(Object.values(DEFAULT_EMOJIS)); + for (const emoji of [api.config.messages?.ackReaction, accountAckReaction]) { + if (emoji?.trim()) { + reserved.add(emoji.trim()); } } - if (terminalRetryTimers.size >= MAX_TRACKED_RUNS) { - return; + for (const agent of api.config.agents?.list ?? []) { + const emoji = agent.identity?.emoji?.trim(); + if (emoji) { + reserved.add(emoji); + } } - const expiresAt = terminalRetryExpiresAt.get(runId) ?? Date.now() + PROGRESS_STORE_TTL_MS; - const attempts = terminalRetryAttempts.get(runId) ?? 0; - if (expiresAt <= Date.now() || attempts >= TERMINAL_RETRY_MAX_ATTEMPTS) { - clearTerminalRetry(runId); - return; - } - terminalRetryExpiresAt.set(runId, expiresAt); - terminalRetryAttempts.set(runId, attempts + 1); - const retryDelayMs = Math.min( - TERMINAL_LOOKUP_RETRY_MS * 2 ** Math.min(attempts, 12), - TERMINAL_RETRY_MAX_DELAY_MS, - ); - const timer = setTimeout(() => { - terminalRetryTimers.delete(runId); - void handleEnded(api, event, owned); - }, retryDelayMs); - timer.unref?.(); - terminalRetryTimers.set(runId, timer); + return reserved; } -async function handleEnded( - api: ProgressApi, - event: Extract, - persistedHint?: PersistedProgressRun, -) { - const runId = event.runId.trim(); - if (!runId) { - return; - } - markRunTerminal(runId, event.outcome); - const lookup = await lookupProgressRun(api, runId); - const persisted = lookup.status === "found" ? lookup.value : persistedHint; - const key = trackerKeyByRunId.get(runId) ?? persisted?.key; - if (!key) { - if (lookup.status === "error") { - scheduleTerminalLookupRetry(api, event); - } else { - clearTerminalRetry(runId); - } - return; - } - cancelTerminalRetryTimer(runId); - await runQueued(key, async () => { - const tracker = trackers.get(key); - const currentLookup = await lookupProgressRun(api, runId); - const currentPersisted = - currentLookup.status === "found" - ? currentLookup.value - : currentLookup.status === "error" - ? (persistedHint ?? persisted) - : (persistedHint ?? (tracker ? undefined : persisted)); - const outcome = - persistedTerminalOutcome(currentPersisted) ?? - persistedTerminalOutcome(persisted) ?? - event.outcome; - const retryEvent = outcome === event.outcome ? event : { ...event, outcome }; - trackerKeyByRunId.delete(runId); - const owned = - tracker?.persistedRunIds.has(runId) && currentPersisted?.status !== "cleanup" - ? persistedProgressRunFromTracker(tracker) - : currentPersisted; - const cleanupMarked = owned - ? owned.status === "cleanup" || (await markProgressRunForCleanup(api, runId, owned, outcome)) - : true; - if (tracker) { - const currentAccount = resolveDiscordAccount({ - cfg: api.config, - accountId: tracker.accountId, - }); - if (!isSubagentProgressEnabled(currentAccount)) { - tracker.reactionsEnabled = false; - stopTyping(tracker); - } else { - await disableTrackerReactionsOnCollision(api, tracker, currentAccount.config.ackReaction); - } - } - if (!tracker) { - const reconciliation = owned - ? await reconcilePersistedTracker(api, owned, outcome, runId) - : { ok: false as const }; - if (reconciliation.ok && owned) { - const consumed = await consumeProgressRun(api, runId); - if (!consumed) { - scheduleTerminalLookupRetry(api, retryEvent, owned); - } else { - clearTerminalRetry(runId); - } - if (reconciliation.typingEnabled && reconciliation.activeRunIds.length > 0) { - const restoredTracker: ProgressTracker = { - accountId: owned.accountId, - channelId: owned.channelId, - messageId: owned.messageId, - activeRunIds: new Set(reconciliation.activeRunIds), - persistedRunIds: new Set(reconciliation.activeRunIds), - runningEmojiConfirmed: Boolean(reconciliation.runningEmoji), - reactionsEnabled: reconciliation.reactionsEnabled, - ...(reconciliation.runningEmoji ? { runningEmoji: reconciliation.runningEmoji } : {}), - typingExpiresAt: 0, - }; - trackers.set(key, restoredTracker); - for (const activeRunId of reconciliation.activeRunIds) { - trackerKeyByRunId.set(activeRunId, key); - } - startTyping(api, restoredTracker); - } - } else if (owned) { - scheduleTerminalLookupRetry(api, retryEvent, owned); - } - return; - } - tracker.activeRunIds.delete(runId); - tracker.persistedRunIds.delete(runId); - const countReconciled = tracker.reactionsEnabled - ? await updateRunningReaction(api, tracker) - : owned - ? (await reconcilePersistedTracker(api, owned, outcome, runId)).ok - : true; - const outcomePresented = - outcome === "ok" || - !tracker.reactionsEnabled || - (countReconciled && (await setReaction(api, tracker, FAILURE_EMOJI))); - const reconciled = countReconciled && outcomePresented; - if (reconciled && owned) { - const consumed = await consumeProgressRun(api, runId); - if (!consumed) { - scheduleTerminalLookupRetry(api, retryEvent, owned); - } else { - clearTerminalRetry(runId); - } - if (!consumed && !cleanupMarked) { - await markProgressRunForCleanup(api, runId, owned, outcome); - } - } else if (owned) { - scheduleTerminalLookupRetry(api, retryEvent, owned); - } - if (tracker.activeRunIds.size === 0) { - stopTyping(tracker); - trackers.delete(key); - } - }); -} - -async function handleDiscordSubagentProgressImpl(api: ProgressApi, event: SubagentProgressEvent) { - if (event.phase === "started") { - await handleStarted(api, event); - return; - } - await handleEnded(api, event); -} - -function clearStartupRecoveryRetry(api: ProgressApi) { - const retry = startupRecoveryRetries.get(api); +function clearRecoveryRetry(api: ProgressCleanupApi) { + const retry = recoveryRetries.get(api); if (retry?.timer) { clearTimeout(retry.timer); } - startupRecoveryRetries.delete(api); + recoveryRetries.delete(api); } -function scheduleStartupRecoveryRetry(api: ProgressApi) { - const retry = startupRecoveryRetries.get(api) ?? { attempts: 0 }; - if (retry.timer || retry.attempts >= STARTUP_RETRY_MAX_ATTEMPTS) { +function scheduleRecoveryRetry(api: ProgressCleanupApi) { + const retry = recoveryRetries.get(api) ?? { attempts: 0 }; + if (retry.timer || retry.attempts >= RETRY_MAX_ATTEMPTS) { return; } - const delayMs = Math.min( - TERMINAL_LOOKUP_RETRY_MS * 2 ** retry.attempts, - TERMINAL_RETRY_MAX_DELAY_MS, - ); + // A later gateway restart starts a fresh bounded recovery window for retained rows. + const delayMs = Math.min(RETRY_BASE_DELAY_MS * 2 ** retry.attempts, RETRY_MAX_DELAY_MS); retry.attempts += 1; retry.timer = setTimeout(() => { retry.timer = undefined; void recoverDiscordSubagentProgress(api); }, delayMs); retry.timer.unref?.(); - startupRecoveryRetries.set(api, retry); + recoveryRetries.set(api, retry); } -export async function recoverDiscordSubagentProgress(api: ProgressApi) { - const store = getProgressStore(api); - if (!store) { - if (api.runtime?.state) { - scheduleStartupRecoveryRetry(api); +async function cleanPersistedReaction( + api: ProgressCleanupApi, + store: PluginStateKeyedStore, + entry: PluginStateEntry, +): Promise { + try { + // Removed named accounts must not inherit today's root token and mutate another bot's reaction. + if ( + entry.value.accountId !== DEFAULT_ACCOUNT_ID && + !resolveDiscordAccountConfig(api.config, entry.value.accountId) + ) { + return false; } + const account = resolveDiscordAccount({ cfg: api.config, accountId: entry.value.accountId }); + if (!account.enabled || account.tokenStatus !== "available") { + return false; + } + const emoji = entry.value.runningEmoji; + // Persisted state can only authorize removal of glyphs the retired feature owned. + if ( + emoji && + HISTORICAL_RUNNING_EMOJIS.has(emoji) && + !reservedReactionEmojis(api, account.config.ackReaction).has(emoji) + ) { + await removeReactionDiscord(entry.value.channelId, entry.value.messageId, emoji, { + cfg: api.config, + accountId: entry.value.accountId, + }); + } + // Persisted ownership must survive until the external reaction is gone. + await store.consume(entry.key); + return true; + } catch (error) { + logFailure(api, "startup cleanup", error); + return false; + } +} + +async function recoverDiscordSubagentProgressImpl(api: ProgressCleanupApi) { + let store: PluginStateKeyedStore; + try { + store = api.runtime.state.openKeyedStore({ + namespace: "subagent-progress", + maxEntries: MAX_TRACKED_RUNS, + overflowPolicy: "reject-new", + defaultTtlMs: PROGRESS_STORE_TTL_MS, + }); + } catch (error) { + logFailure(api, "state store open", error); + scheduleRecoveryRetry(api); return; } - let persistedRuns: Array<{ key: string; value: PersistedProgressRun }>; + + let entries: PluginStateEntry[]; try { - persistedRuns = await store.entries(); + entries = await store.entries(); } catch (error) { logFailure(api, "startup recovery list", error); - scheduleStartupRecoveryRetry(api); + scheduleRecoveryRetry(api); return; } - clearStartupRecoveryRetry(api); - // Subagents share the gateway process, so no active run survives a cold - // start. Replaying every row repairs both interrupted and pending cleanup. - for (const entry of persistedRuns) { - await handleEnded( - api, - { - phase: "ended", - runId: entry.key, - outcome: persistedTerminalOutcome(entry.value) ?? "unknown", - }, - entry.value, - ); + + let retryNeeded = false; + for (const entry of entries) { + if (!(await cleanPersistedReaction(api, store, entry))) { + retryNeeded = true; + } + } + if (retryNeeded) { + scheduleRecoveryRetry(api); + } else { + clearRecoveryRetry(api); } } function resetDiscordSubagentProgressForTest() { - for (const tracker of trackers.values()) { - stopTyping(tracker); - } - trackers.clear(); - trackerKeyByRunId.clear(); - for (const timer of terminalRetryTimers.values()) { - clearTimeout(timer); - } - terminalRetryTimers.clear(); - terminalRetryExpiresAt.clear(); - terminalRetryAttempts.clear(); - for (const retry of startupRecoveryRetries.values()) { + for (const [api, retry] of recoveryRetries) { if (retry.timer) { clearTimeout(retry.timer); } + recoveryRetries.delete(api); } - startupRecoveryRetries.clear(); - resetDiscordSubagentProgressStateForTest(); } -export const handleDiscordSubagentProgress = Object.assign(handleDiscordSubagentProgressImpl, { +export const recoverDiscordSubagentProgress = Object.assign(recoverDiscordSubagentProgressImpl, { resetForTest: resetDiscordSubagentProgressForTest, }); diff --git a/extensions/discord/subagent-hooks-api.ts b/extensions/discord/subagent-hooks-api.ts index b75fd49d5968..16e95a273f3e 100644 --- a/extensions/discord/subagent-hooks-api.ts +++ b/extensions/discord/subagent-hooks-api.ts @@ -16,10 +16,6 @@ export function registerDiscordSubagentHooks(api: OpenClawPluginApi): void { const { recoverDiscordSubagentProgress } = await loadDiscordSubagentProgressModule(); await recoverDiscordSubagentProgress(api); }); - api.on("subagent_progress", async (event) => { - const { handleDiscordSubagentProgress } = await loadDiscordSubagentProgressModule(); - await handleDiscordSubagentProgress(api, event); - }); api.on("subagent_ended", async (event) => { const { handleDiscordSubagentEnded } = await loadDiscordSubagentHooksModule(); handleDiscordSubagentEnded(event); diff --git a/src/plugins/contracts/boundary-invariants.test.ts b/src/plugins/contracts/boundary-invariants.test.ts index e19e812111f9..75fc5153a160 100644 --- a/src/plugins/contracts/boundary-invariants.test.ts +++ b/src/plugins/contracts/boundary-invariants.test.ts @@ -38,7 +38,6 @@ const BUNDLED_TYPED_HOOK_REGISTRATION_GUARDS = { "gateway_start", "subagent_delivery_target", "subagent_ended", - "subagent_progress", ], "extensions/feishu/subagent-hooks-api.ts": ["subagent_delivery_target", "subagent_ended"], "extensions/matrix/subagent-hooks-api.ts": ["subagent_delivery_target", "subagent_ended"],