From 994cdbec9dd6ae6b965b56f4f6816a1dfef8bcc9 Mon Sep 17 00:00:00 2001 From: dax Date: Wed, 5 Aug 2026 20:14:45 +0100 Subject: staging: local drop-folder (watch/) auto-import with startup recovery + failed batches; admin local-files panel; upload 413 pre-check --- server.mjs | 184 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 181 insertions(+), 3 deletions(-) (limited to 'server.mjs') diff --git a/server.mjs b/server.mjs index cfd963c..e44e70e 100644 --- a/server.mjs +++ b/server.mjs @@ -1,5 +1,5 @@ import http from "node:http" -import { open, mkdir, rm, rename, readFile, writeFile } from "node:fs/promises" +import { open, mkdir, rm, rename, readFile, writeFile, readdir, stat } from "node:fs/promises" import { createReadStream, createWriteStream, existsSync } from "node:fs" import { extname, join, normalize } from "node:path" import { execFile } from "node:child_process" @@ -37,6 +37,17 @@ const MAX_UPLOAD_BYTES = Number(process.env.MAX_UPLOAD_BYTES || 1024 * 1024 * 10 const MAX_UPLOAD_FILES = Number(process.env.MAX_UPLOAD_FILES || 100) const IMPORT_TIMEOUT_MS = Number(process.env.IMPORT_TIMEOUT_MS || 30 * 60 * 1000) +const WATCH_ROOT = process.env.WATCH_ROOT || join(STAGING_ROOT, "watch") +const FAILED_ROOT = join(STAGING_ROOT, "failed") +const WATCH_POLL_MS = Number(process.env.WATCH_POLL_MS || 15_000) +const WATCH_SETTLE_MS = Number(process.env.WATCH_SETTLE_MS || 20_000) +const BATCH_MAX_AGE_MS = Number(process.env.BATCH_MAX_AGE_MS || 12 * 3600 * 1000) +const FAILED_MAX_AGE_MS = Number(process.env.FAILED_MAX_AGE_MS || 7 * 24 * 3600 * 1000) +const WATCH_EXTENSIONS = new Set([ + ".jpg", ".jpeg", ".png", ".webp", ".gif", ".tif", ".tiff", + ".heic", ".heif", ".dng", ".cr2", ".nef", ".arw", ".orf", ".rw2", ".raf", ".raw", +]) + function directOrigins() { const lan = [] const tailscale = [] @@ -489,9 +500,162 @@ function photoprismCopy(dir) { return execFileAsync("sudo", args, { cwd: PHOTOPRISM_HOME, timeout: IMPORT_TIMEOUT_MS, maxBuffer: 16 * 1024 * 1024 }) } +const importLog = [] + +function logImport(entry) { + importLog.unshift(entry) + if (importLog.length > 12) importLog.pop() +} + +let watchBusy = false + +async function listFiles(dir, base = "") { + const out = [] + let entries + try { entries = await readdir(dir, { withFileTypes: true }) } catch { return out } + for (const e of entries) { + if (e.name.startsWith(".")) continue + const rel = base ? `${base}/${e.name}` : e.name + if (e.isDirectory()) out.push(...await listFiles(join(dir, e.name), rel)) + else if (e.isFile()) out.push(rel) + } + return out +} + +async function watchFiles() { + const files = [] + for (const rel of await listFiles(WATCH_ROOT)) { + const st = await stat(join(WATCH_ROOT, rel)).catch(() => null) + if (st?.isFile()) files.push({ name: rel, size: st.size, mtime: st.mtimeMs }) + } + return files +} + +async function pruneEmptyDirs(dir) { + let entries = [] + try { entries = await readdir(dir, { withFileTypes: true }) } catch { return } + for (const e of entries) { + if (!e.isDirectory()) continue + const sub = join(dir, e.name) + await pruneEmptyDirs(sub) + try { + if ((await readdir(sub)).length === 0) { + await rm(sub, { recursive: true, force: true }) + console.log(`pathways: watch removed empty folder ${sub}`) + } + } catch { /* ignore */ } + } +} + +async function runWatchImport() { + if (watchBusy) return + watchBusy = true + try { + await mkdir(WATCH_ROOT, { recursive: true }) + const now = Date.now() + const staged = [] + const leftovers = [] + for (const rel of await listFiles(WATCH_ROOT)) { + const st = await stat(join(WATCH_ROOT, rel)).catch(() => null) + if (!st || !st.isFile()) continue + const ext = extname(rel).toLowerCase() + if (ext === ".tmp" || ext === ".part" || ext === ".partial") { + leftovers.push(rel) + continue + } + if (now - st.mtimeMs < WATCH_SETTLE_MS) continue + if (!WATCH_EXTENSIONS.has(ext)) continue + staged.push({ rel, src: join(WATCH_ROOT, rel), size: st.size }) + } + for (const rel of leftovers) { + await rm(join(WATCH_ROOT, rel), { force: true }).catch(() => {}) + console.log(`pathways: watch removed leftover part file ${rel}`) + } + await pruneEmptyDirs(WATCH_ROOT) + if (!staged.length) return + + const batch = randomUUID() + const dir = join(STAGING_ROOT, batch) + await mkdir(dir, { recursive: true }) + const moved = [] + for (const f of staged) { + const safe = f.rel.split("/").pop().replace(/[^A-Za-z0-9._-]/g, "_").slice(0, 120) || "upload" + let out = safe + let n = 1 + while (moved.some((m) => m.name === out)) { + const dot = safe.lastIndexOf(".") + out = `${safe.slice(0, dot)}-${n}${safe.slice(dot)}` + n++ + } + await rename(f.src, join(dir, out)).catch(() => {}) + moved.push({ name: out, size: f.size }) + } + if (!moved.length) { + await rm(dir, { recursive: true, force: true }).catch(() => {}) + return + } + + try { + const { stdout, stderr } = await photoprismCopy(dir) + if (stderr) console.error(`pathways: watch import stderr: ${stderr.slice(-2000)}`) + else if (stdout) console.log(`pathways: watch import: ${stdout.slice(-2000)}`) + await refresh() + await rm(dir, { recursive: true, force: true }).catch(() => {}) + logImport({ at: new Date().toISOString(), count: moved.length, ok: true }) + console.log(`pathways: watch imported ${moved.length} file(s)`) + } catch (err) { + const failDir = join(FAILED_ROOT, batch) + await mkdir(FAILED_ROOT, { recursive: true }) + await rename(dir, failDir).catch(async () => { await rm(dir, { recursive: true, force: true }).catch(() => {}) }) + logImport({ at: new Date().toISOString(), count: moved.length, ok: false, error: (err.stderr || err.message || "photoprism import failed").toString().slice(0, 300) }) + console.error(`pathways: watch import failed (batch ${batch}):`, (err.stderr || err.message || "").toString().slice(-2000)) + } + } catch (err) { + console.error("pathways: watch poll error:", err) + } finally { + watchBusy = false + } +} + +async function recoverStaging() { + await mkdir(WATCH_ROOT, { recursive: true }) + await mkdir(FAILED_ROOT, { recursive: true }) + const now = Date.now() + let entries = [] + try { entries = await readdir(STAGING_ROOT, { withFileTypes: true }) } catch { return } + for (const e of entries) { + if (!e.isDirectory() || e.name === "watch" || e.name === "failed") continue + const dir = join(STAGING_ROOT, e.name) + const st = await stat(dir).catch(() => null) + if (!st) continue + if (now - st.mtimeMs < BATCH_MAX_AGE_MS) { + await rename(dir, join(WATCH_ROOT, `recovered-${e.name.slice(0, 8)}`)).catch(() => {}) + console.log(`pathways: recovered interrupted batch ${e.name}`) + } else { + await rm(dir, { recursive: true, force: true }).catch(() => {}) + console.log(`pathways: purged stale batch ${e.name}`) + } + } + let failed = [] + try { failed = await readdir(FAILED_ROOT, { withFileTypes: true }) } catch { return } + for (const e of failed) { + if (!e.isDirectory()) continue + const dir = join(FAILED_ROOT, e.name) + const st = await stat(dir).catch(() => null) + if (st && now - st.mtimeMs > FAILED_MAX_AGE_MS) { + await rm(dir, { recursive: true, force: true }).catch(() => {}) + } + } +} + async function handleImport(req, res) { if (req.method !== "POST") return json(res, 405, { error: "method not allowed" }) + const contentLength = Number(req.headers["content-length"] || 0) + if (contentLength > MAX_UPLOAD_BYTES) { + return json(res, 413, { error: `upload too large (max ${MAX_UPLOAD_BYTES} bytes)` }) + } + const ctype = req.headers["content-type"] || "" const bm = /boundary=(?:"([^"]+)"|([^;]+))/i.exec(ctype) if (!/^multipart\/form-data/i.test(ctype) || !bm) { @@ -703,6 +867,15 @@ const server = http.createServer(async (req, res) => { return json(res, 200, { id, dirs }) } + if (path === "/api/admin/staging") { + if (req.method === "GET" || req.method === "POST") { + if (req.method === "POST") await runWatchImport() + const files = await watchFiles() + return json(res, 200, { watch: WATCH_ROOT, files, imports: importLog }) + } + return json(res, 405, { error: "method not allowed" }) + } + if (path === "/api/import") { const admin = await loadAdmin() if (!admin) return json(res, 401, { error: "no admin password configured; visit /admin first" }) @@ -734,7 +907,12 @@ 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() + mkdir(WATCH_ROOT, { recursive: true }) + .then(() => mkdir(FAILED_ROOT, { recursive: true })) + .then(recoverStaging) + .then(runWatchImport) + .then(() => refresh()) + .catch((err) => console.error("pathways: watch init:", err.message)) }) setInterval(refresh, 30 * 60 * 1000).unref() +setInterval(runWatchImport, WATCH_POLL_MS).unref() -- cgit v1.3.1