Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **fix(db):** `DELETE /v1/batches/delete-completed` now caps the work it does per request and reports `hasMore` so a caller can resume, and the sweep no longer deletes a file that another batch still references (#13680, #13681)
26 changes: 23 additions & 3 deletions open-sse/services/batchProcessor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import {
getBatch,
getPendingBatches,
getTerminalBatches,
isFileReferencedByOtherBatch,
listBatchItemCheckpoints,
markBatchItemError,
markBatchItemProcessing,
Expand Down Expand Up @@ -253,13 +254,32 @@ async function cleanupExpiredBatches(): Promise<void> {
: null;
const outputExpiresAt = getBatchOutputExpiresAt(batch);

if (batch.inputFileId && inputExpiresAt && now > inputExpiresAt) {
// #13681: skip the soft-delete when some OTHER batch still references
// the same file id (e.g. one input file reused across batches) — a
// terminal batch's own expiry must not null a file a sibling still
// needs.
if (
batch.inputFileId &&
inputExpiresAt &&
now > inputExpiresAt &&
!isFileReferencedByOtherBatch(batch.inputFileId, [batch.id])
) {
deleteFile(batch.inputFileId);
}
if (batch.outputFileId && outputExpiresAt && now > outputExpiresAt) {
if (
batch.outputFileId &&
outputExpiresAt &&
now > outputExpiresAt &&
!isFileReferencedByOtherBatch(batch.outputFileId, [batch.id])
) {
deleteFile(batch.outputFileId);
}
if (batch.errorFileId && outputExpiresAt && now > outputExpiresAt) {
if (
batch.errorFileId &&
outputExpiresAt &&
now > outputExpiresAt &&
!isFileReferencedByOtherBatch(batch.errorFileId, [batch.id])
) {
deleteFile(batch.errorFileId);
}
}
Expand Down
8 changes: 7 additions & 1 deletion src/app/api/v1/batches/delete-completed/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,7 @@ export async function DELETE(request: Request) {
apiKeyId: scope.apiKeyId,
deletedBatches: result.deletedBatches,
deletedFiles: result.deletedFiles,
hasMore: result.hasMore,
};
// A bulk delete is an audit event, not routine chatter: an instance-wide sweep
// and any key-scoped sweep that actually removed rows log at warn so the trail
Expand All @@ -109,7 +110,12 @@ export async function DELETE(request: Request) {
}

return NextResponse.json(
{ deleted: true, deletedBatches: result.deletedBatches, deletedFiles: result.deletedFiles },
{
deleted: true,
deletedBatches: result.deletedBatches,
deletedFiles: result.deletedFiles,
hasMore: result.hasMore,
},
{ headers: CORS_HEADERS }
);
}
93 changes: 84 additions & 9 deletions src/lib/db/batches.ts
Original file line number Diff line number Diff line change
Expand Up @@ -387,22 +387,25 @@ export function deleteBatch(id: string): boolean {

db.prepare("DELETE FROM batch_item_checkpoints WHERE batch_id = ?").run(id);

// Soft-delete associated files (input, output, error)
if (batch.inputFileId) {
// Soft-delete associated files (input, output, error) — but only when no
// OTHER batch still references the same file id (#13681). A file shared
// across batches (e.g. one input file reused for several batch submissions)
// must survive as long as any sibling batch still points at it.
if (batch.inputFileId && !isFileReferencedByOtherBatch(batch.inputFileId, [id])) {
try {
deleteFile(batch.inputFileId);
} catch {
/* ignore */
}
}
if (batch.outputFileId) {
if (batch.outputFileId && !isFileReferencedByOtherBatch(batch.outputFileId, [id])) {
try {
deleteFile(batch.outputFileId);
} catch {
/* ignore */
}
}
if (batch.errorFileId) {
if (batch.errorFileId && !isFileReferencedByOtherBatch(batch.errorFileId, [id])) {
try {
deleteFile(batch.errorFileId);
} catch {
Expand All @@ -424,6 +427,47 @@ export type DeleteCompletedBatchesScope = { apiKeyId: string } | { allTenants: t
/** Both sweep modes commit in chunks of this many batches (SEC-D, LEDGER-4). */
export const INSTANCE_SWEEP_CHUNK = 200;

/**
* Upper bound on the number of `INSTANCE_SWEEP_CHUNK`-sized chunks a single
* `deleteCompletedBatches` call may run (#13680). `sweepLoop` is a synchronous
* `for (;;)` over `better-sqlite3` — with no cap, one request could hold the
* Node.js event loop for as long as it takes to sweep every completed batch on
* the instance. 25 × 200 = 5000 batches/request is a judgment call, not a hard
* constraint; any caller with more to sweep gets `hasMore: true` back and
* resumes by calling again — resumption falls out naturally from rowid
* ordering plus delete-as-you-go (already-swept rows are gone, so the next
* SELECT picks up the next-lowest surviving rowid on its own; no cursor field
* needed).
*/
export const MAX_CHUNKS_PER_REQUEST = 25;

/**
* True when some batch OTHER than one of `excludeBatchIds` still references
* `fileId` as its input/output/error file (#13681). Used before soft-deleting
* a file to avoid nulling content a surviving batch still needs. Binds
* `fileId` three times; when `excludeBatchIds` is empty the `NOT IN (...)`
* clause is dropped entirely rather than emitted empty (`NOT IN ()` is invalid
* SQL, and getting the guard wrong there would silently match everything).
*/
export function isFileReferencedByOtherBatch(fileId: string, excludeBatchIds: string[]): boolean {
const db = getDbInstance();
if (excludeBatchIds.length === 0) {
const row = db
.prepare(
"SELECT 1 FROM batches WHERE input_file_id = ? OR output_file_id = ? OR error_file_id = ? LIMIT 1"
)
.get(fileId, fileId, fileId);
return !!row;
}
const marks = excludeBatchIds.map(() => "?").join(",");
const row = db
.prepare(
`SELECT 1 FROM batches WHERE (input_file_id = ? OR output_file_id = ? OR error_file_id = ?) AND id NOT IN (${marks}) LIMIT 1`
)
.get(fileId, fileId, fileId, ...excludeBatchIds);
return !!row;
}

/**
* Delete completed batches and the files they reference.
*
Expand Down Expand Up @@ -457,10 +501,25 @@ export const INSTANCE_SWEEP_CHUNK = 200;
* so the lowest-privilege caller — any valid API key — cannot hold the
* instance's single writer for the length of its whole sweep. Each chunk stays
* atomic: a failure inside chunk N leaves chunks < N committed, chunk N fully
* rolled back, and rethrows. Inherent to per-chunk commits, in either mode: a
* file shared by batches in two different chunks can be nulled by chunk 1
* before chunk 2 fails; the surviving batch row is swept by the next run. The
* returned totals sum the chunks.
* rolled back, and rethrows. Each chunk's file soft-deletes now check whether
* some batch OUTSIDE that chunk still references the file
* (`isFileReferencedByOtherBatch`, #13681) — a file shared with a
* non-completed sibling batch, or with a completed batch a LATER chunk hasn't
* reached yet, survives this chunk. The only case that check cannot see is
* pure timing: chunk 1 commits and nulls a file, then — before chunk 2 runs —
* a NEW batch is created reusing that same file id. That race is inherent to
* per-chunk commits and stays a known, accepted edge case; everything else
* (a concurrently existing sibling, in any status, in any chunk) is now
* guarded. The returned totals sum the chunks.
*
* A single call commits at most `MAX_CHUNKS_PER_REQUEST` chunks (#13680):
* `better-sqlite3` is synchronous, so an unbounded loop would hold the event
* loop for as long as it takes to sweep the whole key's/instance's backlog.
* When the cap is hit with more rows still pending, the call returns
* `hasMore: true` instead of continuing; the caller (the DELETE route) simply
* calls again. Resumption needs no cursor: swept rows are gone, `rowid` only
* increases, so the next call's `ORDER BY rowid LIMIT ?` picks up exactly
* where the previous call left off.
*
* The loop must make progress: it remembers the first id of the previous chunk
* and throws if the next chunk starts with the same id — the DELETE removed
Expand All @@ -476,6 +535,7 @@ export const INSTANCE_SWEEP_CHUNK = 200;
export function deleteCompletedBatches(scope: DeleteCompletedBatchesScope): {
deletedBatches: number;
deletedFiles: number;
hasMore: boolean;
} {
const scopeObj = scope && typeof scope === "object" ? scope : {};
const allTenants = "allTenants" in scopeObj && scopeObj.allTenants === true;
Expand Down Expand Up @@ -517,6 +577,11 @@ export function deleteCompletedBatches(scope: DeleteCompletedBatchesScope): {

let deletedFiles = 0;
for (const fid of fileIds) {
// #13681: a file referenced by a batch outside this chunk (a different
// chunk not yet processed, or any non-completed batch — completed
// batches outside `ids` cannot exist since the SELECT above IS the
// chunk) must survive this chunk's sweep.
if (isFileReferencedByOtherBatch(fid, ids)) continue;
try {
// Key mode: only the key's OWN files. A batch may reference a file
// another tenant (or nobody) owns; a bulk destructive sweep must not
Expand All @@ -541,11 +606,20 @@ export function deleteCompletedBatches(scope: DeleteCompletedBatchesScope): {
// modes is the SELECT that produces the next chunk. No outer transaction —
// the write lock is released between chunks (LEDGER-4/20/21).
const sweepLoop = (nextIds: () => string[]) => {
const totals = { deletedBatches: 0, deletedFiles: 0 };
const totals = { deletedBatches: 0, deletedFiles: 0, hasMore: false };
let previousFirstId: string | null = null;
let chunkCount = 0;
for (;;) {
const ids = nextIds();
if (ids.length === 0) break;
// Chunk cap (#13680): a single request commits at most
// MAX_CHUNKS_PER_REQUEST chunks. This peek doesn't delete or count
// anything — it only tells the caller whether more work remains so it
// can call again (resumption is natural: already-swept rows are gone).
if (chunkCount >= MAX_CHUNKS_PER_REQUEST) {
totals.hasMore = true;
break;
}
// Forward-progress guard (LEDGER-22): the chunk is re-selected from the
// table after each commit, so a repeated first id means the previous
// DELETE removed nothing and the loop would spin forever. A concurrent
Expand All @@ -557,6 +631,7 @@ export function deleteCompletedBatches(scope: DeleteCompletedBatchesScope): {
const part = sweepIds(ids);
totals.deletedBatches += part.deletedBatches;
totals.deletedFiles += part.deletedFiles;
chunkCount++;
}
return totals;
};
Expand Down
6 changes: 6 additions & 0 deletions tests/unit/batches-delete-completed-route-scope.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,7 @@ async function callDelete(headers: Record<string, string>, url: string = ROUTE_U
deleted?: boolean;
deletedBatches?: number;
deletedFiles?: number;
hasMore?: boolean;
error?: { message: string; type?: string; code?: string };
};
return { res, body };
Expand All @@ -115,6 +116,11 @@ describe("DELETE /api/v1/batches/delete-completed — caller scope (GHSA-wvxc-jp
assert.strictEqual(body.deleted, true);
assert.strictEqual(body.deletedBatches, 0, "key A owns no completed batch — nothing to sweep");
assert.strictEqual(body.deletedFiles, 0);
assert.strictEqual(
body.hasMore,
false,
"the response must surface deleteCompletedBatches' hasMore continuation flag (#13680)"
);
assert.ok(getBatch(victim.batch.id), "key B's completed batch must survive key A's sweep");
assert.strictEqual(
getFileContent(victim.file.id)?.toString(),
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,111 @@
/**
* #13680 — DELETE /v1/batches/delete-completed does unbounded work per request.
*
* `deleteCompletedBatches` commits in chunks of `INSTANCE_SWEEP_CHUNK` (200), but
* the outer `sweepLoop` is a plain synchronous `for (;;)` that only stops when the
* table has no completed batch left — there is no cap on how many chunks a single
* request may run. better-sqlite3 is synchronous, so one request can hold the
* Node.js event loop for as long as it takes to sweep the ENTIRE table, and the
* caller has no way to ask for a bounded amount of work per call (no `hasMore`).
*
* This test seeds one row past the issue's own proposed cap
* (`MAX_CHUNKS_PER_REQUEST = 25` × `INSTANCE_SWEEP_CHUNK` = 5000) and asserts a
* single call stays within that bound and reports a continuation flag.
*/
import { describe, it, after } from "node:test";
import assert from "node:assert";
import fs from "node:fs";
import os from "node:os";
import path from "node:path";

const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "issue13680-"));
process.env.DATA_DIR = TEST_DATA_DIR;

const { createFile } = await import("../../src/lib/db/files.ts");
const { createBatch, deleteCompletedBatches, INSTANCE_SWEEP_CHUNK, MAX_CHUNKS_PER_REQUEST } =
await import("../../src/lib/db/batches.ts");
const { resetDbInstance } = await import("../../src/lib/db/core.ts");

function seedCompletedBatch(label: string, apiKeyId: string | null = null) {
const file = createFile({
bytes: 1,
filename: `${label}.jsonl`,
purpose: "batch",
content: Buffer.from("x"),
apiKeyId,
});
return createBatch({
endpoint: "/v1/chat/completions",
completionWindow: "24h",
inputFileId: file.id,
status: "completed",
apiKeyId,
});
}

describe("#13680 — deleteCompletedBatches has a per-request chunk cap", () => {
after(() => {
resetDbInstance();
fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 });
});

it("stops after MAX_CHUNKS_PER_REQUEST chunks and reports hasMore instead of sweeping the whole table in one synchronous call", () => {
const cap = MAX_CHUNKS_PER_REQUEST * INSTANCE_SWEEP_CHUNK;
const total = cap + 1; // one row past the cap: proves the loop stops at it

for (let i = 0; i < total; i++) seedCompletedBatch(`issue13680-${i}`);

const result = deleteCompletedBatches({ allTenants: true }) as {
deletedBatches: number;
deletedFiles: number;
hasMore?: boolean;
};

assert.ok(
result.deletedBatches <= cap,
`expected a single request to sweep at most ${cap} batches (MAX_CHUNKS_PER_REQUEST=${MAX_CHUNKS_PER_REQUEST} × INSTANCE_SWEEP_CHUNK=${INSTANCE_SWEEP_CHUNK}), ` +
`but one call deleted ${result.deletedBatches} of ${total} in one synchronous pass — no per-request cap exists`
);
assert.strictEqual(
result.hasMore,
true,
"the result carries no continuation signal (`hasMore`), so a caller cannot tell more completed batches remain to sweep"
);
});

it("resumes across repeated calls until hasMore is false, sweeping the entire backlog", () => {
// Key-scoped on purpose: isolates this test's count from the leftover
// unowned batch the previous test's `allTenants` sweep may not have caught
// (its cap+1 seed leaves exactly one row past MAX_CHUNKS_PER_REQUEST), so
// the expected call count here stays exact regardless of test order.
const apiKeyId = "resume-key-13680";
const cap = MAX_CHUNKS_PER_REQUEST * INSTANCE_SWEEP_CHUNK;
const total = cap + 50;
for (let i = 0; i < total; i++) seedCompletedBatch(`issue13680-resume-${i}`, apiKeyId);

let totalDeleted = 0;
let hasMore = true;
let calls = 0;
while (hasMore) {
calls++;
if (calls > 10) throw new Error("resumption did not converge within 10 calls");
const result = deleteCompletedBatches({ apiKeyId }) as {
deletedBatches: number;
hasMore: boolean;
};
totalDeleted += result.deletedBatches;
hasMore = result.hasMore;
}

assert.strictEqual(
calls,
2,
"50 extra rows past one cap should resume in exactly one more call"
);
assert.strictEqual(
totalDeleted,
total,
"every seeded batch must be swept across the resumed calls"
);
});
});
Loading
Loading