cat > "$ANDO_REALTIME_DIR/listener.mjs" <<'EOF'
import { readFile, writeFile } from "node:fs/promises";
import path from "node:path";
import WebSocket from "ws";
const apiBase = requiredEnv("ANDO_API_BASE").replace(/\/$/, "");
const apiKey = requiredEnv("ANDO_API_KEY");
const workingDir = requiredEnv("ANDO_REALTIME_DIR");
const connectionRequestPath = path.join(workingDir, "open-connection.json");
const cursorPath = path.join(workingDir, "cursor.txt");
const reconnectableDisconnectReasons = new Set([
"server_restart",
"deploy_draining",
"connection_max_age",
"backpressure",
"temporary_unavailable",
"refresh_requested",
]);
const nonReconnectableCloseCodes = new Set([4401, 4403, 4406, 4409, 4410]);
let shouldStop = false;
let activeSocket = null;
process.on("SIGINT", () => {
shouldStop = true;
activeSocket?.close(1000, "client_shutdown");
});
await run();
async function run() {
let reconnectAttempt = 0;
while (!shouldStop) {
const cursor = await readCursor();
const ticket = await openConnection(cursor);
const result = await connectOnce(ticket);
if (shouldStop || !result.reconnect) {
break;
}
const waitSeconds =
result.retryAfterSeconds ?? Math.min(30, 2 ** reconnectAttempt);
reconnectAttempt = Math.min(reconnectAttempt + 1, 5);
console.log(`reconnecting in ${waitSeconds}s`);
await sleep(waitSeconds * 1000);
}
}
async function openConnection(cursor) {
const requestBody = await buildConnectionRequest(cursor);
const response = await fetch(`${apiBase}/realtime/connections`, {
method: "POST",
headers: {
"content-type": "application/json",
"x-api-key": apiKey,
},
body: JSON.stringify(requestBody),
});
const text = await response.text();
const body = parseJson(text);
if (!response.ok) {
console.error("ticket request failed", response.status, body ?? text);
throw new Error(`POST /realtime/connections returned ${response.status}`);
}
logWarnings(body.warnings);
return body;
}
async function buildConnectionRequest(cursor) {
const request = JSON.parse(await readFile(connectionRequestPath, "utf8"));
if (cursor == null) {
delete request.resume_from;
return request;
}
request.resume_from = { cursor };
return request;
}
async function connectOnce(ticket) {
return new Promise((resolve) => {
const socket = new WebSocket(ticket.url, ticket.protocol);
activeSocket = socket;
let settled = false;
let reconnect = true;
let retryAfterSeconds = null;
function settle() {
if (settled) {
return;
}
settled = true;
if (activeSocket === socket) {
activeSocket = null;
}
resolve({ reconnect, retryAfterSeconds });
}
socket.on("open", () => {
console.log("connected", {
connection_id: ticket.connection_id,
protocol: ticket.protocol,
});
});
socket.on("message", (data) => {
void handleServerFrame(socket, data).catch((error) => {
console.error("frame handler failed", error);
reconnect = true;
socket.close(1011, "client_error");
});
});
socket.on("error", (error) => {
console.error("socket error", error.message);
});
socket.on("close", (code, reason) => {
const closeReason = reason.toString();
if (nonReconnectableCloseCodes.has(code)) {
reconnect = false;
}
console.log("socket closed", {
code,
reason: closeReason || null,
reconnect,
});
settle();
});
async function handleServerFrame(openSocket, data) {
const frame = JSON.parse(data.toString());
if (frame.type === "hello") {
console.log("hello", {
connection_id: frame.connection_id,
heartbeat_interval_seconds: frame.heartbeat_interval_seconds,
resume_supported: frame.resume_supported,
subscriptions: frame.subscriptions,
});
logWarnings(frame.warnings);
return;
}
if (frame.type === "warning") {
logWarnings(frame.warnings);
return;
}
if (frame.type === "disconnect") {
console.warn("server requested disconnect", frame);
retryAfterSeconds = frame.retry_after_seconds ?? null;
reconnect = reconnectableDisconnectReasons.has(frame.reason);
openSocket.close(1000, frame.reason);
return;
}
if (frame.type === "event") {
try {
await handleProductEvent(frame.payload);
await sendJson(openSocket, {
envelope_id: frame.envelope_id,
});
await writeCursor(frame.cursor);
console.log("acked", {
envelope_id: frame.envelope_id,
cursor: frame.cursor,
});
} catch (error) {
await sendJson(openSocket, {
envelope_id: frame.envelope_id,
error: {
code: "handler_failed",
message: getErrorMessage(error),
},
});
console.error(
"event handler failed; sent error ack without advancing cursor",
error
);
}
return;
}
console.warn("unknown realtime frame", frame);
}
});
}
async function handleProductEvent(event) {
console.log("product event", {
id: event.id,
type: event.type,
created_at: event.created_at,
related: event.related,
});
}
function sendJson(socket, value) {
return new Promise((resolve, reject) => {
socket.send(JSON.stringify(value), (error) => {
if (error != null) {
reject(error);
return;
}
resolve();
});
});
}
async function readCursor() {
try {
const cursor = (await readFile(cursorPath, "utf8")).trim();
return cursor.length > 0 ? cursor : null;
} catch (error) {
if (error?.code === "ENOENT") {
return null;
}
throw error;
}
}
async function writeCursor(cursor) {
await writeFile(cursorPath, `${cursor}\n`, { mode: 0o600 });
}
function logWarnings(warnings) {
for (const warning of warnings ?? []) {
console.warn("realtime warning", warning);
}
}
function parseJson(text) {
try {
return JSON.parse(text);
} catch {
return null;
}
}
function getErrorMessage(error) {
if (error instanceof Error) {
return error.message;
}
return "Unknown handler error.";
}
function requiredEnv(name) {
const value = process.env[name];
if (value == null || value.trim() === "") {
throw new Error(`${name} is required`);
}
return value;
}
function sleep(ms) {
return new Promise((resolve) => setTimeout(resolve, ms));
}
EOF