Files
photon---photo-intelligence/server.py
drjones adbe4d8f24 Initial commit of PHOTON photo intelligence app
Ollama-based photo tagging with exiftool metadata writes.
2026-07-18 09:21:10 -07:00

1296 lines
52 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
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
APP_DIR = os.path.dirname(os.path.abspath(__file__))
OLLAMA = "http://localhost:11434"
PORT = 8765
DEFAULT_FOLDER = "/Volumes/sanD/allphotos from phone"
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.total_images = 0
self.skipped_videos = 0
self.skipped_sidecars = 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")
# ---------------------------------------------------------------- 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, 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 += 1
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 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):
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
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",
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/",
"-EXIF:ImageDescription", "-XMP-dc:Description", "-IPTC:Caption-Abstract",
"-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 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")
if not os.path.exists(thumb_path):
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/search"):
parsed_url = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed_url.query)
q = params.get("q", [""])[0].strip().lower()
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()
results = []
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)
path = rec.get("path", "")
desc = rec.get("desc", "")
category = rec.get("category", "")
name = os.path.basename(path)
if (not q) or (q in path.lower()) or (q in desc.lower()) or (q in category.lower()) or (q in name.lower()):
rec["name"] = name
results.append(rec)
except Exception:
pass
self._json({"results": results})
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/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
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.skipped_videos = videos
S.skipped_sidecars = sidecars
S.status = "idle"
log("ok", f"scan complete: {len(images)} images | {videos} videos skipped | "
f"{sidecars} hidden/sidecar files ignored | {done} already tagged")
broadcast("state", {"status": "idle"})
push_stats()
self._json({"images": len(images), "videos": 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/",
"-EXIF:ImageDescription=", "-IPTC:Caption-Abstract=", "-XMP-dc:Description=",
"-XMP-dc:Subject-=photon-tagged", "-IPTC:Keywords-=photon-tagged"]
for cat in CATEGORIES:
args.append(f"-XMP-dc:Subject-={cat}")
args.append(f"-IPTC: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")
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_")
img = downscale(path, 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
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)
else:
self.send_response(404); self.end_headers()
def main():
n = load_journal()
load_failures()
print(f"PHOTON console → http://localhost:{PORT}")
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()