aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authordax <me@dax.ist>2026-08-04 19:49:03 +0100
committerdax <me@dax.ist>2026-08-04 19:49:03 +0100
commit815a0d1b4b71e12512cef5a8f59aa9d454d779d7 (patch)
treed2193b78b795969e4c84b5a18748a7191525648c
parent247390000a0fdf7eeedb89a2c5500c64f36df1c3 (diff)
stream import uploads to disk; advertise direct-origin fast uploads; fix /api/paths total
-rw-r--r--public/app.css13
-rw-r--r--public/app.js16
-rw-r--r--public/index.html5
-rw-r--r--server.mjs231
4 files changed, 199 insertions, 66 deletions
diff --git a/public/app.css b/public/app.css
index 153c334..a0cde93 100644
--- a/public/app.css
+++ b/public/app.css
@@ -146,6 +146,19 @@ html, body {
}
#import-state[hidden] { display: none; }
+.direct-hint {
+ margin-top: 8px;
+ font-size: 12px;
+ letter-spacing: .03em;
+ color: var(--muted);
+}
+.direct-hint a {
+ color: var(--ring);
+ text-decoration: none;
+ border-bottom: 1px solid rgba(125, 211, 252, .3);
+}
+.direct-hint a:hover { border-bottom-color: var(--ring); }
+
#file {
position: absolute;
width: 1px;
diff --git a/public/app.js b/public/app.js
index 78584bf..780bc98 100644
--- a/public/app.js
+++ b/public/app.js
@@ -661,6 +661,21 @@ function hideImportState() {
stage.classList.remove("blocked")
}
+async function initDirectHint() {
+ try {
+ const r = await fetch("/api/health")
+ const data = await r.json()
+ const d = data.direct || {}
+ const urls = [d.magic, d.lan, d.tailscale].filter(Boolean)
+ if (!urls.length) return
+ if (urls.some((u) => { try { return new URL(u).host === location.host } catch { return false } })) return
+ const link = $("#direct-link")
+ link.textContent = d.magic
+ link.href = d.magic
+ $("#direct-hint").hidden = false
+ } catch { /* ignore */ }
+}
+
async function beginWalk() {
if (started) return
started = true
@@ -695,6 +710,7 @@ async function beginWalk() {
} catch {
showImportState()
}
+ initDirectHint()
})()
let toastTimer = 0
diff --git a/public/index.html b/public/index.html
index ae56eea..736d446 100644
--- a/public/index.html
+++ b/public/index.html
@@ -5,7 +5,7 @@
<meta name="viewport" content="width=device-width, initial-scale=1">
<meta name="color-scheme" content="dark">
<title>pathways</title>
- <link rel="stylesheet" href="/app.css?v=44">
+ <link rel="stylesheet" href="/app.css?v=45">
</head>
<body>
<main id="stage">
@@ -19,6 +19,7 @@
</svg>
<span>import photos</span>
</label>
+ <p id="direct-hint" class="direct-hint" hidden>upload slow? use the direct link <a id="direct-link" href="#"></a></p>
</div>
</main>
<input id="file" type="file" accept="image/jpeg,image/png,image/webp" multiple>
@@ -28,6 +29,6 @@
<span id="progress-text"></span>
</div>
<div id="toast" role="status" aria-live="polite"></div>
- <script src="/app.js?v=44"></script>
+ <script src="/app.js?v=45"></script>
</body>
</html>
diff --git a/server.mjs b/server.mjs
index 2e1d789..e3c4c2a 100644
--- a/server.mjs
+++ b/server.mjs
@@ -1,10 +1,12 @@
import http from "node:http"
-import { readFile, mkdir, writeFile, rm } from "node:fs/promises"
-import { createReadStream, existsSync } from "node:fs"
+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)
@@ -32,6 +34,26 @@ 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)
+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"
@@ -212,48 +234,134 @@ function sniff(buf) {
return SIGNATURES.find((s) => s.test(buf)) || null
}
-function readBody(req, limit) {
- return new Promise((resolve, reject) => {
- const chunks = []
- let size = 0
- req.on("data", (c) => {
- size += c.length
- if (size > limit) {
- const err = new Error("payload too large")
- err.code = "TOO_LARGE"
- reject(err)
- req.destroy()
- return
- }
- chunks.push(c)
- })
- req.on("end", () => resolve(Buffer.concat(chunks)))
- req.on("error", reject)
- })
+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
+ }
}
-function parseMultipart(buf, boundary) {
- const parts = []
- const delim = Buffer.from(`--${boundary}`)
- const sep = Buffer.from("\r\n\r\n")
- let pos = buf.indexOf(delim)
- if (pos < 0) return parts
- pos += delim.length
- while (pos < buf.length) {
- if (buf[pos] === 0x2d && buf[pos + 1] === 0x2d) break
- if (buf[pos] === 0x0d && buf[pos + 1] === 0x0a) pos += 2
- const headEnd = buf.indexOf(sep, pos)
- if (headEnd < 0) break
- const headers = buf.subarray(pos, headEnd).toString("utf8")
- const bodyStart = headEnd + sep.length
- const next = buf.indexOf(delim, bodyStart)
- if (next < 0) break
- let bodyEnd = next
- if (buf[bodyEnd - 2] === 0x0d && buf[bodyEnd - 1] === 0x0a) bodyEnd -= 2
- parts.push({ headers, body: buf.subarray(bodyStart, bodyEnd) })
- pos = next + delim.length
+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 })
}
- return parts
+ if (state !== "done") throw Object.assign(new Error("multipart truncated"), { status: 400 })
+ return files
}
function fileNameOf(headers) {
@@ -307,36 +415,31 @@ async function handleImport(req, res) {
return json(res, 400, { error: "expected multipart/form-data" })
}
- let body
- try {
- body = await readBody(req, MAX_UPLOAD_BYTES)
- } catch (err) {
- if (err.code === "TOO_LARGE") return json(res, 413, { error: `upload exceeds ${MAX_UPLOAD_BYTES} bytes` })
- return json(res, 400, { error: "could not read upload" })
- }
-
- const taken = new Set()
- const files = []
- for (const part of parseMultipart(body, (bm[1] || bm[2]).trim())) {
- const raw = fileNameOf(part.headers)
- if (!raw || part.body.length === 0) continue
- const sig = sniff(part.body)
- if (!sig) return json(res, 415, { error: `unsupported file type: ${raw}` })
- if (files.length >= MAX_UPLOAD_FILES) return json(res, 413, { error: `too many files (max ${MAX_UPLOAD_FILES})` })
- files.push({ name: safeName(raw, sig, taken), data: part.body })
- }
- if (files.length === 0) return json(res, 400, { error: "no image files in request" })
-
const batch = randomUUID()
const dir = join(STAGING_ROOT, batch)
try {
await mkdir(dir, { recursive: true })
- for (const f of files) await writeFile(join(dir, f.name), f.data)
} 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}`)
@@ -360,7 +463,7 @@ const server = http.createServer(async (req, res) => {
try {
if (path === "/api/paths") {
if (!state) await refresh()
- return json(res, 200, { builtAt: state.builtAt, photos: state.photos, facets: state.facets })
+ return json(res, 200, { builtAt: state.builtAt, total: state.total, photos: state.photos, facets: state.facets })
}
if (path === "/api/random") {
if (!state) await refresh()
@@ -385,7 +488,7 @@ const server = http.createServer(async (req, res) => {
})
}
if (path === "/api/import") return await handleImport(req, res)
- if (path === "/api/health") return json(res, 200, { ok: true })
+ if (path === "/api/health") return json(res, 200, { ok: true, direct: DIRECT })
if (path === "/" || path === "") {
return serveFile(res, join(PUBLIC, "index.html"))