| 1 | // Keeps the web node's copy of cache.sqlite in sync with the source of |
| 2 | // truth. The SOT serves a server-sent-events feed at /db/events; this |
| 3 | // module holds it open and pulls /db on every beacon. The server emits a |
| 4 | // beacon on connect, so boot, reconnect, and missed-while-disconnected all |
| 5 | // resolve through the same path — there is no polling and no TTL. |
| 6 | // |
| 7 | // Pulls are ETag-guarded (no-change is a tiny 304; the etag is persisted so |
| 8 | // a reboot does not re-download an unchanged database), land in a temp |
| 9 | // file, and atomically rename over .clover/cache.sqlite before hot-swapping |
| 10 | // the open database handle. |
| 11 | const console = log.scoped("sync"); |
| 12 | |
| 13 | const token = process.env.CLOVER_SOT_KEY ?? null; |
| 14 | const enabled = process.env.CLOVER_DB_SYNC !== "0" && token != null; |
| 15 | const sotUrl = new URL(process.env.CLOVER_SOT_URL ?? "https://db.paperclover.net") |
| 16 | .toString().replace(/\/+$/g, ""); |
| 17 | if (!enabled) { |
| 18 | console.warn( |
| 19 | "database sync disabled" |
| 20 | + (token ? " (CLOVER_DB_SYNC=0)" : " (no CLOVER_SOT_KEY)"), |
| 21 | ); |
| 22 | } |
| 23 | |
| 24 | let inFlight: Promise<void> | null = null; |
| 25 | |
| 26 | export function revalidate(): Promise<void> { |
| 27 | return inFlight ??= revalidateInner().finally(() => { |
| 28 | inFlight = null; |
| 29 | }); |
| 30 | } |
| 31 | |
| 32 | async function revalidateInner(): Promise<void> { |
| 33 | const db = getDb("cache.sqlite"); |
| 34 | const res = await fetch(`${sotUrl}/db`, { |
| 35 | headers: { |
| 36 | Authorization: token!, |
| 37 | ...etag() ? { "If-None-Match": etag()! } : null, |
| 38 | }, |
| 39 | signal: AbortSignal.timeout(120_000), |
| 40 | }); |
| 41 | if (res.status === 304) return; |
| 42 | if (!res.ok || !res.body) { |
| 43 | throw new Error(`source of truth responded ${res.status} ${res.statusText}`); |
| 44 | } |
| 45 | |
| 46 | const tmp = db.file + ".tmp"; |
| 47 | await stream.promises.pipeline( |
| 48 | stream.Readable.fromWeb(res.body as any), |
| 49 | fs.createWriteStream(tmp), |
| 50 | ); |
| 51 | // fsync before the rename so a power cut cannot leave a torn file |
| 52 | const handle = await fsp.open(tmp, "r+"); |
| 53 | await handle.sync(); |
| 54 | await handle.close(); |
| 55 | await fsp.rename(tmp, db.file); |
| 56 | db.reload(); |
| 57 | saveEtag(res.headers.get("ETag")); |
| 58 | console.info(`database updated (${string.formatByteSize(Number(res.headers.get("Content-Length") ?? "0"))})`); |
| 59 | } |
| 60 | |
| 61 | // -- the persisted etag, stored beside the database -- |
| 62 | let cachedEtag: string | null | undefined; |
| 63 | function etagPath() { |
| 64 | return getDb("cache.sqlite").file + ".etag"; |
| 65 | } |
| 66 | function etag(): string | null { |
| 67 | if (cachedEtag === undefined) { |
| 68 | try { |
| 69 | cachedEtag = fs.readFileSync(etagPath(), "utf-8").trim() || null; |
| 70 | } catch { |
| 71 | cachedEtag = null; |
| 72 | } |
| 73 | } |
| 74 | return cachedEtag; |
| 75 | } |
| 76 | function saveEtag(value: string | null) { |
| 77 | cachedEtag = value; |
| 78 | try { |
| 79 | fs.writeFileSync(etagPath(), value ?? ""); |
| 80 | } catch (err) { |
| 81 | console.warn("could not persist database etag:", err); |
| 82 | } |
| 83 | } |
| 84 | |
| 85 | // -- the subscription -- |
| 86 | // a hand-rolled sse reader over fetch (EventSource cannot send the |
| 87 | // Authorization header). any event block triggers a revalidate; a failed |
| 88 | // revalidate aborts the connection so the reconnect path retries it. |
| 89 | async function subscribeLoop() { |
| 90 | let backoff = 1000; |
| 91 | while (true) { |
| 92 | const ctrl = new AbortController(); |
| 93 | let watchdog: NodeJS.Timeout | null = null; |
| 94 | try { |
| 95 | const res = await fetch(`${sotUrl}/db/events`, { |
| 96 | headers: { Authorization: token! }, |
| 97 | signal: ctrl.signal, |
| 98 | }); |
| 99 | if (!res.ok || !res.body) { |
| 100 | throw new Error(`${res.status} ${res.statusText}`); |
| 101 | } |
| 102 | console.info(`subscribed to database changes from ${sotUrl}`); |
| 103 | backoff = 1000; |
| 104 | |
| 105 | // the server keepalives every 25s; 90s of silence means the |
| 106 | // connection died without a FIN |
| 107 | const armWatchdog = () => { |
| 108 | if (watchdog) clearTimeout(watchdog); |
| 109 | watchdog = setTimeout( |
| 110 | () => ctrl.abort(new Error("no keepalive")), |
| 111 | 90_000, |
| 112 | ); |
| 113 | watchdog.unref(); |
| 114 | }; |
| 115 | armWatchdog(); |
| 116 | const reader = res.body.getReader(); |
| 117 | const decoder = new TextDecoder(); |
| 118 | let buffer = ""; |
| 119 | while (true) { |
| 120 | const { value, done } = await reader.read(); |
| 121 | if (done) break; |
| 122 | armWatchdog(); |
| 123 | buffer += decoder.decode(value, { stream: true }); |
| 124 | const blocks = buffer.split("\n\n"); |
| 125 | buffer = blocks.pop()!; |
| 126 | const changed = blocks.some((block) => |
| 127 | block.split("\n").some((line) => line.startsWith("event:") || line.startsWith("data:")) |
| 128 | ); |
| 129 | if (changed) { |
| 130 | void revalidate().catch((err) => { |
| 131 | console.error("revalidate failed:", err); |
| 132 | ctrl.abort(new Error("revalidate failed")); |
| 133 | }); |
| 134 | } |
| 135 | } |
| 136 | throw new Error("stream ended"); |
| 137 | } catch (err: any) { |
| 138 | console.warn( |
| 139 | `database subscription lost (${err?.message ?? err}); ` |
| 140 | + `retrying in ${Math.round(backoff / 1000)}s`, |
| 141 | ); |
| 142 | } finally { |
| 143 | if (watchdog) clearTimeout(watchdog); |
| 144 | ctrl.abort(); |
| 145 | } |
| 146 | await new Promise<void>((resolve) => setTimeout(resolve, backoff).unref()); |
| 147 | backoff = Math.min(backoff * 2, 60_000); |
| 148 | } |
| 149 | } |
| 150 | |
| 151 | if (enabled) void subscribeLoop(); |
| 152 | |
| 153 | import * as fs from "node:fs"; |
| 154 | import * as fsp from "node:fs/promises"; |
| 155 | import * as stream from "node:stream"; |
| 156 | |
| 157 | import { getDb } from "#sitegen/sqlite"; |
| 158 | |
| 159 | import * as log from "@clo/lib/log"; |
| 160 | import * as string from "@clo/lib/string"; |