diff --git a/examples/memory-plugin-shared/sync.mjs b/examples/memory-plugin-shared/sync.mjs index 72d46423b..3d2c58534 100644 --- a/examples/memory-plugin-shared/sync.mjs +++ b/examples/memory-plugin-shared/sync.mjs @@ -52,7 +52,7 @@ const OPENCODE_SHARED_FILES = [...HOOK_SHARED_FILES, ...SETUP_WIZARD_SHARED_FILE // dsh and zcode ship no setup entry point, so nothing there calls the wizard. const ZCODE_SHARED_FILES = [...HOOK_SHARED_FILES, ...MCP_PROXY_SHARED_FILES, ...BATCH_SHARED_FILES, ...ASYNC_WRITE_SHARED_FILES, "agent-hook-runtime.mjs", "agent-uri-guard.mjs"]; const DSH_SHARED_FILES = [...HOOK_SHARED_FILES, ...MCP_PROXY_SHARED_FILES]; -const PI_SHARED_FILES = [...HOOK_SHARED_FILES, ...SETUP_WIZARD_SHARED_FILES]; +const PI_SHARED_FILES = [...HOOK_SHARED_FILES, ...SETUP_WIZARD_SHARED_FILES, ...BATCH_SHARED_FILES]; // Agent Plugins 1.0 has no hooks: it is the proxy and nothing else. const AGENT_PLUGINS_SHARED_FILES = ["credentials.mjs", "debug-log.mjs", ...MCP_PROXY_SHARED_FILES]; // openclaw assembles recall server-side, so it takes the recall pair alone. diff --git a/examples/pi-coding-agent-extension/shared/batch-send.d.mts b/examples/pi-coding-agent-extension/shared/batch-send.d.mts new file mode 100644 index 000000000..3e7bb6739 --- /dev/null +++ b/examples/pi-coding-agent-extension/shared/batch-send.d.mts @@ -0,0 +1,15 @@ +export const BATCH_LIMIT: number; +export function sendSessionMessages( + fetchJSON: (path: string, init?: any) => Promise<{ ok: boolean; status?: number; result?: any; error?: any }>, + sessionId: string, + payloads: Array>, + opts?: { enqueueOnRetryable?: boolean; onSent?: (count: number) => void | Promise }, +): Promise<{ + sent: number; + queued: number; + enqueueFailed: number; + failed: number; + retryable: boolean; + usedBatch: boolean; + lastError: any; +}>; diff --git a/examples/pi-coding-agent-extension/shared/batch-send.mjs b/examples/pi-coding-agent-extension/shared/batch-send.mjs new file mode 100644 index 000000000..d7ec355cc --- /dev/null +++ b/examples/pi-coding-agent-extension/shared/batch-send.mjs @@ -0,0 +1,121 @@ +// GENERATED FROM examples/memory-plugin-shared/lib. DO NOT EDIT. +import { enqueue } from "./pending-queue.mjs"; +import { isRetryableFailure } from "./retryable.mjs"; + +export const BATCH_LIMIT = 100; + +export function isRetryableSendFailure(res) { + return isRetryableFailure(res); +} + +function makeResult() { + return { + sent: 0, + queued: 0, + enqueueFailed: 0, + failed: 0, + retryable: false, + usedBatch: true, + lastError: null, + }; +} + +async function enqueueRemainder(sessionId, payloads, startIndex, result, retryable, lastError) { + result.retryable = retryable; + result.lastError = lastError ?? null; + + if (!retryable) { + result.failed += Math.max(0, payloads.length - startIndex); + return result; + } + + const baseCreatedAt = Date.now(); + for (let i = startIndex; i < payloads.length; i++) { + const queued = await enqueue("addMessage", sessionId, payloads[i], { + createdAt: baseCreatedAt + (i - startIndex), + }); + if (!queued.ok) { + // Stop at the first enqueue failure so queued entries stay a contiguous + // prefix: consumers mark the first sent+queued payloads as captured, so + // skipping a payload here and queueing a later one would silently drop it. + result.enqueueFailed += payloads.length - i; + return result; + } + result.queued++; + } + return result; +} + +async function sendSerial(fetchJSON, sessionId, payloads, startIndex, opts, result) { + result.usedBatch = false; + const encodedSid = encodeURIComponent(sessionId); + for (let i = startIndex; i < payloads.length; i++) { + const res = await fetchJSON(`/api/v1/sessions/${encodedSid}/messages`, { + method: "POST", + body: JSON.stringify(payloads[i]), + }); + if (res?.ok) { + result.sent++; + await opts.onSent?.(1); + continue; + } + + const retryable = isRetryableSendFailure(res); + if (opts.enqueueOnRetryable) { + return enqueueRemainder(sessionId, payloads, i, result, retryable, res?.error ?? res); + } + result.retryable = retryable; + result.lastError = res?.error ?? res ?? null; + result.failed += payloads.length - i; + return result; + } + return result; +} + +/** + * Send add-message payloads in server-sized batches with serial fallback. + * + * @param {Function} fetchJSON - (path, init) => { ok, status, result?, error? } + * @param {string} sessionId - OpenViking session id + * @param {Array} payloads - sanitized add-message request bodies + * @param {object} opts + * @param {boolean} opts.enqueueOnRetryable - enqueue unsent payloads after a retryable failure + * @param {Function} opts.onSent - called with the number of messages durably sent after each success + * @returns {Promise<{sent:number,queued:number,enqueueFailed:number,failed:number,retryable:boolean,usedBatch:boolean,lastError:any}>} + */ +export async function sendSessionMessages(fetchJSON, sessionId, payloads, opts = {}) { + const result = makeResult(); + const messages = Array.isArray(payloads) ? payloads : []; + if (messages.length === 0) return result; + + const encodedSid = encodeURIComponent(sessionId); + for (let start = 0; start < messages.length; start += BATCH_LIMIT) { + const chunk = messages.slice(start, start + BATCH_LIMIT); + const res = await fetchJSON(`/api/v1/sessions/${encodedSid}/messages/batch`, { + method: "POST", + body: JSON.stringify({ messages: chunk }), + }); + + if (res?.ok) { + result.sent += chunk.length; + await opts.onSent?.(chunk.length); + continue; + } + + const status = Number(res?.status || 0); + if (status === 404 || status === 405) { + return sendSerial(fetchJSON, sessionId, messages, start, opts, result); + } + + const retryable = isRetryableSendFailure(res); + if (opts.enqueueOnRetryable) { + return enqueueRemainder(sessionId, messages, start, result, retryable, res?.error ?? res); + } + result.retryable = retryable; + result.lastError = res?.error ?? res ?? null; + result.failed += messages.length - start; + return result; + } + + return result; +} diff --git a/examples/pi-coding-agent-extension/shared/pending-queue.d.mts b/examples/pi-coding-agent-extension/shared/pending-queue.d.mts index 3d898ce40..9051bdcdb 100644 --- a/examples/pi-coding-agent-extension/shared/pending-queue.d.mts +++ b/examples/pi-coding-agent-extension/shared/pending-queue.d.mts @@ -4,4 +4,7 @@ export function replayPending( fetchJSON: (path: string, init?: any) => Promise<{ ok: boolean; status?: number; result?: any; error?: any }>, log: (stage: string, data?: any) => void, ): Promise<{ replayed: number; failed: number; skipped: number; deferred: number }>; +export function claimForReplay(filename: string): Promise; +export function dequeue(filename: string): Promise; +export function incrementRetry(filename: string, entry: Record): Promise; export function cleanStale(): Promise; diff --git a/examples/pi-coding-agent-extension/sync.ts b/examples/pi-coding-agent-extension/sync.ts index fc754ea07..50a66b958 100644 --- a/examples/pi-coding-agent-extension/sync.ts +++ b/examples/pi-coding-agent-extension/sync.ts @@ -2,7 +2,8 @@ import type { OVClient } from "./client.js"; import { createLogger } from "./shared/debug-log.mjs"; import type { OVConfig } from "./config.js"; import { deriveHarnessSessionId } from "./shared/session-model.mjs"; -import { enqueue, listPending, replayPending } from "./shared/pending-queue.mjs"; +import { claimForReplay, dequeue, enqueue, incrementRetry, listPending, replayPending } from "./shared/pending-queue.mjs"; +import { BATCH_LIMIT, sendSessionMessages } from "./shared/batch-send.mjs"; import { extractBranchCapturePayloads } from "./lib/capture-adapter.mjs"; import { countUndeliveredForSession, estimatePayloadTokens } from "./lib/takeover-core.mjs"; @@ -61,26 +62,101 @@ export class SyncManager { async flushForTakeover(): Promise { if (!this.ovSessionId) return false; - await this.replayPending(); + // Drain is bounded (time / max batches). Remaining undelivered entries — + // including any still claimed as `.processing` — keep the barrier closed + // until a later turn finishes draining them. + if (this.client.connected) await this.drainSessionBacklog(); const pending = await listPending(); return countUndeliveredForSession(pending, this.ovSessionId) === 0; } + /** + * Replay this session's queued addMessage entries through the batch + * endpoint, BATCH_LIMIT per request. The shared replayPending() sends one + * request per entry and stops after one replay window, which is what let a + * large offline backlog block the takeover barrier for many turns (#4504). + * Entries are claimed one batch at a time so a failed batch only costs a + * retry for the entries it contained. + * + * Bounded per call via OPENVIKING_PENDING_DRAIN_BUDGET_MS (default 60s) and + * optional OPENVIKING_PENDING_DRAIN_MAX_BATCHES so a huge backlog cannot + * block turn_end for an unbounded wall time; remainder drains on later turns. + */ + private async drainSessionBacklog(): Promise { + const sid = this.ovSessionId; + if (!sid) return; + const budgetRaw = Number(process.env.OPENVIKING_PENDING_DRAIN_BUDGET_MS); + const timeBudgetMs = Number.isFinite(budgetRaw) && budgetRaw >= 0 ? budgetRaw : 60_000; + const maxRaw = Number(process.env.OPENVIKING_PENDING_DRAIN_MAX_BATCHES); + const maxBatches = + Number.isFinite(maxRaw) && maxRaw > 0 ? Math.floor(maxRaw) : Number.POSITIVE_INFINITY; + const started = Date.now(); + let batches = 0; + + const backlog = (await listPending()).filter( + ({ entry }) => entry?.type === "addMessage" && entry.sessionId === sid, + ); + for (let start = 0; start < backlog.length; start += BATCH_LIMIT) { + if (batches >= maxBatches || Date.now() - started >= timeBudgetMs) { + this.logger.log("drain", { + session: sid, + stopped: batches >= maxBatches ? "max-batches" : "time-budget", + batches, + remaining: backlog.length - start, + elapsedMs: Date.now() - started, + }); + return; + } + + const claimed: Array<{ filename: string; entry: any }> = []; + for (const { filename, entry } of backlog.slice(start, start + BATCH_LIMIT)) { + const name = await claimForReplay(filename); + if (name) claimed.push({ filename: name, entry }); + } + if (claimed.length === 0) continue; + batches += 1; + + let delivered = 0; + await sendSessionMessages( + this.fetchJSON, + sid, + claimed.map(({ entry }) => entry.payload), + { + onSent: async (count: number) => { + for (let i = 0; i < count; i++) await dequeue(claimed[delivered++].filename); + }, + }, + ); + if (delivered === claimed.length) continue; + + // Undelivered entries stay queued with one more retry; incrementRetry + // drops them once the retry budget is exhausted. + for (const { filename, entry } of claimed.slice(delivered)) { + await incrementRetry(filename, entry); + } + this.logger.log("drain", { + session: sid, + delivered, + retried: claimed.length - delivered, + }); + return; + } + } + + private fetchJSON = (path: string, init?: any) => this.client.fetchJSON(path, init, 10000); + async syncBranch(branch: any[]): Promise { if (!this.ovSessionId) return { added: 0, tokens: 0, allDelivered: true }; const extracted = extractBranchCapturePayloads(branch, this.syncedEntryCount, this.config); if (extracted.resetWatermark) this.syncedEntryCount = 0; - let added = 0; + const sent = await this.sendPayloads(extracted.payloads); + const added = sent.accepted; let tokens = 0; - let allDelivered = true; - for (const payload of extracted.payloads) { - const result = await this.addPayload(payload); - if (!result.accepted) break; - added++; + for (const payload of extracted.payloads.slice(0, added)) { tokens += estimatePayloadTokens(payload); - allDelivered = allDelivered && result.delivered; } + const allDelivered = sent.delivered === added; if (added === extracted.payloads.length) { this.syncedEntryCount = extracted.nextEntryCount; } @@ -91,11 +167,36 @@ export class SyncManager { } async addPayload(payload: any): Promise { - if (!this.ovSessionId) return { accepted: false, delivered: false }; - const ok = await this.client.addMessagePayload(this.ovSessionId, payload); - if (ok) return { accepted: true, delivered: true }; - await enqueue("addMessage", this.ovSessionId, payload); - return { accepted: true, delivered: false }; + const sent = await this.sendPayloads([payload]); + return { accepted: sent.accepted === 1, delivered: sent.delivered === 1 }; + } + + /** + * Send payloads in one batch request; retryable failures are queued to disk. + * Returns how many payloads were accepted (sent or queued, always a prefix) + * and how many of those were delivered to the server. + */ + private async sendPayloads(payloads: any[]): Promise<{ accepted: number; delivered: number }> { + if (!this.ovSessionId || payloads.length === 0) return { accepted: 0, delivered: 0 }; + const res = await sendSessionMessages(this.fetchJSON, this.ovSessionId, payloads, { + enqueueOnRetryable: true, + }); + if (res.failed > 0 || res.enqueueFailed > 0) { + this.logger.log("send", { + session: this.ovSessionId, + sent: res.sent, + queued: res.queued, + failed: res.failed, + enqueueFailed: res.enqueueFailed, + error: res.lastError?.message || res.lastError?.code || "unknown", + }); + } + // A non-retryable rejection (4xx) drops the remaining payloads, the same + // outcome replayPending() applies to such entries. Count them as accepted + // so the watermark still advances past them; otherwise the next turn would + // re-extract and re-send the payloads that were already delivered. + const dropped = res.retryable ? 0 : res.failed; + return { accepted: res.sent + res.queued + dropped, delivered: res.sent }; } async commitIfNeeded(): Promise { diff --git a/examples/pi-coding-agent-extension/tests/sync-barrier.test.mjs b/examples/pi-coding-agent-extension/tests/sync-barrier.test.mjs index b3a5db4ee..541aaa57f 100644 --- a/examples/pi-coding-agent-extension/tests/sync-barrier.test.mjs +++ b/examples/pi-coding-agent-extension/tests/sync-barrier.test.mjs @@ -179,9 +179,9 @@ test("restoreWatermark prevents pi -c from re-syncing already captured entries", await withPendingDir(async () => { const calls = []; const c = client({ - addMessagePayload: async (_sid, payload) => { - calls.push(payload); - return true; + fetchJSON: async (_path, init) => { + calls.push(...JSON.parse(init.body).messages); + return { ok: true, result: {} }; }, }); const sync = new SyncManager(c, config()); @@ -198,3 +198,118 @@ test("restoreWatermark prevents pi -c from re-syncing already captured entries", assert.match(calls[0].parts[0].text, /Fresh entry/); }); }); + +function batchClient(overrides = {}) { + const calls = []; + const c = client({ + fetchJSON: async (path, init) => { + calls.push({ path: String(path), body: init?.body ? JSON.parse(init.body) : null }); + return overrides.respond ? overrides.respond(calls.length, path) : { ok: true, result: {} }; + }, + }); + return { c, calls }; +} + +test("syncBranch sends the whole turn in one batch request", async () => { + await withPendingDir(async () => { + const { c, calls } = batchClient(); + const sync = new SyncManager(c, config()); + await sync.ensureSession("pi-session"); + + const result = await sync.syncBranch([ + { type: "message", message: { role: "user", content: "First user message for the batch write test." } }, + { type: "message", message: { role: "assistant", content: "Assistant reply for the batch write test." } }, + { type: "message", message: { role: "user", content: "Second user message for the batch write test." } }, + ]); + + assert.equal(result.added, 3); + assert.equal(result.allDelivered, true); + assert.equal(calls.length, 1); + assert.match(calls[0].path, /\/messages\/batch$/); + assert.equal(calls[0].body.messages.length, 3); + }); +}); + +test("takeover barrier drains a large backlog through the batch endpoint", async () => { + await withPendingDir(async () => { + const { c, calls } = batchClient(); + const sync = new SyncManager(c, config()); + await sync.ensureSession("pi-session"); + const t0 = Date.now(); + for (let i = 0; i < 250; i++) { + await enqueue("addMessage", sync.sessionId, { role: "user", content: `m${i}` }, { createdAt: t0 + i }); + } + await enqueue("addMessage", "other-session", { role: "user", content: "other" }, { createdAt: t0 + 999 }); + + assert.equal(await sync.flushForTakeover(), true); + assert.deepEqual(calls.map((call) => call.body.messages.length), [100, 100, 50]); + // Order preserved across batches. + assert.equal(calls[0].body.messages[0].content, "m0"); + assert.equal(calls[2].body.messages[49].content, "m249"); + const left = await listPending(); + assert.equal(left.length, 1); + assert.equal(left[0].entry.sessionId, "other-session"); + }); +}); + +test("failed batch keeps its entries queued with one retry and leaves the rest untouched", async () => { + await withPendingDir(async () => { + const { c, calls } = batchClient({ respond: () => ({ ok: false, status: 500 }) }); + const sync = new SyncManager(c, config()); + await sync.ensureSession("pi-session"); + const t0 = Date.now(); + for (let i = 0; i < 120; i++) { + await enqueue("addMessage", sync.sessionId, { role: "user", content: `m${i}` }, { createdAt: t0 + i }); + } + + assert.equal(await sync.flushForTakeover(), false); + assert.equal(calls.length, 1); + const retries = (await listPending()).map((p) => p.entry.retries); + assert.equal(retries.length, 120); + assert.equal(retries.filter((r) => r === 1).length, 100); + assert.equal(retries.filter((r) => r === 0).length, 20); + }); +}); + +test("non-retryable batch failure drops the payloads but still advances the sync watermark", async () => { + await withPendingDir(async () => { + const { c, calls } = batchClient({ respond: () => ({ ok: false, status: 400, error: { message: "bad" } }) }); + const sync = new SyncManager(c, config()); + await sync.ensureSession("pi-session"); + + const result = await sync.syncBranch([ + { type: "message", message: { role: "user", content: "Poison payload one for watermark test." } }, + { type: "message", message: { role: "user", content: "Poison payload two for watermark test." } }, + ]); + + assert.equal(result.added, 2); + assert.equal(result.allDelivered, false); + assert.equal(sync.syncedCount, 2); + assert.equal(calls.length, 1); + assert.equal((await listPending()).length, 0); + }); +}); + +test("drainSessionBacklog stops after maxBatches and leaves remainder for later turns", async () => { + await withPendingDir(async () => { + const previous = process.env.OPENVIKING_PENDING_DRAIN_MAX_BATCHES; + process.env.OPENVIKING_PENDING_DRAIN_MAX_BATCHES = "1"; + try { + const { c, calls } = batchClient(); + const sync = new SyncManager(c, config()); + await sync.ensureSession("pi-session"); + const t0 = Date.now(); + for (let i = 0; i < 250; i++) { + await enqueue("addMessage", sync.sessionId, { role: "user", content: `m${i}` }, { createdAt: t0 + i }); + } + + assert.equal(await sync.flushForTakeover(), false); + assert.equal(calls.length, 1); + assert.equal(calls[0].body.messages.length, 100); + assert.equal((await listPending()).length, 150); + } finally { + if (previous === undefined) delete process.env.OPENVIKING_PENDING_DRAIN_MAX_BATCHES; + else process.env.OPENVIKING_PENDING_DRAIN_MAX_BATCHES = previous; + } + }); +});