170 lines
6 KiB
TypeScript
170 lines
6 KiB
TypeScript
// Minimal GTFS-Realtime (protobuf) reader: just the fields the app uses, no dependency.
|
|
// Spec: https://gtfs.org/documentation/realtime/proposal/ (field numbers below are from gtfs-realtime.proto).
|
|
|
|
export interface StopEvent { delay?: number; time?: number }
|
|
export interface StopUpdate { stopId: string; seq?: number; arrival?: StopEvent; departure?: StopEvent; skipped?: boolean }
|
|
export interface TripUpdate { tripId: string; routeId?: string; startDate?: string; delay?: number; canceled?: boolean; stops: StopUpdate[] }
|
|
export interface Alert { routeIds: string[]; stopIds: string[]; header: string; description: string; start?: number; end?: number }
|
|
export interface RealtimeFeed { timestamp: number; trips: TripUpdate[]; alerts: Alert[] }
|
|
|
|
class Reader {
|
|
private pos = 0;
|
|
constructor(private readonly buf: Uint8Array, private readonly end = buf.length) {}
|
|
get done() { return this.pos >= this.end; }
|
|
|
|
varint(): number {
|
|
let result = 0, mul = 1, byte: number;
|
|
// up to 10 bytes; values past 2^53 lose precision, which never matters for the fields we read
|
|
for (let i = 0; i < 10; i++) {
|
|
if (this.pos >= this.end) throw new Error('truncated varint');
|
|
byte = this.buf[this.pos++]!;
|
|
result += (byte & 0x7f) * mul;
|
|
mul *= 128;
|
|
if (byte < 0x80) return result;
|
|
}
|
|
throw new Error('varint too long');
|
|
}
|
|
/** int32/int64 as two's complement (negative delays arrive as 10-byte varints). */
|
|
signed(): number {
|
|
const start = this.pos;
|
|
let big = 0n, shift = 0n, byte: number;
|
|
do {
|
|
if (this.pos >= this.end) throw new Error('truncated varint');
|
|
byte = this.buf[this.pos++]!;
|
|
big |= BigInt(byte & 0x7f) << shift;
|
|
shift += 7n;
|
|
} while (byte >= 0x80 && this.pos - start < 10);
|
|
return Number(BigInt.asIntN(64, big));
|
|
}
|
|
bytes(): Reader {
|
|
const len = this.varint();
|
|
const r = new Reader(this.buf, this.pos + len);
|
|
r.pos = this.pos;
|
|
if (this.pos + len > this.end) throw new Error('truncated message');
|
|
this.pos += len;
|
|
return r;
|
|
}
|
|
string(): string {
|
|
const r = this.bytes();
|
|
return new TextDecoder().decode(this.buf.subarray(r.pos, r.end));
|
|
}
|
|
skip(wire: number) {
|
|
if (wire === 0) this.varint();
|
|
else if (wire === 1) this.pos += 8;
|
|
else if (wire === 2) this.bytes();
|
|
else if (wire === 5) this.pos += 4;
|
|
else throw new Error(`unsupported wire type ${wire}`);
|
|
}
|
|
/** Iterate (fieldNumber, wireType) pairs; the caller must consume or skip each value. */
|
|
fields(fn: (field: number, wire: number) => void) {
|
|
while (!this.done) {
|
|
const tag = this.varint();
|
|
fn(tag >>> 3, tag & 7);
|
|
}
|
|
}
|
|
}
|
|
|
|
function event(r: Reader): StopEvent {
|
|
const e: StopEvent = {};
|
|
r.fields((f, w) => {
|
|
if (f === 1 && w === 0) e.delay = r.signed();
|
|
else if (f === 2 && w === 0) e.time = r.signed();
|
|
else r.skip(w);
|
|
});
|
|
return e;
|
|
}
|
|
|
|
function stopUpdate(r: Reader): StopUpdate {
|
|
const u: StopUpdate = { stopId: '' };
|
|
r.fields((f, w) => {
|
|
if (f === 1 && w === 0) u.seq = r.varint();
|
|
else if (f === 2 && w === 2) u.arrival = event(r.bytes());
|
|
else if (f === 3 && w === 2) u.departure = event(r.bytes());
|
|
else if (f === 4 && w === 2) u.stopId = r.string();
|
|
else if (f === 5 && w === 0) { if (r.varint() === 1) u.skipped = true; } // SKIPPED
|
|
else r.skip(w);
|
|
});
|
|
return u;
|
|
}
|
|
|
|
function tripUpdate(r: Reader): TripUpdate {
|
|
const t: TripUpdate = { tripId: '', stops: [] };
|
|
r.fields((f, w) => {
|
|
if (f === 1 && w === 2) {
|
|
const d = r.bytes();
|
|
d.fields((g, x) => {
|
|
if (g === 1 && x === 2) t.tripId = d.string();
|
|
else if (g === 3 && x === 2) t.startDate = d.string();
|
|
else if (g === 4 && x === 0) { if (d.varint() === 3) t.canceled = true; } // CANCELED
|
|
else if (g === 5 && x === 2) t.routeId = d.string();
|
|
else d.skip(x);
|
|
});
|
|
} else if (f === 2 && w === 2) t.stops.push(stopUpdate(r.bytes()));
|
|
else if (f === 5 && w === 0) t.delay = r.signed();
|
|
else r.skip(w);
|
|
});
|
|
return t;
|
|
}
|
|
|
|
function translated(r: Reader): string {
|
|
let text = '';
|
|
r.fields((f, w) => {
|
|
if (f === 1 && w === 2) {
|
|
const tr = r.bytes();
|
|
tr.fields((g, x) => {
|
|
if (g === 1 && x === 2) { const s = tr.string(); if (!text) text = s; } else tr.skip(x);
|
|
});
|
|
} else r.skip(w);
|
|
});
|
|
return text;
|
|
}
|
|
|
|
function alert(r: Reader): Alert {
|
|
const a: Alert = { routeIds: [], stopIds: [], header: '', description: '' };
|
|
r.fields((f, w) => {
|
|
if (f === 1 && w === 2) {
|
|
const p = r.bytes();
|
|
p.fields((g, x) => {
|
|
if (g === 1 && x === 0) a.start = p.varint();
|
|
else if (g === 2 && x === 0) a.end = p.varint();
|
|
else p.skip(x);
|
|
});
|
|
} else if (f === 5 && w === 2) {
|
|
const e = r.bytes();
|
|
e.fields((g, x) => {
|
|
if (g === 2 && x === 2) a.routeIds.push(e.string());
|
|
else if (g === 5 && x === 2) a.stopIds.push(e.string());
|
|
else e.skip(x);
|
|
});
|
|
} else if (f === 10 && w === 2) a.header = translated(r.bytes());
|
|
else if (f === 11 && w === 2) a.description = translated(r.bytes());
|
|
else r.skip(w);
|
|
});
|
|
return a;
|
|
}
|
|
|
|
/** Decodes a TripUpdates / ServiceAlerts feed. Throws when the bytes are not a GTFS-RT FeedMessage (e.g. an HTML error page). */
|
|
export function decodeFeed(bytes: Uint8Array): RealtimeFeed {
|
|
const out: RealtimeFeed = { timestamp: 0, trips: [], alerts: [] };
|
|
const r = new Reader(bytes);
|
|
let sawHeader = false;
|
|
r.fields((f, w) => {
|
|
if (f === 1 && w === 2) {
|
|
const h = r.bytes();
|
|
h.fields((g, x) => {
|
|
if (g === 1 && x === 2) { h.string(); sawHeader = true; }
|
|
else if (g === 3 && x === 0) out.timestamp = h.varint();
|
|
else h.skip(x);
|
|
});
|
|
} else if (f === 2 && w === 2) {
|
|
const e = r.bytes();
|
|
e.fields((g, x) => {
|
|
if (g === 3 && x === 2) out.trips.push(tripUpdate(e.bytes()));
|
|
else if (g === 5 && x === 2) out.alerts.push(alert(e.bytes()));
|
|
else e.skip(x);
|
|
});
|
|
} else r.skip(w);
|
|
});
|
|
if (!sawHeader) throw new Error('not a GTFS-realtime feed');
|
|
return out;
|
|
}
|