Skip to content

Commit 27e1393

Browse files
committed
refactor: share store writer queue
1 parent ec1e27d commit 27e1393

4 files changed

Lines changed: 162 additions & 216 deletions

File tree

src/commitments/store-writer.ts

Lines changed: 15 additions & 99 deletions
Original file line numberDiff line numberDiff line change
@@ -5,20 +5,13 @@
55
import fs from "node:fs/promises";
66
import path from "node:path";
77
import { type FileLockOptions, withFileLock } from "../plugin-sdk/file-lock.js";
8+
import {
9+
clearStoreWriterQueuesForTest,
10+
runQueuedStoreWrite,
11+
type StoreWriterQueue,
12+
} from "../shared/store-writer-queue.js";
813

9-
type CommitmentsStoreWriterTask = {
10-
fn: () => Promise<unknown>;
11-
resolve: (value: unknown) => void;
12-
reject: (reason: unknown) => void;
13-
};
14-
15-
type CommitmentsStoreWriterQueue = {
16-
running: boolean;
17-
pending: CommitmentsStoreWriterTask[];
18-
drainPromise: Promise<void> | null;
19-
};
20-
21-
const WRITER_QUEUES = new Map<string, CommitmentsStoreWriterQueue>();
14+
const WRITER_QUEUES = new Map<string, StoreWriterQueue>();
2215

2316
// Matches src/plugin-sdk/persistent-dedupe.ts so both lock-protected stores share tuning.
2417
const DEFAULT_COMMITMENTS_LOCK_OPTIONS: FileLockOptions = {
@@ -32,67 +25,6 @@ const DEFAULT_COMMITMENTS_LOCK_OPTIONS: FileLockOptions = {
3225
stale: 60_000,
3326
};
3427

35-
function getOrCreateWriterQueue(storePath: string): CommitmentsStoreWriterQueue {
36-
const existing = WRITER_QUEUES.get(storePath);
37-
if (existing) {
38-
return existing;
39-
}
40-
const created: CommitmentsStoreWriterQueue = {
41-
running: false,
42-
pending: [],
43-
drainPromise: null,
44-
};
45-
WRITER_QUEUES.set(storePath, created);
46-
return created;
47-
}
48-
49-
async function drainCommitmentsStoreWriterQueue(storePath: string): Promise<void> {
50-
const queue = WRITER_QUEUES.get(storePath);
51-
if (!queue) {
52-
return;
53-
}
54-
if (queue.drainPromise) {
55-
await queue.drainPromise;
56-
return;
57-
}
58-
queue.running = true;
59-
queue.drainPromise = (async () => {
60-
try {
61-
while (queue.pending.length > 0) {
62-
const task = queue.pending.shift();
63-
if (!task) {
64-
continue;
65-
}
66-
let result: unknown;
67-
let failed: unknown;
68-
let hasFailure = false;
69-
try {
70-
result = await task.fn();
71-
} catch (err) {
72-
hasFailure = true;
73-
failed = err;
74-
}
75-
if (hasFailure) {
76-
task.reject(failed);
77-
continue;
78-
}
79-
task.resolve(result);
80-
}
81-
} finally {
82-
queue.running = false;
83-
queue.drainPromise = null;
84-
if (queue.pending.length === 0) {
85-
WRITER_QUEUES.delete(storePath);
86-
} else {
87-
queueMicrotask(() => {
88-
void drainCommitmentsStoreWriterQueue(storePath);
89-
});
90-
}
91-
}
92-
})();
93-
await queue.drainPromise;
94-
}
95-
9628
// The advisory lockfile lives next to the data file; create the parent dir up
9729
// front so acquireFileLock does not ENOENT before the user fn ever runs.
9830
async function ensureCommitmentsStoreDir(storePath: string): Promise<void> {
@@ -103,33 +35,17 @@ export async function runExclusiveCommitmentsStoreWrite<T>(
10335
storePath: string,
10436
fn: () => Promise<T>,
10537
): Promise<T> {
106-
if (!storePath || typeof storePath !== "string") {
107-
throw new Error(
108-
`runExclusiveCommitmentsStoreWrite: storePath must be a non-empty string, got ${JSON.stringify(
109-
storePath,
110-
)}`,
111-
);
112-
}
113-
const queue = getOrCreateWriterQueue(storePath);
114-
return await new Promise<T>((resolve, reject) => {
115-
const task: CommitmentsStoreWriterTask = {
116-
fn: async () => {
117-
await ensureCommitmentsStoreDir(storePath);
118-
return await withFileLock(storePath, DEFAULT_COMMITMENTS_LOCK_OPTIONS, fn);
119-
},
120-
resolve: (value) => resolve(value as T),
121-
reject,
122-
};
123-
queue.pending.push(task);
124-
void drainCommitmentsStoreWriterQueue(storePath);
38+
return await runQueuedStoreWrite({
39+
queues: WRITER_QUEUES,
40+
storePath,
41+
label: "runExclusiveCommitmentsStoreWrite",
42+
fn: async () => {
43+
await ensureCommitmentsStoreDir(storePath);
44+
return await withFileLock(storePath, DEFAULT_COMMITMENTS_LOCK_OPTIONS, fn);
45+
},
12546
});
12647
}
12748

12849
export function clearCommitmentsStoreWriterQueuesForTest(): void {
129-
for (const queue of WRITER_QUEUES.values()) {
130-
for (const task of queue.pending) {
131-
task.reject(new Error("commitments store writer queue cleared for test"));
132-
}
133-
}
134-
WRITER_QUEUES.clear();
50+
clearStoreWriterQueuesForTest(WRITER_QUEUES, "commitments store writer queue cleared for test");
13551
}

src/config/sessions/store-writer-state.ts

Lines changed: 10 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -1,47 +1,23 @@
1+
import {
2+
clearStoreWriterQueuesForTest,
3+
drainStoreWriterQueuesForTest,
4+
type StoreWriterQueue,
5+
type StoreWriterTask,
6+
} from "../../shared/store-writer-queue.js";
17
import { clearSessionStoreCaches } from "./store-cache.js";
28

3-
export type SessionStoreWriterTask = {
4-
fn: () => Promise<unknown>;
5-
resolve: (value: unknown) => void;
6-
reject: (reason: unknown) => void;
7-
};
8-
9-
export type SessionStoreWriterQueue = {
10-
running: boolean;
11-
pending: SessionStoreWriterTask[];
12-
drainPromise: Promise<void> | null;
13-
};
9+
export type SessionStoreWriterTask = StoreWriterTask;
10+
export type SessionStoreWriterQueue = StoreWriterQueue;
1411

1512
export const WRITER_QUEUES = new Map<string, SessionStoreWriterQueue>();
1613

1714
export function clearSessionStoreCacheForTest(): void {
1815
clearSessionStoreCaches();
19-
for (const queue of WRITER_QUEUES.values()) {
20-
for (const task of queue.pending) {
21-
task.reject(new Error("session store queue cleared for test"));
22-
}
23-
}
24-
WRITER_QUEUES.clear();
16+
clearStoreWriterQueuesForTest(WRITER_QUEUES, "session store queue cleared for test");
2517
}
2618

2719
export async function drainSessionStoreWriterQueuesForTest(): Promise<void> {
28-
while (WRITER_QUEUES.size > 0) {
29-
const queues = [...WRITER_QUEUES.values()];
30-
for (const queue of queues) {
31-
for (const task of queue.pending) {
32-
task.reject(new Error("session store queue cleared for test"));
33-
}
34-
queue.pending.length = 0;
35-
}
36-
const activeDrains = queues.flatMap((queue) =>
37-
queue.drainPromise ? [queue.drainPromise] : [],
38-
);
39-
if (activeDrains.length === 0) {
40-
WRITER_QUEUES.clear();
41-
return;
42-
}
43-
await Promise.allSettled(activeDrains);
44-
}
20+
await drainStoreWriterQueuesForTest(WRITER_QUEUES, "session store queue cleared for test");
4521
}
4622

4723
export function getSessionStoreWriterQueueSizeForTest(): number {
Lines changed: 7 additions & 83 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,5 @@
1-
import {
2-
WRITER_QUEUES,
3-
type SessionStoreWriterQueue,
4-
type SessionStoreWriterTask,
5-
} from "./store-writer-state.js";
1+
import { runQueuedStoreWrite } from "../../shared/store-writer-queue.js";
2+
import { WRITER_QUEUES } from "./store-writer-state.js";
63

74
export async function withSessionStoreWriterForTest<T>(
85
storePath: string,
@@ -11,87 +8,14 @@ export async function withSessionStoreWriterForTest<T>(
118
return await runExclusiveSessionStoreWrite(storePath, fn);
129
}
1310

14-
function getOrCreateWriterQueue(storePath: string): SessionStoreWriterQueue {
15-
const existing = WRITER_QUEUES.get(storePath);
16-
if (existing) {
17-
return existing;
18-
}
19-
const created: SessionStoreWriterQueue = { running: false, pending: [], drainPromise: null };
20-
WRITER_QUEUES.set(storePath, created);
21-
return created;
22-
}
23-
24-
async function drainSessionStoreWriterQueue(storePath: string): Promise<void> {
25-
const queue = WRITER_QUEUES.get(storePath);
26-
if (!queue) {
27-
return;
28-
}
29-
if (queue.drainPromise) {
30-
await queue.drainPromise;
31-
return;
32-
}
33-
queue.running = true;
34-
queue.drainPromise = (async () => {
35-
try {
36-
while (queue.pending.length > 0) {
37-
const task = queue.pending.shift();
38-
if (!task) {
39-
continue;
40-
}
41-
42-
let result: unknown;
43-
let failed: unknown;
44-
let hasFailure = false;
45-
try {
46-
result = await task.fn();
47-
} catch (err) {
48-
hasFailure = true;
49-
failed = err;
50-
}
51-
if (hasFailure) {
52-
task.reject(failed);
53-
continue;
54-
}
55-
task.resolve(result);
56-
}
57-
} finally {
58-
queue.running = false;
59-
queue.drainPromise = null;
60-
if (queue.pending.length === 0) {
61-
WRITER_QUEUES.delete(storePath);
62-
} else {
63-
queueMicrotask(() => {
64-
void drainSessionStoreWriterQueue(storePath);
65-
});
66-
}
67-
}
68-
})();
69-
await queue.drainPromise;
70-
}
71-
7211
export async function runExclusiveSessionStoreWrite<T>(
7312
storePath: string,
7413
fn: () => Promise<T>,
7514
): Promise<T> {
76-
if (!storePath || typeof storePath !== "string") {
77-
throw new Error(
78-
`runExclusiveSessionStoreWrite: storePath must be a non-empty string, got ${JSON.stringify(
79-
storePath,
80-
)}`,
81-
);
82-
}
83-
const queue = getOrCreateWriterQueue(storePath);
84-
85-
const promise = new Promise<T>((resolve, reject) => {
86-
const task: SessionStoreWriterTask = {
87-
fn: async () => await fn(),
88-
resolve: (value) => resolve(value as T),
89-
reject,
90-
};
91-
92-
queue.pending.push(task);
93-
void drainSessionStoreWriterQueue(storePath);
15+
return await runQueuedStoreWrite({
16+
queues: WRITER_QUEUES,
17+
storePath,
18+
label: "runExclusiveSessionStoreWrite",
19+
fn,
9420
});
95-
96-
return await promise;
9721
}

0 commit comments

Comments
 (0)