Provisioner asks the backend, not memory, whether a consumer is still there (hq issue 120) #7
Generated
+2
-2
@@ -1,12 +1,12 @@
|
|||||||
{
|
{
|
||||||
"name": "@novox/mesh-sdk",
|
"name": "@novox/mesh-sdk",
|
||||||
"version": "0.1.0",
|
"version": "0.1.1",
|
||||||
"lockfileVersion": 3,
|
"lockfileVersion": 3,
|
||||||
"requires": true,
|
"requires": true,
|
||||||
"packages": {
|
"packages": {
|
||||||
"": {
|
"": {
|
||||||
"name": "@novox/mesh-sdk",
|
"name": "@novox/mesh-sdk",
|
||||||
"version": "0.1.0",
|
"version": "0.1.1",
|
||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"@types/node": "^22.0.0",
|
"@types/node": "^22.0.0",
|
||||||
"typescript": "^5.6.0"
|
"typescript": "^5.6.0"
|
||||||
|
|||||||
+1
-1
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "@novox/mesh-sdk",
|
"name": "@novox/mesh-sdk",
|
||||||
"version": "0.1.0",
|
"version": "0.1.1",
|
||||||
"description": "The stable spine a Novox Mesh module's own code builds against.",
|
"description": "The stable spine a Novox Mesh module's own code builds against.",
|
||||||
"type": "module",
|
"type": "module",
|
||||||
"exports": {
|
"exports": {
|
||||||
|
|||||||
+100
-4
@@ -35,6 +35,13 @@ export interface Provision {
|
|||||||
export interface Adapter {
|
export interface Adapter {
|
||||||
create(p: Provision): Promise<void>;
|
create(p: Provision): Promise<void>;
|
||||||
remove(p: { readonly as: string }): Promise<void>;
|
remove(p: { readonly as: string }): Promise<void>;
|
||||||
|
/** Optional: whether the backend still holds this consumer's credential exactly as `p` says.
|
||||||
|
* Asked of every consumer already applied, every `verifyEveryMs`. `false` makes the harness apply
|
||||||
|
* it again on the same pass, so a backend that lost what it was given (a server restarted
|
||||||
|
* without persisting its users, a restore, a login removed by hand) is provisioned again instead
|
||||||
|
* of being trusted from memory (novox/hq issue 120). An adapter without it is trusted from memory,
|
||||||
|
* as before. It must only read: it is asked often, and must never change the backend. */
|
||||||
|
holds?(p: Provision): Promise<boolean>;
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface ProvisionerOptions {
|
export interface ProvisionerOptions {
|
||||||
@@ -43,8 +50,18 @@ export interface ProvisionerOptions {
|
|||||||
receives?: string;
|
receives?: string;
|
||||||
/** Reconcile interval in ms. Defaults to 5000. */
|
/** Reconcile interval in ms. Defaults to 5000. */
|
||||||
everyMs?: number;
|
everyMs?: number;
|
||||||
|
/** How often, in ms, an adapter with `holds` is asked whether the backend still holds each
|
||||||
|
* applied consumer. Defaults to 60000: slower than reconciling, because it reads the backend for
|
||||||
|
* every consumer, and fast enough that a lost login is back within a minute. */
|
||||||
|
verifyEveryMs?: number;
|
||||||
|
/** How long one `holds` may take before it counts as "could not ask". Defaults to 30000: the
|
||||||
|
* loop is sequential, so a check that hangs would otherwise stop every consumer's provisioning. */
|
||||||
|
holdsTimeoutMs?: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** The longest a consumer whose `create` keeps failing to satisfy `holds` waits between checks. */
|
||||||
|
const MAX_BACKOFF_MS = 60 * 60_000;
|
||||||
|
|
||||||
/** One entry in the mesh's contributions file: a consumer the provider must serve. */
|
/** One entry in the mesh's contributions file: a consumer the provider must serve. */
|
||||||
interface Contribution {
|
interface Contribution {
|
||||||
readonly as: string;
|
readonly as: string;
|
||||||
@@ -62,6 +79,14 @@ interface Contribution {
|
|||||||
export function runProvisioner(resource: string, adapter: Adapter, opts: ProvisionerOptions = {}): () => void {
|
export function runProvisioner(resource: string, adapter: Adapter, opts: ProvisionerOptions = {}): () => void {
|
||||||
const receives = opts.receives ?? envOrThrow("MESH_RECEIVES");
|
const receives = opts.receives ?? envOrThrow("MESH_RECEIVES");
|
||||||
const everyMs = opts.everyMs ?? 5000;
|
const everyMs = opts.everyMs ?? 5000;
|
||||||
|
const verifyEveryMs = opts.verifyEveryMs ?? 60_000;
|
||||||
|
const holdsTimeoutMs = opts.holdsTimeoutMs ?? 30_000;
|
||||||
|
let verifiedAt = 0;
|
||||||
|
// Consumers the backend reported lost, and how many times in a row. A `create` that succeeded and
|
||||||
|
// is still not held on the next check will never be: something in the adapter disagrees with
|
||||||
|
// itself. Each repeat doubles the wait before asking again, so such a bug costs one re-apply an
|
||||||
|
// hour rather than one a minute, and it is said loudly rather than quietly repeated.
|
||||||
|
const lost = new Map<string, { times: number; nextAt: number }>();
|
||||||
|
|
||||||
const applied = new Map<string, string>(); // login (`as`) -> hash of what was last applied
|
const applied = new Map<string, string>(); // login (`as`) -> hash of what was last applied
|
||||||
let stopped = false;
|
let stopped = false;
|
||||||
@@ -69,6 +94,9 @@ export function runProvisioner(resource: string, adapter: Adapter, opts: Provisi
|
|||||||
async function reconcile(): Promise<void> {
|
async function reconcile(): Promise<void> {
|
||||||
const given = await readContributions(receives, resource);
|
const given = await readContributions(receives, resource);
|
||||||
const wantByAs = new Map(given.map((g) => [g.as, g]));
|
const wantByAs = new Map(given.map((g) => [g.as, g]));
|
||||||
|
// On this pass, ask the backend rather than memory whether each applied consumer is still there.
|
||||||
|
let verifying = adapter.holds !== undefined && Date.now() - verifiedAt >= verifyEveryMs;
|
||||||
|
if (verifying) verifiedAt = Date.now();
|
||||||
|
|
||||||
// Create or update every consumer whose login, password or values changed.
|
// Create or update every consumer whose login, password or values changed.
|
||||||
for (const g of given) {
|
for (const g of given) {
|
||||||
@@ -85,12 +113,58 @@ export function runProvisioner(resource: string, adapter: Adapter, opts: Provisi
|
|||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
const h = hash(g.as, password, g.values ?? {});
|
const h = hash(g.as, password, g.values ?? {});
|
||||||
if (applied.get(g.as) === h) continue;
|
const p: Provision = { as: g.as, password, values: g.values ?? {}, at: g.at, consumer: g.node };
|
||||||
|
// Set when the backend reported this consumer lost: how many successful re-applies in a row
|
||||||
|
// it has now needed, counting this one.
|
||||||
|
let reapplying: number | undefined;
|
||||||
|
if (applied.get(g.as) === h) {
|
||||||
|
if (!verifying) continue;
|
||||||
|
const brake = lost.get(g.as);
|
||||||
|
if (brake && Date.now() < brake.nextAt) continue;
|
||||||
try {
|
try {
|
||||||
await adapter.create({ as: g.as, password, values: g.values ?? {}, at: g.at, consumer: g.node });
|
if (await withTimeout(adapter.holds!(p), holdsTimeoutMs)) {
|
||||||
applied.set(g.as, h);
|
lost.delete(g.as);
|
||||||
|
continue;
|
||||||
|
}
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
console.error(`[provisioner:${resource}] ${g.as}: create failed, will retry: ${err}`);
|
// Unable to ask is not evidence of loss. Kept as applied, asked again next time.
|
||||||
|
console.error(`[provisioner:${resource}] ${g.as}: could not check the backend, will ask again: ${scrub(err, password)}`);
|
||||||
|
// A backend that did not answer in time will not answer for the next consumer either, and
|
||||||
|
// each would cost a timeout, so the rest of this pass is not asked. Any other failure may
|
||||||
|
// be this consumer's alone, and the others are still asked.
|
||||||
|
if (err instanceof HoldsTimeout) verifying = false;
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
reapplying = (brake?.times ?? 0) + 1;
|
||||||
|
if (reapplying === 1) {
|
||||||
|
// Said, because it means the backend lost something while nothing was looking.
|
||||||
|
console.error(`[provisioner:${resource}] ${g.as}: the backend no longer holds it; applying again`);
|
||||||
|
} else {
|
||||||
|
console.error(
|
||||||
|
`[provisioner:${resource}] ${g.as}: still not held after being applied again (${reapplying} times in a row) — ` +
|
||||||
|
`create does not produce what holds checks; applying again`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
await adapter.create(p);
|
||||||
|
applied.set(g.as, h);
|
||||||
|
if (reapplying === undefined) {
|
||||||
|
// Applied for a new login, password or values: whatever was counted before does not carry over.
|
||||||
|
lost.delete(g.as);
|
||||||
|
} else {
|
||||||
|
// Only a create that succeeded counts toward the brake: if the backend still does not hold
|
||||||
|
// it at the next check, create and holds disagree, and each repeat waits twice as long.
|
||||||
|
const waitMs = Math.min(MAX_BACKOFF_MS, verifyEveryMs * 2 ** (reapplying - 1));
|
||||||
|
lost.set(g.as, { times: reapplying, nextAt: Date.now() + waitMs });
|
||||||
|
if (reapplying > 1) {
|
||||||
|
console.error(`[provisioner:${resource}] ${g.as}: next check in ${Math.round(waitMs / 1000)}s`);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
} catch (err) {
|
||||||
|
console.error(`[provisioner:${resource}] ${g.as}: create failed, will retry: ${scrub(err, password)}`);
|
||||||
|
// A failed create is not a disagreement: asked again at the next check, the count unchanged.
|
||||||
|
if (reapplying !== undefined) lost.set(g.as, { times: reapplying - 1, nextAt: 0 });
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -102,6 +176,7 @@ export function runProvisioner(resource: string, adapter: Adapter, opts: Provisi
|
|||||||
try {
|
try {
|
||||||
await adapter.remove({ as });
|
await adapter.remove({ as });
|
||||||
applied.delete(as);
|
applied.delete(as);
|
||||||
|
lost.delete(as);
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
console.error(`[provisioner:${resource}] ${as}: remove failed, will retry: ${err}`);
|
console.error(`[provisioner:${resource}] ${as}: remove failed, will retry: ${err}`);
|
||||||
}
|
}
|
||||||
@@ -151,6 +226,27 @@ async function readContributions(path: string, resource: string): Promise<Contri
|
|||||||
return (doc.given ?? []).filter((g) => g.as && g.secret);
|
return (doc.given ?? []).filter((g) => g.as && g.secret);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** A check that did not answer in time: the backend, not the consumer, is the likely cause. */
|
||||||
|
class HoldsTimeout extends Error {}
|
||||||
|
|
||||||
|
/** Reject with a HoldsTimeout if `p` has not settled within `ms`. */
|
||||||
|
function withTimeout<T>(p: Promise<T>, ms: number): Promise<T> {
|
||||||
|
let timer: NodeJS.Timeout | undefined;
|
||||||
|
const timeout = new Promise<never>((_, reject) => {
|
||||||
|
timer = setTimeout(() => reject(new HoldsTimeout(`no answer within ${ms}ms`)), ms);
|
||||||
|
});
|
||||||
|
return Promise.race([p, timeout]).finally(() => clearTimeout(timer));
|
||||||
|
}
|
||||||
|
|
||||||
|
/** An error as text with the consumer's password removed, raw and URL-encoded: a failed command's
|
||||||
|
* message can carry its arguments, and this log is not a place a password may appear. */
|
||||||
|
function scrub(err: unknown, password: string): string {
|
||||||
|
let text = String(err);
|
||||||
|
if (!password) return text;
|
||||||
|
for (const form of new Set([password, encodeURIComponent(password)])) text = text.split(form).join("***");
|
||||||
|
return text;
|
||||||
|
}
|
||||||
|
|
||||||
function hash(as: string, password: string, values: Readonly<Record<string, unknown>>): string {
|
function hash(as: string, password: string, values: Readonly<Record<string, unknown>>): string {
|
||||||
return JSON.stringify([as, password, values]);
|
return JSON.stringify([as, password, values]);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -74,6 +74,268 @@ test("provisioner creates each consumer with the mesh's login and password, remo
|
|||||||
stop();
|
stop();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test("provisioner applies again what the backend no longer holds, and trusts memory without holds", async () => {
|
||||||
|
const dir = await mkdtemp(join(tmpdir(), "prov-holds-"));
|
||||||
|
await writeFile(join(dir, "webapp.secret"), "minted-pw\n");
|
||||||
|
const receives = join(dir, "cache.json");
|
||||||
|
await writeFile(
|
||||||
|
receives,
|
||||||
|
JSON.stringify({ contributions: 1, requirement: "cache", given: [{ from: "webapp", node: "anchor", as: "webapp-anchor", secret: join(dir, "webapp.secret") }] }),
|
||||||
|
);
|
||||||
|
|
||||||
|
// A backend that forgets: what was created is held until it "restarts".
|
||||||
|
const backend = new Set<string>();
|
||||||
|
let creates = 0;
|
||||||
|
let asked = 0;
|
||||||
|
const stop = runProvisioner(
|
||||||
|
"cache",
|
||||||
|
{
|
||||||
|
async create(p) {
|
||||||
|
creates++;
|
||||||
|
backend.add(`${p.as}:${p.password}`);
|
||||||
|
},
|
||||||
|
async remove(p) {
|
||||||
|
for (const k of backend) if (k.startsWith(`${p.as}:`)) backend.delete(k);
|
||||||
|
},
|
||||||
|
async holds(p) {
|
||||||
|
asked++;
|
||||||
|
return backend.has(`${p.as}:${p.password}`);
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{ receives, everyMs: 10, verifyEveryMs: 30 },
|
||||||
|
);
|
||||||
|
|
||||||
|
await waitFor(() => creates === 1, 2000);
|
||||||
|
// While the backend holds it, asking changes nothing: no second create.
|
||||||
|
await waitFor(() => asked >= 2, 2000);
|
||||||
|
assert.equal(creates, 1);
|
||||||
|
|
||||||
|
// The backend restarts and forgets. The contributions did not change; only asking can notice.
|
||||||
|
backend.clear();
|
||||||
|
await waitFor(() => creates === 2, 2000);
|
||||||
|
assert.ok(backend.has("webapp-anchor:minted-pw"));
|
||||||
|
stop();
|
||||||
|
|
||||||
|
// An adapter without holds is trusted from memory, as before: a forgotten backend stays forgotten.
|
||||||
|
const forgetful = new Set<string>();
|
||||||
|
let plainCreates = 0;
|
||||||
|
const stopPlain = runProvisioner(
|
||||||
|
"cache",
|
||||||
|
{
|
||||||
|
async create(p) {
|
||||||
|
plainCreates++;
|
||||||
|
forgetful.add(p.as);
|
||||||
|
},
|
||||||
|
async remove() {},
|
||||||
|
},
|
||||||
|
{ receives, everyMs: 10, verifyEveryMs: 10 },
|
||||||
|
);
|
||||||
|
await waitFor(() => plainCreates === 1, 2000);
|
||||||
|
forgetful.clear();
|
||||||
|
await new Promise((r) => setTimeout(r, 100));
|
||||||
|
assert.equal(plainCreates, 1);
|
||||||
|
stopPlain();
|
||||||
|
});
|
||||||
|
|
||||||
|
test("provisioner keeps a consumer applied when the backend cannot be asked", async () => {
|
||||||
|
const dir = await mkdtemp(join(tmpdir(), "prov-unreachable-"));
|
||||||
|
await writeFile(join(dir, "webapp.secret"), "minted-pw");
|
||||||
|
const receives = join(dir, "cache.json");
|
||||||
|
await writeFile(
|
||||||
|
receives,
|
||||||
|
JSON.stringify({ requirement: "cache", given: [{ as: "webapp-anchor", secret: join(dir, "webapp.secret") }] }),
|
||||||
|
);
|
||||||
|
let creates = 0;
|
||||||
|
let asked = 0;
|
||||||
|
const stop = runProvisioner(
|
||||||
|
"cache",
|
||||||
|
{
|
||||||
|
async create() {
|
||||||
|
creates++;
|
||||||
|
},
|
||||||
|
async remove() {},
|
||||||
|
async holds() {
|
||||||
|
asked++;
|
||||||
|
throw new Error("connection refused");
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{ receives, everyMs: 10, verifyEveryMs: 20 },
|
||||||
|
);
|
||||||
|
await waitFor(() => asked >= 3, 2000);
|
||||||
|
// Unable to ask is not evidence of loss: nothing is applied again.
|
||||||
|
assert.equal(creates, 1);
|
||||||
|
stop();
|
||||||
|
});
|
||||||
|
|
||||||
|
test("provisioner backs off when create does not satisfy holds, and says so", async () => {
|
||||||
|
const dir = await mkdtemp(join(tmpdir(), "prov-brake-"));
|
||||||
|
await writeFile(join(dir, "webapp.secret"), "minted-pw");
|
||||||
|
const receives = join(dir, "cache.json");
|
||||||
|
await writeFile(receives, JSON.stringify({ requirement: "cache", given: [{ as: "webapp-anchor", secret: join(dir, "webapp.secret") }] }));
|
||||||
|
let creates = 0;
|
||||||
|
let asked = 0;
|
||||||
|
const logged: string[] = [];
|
||||||
|
const original = console.error;
|
||||||
|
console.error = (...a: unknown[]) => logged.push(a.join(" "));
|
||||||
|
try {
|
||||||
|
const stop = runProvisioner(
|
||||||
|
"cache",
|
||||||
|
{
|
||||||
|
async create() {
|
||||||
|
creates++;
|
||||||
|
},
|
||||||
|
async remove() {},
|
||||||
|
async holds() {
|
||||||
|
asked++;
|
||||||
|
return false; // an adapter that disagrees with itself: nothing create does is ever held
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{ receives, everyMs: 5, verifyEveryMs: 20 },
|
||||||
|
);
|
||||||
|
await new Promise((r) => setTimeout(r, 400));
|
||||||
|
stop();
|
||||||
|
} finally {
|
||||||
|
console.error = original;
|
||||||
|
}
|
||||||
|
// Without the brake this would be ~20 creates (one per 20ms check). With doubling waits
|
||||||
|
// (20, 40, 80, 160ms…) it is a handful.
|
||||||
|
assert.ok(creates >= 3 && creates <= 7, `creates: ${creates}`);
|
||||||
|
assert.equal(asked, creates - 1);
|
||||||
|
assert.ok(logged.some((l) => l.includes("create does not produce what holds checks")));
|
||||||
|
});
|
||||||
|
|
||||||
|
test("provisioner treats a failed or hung holds as could-not-ask, and never logs the password", async () => {
|
||||||
|
const dir = await mkdtemp(join(tmpdir(), "prov-timeout-"));
|
||||||
|
await writeFile(join(dir, "a.secret"), "s3cret-pw");
|
||||||
|
await writeFile(join(dir, "b.secret"), "other-pw");
|
||||||
|
const receives = join(dir, "cache.json");
|
||||||
|
await writeFile(
|
||||||
|
receives,
|
||||||
|
JSON.stringify({ requirement: "cache", given: [{ as: "a", secret: join(dir, "a.secret") }, { as: "b", secret: join(dir, "b.secret") }] }),
|
||||||
|
);
|
||||||
|
let creates = 0;
|
||||||
|
const askedFor: string[] = [];
|
||||||
|
const logged: string[] = [];
|
||||||
|
const original = console.error;
|
||||||
|
console.error = (...a: unknown[]) => logged.push(a.join(" "));
|
||||||
|
try {
|
||||||
|
const stop = runProvisioner(
|
||||||
|
"cache",
|
||||||
|
{
|
||||||
|
async create() {
|
||||||
|
creates++;
|
||||||
|
},
|
||||||
|
async remove() {},
|
||||||
|
holds(p) {
|
||||||
|
askedFor.push(p.as);
|
||||||
|
// "a" fails the way a shelled-out tool does, with its arguments in the message.
|
||||||
|
if (p.as === "a") return Promise.reject(new Error(`Command failed: tool --password ${p.password}`));
|
||||||
|
return Promise.resolve(true);
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{ receives, everyMs: 5, verifyEveryMs: 30, holdsTimeoutMs: 20 },
|
||||||
|
);
|
||||||
|
await new Promise((r) => setTimeout(r, 200));
|
||||||
|
stop();
|
||||||
|
} finally {
|
||||||
|
console.error = original;
|
||||||
|
}
|
||||||
|
assert.equal(creates, 2); // the first pass only; nothing is applied again on "could not ask"
|
||||||
|
// A failure that is not a timeout may be one consumer's alone: the consumers after it are still
|
||||||
|
// asked on every pass.
|
||||||
|
assert.ok(askedFor.filter((as) => as === "b").length >= 2, askedFor.join(","));
|
||||||
|
assert.ok(logged.some((l) => l.includes("Command failed: tool --password ***")));
|
||||||
|
assert.ok(!logged.some((l) => l.includes("s3cret-pw")));
|
||||||
|
|
||||||
|
// A check that never answers counts as could-not-ask once its time is up.
|
||||||
|
let hungCreates = 0;
|
||||||
|
const hungLogged: string[] = [];
|
||||||
|
console.error = (...a: unknown[]) => hungLogged.push(a.join(" "));
|
||||||
|
try {
|
||||||
|
const stopHung = runProvisioner(
|
||||||
|
"cache",
|
||||||
|
{
|
||||||
|
async create() {
|
||||||
|
hungCreates++;
|
||||||
|
},
|
||||||
|
async remove() {},
|
||||||
|
holds: () => new Promise<boolean>(() => {}),
|
||||||
|
},
|
||||||
|
{ receives, everyMs: 5, verifyEveryMs: 30, holdsTimeoutMs: 20 },
|
||||||
|
);
|
||||||
|
await new Promise((r) => setTimeout(r, 200));
|
||||||
|
stopHung();
|
||||||
|
} finally {
|
||||||
|
console.error = original;
|
||||||
|
}
|
||||||
|
assert.equal(hungCreates, 2);
|
||||||
|
assert.ok(hungLogged.some((l) => l.includes("no answer within 20ms")));
|
||||||
|
});
|
||||||
|
|
||||||
|
test("provisioner stops a verify pass at a timeout, and retries a failed re-apply at the next check", async () => {
|
||||||
|
const dir = await mkdtemp(join(tmpdir(), "prov-retry-"));
|
||||||
|
await writeFile(join(dir, "a.secret"), "pa");
|
||||||
|
await writeFile(join(dir, "b.secret"), "pb");
|
||||||
|
const receives = join(dir, "cache.json");
|
||||||
|
await writeFile(
|
||||||
|
receives,
|
||||||
|
JSON.stringify({ requirement: "cache", given: [{ as: "a", secret: join(dir, "a.secret") }, { as: "b", secret: join(dir, "b.secret") }] }),
|
||||||
|
);
|
||||||
|
const original = console.error;
|
||||||
|
console.error = () => {};
|
||||||
|
try {
|
||||||
|
// A timeout on "a" means "b" is not asked in that pass.
|
||||||
|
const askedFor: string[] = [];
|
||||||
|
const stopHung = runProvisioner(
|
||||||
|
"cache",
|
||||||
|
{
|
||||||
|
async create() {},
|
||||||
|
async remove() {},
|
||||||
|
holds(p) {
|
||||||
|
askedFor.push(p.as);
|
||||||
|
return p.as === "a" ? new Promise<boolean>(() => {}) : Promise.resolve(true);
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{ receives, everyMs: 5, verifyEveryMs: 30, holdsTimeoutMs: 10 },
|
||||||
|
);
|
||||||
|
await new Promise((r) => setTimeout(r, 150));
|
||||||
|
stopHung();
|
||||||
|
assert.ok(askedFor.length >= 2 && askedFor.every((as) => as === "a"), askedFor.join(","));
|
||||||
|
|
||||||
|
// A lost consumer whose re-apply fails is retried at the next check, not held back by the brake,
|
||||||
|
// and the retry that succeeds counts as the first re-apply, not the second.
|
||||||
|
let held = true;
|
||||||
|
let createCalls = 0;
|
||||||
|
const logged: string[] = [];
|
||||||
|
console.error = (...a: unknown[]) => logged.push(a.join(" "));
|
||||||
|
const stop = runProvisioner(
|
||||||
|
"cache",
|
||||||
|
{
|
||||||
|
async create(p) {
|
||||||
|
if (p.as !== "a") return;
|
||||||
|
createCalls++;
|
||||||
|
if (createCalls === 2) throw new Error("write refused"); // the first re-apply fails
|
||||||
|
if (createCalls === 3) held = true; // the retry works
|
||||||
|
},
|
||||||
|
async remove() {},
|
||||||
|
async holds(p) {
|
||||||
|
return p.as === "a" ? held : true;
|
||||||
|
},
|
||||||
|
},
|
||||||
|
{ receives, everyMs: 5, verifyEveryMs: 20 },
|
||||||
|
);
|
||||||
|
await waitFor(() => createCalls === 1, 2000);
|
||||||
|
held = false; // the backend loses "a"
|
||||||
|
await waitFor(() => createCalls === 3, 2000);
|
||||||
|
await new Promise((r) => setTimeout(r, 100));
|
||||||
|
stop();
|
||||||
|
assert.equal(createCalls, 3);
|
||||||
|
assert.ok(!logged.some((l) => l.includes("create does not produce what holds checks")), logged.join("\n"));
|
||||||
|
} finally {
|
||||||
|
console.error = original;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
test("modules can SERVE: a real async tool, loaded and invoked over the broker", async () => {
|
test("modules can SERVE: a real async tool, loaded and invoked over the broker", async () => {
|
||||||
resetTools();
|
resetTools();
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user