|
8 | 8 | getSock, |
9 | 9 | installWebMonitorInboxUnitTestHooks, |
10 | 10 | mockLoadConfig, |
11 | | - settleInboundWork, |
12 | | - waitForMessageCalls, |
13 | 11 | } from "./monitor-inbox.test-harness.js"; |
14 | 12 | let monitorWebInbox: typeof import("./inbound.js").monitorWebInbox; |
15 | 13 | const inboundLoggerInfoMock = vi.hoisted(() => vi.fn()); |
@@ -37,23 +35,45 @@ describe("web monitor inbox", () => { |
37 | 35 | monitorWebInbox = getMonitorWebInbox(); |
38 | 36 | }); |
39 | 37 |
|
40 | | - async function openMonitor(onMessage = vi.fn()) { |
| 38 | + async function openMonitor( |
| 39 | + onMessage = vi.fn(), |
| 40 | + extraOptions: Partial<Parameters<typeof monitorWebInbox>[0]> = {}, |
| 41 | + ) { |
41 | 42 | return await monitorWebInbox({ |
42 | 43 | cfg: mockLoadConfig() as never, |
43 | 44 | verbose: false, |
44 | 45 | accountId: DEFAULT_ACCOUNT_ID, |
45 | 46 | authDir: getAuthDir(), |
46 | 47 | onMessage, |
| 48 | + ...extraOptions, |
47 | 49 | }); |
48 | 50 | } |
49 | 51 |
|
50 | 52 | async function runSingleUpsertAndCapture(upsert: unknown) { |
51 | 53 | const onMessage = vi.fn(); |
52 | | - const listener = await openMonitor(onMessage); |
| 54 | + let armed = false; |
| 55 | + let observedPendingWork = false; |
| 56 | + let resolvePendingWorkDrained!: () => void; |
| 57 | + const pendingWorkDrained = new Promise<void>((resolve) => { |
| 58 | + resolvePendingWorkDrained = resolve; |
| 59 | + }); |
| 60 | + const listener = await openMonitor(onMessage, { |
| 61 | + onPendingWorkChanged: (pendingWorkCount) => { |
| 62 | + if (!armed) { |
| 63 | + return; |
| 64 | + } |
| 65 | + if (pendingWorkCount > 0) { |
| 66 | + observedPendingWork = true; |
| 67 | + } else if (observedPendingWork) { |
| 68 | + resolvePendingWorkDrained(); |
| 69 | + } |
| 70 | + }, |
| 71 | + }); |
53 | 72 | const sock = getSock(); |
| 73 | + // The monitor owns async media and delivery work; wait for its drain signal instead of polling. |
| 74 | + armed = true; |
54 | 75 | sock.ev.emit("messages.upsert", upsert); |
55 | | - await waitForMessageCalls(onMessage, 1); |
56 | | - await settleInboundWork(); |
| 76 | + await pendingWorkDrained; |
57 | 77 | return { onMessage, listener, sock }; |
58 | 78 | } |
59 | 79 |
|
|
0 commit comments