forked from snxraven/p2ns
Updates
This commit is contained in:
@@ -592,6 +592,11 @@ async function doAutoVotes() {
|
||||
return;
|
||||
}
|
||||
|
||||
if (state.dnsPassWriteInProgress > 0) {
|
||||
logDebug('Core', 'Skipping doAutoVotes — invite or dnsPass write in progress');
|
||||
return;
|
||||
}
|
||||
|
||||
const pass = state.dnsPass;
|
||||
if (!isDnsPassUsable(pass)) {
|
||||
logDebug('Core', 'dnsPass is not usable, skipping doAutoVotes');
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
/**
|
||||
* Serializes all Autopass / Autobase operations on the master corestore.
|
||||
* Autopass tests use base.replicate(connection), not Corestore.replicate() on every core.
|
||||
* createInvite() races HyperDB autoUpdate unless Autopass replication is suspended briefly.
|
||||
* Use base.replicate(connection) for replication — not Corestore.replicate() on every core.
|
||||
* Do NOT call pass.suspend() around createInvite — member.flushed() needs the swarm.
|
||||
*/
|
||||
|
||||
const state = require('../infrastructure/state');
|
||||
@@ -12,6 +12,7 @@ let chain = Promise.resolve();
|
||||
const ATOMIC_FLUSH_RE = /Atomic state must flush/i;
|
||||
const TERMINAL_ERR_RE = /SESSION_CLOSED|closing core/i;
|
||||
const OP_GAP_MS = parseInt(process.env.DNS_PASS_OP_GAP_MS || '50', 10);
|
||||
const CREATE_INVITE_TIMEOUT_MS = parseInt(process.env.CREATE_INVITE_TIMEOUT_MS || '20000', 10);
|
||||
|
||||
function isTerminalDnsPassError(err) {
|
||||
const msg = err && err.message ? err.message : String(err);
|
||||
@@ -27,6 +28,15 @@ function sleep(ms) {
|
||||
return new Promise((resolve) => setTimeout(resolve, ms));
|
||||
}
|
||||
|
||||
function withTimeout(promise, ms, label) {
|
||||
return Promise.race([
|
||||
promise,
|
||||
sleep(ms).then(() => {
|
||||
throw new Error(`${label} timed out after ${ms}ms`);
|
||||
})
|
||||
]);
|
||||
}
|
||||
|
||||
async function ensureDnsPassOpen(pass) {
|
||||
if (!pass || pass.closed) return false;
|
||||
await pass.ready();
|
||||
@@ -63,30 +73,8 @@ function whenDnsPassIdle() {
|
||||
return chain;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
function isInviteCreationActive() {
|
||||
return (state.dnsPassWriteInProgress || 0) > 0;
|
||||
}
|
||||
|
||||
async function createInvite(pass, opts) {
|
||||
@@ -94,21 +82,28 @@ async function createInvite(pass, opts) {
|
||||
if (!(await ensureDnsPassOpen(pass))) {
|
||||
throw new Error('dnsPass not open');
|
||||
}
|
||||
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);
|
||||
logDebug('DnsPassQueue', 'createInvite: starting');
|
||||
try {
|
||||
return await withTimeout(
|
||||
pass.createInvite(opts),
|
||||
CREATE_INVITE_TIMEOUT_MS,
|
||||
'createInvite'
|
||||
);
|
||||
} 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;
|
||||
}
|
||||
throw err;
|
||||
return await withTimeout(
|
||||
pass.createInvite(opts),
|
||||
CREATE_INVITE_TIMEOUT_MS,
|
||||
'createInvite retry'
|
||||
);
|
||||
}
|
||||
});
|
||||
throw err;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
@@ -171,6 +166,7 @@ module.exports = {
|
||||
whenDnsPassIdle,
|
||||
ensureDnsPassOpen,
|
||||
syncDnsPassView,
|
||||
isInviteCreationActive,
|
||||
isAtomicDnsPassError,
|
||||
isTerminalDnsPassError,
|
||||
createInvite,
|
||||
|
||||
Reference in New Issue
Block a user