1// Uncompressed files are read directly from the media store root. Derivied
2// assets like compressed files, optimized images, and streamable video are
3// stored in the `derived` folder. After scanning, the derived assets are
4// uploaded into the store (storage1/clofi-derived dataset on NAS). Since
5// multiple files can share the same hash, the number of references is
6// tracked, and the derived content is only produced once. This means if a
7// file is deleted, it should only decrement a reference count; deleting it
8// once all references are removed.
9
10const db = getDb("cache.sqlite");
11db.table(
12 "asset_refs",
13 /* SQL */ `
14 create table if not exists derived_roots (
15 id integer primary key autoincrement,
16 date integer not null, -- milliseconds
17 key text not null unique
18 );
19 create table if not exists derived_files (
20 id integer primary key autoincrement,
21 file text not null,
22 root integer not null,
23 size integer not null, -- bytes
24 hash text not null, -- sha1
25 foreign key(root) references derived_roots(id) on delete cascade
26 );
27 create table if not exists derived_refs (
28 id integer primary key autoincrement,
29 file integer not null,
30 root integer not null,
31 foreign key(root) references derived_roots(id) on delete cascade,
32 foreign key(file) references media_files(id) on delete cascade,
33 unique(file, root) on conflict replace
34 );
35 create index derived_roots_key on derived_roots(key);
36 create index derived_files_root on derived_files(root);
37 create index derived_refs_file on derived_refs(file);
38 create index derived_refs_root on derived_refs(root);
39 `,
40);
41
42// derived assets are written directly into the derived file store. files are
43// produced into a `tmp.`-prefixed sibling directory and renamed into place,
44// so a crash never leaves a partially-written asset at its final path.
45export const workDir = Path.resolve(derivedFileRoot);
46
47let ongoing = new Map<string, Promise<number>>();
48
49/** produce a derived */
50export async function produce(
51 { mediaFile: file, node, subkey, cores, producer }: {
52 mediaFile: MediaFile;
53 node: progress.Node;
54 subkey: string;
55 cores: number;
56 producer: (root: Path) => Promise<void>;
57 },
58) {
59 const key = `${file.hash}/${subkey}`;
60 const current = ongoing.get(key);
61 let root: number | null = null;
62 if (current) {
63 root = await current;
64 } else {brk: {
65 root = getRootQuery.get({ key })?.id ?? null;
66 if (root) break brk;
67 const { promise, resolve, reject } = Promise.withResolvers<number>();
68 ongoing.set(key, promise);
69 // the rejection is rethrown to this caller; concurrent callers await
70 // through `ongoing`, but when there are none the rejection would
71 // otherwise be unobserved
72 promise.catch(() => {});
73 const finalDir = workDir.join(key);
74 const tmp = workDir.join(`${file.hash}/tmp.${subkey}`);
75 try {
76 node.hidden = true;
77 const filesWithStats = await queue.run({
78 cores,
79 async run() {
80 node.hidden = false;
81 await tmp.delete({ recursive: true, force: true });
82 await tmp.makeOrEmptyDir();
83 await producer(tmp);
84 const sha1queue = new queue.PriorityQueue(cores);
85 return await Promise.all(
86 (await tmp.readDir({ recursive: true }))
87 .filter((path) => !path.base.startsWith("tmp."))
88 .map(async (path) =>
89 sha1queue.run({
90 async run() {
91 using _ = node.start("hash " + path.relative(tmp));
92 const stats = await path.stat();
93 if (stats.isDirectory()) {
94 return [path, stats.size, null] as const;
95 }
96 const hash = await new Promise<string>(
97 (resolve, reject) => {
98 const reader = fs.createReadStream(path.toString());
99 reader.on("error", reject);
100
101 const hasher = crypto.createHash("sha1").setEncoding(
102 "hex",
103 );
104 hasher.on("error", reject);
105 hasher.on("readable", () => resolve(hasher.read()));
106
107 reader.pipe(hasher);
108 },
109 );
110 return [path, stats.size, hash] as const;
111 },
112 })
113 ),
114 );
115 },
116 });
117 // move into the final location before recording rows; readers only
118 // discover the asset through the database, so this is safe.
119 await finalDir.delete({ recursive: true, force: true });
120 await fsp.rename(tmp.toString(), finalDir.toString());
121
122 db.node.exec("BEGIN");
123 try {
124 root = insertRootQuery.getNonNull({ key, date: Date.now() }).id;
125 for (const [path, size, hash] of filesWithStats) {
126 if (hash === null) continue;
127 insertFileQuery.run({
128 root,
129 file: path.relative(tmp),
130 size,
131 hash,
132 });
133 }
134 } catch (e) {
135 db.node.exec("ROLLBACK");
136 throw e;
137 }
138 db.node.exec("COMMIT");
139 resolve(root);
140 } catch (e) {
141 if (root) deleteRootQuery.run({ root });
142 // remove partial output from disk: the tmp dir if the producer died,
143 // or the renamed final dir if the database writes failed (e.g. the
144 // source file was deleted while this asset was being produced)
145 await tmp.delete({ recursive: true, force: true }).catch(() => {});
146 await finalDir.delete({ recursive: true, force: true }).catch(() => {});
147 reject(e);
148 throw e;
149 } finally {
150 ongoing.delete(key);
151 }
152 }}
153 ASSERT(root != null);
154 insertRefQuery.run({ root, file: file.id });
155 return {
156 [Symbol.dispose]: () => {
157 deleteRefQuery.run({ root, file: file.id });
158 },
159 };
160}
161
162const getDerivedAssetQuery = db.prepare<
163 [number, string],
164 { size: number; key: string; hash: string; date: number }
165>(/* SQL */ `
166 select df.size, df.hash, rt.date, rt.key
167 from derived_refs dr
168 join derived_files df on df.root = dr.root
169 join derived_roots rt on df.root = rt.id
170 where dr.file = ? and df.file = ?
171`);
172
173export function get(file: MediaFile, subPath: string) {
174 const row = getDerivedAssetQuery.get(file.id, subPath) ?? null;
175 if (row) {
176 return {
177 path: row.key + "/" + subPath,
178 size: row.size,
179 hash: row.hash,
180 date: new Date(row.date),
181 };
182 }
183 return null;
184}
185
186export function findOrphanedRoots() {
187 return findOrphanedRootsQuery.array();
188}
189
190export function deleteRoot(root: { id: number }) {
191 return deleteRootQuery.run({ root: root.id });
192}
193
194/** remove an orphaned root's files from the derived store */
195export async function deleteRootFiles(root: { key: string }) {
196 const dir = workDir.join(root.key);
197 await dir.delete({ recursive: true, force: true });
198 // reclaim the per-hash parent directory once its last root is gone
199 // (rmdir refuses non-empty directories, which is exactly what we want)
200 await fsp.rmdir(UNWRAP(dir.parent).toString()).catch(() => {});
201}
202
203/**
204 * delete `tmp.*` directories left behind by crashed producers. only removes
205 * directories untouched for a day, to never race an ongoing producer.
206 *
207 * `tmp.*` FILES are never touched: encode intermediates inside completed
208 * roots (`av1-au/tmp.av1.mp4`, ...) are intentionally named that way to be
209 * excluded from `derived_files`, but the dash producer reads them across
210 * root directories, so they must persist.
211 */
212export async function cleanAbandonedTmp(): Promise<number> {
213 let cleaned = 0;
214 const dayAgo = Date.now() - 24 * 60 * 60 * 1000;
215 for (const hashDir of await workDir.readDir().catch(() => [])) {
216 for (const sub of await hashDir.readDir().catch(() => [])) {
217 if (!sub.base.startsWith("tmp.")) continue;
218 try {
219 const stat = await sub.stat();
220 if (!stat.isDirectory()) continue;
221 if (stat.mtime.getTime() > dayAgo) continue;
222 await sub.delete({ recursive: true, force: true });
223 cleaned += 1;
224 } catch {}
225 }
226 }
227 return cleaned;
228}
229
230const insertRootQuery = db.prepare<
231 [{ key: string; date: number }],
232 { id: number }
233>(
234 /* SQL */ `insert into derived_roots (key, date) values ($key, $date) returning id;`,
235);
236
237const getRootQuery = db.prepare<[{ key: string }], { id: number }>(/* SQL */ `
238 select id from derived_roots where key = $key;
239`);
240
241const insertFileQuery = db.prepare<
242 [{ file: string; root: number; size: number; hash: string }]
243>(/* SQL */ `
244 insert into derived_files (file, root, size, hash) values ($file, $root, $size, $hash);
245`);
246
247const insertRefQuery = db.prepare<[{ root: number; file: number }]>(/* SQL */ `
248 insert into derived_refs (file, root) values ($file, $root);
249`);
250
251const deleteRootQuery = db.prepare<[{ root: number }]>(/* SQL */ `
252 delete from derived_roots where id = $root;
253`);
254
255const deleteRefQuery = db.prepare<[{ root: number; file: number }]>(/* SQL */ `
256 delete from derived_refs where root = $root and file = $file;
257`);
258
259const findOrphanedRootsQuery = db.prepare<
260 [],
261 { id: number; key: string; date: number }
262>(/* SQL */ `
263 select dr.id, dr.key, dr.date
264 from derived_roots dr
265 left join derived_refs refs on refs.root = dr.id
266 where refs.id is null
267`);
268
269import { Path } from "#sitegen/path";
270import { getDb } from "#sitegen/sqlite";
271import { ASSERT, UNWRAP } from "@clo/lib/assert";
272import * as progress from "@clo/lib/progress";
273import * as queue from "@clo/lib/queue";
274import * as crypto from "node:crypto";
275import * as fs from "node:fs";
276import * as fsp from "node:fs/promises";
277import { derivedFileRoot } from "../paths.ts";
278import type { MediaFile } from "./MediaFile.ts";