Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 57 additions & 15 deletions src/server.js
Original file line number Diff line number Diff line change
Expand Up @@ -39,20 +39,59 @@ app.use(express.static(path.join(__dirname, '..', 'public')));
app.use(requestId);

// ─── SSE Event Stream ────────────────────────────────────────
const sseClients = [];
const HEARTBEAT_INTERVAL_MS = 30_000;
const STALE_CLIENT_THRESHOLD_MS = 90_000;
const sseClients = new Map();
let x402MiddlewareReady = false;

function removeClient(clientId) {
const client = sseClients.get(clientId);
if (!client) return;
sseClients.delete(clientId);
client.res.end();
}
Comment thread
Kingajong marked this conversation as resolved.

function writeSse(clientId, event) {
const client = sseClients.get(clientId);
if (!client) return false;

try {
const ok = client.res.write(`data: ${JSON.stringify(event)}\n\n`);
client.lastWriteAt = Date.now();
if (!ok) {
client.res.once('drain', () => {
const activeClient = sseClients.get(clientId);
if (activeClient) activeClient.lastWriteAt = Date.now();
});
}
return true;
} catch {
removeClient(clientId);
return false;
}
}

function broadcast(event) {
const data = JSON.stringify(event);
sseClients.forEach(res => {
res.write(`data: ${data}\n\n`);
});
for (const clientId of sseClients.keys()) {
writeSse(clientId, event);
}
}

setInterval(() => {
const now = Date.now();
for (const [clientId, client] of sseClients.entries()) {
if (now - client.lastWriteAt > STALE_CLIENT_THRESHOLD_MS) {
removeClient(clientId);
continue;
}

writeSse(clientId, { type: 'heartbeat', timestamp: new Date(now).toISOString() });
}
}, HEARTBEAT_INTERVAL_MS);

function buildReadinessPayload() {
const anthropicConfigured = !!config.anthropicApiKey;
const x402Enabled = !!config.serverAddress;

const components = {
app: {
ready: true,
Expand All @@ -69,14 +108,13 @@ function buildReadinessPayload() {
enabled: x402Enabled,
ready: !x402Enabled || (x402MiddlewareReady && !!config.facilitatorUrl),
description: x402Enabled
? (x402MiddlewareReady
? 'x402 payment middleware initialized'
: 'x402 is enabled but middleware initialization failed')
? x402MiddlewareReady
? 'x402 payment middleware initialized'
: 'x402 is enabled but middleware initialization failed'
: 'x402 paywall is disabled; premium endpoints are not protected by payments',
},
};

const ready = Object.values(components).every(component => component.ready);
const ready = Object.values(components).every((component) => component.ready);
return {
status: ready ? 'ready' : 'not_ready',
ready,
Expand Down Expand Up @@ -106,12 +144,16 @@ app.get('/api/events', (req, res) => {
res.flushHeaders();

// Send initial connection event
res.write(`data: ${JSON.stringify({ type: 'connected', timestamp: new Date().toISOString() })}\n\n`);
const clientId = `${Date.now()}-${Math.random().toString(36).slice(2)}`;
sseClients.set(clientId, { res, lastWriteAt: Date.now() });
writeSse(clientId, { type: 'connected', timestamp: new Date().toISOString() });

sseClients.push(res);
req.on('close', () => {
const idx = sseClients.indexOf(res);
if (idx !== -1) sseClients.splice(idx, 1);
removeClient(clientId);
});

req.on('error', () => {
removeClient(clientId);
});
});

Expand Down