- 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>
2281 lines
95 KiB
Python
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()
|