Updates
This commit is contained in:
@@ -1,7 +1,7 @@
|
||||
/**
|
||||
* Serializes all Autopass / Autobase operations on the master corestore.
|
||||
* Concurrent list(), createInvite(), add(), and remove() — plus live replication
|
||||
* during store.replicate() — cause "Atomic state must flush" / SESSION_CLOSED.
|
||||
* Autopass tests use base.replicate(connection), not Corestore.replicate() on every core.
|
||||
* createInvite() races HyperDB autoUpdate unless Autopass replication is suspended briefly.
|
||||
*/
|
||||
|
||||
const state = require('../infrastructure/state');
|
||||
@@ -27,10 +27,6 @@ function sleep(ms) {
|
||||
return new Promise((resolve) => setTimeout(resolve, ms));
|
||||
}
|
||||
|
||||
/**
|
||||
* Wait for Autopass to be open. Does not call base.update() (that races with replication).
|
||||
* @param {object} pass
|
||||
*/
|
||||
async function ensureDnsPassOpen(pass) {
|
||||
if (!pass || pass.closed) return false;
|
||||
await pass.ready();
|
||||
@@ -63,36 +59,59 @@ function enqueueDnsPass(operation) {
|
||||
return run;
|
||||
}
|
||||
|
||||
/** @deprecated Use enqueueDnsPass; kept for compatibility */
|
||||
function whenDnsPassIdle() {
|
||||
return chain;
|
||||
}
|
||||
|
||||
/**
|
||||
* createInvite must not be retried after SESSION_CLOSED — retries worsen core state.
|
||||
* Pause Autopass swarm/store replication while mutating the view (createInvite).
|
||||
* Matches the isolation Autopass expects when not competing with live replication.
|
||||
*/
|
||||
async function withAutopassReplicationPaused(pass, fn) {
|
||||
const canSuspend = typeof pass.suspend === 'function' && pass.swarm;
|
||||
if (canSuspend) {
|
||||
try {
|
||||
await pass.suspend();
|
||||
} catch (err) {
|
||||
logDebug('DnsPassQueue', `suspend skipped: ${err.message}`);
|
||||
}
|
||||
}
|
||||
try {
|
||||
return await fn();
|
||||
} finally {
|
||||
if (canSuspend) {
|
||||
try {
|
||||
await pass.resume();
|
||||
} catch (err) {
|
||||
logWarn('DnsPassQueue', `resume after invite op failed: ${err.message}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async function createInvite(pass, opts) {
|
||||
return enqueueDnsPass(async () => {
|
||||
if (!(await ensureDnsPassOpen(pass))) {
|
||||
throw new Error('dnsPass not open');
|
||||
}
|
||||
try {
|
||||
return await pass.createInvite(opts);
|
||||
} catch (err) {
|
||||
if (isAtomicFlushError(err) && !isTerminalDnsPassError(err)) {
|
||||
logWarn('DnsPassQueue', `createInvite atomic race, one delayed retry: ${err.message}`);
|
||||
await sleep(400);
|
||||
if (!(await ensureDnsPassOpen(pass))) {
|
||||
throw err;
|
||||
}
|
||||
return withAutopassReplicationPaused(pass, async () => {
|
||||
try {
|
||||
return await pass.createInvite(opts);
|
||||
} catch (err) {
|
||||
if (isAtomicFlushError(err) && !isTerminalDnsPassError(err)) {
|
||||
logWarn('DnsPassQueue', `createInvite atomic race, one delayed retry: ${err.message}`);
|
||||
await sleep(500);
|
||||
if (!(await ensureDnsPassOpen(pass))) {
|
||||
throw err;
|
||||
}
|
||||
return await pass.createInvite(opts);
|
||||
}
|
||||
throw err;
|
||||
}
|
||||
throw err;
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
/** @deprecated No-op — avoid base.update() during replication */
|
||||
async function syncDnsPassView(pass) {
|
||||
return ensureDnsPassOpen(pass);
|
||||
}
|
||||
@@ -101,9 +120,6 @@ function isAtomicDnsPassError(err) {
|
||||
return isAtomicFlushError(err) || isTerminalDnsPassError(err);
|
||||
}
|
||||
|
||||
/**
|
||||
* Autopass createInvite() returns a z32 invite string — send as-is, not .toString('hex').
|
||||
*/
|
||||
function inviteToWire(inv) {
|
||||
if (typeof inv === 'string') return inv;
|
||||
if (Buffer.isBuffer(inv)) return inv.toString('utf8');
|
||||
|
||||
Reference in New Issue
Block a user