mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
fix(imessage): recover oldest messages missed during downtime (#147395)
Supersedes #139693 by @zhangguiping-xydt. Credit to @yetval for #139669 and the earlier repair. ## What Problem This Solves With iMessage catchup enabled, the oldest messages sent during downtime can be silently skipped. A 100-message backlog with a replay limit of 50 previously recovered messages 51–100 and saved cursor 100. The next startup recovered none of messages 1–50. ## Why This Change Was Made Request the existing 500-row history window per chat before selecting the oldest rows across chats. The contributor's production fix is unchanged. The two regressions now enter through registered iMessage account startup and check the real durable ingress and persisted cursor after each start. ## User Impact Backlogs within the existing history window recover oldest first while the per-startup replay limit stays unchanged. The existing 500-row per-chat ceiling and age window still apply. This does not guarantee recovery of arbitrarily large backlogs, automatically repeat catchup passes, or rewind a cursor that already skipped older messages. ## Consumers - Registered account startup calls the production monitor, which sends catchup rows through the same durable queue as live messages. - Cursor persistence keeps its existing owner. Live advances wait during startup and resume only after a successful complete catchup pass. A partial pass keeps the cursor at its replay boundary. - Durable message records retain catchup provenance and duplicate protection. The outgoing-message observer and self-chat duplicate checks remain active. - Doctor imports and the alternate recovery cursor keep their existing rules. Enabled catchup prevents the alternate path from consuming its cursor. - Recent direct-message context, conversation repair, and reaction polling use separate history requests and unchanged limits. No other channel reads the changed private budget. - No configuration key, stored format, migration, protocol, or SDK contract changes. ## Evidence - Both replacement startup regressions fail on main: SQLite contains rows 51–100 instead of 1–50 after the first start. Both pass with the fix. - The one-chat case checks two persisted starts through 100 rows. The two-chat case checks twelve starts through 600 rows, with exact oldest-first admission, completed status, cursor row/time, and the 500-row history request after every pass. - A fresh merge with main passed 419 tests across six startup, catchup, monitor, ingress, and self-chat files. The test-routing suite passed all 602 tests. Changed-file lint and formatting passed. - Retained full-Gateway proof completed 50/100 synthetic requests on main and 100/100 with the production fix, with cursors 50 then 100 on the fixed starts and 100 status reply requests total. The committed regressions prove durable admission/completion and suppress unrelated agent replies through normal direct-message policy. The separate full-Gateway proof used a synthetic external bridge. Physical-device delivery was not tested. No local typecheck or full package build ran. Co-authored-by: zhang-guiping <zhang.guiping@xydigit.com> Co-authored-by: Ayaan Zaidi <hi@obviy.us>
This commit is contained in:
parent
db5754328a
commit
be8f4ebff9
2 changed files with 182 additions and 7 deletions
|
|
@ -13,11 +13,12 @@ import {
|
|||
import { parseIMessageNotification } from "./parse-notification.js";
|
||||
import type { IMessagePayload } from "./types.js";
|
||||
|
||||
// Per-chat history fetch budget. messages.history is per-chat; we cap each
|
||||
// chat's fetch to the global perRunLimit so a single noisy group cannot
|
||||
// dominate the cursor advance — the cross-chat sort + final slice still
|
||||
// caps the global pass at perRunLimit.
|
||||
const PER_CHAT_HISTORY_LIMIT_CAP = 500;
|
||||
// Per-chat history fetch budget. Upstream `messages.history` serves rows
|
||||
// `ORDER BY date DESC LIMIT ?`, so a smaller limit trims the OLDEST rows
|
||||
// server-side before we ever see them. Always request the full budget: the
|
||||
// cross-chat sort plus the perRunLimit slice below can only pick the true
|
||||
// oldest rows if the per-chat page reached them.
|
||||
const PER_CHAT_HISTORY_LIMIT = 500;
|
||||
|
||||
// chats.list page size used during catchup. 200 covers far more than any
|
||||
// realistic offline window worth of distinct chats while staying well under
|
||||
|
|
@ -113,7 +114,6 @@ export async function runIMessageCatchup(
|
|||
}
|
||||
const chats = chatsResult?.chats ?? [];
|
||||
const collected: IMessageCatchupRow[] = [];
|
||||
const perChatLimit = Math.min(limit, PER_CHAT_HISTORY_LIMIT_CAP);
|
||||
let historyFetchFailed = false;
|
||||
// Track the highest rowid / date the imsg bridge actually returned across
|
||||
// all chats, regardless of whether each row passed the parser. The catchup
|
||||
|
|
@ -143,7 +143,7 @@ export async function runIMessageCatchup(
|
|||
"messages.history",
|
||||
{
|
||||
chat_id: chatId,
|
||||
limit: perChatLimit,
|
||||
limit: PER_CHAT_HISTORY_LIMIT,
|
||||
start: sinceISO,
|
||||
attachments: includeAttachments,
|
||||
},
|
||||
|
|
|
|||
175
extensions/imessage/src/monitor/catchup-startup.test.ts
Normal file
175
extensions/imessage/src/monitor/catchup-startup.test.ts
Normal file
|
|
@ -0,0 +1,175 @@
|
|||
import fs from "node:fs";
|
||||
import path from "node:path";
|
||||
import { DatabaseSync } from "node:sqlite";
|
||||
import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts";
|
||||
import { closeOpenClawStateDatabaseForTest } from "openclaw/plugin-sdk/plugin-state-test-runtime";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { resolveIMessageAccount } from "../accounts.js";
|
||||
import { imessagePlugin } from "../channel.js";
|
||||
import { IMessageRpcClient, type createIMessageRpcClient } from "../client.js";
|
||||
import { getIMessageRuntime } from "../runtime.js";
|
||||
import {
|
||||
IMESSAGE_CATCHUP_CURSOR_MAX_ENTRIES,
|
||||
IMESSAGE_CATCHUP_CURSOR_NAMESPACE,
|
||||
resolveIMessageCatchupCursorKey,
|
||||
type IMessageCatchupCursor,
|
||||
} from "../state-contract.js";
|
||||
import { installIMessageStateRuntimeForTest } from "../test-support/runtime.js";
|
||||
|
||||
const createClient = vi.hoisted(() => vi.fn<typeof createIMessageRpcClient>());
|
||||
|
||||
vi.mock("../client.js", async (importOriginal) => ({
|
||||
...(await importOriginal<typeof import("../client.js")>()),
|
||||
createIMessageRpcClient: createClient,
|
||||
}));
|
||||
|
||||
vi.mock("../probe.js", async (importOriginal) => ({
|
||||
...(await importOriginal<typeof import("../probe.js")>()),
|
||||
probeIMessage: vi.fn(async () => ({ ok: true })),
|
||||
}));
|
||||
|
||||
describe("registered iMessage account startup catchup", () => {
|
||||
let stateDir: string;
|
||||
|
||||
beforeEach(() => {
|
||||
installIMessageStateRuntimeForTest();
|
||||
stateDir = getIMessageRuntime().state.resolveStateDir();
|
||||
vi.stubEnv("OPENCLAW_STATE_DIR", stateDir);
|
||||
createClient.mockReset();
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
vi.restoreAllMocks();
|
||||
vi.unstubAllEnvs();
|
||||
closeOpenClawStateDatabaseForTest();
|
||||
fs.rmSync(stateDir, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
it.each([
|
||||
{ name: "recovers the oldest rows from a 100-message chat", total: 100, twoChats: false },
|
||||
{ name: "caps each pass when a 500-message chat dominates", total: 600, twoChats: true },
|
||||
])(
|
||||
"imessagePlugin.gateway.startAccount $name across persisted starts",
|
||||
async ({ total, twoChats }) => {
|
||||
const backlogStartMs = Date.now() - 20 * 60_000;
|
||||
const rows = Array.from({ length: total }, (_, index) => ({
|
||||
id: index + 1,
|
||||
guid: `catchup-startup-${index + 1}`,
|
||||
chat_id: twoChats && index < 100 ? 2 : 1,
|
||||
chat_identifier: "+15555550123",
|
||||
chat_guid: "iMessage;-;+15555550123",
|
||||
sender: "+15555550123",
|
||||
is_from_me: false,
|
||||
is_group: false,
|
||||
text: "missed during downtime",
|
||||
created_at: new Date(backlogStartMs + (index + 1) * 1_000).toISOString(),
|
||||
}));
|
||||
const cfg: OpenClawConfig = {
|
||||
channels: {
|
||||
imessage: {
|
||||
cliPath: path.join(stateDir, "synthetic-imsg"),
|
||||
// Admission is the boundary under test; ordinary policy completes the rows
|
||||
// without starting unrelated agent turns.
|
||||
dmPolicy: "disabled",
|
||||
catchup: { enabled: true, perRunLimit: 50, maxAgeMinutes: 60 },
|
||||
},
|
||||
},
|
||||
};
|
||||
const account = resolveIMessageAccount({ cfg, accountId: "default" });
|
||||
const cursorStore = getIMessageRuntime().state.openSyncKeyedStore<IMessageCatchupCursor>({
|
||||
namespace: IMESSAGE_CATCHUP_CURSOR_NAMESPACE,
|
||||
maxEntries: IMESSAGE_CATCHUP_CURSOR_MAX_ENTRIES,
|
||||
});
|
||||
const cursorKey = resolveIMessageCatchupCursorKey(account.accountId);
|
||||
cursorStore.register(cursorKey, {
|
||||
lastSeenMs: backlogStartMs,
|
||||
lastSeenRowid: 0,
|
||||
updatedAt: Date.now(),
|
||||
});
|
||||
const database = new DatabaseSync(path.join(stateDir, "state", "openclaw.sqlite"), {
|
||||
readOnly: true,
|
||||
});
|
||||
try {
|
||||
for (let end = 50; end <= total; end += 50) {
|
||||
const client = new IMessageRpcClient({ cliPath: cfg.channels?.imessage?.cliPath });
|
||||
const request = vi.spyOn(client, "request").mockImplementation(async (method, params) => {
|
||||
if (method === "watch.subscribe") {
|
||||
return { subscription: 1 };
|
||||
}
|
||||
if (method === "chats.list") {
|
||||
return {
|
||||
chats: (twoChats ? [1, 2] : [1]).map((id) => ({
|
||||
id,
|
||||
last_message_at: rows.at(-1)!.created_at,
|
||||
})),
|
||||
};
|
||||
}
|
||||
if (method === "messages.history") {
|
||||
const { chat_id, limit, start } = params!;
|
||||
if (typeof limit !== "number" || typeof start !== "string") {
|
||||
throw new Error("history request must carry its limit and lower date bound");
|
||||
}
|
||||
// Match the external bridge: filter by date, then select newest first.
|
||||
return {
|
||||
messages: rows
|
||||
.filter(
|
||||
(row) =>
|
||||
row.chat_id === chat_id && Date.parse(row.created_at) >= Date.parse(start),
|
||||
)
|
||||
.toSorted(
|
||||
(a, b) => Date.parse(b.created_at) - Date.parse(a.created_at) || b.id - a.id,
|
||||
)
|
||||
.slice(0, limit),
|
||||
};
|
||||
}
|
||||
throw new Error(`unexpected bridge method ${method}`);
|
||||
});
|
||||
// The monitor calls this after catchup finishes, then drains real ingress
|
||||
// during shutdown before the next registered account start.
|
||||
vi.spyOn(client, "waitForClose").mockResolvedValue(undefined);
|
||||
createClient.mockResolvedValue(client);
|
||||
const runtime = {
|
||||
log: vi.fn(),
|
||||
error: vi.fn(),
|
||||
exit: vi.fn((code: number) => {
|
||||
throw new Error(`unexpected exit ${code}`);
|
||||
}),
|
||||
};
|
||||
await imessagePlugin.gateway!.startAccount!({
|
||||
cfg,
|
||||
accountId: account.accountId,
|
||||
account,
|
||||
runtime,
|
||||
abortSignal: new AbortController().signal,
|
||||
getStatus: () => ({ accountId: account.accountId }),
|
||||
setStatus: vi.fn(),
|
||||
});
|
||||
|
||||
const admitted = database
|
||||
.prepare(
|
||||
"SELECT event_id, status FROM channel_ingress_events WHERE channel_id = 'imessage' AND account_id = 'default' ORDER BY rowid",
|
||||
)
|
||||
.all();
|
||||
expect(admitted.map((row) => row.event_id)).toEqual(
|
||||
Array.from({ length: end }, (_, index) => `catchup-startup-${index + 1}`),
|
||||
);
|
||||
expect(admitted.every((row) => row.status === "completed")).toBe(true);
|
||||
expect(cursorStore.lookup(cursorKey)).toMatchObject({
|
||||
lastSeenRowid: end,
|
||||
lastSeenMs: backlogStartMs + end * 1_000,
|
||||
});
|
||||
const historyCalls = request.mock.calls.filter(
|
||||
([method]) => method === "messages.history",
|
||||
);
|
||||
expect(historyCalls).toHaveLength(twoChats ? 2 : 1);
|
||||
for (const [, params] of historyCalls) {
|
||||
expect(params).toMatchObject({ limit: 500 });
|
||||
}
|
||||
expect(runtime.error).not.toHaveBeenCalled();
|
||||
}
|
||||
} finally {
|
||||
database.close();
|
||||
}
|
||||
},
|
||||
);
|
||||
});
|
||||
Loading…
Add table
Add a link
Reference in a new issue