forked from snxraven/p2ns
Updates
This commit is contained in:
@@ -64,7 +64,7 @@ function isDnsPassUsable(pass) {
|
||||
const base = pass.base;
|
||||
if (!base) return false;
|
||||
if (base.closed) return false;
|
||||
// base.closing is transient during replication; only treat as unusable when fully closed
|
||||
if (base.closing) return false;
|
||||
return true;
|
||||
}
|
||||
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
/**
|
||||
* Serializes all Autopass / Autobase operations on the master corestore.
|
||||
* Concurrent list(), createInvite(), add(), and remove() — plus HyperDB autoUpdate
|
||||
* during replication — cause "Atomic state must flush to parent" on Hypercore 11.
|
||||
* Concurrent list(), createInvite(), add(), and remove() — plus live replication
|
||||
* during store.replicate() — cause "Atomic state must flush" / SESSION_CLOSED.
|
||||
*/
|
||||
|
||||
const state = require('../infrastructure/state');
|
||||
@@ -9,13 +9,18 @@ const { logDebug, logWarn } = require('../infrastructure/logger');
|
||||
|
||||
let chain = Promise.resolve();
|
||||
|
||||
const ATOMIC_ERR_RE = /Atomic state must flush|SESSION_CLOSED|closing core/i;
|
||||
const OP_GAP_MS = parseInt(process.env.DNS_PASS_OP_GAP_MS || '75', 10);
|
||||
const SYNC_RETRIES = parseInt(process.env.DNS_PASS_SYNC_RETRIES || '4', 10);
|
||||
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);
|
||||
|
||||
function isAtomicDnsPassError(err) {
|
||||
function isTerminalDnsPassError(err) {
|
||||
const msg = err && err.message ? err.message : String(err);
|
||||
return ATOMIC_ERR_RE.test(msg);
|
||||
return TERMINAL_ERR_RE.test(msg);
|
||||
}
|
||||
|
||||
function isAtomicFlushError(err) {
|
||||
const msg = err && err.message ? err.message : String(err);
|
||||
return ATOMIC_FLUSH_RE.test(msg);
|
||||
}
|
||||
|
||||
function sleep(ms) {
|
||||
@@ -23,49 +28,19 @@ function sleep(ms) {
|
||||
}
|
||||
|
||||
/**
|
||||
* Bring Autobase + HyperDB view in sync before a read/write (reduces atomic flush races).
|
||||
* @param {object} pass - Autopass instance
|
||||
* @returns {Promise<boolean>} false if sync could not complete (busy)
|
||||
* Wait for Autopass to be open. Does not call base.update() (that races with replication).
|
||||
* @param {object} pass
|
||||
*/
|
||||
async function syncDnsPassView(pass) {
|
||||
async function ensureDnsPassOpen(pass) {
|
||||
if (!pass || pass.closed) return false;
|
||||
try {
|
||||
await pass.ready();
|
||||
const base = pass.base;
|
||||
if (!base || base.closed || base.closing) return false;
|
||||
await base.update();
|
||||
const view = base.view;
|
||||
if (view && !view.closed && typeof view.update === 'function') {
|
||||
view.update();
|
||||
}
|
||||
return true;
|
||||
} catch (err) {
|
||||
if (isAtomicDnsPassError(err)) {
|
||||
logDebug('DnsPassQueue', `View sync deferred (core busy): ${err.message}`);
|
||||
return false;
|
||||
}
|
||||
throw err;
|
||||
await pass.ready();
|
||||
const base = pass.base;
|
||||
if (!base || base.closed) return false;
|
||||
if (base.closing) {
|
||||
await sleep(200);
|
||||
if (base.closed || base.closing) return false;
|
||||
}
|
||||
}
|
||||
|
||||
async function runWithDnsPassRetry(pass, operation) {
|
||||
let lastErr;
|
||||
for (let attempt = 0; attempt < SYNC_RETRIES; attempt++) {
|
||||
if (attempt > 0) {
|
||||
await sleep(100 * attempt);
|
||||
}
|
||||
await syncDnsPassView(pass);
|
||||
try {
|
||||
return await operation();
|
||||
} catch (err) {
|
||||
lastErr = err;
|
||||
if (!isAtomicDnsPassError(err) || attempt === SYNC_RETRIES - 1) {
|
||||
throw err;
|
||||
}
|
||||
logWarn('DnsPassQueue', `dnsPass op retry ${attempt + 1}/${SYNC_RETRIES}: ${err.message}`);
|
||||
}
|
||||
}
|
||||
throw lastErr;
|
||||
return true;
|
||||
}
|
||||
|
||||
function enqueueDnsPass(operation) {
|
||||
@@ -93,16 +68,41 @@ function whenDnsPassIdle() {
|
||||
return chain;
|
||||
}
|
||||
|
||||
/**
|
||||
* createInvite must not be retried after SESSION_CLOSED — retries worsen core state.
|
||||
*/
|
||||
async function createInvite(pass, opts) {
|
||||
return enqueueDnsPass(() =>
|
||||
runWithDnsPassRetry(pass, () => pass.createInvite(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 await pass.createInvite(opts);
|
||||
}
|
||||
throw err;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/** @deprecated No-op — avoid base.update() during replication */
|
||||
async function syncDnsPassView(pass) {
|
||||
return ensureDnsPassOpen(pass);
|
||||
}
|
||||
|
||||
function isAtomicDnsPassError(err) {
|
||||
return isAtomicFlushError(err) || isTerminalDnsPassError(err);
|
||||
}
|
||||
|
||||
/**
|
||||
* Autopass createInvite() returns a z32 invite string — send as-is, not .toString('hex').
|
||||
* @param {string|Buffer} inv
|
||||
* @returns {string}
|
||||
*/
|
||||
function inviteToWire(inv) {
|
||||
if (typeof inv === 'string') return inv;
|
||||
@@ -115,50 +115,48 @@ function invitePreview(inv) {
|
||||
return s.length > 24 ? `${s.slice(0, 24)}...` : s;
|
||||
}
|
||||
|
||||
async function runSerialized(pass, fn) {
|
||||
return enqueueDnsPass(async () => {
|
||||
if (!(await ensureDnsPassOpen(pass))) {
|
||||
throw new Error('dnsPass not open');
|
||||
}
|
||||
return fn();
|
||||
});
|
||||
}
|
||||
|
||||
async function dnsPassAdd(pass, key, value, file) {
|
||||
return enqueueDnsPass(() =>
|
||||
runWithDnsPassRetry(pass, () => pass.add(key, value, file))
|
||||
);
|
||||
return runSerialized(pass, () => pass.add(key, value, file));
|
||||
}
|
||||
|
||||
async function dnsPassRemove(pass, key) {
|
||||
return enqueueDnsPass(() =>
|
||||
runWithDnsPassRetry(pass, () => pass.remove(key))
|
||||
);
|
||||
return runSerialized(pass, () => pass.remove(key));
|
||||
}
|
||||
|
||||
async function dnsPassGet(pass, key) {
|
||||
return enqueueDnsPass(() =>
|
||||
runWithDnsPassRetry(pass, () => pass.get(key))
|
||||
);
|
||||
return runSerialized(pass, () => pass.get(key));
|
||||
}
|
||||
|
||||
/**
|
||||
* List all Autopass records (serialized with other dnsPass ops).
|
||||
* @param {object} pass - Autopass instance
|
||||
* @returns {Promise<Array<{key: string, value: string}>>}
|
||||
*/
|
||||
async function listAllEntries(pass) {
|
||||
return enqueueDnsPass(async () =>
|
||||
runWithDnsPassRetry(pass, async () => {
|
||||
const entries = [];
|
||||
const stream = pass.list();
|
||||
for await (const entry of stream) {
|
||||
entries.push({
|
||||
key: entry.key.toString('utf8'),
|
||||
value: entry.value.toString('utf8')
|
||||
});
|
||||
}
|
||||
return entries;
|
||||
})
|
||||
);
|
||||
return runSerialized(pass, async () => {
|
||||
const entries = [];
|
||||
const stream = pass.list();
|
||||
for await (const entry of stream) {
|
||||
entries.push({
|
||||
key: entry.key.toString('utf8'),
|
||||
value: entry.value.toString('utf8')
|
||||
});
|
||||
}
|
||||
return entries;
|
||||
});
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
enqueueDnsPass,
|
||||
whenDnsPassIdle,
|
||||
ensureDnsPassOpen,
|
||||
syncDnsPassView,
|
||||
isAtomicDnsPassError,
|
||||
isTerminalDnsPassError,
|
||||
createInvite,
|
||||
inviteToWire,
|
||||
invitePreview,
|
||||
|
||||
Reference in New Issue
Block a user