aboutsummaryrefslogtreecommitdiff
path: root/server.mjs
diff options
context:
space:
mode:
authordax <me@dax.ist>2026-08-05 20:14:45 +0100
committerdax <me@dax.ist>2026-08-05 20:14:45 +0100
commit994cdbec9dd6ae6b965b56f4f6816a1dfef8bcc9 (patch)
tree0d7b0d1c5bec8d34b5ebe731af34ce3e8b1706ca /server.mjs
parentd2a8b2abc908311fc249454086d0d0eacfe861e1 (diff)
staging: local drop-folder (watch/) auto-import with startup recovery + failed batches; admin local-files panel; upload 413 pre-checkHEADmain
Diffstat (limited to 'server.mjs')
-rw-r--r--server.mjs184
1 files changed, 181 insertions, 3 deletions
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()