1/**
2 * The "source of truth" server is the canonical storage for
3 * paper clover's files. This is technically needed because
4 * the VPS she uses can only store about 20gb of content, where
5 * the contents of /file is about 48gb as of writing; only a
6 * limited amount of data can be cached.
7 *
8 * What's great about this system is it also allows scaling the
9 * website up into multiple servers, if that is ever desired.
10 * When that happens, mutations to "q+a" will be moved here,
11 * and the SQLite synchronization mechanism will apply to both
12 * of those databases.
13 *
14 * An alternative to this would have been to use a SMB client,
15 * but I've read that the systems used for caching don't work
16 * like the HTTP Cache-Control header, where you can say a file
17 * is valid for up to a certain amount of time. If we seriously
18 * need cache busts (paper clover does not), the proper way
19 * would be to push a message to all VPS nodes instead of
20 * checking upstream if a file changed every time.
21 * @module
22 */
23const app = new Hono();
24export default app;
25export const port = Number(process.env.PORT ?? 4000);
26
27const token = UNWRAP(process.env.CLOVER_SOT_KEY);
28
29if (!fs.existsSync(rawFileRoot)) {
30 throw new Error(`${rawFileRoot} does not exist`);
31}
32if (!fs.existsSync(derivedFileRoot)) {
33 throw new Error(`${derivedFileRoot} does not exist`);
34}
35
36/**
37 * `node run source-of-truth` boots the full source of truth: the indexer
38 * service (file watcher, periodic sweeps, processors) plus this http server.
39 * set CLOVER_INDEXER=0 for a serve-only process.
40 */
41export async function main() {
42 if (process.env.CLOVER_INDEXER !== "0") {
43 service = indexer.start();
44 }
45 const server = await http.serve({
46 respond: async (request) => app.fetch(request),
47 port,
48 });
49 console.info(`source of truth live at ${server.url}`);
50 const close = () => void server.close().then(() => process.exit(0));
51 process.on("SIGINT", close);
52 process.on("SIGTERM", close);
53}
54let service: indexer.Service | null = null;
55
56// Re-use file descriptors if the same file is being read twice.
57const fds = new Map<string, Awaitable<{ fd: number; refs: number }>>();
58
59app.get("/", (c) => {
60 return c.text("everything you love, loves you back.");
61});
62
63app.get("/file/*", async (c) => {
64 const fullQuery = c.req.path.slice("/file".length);
65 const [filePath, derivedAsset, ...invalid] = fullQuery.split(":/");
66 ASSERT(filePath != null);
67 if (invalid.length > 0) return c.notFound();
68 if (filePath.length <= 1) return c.notFound();
69 const permissions = FilePermissions.getByPrefix(filePath);
70 if (permissions !== 0 && c.req.header("Authorization") !== token) {
71 return c.json({ error: "invalid authorization header" }, 401);
72 }
73 const file = MediaFile.getByPath(filePath);
74 if (!file || file.kind === MediaFileKind.directory) return c.notFound();
75 // derived asset names are plain relative paths under the hash directory
76 if (
77 derivedAsset
78 && derivedAsset.split("/").some((part) => part === "" || part === "." || part === "..")
79 ) {
80 return c.notFound();
81 }
82 const fullPath = derivedAsset
83 ? path.join(derivedFileRoot, file.hash, derivedAsset)
84 : path.join(rawFileRoot, file.path);
85 let handle: { fd: number; refs: number } | null = null;
86 try {
87 handle = await fds.get(fullPath) ?? null;
88 if (!handle) {
89 const promise = openFile(fullPath, "r")
90 .then((fd) => ({ fd, refs: 0 }))
91 .catch((err) => {
92 fds.delete(fullPath);
93 throw err;
94 });
95 fds.set(fullPath, promise);
96 fds.set(fullPath, handle = await promise);
97 }
98 handle.refs += 1;
99 } catch (err: any) {
100 if (err.code === "ENOENT") return c.notFound();
101 throw err;
102 }
103 c.header("Content-Length", (await fstat(handle.fd)).size.toString());
104 const nodeStream = fs.createReadStream(fullPath, {
105 fd: handle.fd,
106 fs: {
107 close: util.callbackify(async () => {
108 ASSERT(handle);
109 if ((handle.refs -= 1) <= 0) {
110 fds.delete(fullPath);
111 await closeFile(handle.fd);
112 }
113 handle = null;
114 }),
115 read: fsCallbacks.read,
116 },
117 });
118 c.header("Content-Type", mime.for(fullPath));
119 return c.body(stream.Readable.toWeb(nodeStream) as ReadableStream);
120});
121
122app.post("/reload", async (c) => {
123 if (c.req.header("Authorization") !== token) {
124 return c.json({ error: "invalid authorization header" }, 401);
125 }
126 MediaFile.db.reload();
127 indexer.emitDbChange();
128 return c.body(null, 204);
129});
130
131// a tiny server-sent-events feed of database changes. consumers (web nodes,
132// dev servers) hold this open and pull /db whenever a beacon arrives; the
133// connection itself is the subscription, so the source of truth never needs
134// a list of its consumers. events carry the generation, but clients can
135// treat any event as "something changed" — /db pulls are ETag-guarded.
136app.get("/db/events", (c) => {
137 if (c.req.header("Authorization") !== token) {
138 return c.json({ error: "invalid authorization header" }, 401);
139 }
140 const encoder = new TextEncoder();
141 let cleanup: (() => void) | null = null;
142 const stream = new ReadableStream<Uint8Array>({
143 start(controller) {
144 const send = (text: string) => controller.enqueue(encoder.encode(text));
145 // an event on connect makes clients revalidate immediately, covering
146 // anything they missed while disconnected
147 send(`retry: 5000\n\n`);
148 send(`event: change\ndata: ${indexer.generation}\n\n`);
149 const unsubscribe = indexer.subscribeDbChange((generation) => {
150 send(`event: change\ndata: ${generation}\n\n`);
151 });
152 // keepalive comments so proxies and clients hold the connection open
153 const keepalive = setInterval(() => send(`: keepalive\n\n`), 25_000);
154 keepalive.unref();
155 cleanup = () => {
156 unsubscribe[Symbol.dispose]();
157 clearInterval(keepalive);
158 };
159 },
160 cancel() {
161 cleanup?.();
162 },
163 });
164 c.header("Content-Type", "text/event-stream");
165 c.header("Cache-Control", "no-store");
166 return c.body(stream);
167});
168
169// live progress of all indexing operations, as a clover progress stream.
170// each request snapshots the current tree state and then follows along;
171// connect any time with `node run tail-progress`.
172app.get("/progress", (c) => {
173 if (c.req.header("Authorization") !== token) {
174 return c.json({ error: "invalid authorization header" }, 401);
175 }
176 c.header("Content-Type", progress.contentType);
177 return c.body(progress.encodeByteStream(indexer.root) as ReadableStream);
178});
179
180// download a consistent snapshot of cache.sqlite. web nodes call this with
181// `If-None-Match` on a stale-while-revalidate schedule, and immediately
182// when pinged at /internal/db-changed.
183app.get("/db", async (c) => {
184 if (c.req.header("Authorization") !== token) {
185 return c.json({ error: "invalid authorization header" }, 401);
186 }
187 const snap = await getDbSnapshot();
188 if (c.req.header("If-None-Match") === snap.etag) {
189 return c.body(null, 304);
190 }
191 c.header("Content-Type", "application/vnd.sqlite3");
192 c.header("Content-Length", String(snap.size));
193 c.header("ETag", snap.etag);
194 return c.body(
195 stream.Readable.toWeb(fs.createReadStream(snap.path)) as ReadableStream,
196 );
197});
198
199// trigger a sweep (or a subtree scan with {"path": "/2025"}). progress is
200// visible on /progress; this returns immediately.
201app.post("/scan", async (c) => {
202 if (c.req.header("Authorization") !== token) {
203 return c.json({ error: "invalid authorization header" }, 401);
204 }
205 if (!service) return c.json({ error: "indexer is disabled" }, 409);
206 const body = await c.req.json().catch(() => ({}));
207 const subPath = typeof body.path === "string" ? body.path : undefined;
208 void service.sweep("manual", subPath);
209 return c.json({ ok: true }, 202);
210});
211
212// -- database snapshotting --
213// `VACUUM INTO` writes a compact, transactionally-consistent copy that does
214// not depend on -wal sidecar files. regenerated lazily when the indexer
215// reports changes, at most every 30 seconds.
216interface DbSnapshot {
217 gen: number;
218 time: number;
219 path: string;
220 etag: string;
221 size: number;
222}
223let snapshot: DbSnapshot | null = null;
224let snapshotInFlight: Promise<DbSnapshot> | null = null;
225
226function getDbSnapshot() {
227 const gen = indexer.generation;
228 // a stale-generation snapshot is never served: subscribers pull exactly
229 // once per beacon, so a 304 here would leave them out of date until the
230 // ttl backstop. beacons are debounced upstream, which bounds how often
231 // the vacuum can actually run.
232 if (snapshot && snapshot.gen === gen && fs.existsSync(snapshot.path)) {
233 return Promise.resolve(snapshot);
234 }
235 return snapshotInFlight ??= (async () => {
236 try {
237 const dir = process.env.CLOVER_DB ?? ".clover";
238 const file = path.join(dir, "db-snapshot.sqlite");
239 await fs.rm(file, { force: true });
240 MediaFile.db.node.exec(
241 `vacuum into '${file.replaceAll("'", "''")}';`,
242 );
243 const [hash, { size }] = await Promise.all([
244 hashFile(new Path(file)),
245 fs.stat(file),
246 ]);
247 return snapshot = {
248 gen,
249 time: Date.now(),
250 path: file,
251 etag: `"${hash}"`,
252 size,
253 };
254 } finally {
255 snapshotInFlight = null;
256 }
257 })();
258}
259
260const openFile = util.promisify(fsCallbacks.open);
261const closeFile = util.promisify(fsCallbacks.close);
262const fstat = util.promisify(fsCallbacks.fstat);
263
264import * as fs from "#sitegen/fs";
265import { Path } from "#sitegen/path";
266import { hashFile } from "#src/file-viewer/indexer/scan.ts";
267import * as indexer from "#src/file-viewer/indexer/service.ts";
268import { FilePermissions } from "#src/file-viewer/models/FilePermissions.ts";
269import { MediaFile, MediaFileKind } from "#src/file-viewer/models/MediaFile.ts";
270import { ASSERT, UNWRAP } from "@clo/lib/assert";
271import * as http from "@clo/lib/http";
272import * as mime from "@clo/lib/mime";
273import * as progress from "@clo/lib/progress";
274import { Hono } from "hono";
275import * as fsCallbacks from "node:fs";
276import * as path from "node:path";
277import * as stream from "node:stream";
278import * as util from "node:util";
279import { derivedFileRoot, rawFileRoot } from "./file-viewer/paths.ts";