File size: 13,021 Bytes
3144483
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
/**
 * Atomic lane publication must never admit work before the group budget exists.
 * Park admitted tasks and sample peak concurrency at task entry so an excess
 * admission cannot disappear before the assertion observes it.
 */
import { afterEach, beforeEach, describe, expect, test } from "vitest";
import { createDeferred, withTestTimeout } from "../../test/helpers/promise.js";
import {
  clearCommandLane,
  enqueueCommandInLane,
  getCommandLaneSnapshot,
  publishLaneConfiguration,
  resetAllLanes,
  setCommandLaneConcurrency,
} from "./command-queue.js";

const CRON = "cron-nested";
const HOOK = "hook-dispatch";
const DELIVERY = "delivery-dispatch";
const GROUP = "cron-hooks";
const MOVED_GROUP = "cron-delivery";

type LaneGroupSpec = NonNullable<Parameters<typeof publishLaneConfiguration>[0]["groups"]>[string];

function setCommandLaneGroup(group: string, spec: LaneGroupSpec): void {
  publishLaneConfiguration({ groups: { [group]: spec } });
}

function clearCommandLaneGroup(group: string): void {
  publishLaneConfiguration({ clearGroups: [group] });
}

beforeEach(() => {
  resetAllLanes();
  clearCommandLaneGroup(GROUP);
  clearCommandLaneGroup(MOVED_GROUP);
});

afterEach(() => {
  clearCommandLaneGroup(GROUP);
  clearCommandLaneGroup(MOVED_GROUP);
  resetAllLanes();
});

describe("publishLaneConfiguration", () => {
  test("no member dispatches above budget DURING publication", async () => {
    // Both lanes start closed with work already queued, so the only thing that
    // can release them is publication itself. If publication widened a lane and
    // drained it before installing the group — what the sequential per-lane
    // setter does — the two lanes would admit up to 8 + 4 = 12 tasks.
    setCommandLaneConcurrency(CRON, 0);
    setCommandLaneConcurrency(HOOK, 0);

    let active = 0;
    let peak = 0;
    const publishedStarts = createDeferred();
    const gates: Array<{ resolve: () => void }> = [];
    const runs: Array<Promise<unknown>> = [];
    const park = (lane: string) => {
      const g = createDeferred();
      gates.push(g);
      runs.push(
        enqueueCommandInLane(lane, async () => {
          active += 1;
          // Peak is sampled on entry, before anything can retire, so work
          // admitted inside the publication window cannot escape the count.
          peak = Math.max(peak, active);
          if (active >= 8) {
            publishedStarts.resolve();
          }
          await g.promise;
          active -= 1;
        }),
      );
    };
    for (let i = 0; i < 12; i++) {
      park(CRON);
    }
    for (let i = 0; i < 6; i++) {
      park(HOOK);
    }
    expect(active).toBe(0); // nothing may run before publication
    expect(getCommandLaneSnapshot(CRON).activeCount).toBe(0);
    expect(getCommandLaneSnapshot(HOOK).activeCount).toBe(0);

    try {
      publishLaneConfiguration({
        lanes: { [CRON]: 8, [HOOK]: 4 },
        groups: {
          [GROUP]: {
            budget: 8,
            members: [CRON, HOOK],
            reservations: { [HOOK]: 1 },
          },
        },
      });

      await withTestTimeout(
        publishedStarts.promise,
        1_000,
        "publication did not start the shared group budget",
      );

      // Sample task entry while all admitted tasks remain parked.
      expect(peak).toBeLessThanOrEqual(8);
      // And not vacuous — publication must actually have dispatched to the cap.
      expect(peak).toBe(8);
    } finally {
      for (const g of gates) {
        g.resolve();
      }
      clearCommandLane(CRON);
      clearCommandLane(HOOK);
      await Promise.allSettled(runs);
    }
  });

  test("commit dispatch uses group order rather than publication object order", async () => {
    setCommandLaneConcurrency(CRON, 0);
    setCommandLaneConcurrency(HOOK, 0);

    const starts: string[] = [];
    const cronStarted = createDeferred();
    const hookStarted = createDeferred();
    const cronGate = createDeferred();
    const hookGate = createDeferred();
    const olderCron = enqueueCommandInLane(
      CRON,
      async () => {
        starts.push(CRON);
        cronStarted.resolve();
        await cronGate.promise;
      },
      { priority: "background" },
    );
    const newerHook = enqueueCommandInLane(
      HOOK,
      async () => {
        starts.push(HOOK);
        hookStarted.resolve();
        await hookGate.promise;
      },
      { priority: "background" },
    );

    try {
      // Deliberately publish HOOK first in both objects. The older CRON head must
      // still own the single shared slot.
      publishLaneConfiguration({
        lanes: { [HOOK]: 1, [CRON]: 1 },
        groups: { [GROUP]: { budget: 1, members: [HOOK, CRON] } },
      });
      await withTestTimeout(
        cronStarted.promise,
        1_000,
        "publication did not start the older cron task",
      );
      expect(starts).toEqual([CRON]);

      cronGate.resolve();
      await olderCron;
      await withTestTimeout(
        hookStarted.promise,
        1_000,
        "cron completion did not start the queued hook",
      );
      expect(starts).toEqual([CRON, HOOK]);
    } finally {
      cronGate.resolve();
      hookGate.resolve();
      clearCommandLane(CRON);
      clearCommandLane(HOOK);
      await Promise.allSettled([olderCron, newerHook]);
    }
  });

  test("moving a busy member wakes queued work in its previous group", async () => {
    publishLaneConfiguration({
      lanes: { [CRON]: 1, [HOOK]: 1, [DELIVERY]: 1 },
      groups: { [GROUP]: { budget: 1, members: [CRON, HOOK] } },
    });

    const cronGate = createDeferred();
    const cronRun = enqueueCommandInLane(CRON, async () => await cronGate.promise);
    const hookGate = createDeferred();
    const hookRun = enqueueCommandInLane(HOOK, async () => await hookGate.promise);
    expect(getCommandLaneSnapshot(HOOK)).toMatchObject({ activeCount: 0, queuedCount: 1 });

    // CRON's active task stops counting against the old group as soon as it is
    // moved. That newly free old-group capacity must wake HOOK immediately.
    publishLaneConfiguration({
      groups: { [MOVED_GROUP]: { budget: 1, members: [CRON, DELIVERY] } },
    });
    expect(getCommandLaneSnapshot(HOOK)).toMatchObject({
      group: GROUP,
      activeCount: 1,
      queuedCount: 0,
    });
    expect(getCommandLaneSnapshot(CRON).group).toBe(MOVED_GROUP);

    cronGate.resolve();
    hookGate.resolve();
    await Promise.all([cronRun, hookRun]);
  });

  test("a rejected configuration does not leave lanes widened and dispatching", async () => {
    setCommandLaneConcurrency(CRON, 0);
    const gates = Array.from({ length: 4 }, () => createDeferred());
    const runs = gates.map((g) => enqueueCommandInLane(CRON, async () => await g.promise));

    // sum(reservations) > budget is rejected. Validation must happen before any
    // drain, or the lane is left open at width 8 governed by no group at all.
    expect(() =>
      publishLaneConfiguration({
        lanes: { [CRON]: 8 },
        groups: {
          [GROUP]: {
            budget: 2,
            members: [CRON, HOOK],
            reservations: { [CRON]: 2, [HOOK]: 1 },
          },
        },
      }),
    ).toThrow(/reserves 3 slots but its budget is 2/);

    expect(getCommandLaneSnapshot(CRON).activeCount).toBe(0);

    for (const g of gates) {
      g.resolve();
    }
    // The lane never opened, so this work is still queued. resetAllLanes
    // PRESERVES queued entries by design, so it would never settle these —
    // clearCommandLane rejects them instead.
    clearCommandLane(CRON);
    await Promise.allSettled(runs);
  });

  test("a rejected configuration does not leave lane maxima mutated", async () => {
    // Stronger than asserting activeCount === 0 after the throw: that only
    // proves no commit-time drain ran, not that the lane was left alone. If
    // phase 1 widens a lane and group validation then throws, the lane sits at
    // the new width governed by NO group, and the next unrelated drain trigger
    // dispatches the preserved queue ungoverned.
    setCommandLaneConcurrency(CRON, 0);
    const gates = Array.from({ length: 4 }, () => createDeferred());
    const runs = gates.map((g) => enqueueCommandInLane(CRON, async () => await g.promise));
    expect(getCommandLaneSnapshot(CRON).maxConcurrent).toBe(0);

    expect(() =>
      publishLaneConfiguration({
        lanes: { [CRON]: 8 },
        groups: {
          [GROUP]: {
            budget: 2,
            members: [CRON, HOOK],
            reservations: { [CRON]: 2, [HOOK]: 1 },
          },
        },
      }),
    ).toThrow(/reserves 3 slots but its budget is 2/);

    // The lane must be exactly as it was before the rejected publish.
    expect(getCommandLaneSnapshot(CRON).maxConcurrent).toBe(0);
    expect(getCommandLaneSnapshot(CRON).group).toBeUndefined();

    // And a later drain trigger must not dispatch the queue that was preserved
    // across the failed publish.
    const extra = createDeferred();
    const extraRun = enqueueCommandInLane(CRON, async () => await extra.promise);
    expect(getCommandLaneSnapshot(CRON).activeCount).toBe(0);

    for (const g of gates) {
      g.resolve();
    }
    extra.resolve();
    clearCommandLane(CRON);
    await Promise.allSettled([...runs, extraRun]);
  });

  test("a rejected replacement does not tear down the existing group first", async () => {
    // Combining clearGroups with an invalid replacement is the
    // worst case — the old group could be removed before the new one throws,
    // leaving BOTH lane width and group membership partially committed. Phase 0
    // validation has to run before the clear, not just before the install.
    publishLaneConfiguration({
      lanes: { [CRON]: 8, [HOOK]: 1 },
      groups: {
        [GROUP]: { budget: 8, members: [CRON, HOOK], reservations: { [HOOK]: 1 } },
      },
    });
    expect(getCommandLaneSnapshot(CRON).group).toBe(GROUP);

    expect(() =>
      publishLaneConfiguration({
        lanes: { [CRON]: 99 },
        clearGroups: [GROUP],
        groups: {
          "replacement-group": {
            budget: 1,
            members: [CRON, HOOK],
            reservations: { [CRON]: 1, [HOOK]: 1 },
          },
        },
      }),
    ).toThrow(/reserves 2 slots but its budget is 1/);

    // Everything must be exactly as before: group intact, width untouched.
    expect(getCommandLaneSnapshot(CRON).group).toBe(GROUP);
    expect(getCommandLaneSnapshot(CRON).groupBudget).toBe(8);
    expect(getCommandLaneSnapshot(CRON).maxConcurrent).toBe(8);
    expect(getCommandLaneSnapshot(HOOK).reservedForLane).toBe(1);
  });

  test("publication wakes members when a replacement frees capacity", async () => {
    // Replacing a group must wake queued members when its new budget has room,
    // without waiting for an unrelated enqueue to trigger another drain.
    setCommandLaneConcurrency(CRON, 8);
    setCommandLaneConcurrency(HOOK, 1);
    setCommandLaneGroup(GROUP, { budget: 2, members: [CRON, HOOK] });

    const gates = Array.from({ length: 5 }, () => createDeferred());
    const runs = gates.map((g) => enqueueCommandInLane(CRON, async () => await g.promise));
    expect(getCommandLaneSnapshot(CRON).activeCount).toBe(2);
    expect(getCommandLaneSnapshot(CRON).queuedCount).toBe(3);

    // Publish a wider group budget without changing individual lane widths.
    setCommandLaneGroup(GROUP, { budget: 5, members: [CRON, HOOK] });

    // The queued work must start on the replacement itself.
    expect(getCommandLaneSnapshot(CRON).activeCount).toBe(5);
    expect(getCommandLaneSnapshot(CRON).queuedCount).toBe(0);

    for (const g of gates) {
      g.resolve();
    }
    await Promise.all(runs);
  });

  test("republishing a narrower budget does not admit beyond the new cap", async () => {
    publishLaneConfiguration({
      lanes: { [CRON]: 8, [HOOK]: 1 },
      groups: {
        [GROUP]: { budget: 8, members: [CRON, HOOK], reservations: { [HOOK]: 1 } },
      },
    });

    const gates = Array.from({ length: 3 }, () => createDeferred());
    const runs = gates.map((g) => enqueueCommandInLane(CRON, async () => await g.promise));
    expect(getCommandLaneSnapshot(CRON).activeCount).toBe(3);

    // Narrowing mid-flight cannot evict running work, but it must not admit
    // more: the group is already over its new budget.
    publishLaneConfiguration({
      lanes: { [CRON]: 8, [HOOK]: 1 },
      groups: {
        [GROUP]: { budget: 2, members: [CRON, HOOK], reservations: { [HOOK]: 1 } },
      },
    });
    const extra = createDeferred();
    const blocked = enqueueCommandInLane(CRON, async () => await extra.promise);

    expect(getCommandLaneSnapshot(CRON).activeCount).toBe(3);
    expect(getCommandLaneSnapshot(CRON).blockedBy).toBe("group-budget");

    for (const g of gates) {
      g.resolve();
    }
    extra.resolve();
    clearCommandLane(CRON);
    await Promise.allSettled([...runs, blocked]);
  });
});