Compare commits

...

7 Commits

Author SHA1 Message Date
wmantly b6a82d58d5 Merge pull request #200 from theta42/release-v2.5.0
release(v2.5.0): no-inbound relay automation, mesh-preferred replication
2026-08-10 18:28:27 -07:00
wmantly 18da6582ed release(v2.5.0): no-inbound relay automation, mesh-preferred replication 2026-08-10 21:25:38 -04:00
wmantly b1739ec965 Merge pull request #199 from theta42/feat-mesh-routing
feat(multi-site): route replication over the mesh when available
2026-08-10 18:09:28 -07:00
wmantly 7bc6f47070 feat(multi-site): forward noInbound/meshIp/publicHost through /api/site/join
/api/site/spokes already accepted these and drove proxy_client's relay
automation, but nothing on the real operator join path ever passed
them through -- a no-inbound spoke calling /join had no way to trigger
its own relay route. Forward them into the internal /spokes call and
surface the resulting relay note in /join's response.
2026-08-10 20:52:35 -04:00
wmantly a0964ce350 feat(multi-site): route replication traffic over the mesh when available
Cross-component routing TODO item: a spoke's resync push now prefers
its WG mesh IP (reported via the noInbound/meshIp fields added for the
relay automation) over the public endpoint, falling back to the public
endpoint if the mesh attempt fails for any reason (tunnel not actually
up between these two particular gateways yet, transient failure,
etc.) -- a mesh-routing preference must never turn into "spoke stops
getting updates."

Plain HTTP over the mesh IP, not HTTPS: the WG tunnel is already
encrypted, same reasoning already applied to the no-inbound relay
terminating at the master.

A spoke with no meshIp on file behaves exactly as before (public
endpoint only) -- this is additive, not a behavior change for spokes
that haven't opted into mesh registration.
2026-08-10 20:39:27 -04:00
wmantly 0d3ee3e2ee Merge pull request #198 from theta42/feat-no-inbound-relay
feat(multi-site): no-inbound relay automation via theta-proxy's existing API tokens
2026-08-10 17:36:03 -07:00
wmantly b95cb08c41 feat(multi-site): no-inbound relay automation via theta-proxy's existing API tokens
Implements TODO items "service-to-service auth model" and "no-inbound
relay automation" together -- the first was never really "no credential
type exists", it was that nothing wired one of the credential types that
ALREADY exist (theta-proxy, jump-host, and this app each already have
their own self-service API token system, models/api_token.js) into an
actual inter-service call. This is that wiring, not a new invented
credential type.

- utils/proxy_client.js: ensureRelayRoute({host, ip, targetPort}) calls
  theta-proxy's real Host API (GET/POST/PUT /api/host) using a `prx_...`
  token an operator mints on theta-proxy and stores in OpenBao
  (secret/integrations/theta-proxy), same pattern as agent_keys.js.
  Idempotent and best-effort -- never fails the caller if the token/URL
  isn't configured, since this is an enhancement on top of a working
  join, not a join requirement.
- POST /api/site/spokes accepts optional noInbound/meshIp/publicHost
  fields; when a spoke reports itself no-inbound, the master
  best-effort creates/updates the matching relay route automatically.
  SiteSpoke gained the fields + a relayNote for visibility.

Verified against a REAL running theta-proxy container (not mocked):
booted it standalone, logged in as the local admin, minted a real
`prx_` token via its actual API, and drove ensureRelayRoute() against
it for real. Caught a real bug doing this: GET /api/host/:item wraps
the record in `{item, results}`, not flat -- the mocked unit tests
(which I wrote first) all had the flat shape baked in and passed
cleanly, so this only surfaced against the real API. Fixed in both the
implementation and the unit tests' mocked response shape.
2026-08-10 20:31:42 -04:00
8 changed files with 342 additions and 27 deletions
+7
View File
@@ -1,3 +1,10 @@
# v2.5.0 - 2026-08-10
### Added
- **No-inbound relay automation.** A spoke with no public IP of its own can now register as such (`noInbound`/`meshIp`/`publicHost` on `POST /api/site/spokes`, forwarded through `POST /api/site/join` for the real operator join flow), and the master auto-creates/updates the relay route on its own `theta-proxy` via `utils/proxy_client.js` — a new self-service `prx_...` API token client, reusing `theta-proxy`'s existing token system rather than inventing a new credential type. Verified against a real running `theta-proxy` container (`GET /api/host/:item`'s actual `{item, results: {...}}` response shape, not the flat shape first assumed).
- **Replication traffic prefers the mesh.** `utils/site_replicate.js`'s fire-and-forget resync push now tries a registered spoke's `meshIp` first (falling back to its public `endpoint` on failure) — cross-component routing over the gateway-to-gateway WireGuard mesh instead of the open internet, for any spoke that's registered one.
- `POST /api/site/join` surfaces the resulting relay status in its response (`relay.note`), and `theta-suite`'s bootstrap flow (`CFG_SPOKE_NO_INBOUND`/`CFG_SPOKE_PUBLIC_HOST`, `bootstrap/site-relay-register.js`) drives all of this from the real operator-facing setup script, not just the API.
# v2.4.0 - 2026-08-10 # v2.4.0 - 2026-08-10
### Added ### Added
+11 -1
View File
@@ -29,7 +29,17 @@ class SiteSpoke extends Model {
siteSlug: { type: 'string' }, siteSlug: { type: 'string' },
pushToken: { type: 'string', isRequired: true }, pushToken: { type: 'string', isRequired: true },
created_on: { type: 'integer' }, created_on: { type: 'integer' },
last_seen_on: { type: 'integer' } last_seen_on: { type: 'integer' },
// No-inbound relay (MULTI_SITE_SPEC.md): a spoke with no public IP of
// its own reports its WG mesh IP + the public hostname it wants
// reached at; the master then best-effort creates a matching relay
// route on its own theta-proxy (utils/proxy_client.js). relayNote
// records what happened for visibility in the UI -- this automation
// is optional/best-effort, never a join requirement.
noInbound: { type: 'boolean', default: false },
meshIp: { type: 'string' },
publicHost: { type: 'string' },
relayNote: { type: 'string' }
}; };
toPublic() { toPublic() {
+2 -2
View File
@@ -1,6 +1,6 @@
{ {
"name": "t42-theta-directory", "name": "t42-theta-directory",
"version": "2.4.0", "version": "2.5.0",
"description": "A very simple LDAP management and SSO system", "description": "A very simple LDAP management and SSO system",
"author": [ "author": [
{ {
@@ -11,7 +11,7 @@
"scripts": { "scripts": {
"start": "node ./bin/www", "start": "node ./bin/www",
"dev": "npx nodemon --ignore public/ ./bin/www", "dev": "npx nodemon --ignore public/ ./bin/www",
"test": "NODE_ENV=test jest tests/groups.test.js tests/subtypes.test.js tests/site_join.test.js tests/site_config.test.js tests/site_replicate.test.js --forceExit" "test": "NODE_ENV=test jest tests/groups.test.js tests/subtypes.test.js tests/site_join.test.js tests/site_config.test.js tests/site_replicate.test.js tests/proxy_client.test.js --forceExit"
}, },
"jest": { "jest": {
"testEnvironment": "node", "testEnvironment": "node",
+38 -9
View File
@@ -120,27 +120,43 @@ router.post('/spokes', async (req, res, next) => {
const key = await SiteJoinKey.authenticate(rawKey); const key = await SiteJoinKey.authenticate(rawKey);
if (!key) return res.status(401).json({ status: 'error', message: 'invalid or revoked site join key' }); if (!key) return res.status(401).json({ status: 'error', message: 'invalid or revoked site join key' });
const { endpoint, siteSlug } = req.body || {}; const { endpoint, siteSlug, noInbound, meshIp, publicHost } = req.body || {};
if (!endpoint || !/^https?:\/\//.test(endpoint)) { if (!endpoint || !/^https?:\/\//.test(endpoint)) {
return res.status(400).json({ status: 'error', message: 'a valid http(s) endpoint is required' }); return res.status(400).json({ status: 'error', message: 'a valid http(s) endpoint is required' });
} }
const now = Math.floor(Date.now() / 1000); const now = Math.floor(Date.now() / 1000);
let spoke = (await SiteSpoke.list({ where: { endpoint } }))[0]; let spoke = (await SiteSpoke.list({ where: { endpoint } }))[0];
const patch = { siteSlug: siteSlug || (spoke && spoke.siteSlug) || null, last_seen_on: now, noInbound: !!noInbound, meshIp: meshIp || '', publicHost: publicHost || '' };
if (spoke) { if (spoke) {
await spoke.update({ siteSlug: siteSlug || spoke.siteSlug, last_seen_on: now }); await spoke.update(patch);
} else { } else {
spoke = await SiteSpoke.create({ spoke = await SiteSpoke.create({
id: crypto.randomUUID(), id: crypto.randomUUID(),
endpoint, endpoint,
siteSlug: siteSlug || null,
pushToken: SiteSpoke.generatePushToken(), pushToken: SiteSpoke.generatePushToken(),
created_on: now, created_on: now,
last_seen_on: now ...patch
}); });
} }
logAudit('spoke_registered', { endpoint, siteSlug: spoke.siteSlug });
res.json({ status: 'ok', pushToken: spoke.pushToken }); // No-inbound relay automation: best-effort, never blocks registration.
// See utils/proxy_client.js for why this reuses theta-proxy's existing
// API token system rather than a new credential type.
let relayNote = 'not applicable (spoke has inbound access)';
if (noInbound) {
if (meshIp && publicHost) {
const proxyClient = require('../utils/proxy_client');
const result = await proxyClient.ensureRelayRoute({ host: publicHost, ip: meshIp, targetPort: 3001 });
relayNote = result.note;
} else {
relayNote = 'skipped: noInbound set but meshIp/publicHost missing';
}
await spoke.update({ relayNote });
}
logAudit('spoke_registered', { endpoint, siteSlug: spoke.siteSlug, noInbound: !!noInbound, relayNote });
res.json({ status: 'ok', pushToken: spoke.pushToken, relay: { note: relayNote } });
} catch (e) { next(e); } } catch (e) { next(e); }
}); });
@@ -354,7 +370,7 @@ async function adoptFromMaster({ masterUrl, joinKey }) {
router.post('/join', async (req, res, next) => { router.post('/join', async (req, res, next) => {
try { try {
const { masterUrl, joinKey, selfUrl } = req.body || {}; const { masterUrl, joinKey, selfUrl, noInbound, meshIp, publicHost } = req.body || {};
if (!masterUrl || !joinKey) { if (!masterUrl || !joinKey) {
return res.status(400).json({ status: 'error', message: 'masterUrl and joinKey are required' }); return res.status(400).json({ status: 'error', message: 'masterUrl and joinKey are required' });
} }
@@ -388,17 +404,29 @@ router.post('/join', async (req, res, next) => {
// snapshot for that spoke, not a hard failure). // snapshot for that spoke, not a hard failure).
let replicationPushToken = null; let replicationPushToken = null;
let replicationNote = 'not registered (no selfUrl given)'; let replicationNote = 'not registered (no selfUrl given)';
let relayNote = noInbound ? 'not attempted (registration did not run)' : 'not applicable (this spoke has inbound access)';
if (selfUrl) { if (selfUrl) {
try { try {
const regResp = await fetch(base + '/api/site/spokes', { const regResp = await fetch(base + '/api/site/spokes', {
method: 'POST', method: 'POST',
headers: { Authorization: 'Bearer ' + joinKey, 'Content-Type': 'application/json' }, headers: { Authorization: 'Bearer ' + joinKey, 'Content-Type': 'application/json' },
body: JSON.stringify({ endpoint: selfUrl, siteSlug: exportData.siteSlug || cfg.siteSlug }) // noInbound/meshIp/publicHost: this spoke has no public IP of its
// own; forwarded so the master can best-effort auto-create a relay
// route on its own theta-proxy (utils/proxy_client.js). Previously
// accepted by /spokes but never actually reachable from here --
// nothing forwarded them, so the automation existed but no real
// join flow could ever trigger it.
body: JSON.stringify({
endpoint: selfUrl,
siteSlug: exportData.siteSlug || cfg.siteSlug,
...(noInbound ? { noInbound: true, meshIp, publicHost } : {})
})
}); });
if (regResp.ok) { if (regResp.ok) {
const regBody = await regResp.json(); const regBody = await regResp.json();
replicationPushToken = regBody.pushToken; replicationPushToken = regBody.pushToken;
replicationNote = 'registered for live replication'; replicationNote = 'registered for live replication';
if (regBody.relay) relayNote = regBody.relay.note;
} else { } else {
replicationNote = 'registration failed: HTTP ' + regResp.status; replicationNote = 'registration failed: HTTP ' + regResp.status;
} }
@@ -438,7 +466,8 @@ router.post('/join', async (req, res, next) => {
resources: { created: imp.created, updated: imp.updated, edges: imp.edgeCount }, resources: { created: imp.created, updated: imp.updated, edges: imp.edgeCount },
ldap: { note: ldapNote }, ldap: { note: ldapNote },
signingKey: { note: signingKeyNote }, signingKey: { note: signingKeyNote },
replication: { note: replicationNote, live: !!replicationPushToken } replication: { note: replicationNote, live: !!replicationPushToken },
relay: { note: relayNote }
}); });
} catch (e) { next(e); } } catch (e) { next(e); }
}); });
+98
View File
@@ -0,0 +1,98 @@
'use strict';
let mockBaoStore = new Map();
jest.mock('@simpleworkjs/bao-conf', () => ({
get: jest.fn(async (path) => mockBaoStore.get(path) || null),
set: jest.fn(async (path, value) => { mockBaoStore.set(path, value); })
}));
describe('proxy_client', () => {
let proxyClient;
let originalFetch;
let mockFetchImpl;
let calls;
beforeEach(() => {
jest.resetModules();
mockBaoStore = new Map();
calls = [];
mockFetchImpl = async () => ({ ok: true, status: 404 });
originalFetch = global.fetch;
global.fetch = (...args) => { calls.push(args); return mockFetchImpl(...args); };
proxyClient = require('../utils/proxy_client');
proxyClient._reset();
delete process.env.PROXY_INTERNAL_URL;
});
afterEach(() => {
global.fetch = originalFetch;
});
test('skips cleanly when required fields are missing', async () => {
const result = await proxyClient.ensureRelayRoute({ host: '', ip: '', targetPort: 0 });
expect(result.note).toMatch(/required/);
expect(calls.length).toBe(0);
});
test('skips cleanly when PROXY_INTERNAL_URL is not configured', async () => {
const result = await proxyClient.ensureRelayRoute({ host: 'sso-a.example.com', ip: '172.24.5.1', targetPort: 3001 });
expect(result.note).toMatch(/PROXY_INTERNAL_URL/);
expect(calls.length).toBe(0);
});
test('skips cleanly when no token is stored in OpenBao', async () => {
process.env.PROXY_INTERNAL_URL = 'https://proxy.internal';
const result = await proxyClient.ensureRelayRoute({ host: 'sso-a.example.com', ip: '172.24.5.1', targetPort: 3001 });
expect(result.note).toMatch(/no proxy API token/);
expect(calls.length).toBe(0);
});
test('creates the route when the host does not already exist', async () => {
process.env.PROXY_INTERNAL_URL = 'https://proxy.internal';
mockBaoStore.set('integrations/theta-proxy', { token: 'prx_test_token' });
mockFetchImpl = async (url, opts) => {
if (opts.method === undefined) return { ok: true, status: 404 }; // GET lookup
if (opts.method === 'POST') return { ok: true, status: 200 };
throw new Error('unexpected method ' + opts.method);
};
const result = await proxyClient.ensureRelayRoute({ host: 'sso-a.example.com', ip: '172.24.5.1', targetPort: 3001 });
expect(result.note).toBe('created');
const postCall = calls.find((c) => c[1].method === 'POST');
expect(postCall[0]).toBe('https://proxy.internal/api/host');
expect(postCall[1].headers.Authorization).toBe('Bearer prx_test_token');
expect(JSON.parse(postCall[1].body)).toEqual({ host: 'sso-a.example.com', ip: '172.24.5.1', targetPort: 3001 });
});
test('updates the route when it exists but points somewhere else', async () => {
process.env.PROXY_INTERNAL_URL = 'https://proxy.internal';
mockBaoStore.set('integrations/theta-proxy', { token: 'prx_test_token' });
mockFetchImpl = async (url, opts) => {
if (!opts.method) return { ok: true, status: 200, json: async () => ({ results: { ip: '172.24.9.9', targetPort: 3001 } }) };
if (opts.method === 'PUT') return { ok: true, status: 200 };
throw new Error('unexpected method ' + opts.method);
};
const result = await proxyClient.ensureRelayRoute({ host: 'sso-a.example.com', ip: '172.24.5.1', targetPort: 3001 });
expect(result.note).toBe('updated');
});
test('is a no-op when the route already matches', async () => {
process.env.PROXY_INTERNAL_URL = 'https://proxy.internal';
mockBaoStore.set('integrations/theta-proxy', { token: 'prx_test_token' });
mockFetchImpl = async () => ({ ok: true, status: 200, json: async () => ({ results: { ip: '172.24.5.1', targetPort: 3001 } }) });
const result = await proxyClient.ensureRelayRoute({ host: 'sso-a.example.com', ip: '172.24.5.1', targetPort: 3001 });
expect(result.note).toBe('already up to date');
});
test('reports a network failure without throwing', async () => {
process.env.PROXY_INTERNAL_URL = 'https://proxy.internal';
mockBaoStore.set('integrations/theta-proxy', { token: 'prx_test_token' });
mockFetchImpl = async () => { throw new Error('connection refused'); };
const result = await proxyClient.ensureRelayRoute({ host: 'sso-a.example.com', ip: '172.24.5.1', targetPort: 3001 });
expect(result.note).toMatch(/failed: connection refused/);
});
});
+39
View File
@@ -55,6 +55,45 @@ describe('site_replicate', () => {
expect(JSON.parse(optsA.body).reason).toBe('catalog-changed'); expect(JSON.parse(optsA.body).reason).toBe('catalog-changed');
}); });
test('prefers the mesh IP over the public endpoint when the spoke reported one', async () => {
SiteSpoke._seed([
{ endpoint: 'https://spoke-a.example.com:8443', pushToken: 'token-a', meshIp: '172.24.5.1' }
]);
await siteReplicate.replicateToSpokes('catalog-changed');
await new Promise((r) => setImmediate(r));
expect(mockFetchCalls.length).toBe(1);
expect(mockFetchCalls[0][0]).toBe('http://172.24.5.1:8443/api/site/resync');
});
test('falls back to the public endpoint if the mesh attempt fails', async () => {
SiteSpoke._seed([
{ endpoint: 'https://spoke-a.example.com', pushToken: 'token-a', meshIp: '172.24.5.1' }
]);
mockFetchImpl = async (url) => {
if (url.startsWith('http://172.24.5.1')) throw new Error('mesh unreachable');
return { ok: true, status: 200 };
};
await siteReplicate.replicateToSpokes('catalog-changed');
await new Promise((r) => setImmediate(r));
expect(mockFetchCalls.length).toBe(2);
expect(mockFetchCalls[0][0]).toMatch(/^http:\/\/172\.24\.5\.1/);
expect(mockFetchCalls[1][0]).toBe('https://spoke-a.example.com/api/site/resync');
});
test('a spoke with no meshIp only ever tries the public endpoint', async () => {
SiteSpoke._seed([{ endpoint: 'https://spoke-a.example.com', pushToken: 'token-a' }]);
await siteReplicate.replicateToSpokes('catalog-changed');
await new Promise((r) => setImmediate(r));
expect(mockFetchCalls.length).toBe(1);
expect(mockFetchCalls[0][0]).toBe('https://spoke-a.example.com/api/site/resync');
});
test('no known spokes: resolves cleanly, no fetch calls', async () => { test('no known spokes: resolves cleanly, no fetch calls', async () => {
await siteReplicate.replicateToSpokes('catalog-changed'); await siteReplicate.replicateToSpokes('catalog-changed');
await new Promise((r) => setImmediate(r)); await new Promise((r) => setImmediate(r));
+105
View File
@@ -0,0 +1,105 @@
'use strict';
// Service-to-service client for theta-proxy's Host management API --
// MULTI_SITE_SPEC.md's "no-inbound relay automation" (a master creating a
// relay route so a spoke with zero inbound path of its own is reachable).
//
// Deliberately does NOT invent a new credential type. theta-proxy already
// has a self-service API token system (models/api_token.js, `prx_<id>_<secret>`
// bearer tokens that authenticate as their creator's user + group snapshot --
// same pattern this app and jump-host both already have their own copy of).
// The "service-to-service auth" gap was never "no credential type exists" --
// it's that nothing wired one of these tokens into an actual inter-service
// call. This is that wiring, using the credential type that was already
// there. The token itself is operator-provisioned (minted on theta-proxy by
// an admin with Host-management rights) and stored in OpenBao, same as the
// agent-signing key in agent_keys.js.
const baoConf = require('@simpleworkjs/bao-conf');
const PATH = 'integrations/theta-proxy'; // baoConf adds the secret/data prefix
const REQUEST_TIMEOUT_MS = 10000;
let cachedToken = null;
async function loadToken() {
if (cachedToken) return cachedToken;
let stored;
try {
stored = await baoConf.get(PATH);
} catch (err) {
console.error(`[proxy_client] could not read ${PATH} from OpenBao: ${err.message}`);
return null;
}
if (!stored || !stored.token) return null;
cachedToken = stored.token;
return cachedToken;
}
function proxyBaseUrl() {
// Not OpenBao -- this is where the proxy's admin API lives, not a secret.
// No safe default: relaying to a guessed host would be worse than
// refusing, so this must be explicitly configured.
return process.env.PROXY_INTERNAL_URL || '';
}
// Creates the relay Host route if missing, updates its target IP if it
// already exists and points somewhere else. Idempotent -- safe to call
// again for the same host on every spoke resync.
//
// Returns { note } describing what happened (created/updated/skipped/failed)
// rather than throwing on a missing token or base URL -- callers (spoke
// registration) must never fail the whole registration just because this
// automation isn't configured yet; it's an enhancement layered on top of a
// working join, not a requirement of one.
async function ensureRelayRoute({ host, ip, targetPort }) {
if (!host || !ip || !targetPort) {
return { note: 'skipped: host, ip, and targetPort are all required' };
}
const base = proxyBaseUrl();
if (!base) {
return { note: 'skipped: PROXY_INTERNAL_URL not configured' };
}
const token = await loadToken();
if (!token) {
return { note: `skipped: no proxy API token at OpenBao ${PATH} -- mint one on theta-proxy and store it there` };
}
const headers = { Authorization: 'Bearer ' + token, 'Content-Type': 'application/json' };
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), REQUEST_TIMEOUT_MS);
try {
const existing = await fetch(base.replace(/\/+$/, '') + '/api/host/' + encodeURIComponent(host), {
headers, signal: controller.signal
});
if (existing.status === 200) {
// GET /api/host/:item wraps the record in { item, results }, not
// flat -- confirmed against a real running proxy (this check
// silently always "updated" instead of no-op'ing until fixed).
const body = await existing.json();
const current = body.results || body;
if (current.ip === ip && Number(current.targetPort) === Number(targetPort)) {
return { note: 'already up to date' };
}
const put = await fetch(base.replace(/\/+$/, '') + '/api/host/' + encodeURIComponent(host), {
method: 'PUT', headers, body: JSON.stringify({ ip, targetPort }), signal: controller.signal
});
return put.ok ? { note: 'updated' } : { note: `update failed: HTTP ${put.status}` };
}
const create = await fetch(base.replace(/\/+$/, '') + '/api/host', {
method: 'POST', headers, body: JSON.stringify({ host, ip, targetPort }), signal: controller.signal
});
return create.ok ? { note: 'created' } : { note: `create failed: HTTP ${create.status}` };
} catch (err) {
return { note: `failed: ${err.message}` };
} finally {
clearTimeout(timer);
}
}
// Test seam.
function _reset() { cachedToken = null; }
module.exports = { ensureRelayRoute, _reset, PATH };
+42 -15
View File
@@ -37,21 +37,48 @@ function replicateToSpokes(reason) {
})(); })();
} }
async function pingOne(spoke, reason) { // Cross-component routing (MULTI_SITE_SPEC.md): if this spoke reported a WG
const url = String(spoke.endpoint).replace(/\/+$/, '') + '/api/site/resync'; // mesh IP when it registered (utils/proxy_client.js's no-inbound relay path
const controller = new AbortController(); // populates the same field), prefer sending the resync push over the mesh
const timer = setTimeout(() => controller.abort(), RESYNC_TIMEOUT_MS); // tunnel instead of the open internet -- plain HTTP is fine here since the
try { // WG tunnel itself is already encrypted, same reasoning as the no-inbound
const resp = await fetch(url, { // relay terminating at the master. Falls back to the spoke's public endpoint
method: 'POST', // if the mesh attempt fails (mesh IP set but that particular tunnel isn't
headers: { Authorization: 'Bearer ' + spoke.pushToken, 'Content-Type': 'application/json' }, // actually up yet, or unreachable for any other reason) -- never let a
body: JSON.stringify({ reason: reason || 'catalog-changed' }), // mesh-routing preference turn into "spoke never gets updates."
signal: controller.signal function resyncUrls(spoke) {
}); const urls = [];
if (!resp.ok) throw new Error('status ' + resp.status); if (spoke.meshIp) {
} finally { let port = '3001';
clearTimeout(timer); try { port = new URL(spoke.endpoint).port || port; } catch (_) { /* keep default */ }
urls.push(`http://${spoke.meshIp}:${port}/api/site/resync`);
} }
urls.push(String(spoke.endpoint).replace(/\/+$/, '') + '/api/site/resync');
return urls;
} }
module.exports = { replicateToSpokes }; async function pingOne(spoke, reason) {
const urls = resyncUrls(spoke);
let lastErr;
for (const url of urls) {
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), RESYNC_TIMEOUT_MS);
try {
const resp = await fetch(url, {
method: 'POST',
headers: { Authorization: 'Bearer ' + spoke.pushToken, 'Content-Type': 'application/json' },
body: JSON.stringify({ reason: reason || 'catalog-changed' }),
signal: controller.signal
});
if (!resp.ok) throw new Error('status ' + resp.status);
return; // success -- don't try the next (fallback) URL
} catch (err) {
lastErr = err;
} finally {
clearTimeout(timer);
}
}
throw lastErr;
}
module.exports = { replicateToSpokes, resyncUrls };