feat(app): expose local sharing and verified transfers
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
This commit is contained in:
@@ -0,0 +1,776 @@
|
||||
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");
|
||||
});
|
||||
Reference in New Issue
Block a user