// REST client in plain Node 18+ (no dependencies): mint a token, subscribe by long polling, publish. // CYBERPLEX_URL=https://cyberplex.replit.app CYBERPLEX_API_KEY=ak_... node node-rest.mjs // Optional: RUN_SECONDS (how long to keep listening; default 2) and CYBERPLEX_TOKEN_TTL_SEC (default 600). const URL_ = process.env.CYBERPLEX_URL ?? 'https://cyberplex.replit.app'; const API_KEY = process.env.CYBERPLEX_API_KEY; const TTL_SEC = Number(process.env.CYBERPLEX_TOKEN_TTL_SEC ?? 600); const RUN_SECONDS = Number(process.env.RUN_SECONDS ?? 2); if (!API_KEY) throw new Error('set CYBERPLEX_API_KEY to your tenant API key'); const CHANNEL = 'node-demo'; async function call(token, method, path, body) { const res = await fetch(`${URL_}${path}`, { method, headers: { authorization: `Bearer ${token}`, ...(body ? { 'content-type': 'application/json' } : {}) }, body: body ? JSON.stringify(body) : undefined }); const json = await res.json().catch(() => ({})); if (!res.ok) { const retry = res.headers.get('retry-after'); throw Object.assign(new Error(`${res.status} ${json?.error?.code ?? ''} ${json?.error?.message ?? JSON.stringify(json)}`), { status: res.status, retryAfterSec: retry ? Number(retry) : undefined }); } return json; } // Agent tokens are short-lived, so a long-running client must renew them. Keep one token, renew it at 80% of its life, and if the // server still answers 401 (clock skew, or the service rotated its signing secret) mint a fresh one and retry once. The API key does // not expire on its own, so it can always mint again. Do this minting on your SERVER, never in a browser (it needs the secret key). let token; let renewAt = 0; async function mint() { if (token) console.log('(renewing token)'); const r = await call(API_KEY, 'POST', '/v1/agent-token', { agentId: 'node-agent', ttlSec: TTL_SEC }); token = r.token; renewAt = Date.now() + r.expiresInSec * 800; // 80% of the lifetime, in ms } async function api(method, path, body) { if (!token || Date.now() >= renewAt) await mint(); try { return await call(token, method, path, body); } catch (e) { if (e.status !== 401) throw e; console.log('(token rejected: minting a new one)'); await mint(); return call(token, method, path, body); } } // Subscribe: ask for the head cursor once, then loop on long polls, resuming from the last cursor seen. async function subscribe(channel, onMessage, signal) { let { cursor } = await api('GET', `/v1/channels/${channel}/messages?limit=1`); while (!signal.aborted) { try { const page = await api('GET', `/v1/channels/${channel}/messages?cursor=${encodeURIComponent(cursor)}&wait=20&limit=100`); for (const m of page.publications) { onMessage(m.data); cursor = m.cursor; // advance per message so a failure never skips one } if (!page.publications.length) cursor = page.cursor; if (page.truncated) console.warn('gap: some messages were evicted before we could read them'); } catch (e) { if (signal.aborted) return; // 429 => wait as told; anything else (a network blip, a restart) => brief backoff. The cursor is kept, so nothing is lost. await new Promise((r) => setTimeout(r, (e.retryAfterSec ?? 2) * 1000)); } } } const stop = new AbortController(); const reader = subscribe(CHANNEL, (m) => console.log('received:', JSON.stringify(m)), stop.signal); // Publish a few messages. await new Promise((r) => setTimeout(r, 500)); for (const n of [1, 2, 3]) { const { cursor } = await api('POST', `/v1/channels/${CHANNEL}/messages`, { data: { n } }); console.log(`published n=${n} -> cursor ${cursor}`); } await new Promise((r) => setTimeout(r, RUN_SECONDS * 1000)); stop.abort(); await Promise.race([reader, new Promise((r) => setTimeout(r, 100))]); process.exit(0); // the in-flight long poll would otherwise keep the process alive for up to 20s