Files
photon---photo-intelligence/server.py
drjones 52442d6a6d Fix pipeline shutdown hang, delete path safety, organize folder scoping, journal duplication; add PHOTON.app menu-bar launcher
- Stop/pause could block the downscale thread for up to 60s and throw an
  unhandled exception on shutdown; queue puts are now stop-aware
- /api/delete_photo(s) now refuse to delete anything outside the
  configured photo folder
- organize_alias used a stale startup-time default folder instead of the
  folder actually being scanned/tagged
- journal_write is append-only, so retagging a photo left stale duplicate
  lines that inflated category counts and search results; reads now dedupe
  by path and the journal is compacted once on startup
- batch delete now does one journal rewrite for the whole set instead of
  one per file
- add PHOTON.app: a signed, dependency-free Swift menu-bar app that starts
  server.py in the background (or attaches to an already-running instance),
  opens the console, and exposes Open/Restart/Stop/Show Log/Quit

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-07-23 03:43:41 -07:00

2281 lines
95 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 resolve_portable_path(path):
if not path:
return path
if os.path.exists(path):
return path
with S.lock:
curr_folder = S.folder
if curr_folder and os.path.isdir(curr_folder):
cand = os.path.join(curr_folder, os.path.basename(path))
if os.path.exists(cand):
return cand
return path
def read_journal_deduped():
"""Read the journal into {abs_path: latest_record}. journal_write only
ever appends (fast during a tagging run), so a photo retagged more than
once leaves stale earlier lines behind — this keeps just the newest
record per path so counts/search never double up a single photo."""
latest = {}
if not os.path.exists(JOURNAL):
return latest
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 = resolve_portable_path(rec.get("path", ""))
latest[path] = rec
except Exception:
pass
return latest
def compact_journal():
"""Rewrite the journal keeping only the latest record per path, so
on-disk duplicate lines from past retags don't linger forever."""
latest = read_journal_deduped()
if not latest:
return 0
with open(JOURNAL, "r", encoding="utf-8") as f:
line_count = sum(1 for l in f if l.strip())
removed = line_count - len(latest)
if removed > 0:
with open(JOURNAL, "w", encoding="utf-8") as f:
for rec in latest.values():
f.write(json.dumps(rec, ensure_ascii=False) + "\n")
return removed
def load_journal():
with S.lock:
S.cat_counts = {c: 0 for c in CATEGORIES}
S.done_paths = set()
latest = read_journal_deduped()
with S.lock:
for path, rec in latest.items():
S.done_paths.add(path)
cat = rec.get("category")
if cat not in S.cat_counts:
S.cat_counts[cat] = 0
S.cat_counts[cat] += 1
return len(latest)
def journal_write(rec):
with open(JOURNAL, "a", encoding="utf-8") as f:
f.write(json.dumps(rec, ensure_ascii=False) + "\n")
def trash_or_remove_file(path):
if not os.path.exists(path):
return True
try:
cmd = ["osascript", "-e", f'tell application "Finder" to delete POSIX file "{os.path.abspath(path)}"']
r = subprocess.run(cmd, capture_output=True, timeout=10)
if r.returncode == 0:
return True
except Exception:
pass
try:
os.remove(path)
return True
except Exception as e:
log("error", f"Failed to delete file {path}: {e}")
return False
def cleanup_organized_symlinks(abs_paths):
"""Remove _organized alias symlinks pointing at any of the given files.
One folder walk covers the whole batch instead of one walk per file."""
if not abs_paths:
return
with S.lock:
folder = S.folder
org_root = os.path.join(folder, "_organized")
if not os.path.exists(org_root):
return
targets = set(abs_paths)
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) in targets:
os.remove(p)
except Exception:
pass
def cleanup_thumb_cache(abs_paths):
for abs_path in abs_paths:
try:
h = hashlib.md5(abs_path.encode()).hexdigest()
for sub in ("photon_thumbs", "photon_view"):
t = os.path.join(tempfile.gettempdir(), sub, f"{h}.jpg")
if os.path.exists(t):
try:
os.remove(t)
except Exception:
pass
except Exception:
pass
def remove_journal_entries(abs_paths):
"""Strip every journal line matching any path in abs_paths in a single
read+write pass, instead of rewriting the whole journal once per file."""
if not os.path.exists(JOURNAL) or not abs_paths:
return
targets = set(abs_paths)
lines_to_keep = []
with open(JOURNAL, "r", encoding="utf-8") as f:
for line in f:
line_str = line.strip()
if not line_str:
continue
try:
rec = json.loads(line_str)
if os.path.abspath(rec.get("path", "")) not in targets:
lines_to_keep.append(line_str)
except Exception:
lines_to_keep.append(line_str)
with open(JOURNAL, "w", encoding="utf-8") as f:
for l in lines_to_keep:
f.write(l + "\n")
def remove_photo_records_batch(paths):
"""Purge journal/cache/in-memory records for a whole batch of files at
once. Used by both single and bulk delete so a duplicate-cleanup of
hundreds of files does one journal pass instead of hundreds."""
abs_paths = [os.path.abspath(p) for p in paths]
if not abs_paths:
return
cleanup_organized_symlinks(abs_paths)
cleanup_thumb_cache(abs_paths)
remove_journal_entries(abs_paths)
abs_set = set(abs_paths)
removed_failed = None
with S.lock:
S.done_paths -= abs_set
S.files = [p for p in S.files if p not in abs_set]
S.video_files = [p for p in S.video_files if p not in abs_set]
removed_failed = S.failed_paths & abs_set
if removed_failed:
S.failed_paths -= removed_failed
if removed_failed:
write_failures()
global SEARCH_CACHE
SEARCH_CACHE = {"mtime": None, "records": []}
load_journal()
push_stats()
def remove_photo_records(path):
remove_photo_records_batch([path])
def delete_photo_file_and_record(path, skip_records=False):
"""skip_records=True defers journal/cache cleanup to a single batched
call afterwards — used by bulk delete so N files cost one journal pass,
not N. The single-delete endpoint uses the default (clean up right away)."""
if not path:
return False, "path is required"
abs_path = os.path.abspath(path)
with S.lock:
folder = os.path.abspath(S.folder) if S.folder else None
try:
inside_folder = bool(folder) and os.path.commonpath([folder, abs_path]) == folder
except ValueError:
inside_folder = False
if not inside_folder:
return False, f"refusing to delete outside the configured photo folder: {abs_path}"
deleted_disk = False
if os.path.exists(abs_path):
deleted_disk = trash_or_remove_file(abs_path)
else:
deleted_disk = True
if not skip_records:
remove_photo_records(abs_path)
return deleted_disk, None
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 = []
for path, rec in read_journal_deduped().items():
try:
rec["path"] = 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 and size/filename matches."""
records = {r["path"]: r for r in load_search_index()}
cache = load_phash_cache()
groups_by_key = {}
# 1. Group by pHash if available
for path, h in cache.items():
if path in records and os.path.exists(path):
groups_by_key.setdefault(f"phash:{h}", []).append(records[path])
# 2. Also group by (file size, filename) for non-hashed or exact size matches
by_size_name = {}
for path, r in records.items():
if path not in cache and os.path.exists(path):
try:
sz = os.path.getsize(path)
name = os.path.basename(path).lower()
key = f"size:{sz}_{name}"
by_size_name.setdefault(key, []).append(r)
except Exception:
pass
for key, g in by_size_name.items():
if len(g) > 1:
groups_by_key[key] = g
groups = [g for g in groups_by_key.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))
with S.lock:
organize_folder = S.folder or DEFAULT_FOLDER
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
def should_halt():
return S.stop_evt.is_set() or stop_pipeline.is_set()
def put_until_halt(q, item):
"""Blocking put that keeps re-checking stop flags instead of hanging on
a full queue for up to 60s when the consumer has already exited."""
while not should_halt():
try:
q.put(item, timeout=0.5)
return True
except queue.Full:
continue
return False
# Stage 1: Downscaler Thread
def downscale_worker():
for i, path in enumerate(pending):
if should_halt():
break
while not S.pause_evt.is_set():
if should_halt():
break
time.sleep(0.2)
if should_halt():
break
try:
t0 = time.time()
img_bytes = downscale(path, max_px, tmpdir)
b64 = base64.b64encode(img_bytes).decode()
if not put_until_halt(downscale_queue, (i, path, img_bytes, b64, t0)):
break
except Exception as e:
log("error", f"Downscaling failed for {os.path.basename(path)} :: {e}")
if not put_until_halt(downscale_queue, (i, path, None, str(e), time.time())):
break
put_until_halt(downscale_queue, None)
# Stage 2: Inference Thread
def inference_worker():
while True:
if should_halt():
break
while not S.pause_evt.is_set():
if should_halt():
break
time.sleep(0.2)
if should_halt():
break
try:
item = downscale_queue.get(timeout=1.0)
except queue.Empty:
continue
if item is None:
put_until_halt(write_queue, None)
break
i, path, img_bytes, b64_or_err, t0 = item
name = os.path.basename(path)
if img_bytes is None:
if not put_until_halt(write_queue, (i, path, None, None, None, f"Downscale error: {b64_or_err}", t0)):
break
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
put_until_halt(write_queue, (i, path, desc, category, route, None, t0))
except Exception as e:
put_until_halt(write_queue, (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(organize_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/documents"):
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
records = load_search_index()
docs = [r for r in records if r.get("category") == "Screenshots & Documents" or r.get("route") == "screenshot" or "text:" in (r.get("desc") or "").lower()]
total = len(docs)
page = [{k: v for k, v in r.items() if k not in ("_hay", "_mtime", "_ext")} for r in docs[offset:offset + limit]]
self._json({"results": page, "total": total, "offset": offset, "limit": limit})
elif self.path.startswith("/api/convert_pdf"):
parsed_url = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed_url.query)
paths_param = params.get("paths", [""])[0] or params.get("path", [""])[0]
if not paths_param:
self._json({"error": "path parameter is required"}, 400); return
paths = [p for p in paths_param.split(",") if p and os.path.exists(p)]
if not paths:
self._json({"error": "files not found"}, 404); return
if len(paths) == 1:
src = paths[0]
base_name = os.path.splitext(os.path.basename(src))[0] + ".pdf"
tmp_pdf = tempfile.NamedTemporaryFile(suffix=".pdf", delete=False)
tmp_pdf.close()
try:
cmd = ["sips", "-s", "format", "pdf", src, "--out", tmp_pdf.name]
r = subprocess.run(cmd, capture_output=True, timeout=30)
if r.returncode != 0 or not os.path.exists(tmp_pdf.name):
raise RuntimeError(f"sips failed: {r.stderr.decode(errors='replace')}")
with open(tmp_pdf.name, "rb") as f:
pdf_data = f.read()
self.send_response(200)
self.send_header("Content-Type", "application/pdf")
self.send_header("Content-Disposition", f'attachment; filename="{base_name}"')
self.send_header("Content-Length", str(len(pdf_data)))
self.end_headers()
self.wfile.write(pdf_data)
except Exception as e:
self._json({"error": f"Failed converting image to PDF: {e}"}, 500)
finally:
try: os.remove(tmp_pdf.name)
except Exception: pass
else:
tmp_zip = tempfile.NamedTemporaryFile(suffix=".zip", delete=False)
tmp_zip.close()
tmp_dir = tempfile.mkdtemp(prefix="photon_pdf_zip_")
try:
used_names = set()
with zipfile.ZipFile(tmp_zip.name, "w", zipfile.ZIP_STORED, allowZip64=True) as zf:
for src in paths:
base = os.path.splitext(os.path.basename(src))[0]
pdf_name = base + ".pdf"
i = 1
while pdf_name in used_names:
pdf_name = f"{base}_{i}.pdf"
i += 1
used_names.add(pdf_name)
out_pdf = os.path.join(tmp_dir, pdf_name)
cmd = ["sips", "-s", "format", "pdf", src, "--out", out_pdf]
r = subprocess.run(cmd, capture_output=True, timeout=30)
if r.returncode == 0 and os.path.exists(out_pdf):
zf.write(out_pdf, arcname=pdf_name)
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_documents_pdf.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
self.wfile.write(chunk)
except Exception as e:
self._json({"error": f"Failed converting batch PDFs: {e}"}, 500)
finally:
try: os.remove(tmp_zip.name)
except Exception: pass
try: shutil.rmtree(tmp_dir)
except Exception: pass
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/choose_folder":
try:
cmd = ["osascript", "-e", 'tell application "System Events" to set frontmost of (first process whose background only is false) to true', "-e", 'POSIX path of (choose folder with prompt "Select Photo Folder to Organize:")']
r = subprocess.run(cmd, capture_output=True, timeout=60)
if r.returncode == 0:
chosen = r.stdout.decode("utf-8").strip()
if chosen and os.path.isdir(chosen):
with S.lock:
S.folder = chosen
save_folder_config(chosen)
images, videos, sidecars = scan_folder(chosen)
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"
sync_metadata_from_folder(chosen)
load_journal()
push_stats()
self._json({"ok": True, "folder": chosen, "images": len(images), "videos": len(videos), "alreadyDone": done})
return
self._json({"error": "No folder selected or dialog cancelled"}, 400)
except Exception as e:
self._json({"error": str(e)}, 500)
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/delete_photo":
path = body.get("path")
if not path:
self._json({"error": "path is required"}, 400); return
success, err = delete_photo_file_and_record(path)
if success:
log("ok", f"Deleted photo from disk: {os.path.basename(path)}")
self._json({"ok": True, "path": path})
else:
self._json({"error": f"Failed deleting file: {err}"}, 500)
elif self.path == "/api/delete_photos":
paths = body.get("paths") or []
if not paths:
self._json({"error": "paths array is required"}, 400); return
deleted_paths = []
failed = []
for p in paths:
success, err = delete_photo_file_and_record(p, skip_records=True)
if success:
deleted_paths.append(p)
else:
failed.append({"path": p, "error": err})
# One batched journal/cache pass for the whole set, instead of
# one full journal rewrite per file.
remove_photo_records_batch(deleted_paths)
log("ok", f"Batch deleted {len(deleted_paths)} photos from disk" + (f" ({len(failed)} failed)" if failed else ""))
self._json({"ok": True, "deleted": len(deleted_paths), "failed": failed})
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():
removed = compact_journal()
if removed:
print(f"journal compacted: removed {removed} duplicate/stale entries from past retags")
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()