mirror of
https://github.com/volcengine/OpenViking.git
synced 2026-09-29 16:58:31 +08:00
fix(pi-extension): replay offline backlog through the batch endpoint (#4692)
* fix(pi-extension): replay offline backlog through the batch endpoint After a server outage pi's local pending queue can hold hundreds of addMessage entries. The takeover barrier requires that queue to be empty for the session, but replayed it with the shared replayPending(): one POST per entry, at most one replay window (50) per attempt. At the ~0.67s/entry measured in #4504 a 670-entry backlog needed 14 takeover attempts of ~33s each, and pendingTokens kept growing meanwhile. - flushForTakeover drains the session's backlog via /messages/batch (sendSessionMessages, BATCH_LIMIT=100), so the whole backlog clears in a handful of requests within one attempt. Entries are claimed one batch at a time; a failed batch costs a retry only for the entries it contained, the rest stay untouched. - syncBranch sends the whole turn in one batch request instead of one POST per payload, with retryable failures queued as before. - pi's shared dir now vendors batch-send.mjs (already used by the claude-code, codex, opencode and zcode plugins). Fixes #4504 Co-Authored-By: ktz03 <2484593937@qq.com> Co-Authored-By: jiale li <2946192893@qq.com> Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Uued6mAGJ4jfTW4XRmNWwL * fix(pi-extension): enqueue non-retryable batch failures and bound drain (#4702) Address #4692 review nits from now-ing and jiale-li-orion: - enqueueRemainder always queues on enqueue-on-failure (incl. 400/403) so the sync watermark advances - drainSessionBacklog soft-bounded by OPENVIKING_PENDING_DRAIN_BUDGET_MS / MAX_BATCHES - tests for watermark enqueue and maxBatches stop * fix(pi-extension): keep batch-send drop policy, advance watermark in pi #4702 made sendSessionMessages enqueue payloads after a non-retryable rejection so pi's sync watermark would advance. That reverses a policy pinned by the shared batch-send tests (poison payloads are dropped, not queued) for every harness, and the other vendored copies were not regenerated. Keep the shared policy and fix the watermark where the need is: pi's sendPayloads counts non-retryable drops as accepted, matching the outcome replayPending() applies to such entries. Co-Authored-By: ktz03 <2484593937@qq.com> Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Uued6mAGJ4jfTW4XRmNWwL --------- Co-authored-by: ktz03 <2484593937@qq.com> Co-authored-by: jiale li <2946192893@qq.com> Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
co-authored by
ktz03
Claude Fable 5.1
jiale li
parent
37db9c834e
commit
854ff4ceb0
@@ -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.
|
||||
|
||||
@@ -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<Record<string, any>>,
|
||||
opts?: { enqueueOnRetryable?: boolean; onSent?: (count: number) => void | Promise<void> },
|
||||
): Promise<{
|
||||
sent: number;
|
||||
queued: number;
|
||||
enqueueFailed: number;
|
||||
failed: number;
|
||||
retryable: boolean;
|
||||
usedBatch: boolean;
|
||||
lastError: any;
|
||||
}>;
|
||||
@@ -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<object>} 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;
|
||||
}
|
||||
@@ -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<string | null>;
|
||||
export function dequeue(filename: string): Promise<boolean>;
|
||||
export function incrementRetry(filename: string, entry: Record<string, any>): Promise<boolean>;
|
||||
export function cleanStale(): Promise<number>;
|
||||
|
||||
@@ -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<boolean> {
|
||||
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<void> {
|
||||
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<SyncBranchResult> {
|
||||
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<AddPayloadResult> {
|
||||
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<void> {
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user