| 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 | |
| 10 | const db = getDb("cache.sqlite"); |
| 11 | db.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. |
| 45 | export const workDir = Path.resolve(derivedFileRoot); |
| 46 | |
| 47 | let ongoing = new Map<string, Promise<number>>(); |
| 48 | |
| 49 | /** produce a derived */ |
| 50 | export 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 | |
| 162 | const 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 | |
| 173 | export 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 | |
| 186 | export function findOrphanedRoots() { |
| 187 | return findOrphanedRootsQuery.array(); |
| 188 | } |
| 189 | |
| 190 | export 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 */ |
| 195 | export 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 | */ |
| 212 | export 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 | |
| 230 | const 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 | |
| 237 | const getRootQuery = db.prepare<[{ key: string }], { id: number }>(/* SQL */ ` |
| 238 | select id from derived_roots where key = $key; |
| 239 | `); |
| 240 | |
| 241 | const 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 | |
| 247 | const insertRefQuery = db.prepare<[{ root: number; file: number }]>(/* SQL */ ` |
| 248 | insert into derived_refs (file, root) values ($file, $root); |
| 249 | `); |
| 250 | |
| 251 | const deleteRootQuery = db.prepare<[{ root: number }]>(/* SQL */ ` |
| 252 | delete from derived_roots where id = $root; |
| 253 | `); |
| 254 | |
| 255 | const deleteRefQuery = db.prepare<[{ root: number; file: number }]>(/* SQL */ ` |
| 256 | delete from derived_refs where root = $root and file = $file; |
| 257 | `); |
| 258 | |
| 259 | const 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 | |
| 269 | import { Path } from "#sitegen/path"; |
| 270 | import { getDb } from "#sitegen/sqlite"; |
| 271 | import { ASSERT, UNWRAP } from "@clo/lib/assert"; |
| 272 | import * as progress from "@clo/lib/progress"; |
| 273 | import * as queue from "@clo/lib/queue"; |
| 274 | import * as crypto from "node:crypto"; |
| 275 | import * as fs from "node:fs"; |
| 276 | import * as fsp from "node:fs/promises"; |
| 277 | import { derivedFileRoot } from "../paths.ts"; |
| 278 | import type { MediaFile } from "./MediaFile.ts"; |