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