Add a durable, acknowledged Local network sharing switch with fail-closed hydration, serialized mutation, and redacted ephemeral-identity diagnostics. Keep local Call-to-Play state available while gating every network action on the effective sharing generation. Render revisioned verification, invalid-source retry, and sticky source exhaustion states. Preserve opaque attempt IDs through progress delivery so out-of-order webview events cannot attach stale bytes to a successor transfer, and keep terminal exhaustion visible after the last source departs. Own listeners, native invokes, persistence, dialogs, and companion-window creation through webview close. Late creation is settled and cleaned before the parent realm is destroyed. Test Plan: - `just frontend-test` -- passed (91/91) - `just build` -- passed with TypeScript, Vite, and release Tauri compilation - `just test` -- passed on the completed stack (708 workspace tests) - `just clippy` -- passed on the completed stack - `git diff --cached --check` -- passed
777 lines
20 KiB
TypeScript
777 lines
20 KiB
TypeScript
import {
|
|
AsyncAdoptionScope,
|
|
type AsyncCleanup,
|
|
AsyncOwner,
|
|
createSerializedAsyncWriter,
|
|
mergeHydratedState,
|
|
ownCompanionWindowCreation,
|
|
registerSequentially,
|
|
} from "../src/lib/asyncOwnership.ts";
|
|
|
|
const assertEquals = <T>(actual: T, expected: T, message: string) => {
|
|
if (actual !== expected) {
|
|
throw new Error(`${message}: expected ${expected}, got ${actual}`);
|
|
}
|
|
};
|
|
|
|
const assertArrayEquals = <T>(actual: T[], expected: T[], message: string) => {
|
|
assertEquals(JSON.stringify(actual), JSON.stringify(expected), message);
|
|
};
|
|
|
|
const deferred = <T>() => {
|
|
let resolve!: (value: T) => void;
|
|
let reject!: (reason: unknown) => void;
|
|
const promise = new Promise<T>((resolvePromise, rejectPromise) => {
|
|
resolve = resolvePromise;
|
|
reject = rejectPromise;
|
|
});
|
|
return { promise, resolve, reject };
|
|
};
|
|
|
|
Deno.test("unmount before listener resolution cleans the late listener and stops registration", async () => {
|
|
const owner = new AsyncOwner();
|
|
const first = deferred<AsyncCleanup>();
|
|
let firstRegistrations = 0;
|
|
let secondRegistrations = 0;
|
|
let firstCleanups = 0;
|
|
|
|
const registration = registerSequentially(owner, [
|
|
() => {
|
|
firstRegistrations += 1;
|
|
return first.promise;
|
|
},
|
|
() => {
|
|
secondRegistrations += 1;
|
|
return Promise.resolve(() => {});
|
|
},
|
|
]);
|
|
|
|
assertEquals(
|
|
firstRegistrations,
|
|
1,
|
|
"the first listener should start immediately",
|
|
);
|
|
const disposal = owner.dispose();
|
|
first.resolve(() => {
|
|
firstCleanups += 1;
|
|
});
|
|
|
|
assertEquals(
|
|
await registration,
|
|
false,
|
|
"disposed registration should report cancellation",
|
|
);
|
|
assertEquals(
|
|
firstCleanups,
|
|
1,
|
|
"the late listener should unlisten immediately",
|
|
);
|
|
assertEquals(
|
|
secondRegistrations,
|
|
0,
|
|
"no later listener should start after disposal",
|
|
);
|
|
await disposal;
|
|
});
|
|
|
|
Deno.test("partial listener registration is fully cleaned across an in-flight listener", async () => {
|
|
const owner = new AsyncOwner();
|
|
const second = deferred<AsyncCleanup>();
|
|
const secondStarted = deferred<void>();
|
|
let firstCleanups = 0;
|
|
let secondCleanups = 0;
|
|
let thirdRegistrations = 0;
|
|
|
|
const registration = registerSequentially(owner, [
|
|
() =>
|
|
Promise.resolve(() => {
|
|
firstCleanups += 1;
|
|
}),
|
|
() => {
|
|
secondStarted.resolve();
|
|
return second.promise;
|
|
},
|
|
() => {
|
|
thirdRegistrations += 1;
|
|
return Promise.resolve(() => {});
|
|
},
|
|
]);
|
|
|
|
await secondStarted.promise;
|
|
const disposal = owner.dispose();
|
|
assertEquals(
|
|
firstCleanups,
|
|
1,
|
|
"an owned listener should clean up during disposal",
|
|
);
|
|
|
|
second.resolve(() => {
|
|
secondCleanups += 1;
|
|
});
|
|
assertEquals(
|
|
await registration,
|
|
false,
|
|
"the in-flight registration should stop the sequence",
|
|
);
|
|
assertEquals(
|
|
secondCleanups,
|
|
1,
|
|
"the in-flight listener should clean up when it resolves",
|
|
);
|
|
assertEquals(thirdRegistrations, 0, "the third listener should never start");
|
|
await disposal;
|
|
|
|
await owner.dispose();
|
|
assertEquals(
|
|
firstCleanups,
|
|
1,
|
|
"repeated disposal must not duplicate owned cleanup",
|
|
);
|
|
assertEquals(
|
|
secondCleanups,
|
|
1,
|
|
"repeated disposal must not duplicate late cleanup",
|
|
);
|
|
});
|
|
|
|
Deno.test("listener registration failure cleans the partial scope and skips later listeners", async () => {
|
|
const owner = new AsyncOwner();
|
|
const failure = new Error("registration failed");
|
|
let firstCleanups = 0;
|
|
let thirdRegistrations = 0;
|
|
let reported: unknown;
|
|
|
|
try {
|
|
await registerSequentially(owner, [
|
|
() =>
|
|
Promise.resolve(() => {
|
|
firstCleanups += 1;
|
|
}),
|
|
() => Promise.reject(failure),
|
|
() => {
|
|
thirdRegistrations += 1;
|
|
return Promise.resolve(() => {});
|
|
},
|
|
]);
|
|
} catch (error) {
|
|
reported = error;
|
|
}
|
|
|
|
assertEquals(reported, failure, "the registration error should be preserved");
|
|
assertEquals(
|
|
firstCleanups,
|
|
1,
|
|
"earlier listeners should clean up on setup failure",
|
|
);
|
|
assertEquals(
|
|
thirdRegistrations,
|
|
0,
|
|
"later listeners should not start after setup failure",
|
|
);
|
|
assertEquals(
|
|
owner.isActive(),
|
|
false,
|
|
"a failed registration scope should stay disposed",
|
|
);
|
|
});
|
|
|
|
Deno.test("late refresh results and post-disposal refreshes cannot publish", async () => {
|
|
const owner = new AsyncOwner();
|
|
const refresh = deferred<string>();
|
|
let refreshStarts = 0;
|
|
const applied: string[] = [];
|
|
|
|
const pending = owner.applyIfActive(
|
|
() => {
|
|
refreshStarts += 1;
|
|
return refresh.promise;
|
|
},
|
|
(value) => applied.push(value),
|
|
);
|
|
|
|
const disposal = owner.dispose();
|
|
refresh.resolve("stale");
|
|
assertEquals(
|
|
await pending,
|
|
false,
|
|
"a late refresh should report cancellation",
|
|
);
|
|
assertArrayEquals(applied, [], "a late refresh must not publish state");
|
|
|
|
const afterDispose = await owner.applyIfActive(
|
|
() => {
|
|
refreshStarts += 1;
|
|
return Promise.resolve("newer");
|
|
},
|
|
(value) => applied.push(value),
|
|
);
|
|
assertEquals(
|
|
afterDispose,
|
|
false,
|
|
"a disposed owner should reject new refresh work",
|
|
);
|
|
assertEquals(refreshStarts, 1, "refresh work must not start after disposal");
|
|
assertArrayEquals(
|
|
applied,
|
|
[],
|
|
"post-disposal refreshes must not publish state",
|
|
);
|
|
await disposal;
|
|
});
|
|
|
|
Deno.test("disposal joins a late asynchronous listener cleanup", async () => {
|
|
const owner = new AsyncOwner();
|
|
const listener = deferred<AsyncCleanup>();
|
|
const cleanupDone = deferred<void>();
|
|
let cleanupStarts = 0;
|
|
let disposalSettled = false;
|
|
|
|
const registration = owner.register(() => listener.promise);
|
|
const disposal = owner.dispose().then(() => {
|
|
disposalSettled = true;
|
|
});
|
|
listener.resolve(async () => {
|
|
cleanupStarts += 1;
|
|
await cleanupDone.promise;
|
|
});
|
|
|
|
await Promise.resolve();
|
|
assertEquals(
|
|
cleanupStarts,
|
|
1,
|
|
"the late cleanup should start immediately on resolution",
|
|
);
|
|
assertEquals(
|
|
disposalSettled,
|
|
false,
|
|
"disposal must wait for the asynchronous cleanup",
|
|
);
|
|
|
|
cleanupDone.resolve();
|
|
assertEquals(
|
|
await registration,
|
|
false,
|
|
"the disposed registration should remain cancelled",
|
|
);
|
|
await disposal;
|
|
assertEquals(
|
|
disposalSettled,
|
|
true,
|
|
"disposal should settle after cleanup completion",
|
|
);
|
|
});
|
|
|
|
Deno.test("asynchronous cleanup rejection is reported and disposal still drains", async () => {
|
|
const cleanupDone = deferred<void>();
|
|
const failure = new Error("unlisten failed");
|
|
const reported: unknown[] = [];
|
|
const owner = new AsyncOwner((error) => reported.push(error));
|
|
|
|
await owner.register(() =>
|
|
Promise.resolve(async () => {
|
|
await cleanupDone.promise;
|
|
throw failure;
|
|
})
|
|
);
|
|
const disposal = owner.dispose();
|
|
cleanupDone.resolve();
|
|
await disposal;
|
|
|
|
assertEquals(
|
|
reported.length,
|
|
1,
|
|
"the cleanup failure should be reported once",
|
|
);
|
|
assertEquals(
|
|
reported[0],
|
|
failure,
|
|
"the original cleanup failure should be reported",
|
|
);
|
|
});
|
|
|
|
Deno.test("a root adoption scope drains work and observes rejection", async () => {
|
|
const first = deferred<void>();
|
|
const failure = new Error("adopted cleanup failed");
|
|
const reported: unknown[] = [];
|
|
const scope = new AsyncAdoptionScope((error) => reported.push(error));
|
|
let drained = false;
|
|
|
|
scope.adopt(first.promise);
|
|
scope.adopt(Promise.reject(failure));
|
|
const drain = scope.drain().then(() => {
|
|
drained = true;
|
|
});
|
|
await Promise.resolve();
|
|
|
|
assertEquals(drained, false, "the adoption scope should retain pending work");
|
|
assertEquals(
|
|
reported.length,
|
|
1,
|
|
"an adopted rejection should be observed once",
|
|
);
|
|
assertEquals(
|
|
reported[0],
|
|
failure,
|
|
"the adoption scope should report the original error",
|
|
);
|
|
|
|
first.resolve();
|
|
await drain;
|
|
assertEquals(
|
|
drained,
|
|
true,
|
|
"the adoption scope should drain after all work settles",
|
|
);
|
|
});
|
|
|
|
Deno.test("window close drains every admitted operation and rejects late owner registration", async () => {
|
|
const scope = new AsyncAdoptionScope();
|
|
const cleanupDone = deferred<void>();
|
|
const unrelatedInvoke = deferred<string>();
|
|
const owner = new AsyncOwner(() => {}, scope);
|
|
let cleanupStarts = 0;
|
|
|
|
await owner.register(() =>
|
|
Promise.resolve(async () => {
|
|
cleanupStarts += 1;
|
|
await cleanupDone.promise;
|
|
})
|
|
);
|
|
const unrelated = owner.applyIfActive(
|
|
() => unrelatedInvoke.promise,
|
|
() => {},
|
|
);
|
|
let closeDrained = false;
|
|
const close = scope.disposeOwned().then(() => {
|
|
closeDrained = true;
|
|
});
|
|
|
|
const lateOwner = new AsyncOwner(() => {}, scope);
|
|
assertEquals(
|
|
lateOwner.isActive(),
|
|
false,
|
|
"owner admission must close synchronously",
|
|
);
|
|
assertEquals(cleanupStarts, 1, "listener cleanup should start during close");
|
|
assertEquals(closeDrained, false, "window close must await listener cleanup");
|
|
|
|
cleanupDone.resolve();
|
|
await Promise.resolve();
|
|
assertEquals(
|
|
closeDrained,
|
|
false,
|
|
"an admitted invoke must retain the webview after listener cleanup",
|
|
);
|
|
|
|
unrelatedInvoke.resolve("settled during close");
|
|
assertEquals(
|
|
await unrelated,
|
|
false,
|
|
"unrelated late invoke publication stays suppressed",
|
|
);
|
|
await close;
|
|
assertEquals(
|
|
closeDrained,
|
|
true,
|
|
"window close must join every admitted operation",
|
|
);
|
|
});
|
|
|
|
Deno.test("closing during a dialog prevents its chained invoke from starting", async () => {
|
|
const scope = new AsyncAdoptionScope();
|
|
const owner = new AsyncOwner(() => {}, scope);
|
|
const dialog = deferred<boolean>();
|
|
let invokes = 0;
|
|
|
|
const action = (async () => {
|
|
let confirmed = false;
|
|
if (
|
|
!await owner.applyIfActive(
|
|
() => dialog.promise,
|
|
(answer) => {
|
|
confirmed = answer;
|
|
},
|
|
) || !confirmed
|
|
) {
|
|
return;
|
|
}
|
|
await owner.applyIfActive(
|
|
() => {
|
|
invokes += 1;
|
|
return Promise.resolve();
|
|
},
|
|
() => {},
|
|
);
|
|
})();
|
|
|
|
const close = scope.disposeOwned();
|
|
dialog.resolve(true);
|
|
await action;
|
|
await close;
|
|
|
|
assertEquals(
|
|
invokes,
|
|
0,
|
|
"a dialog result must not start work after close admission",
|
|
);
|
|
});
|
|
|
|
Deno.test("closing during a helper import prevents later helper side effects", async () => {
|
|
const scope = new AsyncAdoptionScope();
|
|
const owner = new AsyncOwner(() => {}, scope);
|
|
const imported = deferred<string>();
|
|
let helperStarts = 0;
|
|
|
|
const helper = (async () => {
|
|
let moduleName: string | undefined;
|
|
if (
|
|
!await owner.applyIfActive(
|
|
() => imported.promise,
|
|
(value) => {
|
|
moduleName = value;
|
|
},
|
|
) || moduleName === undefined
|
|
) {
|
|
return;
|
|
}
|
|
await owner.applyIfActive(
|
|
() => {
|
|
helperStarts += 1;
|
|
return Promise.resolve();
|
|
},
|
|
() => {},
|
|
);
|
|
})();
|
|
|
|
const close = scope.disposeOwned();
|
|
imported.resolve("window helper");
|
|
await helper;
|
|
await close;
|
|
|
|
assertEquals(
|
|
helperStarts,
|
|
0,
|
|
"a late import must not focus or create a window",
|
|
);
|
|
});
|
|
|
|
Deno.test("close after companion construction drains creation and both listeners before destroy", async () => {
|
|
const scope = new AsyncAdoptionScope();
|
|
const owner = new AsyncOwner(() => {}, scope);
|
|
const createdCleanupDone = deferred<void>();
|
|
const errorCleanupDone = deferred<void>();
|
|
const bothCleanupsStarted = deferred<void>();
|
|
const companionDestroyStarted = deferred<void>();
|
|
const companionDestroyDone = deferred<void>();
|
|
const order: string[] = [];
|
|
let createdHandler: (() => void) | undefined;
|
|
let errorHandler: ((payload: unknown) => void) | undefined;
|
|
let cleanupStarts = 0;
|
|
let latePublications = 0;
|
|
let rootDestroyed = false;
|
|
|
|
const noteCleanupStart = () => {
|
|
cleanupStarts += 1;
|
|
if (cleanupStarts === 2) bothCleanupsStarted.resolve();
|
|
};
|
|
const creation = ownCompanionWindowCreation(owner, {
|
|
registerCreated: (handler) => {
|
|
createdHandler = handler;
|
|
return Promise.resolve(async () => {
|
|
order.push("created-listener-cleanup");
|
|
noteCleanupStart();
|
|
await createdCleanupDone.promise;
|
|
});
|
|
},
|
|
registerError: (handler) => {
|
|
errorHandler = handler;
|
|
return Promise.resolve(async () => {
|
|
order.push("error-listener-cleanup");
|
|
noteCleanupStart();
|
|
await errorCleanupDone.promise;
|
|
});
|
|
},
|
|
destroy: async () => {
|
|
order.push("companion-destroy");
|
|
companionDestroyStarted.resolve();
|
|
await companionDestroyDone.promise;
|
|
},
|
|
}).then((result) => {
|
|
if (owner.isActive()) latePublications += 1;
|
|
return result;
|
|
});
|
|
scope.adopt(creation);
|
|
|
|
if (createdHandler === undefined || errorHandler === undefined) {
|
|
throw new Error(
|
|
"construction must synchronously begin both creation listeners",
|
|
);
|
|
}
|
|
const close = scope.disposeOwned().then(() => {
|
|
order.push("root-destroy");
|
|
rootDestroyed = true;
|
|
});
|
|
await Promise.resolve();
|
|
assertEquals(
|
|
rootDestroyed,
|
|
false,
|
|
"close must wait for the native creation outcome",
|
|
);
|
|
|
|
order.push("created-event");
|
|
createdHandler();
|
|
errorHandler(new Error("late losing outcome"));
|
|
await bothCleanupsStarted.promise;
|
|
assertEquals(
|
|
rootDestroyed,
|
|
false,
|
|
"root destruction must wait for both native unlisten acknowledgements",
|
|
);
|
|
|
|
createdCleanupDone.resolve();
|
|
errorCleanupDone.resolve();
|
|
await companionDestroyStarted.promise;
|
|
assertEquals(
|
|
rootDestroyed,
|
|
false,
|
|
"a companion created after close must be destroyed before its parent realm",
|
|
);
|
|
companionDestroyDone.resolve();
|
|
|
|
assertEquals(
|
|
(await creation).kind,
|
|
"created",
|
|
"the first native outcome must win exactly once",
|
|
);
|
|
await close;
|
|
assertEquals(
|
|
latePublications,
|
|
0,
|
|
"late creation must not publish into disposed React state",
|
|
);
|
|
assertArrayEquals(
|
|
order,
|
|
[
|
|
"created-event",
|
|
"created-listener-cleanup",
|
|
"error-listener-cleanup",
|
|
"companion-destroy",
|
|
"root-destroy",
|
|
],
|
|
"creation, listener cleanup, companion destruction, and root destruction stay ordered",
|
|
);
|
|
});
|
|
|
|
Deno.test("listener-first bootstrap keeps an update that arrives during the initial refresh", async () => {
|
|
const owner = new AsyncOwner();
|
|
const initial = deferred<string>();
|
|
const update = deferred<string>();
|
|
const published: string[] = [];
|
|
let emitUpdate: (() => Promise<boolean>) | undefined;
|
|
const order: string[] = [];
|
|
|
|
await owner.register(() => {
|
|
order.push("listener");
|
|
emitUpdate = () =>
|
|
owner.applyLatestIfActive(
|
|
() => update.promise,
|
|
(value) => published.push(value),
|
|
);
|
|
return Promise.resolve(() => {});
|
|
});
|
|
order.push("snapshot");
|
|
const initialRefresh = owner.applyLatestIfActive(
|
|
() => initial.promise,
|
|
(value) => published.push(value),
|
|
);
|
|
const eventRefresh = emitUpdate?.();
|
|
if (eventRefresh === undefined) throw new Error("listener was not installed");
|
|
|
|
initial.resolve("stale-initial");
|
|
assertEquals(
|
|
await initialRefresh,
|
|
false,
|
|
"the event should supersede the bootstrap snapshot",
|
|
);
|
|
update.resolve("event-update");
|
|
assertEquals(await eventRefresh, true, "the event refresh should publish");
|
|
|
|
assertArrayEquals(
|
|
order,
|
|
["listener", "snapshot"],
|
|
"the listener must precede the snapshot",
|
|
);
|
|
assertArrayEquals(
|
|
published,
|
|
["event-update"],
|
|
"the bootstrap interval must not lose updates",
|
|
);
|
|
await owner.dispose();
|
|
});
|
|
|
|
Deno.test("out-of-order refresh completion publishes only the latest request", async () => {
|
|
const owner = new AsyncOwner();
|
|
const older = deferred<string>();
|
|
const newer = deferred<string>();
|
|
const published: string[] = [];
|
|
|
|
const olderRefresh = owner.applyLatestIfActive(
|
|
() => older.promise,
|
|
(value) => published.push(value),
|
|
);
|
|
const newerRefresh = owner.applyLatestIfActive(
|
|
() => newer.promise,
|
|
(value) => published.push(value),
|
|
);
|
|
|
|
newer.resolve("newer");
|
|
assertEquals(await newerRefresh, true, "the newest refresh should publish");
|
|
older.resolve("older");
|
|
assertEquals(
|
|
await olderRefresh,
|
|
false,
|
|
"the older refresh should be suppressed",
|
|
);
|
|
assertArrayEquals(
|
|
published,
|
|
["newer"],
|
|
"older completion must not overwrite newer state",
|
|
);
|
|
await owner.dispose();
|
|
});
|
|
|
|
Deno.test("settings hydration keeps saved fields while applying newer user edits", () => {
|
|
const restored = {
|
|
accent: "saved-accent",
|
|
density: "saved-density",
|
|
username: "saved-user",
|
|
};
|
|
const pendingEdits = { accent: "user-edit" };
|
|
|
|
const merged = mergeHydratedState(restored, pendingEdits);
|
|
|
|
assertEquals(
|
|
merged.accent,
|
|
"user-edit",
|
|
"the newer edit should win its field",
|
|
);
|
|
assertEquals(
|
|
merged.density,
|
|
"saved-density",
|
|
"unrelated saved fields must survive",
|
|
);
|
|
assertEquals(
|
|
merged.username,
|
|
"saved-user",
|
|
"the full saved baseline must be retained",
|
|
);
|
|
});
|
|
|
|
Deno.test("serialized writes do not start a newer value before the older value settles", async () => {
|
|
const firstDone = deferred<void>();
|
|
const firstStarted = deferred<void>();
|
|
const secondDone = deferred<void>();
|
|
const secondStarted = deferred<void>();
|
|
const starts: number[] = [];
|
|
const writer = createSerializedAsyncWriter<number>((value) => {
|
|
starts.push(value);
|
|
if (value === 1) {
|
|
firstStarted.resolve();
|
|
return firstDone.promise;
|
|
}
|
|
secondStarted.resolve();
|
|
return secondDone.promise;
|
|
}, () => {});
|
|
|
|
const firstWrite = writer.enqueue(1);
|
|
const secondWrite = writer.enqueue(2);
|
|
await firstStarted.promise;
|
|
assertArrayEquals(
|
|
starts,
|
|
[1],
|
|
"only the oldest write should start initially",
|
|
);
|
|
|
|
firstDone.resolve();
|
|
await firstWrite;
|
|
await secondStarted.promise;
|
|
assertArrayEquals(
|
|
starts,
|
|
[1, 2],
|
|
"the newer write should start after the older write",
|
|
);
|
|
|
|
secondDone.resolve();
|
|
await secondWrite;
|
|
await writer.waitForIdle();
|
|
});
|
|
|
|
Deno.test("a failed serialized write is reported and does not block the newer value", async () => {
|
|
const firstDone = deferred<void>();
|
|
const firstStarted = deferred<void>();
|
|
const secondStarted = deferred<void>();
|
|
const failure = new Error("first write failed");
|
|
const errors: unknown[] = [];
|
|
const starts: number[] = [];
|
|
const writer = createSerializedAsyncWriter<number>(async (value) => {
|
|
starts.push(value);
|
|
if (value === 1) {
|
|
firstStarted.resolve();
|
|
await firstDone.promise;
|
|
throw failure;
|
|
}
|
|
secondStarted.resolve();
|
|
}, (error) => errors.push(error));
|
|
|
|
const firstWrite = writer.enqueue(1);
|
|
const secondWrite = writer.enqueue(2);
|
|
await firstStarted.promise;
|
|
assertArrayEquals(
|
|
starts,
|
|
[1],
|
|
"the failed write should still own the queue first",
|
|
);
|
|
|
|
firstDone.resolve();
|
|
await firstWrite;
|
|
await secondStarted.promise;
|
|
await secondWrite;
|
|
|
|
assertArrayEquals(
|
|
starts,
|
|
[1, 2],
|
|
"a newer write should run after the failure is handled",
|
|
);
|
|
assertEquals(errors.length, 1, "the failed write should be reported once");
|
|
assertEquals(
|
|
errors[0],
|
|
failure,
|
|
"the original write failure should be reported",
|
|
);
|
|
});
|
|
|
|
Deno.test("closing a writer stops admission and drains the pending write", async () => {
|
|
const pending = deferred<void>();
|
|
const started = deferred<void>();
|
|
const writes: number[] = [];
|
|
const writer = createSerializedAsyncWriter<number>(async (value) => {
|
|
writes.push(value);
|
|
started.resolve();
|
|
await pending.promise;
|
|
}, () => {});
|
|
|
|
void writer.enqueue(1);
|
|
await started.promise;
|
|
let drained = false;
|
|
const drain = writer.closeAndWait().then(() => {
|
|
drained = true;
|
|
});
|
|
await writer.enqueue(2);
|
|
|
|
assertEquals(drained, false, "close must wait for the admitted write");
|
|
assertArrayEquals(writes, [1], "close must reject newer writes");
|
|
pending.resolve();
|
|
await drain;
|
|
assertEquals(drained, true, "close should settle after the pending write");
|
|
});
|