Files
photon---photo-intelligence/server.py
drjones 5208c5ecef Add full-text photo library, bulk select/download, duplicate detection, video support, and zero-config folder detection
- Search & Edit becomes a multi-view library: Tagged/Untagged/Failed/Videos/Duplicates
- Bulk select (click, shift-click range, select-all-matching) with bulk recategorize,
  CSV/JSON export, and original-quality ZIP download (verified byte-identical)
- Video browsing/tagging via ffmpeg frame extraction, native playback with range support
- Duplicate detection via a lightweight perceptual-hash index (no new dependencies)
- Full-resolution photo/video viewer with editable metadata, prev/next navigation
- Date range, file-type, and category quick-filters; saved searches (localStorage)
- Appearance settings: accent color picker, grid density, results-per-page; full dark mode
- Auto-detects its photo folder (CLI arg > env var > last used > parent directory),
  so the app can be dropped into any photo collection and just work
2026-07-18 19:43:26 -07:00

1930 lines
80 KiB
Python

#!/usr/bin/env python3
"""
PHOTON // local photo intelligence console
- Scans a folder of photos (never deletes/moves anything)
- Sends downscaled copies to a local Ollama vision model
- Writes description + category metadata back into the photo via exiftool
- Streams live telemetry to the web GUI over SSE
Stdlib only. Requires: ollama (running), exiftool, sips (macOS built-in).
"""
import base64
import hashlib
import json
import os
import queue
import re
import shutil
import subprocess
import sys
import tempfile
import threading
import time
import urllib.request
import urllib.error
import urllib.parse
import zipfile
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
APP_DIR = os.path.dirname(os.path.abspath(__file__))
OLLAMA = "http://localhost:11434"
PORT = int(os.environ.get("PHOTON_PORT", 8765))
FOLDER_CONFIG_FILE = os.path.join(APP_DIR, "photon_folder.json")
def load_saved_folder():
if os.path.exists(FOLDER_CONFIG_FILE):
try:
with open(FOLDER_CONFIG_FILE, "r", encoding="utf-8") as f:
folder = json.load(f).get("folder")
if folder and os.path.isdir(folder):
return folder
except Exception:
pass
return None
def save_folder_config(folder):
try:
with open(FOLDER_CONFIG_FILE, "w", encoding="utf-8") as f:
json.dump({"folder": folder}, f)
except Exception:
pass
def resolve_default_folder():
"""Zero-config folder detection, in priority order:
1. `python3 server.py /path/to/photos` — explicit CLI argument
2. PHOTON_FOLDER env var
3. the last folder scanned in a previous session (remembered automatically)
4. the parent directory of wherever this app folder lives — the intended
workflow is dropping the PHOTON folder directly inside a photo
collection, so its parent IS that collection.
Always overridable from the Console tab's folder field + Scan button.
"""
if len(sys.argv) > 1:
arg = os.path.abspath(sys.argv[1])
if os.path.isdir(arg):
return arg
env = os.environ.get("PHOTON_FOLDER")
if env and os.path.isdir(env):
return os.path.abspath(env)
saved = load_saved_folder()
if saved:
return saved
return os.path.dirname(APP_DIR)
DEFAULT_FOLDER = resolve_default_folder()
JOURNAL = os.path.join(APP_DIR, "photon_journal.jsonl")
CATEGORIES_FILE = os.path.join(APP_DIR, "photon_categories.json")
IMAGE_EXTS = {".jpg", ".jpeg", ".png", ".heic", ".heif", ".tif", ".tiff", ".webp"}
VIDEO_EXTS = {".mov", ".mp4", ".m4v", ".avi", ".3gp", ".mkv"}
CATEGORIES = [
"People",
"Animals",
"Food & Drink",
"Nature & Outdoors",
"City & Buildings",
"Vehicles",
"Screenshots & Documents",
"Events & Parties",
"Objects & Stuff",
"Art & Miscellaneous",
]
SCHEMA = {
"type": "object",
"properties": {
"description": {"type": "string"},
"category": {"type": "string", "enum": CATEGORIES},
},
"required": ["description", "category"],
}
def load_categories():
global CATEGORIES, SCHEMA
if os.path.exists(CATEGORIES_FILE):
try:
with open(CATEGORIES_FILE, "r", encoding="utf-8") as f:
cats = json.load(f)
if isinstance(cats, list) and len(cats) > 0:
CATEGORIES = cats
except Exception as e:
print(f"Error loading categories, using defaults: {e}")
SCHEMA["properties"]["category"]["enum"] = CATEGORIES
def save_categories(cats):
global CATEGORIES, SCHEMA
if isinstance(cats, list) and len(cats) > 0:
CATEGORIES = cats
SCHEMA["properties"]["category"]["enum"] = CATEGORIES
try:
with open(CATEGORIES_FILE, "w", encoding="utf-8") as f:
json.dump(CATEGORIES, f, indent=2)
return True
except Exception as e:
print(f"Error saving categories: {e}")
return False
# Initialize categories
load_categories()
SCHEMA = {
"type": "object",
"properties": {
"description": {"type": "string"},
"category": {"type": "string", "enum": CATEGORIES},
},
"required": ["description", "category"],
}
KIND_SCHEMA = {
"type": "object",
"properties": {"kind": {"type": "string", "enum": ["screenshot", "photo"]}},
"required": ["kind"],
}
OCR_PROMPT = "Transcribe the readable text in this image. Output only the text itself."
ROUTER_PROMPT = (
"Classify this image. Is it a SCREENSHOT (a capture OF a phone/computer screen: "
"app, website, chat, map, or a scanned/photographed page of a document) or a "
"PHOTO (a camera photograph of the real world)? A photo of a physical object "
"that happens to have printed text or labels on it is still a photo. "
'Answer as JSON: {"kind": "screenshot"} or {"kind": "photo"}'
)
LENGTH_PRESETS = {
"brief": {"words": 12, "num_predict": 80},
"standard": {"words": 22, "num_predict": 140},
"detailed": {"words": 45, "num_predict": 260},
}
# ---------------------------------------------------------------- state
class State:
def __init__(self):
self.lock = threading.RLock()
self.status = "idle" # idle | scanning | running | paused | stopping | done
self.folder = DEFAULT_FOLDER
self.files = [] # pending image paths (after scan)
self.video_files = [] # video paths found on last scan
self.total_images = 0
self.skipped_videos = 0
self.skipped_sidecars = 0
self.dedupe_status = "idle" # idle | running
self.dedupe_progress = {"done": 0, "total": 0}
self.already_done = 0
self.processed_session = 0
self.failed_session = 0
self.current = None # dict about photo in flight
self.cat_counts = {c: 0 for c in CATEGORIES}
self.times = [] # rolling per-photo seconds
self.preview = None # (bytes, seq) last downscaled jpeg
self.preview_seq = 0
self.settings = {}
self.done_paths = set()
self.failed_paths = set() # pending retry paths
self.failures_file = os.path.join(APP_DIR, "photon_failures.jsonl")
self.log_ring = []
self.clients = [] # SSE queues
self.pause_evt = threading.Event()
self.stop_evt = threading.Event()
self.worker = None
S = State()
def load_failures():
S.failed_paths.clear()
if os.path.exists(S.failures_file):
try:
with open(S.failures_file, "r", encoding="utf-8") as f:
for line in f:
p = line.strip()
if p:
S.failed_paths.add(p)
except Exception:
pass
def write_failures():
try:
with open(S.failures_file, "w", encoding="utf-8") as f:
for p in sorted(list(S.failed_paths)):
f.write(p + "\n")
except Exception:
pass
def load_journal():
n = 0
with S.lock:
S.cat_counts = {c: 0 for c in CATEGORIES}
S.done_paths = set()
if os.path.exists(JOURNAL):
with open(JOURNAL, "r", encoding="utf-8") as f:
for line in f:
line = line.strip()
if not line:
continue
try:
rec = json.loads(line)
with S.lock:
S.done_paths.add(rec["path"])
cat = rec.get("category")
if cat not in S.cat_counts:
S.cat_counts[cat] = 0
S.cat_counts[cat] += 1
n += 1
except Exception:
pass
return n
def journal_write(rec):
with open(JOURNAL, "a", encoding="utf-8") as f:
f.write(json.dumps(rec, ensure_ascii=False) + "\n")
SEARCH_CACHE = {"mtime": None, "records": []}
def load_search_index():
"""Journal records with a precomputed lowercase haystack, cached by file mtime."""
global SEARCH_CACHE
if not os.path.exists(JOURNAL):
SEARCH_CACHE = {"mtime": None, "records": []}
return []
mtime = os.path.getmtime(JOURNAL)
if SEARCH_CACHE["mtime"] == mtime:
return SEARCH_CACHE["records"]
records = []
with open(JOURNAL, "r", encoding="utf-8") as f:
for line in f:
line = line.strip()
if not line:
continue
try:
rec = json.loads(line)
path = rec.get("path", "")
name = os.path.basename(path)
rec["name"] = name
rec["_ext"] = os.path.splitext(name)[1].lower()
rec["_hay"] = " ".join([
path.lower(), (rec.get("desc") or "").lower(),
(rec.get("category") or "").lower(), name.lower()
])
try:
st = os.stat(path)
rec["_mtime"] = st.st_mtime
rec["sizeBytes"] = st.st_size
except OSError:
rec["_mtime"] = rec.get("ts", 0)
rec["sizeBytes"] = None
records.append(rec)
except Exception:
pass
SEARCH_CACHE = {"mtime": mtime, "records": records}
return records
def parse_date_bound(s, end_of_day=False):
if not s:
return None
try:
t = time.strptime(s, "%Y-%m-%d")
epoch = time.mktime(t)
return epoch + 86399 if end_of_day else epoch
except Exception:
return None
def filtered_search_results(q, cat, date_from, date_to, sort, ftype=""):
records = load_search_index()
tokens = (q or "").strip().lower().split()
df = parse_date_bound(date_from)
dt = parse_date_bound(date_to, end_of_day=True)
def matches(rec):
if cat and rec.get("category") != cat:
return False
if ftype == "photo" and rec.get("_ext") not in IMAGE_EXTS:
return False
if ftype == "video" and rec.get("_ext") not in VIDEO_EXTS:
return False
if df is not None and rec.get("_mtime", 0) < df:
return False
if dt is not None and rec.get("_mtime", 0) > dt:
return False
if not tokens:
return True
hay = rec.get("_hay", "")
return all(t in hay for t in tokens)
filtered = [r for r in records if matches(r)]
if sort == "name":
filtered.sort(key=lambda r: r.get("name", "").lower())
elif sort == "category":
filtered.sort(key=lambda r: ((r.get("category") or ""), r.get("name", "").lower()))
else:
filtered.sort(key=lambda r: r.get("ts", 0), reverse=True)
return filtered
# ---------------------------------------------------------------- duplicate detection
PHASH_FILE = os.path.join(APP_DIR, "photon_phash.json")
PHASH_CACHE = {}
PHASH_LOCK = threading.Lock()
def load_phash_cache():
global PHASH_CACHE
with PHASH_LOCK:
if PHASH_CACHE:
return PHASH_CACHE
if os.path.exists(PHASH_FILE):
try:
with open(PHASH_FILE, "r", encoding="utf-8") as f:
PHASH_CACHE = json.load(f)
except Exception:
PHASH_CACHE = {}
return PHASH_CACHE
def save_phash_cache():
with PHASH_LOCK:
try:
with open(PHASH_FILE, "w", encoding="utf-8") as f:
json.dump(PHASH_CACHE, f)
except Exception:
pass
def compute_phash(path):
"""8x8 grayscale average-hash via sips + manual BMP parsing. No external deps."""
tmp_path = os.path.join(tempfile.gettempdir(), f"photon_ph_{os.getpid()}_{threading.get_ident()}.bmp")
try:
cmd = ["sips", "-s", "format", "bmp", "-z", "8", "8", path, "--out", tmp_path]
r = subprocess.run(cmd, capture_output=True, timeout=20)
if r.returncode != 0 or not os.path.exists(tmp_path):
return None
with open(tmp_path, "rb") as f:
data = f.read()
if data[0:2] != b"BM":
return None
pixel_offset = int.from_bytes(data[10:14], "little")
width = int.from_bytes(data[18:22], "little")
height = int.from_bytes(data[22:26], "little")
bpp = int.from_bytes(data[28:30], "little")
if bpp not in (24, 32) or width <= 0 or height <= 0:
return None
row_bytes = width * (bpp // 8)
row_padded = (row_bytes + 3) & ~3
grays = []
for y in range(height):
row_start = pixel_offset + y * row_padded
for x in range(width):
px = row_start + x * (bpp // 8)
if px + 2 >= len(data):
return None
b, g, rr = data[px], data[px + 1], data[px + 2]
grays.append((rr * 299 + g * 587 + b * 114) // 1000)
if not grays:
return None
avg = sum(grays) / len(grays)
bits = "".join("1" if p >= avg else "0" for p in grays)
return f"{int(bits, 2):0{len(bits)//4}x}"
except Exception:
return None
finally:
try:
os.remove(tmp_path)
except OSError:
pass
def build_dedupe_index():
"""Background worker: hash every tagged image not yet in the phash cache."""
with S.lock:
if S.dedupe_status == "running":
return
S.dedupe_status = "running"
try:
records = load_search_index()
paths = [r["path"] for r in records
if os.path.splitext(r["path"])[1].lower() in IMAGE_EXTS]
cache = load_phash_cache()
todo = [p for p in paths if p not in cache]
total = len(todo)
with S.lock:
S.dedupe_progress = {"done": 0, "total": total}
log("info", f"duplicate scan: hashing {total} photos …")
for i, p in enumerate(todo):
if not os.path.exists(p):
continue
h = compute_phash(p)
if h:
cache[p] = h
if i % 50 == 0:
save_phash_cache()
with S.lock:
S.dedupe_progress = {"done": i + 1, "total": total}
if i % 25 == 0:
broadcast("dedupe", {"done": i + 1, "total": total})
save_phash_cache()
broadcast("dedupe", {"done": total, "total": total, "finished": True})
log("ok", f"duplicate scan complete: {total} photos hashed")
finally:
with S.lock:
S.dedupe_status = "idle"
def find_duplicate_groups():
"""Groups of exact perceptual-hash matches — visually identical / re-saved copies."""
records = {r["path"]: r for r in load_search_index()}
cache = load_phash_cache()
by_hash = {}
for path, h in cache.items():
if path in records:
by_hash.setdefault(h, []).append(records[path])
groups = [g for g in by_hash.values() if len(g) > 1]
groups.sort(key=len, reverse=True)
return groups
# ---------------------------------------------------------------- SSE
def broadcast(kind, data):
msg = json.dumps({"type": kind, "data": data, "ts": time.time()})
with S.lock:
dead = []
for q in S.clients:
try:
q.put_nowait(msg)
except queue.Full:
dead.append(q)
for q in dead:
S.clients.remove(q)
def log(level, msg):
entry = {"level": level, "msg": msg, "ts": time.strftime("%H:%M:%S")}
with S.lock:
S.log_ring.append(entry)
if len(S.log_ring) > 400:
S.log_ring = S.log_ring[-400:]
broadcast("log", entry)
print(f"[{entry['ts']}] {level.upper():5} {msg}")
def stats_payload():
with S.lock:
avg = sum(S.times[-40:]) / len(S.times[-40:]) if S.times else 0
remaining = max(0, S.total_images - S.processed_session - S.failed_session)
return {
"status": S.status,
"total": S.total_images,
"processed": S.processed_session,
"failed": S.failed_session,
"alreadyDone": S.already_done,
"avgSec": round(avg, 2),
"perMin": round(60 / avg, 1) if avg else 0,
"etaSec": int(remaining * avg) if avg else None,
"catCounts": S.cat_counts,
"journalTotal": len(S.done_paths),
}
def push_stats():
broadcast("stats", stats_payload())
# ---------------------------------------------------------------- ollama
def ollama_get(path):
with urllib.request.urlopen(OLLAMA + path, timeout=15) as r:
return json.loads(r.read())
def ollama_models():
out = []
try:
tags = ollama_get("/api/tags").get("models", [])
except Exception as e:
return {"error": f"Ollama unreachable: {e}", "models": []}
for m in tags:
name = m["name"]
caps = []
try:
req = urllib.request.Request(
OLLAMA + "/api/show", data=json.dumps({"model": name}).encode(),
headers={"Content-Type": "application/json"})
with urllib.request.urlopen(req, timeout=15) as r:
caps = json.loads(r.read()).get("capabilities", [])
except Exception:
pass
out.append({
"name": name,
"sizeGB": round(m.get("size", 0) / 1e9, 1),
"vision": "vision" in caps,
})
return {"models": out}
def model_thinks(model):
"""True if the model has the 'thinking' capability (must be disabled,
otherwise the whole token budget is spent reasoning and the JSON is empty)."""
try:
req = urllib.request.Request(
OLLAMA + "/api/show", data=json.dumps({"model": model}).encode(),
headers={"Content-Type": "application/json"})
with urllib.request.urlopen(req, timeout=15) as r:
return "thinking" in json.loads(r.read()).get("capabilities", [])
except Exception:
return False
def ollama_generate(model, prompt, img_b64, opts, keep_alive, think=None, schema=SCHEMA):
body = {
"model": model,
"prompt": prompt,
"images": [img_b64],
"stream": False,
"options": opts,
"keep_alive": keep_alive,
}
if schema is not None:
body["format"] = schema
if think is not None:
body["think"] = think
req = urllib.request.Request(
OLLAMA + "/api/generate", data=json.dumps(body).encode(),
headers={"Content-Type": "application/json"})
with urllib.request.urlopen(req, timeout=600) as r:
return json.loads(r.read())
# ---------------------------------------------------------------- pipeline
def scan_folder(folder):
images, videos, sidecars = [], [], 0
for root, dirs, files in os.walk(folder):
dirs[:] = [d for d in dirs if not d.startswith(".") and not d.startswith("_")]
for name in sorted(files):
if name.startswith("."): # ._AppleDouble & hidden files: never touch
sidecars += 1
continue
ext = os.path.splitext(name)[1].lower()
if ext in IMAGE_EXTS:
images.append(os.path.join(root, name))
elif ext in VIDEO_EXTS:
videos.append(os.path.join(root, name))
return images, videos, sidecars
def organize_alias(folder, file_path, category):
folder = os.path.abspath(folder)
file_path = os.path.abspath(file_path)
org_root = os.path.join(folder, "_organized")
# 1. Remove existing symlinks for this file in _organized
if os.path.exists(org_root):
for root, dirs, files in os.walk(org_root):
dirs[:] = [d for d in dirs if not d.startswith(".")]
for name in files:
p = os.path.join(root, name)
if os.path.islink(p):
try:
if os.path.realpath(p) == file_path:
os.remove(p)
except Exception:
pass
# 2. Create the new symlink
cat_dir = os.path.join(org_root, category)
os.makedirs(cat_dir, exist_ok=True)
sym_path = os.path.join(cat_dir, os.path.basename(file_path))
if os.path.exists(sym_path) or os.path.islink(sym_path):
try:
os.remove(sym_path)
except Exception:
pass
try:
rel_target = os.path.relpath(file_path, cat_dir)
os.symlink(rel_target, sym_path)
except Exception:
try:
os.symlink(file_path, sym_path)
except Exception as e:
log("error", f"Failed to create symlink for {file_path}: {e}")
def verify_pixel_integrity(path, write_fn, *args, **kwargs):
h = hashlib.md5(path.encode()).hexdigest()
tmpdir = tempfile.gettempdir()
bmp_before = os.path.join(tmpdir, f"photon_int_before_{h}.bmp")
bmp_after = os.path.join(tmpdir, f"photon_int_after_{h}.bmp")
backup_file = os.path.join(tmpdir, f"photon_int_backup_{h}{os.path.splitext(path)[1]}")
shutil.copy2(path, backup_file)
try:
# Convert to BMP before
cmd_before = ["sips", "-s", "format", "bmp", path, "--out", bmp_before]
r = subprocess.run(cmd_before, capture_output=True, timeout=30)
if r.returncode != 0 or not os.path.exists(bmp_before):
raise RuntimeError(f"Pre-write BMP conversion failed: {r.stderr.decode(errors='replace')[:150]}")
with open(bmp_before, "rb") as f:
hash_before = hashlib.sha256(f.read()).hexdigest()
# Write
write_fn(*args, **kwargs)
# Convert to BMP after
cmd_after = ["sips", "-s", "format", "bmp", path, "--out", bmp_after]
r = subprocess.run(cmd_after, capture_output=True, timeout=30)
if r.returncode != 0 or not os.path.exists(bmp_after):
raise RuntimeError(f"Post-write BMP conversion failed: {r.stderr.decode(errors='replace')[:150]}")
with open(bmp_after, "rb") as f:
hash_after = hashlib.sha256(f.read()).hexdigest()
if hash_before != hash_after:
raise RuntimeError("Image pixels were modified or damaged!")
for fpath in (bmp_before, bmp_after, backup_file):
try:
os.remove(fpath)
except Exception:
pass
except Exception as e:
try:
shutil.copy2(backup_file, path)
except Exception as re:
print(f"CRITICAL: Failed to restore backup: {re}")
for fpath in (bmp_before, bmp_after, backup_file):
try:
os.remove(fpath)
except Exception:
pass
raise e
def downscale(path, max_px, tmpdir):
"""Resized jpeg copy via sips (read-only on the source). Returns bytes."""
out = os.path.join(tmpdir, "photon_frame.jpg")
cmd = ["sips", "-s", "format", "jpeg", "-s", "formatOptions", "82"]
if max_px:
cmd += ["-Z", str(max_px)]
cmd += [path, "--out", out]
r = subprocess.run(cmd, capture_output=True, timeout=120)
if r.returncode != 0 or not os.path.exists(out):
raise RuntimeError(f"sips failed: {r.stderr.decode(errors='replace')[:200]}")
with open(out, "rb") as f:
return f.read()
def extract_video_frame(path, tmpdir):
"""Grab one representative frame from a video via ffmpeg (read-only). Returns a jpeg file path."""
out = os.path.join(tmpdir, "photon_vframe.jpg")
dur = 3.0
try:
pr = subprocess.run(
["ffprobe", "-v", "error", "-show_entries", "format=duration",
"-of", "default=noprint_wrappers=1:nokey=1", path],
capture_output=True, timeout=15, text=True)
dur = float(pr.stdout.strip())
except Exception:
pass
ts = max(0.5, min(dur * 0.3, max(dur - 0.2, 0.5))) if dur > 1 else 0.1
cmd = ["ffmpeg", "-y", "-ss", str(ts), "-i", path, "-frames:v", "1", "-q:v", "3", out]
r = subprocess.run(cmd, capture_output=True, timeout=60)
if r.returncode != 0 or not os.path.exists(out):
raise RuntimeError(f"ffmpeg frame extraction failed: {r.stderr.decode(errors='replace')[:200]}")
return out
def video_thumbnail(path, out_path):
"""Cheap ffmpeg frame-grab thumbnail for the browsing grid (read-only)."""
cmd = ["ffmpeg", "-y", "-ss", "1", "-i", path, "-frames:v", "1",
"-vf", "scale=256:-1", "-q:v", "5", out_path]
r = subprocess.run(cmd, capture_output=True, timeout=30)
if r.returncode != 0 or not os.path.exists(out_path):
# very short clips: fall back to the first frame
cmd2 = ["ffmpeg", "-y", "-i", path, "-frames:v", "1",
"-vf", "scale=256:-1", "-q:v", "5", out_path]
r = subprocess.run(cmd2, capture_output=True, timeout=30)
if r.returncode != 0 or not os.path.exists(out_path):
raise RuntimeError(r.stderr.decode(errors="replace")[:200])
def salvage_json(raw):
"""Parse model output, surviving truncated/unterminated JSON."""
try:
return json.loads(raw)
except Exception:
pass
m = re.search(r"\{.*\}", raw, re.S)
if m:
try:
return json.loads(m.group(0))
except Exception:
pass
out = {}
dm = re.search(r'"description"\s*:\s*"((?:[^"\\]|\\.)*)', raw)
if dm:
# keep only complete sentences if the text was cut off mid-stream
txt = dm.group(1).replace('\\"', '"').replace("\\n", " ").strip()
cut = txt.rfind(". ")
if not txt.endswith(".") and cut > 20:
txt = txt[:cut + 1]
out["description"] = txt
cm = re.search(r'"category"\s*:\s*"((?:[^"\\]|\\.)*)"?', raw)
if cm:
out["category"] = cm.group(1).strip()
return out
def build_prompt(length_key, mode="photo"):
words = LENGTH_PRESETS[length_key]["words"]
cats = "\n".join(f"- {c}" for c in CATEGORIES)
if mode == "screenshot":
task = (f"1. This image is a screenshot or document. In ONE sentence of at most "
f"{words} words, say what app/website/document it is and what it shows "
"(the topic, not the exact words). NEVER copy the text verbatim.\n")
else:
task = (f"1. Describe this photo in ONE sentence, at most {words} words. "
"Mention the main subject, setting, and any clearly readable text.\n")
return (
"You are a photo cataloging assistant.\n" + task +
"2. Pick EXACTLY ONE category that best fits, from this list:\n"
f"{cats}\n"
"Rules: screenshots of apps/text/websites are 'Screenshots & Documents' — "
"for those, say what app/site it is and what it shows; NEVER copy the text verbatim. "
"If people are the main subject use 'People'. If unsure, use 'Art & Miscellaneous'.\n"
'Answer as JSON: {"description": "...", "category": "..."}'
)
def write_metadata(path, desc, category, keep_backup, preserve_date):
is_video = os.path.splitext(path)[1].lower() in VIDEO_EXTS
args = ["exiftool", "-m", "-q", "-codedcharacterset=utf8"]
if preserve_date:
args.append("-P")
if not keep_backup:
args.append("-overwrite_original") # writes temp file then atomic rename
if is_video:
# mp4/mov containers don't carry EXIF/IPTC — use exiftool's generic
# tag names so it resolves to QuickTime/Keys groups automatically.
args += [
f"-Description={desc}",
f"-Keywords+={category}",
"-Keywords+=photon-tagged",
]
else:
args += [
f"-EXIF:ImageDescription={desc}",
f"-IPTC:Caption-Abstract={desc}",
f"-XMP-dc:Description={desc}",
f"-XMP-dc:Subject+={category}",
f"-IPTC:Keywords+={category}",
"-XMP-dc:Subject+=photon-tagged",
]
args.append(path)
r = subprocess.run(args, capture_output=True, timeout=120)
if r.returncode != 0:
raise RuntimeError(f"exiftool: {r.stderr.decode(errors='replace')[:300]}")
def sync_metadata_from_folder(folder):
log("info", f"Syncing journal with folder metadata: {folder} ...")
cmd = ["exiftool", "-r", "-json", "-if",
"$Subject =~ /photon-tagged/ or $Keywords =~ /photon-tagged/",
"-EXIF:ImageDescription", "-XMP-dc:Description", "-IPTC:Caption-Abstract", "-Description",
"-XMP-dc:Subject", "-IPTC:Keywords", folder]
r = subprocess.run(cmd, capture_output=True, timeout=300)
if r.returncode != 0:
err_str = r.stderr.decode(errors='replace')
if "failed condition" in err_str or "No matching files" in err_str:
log("info", "Sync complete: No files found with photon-tagged metadata.")
return 0
log("error", f"Exiftool sync failed: {err_str[:200]}")
return 0
try:
data = json.loads(r.stdout)
except Exception as e:
log("error", f"Failed to parse exiftool sync JSON: {e}")
return 0
synced_count = 0
existing_recs = {}
if os.path.exists(JOURNAL):
with open(JOURNAL, "r", encoding="utf-8") as f:
for line in f:
try:
rec = json.loads(line)
existing_recs[rec["path"]] = rec
except Exception:
pass
for item in data:
path = item.get("SourceFile")
if not path:
continue
path = os.path.abspath(path)
desc = item.get("ImageDescription") or item.get("Description") or item.get("Caption-Abstract") or ""
if isinstance(desc, list) and desc:
desc = desc[0]
subj = item.get("Subject") or item.get("Keywords") or []
if isinstance(subj, str):
subj = [subj]
category = "Art & Miscellaneous"
for cat in CATEGORIES:
if cat in subj:
category = cat
break
existing = existing_recs.get(path)
if existing:
if existing["desc"] != desc or existing["category"] != category:
existing["desc"] = desc
existing["category"] = category
synced_count += 1
else:
existing_recs[path] = {
"path": path, "desc": desc, "category": category,
"model": "sync", "route": None, "sec": 0.0, "ts": time.time()
}
synced_count += 1
with open(JOURNAL, "w", encoding="utf-8") as f:
for rec in existing_recs.values():
f.write(json.dumps(rec, ensure_ascii=False) + "\n")
load_journal()
log("ok", f"Sync complete: {synced_count} entries added or updated in journal.")
return synced_count
def process_loop(settings):
model = settings["model"]
length_key = settings.get("length", "standard")
max_px = int(settings.get("maxPx", 896)) or None
temp = float(settings.get("temperature", 0.1))
keep_alive = settings.get("keepAlive", "10m")
try:
keep_alive = int(keep_alive)
except (ValueError, TypeError):
pass
dry = bool(settings.get("dryRun", False))
keep_backup = bool(settings.get("keepBackup", False))
preserve_date = bool(settings.get("preserveDate", True))
skip_done = bool(settings.get("skipDone", True))
router = bool(settings.get("router", False))
router_model = settings.get("routerModel") or "glm-ocr:latest"
shot_model = settings.get("shotModel") or model
photo_model = settings.get("photoModel") or model
ocr_text = bool(settings.get("ocrText", False)) and router
ocr_model = settings.get("ocrModel") or "glm-ocr:latest"
organize = bool(settings.get("organize", False))
integrity = bool(settings.get("integrity", False))
photo_prompt = build_prompt(length_key, "photo")
shot_prompt = build_prompt(length_key, "screenshot")
num_predict = LENGTH_PRESETS[length_key]["num_predict"]
opts = {"temperature": temp, "num_predict": num_predict}
involved = {model} if not router else {router_model, shot_model, photo_model}
if ocr_text:
involved.add(ocr_model)
think_map = {}
for m in involved:
think_map[m] = False if model_thinks(m) else None
if think_map[m] is False:
log("info", f"{m} is a thinking model — thinking disabled for speed")
tmpdir = tempfile.mkdtemp(prefix="photon_")
log("info", f"engine online :: {'router mode' if router else 'model=' + model} "
f"resize={max_px or 'off'}px len={length_key} temp={temp} dry_run={dry} integrity_check={integrity} organize_aliases={organize}")
if dry:
log("warn", "DRY RUN — no metadata will be written")
with S.lock:
pending = list(S.files)
idx_offset = 0
if skip_done:
before = len(pending)
pending = [p for p in pending if p not in S.done_paths]
idx_offset = before - len(pending)
with S.lock:
S.already_done = idx_offset
if idx_offset:
log("info", f"resume: {idx_offset} photos already in journal, skipping them")
total = len(pending) + idx_offset
if total == idx_offset:
log("info", "All photos in this folder are already tagged.")
with S.lock:
S.status = "done"
push_stats()
broadcast("state", {"status": "done"})
return
# Multi-threaded Queues
downscale_queue = queue.Queue(maxsize=2)
write_queue = queue.Queue(maxsize=2)
stop_pipeline = threading.Event()
consec_fail = 0
# Stage 1: Downscaler Thread
def downscale_worker():
for i, path in enumerate(pending):
if S.stop_evt.is_set() or stop_pipeline.is_set():
break
while not S.pause_evt.is_set():
if S.stop_evt.is_set() or stop_pipeline.is_set():
break
time.sleep(0.2)
if S.stop_evt.is_set() or stop_pipeline.is_set():
break
try:
t0 = time.time()
img_bytes = downscale(path, max_px, tmpdir)
b64 = base64.b64encode(img_bytes).decode()
downscale_queue.put((i, path, img_bytes, b64, t0), timeout=60)
except Exception as e:
log("error", f"Downscaling failed for {os.path.basename(path)} :: {e}")
downscale_queue.put((i, path, None, str(e), time.time()), timeout=60)
downscale_queue.put(None)
# Stage 2: Inference Thread
def inference_worker():
while True:
if S.stop_evt.is_set() or stop_pipeline.is_set():
break
while not S.pause_evt.is_set():
if S.stop_evt.is_set() or stop_pipeline.is_set():
break
time.sleep(0.2)
if S.stop_evt.is_set() or stop_pipeline.is_set():
break
try:
item = downscale_queue.get(timeout=1.0)
except queue.Empty:
continue
if item is None:
write_queue.put(None)
break
i, path, img_bytes, b64_or_err, t0 = item
name = os.path.basename(path)
if img_bytes is None:
write_queue.put((i, path, None, None, None, f"Downscale error: {b64_or_err}", t0))
continue
with S.lock:
S.current = {"path": path, "name": name, "idx": idx_offset + i + 1, "total": total}
S.preview_seq += 1
S.preview = (img_bytes, S.preview_seq)
broadcast("photo_start", S.current)
broadcast("preview", {"seq": S.preview_seq})
try:
route = None
use_model, use_prompt = model, photo_prompt
if router:
r = ollama_generate(router_model, ROUTER_PROMPT, b64_or_err,
{"temperature": 0, "num_predict": 30},
keep_alive, think_map.get(router_model), KIND_SCHEMA)
route = salvage_json(r.get("response", "")).get("kind")
if route not in ("screenshot", "photo"):
route = "photo"
use_model = shot_model if route == "screenshot" else photo_model
use_prompt = shot_prompt if route == "screenshot" else photo_prompt
resp = ollama_generate(use_model, use_prompt, b64_or_err,
opts, keep_alive, think_map.get(use_model))
raw = resp.get("response", "")
parsed = salvage_json(raw)
desc = (parsed.get("description") or "").strip()[:500]
category = parsed.get("category")
if category not in CATEGORIES:
category = "Art & Miscellaneous"
if not desc:
raise RuntimeError("model returned empty description")
if ocr_text and route == "screenshot":
try:
o = ollama_generate(ocr_model, OCR_PROMPT, b64_or_err,
{"temperature": 0, "num_predict": 300},
keep_alive, think_map.get(ocr_model), schema=None)
words = re.sub(r"\s+", " ", o.get("response", "")).strip()
if words:
desc = f"{desc} | text: {words[:300]}"
except Exception as e:
pass
write_queue.put((i, path, desc, category, route, None, t0))
except Exception as e:
write_queue.put((i, path, None, None, None, f"Inference error: {e}", t0))
t_downscale = threading.Thread(target=downscale_worker, daemon=True)
t_inference = threading.Thread(target=inference_worker, daemon=True)
t_downscale.start()
t_inference.start()
# Stage 3: Writer Thread (main worker thread context)
while True:
if S.stop_evt.is_set() or stop_pipeline.is_set():
break
try:
item = write_queue.get(timeout=1.0)
except queue.Empty:
continue
if item is None:
break
i, path, desc, category, route, err, t0 = item
name = os.path.basename(path)
if err:
with S.lock:
S.failed_session += 1
S.failed_paths.add(path)
write_failures()
log("error", f"[{idx_offset+i+1}/{total}] {name} FAILED :: {err}")
consec_fail += 1
if consec_fail >= 12:
with S.lock:
S.pause_evt.clear()
S.status = "paused"
log("error", f"CIRCUIT BREAKER — {consec_fail} consecutive failures, "
"auto-paused. Fix the issue and press Resume.")
broadcast("state", {"status": "paused"})
consec_fail = 0
push_stats()
continue
try:
if not dry:
if integrity:
verify_pixel_integrity(path, write_metadata, path, desc, category, keep_backup, preserve_date)
else:
write_metadata(path, desc, category, keep_backup, preserve_date)
if organize:
organize_alias(settings.get("folder") or DEFAULT_FOLDER, path, category)
dt = time.time() - t0
with S.lock:
S.processed_session += 1
S.times.append(dt)
if len(S.times) > 200:
S.times = S.times[-200:]
S.cat_counts[category] += 1
S.done_paths.add(path)
if path in S.failed_paths:
S.failed_paths.remove(path)
write_failures()
if not dry:
journal_write({"path": path, "desc": desc, "category": category,
"model": model, "route": route,
"sec": round(dt, 2), "ts": time.time()})
broadcast("photo_done", {"name": name, "path": path, "desc": desc, "category": category,
"route": route, "sec": round(dt, 2)})
tag = f" ⌁{route}" if route else ""
log("ok", f"[{idx_offset+i+1}/{total}] {name}{tag} → {category} :: {desc}")
consec_fail = 0
except Exception as e:
with S.lock:
S.failed_session += 1
S.failed_paths.add(path)
write_failures()
log("error", f"[{idx_offset+i+1}/{total}] {name} WRITE FAILED :: {e}")
consec_fail += 1
if consec_fail >= 12:
with S.lock:
S.pause_evt.clear()
S.status = "paused"
log("error", f"CIRCUIT BREAKER — {consec_fail} consecutive failures, "
"auto-paused. Fix the issue and press Resume.")
broadcast("state", {"status": "paused"})
consec_fail = 0
push_stats()
stop_pipeline.set()
t_downscale.join(timeout=2)
t_inference.join(timeout=2)
try:
shutil.rmtree(tmpdir)
except Exception:
pass
with S.lock:
S.status = "done" if not S.stop_evt.is_set() else "idle"
S.current = None
log("info", "run halted by operator" if S.stop_evt.is_set()
else f"MISSION COMPLETE — {S.processed_session} photos tagged, "
f"{S.failed_session} failed")
push_stats()
broadcast("state", {"status": S.status})
# ---------------------------------------------------------------- http
def read_body(handler):
n = int(handler.headers.get("Content-Length", 0))
return json.loads(handler.rfile.read(n)) if n else {}
class Handler(BaseHTTPRequestHandler):
def log_message(self, *a): # silence default request logging
pass
def _json(self, obj, code=200):
body = json.dumps(obj).encode()
self.send_response(code)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)
def _serve_file_with_range(self, path, ctype):
"""Basic HTTP Range support, needed for native <video> seeking."""
try:
size = os.path.getsize(path)
range_hdr = self.headers.get("Range")
start, end = 0, size - 1
status = 200
if range_hdr and range_hdr.startswith("bytes="):
status = 206
rng = range_hdr[6:].split("-")
if rng[0]:
start = int(rng[0])
if len(rng) > 1 and rng[1]:
end = int(rng[1])
end = min(end, size - 1)
length = end - start + 1
with open(path, "rb") as f:
f.seek(start)
data = f.read(length)
self.send_response(status)
self.send_header("Content-Type", ctype)
self.send_header("Accept-Ranges", "bytes")
self.send_header("Content-Length", str(len(data)))
if status == 206:
self.send_header("Content-Range", f"bytes {start}-{end}/{size}")
self.send_header("Cache-Control", "no-store")
self.end_headers()
self.wfile.write(data)
except (BrokenPipeError, ConnectionResetError):
pass
except Exception as e:
self.send_response(500); self.end_headers()
self.wfile.write(str(e).encode())
def do_GET(self):
if self.path == "/" or self.path.startswith("/index"):
with open(os.path.join(APP_DIR, "index.html"), "rb") as f:
body = f.read()
self.send_response(200)
self.send_header("Content-Type", "text/html; charset=utf-8")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)
elif self.path == "/api/models":
self._json(ollama_models())
elif self.path == "/api/state":
with S.lock:
payload = stats_payload()
payload.update({"folder": S.folder, "log": S.log_ring[-200:],
"categories": CATEGORIES, "current": S.current,
"settings": S.settings, "failuresCount": len(S.failed_paths)})
self._json(payload)
elif self.path.startswith("/api/preview"):
with S.lock:
pv = S.preview
if not pv:
self.send_response(404); self.end_headers(); return
self.send_response(200)
self.send_header("Content-Type", "image/jpeg")
self.send_header("Cache-Control", "no-store")
self.send_header("Content-Length", str(len(pv[0])))
self.end_headers()
self.wfile.write(pv[0])
elif self.path.startswith("/api/thumbnail"):
parsed_url = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed_url.query)
img_path = params.get("path", [""])[0]
if not img_path or not os.path.exists(img_path):
self.send_response(404); self.end_headers(); return
abs_path = os.path.abspath(img_path)
with S.lock:
folder_ok = abs_path.startswith(os.path.abspath(S.folder))
journal_ok = abs_path in S.done_paths
if not (folder_ok or journal_ok):
self.send_response(403); self.end_headers(); return
try:
cache_dir = os.path.join(tempfile.gettempdir(), "photon_thumbs")
os.makedirs(cache_dir, exist_ok=True)
h = hashlib.md5(abs_path.encode()).hexdigest()
thumb_path = os.path.join(cache_dir, f"{h}.jpg")
is_video = os.path.splitext(abs_path)[1].lower() in VIDEO_EXTS
if not os.path.exists(thumb_path):
if is_video:
video_thumbnail(abs_path, thumb_path)
else:
cmd = ["sips", "-s", "format", "jpeg", "-s", "formatOptions", "70", "-Z", "256", abs_path, "--out", thumb_path]
r = subprocess.run(cmd, capture_output=True, timeout=15)
if r.returncode != 0:
raise RuntimeError(f"sips failed: {r.stderr.decode(errors='replace')}")
with open(thumb_path, "rb") as f:
data = f.read()
self.send_response(200)
self.send_header("Content-Type", "image/jpeg")
self.send_header("Content-Length", str(len(data)))
self.send_header("Cache-Control", "max-age=86400")
self.end_headers()
self.wfile.write(data)
except Exception as e:
self.send_response(500); self.end_headers()
self.wfile.write(str(e).encode())
elif self.path.startswith("/api/photo"):
parsed_url = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed_url.query)
img_path = params.get("path", [""])[0]
if not img_path or not os.path.exists(img_path):
self.send_response(404); self.end_headers(); return
abs_path = os.path.abspath(img_path)
with S.lock:
folder_ok = abs_path.startswith(os.path.abspath(S.folder))
journal_ok = abs_path in S.done_paths
if not (folder_ok or journal_ok):
self.send_response(403); self.end_headers(); return
ext = os.path.splitext(abs_path)[1].lower()
serve_path = abs_path
ctype = {
".jpg": "image/jpeg", ".jpeg": "image/jpeg", ".png": "image/png",
".webp": "image/webp", ".heic": "image/heic", ".heif": "image/heif",
".tif": "image/tiff", ".tiff": "image/tiff",
".mov": "video/quicktime", ".mp4": "video/mp4", ".m4v": "video/mp4",
".avi": "video/x-msvideo", ".3gp": "video/3gpp", ".mkv": "video/x-matroska",
}.get(ext, "application/octet-stream")
if ext in VIDEO_EXTS:
self._serve_file_with_range(abs_path, ctype)
return
if ext in (".heic", ".heif", ".tif", ".tiff"):
# browsers other than Safari can't render these — serve a read-only
# full-size jpeg conversion via sips, cached, original untouched.
try:
cache_dir = os.path.join(tempfile.gettempdir(), "photon_view")
os.makedirs(cache_dir, exist_ok=True)
h = hashlib.md5(abs_path.encode()).hexdigest()
view_path = os.path.join(cache_dir, f"{h}.jpg")
if not os.path.exists(view_path) or os.path.getmtime(abs_path) > os.path.getmtime(view_path):
cmd = ["sips", "-s", "format", "jpeg", "-s", "formatOptions", "92",
"-Z", "2200", abs_path, "--out", view_path]
r = subprocess.run(cmd, capture_output=True, timeout=30)
if r.returncode != 0 or not os.path.exists(view_path):
raise RuntimeError(r.stderr.decode(errors="replace")[:200])
serve_path = view_path
ctype = "image/jpeg"
except Exception as e:
self.send_response(500); self.end_headers()
self.wfile.write(str(e).encode())
return
try:
with open(serve_path, "rb") as f:
data = f.read()
self.send_response(200)
self.send_header("Content-Type", ctype)
self.send_header("Content-Length", str(len(data)))
self.send_header("Cache-Control", "no-store")
self.end_headers()
self.wfile.write(data)
except Exception as e:
self.send_response(500); self.end_headers()
self.wfile.write(str(e).encode())
elif self.path.startswith("/api/search/paths"):
parsed_url = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed_url.query)
q = params.get("q", [""])[0]
cat = params.get("cat", [""])[0]
sort = params.get("sort", ["recent"])[0]
date_from = params.get("date_from", [""])[0]
date_to = params.get("date_to", [""])[0]
ftype = params.get("ftype", [""])[0]
filtered = filtered_search_results(q, cat, date_from, date_to, sort, ftype)
MAX_SELECT = 3000
paths = [r["path"] for r in filtered[:MAX_SELECT]]
self._json({"paths": paths, "total": len(filtered), "capped": len(filtered) > MAX_SELECT})
elif self.path.startswith("/api/search"):
parsed_url = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed_url.query)
q = params.get("q", [""])[0]
cat = params.get("cat", [""])[0]
sort = params.get("sort", ["recent"])[0]
date_from = params.get("date_from", [""])[0]
date_to = params.get("date_to", [""])[0]
ftype = params.get("ftype", [""])[0]
try:
limit = max(1, min(200, int(params.get("limit", ["60"])[0])))
except ValueError:
limit = 60
try:
offset = max(0, int(params.get("offset", ["0"])[0]))
except ValueError:
offset = 0
re_read = params.get("re_read", ["0"])[0] == "1"
if re_read:
with S.lock:
folder = S.folder
sync_metadata_from_folder(folder)
push_stats()
filtered = filtered_search_results(q, cat, date_from, date_to, sort, ftype)
total = len(filtered)
page = [{k: v for k, v in r.items() if k not in ("_hay", "_mtime", "_ext")} for r in filtered[offset:offset + limit]]
self._json({"results": page, "total": total, "offset": offset, "limit": limit})
elif self.path.startswith("/api/export"):
parsed_url = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed_url.query)
q = params.get("q", [""])[0]
cat = params.get("cat", [""])[0]
sort = params.get("sort", ["recent"])[0]
date_from = params.get("date_from", [""])[0]
date_to = params.get("date_to", [""])[0]
ftype = params.get("ftype", [""])[0]
fmt = params.get("format", ["json"])[0]
paths_only = params.get("paths", [""])[0].split(",") if params.get("paths", [""])[0] else None
filtered = filtered_search_results(q, cat, date_from, date_to, sort, ftype)
if paths_only:
wanted = set(paths_only)
filtered = [r for r in filtered if r["path"] in wanted]
if fmt == "csv":
import csv, io
buf = io.StringIO()
w = csv.writer(buf)
w.writerow(["path", "name", "category", "description", "model", "route", "tagged_at"])
for r in filtered:
w.writerow([r.get("path", ""), r.get("name", ""), r.get("category", ""),
r.get("desc", ""), r.get("model", ""), r.get("route", ""),
time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(r.get("ts", 0))) if r.get("ts") else ""])
body = buf.getvalue().encode("utf-8")
ctype = "text/csv"
fname = "photon_export.csv"
else:
page = [{k: v for k, v in r.items() if k not in ("_hay", "_mtime", "_ext")} for r in filtered]
body = json.dumps(page, ensure_ascii=False, indent=2).encode("utf-8")
ctype = "application/json"
fname = "photon_export.json"
self.send_response(200)
self.send_header("Content-Type", ctype)
self.send_header("Content-Disposition", f'attachment; filename="{fname}"')
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)
elif self.path.startswith("/api/untagged"):
parsed_url = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed_url.query)
try:
limit = max(1, min(200, int(params.get("limit", ["60"])[0])))
except ValueError:
limit = 60
try:
offset = max(0, int(params.get("offset", ["0"])[0]))
except ValueError:
offset = 0
with S.lock:
pending = [p for p in S.files if p not in S.done_paths]
total = len(pending)
page = [{"path": p, "name": os.path.basename(p)} for p in pending[offset:offset + limit]]
self._json({"results": page, "total": total, "offset": offset, "limit": limit})
elif self.path.startswith("/api/failed"):
with S.lock:
paths = sorted(S.failed_paths)
results = [{"path": p, "name": os.path.basename(p)} for p in paths]
self._json({"results": results, "total": len(results)})
elif self.path.startswith("/api/videos"):
parsed_url = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed_url.query)
try:
limit = max(1, min(200, int(params.get("limit", ["40"])[0])))
except ValueError:
limit = 40
try:
offset = max(0, int(params.get("offset", ["0"])[0]))
except ValueError:
offset = 0
with S.lock:
vids = list(S.video_files)
done = set(S.done_paths)
total = len(vids)
page = [{"path": p, "name": os.path.basename(p), "tagged": p in done}
for p in vids[offset:offset + limit]]
self._json({"results": page, "total": total, "offset": offset, "limit": limit})
elif self.path.startswith("/api/dedupe/groups"):
groups = find_duplicate_groups()
payload = [
{"hash": None, "count": len(g),
"items": [{k: v for k, v in r.items() if k not in ("_hay", "_mtime", "_ext")} for r in g]}
for g in groups
]
cache = load_phash_cache()
with S.lock:
total_images = len([p for p in S.done_paths if os.path.splitext(p)[1].lower() in IMAGE_EXTS])
self._json({"groups": payload, "hashed": len(cache), "totalImages": total_images,
"dedupeStatus": S.dedupe_status, "dedupeProgress": S.dedupe_progress})
elif self.path.startswith("/api/download_zip"):
parsed_url = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed_url.query)
paths_param = params.get("paths", [""])[0]
if paths_param:
paths = paths_param.split(",")
else:
q = params.get("q", [""])[0]
cat = params.get("cat", [""])[0]
sort = params.get("sort", ["recent"])[0]
date_from = params.get("date_from", [""])[0]
date_to = params.get("date_to", [""])[0]
ftype = params.get("ftype", [""])[0]
paths = [r["path"] for r in filtered_search_results(q, cat, date_from, date_to, sort, ftype)]
paths = [p for p in paths if p and os.path.exists(p)]
if not paths:
self._json({"error": "no files to download"}, 400); return
if len(paths) > 3000:
self._json({"error": f"too many files for one zip ({len(paths)}, max 3000) — narrow your selection"}, 400); return
tmp_zip = tempfile.NamedTemporaryFile(suffix=".zip", delete=False)
tmp_zip.close()
try:
used_names = set()
with zipfile.ZipFile(tmp_zip.name, "w", zipfile.ZIP_STORED, allowZip64=True) as zf:
for p in paths:
name = os.path.basename(p)
base, ext = os.path.splitext(name)
n, i = name, 1
while n in used_names:
n = f"{base}_{i}{ext}"; i += 1
used_names.add(n)
# original bytes, untouched — this is a read-only archive of the source files
zf.write(p, arcname=n)
size = os.path.getsize(tmp_zip.name)
self.send_response(200)
self.send_header("Content-Type", "application/zip")
self.send_header("Content-Disposition", f'attachment; filename="photon_{len(paths)}_photos.zip"')
self.send_header("Content-Length", str(size))
self.end_headers()
with open(tmp_zip.name, "rb") as f:
while True:
chunk = f.read(1024 * 1024)
if not chunk:
break
try:
self.wfile.write(chunk)
except (BrokenPipeError, ConnectionResetError):
break
finally:
try:
os.remove(tmp_zip.name)
except OSError:
pass
elif self.path == "/api/categories":
self._json({"categories": CATEGORIES})
elif self.path == "/api/audit/start":
import random
records = []
if os.path.exists(JOURNAL):
with open(JOURNAL, "r", encoding="utf-8") as f:
for line in f:
try:
records.append(json.loads(line))
except Exception:
pass
if len(records) > 100:
sampled = random.sample(records, 100)
else:
sampled = records
random.shuffle(sampled)
self._json({"results": sampled})
elif self.path == "/api/audit/history":
audit_file = os.path.join(APP_DIR, "photon_audits.jsonl")
entries = []
if os.path.exists(audit_file):
with open(audit_file, "r", encoding="utf-8") as f:
for line in f:
try:
entries.append(json.loads(line))
except Exception:
pass
passed = sum(1 for e in entries if e.get("grade") == "pass")
failed = sum(1 for e in entries if e.get("grade") == "fail")
by_day = {}
for e in entries:
day = time.strftime("%Y-%m-%d", time.localtime(e.get("ts", 0)))
d = by_day.setdefault(day, {"pass": 0, "fail": 0})
d[e.get("grade", "fail")] += 1
days = sorted(by_day.keys())[-14:]
self._json({
"total": len(entries), "passed": passed, "failed": failed,
"accuracy": round(passed / len(entries) * 100, 1) if entries else 0,
"byDay": [{"day": d, **by_day[d]} for d in days],
})
elif self.path == "/api/events":
self.send_response(200)
self.send_header("Content-Type", "text/event-stream")
self.send_header("Cache-Control", "no-cache")
self.end_headers()
q = queue.Queue(maxsize=500)
with S.lock:
S.clients.append(q)
try:
while True:
try:
msg = q.get(timeout=15)
self.wfile.write(f"data: {msg}\n\n".encode())
except queue.Empty:
self.wfile.write(b": ping\n\n")
self.wfile.flush()
except (BrokenPipeError, ConnectionResetError, OSError):
pass
finally:
with S.lock:
if q in S.clients:
S.clients.remove(q)
else:
self.send_response(404); self.end_headers()
def do_POST(self):
try:
body = read_body(self)
except Exception:
self._json({"error": "bad json"}, 400); return
if self.path == "/api/scan":
folder = body.get("folder") or DEFAULT_FOLDER
if not os.path.isdir(folder):
self._json({"error": f"not a folder: {folder}"}, 400); return
with S.lock:
if S.status == "running":
self._json({"error": "stop the run before rescanning"}, 400); return
S.status = "scanning"; S.folder = folder
save_folder_config(folder)
broadcast("state", {"status": "scanning"})
log("info", f"scanning {folder} …")
images, videos, sidecars = scan_folder(folder)
done = sum(1 for p in images if p in S.done_paths)
with S.lock:
S.files = images
S.total_images = len(images)
S.video_files = videos
S.skipped_videos = len(videos)
S.skipped_sidecars = sidecars
S.status = "idle"
log("ok", f"scan complete: {len(images)} images | {len(videos)} videos found | "
f"{sidecars} hidden/sidecar files ignored | {done} already tagged")
broadcast("state", {"status": "idle"})
push_stats()
self._json({"images": len(images), "videos": len(videos),
"sidecars": sidecars, "alreadyDone": done})
elif self.path == "/api/start":
with S.lock:
if S.status == "running":
self._json({"error": "already running"}, 400); return
if not S.files:
self._json({"error": "scan a folder first"}, 400); return
if not body.get("model"):
self._json({"error": "pick a model"}, 400); return
S.settings = body
S.status = "running"
S.processed_session = 0
S.failed_session = 0
S.times = []
S.stop_evt.clear()
S.pause_evt.set()
S.worker = threading.Thread(target=process_loop, args=(body,), daemon=True)
S.worker.start()
broadcast("state", {"status": "running"})
self._json({"ok": True})
elif self.path == "/api/pause":
with S.lock:
if S.status == "running":
S.pause_evt.clear(); S.status = "paused"
elif S.status == "paused":
S.pause_evt.set(); S.status = "running"
st = S.status
log("warn", "PAUSED — engine idling" if st == "paused" else "RESUMED")
broadcast("state", {"status": st})
self._json({"status": st})
elif self.path == "/api/stop":
with S.lock:
S.stop_evt.set(); S.pause_evt.set()
log("warn", "stop signal sent — finishing current photo")
self._json({"ok": True})
elif self.path == "/api/reveal":
path = body.get("path")
if path and os.path.exists(path):
subprocess.run(["open", "-R", path])
self._json({"ok": True})
else:
self._json({"error": "file not found"}, 400)
elif self.path == "/api/undo_all":
with S.lock:
folder = S.folder
log("warn", f"UNDO ALL: Removing all PHOTON metadata from files in {folder} ...")
args = ["exiftool", "-r", "-P", "-overwrite_original", "-if",
"$Subject =~ /photon-tagged/ or $Keywords =~ /photon-tagged/",
"-EXIF:ImageDescription=", "-IPTC:Caption-Abstract=", "-XMP-dc:Description=", "-Description=",
"-XMP-dc:Subject-=photon-tagged", "-Keywords-=photon-tagged"]
for cat in CATEGORIES:
args.append(f"-XMP-dc:Subject-={cat}")
args.append(f"-Keywords-={cat}")
args.append(folder)
r = subprocess.run(args, capture_output=True, timeout=600)
log("info", f"Exiftool undo completed: {r.stdout.decode(errors='replace')[:200]}")
org_root = os.path.join(folder, "_organized")
if os.path.exists(org_root):
try:
shutil.rmtree(org_root)
log("info", "Deleted organized symlink directory.")
except Exception as e:
log("error", f"Failed to delete _organized folder: {e}")
if os.path.exists(JOURNAL):
try:
os.remove(JOURNAL)
except Exception:
pass
if os.path.exists(S.failures_file):
try:
os.remove(S.failures_file)
except Exception:
pass
with S.lock:
S.done_paths.clear()
S.failed_paths.clear()
S.cat_counts = {c: 0 for c in CATEGORIES}
S.already_done = 0
S.processed_session = 0
S.failed_session = 0
S.times = []
S.status = "idle"
push_stats()
broadcast("state", {"status": "idle"})
log("ok", "Undo complete! All tags removed and journal reset.")
self._json({"ok": True})
elif self.path == "/api/retry_failed":
with S.lock:
if S.status == "running":
self._json({"error": "already running"}, 400); return
if not S.failed_paths:
self._json({"error": "no failed photos to retry"}, 400); return
S.files = list(S.failed_paths)
S.total_images = len(S.files)
S.processed_session = 0
S.failed_session = 0
S.times = []
S.status = "running"
S.stop_evt.clear()
S.pause_evt.set()
S.worker = threading.Thread(target=process_loop, args=(S.settings,), daemon=True)
S.worker.start()
broadcast("state", {"status": "running"})
push_stats()
self._json({"ok": True})
elif self.path == "/api/write_single":
path = body.get("path")
desc = body.get("desc", "").strip()[:500]
category = body.get("category")
if not path or not os.path.exists(path):
self._json({"error": "file not found"}, 400); return
if category not in CATEGORIES:
self._json({"error": f"invalid category: {category}"}, 400); return
try:
with S.lock:
keep_backup = bool(S.settings.get("keepBackup", False))
preserve_date = bool(S.settings.get("preserveDate", True))
organize = bool(S.settings.get("organize", False))
folder = S.folder
write_metadata(path, desc, category, keep_backup, preserve_date)
if organize:
organize_alias(folder, path, category)
existing_recs = []
found = False
if os.path.exists(JOURNAL):
with open(JOURNAL, "r", encoding="utf-8") as f:
for line in f:
try:
rec = json.loads(line)
if rec["path"] == path:
rec["desc"] = desc
rec["category"] = category
rec["ts"] = time.time()
found = True
existing_recs.append(rec)
except Exception:
pass
if not found:
existing_recs.append({
"path": path, "desc": desc, "category": category,
"model": "manual", "route": "manual", "sec": 0.0, "ts": time.time()
})
with open(JOURNAL, "w", encoding="utf-8") as f:
for rec in existing_recs:
f.write(json.dumps(rec, ensure_ascii=False) + "\n")
with S.lock:
if path in S.failed_paths:
S.failed_paths.remove(path)
write_failures()
load_journal()
push_stats()
self._json({"ok": True, "desc": desc, "category": category})
except Exception as e:
self._json({"error": f"Failed writing metadata: {e}"}, 500)
elif self.path == "/api/redo_single":
path = body.get("path")
settings = body.get("settings") or S.settings
if not path or not os.path.exists(path):
self._json({"error": "file not found"}, 400); return
try:
model = settings.get("model")
length_key = settings.get("length", "standard")
max_px = int(settings.get("maxPx", 896)) or None
temp = float(settings.get("temperature", 0.1))
keep_alive = settings.get("keepAlive", "10m")
try:
keep_alive = int(keep_alive)
except (ValueError, TypeError):
pass
keep_backup = bool(settings.get("keepBackup", False))
preserve_date = bool(settings.get("preserveDate", True))
router = bool(settings.get("router", False))
router_model = settings.get("routerModel") or "glm-ocr:latest"
shot_model = settings.get("shotModel") or model
photo_model = settings.get("photoModel") or model
ocr_text = bool(settings.get("ocrText", False)) and router
ocr_model = settings.get("ocrModel") or "glm-ocr:latest"
organize = bool(settings.get("organize", False))
photo_prompt = build_prompt(length_key, "photo")
shot_prompt = build_prompt(length_key, "screenshot")
num_predict = LENGTH_PRESETS[length_key]["num_predict"]
opts = {"temperature": temp, "num_predict": num_predict}
think = None
if model_thinks(model):
think = False
tmpdir = tempfile.mkdtemp(prefix="photon_redo_")
is_video = os.path.splitext(path)[1].lower() in VIDEO_EXTS
frame_source = extract_video_frame(path, tmpdir) if is_video else path
img = downscale(frame_source, max_px, tmpdir)
b64 = base64.b64encode(img).decode()
route = None
use_model, use_prompt = model, photo_prompt
if router:
r = ollama_generate(router_model, ROUTER_PROMPT, b64,
{"temperature": 0, "num_predict": 30},
keep_alive, think, KIND_SCHEMA)
route = salvage_json(r.get("response", "")).get("kind")
if route not in ("screenshot", "photo"):
route = "photo"
use_model = shot_model if route == "screenshot" else photo_model
use_prompt = shot_prompt if route == "screenshot" else photo_prompt
resp = ollama_generate(use_model, use_prompt, b64,
opts, keep_alive, think)
raw = resp.get("response", "")
parsed = salvage_json(raw)
desc = (parsed.get("description") or "").strip()[:500]
category = parsed.get("category")
if category not in CATEGORIES:
category = "Art & Miscellaneous"
if not desc:
raise RuntimeError("model returned empty description")
if ocr_text and route == "screenshot":
try:
o = ollama_generate(ocr_model, OCR_PROMPT, b64,
{"temperature": 0, "num_predict": 300},
keep_alive, think, schema=None)
words = re.sub(r"\s+", " ", o.get("response", "")).strip()
if words:
desc = f"{desc} | text: {words[:300]}"
except Exception as e:
pass
write_metadata(path, desc, category, keep_backup, preserve_date)
if organize:
with S.lock:
folder = S.folder
organize_alias(folder, path, category)
existing_recs = []
found = False
if os.path.exists(JOURNAL):
with open(JOURNAL, "r", encoding="utf-8") as f:
for line in f:
try:
rec = json.loads(line)
if rec["path"] == path:
rec["desc"] = desc
rec["category"] = category
rec["model"] = use_model
rec["route"] = route
rec["ts"] = time.time()
found = True
existing_recs.append(rec)
except Exception:
pass
if not found:
existing_recs.append({
"path": path, "desc": desc, "category": category,
"model": use_model, "route": route, "sec": 0.0, "ts": time.time()
})
with open(JOURNAL, "w", encoding="utf-8") as f:
for rec in existing_recs:
f.write(json.dumps(rec, ensure_ascii=False) + "\n")
try:
shutil.rmtree(tmpdir)
except Exception:
pass
with S.lock:
if path in S.failed_paths:
S.failed_paths.remove(path)
write_failures()
load_journal()
push_stats()
self._json({"ok": True, "desc": desc, "category": category, "route": route})
except Exception as e:
self._json({"error": str(e)}, 500)
elif self.path == "/api/categories":
cats = body.get("categories")
if save_categories(cats):
load_journal()
push_stats()
self._json({"ok": True, "categories": CATEGORIES})
else:
self._json({"error": "invalid categories list"}, 400)
elif self.path == "/api/audit/grade":
path = body.get("path")
grade = body.get("grade")
notes = body.get("notes", "")
if not path or grade not in ("pass", "fail"):
self._json({"error": "bad request"}, 400); return
audit_file = os.path.join(APP_DIR, "photon_audits.jsonl")
try:
with open(audit_file, "a", encoding="utf-8") as f:
f.write(json.dumps({
"path": path, "grade": grade, "notes": notes, "ts": time.time()
}, ensure_ascii=False) + "\n")
self._json({"ok": True})
except Exception as e:
self._json({"error": str(e)}, 500)
elif self.path == "/api/dedupe/build":
with S.lock:
already_running = S.dedupe_status == "running"
if already_running:
self._json({"error": "duplicate scan already running"}, 400); return
threading.Thread(target=build_dedupe_index, daemon=True).start()
self._json({"ok": True})
elif self.path == "/api/bulk_recategorize":
paths = body.get("paths") or []
category = body.get("category")
if not paths or category not in CATEGORIES:
self._json({"error": "bad request"}, 400); return
records = {r["path"]: r for r in load_search_index()}
with S.lock:
keep_backup = bool(S.settings.get("keepBackup", False))
preserve_date = bool(S.settings.get("preserveDate", True))
ok, failed = 0, []
for p in paths:
rec = records.get(p)
desc = rec.get("desc", "") if rec else ""
try:
write_metadata(p, desc, category, keep_backup, preserve_date)
ok += 1
except Exception as e:
failed.append({"path": p, "error": str(e)})
existing_recs = []
wanted = set(paths)
if os.path.exists(JOURNAL):
with open(JOURNAL, "r", encoding="utf-8") as f:
for line in f:
try:
rec = json.loads(line)
if rec["path"] in wanted and not any(fe["path"] == rec["path"] for fe in failed):
rec["category"] = category
rec["ts"] = time.time()
existing_recs.append(rec)
except Exception:
pass
with open(JOURNAL, "w", encoding="utf-8") as f:
for rec in existing_recs:
f.write(json.dumps(rec, ensure_ascii=False) + "\n")
load_journal()
push_stats()
log("ok", f"Bulk recategorized {ok} photos to '{category}'" + (f" ({len(failed)} failed)" if failed else ""))
self._json({"ok": True, "updated": ok, "failed": failed})
else:
self.send_response(404); self.end_headers()
def main():
n = load_journal()
load_failures()
load_phash_cache()
print(f"PHOTON console → http://localhost:{PORT}")
print(f"photo folder → {DEFAULT_FOLDER}"
+ (" (auto-detected — change it anytime from the Console tab)" if len(sys.argv) <= 1 and not os.environ.get("PHOTON_FOLDER") else ""))
if n:
print(f"journal loaded: {n} photos already tagged (will be skipped on resume)")
ThreadingHTTPServer(("127.0.0.1", PORT), Handler).serve_forever()
if __name__ == "__main__":
main()