From 685a3dbe5d725fa8479bb70628de8caebd035826 Mon Sep 17 00:00:00 2001 From: zhangLei99586 Date: Sun, 2 Aug 2026 01:40:40 +0800 Subject: [PATCH] fix(proxy-capture): make path-based session cleanup atomic (#98852) --- src/proxy-capture/store.sqlite.test.ts | 61 ++++++++++++++++++++++++++ src/proxy-capture/store.sqlite.ts | 18 +++++--- 2 files changed, 72 insertions(+), 7 deletions(-) diff --git a/src/proxy-capture/store.sqlite.test.ts b/src/proxy-capture/store.sqlite.test.ts index 719b2a5a3a94..486ea164e548 100644 --- a/src/proxy-capture/store.sqlite.test.ts +++ b/src/proxy-capture/store.sqlite.test.ts @@ -1,6 +1,7 @@ // Proxy capture SQLite store tests cover persisted capture reads and writes. import fs from "node:fs"; import path from "node:path"; +import { constants } from "node:sqlite"; import { afterEach, describe, expect, it, vi } from "vitest"; import { cleanupTempDirs, makeTempDir } from "../../test/helpers/temp-dir.js"; import { resolveSqliteDatabaseFilePaths } from "../infra/sqlite-files.js"; @@ -161,6 +162,66 @@ describe("DebugProxyCaptureStore", () => { expect(lease.store.isClosed).toBe(true); }); + it.each(["deleteSessions", "purgeAll"] as const)( + "rolls back path-based %s when session deletion fails", + (operation) => { + const root = makeTempDir(cleanupDirs, "openclaw-proxy-capture-rollback-"); + const dbPath = path.join(root, "capture.sqlite"); + const blobDir = path.join(root, "blobs"); + const lease = acquireDebugProxyCaptureStore(dbPath, blobDir); + const sessionId = "path-based-rollback-session"; + + try { + lease.store.upsertSession({ + id: sessionId, + startedAt: 1, + mode: "sdk", + sourceScope: "openclaw", + sourceProcess: "plugin", + dbPath, + blobDir, + }); + const blob = lease.store.persistPayload(Buffer.from("rollback payload"), "text/plain"); + lease.store.recordEvent({ + sessionId, + ts: 2, + sourceScope: "openclaw", + sourceProcess: "plugin", + protocol: "https", + direction: "outbound", + kind: "request", + flowId: "path-based-rollback-flow", + dataBlobId: blob.blobId, + dataSha256: blob.sha256, + }); + + const cleanup = () => + operation === "deleteSessions" + ? lease.store.deleteSessions([sessionId]) + : lease.store.purgeAll(); + lease.store.db.setAuthorizer((action, table) => + action === constants.SQLITE_DELETE && table === "capture_sessions" + ? constants.SQLITE_DENY + : constants.SQLITE_OK, + ); + + expect(cleanup).toThrow(/not authorized/u); + lease.store.db.setAuthorizer(null); + expect(lease.store.listSessions()).toHaveLength(1); + expect(lease.store.getSessionEvents(sessionId)).toHaveLength(1); + expect(fs.existsSync(blob.path)).toBe(true); + + expect(cleanup()).toEqual({ sessions: 1, events: 1, blobs: 1 }); + expect(lease.store.listSessions()).toEqual([]); + expect(lease.store.getSessionEvents(sessionId)).toEqual([]); + expect(fs.existsSync(blob.path)).toBe(false); + } finally { + lease.store.db.setAuthorizer(null); + lease.release(); + } + }, + ); + it("uses rollback journaling for captures on NFS-backed volumes", () => { vi.spyOn(fs, "statfsSync").mockReturnValue({ type: 0x6969, diff --git a/src/proxy-capture/store.sqlite.ts b/src/proxy-capture/store.sqlite.ts index 074926061a5e..4f5ecbce5e1c 100644 --- a/src/proxy-capture/store.sqlite.ts +++ b/src/proxy-capture/store.sqlite.ts @@ -616,7 +616,9 @@ class DebugProxyCaptureStoreImpl { const eventCount = (this.db.prepare(`SELECT COUNT(*) AS count FROM capture_events`).get() as { count: number }) .count ?? 0; - this.db.exec(`DELETE FROM capture_events; DELETE FROM capture_sessions;`); + runSqliteImmediateTransactionSync(this.db, () => { + this.db.exec(`DELETE FROM capture_events; DELETE FROM capture_sessions;`); + }); let blobs = 0; if (fs.existsSync(this.pathBased.blobDir)) { for (const entry of fs.readdirSync(this.pathBased.blobDir)) { @@ -764,12 +766,14 @@ class DebugProxyCaptureStoreImpl { ) .get(...sessionIds) as { count: number } ).count ?? 0; - this.db - .prepare(`DELETE FROM capture_events WHERE session_id IN (${placeholders})`) - .run(...sessionIds); - this.db - .prepare(`DELETE FROM capture_sessions WHERE id IN (${placeholders})`) - .run(...sessionIds); + runSqliteImmediateTransactionSync(this.db, () => { + this.db + .prepare(`DELETE FROM capture_events WHERE session_id IN (${placeholders})`) + .run(...sessionIds); + this.db + .prepare(`DELETE FROM capture_sessions WHERE id IN (${placeholders})`) + .run(...sessionIds); + }); const candidateBlobIds = blobRows .map((row) => row.blobId?.trim()) .filter((blobId): blobId is string => Boolean(blobId));