Skip to content
Merged
Show file tree
Hide file tree
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
67 changes: 67 additions & 0 deletions scripts/natural/source-request-diagnostics.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
import { AsyncLocalStorage } from 'node:async_hooks';
import { channel } from 'node:diagnostics_channel';

const scope = new AsyncLocalStorage();
const requests = new WeakMap();
let activeScopes = 0;

function update(request, change) {
const state = requests.get(request);
if (state?.active && state.request === request) change(state.progress);
}

const listeners = [
['undici:request:create', ({ request }) => {
const state = scope.getStore();
if (!state?.active) return;
state.request = request;
state.progress = {
observed: true,
requestCount: state.progress.requestCount + 1,
requestSendObserved: false,
responseHeadersObserved: false,
wireBodyBytes: null,
firstBodyByteMs: null,
lastBodyByteMs: null,
wireBodyComplete: false,
};
requests.set(request, state);
}],
['undici:client:sendHeaders', ({ request }) => update(request, progress => {
progress.requestSendObserved = true;
})],
['undici:request:headers', ({ request }) => update(request, progress => {
progress.responseHeadersObserved = true;
})],
['undici:request:bodyChunkReceived', ({ request, chunk }) => {
const state = requests.get(request);
if (!state?.active || state.request !== request) return;
const elapsed = Math.round(performance.now() - state.started);
state.progress.wireBodyBytes = (state.progress.wireBodyBytes ?? 0) + chunk.byteLength;
state.progress.firstBodyByteMs ??= elapsed;
state.progress.lastBodyByteMs = elapsed;
}],
['undici:request:trailers', ({ request }) => update(request, progress => {
progress.wireBodyComplete = true;
})],
].map(([name, listener]) => [channel(name), listener]);

// Native request identities keep pooled connections and concurrent fetches separate.
// Counters cover the latest redirect hop and encoded payload bytes. A local send
// observation does not prove remote receipt; wire completion does not prove JSON validity.
// Bytes remain unknown until a chunk is observed: not all transports emit chunk events.
export async function withSourceRequestDiagnostics(operation) {
const state = { active: true, started: performance.now(), request: null, progress: { observed: false, requestCount: 0 } };
if (activeScopes++ === 0) {
for (const [event, listener] of listeners) event.subscribe(listener);
}
try {
return await scope.run(state, () => operation(() => ({ ...state.progress })));
} finally {
state.active = false;
state.request = null;
if (--activeScopes === 0) {
for (const [event, listener] of listeners) event.unsubscribe(listener);
}
}
}
1 change: 1 addition & 0 deletions scripts/railway-services.json
Original file line number Diff line number Diff line change
Expand Up @@ -481,6 +481,7 @@
"scripts/_seed-envelope-source.mjs",
"scripts/_seed-utils.mjs",
"scripts/lib/llm-telemetry.cjs",
"scripts/natural/source-request-diagnostics.mjs",
"scripts/natural/western-pacific-cyclones.mjs",
"scripts/nixpacks.toml",
"scripts/package-lock.json",
Expand Down
16 changes: 13 additions & 3 deletions scripts/seed-natural-events.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import { fileURLToPath } from 'node:url';
import { resolve } from 'node:path';
import { getDefaultAutoSelectFamily, getDefaultAutoSelectFamilyAttemptTimeout, isIP } from 'node:net';
import { Agent } from 'undici';
import { withSourceRequestDiagnostics } from './natural/source-request-diagnostics.mjs';

import {
loadEnvFile,
Expand Down Expand Up @@ -187,7 +188,7 @@ async function fetchEventSourceJson(source, url, fetchFn) {
const deadline = started + SOURCE_REQUEST_BUDGET_MS;
let attempt = 0;
let ipv4Retry = false;
return withRetry(async () => {
return withRetry(() => withSourceRequestDiagnostics(async progress => {
const remaining = Math.floor(deadline - performance.now());
if (remaining <= 0) throw Object.assign(new Error(`${source} request budget exhausted`), { nonRetryable: true });
attempt++;
Expand Down Expand Up @@ -230,7 +231,7 @@ async function fetchEventSourceJson(source, url, fetchFn) {
details = ' details={"unavailable":true}';
}
}
const error = Object.assign(new Error(`${source} ${stage} ${kind} attempt=${attempt} elapsedMs=${Math.round(finished - started)} attemptElapsedMs=${Math.round(finished - attemptStarted)}${phaseTimings}${details}`), {
const error = Object.assign(new Error(`${source} ${stage} ${kind} attempt=${attempt} elapsedMs=${Math.round(finished - started)} attemptElapsedMs=${Math.round(finished - attemptStarted)}${phaseTimings} progress=${JSON.stringify(progress())}${details}`), {
nonRetryable: stage === 'http' ? cause.nonRetryable : !transport,
retryAfterMs: cause.retryAfterMs,
});
Expand All @@ -239,7 +240,7 @@ async function fetchEventSourceJson(source, url, fetchFn) {
} finally {
await dispatcher?.destroy();
}
}, 1, 500);
}), 1, 500);
}

async function fetchEonet(days, fetchFn = globalThis.fetch, now = Date.now()) {
Expand Down Expand Up @@ -1056,6 +1057,15 @@ export function naturalEventsAfterPublish(data) {
if (!failedSources.length) return { freshnessMetaPatch: patch };

console.warn(`[natural-events] DEGRADED: ${failedSources.join(', ')}`);
const failedSourceHealth = Object.fromEntries(failedSources.filter(source => sourceHealth[source]).map(source => {
const health = sourceHealth[source];
return [source, {
...health,
remainingRetentionMs: health.retainedUntil === null ? null
: Math.max(0, health.retainedUntil - (health.lastAttemptAt ?? Date.now())),
}];
}));
if (Object.keys(failedSourceHealth).length) console.warn(`[natural-events] failed source health=${JSON.stringify(failedSourceHealth)}`);
if (nhcFailed) Object.assign(patch, {
errorCode: nhc.errorCode,
skipReason: 'nhc-required-point-coverage-incomplete',
Expand Down
114 changes: 114 additions & 0 deletions tests/natural-events-request-diagnostics.test.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,114 @@
import assert from 'node:assert/strict';
import { test } from 'node:test';
import { createServer } from 'node:http';
import { once } from 'node:events';
import { channel } from 'node:diagnostics_channel';
import { gzipSync } from 'node:zlib';
import { withSourceRequestDiagnostics } from '../scripts/natural/source-request-diagnostics.mjs';

async function fixture(t, handler) {
const server = createServer(handler);
server.listen(0, '127.0.0.1');
await once(server, 'listening');
t.after(() => { server.closeAllConnections(); server.close(); });
return `http://127.0.0.1:${server.address().port}`;
}

async function observe(url) {
return withSourceRequestDiagnostics(async snapshot => {
try {
const response = await fetch(url, { signal: AbortSignal.timeout(150) });
const value = await response.json();
return { value, progress: snapshot() };
} catch (error) {
return { error: error.name, progress: snapshot() };
}
});
}

test('concurrent requests to one origin keep complete, incomplete and header-wait facts separate', async t => {
const origin = await fixture(t, (req, res) => {
if (req.url === '/headers') return;
if (req.url === '/partial') { res.writeHead(200); res.write('{"private":'); return; }
res.end('{"ok":true}');
});
const [complete, partial, headers] = await Promise.all([
observe(`${origin}/complete`), observe(`${origin}/partial`), observe(`${origin}/headers`),
]);
assert.deepEqual(complete.value, { ok: true });
assert.equal(complete.progress.wireBodyComplete, true);
assert.equal(complete.progress.wireBodyBytes, 11);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
assert.equal(partial.error, 'TimeoutError');
assert.equal(partial.progress.wireBodyBytes, 11);
assert.equal(partial.progress.responseHeadersObserved, true);
assert.equal(partial.progress.wireBodyComplete, false);
assert.equal(headers.error, 'TimeoutError');
assert.equal(headers.progress.requestSendObserved, true);
assert.equal(headers.progress.responseHeadersObserved, false);
assert.equal(headers.progress.wireBodyBytes, null);
assert.equal(headers.progress.firstBodyByteMs, null);
assert.equal(headers.progress.lastBodyByteMs, null);
assert.equal(channel('undici:request:create').hasSubscribers, false);
});

test('wire counters measure compressed bytes and keep complete malformed JSON distinct from incomplete bodies', async t => {
const body = gzipSync('{"private":"data",');
const origin = await fixture(t, (_req, res) => {
res.writeHead(200, { 'Content-Encoding': 'gzip' });
res.end(body);
});
const result = await observe(origin);
assert.equal(result.error, 'SyntaxError');
assert.equal(result.progress.wireBodyBytes, body.byteLength);
assert.equal(result.progress.wireBodyComplete, true);
assert.doesNotMatch(JSON.stringify(result.progress), /private|data|127\.0\.0\.1/);
});

test('a complete gzip payload without HTTP completion remains incomplete', async t => {
const body = gzipSync('{"ok":true}');
const origin = await fixture(t, (_req, res) => {
res.writeHead(200, { 'Content-Encoding': 'gzip' });
res.write(body);
});
const result = await observe(origin);
assert.equal(result.error, 'TimeoutError');
assert.equal(result.progress.wireBodyBytes, body.byteLength);
assert.equal(result.progress.wireBodyComplete, false);
});

test('redirects describe the final hop without adding earlier response bytes', async t => {
const origin = await fixture(t, (req, res) => {
if (req.url === '/redirect') { res.writeHead(302, { Location: '/final' }); res.end('redirect body'); return; }
res.writeHead(200); res.write('{');
});
const result = await observe(`${origin}/redirect`);
assert.equal(result.error, 'TimeoutError');
assert.equal(result.progress.requestCount, 2);
assert.equal(result.progress.wireBodyBytes, 1);
assert.equal(result.progress.wireBodyComplete, false);
});

test('unobserved transports remain unknown and throwing operations release subscriptions', async () => {
assert.deepEqual(await withSourceRequestDiagnostics(snapshot => snapshot()), { observed: false, requestCount: 0 });
const error = new Error('fixture');
await assert.rejects(withSourceRequestDiagnostics(() => { throw error; }), value => value === error);
for (const name of ['request:create', 'client:sendHeaders', 'request:headers', 'request:bodyChunkReceived', 'request:trailers']) {
assert.equal(channel(`undici:${name}`).hasSubscribers, false, name);
}
});

test('request events without body-chunk telemetry leave byte measurements unknown', async () => {
const progress = await withSourceRequestDiagnostics(snapshot => {
const request = {};
channel('undici:request:create').publish({ request });
channel('undici:request:headers').publish({ request });
channel('undici:request:trailers').publish({ request });
return snapshot();
});
assert.equal(progress.observed, true);
assert.equal(progress.responseHeadersObserved, true);
assert.equal(progress.wireBodyComplete, true);
assert.equal(progress.wireBodyBytes, null);
assert.equal(progress.firstBodyByteMs, null);
assert.equal(progress.lastBodyByteMs, null);
});
25 changes: 25 additions & 0 deletions tests/natural-events-source-recovery.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -178,3 +178,28 @@ test('all GDACS type requests can fail without discarding valid uncapped last-go
assert.equal(second._sourceSnapshots['gdacs:EQ'].records.length, 101);
assert.equal(naturalEventsAfterPublish(second).freshnessMetaPatch.failedSources.length, 6);
});

test('failure logs report fixed retention clocks without changing published health', async (t) => {
const first = await run();
const retained = await run({ previousSources: first._sourceSnapshots, now: NOW + HOUR, failures: ['eonet'] });
const expired = await run({ previousSources: first._sourceSnapshots, now: NOW + 9 * HOUR, failures: ['eonet'] });
const warnings = [];
t.mock.method(console, 'warn', message => warnings.push(message));
const before = structuredClone(retained);
const result = naturalEventsAfterPublish(retained);
const prefix = '[natural-events] failed source health=';
const logged = JSON.parse(warnings.find(message => message.startsWith(prefix)).slice(prefix.length));
assert.deepEqual(Object.keys(logged), ['eonet']);
assert.deepEqual(logged.eonet, { ...result.freshnessMetaPatch.sourceHealth.eonet, remainingRetentionMs: 8 * HOUR });
assert.deepEqual(retained, before);
assert.equal('remainingRetentionMs' in result.freshnessMetaPatch.sourceHealth.eonet, false);
warnings.length = 0;
naturalEventsAfterPublish(expired);
const unavailable = JSON.parse(warnings.find(message => message.startsWith(prefix)).slice(prefix.length));
assert.equal(unavailable.eonet.status, 'unavailable');
assert.equal(unavailable.eonet.lastSuccessAt, null);
assert.equal(unavailable.eonet.remainingRetentionMs, null);
warnings.length = 0;
naturalEventsAfterPublish(first);
assert.deepEqual(warnings, []);
});
32 changes: 32 additions & 0 deletions tests/natural-events-transport.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -431,3 +431,35 @@ test('request failures omit unobserved phase timings and successful responses ad
assert.equal(result._sourceSnapshots.eonet.fetchedAt, NOW);
assert.doesNotMatch(logs.join('\n'), /headersElapsedMs|bodyElapsedMs|attempt=2/);
});

test('failed native attempts log isolated wire progress without logging response content', async t => {
const logs = [];
t.mock.method(console, 'warn', (...args) => logs.push(args.join(' ')));
t.mock.method(console, 'log', (...args) => logs.push(args.join(' ')));
const timeout = AbortSignal.timeout;
t.mock.method(AbortSignal, 'timeout', () => timeout(100));
const body = '{"private":"secret';
const server = createServer((_req, res) => { res.writeHead(200); res.write(body); });
server.listen(0, '127.0.0.1');
await once(server, 'listening');
t.after(() => { server.closeAllConnections(); server.close(); });
const transport = fixture((source, _attempt, options) => source === 'eonet'
? fetch(`http://127.0.0.1:${server.address().port}/private?token=secret`, options) : undefined);
const result = await run(transport);
const progress = logs.flatMap(line => [...line.matchAll(/ progress=(\{[^}]+\})/g)].map(match => JSON.parse(match[1])));
assert.equal(progress.length, 2);
for (const item of progress) {
assert.equal(item.observed, true);
assert.equal(item.requestCount, 1);
assert.equal(item.requestSendObserved, true);
assert.equal(item.responseHeadersObserved, true);
assert.equal(item.wireBodyBytes, Buffer.byteLength(body));
assert.equal(item.wireBodyComplete, false);
assert.ok(item.firstBodyByteMs >= 0);
assert.ok(item.lastBodyByteMs >= item.firstBodyByteMs);
}
assert.doesNotMatch(logs.join('\n'), /private|secret|127\.0\.0\.1/);
assert.equal(transport.calls.get('eonet'), 2);
for (const [source, count] of transport.calls) if (source !== 'eonet') assert.equal(count, 1, source);
assert.deepEqual(naturalEventsAfterPublish(result).freshnessMetaPatch.failedSources, ['eonet']);
});
Loading