import http from "node:http" import { open, mkdir, rm, rename } from "node:fs/promises" import { createReadStream, createWriteStream, existsSync } from "node:fs" import { extname, join, normalize } from "node:path" import { execFile } from "node:child_process" import { randomUUID } from "node:crypto" import { promisify } from "node:util" import { once } from "node:events" import { networkInterfaces, hostname } from "node:os" import mysql from "mysql2/promise" const execFileAsync = promisify(execFile) const DB = { host: process.env.DB_HOST || "127.0.0.1", port: Number(process.env.DB_PORT || 3306), user: process.env.DB_USER || "photoprism", password: process.env.DB_PASSWORD || "", database: process.env.DB_NAME || "photoprism", connectionLimit: 4, } const PORT = Number(process.env.PORT || 3100) const PUBLIC = join(process.cwd(), "public") const MIN_PATH = 3 const STAGING_ROOT = process.env.STAGING_ROOT || join(process.cwd(), "staging") const PHOTOPRISM_BIN = process.env.PHOTOPRISM_BIN || "photoprism" const PHOTOPRISM_USER = process.env.PHOTOPRISM_USER || "photoprism" const PHOTOPRISM_HOME = process.env.PHOTOPRISM_HOME || "/srv/photoprism" const ORIGINALS_PATH = process.env.PHOTOPRISM_ORIGINALS_PATH || "/srv/photoprism/originals" const STORAGE_PATH = process.env.PHOTOPRISM_STORAGE_PATH || "/srv/photoprism/storage" const MAX_UPLOAD_BYTES = Number(process.env.MAX_UPLOAD_BYTES || 1024 * 1024 * 1024) const MAX_UPLOAD_FILES = Number(process.env.MAX_UPLOAD_FILES || 100) const IMPORT_TIMEOUT_MS = Number(process.env.IMPORT_TIMEOUT_MS || 30 * 60 * 1000) function directOrigins() { const lan = [] const tailscale = [] for (const list of Object.values(networkInterfaces())) { for (const a of list || []) { if ((a.family !== "IPv4" && a.family !== 4) || a.internal) continue if (a.address.startsWith("100.")) tailscale.push(a.address) else lan.push(a.address) } } const pick = lan.find((a) => /^(10\.|192\.168\.|172\.(1[6-9]|2\d|3[01])\.)/.test(a)) return { magic: `http://${hostname()}:8080`, lan: pick ? `http://${pick}:8080` : null, tailscale: tailscale.length ? `http://${tailscale[0]}:8080` : null, } } const DIRECT = directOrigins() const pool = mysql.createPool(DB) const THUMB_SIZE = "1280x1024_fit" const BIG_SIZE = "1920x1200_fit" const thumb = (hash, size = THUMB_SIZE) => { if (!hash) return null const s = hash.toString() return `/thumbs/${s[0]}/${s[1]}/${s[2]}/${s}_${size}.jpg` } let state = null async function load() { const [photos] = await pool.query( `SELECT p.id, p.photo_title, p.photo_year, p.photo_month, p.photo_day, f.file_hash, f.file_width, f.file_height FROM photos p JOIN files f ON f.photo_id = p.id AND f.file_primary = 1 AND f.file_missing = 0 WHERE p.deleted_at IS NULL ORDER BY p.taken_at_local, p.id` ) const byId = new Map() const photoOrder = [] for (const p of photos) { byId.set(p.id, { id: p.id, thumb: thumb(p.file_hash), full: thumb(p.file_hash, BIG_SIZE), title: p.photo_title || "", date: p.photo_year ? `${p.photo_year}-${String(p.photo_month || 1).padStart(2, "0")}-${String(p.photo_day || 1).padStart(2, "0")}` : "", w: p.file_width || 0, h: p.file_height || 0, }) photoOrder.push(p.id) } const rank = new Map(photoOrder.map((id, i) => [id, i])) const paths = new Map() const photoPaths = new Map() function addRows(rows, type, label) { const grouped = new Map() for (const r of rows) { if (!byId.has(r.pid)) continue if (!grouped.has(r.k)) grouped.set(r.k, { key: `${type}:${r.k}`, type, label, name: r.name, ids: [] }) grouped.get(r.k).ids.push(r.pid) } for (const g of grouped.values()) { g.ids.sort((a, b) => rank.get(a) - rank.get(b)) g.count = g.ids.length paths.set(g.key, g) for (const pid of g.ids) { if (!photoPaths.has(pid)) photoPaths.set(pid, []) photoPaths.get(pid).push(g) } } } const [lab] = await pool.query( `SELECT pl.label_id AS k, l.label_name AS name, pl.photo_id AS pid FROM photos_labels pl JOIN labels l ON l.id = pl.label_id`) addRows(lab, "label", "Labels") const [kw] = await pool.query( `SELECT pk.keyword_id AS k, kw.keyword AS name, pk.photo_id AS pid FROM photos_keywords pk JOIN keywords kw ON kw.id = pk.keyword_id`) addRows(kw, "keyword", "Keywords") const [cam] = await pool.query( `SELECT p.camera_id AS k, CONCAT(c.camera_make, ' ', c.camera_model) AS name, p.id AS pid FROM photos p JOIN cameras c ON c.id = p.camera_id WHERE p.camera_id <> 1 AND p.camera_id IS NOT NULL`) addRows(cam, "camera", "Cameras") const [len] = await pool.query( `SELECT p.lens_id AS k, l.lens_model AS name, p.id AS pid FROM photos p JOIN lenses l ON l.id = p.lens_id WHERE p.lens_id <> 1 AND p.lens_id IS NOT NULL`) addRows(len, "lens", "Lenses") const facets = [] for (const type of ["label", "keyword", "camera", "lens"]) { const list = Array.from(paths.values()) .filter((g) => g.type === type && g.count >= MIN_PATH) .sort((a, b) => b.count - a.count || a.name.localeCompare(b.name)) for (const g of list) facets.push({ type: g.type, label: g.label, name: g.name, count: g.count, ids: g.ids }) } const library = { key: "library:all", type: "library", label: "Library", name: "library", ids: photoOrder, count: photoOrder.length } const DIRS = ["N", "E", "S", "W"] const exits = new Map() for (const pid of photoOrder) { const own = (photoPaths.get(pid) || []) .filter((p) => p.count >= 2) .sort((a, b) => b.count - a.count || a.name.localeCompare(b.name) || a.type.localeCompare(b.type)) const chosen = own.slice(0, 4) while (chosen.length < 4) chosen.push(library) const used = new Map() const usedNext = new Set() const list = [] for (let i = 0; i < 4; i++) { const p = chosen[i] const step = used.get(p.key) || 0 used.set(p.key, step + 1) const idx = p.ids.indexOf(pid) const cyc = p.ids const start = (idx + 1 + step) % cyc.length const walk = [] for (let k = 0; k < 60 && walk.length < cyc.length; k++) walk.push(cyc[(start + k) % cyc.length]) let nextId = cyc[start] for (const wid of walk) { if (wid !== pid && !usedNext.has(wid)) { nextId = wid; break } } usedNext.add(nextId) const wi = walk.indexOf(nextId) const rot = wi > 0 ? walk.slice(wi).concat(walk.slice(0, wi)) : walk const li = library.ids.indexOf(pid) const ext = rot.slice() for (let k = 1; ext.length < 60 && k < library.ids.length; k++) { const wid = library.ids[(li + k) % library.ids.length] if (wid !== pid && !ext.includes(wid)) ext.push(wid) } list.push({ dir: DIRS[i], path: { type: p.type, name: p.name, count: p.count }, nextId, walk: ext, }) } exits.set(pid, list) } return { builtAt: new Date().toISOString(), total: photoOrder.length, photos: Object.fromEntries(byId), facets, photoOrder, byId, exits } } async function refresh() { try { state = await load() } catch (err) { console.error("pathways: refresh failed:", err.message) } } const MIME = { ".html": "text/html; charset=utf-8", ".css": "text/css", ".js": "text/javascript", ".svg": "image/svg+xml", ".ico": "image/x-icon" } function serveFile(res, path) { if (!existsSync(path)) { res.writeHead(404).end("not found") return } const type = MIME[extname(path)] || "application/octet-stream" const cache = path.endsWith(".html") ? "no-cache" : "max-age=300" res.writeHead(200, { "Content-Type": type, "Cache-Control": cache }) createReadStream(path).pipe(res) } function json(res, code, data) { res.writeHead(code, { "Content-Type": "application/json; charset=utf-8", "Cache-Control": "no-store" }) res.end(JSON.stringify(data)) } const PNG_MAGIC = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a]) const SIGNATURES = [ { ext: "jpg", accepts: /\.jpe?g$/i, test: (b) => b.length > 3 && b[0] === 0xff && b[1] === 0xd8 && b[2] === 0xff }, { ext: "png", accepts: /\.png$/i, test: (b) => b.length > 8 && b.subarray(0, 8).equals(PNG_MAGIC) }, { ext: "webp", accepts: /\.webp$/i, test: (b) => b.length > 12 && b.subarray(0, 4).toString("latin1") === "RIFF" && b.subarray(8, 12).toString("latin1") === "WEBP", }, ] function sniff(buf) { return SIGNATURES.find((s) => s.test(buf)) || null } async function finalizePart(dir, cur, taken) { const path = join(dir, cur.name) if (cur.written === 0) { await rm(path, { force: true }).catch(() => {}) return null } const head = Buffer.alloc(16) try { const fh = await open(path, "r") let bytesRead = 0 try { ({ bytesRead } = await fh.read(head, 0, head.length, 0)) } finally { await fh.close() } const sig = sniff(head.subarray(0, bytesRead)) if (!sig) throw Object.assign(new Error(`unsupported file type: ${cur.raw}`), { status: 415 }) const final = safeName(cur.raw, sig, taken) await rename(path, join(dir, final)) return { name: final } } catch (err) { await rm(path, { force: true }).catch(() => {}) throw err } } async function ingestMultipart(req, boundary, dir) { const DELIM = Buffer.from(`--${boundary}`) const CLOSE = Buffer.from(`\r\n--${boundary}`) const SEP = Buffer.from("\r\n\r\n") const HEADER_LIMIT = 64 * 1024 const taken = new Set() const files = [] let pending = Buffer.alloc(0) let state = "init" let current = null let totalBytes = 0 const writeChunk = async (cur, buf) => { cur.written += buf.length if (!cur.stream.write(buf)) await once(cur.stream, "drain") } const closePart = async () => { current.stream.end() await once(current.stream, "finish") const rec = await finalizePart(dir, current, taken) if (rec) files.push(rec) current = null } try { for await (const chunk of req) { pending = pending.length ? Buffer.concat([pending, chunk]) : chunk while (pending.length && state !== "done") { if (state === "init") { const i = pending.indexOf(DELIM) if (i < 0) { pending = pending.subarray(Math.max(0, pending.length - (DELIM.length - 1))) break } pending = pending.subarray(i + DELIM.length) state = "postBoundary" continue } if (state === "postBoundary") { if (pending.length < 2) break if (pending[0] === 0x2d && pending[1] === 0x2d) { state = "done"; break } if (pending[0] === 0x0d && pending[1] === 0x0a) { pending = pending.subarray(2) state = "headers" continue } throw Object.assign(new Error("malformed multipart"), { status: 400 }) } if (state === "headers") { const h = pending.indexOf(SEP) if (h < 0) { if (pending.length > HEADER_LIMIT) throw Object.assign(new Error("part headers too long"), { status: 400 }) break } const headerBuf = pending.subarray(0, h) pending = pending.subarray(h + 4) const raw = fileNameOf(headerBuf.toString("utf8")) if (raw) { if (files.length >= MAX_UPLOAD_FILES) throw Object.assign(new Error(`too many files (max ${MAX_UPLOAD_FILES})`), { status: 413 }) current = { raw, name: `part-${files.length}.tmp`, stream: createWriteStream(join(dir, `part-${files.length}.tmp`)), written: 0 } } state = "body" continue } if (state === "body") { const j = pending.indexOf(CLOSE) if (j < 0) { const tail = Math.min(pending.length, CLOSE.length - 1) const safe = pending.length - tail if (safe > 0) { totalBytes += safe if (totalBytes > MAX_UPLOAD_BYTES) throw Object.assign(new Error(`upload exceeds ${MAX_UPLOAD_BYTES} bytes`), { status: 413 }) if (current) await writeChunk(current, pending.subarray(0, safe)) pending = pending.subarray(safe) } break } const body = pending.subarray(0, j) if (body.length) { totalBytes += body.length if (totalBytes > MAX_UPLOAD_BYTES) throw Object.assign(new Error(`upload exceeds ${MAX_UPLOAD_BYTES} bytes`), { status: 413 }) if (current) await writeChunk(current, body) } pending = pending.subarray(j + CLOSE.length) if (current) await closePart() state = "postBoundary" continue } } } } catch (err) { if (current) current.stream.destroy() throw err } if (current) { current.stream.destroy() throw Object.assign(new Error("multipart truncated"), { status: 400 }) } if (state !== "done") throw Object.assign(new Error("multipart truncated"), { status: 400 }) return files } function fileNameOf(headers) { const m = /filename\*?=(?:UTF-8''([^\r\n;]+)|"([^"]*)"|([^\r\n;]+))/i.exec(headers) if (!m) return null const raw = m[1] ? decodeURIComponent(m[1]) : (m[2] ?? m[3] ?? "") return raw.trim() } function safeName(raw, sig, taken) { const base = (raw.split(/[\\/]/).pop() || "").replace(/[^A-Za-z0-9._-]/g, "_").replace(/^\.+/, "").slice(0, 120) let name = base || "upload" if (!sig.accepts.test(name)) name += `.${sig.ext}` let out = name let n = 1 while (taken.has(out)) { const dot = name.lastIndexOf(".") out = `${name.slice(0, dot)}-${n}${name.slice(dot)}` n++ } taken.add(out) return out } function photoprismCopy(dir) { const env = { PHOTOPRISM_DATABASE_DRIVER: "mysql", PHOTOPRISM_DATABASE_SERVER: `${DB.host}:${DB.port}`, PHOTOPRISM_DATABASE_NAME: DB.database, PHOTOPRISM_DATABASE_USER: DB.user, PHOTOPRISM_DATABASE_PASSWORD: DB.password, PHOTOPRISM_ORIGINALS_PATH: ORIGINALS_PATH, PHOTOPRISM_STORAGE_PATH: STORAGE_PATH, PHOTOPRISM_READONLY: "false", PHOTOPRISM_DISABLE_TENSORFLOW: "true", PHOTOPRISM_DISABLE_CLASSIFICATION: "true", PHOTOPRISM_DISABLE_FACES: "true", PHOTOPRISM_LOG_LEVEL: "info", } const assignments = Object.entries(env).map(([k, v]) => `${k}=${v}`) const args = ["-n", "-u", PHOTOPRISM_USER, "env", ...assignments, PHOTOPRISM_BIN, "cp", dir] return execFileAsync("sudo", args, { cwd: PHOTOPRISM_HOME, timeout: IMPORT_TIMEOUT_MS, maxBuffer: 16 * 1024 * 1024 }) } async function handleImport(req, res) { if (req.method !== "POST") return json(res, 405, { error: "method not allowed" }) const ctype = req.headers["content-type"] || "" const bm = /boundary=(?:"([^"]+)"|([^;]+))/i.exec(ctype) if (!/^multipart\/form-data/i.test(ctype) || !bm) { return json(res, 400, { error: "expected multipart/form-data" }) } const batch = randomUUID() const dir = join(STAGING_ROOT, batch) try { await mkdir(dir, { recursive: true }) } catch (err) { console.error("pathways: staging failed:", err.message) return json(res, 500, { error: "could not stage upload" }) } let files try { files = await ingestMultipart(req, (bm[1] || bm[2]).trim(), dir) } catch (err) { await rm(dir, { recursive: true, force: true }).catch(() => {}) const status = err.status === 415 ? 415 : err.status === 413 ? 413 : 400 console.error(`pathways: ingest failed (batch ${batch}):`, err.message) try { json(res, status, { error: err.message }) } catch { /* socket may be gone */ } req.destroy() return } if (files.length === 0) { await rm(dir, { recursive: true, force: true }).catch(() => {}) return json(res, 400, { error: "no image files in request" }) } try { const { stdout, stderr } = await photoprismCopy(dir) console.log(`pathways: imported ${files.length} file(s) via batch ${batch}`) if (stderr) console.error(`pathways: photoprism cp stderr: ${stderr.slice(-2000)}`) else if (stdout) console.log(`pathways: photoprism cp: ${stdout.slice(-2000)}`) } catch (err) { const detail = (err.stderr || err.stdout || err.message || "").slice(-2000) console.error(`pathways: import failed (batch ${batch} kept at ${dir}):`, detail) return json(res, 500, { error: "photoprism import failed", batch, detail }) } await refresh() await rm(dir, { recursive: true, force: true }).catch(() => {}) return json(res, 200, { ok: true, imported: files.length, batch }) } const server = http.createServer(async (req, res) => { const url = new URL(req.url, `http://${req.headers.host || "localhost"}`) const path = url.pathname try { if (path === "/api/paths") { if (!state) await refresh() return json(res, 200, { builtAt: state.builtAt, total: state.total, photos: state.photos, facets: state.facets }) } if (path === "/api/random") { if (!state) await refresh() const id = state.photoOrder[Math.floor(Math.random() * state.photoOrder.length)] return json(res, 200, { photo: state.byId.get(id) }) } const exitMatch = path.match(/^\/api\/exits\/(\d+)$/) if (exitMatch) { if (!state) await refresh() const id = Number(exitMatch[1]) const e = state.exits.get(id) if (!e) return json(res, 404, { error: "not found" }) const photos = {} for (const x of e) for (const wid of x.walk) { const ph = state.byId.get(wid) if (ph) photos[wid] = ph } return json(res, 200, { photo: state.byId.get(id), exits: e.map((x) => ({ dir: x.dir, path: x.path, next: state.byId.get(x.nextId), walk: x.walk })), photos, }) } if (path === "/api/import") return await handleImport(req, res) if (path === "/api/health") return json(res, 200, { ok: true, direct: DIRECT }) if (path === "/" || path === "") { return serveFile(res, join(PUBLIC, "index.html")) } const safe = normalize(path).replace(/^(\.\.[/\\])+/, "") const file = join(PUBLIC, safe) if (file.startsWith(PUBLIC)) return serveFile(res, file) res.writeHead(404).end("not found") } catch (err) { console.error("pathways:", err) json(res, 500, { error: "internal error" }) } }) server.requestTimeout = 0 server.headersTimeout = 65_000 server.keepAliveTimeout = 65_000 server.listen(PORT, "127.0.0.1", () => { console.log(`pathways listening on 127.0.0.1:${PORT}`) mkdir(STAGING_ROOT, { recursive: true }).catch((err) => console.error("pathways: staging dir:", err.message)) refresh() }) setInterval(refresh, 30 * 60 * 1000).unref()