diff --git a/docker-compose.yml b/docker-compose.yml index a7749486..98b448af 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -23,8 +23,10 @@ services: - LONG=${RECEIVER_LON:-0} - TAR1090_DEFAULTCENTERLAT=${RECEIVER_LAT:-0} - TAR1090_DEFAULTCENTERLON=${RECEIVER_LON:-0} - # adsb.lol proxy config + # Remote ADS-B proxy config. ADSB_UPSTREAMS is the ordered, + # comma-separated list of sources tried until one answers. - ADSBLOL_ENABLED=${ADSBLOL_ENABLED:-false} + - ADSB_UPSTREAMS=${ADSB_UPSTREAMS:-https://api.adsb.lol} - RECEIVER_LAT=${RECEIVER_LAT:-0} - RECEIVER_LON=${RECEIVER_LON:-0} - ADSBLOL_RADIUS=${ADSBLOL_RADIUS:-40} diff --git a/proxy/server.js b/proxy/server.js index 609da7d7..797f5c48 100644 --- a/proxy/server.js +++ b/proxy/server.js @@ -27,9 +27,28 @@ const MAX_STALE_MS = parseInt(process.env.ADSBLOL_MAX_STALE_MS || '60000'); const USER_AGENT = process.env.ADSBLOL_USER_AGENT || 'retina-node/1.0 (+https://github.com/offworldlabs/tar1090-node)'; -const ADSBLOL_API = `https://api.adsb.lol/v2/lat/${RECEIVER_LAT}/lon/${RECEIVER_LON}/dist/${ADSBLOL_RADIUS}`; +// Ordered, comma-separated list of adsb.lol-format sources, tried in order +// until one answers. RETINA's own service at adsb.retina.fm serves the same v2 +// envelope, so it needs no conversion changes. Default keeps a node that is +// told nothing on adsb.lol alone. +const ADSB_UPSTREAMS = (process.env.ADSB_UPSTREAMS || 'https://api.adsb.lol') + .split(',') + .map(entry => entry.trim().replace(/\/+$/, '')) + .filter(Boolean); -// Last good adsb.lol response: { payload, fetchedAt }. `payload.now` is the +function upstreamUrl(base) { + return `${base}/v2/lat/${RECEIVER_LAT}/lon/${RECEIVER_LON}/dist/${ADSBLOL_RADIUS}`; +} + +function hostOf(base) { + try { + return new URL(base).host; + } catch { + return base; + } +} + +// Last good remote response: { payload, fetchedAt, source }. `payload.now` is the // fetch time and is never restamped on serve - consumers rely on it to work out // how stale each position is. let cache = null; @@ -45,12 +64,12 @@ function emptyPayload() { return { now: Date.now() / 1000, messages: 0, aircraft: [] }; } -function fetchUrl(url) { +function fetchUrl(url, timeoutMs = UPSTREAM_TIMEOUT_MS) { return new Promise((resolve, reject) => { const client = url.startsWith('https') ? https : http; const req = client.get(url, { - timeout: UPSTREAM_TIMEOUT_MS, + timeout: timeoutMs, headers: { 'User-Agent': USER_AGENT } }, (res) => { if (res.statusCode !== 200) { @@ -128,25 +147,44 @@ function convertAdsbLolToReadsb(adsbLolData) { }; } -// Refreshes the cache, collapsing concurrent callers onto one upstream request. -// Never rejects - a failed refresh leaves the previous cache in place. +// Tries each upstream in turn and caches the first that answers, collapsing +// concurrent callers onto one refresh. Never rejects - a refresh where nothing +// answers leaves the previous cache in place, as the single-upstream version +// did. +// +// The chain shares ONE time budget rather than a timeout per source: two +// sources at UPSTREAM_TIMEOUT_MS each would be 6 s worst case, past blah2-api's +// 5 s client timeout, turning a slow upstream into a consumer-side failure +// instead of the stale-but-served degradation that timeout is chosen to give. function refreshRemote() { - if (!inFlight) { - lastAttemptAt = Date.now(); - console.log('Fetching from adsb.lol...'); - inFlight = fetchUrl(ADSBLOL_API) - .then((raw) => { - const payload = convertAdsbLolToReadsb(raw); - cache = { payload, fetchedAt: Date.now() }; - console.log(`adsb.lol: ${payload.aircraft.length} aircraft`); - }) - .catch((err) => { - console.log(`adsb.lol fetch failed: ${err.message}`); - }) - .finally(() => { - inFlight = null; - }); + if (inFlight) { + return inFlight; } + + lastAttemptAt = Date.now(); + const deadline = lastAttemptAt + UPSTREAM_TIMEOUT_MS; + + inFlight = (async () => { + for (const base of ADSB_UPSTREAMS) { + const host = hostOf(base); + const remaining = deadline - Date.now(); + if (remaining < 250) { + break; + } + + try { + const payload = convertAdsbLolToReadsb(await fetchUrl(upstreamUrl(base), remaining)); + cache = { payload, fetchedAt: Date.now(), source: host }; + console.log(`${host}: ${payload.aircraft.length} aircraft`); + return; + } catch (err) { + console.log(`${host} fetch failed: ${err.message}`); + } + } + })().finally(() => { + inFlight = null; + }); + return inFlight; } @@ -190,7 +228,7 @@ async function getAircraftData() { if (servedAge <= MAX_STALE_MS) { return { data: cache.payload, - source: 'adsb.lol', + source: cache.source, ageMs: servedAge, stale: servedAge >= CACHE_TTL_MS }; @@ -237,9 +275,9 @@ const server = http.createServer(async (req, res) => { server.listen(PORT, () => { console.log(`Aircraft data proxy listening on port ${PORT}`); console.log(`Local data file: ${LOCAL_DATA_PATH}`); - console.log(`adsb.lol fallback: ${ADSBLOL_ENABLED ? 'enabled' : 'disabled'}`); + console.log(`Remote fallback: ${ADSBLOL_ENABLED ? 'enabled' : 'disabled'}`); if (ADSBLOL_ENABLED) { - console.log(`adsb.lol API: ${ADSBLOL_API}`); - console.log(`cache TTL: ${CACHE_TTL_MS} ms, upstream timeout: ${UPSTREAM_TIMEOUT_MS} ms`); + ADSB_UPSTREAMS.forEach((base, i) => console.log(` upstream ${i + 1}: ${upstreamUrl(base)}`)); + console.log(`cache TTL: ${CACHE_TTL_MS} ms, chain budget: ${UPSTREAM_TIMEOUT_MS} ms`); } }); diff --git a/test/proxy-chain.test.js b/test/proxy-chain.test.js new file mode 100644 index 00000000..f64335e1 --- /dev/null +++ b/test/proxy-chain.test.js @@ -0,0 +1,158 @@ +// Tests for the upstream chain in proxy/server.js. +// +// Run: node --test test/proxy-chain.test.js +// +// The proxy is a script with no exports, so each case runs it as a child +// process against local fake upstreams and asserts over its HTTP surface. + +const { test } = require('node:test'); +const assert = require('node:assert'); +const http = require('node:http'); +const fs = require('node:fs'); +const os = require('node:os'); +const path = require('node:path'); +const { spawn } = require('node:child_process'); + +const SERVER = path.join(__dirname, '..', 'proxy', 'server.js'); + +function aircraft(hex) { + return { hex, flight: 'TEST123 ', lat: 42.2, lon: -72.7, alt_baro: 30000, seen_pos: 0.4 }; +} + +// A fake adsb.lol-format upstream that records every hit, so a test can prove +// a source was or was not consulted. +async function fakeUpstream(handler) { + const hits = []; + const server = http.createServer((req, res) => { + hits.push(req.url); + const { status = 200, body = { ac: [], total: 0 } } = handler() || {}; + res.writeHead(status, { 'Content-Type': 'application/json' }); + res.end(typeof body === 'string' ? body : JSON.stringify(body)); + }); + await new Promise(resolve => server.listen(0, '127.0.0.1', resolve)); + const { port } = server.address(); + return { base: `http://127.0.0.1:${port}`, host: `127.0.0.1:${port}`, hits, close: () => server.close() }; +} + +async function startProxy(env) { + const port = 34000 + Math.floor(Math.random() * 1000); + const child = spawn(process.execPath, [SERVER], { + env: { + ...process.env, + PROXY_PORT: String(port), + ADSBLOL_ENABLED: 'true', + RECEIVER_LAT: '42.2', + RECEIVER_LON: '-72.7', + LOCAL_DATA_PATH: '/nonexistent/aircraft.json', + ...env + }, + stdio: ['ignore', 'pipe', 'pipe'] + }); + + let log = ''; + child.stdout.on('data', d => { log += d; }); + child.stderr.on('data', d => { log += d; }); + + const deadline = Date.now() + 5000; + for (;;) { + try { + await get(port, '/health'); + break; + } catch { + if (Date.now() > deadline) throw new Error(`proxy did not start: ${log}`); + } + } + return { port, stop: () => child.kill() }; +} + +function get(port, urlPath) { + return new Promise((resolve, reject) => { + const req = http.get({ host: '127.0.0.1', port, path: urlPath, timeout: 8000 }, res => { + let data = ''; + res.on('data', c => data += c); + res.on('end', () => resolve({ headers: res.headers, body: JSON.parse(data) })); + }); + req.on('timeout', () => { req.destroy(); reject(new Error('timeout')); }); + req.on('error', reject); + }); +} + +test('the first upstream that answers wins, and the second is never consulted', async () => { + const primary = await fakeUpstream(() => ({ body: { ac: [aircraft('abc123')], total: 1 } })); + const secondary = await fakeUpstream(() => ({ body: { ac: [aircraft('def456')], total: 1 } })); + const proxy = await startProxy({ ADSB_UPSTREAMS: `${primary.base},${secondary.base}` }); + + try { + const res = await get(proxy.port, '/data/aircraft.json'); + assert.strictEqual(res.body.aircraft[0].hex, 'abc123'); + assert.strictEqual(res.headers['x-data-source'], primary.host); + assert.strictEqual(secondary.hits.length, 0); + } finally { + proxy.stop(); primary.close(); secondary.close(); + } +}); + +test('a failing first upstream falls through to the second', async () => { + // adsb.lol refuses most of our requests with 429; that is the case this + // whole change exists to cover. + const primary = await fakeUpstream(() => ({ status: 429, body: 'Too Many Requests' })); + const secondary = await fakeUpstream(() => ({ body: { ac: [aircraft('def456')], total: 1 } })); + const proxy = await startProxy({ ADSB_UPSTREAMS: `${primary.base},${secondary.base}` }); + + try { + const res = await get(proxy.port, '/data/aircraft.json'); + assert.strictEqual(res.body.aircraft[0].hex, 'def456'); + assert.strictEqual(res.headers['x-data-source'], secondary.host); + } finally { + proxy.stop(); primary.close(); secondary.close(); + } +}); + +test('a local receiver wins over every remote source', async () => { + const localPath = path.join(os.tmpdir(), `local-${process.pid}-${Date.now()}.json`); + fs.writeFileSync(localPath, JSON.stringify({ + now: Date.now() / 1000, + messages: 42, + aircraft: [{ hex: 'local01', lat: 42.2, lon: -72.7, gs: 200, track: 90, seen_pos: 0.2 }] + })); + const primary = await fakeUpstream(() => ({ body: { ac: [aircraft('abc123')], total: 1 } })); + const proxy = await startProxy({ ADSB_UPSTREAMS: primary.base, LOCAL_DATA_PATH: localPath }); + + try { + const res = await get(proxy.port, '/data/aircraft.json'); + assert.strictEqual(res.body.aircraft[0].hex, 'local01'); + assert.strictEqual(res.headers['x-data-source'], 'local'); + assert.strictEqual(primary.hits.length, 0, 'no upstream may be polled when a receiver has data'); + } finally { + proxy.stop(); primary.close(); fs.unlinkSync(localPath); + } +}); + +test('the chain shares one time budget rather than one timeout per source', async () => { + // Two sources must not add up past blah2-api's 5 s client timeout. + const hang = http.createServer(() => {}); + await new Promise(r => hang.listen(0, '127.0.0.1', r)); + const hangBase = `http://127.0.0.1:${hang.address().port}`; + const proxy = await startProxy({ + ADSB_UPSTREAMS: `${hangBase},${hangBase}`, + ADSBLOL_TIMEOUT_MS: '1000' + }); + + try { + const started = Date.now(); + await get(proxy.port, '/data/aircraft.json'); + const elapsed = Date.now() - started; + assert.ok(elapsed < 2000, `chain took ${elapsed} ms, budget was 1000 ms`); + } finally { + proxy.stop(); hang.close(); + } +}); + +test('a node told nothing stays on adsb.lol alone', async () => { + const proxy = await startProxy({}); + try { + await get(proxy.port, '/health'); + } finally { + proxy.stop(); + } +});