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
4 changes: 3 additions & 1 deletion docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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}
Expand Down
88 changes: 63 additions & 25 deletions proxy/server.js
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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) {
Expand Down Expand Up @@ -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;
}

Expand Down Expand Up @@ -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
};
Expand Down Expand Up @@ -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`);
}
});
158 changes: 158 additions & 0 deletions test/proxy-chain.test.js
Original file line number Diff line number Diff line change
@@ -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();
}
});
Loading