#!/usr/bin/env python3
"""
Robust Multi-Provider & Local-LLM Proxy Server + SMB Direct File Manager + SQLite Cache
for qBittorrent AI Performer Sorter & SMB Video Library Organizer
"""
import http.server
import urllib.request
import urllib.error
import urllib.parse
import base64
import gzip
import json
import os
import re
import shutil
import socket
import sqlite3
import subprocess
import sys
import time
PORT = 8000
DEFAULT_QB_HOST = "http://192.168.6.254:8080"
DEFAULT_SMB_MOUNT = "/mnt/isolation"
DEFAULT_SMB_SOURCE = "/mnt/isolation/videos/videodownloader"
DEFAULT_SMB_TARGET_ROOT = "/mnt/isolation/videos"
DB_PATH = "/home/david/performer_sorter.db"
# Video extensions filter
VIDEO_EXTENSIONS = {'.mp4', '.mkv', '.avi', '.mov', '.wmv', '.flv', '.webm', '.m4v', '.ts', '.iso', '.m2ts', '.mpg', '.mpeg'}
# Global session tracking
latest_sid = None
def init_db():
conn = sqlite3.connect(DB_PATH)
cursor = conn.cursor()
cursor.execute('''
CREATE TABLE IF NOT EXISTS performer_analysis_cache (
id INTEGER PRIMARY KEY AUTOINCREMENT,
item_type TEXT NOT NULL,
item_key TEXT UNIQUE NOT NULL,
name TEXT NOT NULL,
file_size INTEGER DEFAULT 0,
file_mtime REAL DEFAULT 0,
primary_performer TEXT,
all_performers TEXT,
selected_performer TEXT,
target_path TEXT,
confidence TEXT,
reasoning TEXT,
move_status TEXT DEFAULT 'ready',
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
)
''')
cursor.execute('CREATE INDEX IF NOT EXISTS idx_item_key ON performer_analysis_cache(item_key)')
cursor.execute('CREATE INDEX IF NOT EXISTS idx_item_type ON performer_analysis_cache(item_type)')
cursor.execute('''
CREATE TABLE IF NOT EXISTS relocation_history (
id INTEGER PRIMARY KEY AUTOINCREMENT,
source_type TEXT NOT NULL, -- 'qb' or 'smb'
item_key TEXT NOT NULL, -- torrent hash or file path
item_name TEXT NOT NULL,
source_path TEXT NOT NULL,
target_path TEXT NOT NULL,
performer TEXT,
file_size INTEGER DEFAULT 0,
status TEXT DEFAULT 'completed', -- 'completed', 'undone', 'failed'
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
undone_at TIMESTAMP
)
''')
cursor.execute('CREATE INDEX IF NOT EXISTS idx_history_key ON relocation_history(item_key)')
cursor.execute('CREATE INDEX IF NOT EXISTS idx_history_time ON relocation_history(created_at)')
conn.commit()
conn.close()
init_db()
def record_relocation_history(source_type, item_key, item_name, source_path, target_path, performer="", file_size=0):
try:
conn = get_db_connection()
cursor = conn.cursor()
cursor.execute('''
INSERT INTO relocation_history (source_type, item_key, item_name, source_path, target_path, performer, file_size, status)
VALUES (?, ?, ?, ?, ?, ?, ?, 'completed')
''', (source_type, item_key, item_name, source_path, target_path, performer, file_size))
conn.commit()
history_id = cursor.lastrowid
conn.close()
return history_id
except Exception as e:
print(f"Error recording relocation history: {e}")
return None
def get_db_connection():
conn = sqlite3.connect(DB_PATH)
conn.row_factory = sqlite3.Row
return conn
# Stash Server Integration
DEFAULT_STASH_URL = "http://servervm.local:9999/graphql"
DEFAULT_STASH_SQLITE = "/mnt/stash_appdata/stashapp/config/stash-go.sqlite"
def normalize_stash_url(url_str):
if not url_str:
return DEFAULT_STASH_URL
url_str = url_str.strip().rstrip('/')
if not url_str.startswith('http://') and not url_str.startswith('https://'):
url_str = 'http://' + url_str
if not url_str.endswith('/graphql'):
url_str = url_str + '/graphql'
if 'servervm.local' in url_str:
try:
ip = socket.gethostbyname('servervm.local')
url_str = url_str.replace('servervm.local', ip)
except Exception:
pass
return url_str
# In-memory performer image cache
PERFORMER_IMAGE_CACHE = {}
def get_fallback_avatar_svg(name):
clean = (name or '?').replace('.', ' ').strip().title()
initials = ''.join([part[0].upper() for part in clean.split() if part][:2]) or '?'
svg = f''''''
return svg.encode('utf-8')
def fetch_performer_image_data(performer_name=None, performer_id=None, stash_url=None):
cache_key = str(performer_id) if performer_id else str(performer_name).lower()
if cache_key in PERFORMER_IMAGE_CACHE:
return PERFORMER_IMAGE_CACHE[cache_key]
img_url = None
if performer_id:
base = normalize_stash_url(stash_url).replace('/graphql', '')
img_url = f"{base}/performer/{performer_id}/image"
elif performer_name:
clean_name = performer_name.replace('.', ' ')
q = f'''
{{
findPerformers(filter: {{q: {json.dumps(clean_name)}, per_page: 1}}) {{
performers {{
id
name
image_path
}}
}}
}}
'''
try:
res = execute_stash_graphql(q, stash_url, timeout=3)
performers = res.get('data', {}).get('findPerformers', {}).get('performers', [])
if performers and performers[0].get('image_path'):
img_url = performers[0]['image_path']
except Exception:
pass
if img_url:
try:
req = urllib.request.Request(img_url)
with urllib.request.urlopen(req, timeout=4) as res:
content_type = res.headers.get('Content-Type', 'image/jpeg')
img_bytes = res.read()
result = (content_type, img_bytes)
PERFORMER_IMAGE_CACHE[cache_key] = result
return result
except Exception as e:
print(f"Error downloading performer image {img_url}: {e}")
# Fallback to SVG avatar
fallback_data = ('image/svg+xml', get_fallback_avatar_svg(performer_name or 'Performer'))
PERFORMER_IMAGE_CACHE[cache_key] = fallback_data
return fallback_data
def format_performer_slug(raw_name):
if not raw_name:
return ''
clean = re.sub(r'[/\\?%*:|"<>]+', '', raw_name.strip()).lower()
parts = [p for p in re.split(r'[\s_\-\.]+', clean) if p]
return '.'.join(parts)
def execute_stash_graphql(query_str, stash_url=None, variables=None, timeout=4):
endpoint = normalize_stash_url(stash_url)
payload = {'query': query_str}
if variables:
payload['variables'] = variables
req = urllib.request.Request(
endpoint,
data=json.dumps(payload).encode('utf-8'),
headers={'Content-Type': 'application/json'}
)
with urllib.request.urlopen(req, timeout=timeout) as res:
return json.loads(res.read().decode())
def build_stash_mutation(mutation_name, input_obj=None, fields=None, args_dict=None):
call_parts = []
if input_obj is not None:
raw_json = json.dumps(input_obj)
gql_input = re.sub(r'\"([a-zA-Z_][a-zA-Z0-9_]*)\"\s*:', r'\1:', raw_json)
call_parts.append(f"input: {gql_input}")
if args_dict:
for k, v in args_dict.items():
if isinstance(v, bool):
call_parts.append(f"{k}: {'true' if v else 'false'}")
elif isinstance(v, (int, float)):
call_parts.append(f"{k}: {v}")
else:
v_str = str(v).replace('"', '\\"')
call_parts.append(f'{k}: "{v_str}"')
call_str = f"{mutation_name}({', '.join(call_parts)})" if call_parts else mutation_name
if fields:
return f"mutation {{ {call_str} {{ {fields} }} }}"
return f"mutation {{ {call_str} }}"
def parse_stash_scene_data(scene, match_strategy='exact_path'):
if not scene:
return None
performers = scene.get('performers', [])
if not performers:
return None
females = [p for p in performers if p.get('gender') == 'FEMALE']
non_males = [p for p in performers if p.get('gender') not in ('FEMALE', 'MALE')]
# Priority: NEVER include male performers if any female or non-male performers are present
eligible_performers = females if females else (non_males if non_males else performers)
formatted_candidates = []
seen = set()
for p in eligible_performers:
name = p.get('name', '')
fmt = format_performer_slug(name)
if fmt and fmt not in seen:
seen.add(fmt)
formatted_candidates.append({
'raw': name,
'formatted': fmt,
'id': p.get('id'),
'gender': p.get('gender'),
'image': p.get('image_path')
})
if not formatted_candidates:
return None
primary_formatted = formatted_candidates[0]['formatted']
studio_name = scene.get('studio', {}).get('name') if scene.get('studio') else None
title = scene.get('title') or ''
scene_id = scene.get('id') or ''
reasoning = f'Stash Verified: Scene #{scene_id}'
if title:
reasoning += f' "{title}"'
if studio_name:
reasoning += f' ({studio_name})'
return {
'stash_id': scene_id,
'title': title,
'studio': studio_name,
'primaryPerformer': primary_formatted,
'allPerformers': formatted_candidates,
'femalePerformers': [p['formatted'] for p in formatted_candidates if p.get('gender') == 'FEMALE'],
'malePerformers': [],
'confidence': 'verified',
'reasoning': reasoning,
'matchStrategy': match_strategy
}
def clean_torrent_title_for_stash(raw_title):
t = re.sub(r'\.(mp4|mkv|avi|mov|wmv|ts|flv)$', '', raw_title or '', flags=re.IGNORECASE)
t = re.sub(r'\[.*?\]|\(.*?\)', ' ', t)
t = re.sub(r'\b(1080p|720p|2160p|4k|hd|hevc|x265|x264|rarbg|yify|vixen|brazzers|puretaboo|brattysis|team.skeet)\b', ' ', t, flags=re.IGNORECASE)
t = re.sub(r'[\.\-_]+', ' ', t)
return ' '.join(t.split())
def query_stash_single(file_name, file_path=None, item_type='smb', stash_url=None):
raw_name = file_name or os.path.basename(file_path or '')
basename = os.path.basename(file_path or file_name or '')
stash_paths = [
f'/data/videodownloader/{basename}',
f'/data/downloads/{basename}',
f'/data/{basename}'
]
# 1. Exact path match in Stash
for p in stash_paths:
q = f'''
{{
findScenes(scene_filter: {{path: {{value: {json.dumps(p)}, modifier: EQUALS}}}}) {{
scenes {{
id
title
performers {{
id
name
gender
image_path
}}
studio {{
name
}}
}}
}}
}}
'''
try:
res = execute_stash_graphql(q, stash_url, timeout=3)
scenes = res.get('data', {}).get('findScenes', {}).get('scenes', [])
if scenes and scenes[0].get('performers'):
parsed = parse_stash_scene_data(scenes[0], 'exact_path')
if parsed:
return parsed
except Exception:
pass
# 2. Cleaned Title / Keyword search in Stash (especially effective for torrent names)
cleaned_title = clean_torrent_title_for_stash(raw_name)
search_queries = [cleaned_title[:40], os.path.splitext(basename)[0][:35]]
for sq in search_queries:
if not sq or len(sq) < 3:
continue
q = f'''
{{
findScenes(filter: {{q: {json.dumps(sq)}, per_page: 3}}) {{
scenes {{
id
title
files {{
basename
}}
performers {{
id
name
gender
image_path
}}
studio {{
name
}}
}}
}}
}}
'''
try:
res = execute_stash_graphql(q, stash_url, timeout=3)
scenes = res.get('data', {}).get('findScenes', {}).get('scenes', [])
if scenes:
for s in scenes:
if s.get('performers'):
parsed = parse_stash_scene_data(s, 'title_match')
if parsed:
return parsed
except Exception:
pass
return None
def query_stash_batch(items, stash_url=None):
if not items:
return []
results = [None] * len(items)
unmatched_indices = []
# 1. First pass: Batch exact path queries in chunks of 30
chunk_size = 30
for c_start in range(0, len(items), chunk_size):
c_items = items[c_start:c_start + chunk_size]
queries = []
for i_idx, item in enumerate(c_items):
file_path = item.get('path') or item.get('name') or ''
basename = os.path.basename(file_path)
p = f'/data/videodownloader/{basename}'
queries.append(f'item_{i_idx}: findScenes(scene_filter: {{path: {{value: {json.dumps(p)}, modifier: EQUALS}}}}) {{ scenes {{ id title performers {{ id name gender image_path }} studio {{ name }} }} }}')
batch_q = '{\n' + '\n'.join(queries) + '\n}'
try:
res = execute_stash_graphql(batch_q, stash_url, timeout=4)
data = res.get('data', {})
for i_idx in range(len(c_items)):
global_idx = c_start + i_idx
scenes = data.get(f'item_{i_idx}', {}).get('scenes', [])
if scenes and scenes[0].get('performers'):
parsed = parse_stash_scene_data(scenes[0], 'exact_path')
if parsed:
results[global_idx] = parsed
if not results[global_idx]:
unmatched_indices.append(global_idx)
except Exception as e:
print(f"Batch exact Stash query failed: {e}")
for i_idx in range(len(c_items)):
unmatched_indices.append(c_start + i_idx)
# 2. Second pass: Batch cleaned title queries for all unmatched items in chunks of 30
if unmatched_indices:
for c_start in range(0, len(unmatched_indices), chunk_size):
c_indices = unmatched_indices[c_start:c_start + chunk_size]
queries = []
for sub_idx, g_idx in enumerate(c_indices):
item = items[g_idx]
raw_name = item.get('name') or os.path.basename(item.get('path') or '')
cleaned_title = clean_torrent_title_for_stash(raw_name)
queries.append(f'title_{sub_idx}: findScenes(filter: {{q: {json.dumps(cleaned_title[:35])}, per_page: 2}}) {{ scenes {{ id title performers {{ id name gender image_path }} studio {{ name }} }} }}')
batch_q = '{\n' + '\n'.join(queries) + '\n}'
try:
res = execute_stash_graphql(batch_q, stash_url, timeout=4)
data = res.get('data', {})
for sub_idx, g_idx in enumerate(c_indices):
scenes = data.get(f'title_{sub_idx}', {}).get('scenes', [])
if scenes:
for s in scenes:
if s.get('performers'):
parsed = parse_stash_scene_data(s, 'title_match')
if parsed:
results[g_idx] = parsed
break
except Exception as e:
print(f"Batch title Stash query failed: {e}")
return results
def scrape_babepedia_images(name, limit=4):
slug = name.strip().replace(' ', '_')
url = f"https://www.babepedia.com/babe/{urllib.parse.quote(slug)}"
req = urllib.request.Request(
url,
headers={
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/122.0.0.0 Safari/537.36',
'Accept-Language': 'en-US,en;q=0.9',
'Accept': 'text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8'
}
)
results = []
seen = set()
try:
with urllib.request.urlopen(req, timeout=6) as res:
html = res.read().decode('utf-8', errors='ignore')
m = re.search(r'id=\"profimg\"[^\>]*href=\"([^\"]+)\"', html)
if not m:
m = re.search(r'
]+src=[\"\'](/babeimg/[^\"\']+)[\"\']', html)
if m:
u = m.group(1)
if not u.startswith('http'):
u = 'https://www.babepedia.com' + u
seen.add(u)
results.append({'url': u, 'thumbnail': u, 'source': 'Babepedia'})
gallery = re.findall(r'href=[\"\'](/pics/[^\"\']+)[\"\']', html)
for g in gallery:
if not g.startswith('http'):
g = 'https://www.babepedia.com' + g
if g not in seen:
seen.add(g)
results.append({'url': g, 'thumbnail': g, 'source': 'Babepedia'})
if len(results) >= limit:
break
except Exception:
pass
return results
def clean_performer_search_name(raw_name):
cleaned = re.sub(r'\s*[-–—]\s*.*$', '', raw_name or '')
cleaned = re.sub(r'\(.*?\)', '', cleaned)
cleaned = re.sub(r'[\(\)\[\]\{\}]', '', cleaned)
cleaned = re.sub(r'[\s\.\-_]+', ' ', cleaned).strip()
return cleaned
def is_valid_adult_candidate(img_url, title, desc, performer_name):
clean_name = clean_performer_search_name(performer_name).lower()
name_tokens = [tok for tok in clean_name.split() if len(tok) > 2]
if not name_tokens:
return False
url_lower = (img_url or '').lower()
text_lower = f"{img_url} {title} {desc}".lower()
# 1. Blacklist non-person, ecommerce, stock photo, and unrelated domains
blacklist_domains = [
'pxhere.com', 'pixabay.com', 'freepik.com', 'alamy.com', 'wikipedia.org',
'wikimedia.org', 'infoescola.com', 'ufrgs.br', 'todojujuy.com', 'gettyimages.com',
'shutterstock.com', 'tripadvisor.com', 'etsy.com', 'etsystatic.com', 'pinterest.com',
'pinimg.com', 'ebay.com', 'amazon.com', 'indiamike.com', 'toolstop.co.uk',
'homedepot.com', 'youtube.com', 'disney', 'marvel', 'cartoon', 'anime', 'uhrcenter.de',
'swimxwin.com', 'cdiscount.com', 'dailymail.co.uk', 'dreamstime.com'
]
if any(b in url_lower for b in blacklist_domains):
return False
# 2. Blacklist non-human junk keywords
junk_keywords = [
'sapo', 'frog', 'toad', 'shrine', 'temple', 'power-tool', 'drill', 'tractor',
'boat', 'yacht', 'car-parts', 'engine', 'railroad', 'railway', 'hardware',
'wildlife', 'amphibian', 'reptile', 'jewelry', 'schmuck', 'panties', 'costume',
'tonsils', 'throat', 'golf-course', 'clubhouse', 'drone'
]
if any(j in text_lower for j in junk_keywords):
return False
# 3. Verified adult biography / media domains
adult_domains = [
'babepedia.com', 'freeones.com', 'iafd.com', 'boobpedia.com',
'thenewsgod.com', 'famousbio.wiki', 'globalzonetoday.com', 'adultdvdempire.com',
'stashdb.org', 'theporndb.net', 'wikistarbio.com', 'babecelebs.com',
'pornstarbio.com', 'babe.today', 'indexxx.com', 'scrolller.com'
]
is_adult_domain = any(ad in url_lower for ad in adult_domains)
# 4. Strict name token matching:
# If from an adult domain, at least 1 significant token must match.
# If from a general web domain, ALL tokens (or the full name phrase) must match.
if is_adult_domain:
return any(tok in text_lower for tok in name_tokens)
else:
# Full name phrase or all tokens must match
if clean_name in text_lower:
return True
return all(tok in text_lower for tok in name_tokens)
def search_performer_candidate_images(name, aliases=None, limit=6):
clean_name = clean_performer_search_name(name)
if not clean_name or len(clean_name) < 3:
return []
results = []
seen = set()
# 1. Primary: Direct Babepedia Adult Database Scrape
babe_imgs = scrape_babepedia_images(clean_name, limit=4)
for it in babe_imgs:
if it['url'] not in seen:
seen.add(it['url'])
results.append(it)
if len(results) >= limit:
return results
# 2. Secondary: Adult Database Targeted Queries on Bing
queries = [
f'"{clean_name}" site:babepedia.com',
f'"{clean_name}" site:freeones.com',
f'"{clean_name}" site:iafd.com',
f'"{clean_name}" site:boobpedia.com',
f'"{clean_name}" adult star portrait',
f'"{clean_name}" adult model photoshoot'
]
if aliases and isinstance(aliases, list):
for a in aliases[:2]:
clean_a = clean_performer_search_name(a)
if clean_a and clean_a.lower() != clean_name.lower():
queries.append(f'"{clean_a}" site:babepedia.com')
queries.append(f'"{clean_a}" adult star portrait')
for q in queries:
url = f"https://www.bing.com/images/search?q={urllib.parse.quote(q)}&FORM=HDRSC2"
req = urllib.request.Request(
url,
headers={
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/122.0.0.0 Safari/537.36',
'Accept-Language': 'en-US,en;q=0.9',
'Accept': 'text/html,application/xhtml+xml,application/xml;q=0.9,image/webp,*/*;q=0.8'
}
)
try:
with urllib.request.urlopen(req, timeout=6) as res:
page_html = res.read().decode('utf-8', errors='ignore')
raw_blocks = re.findall(r'class=\"iusc\"[^\>]*m=\"([^\"]+)\"', page_html)
for raw_b in raw_blocks:
unesc = urllib.parse.unquote(raw_b).replace('"', '"')
murl_m = re.search(r'\"murl\":\"([^\"]+)\"', unesc)
turl_m = re.search(r'\"turl\":\"([^\"]+)\"', unesc)
title_m = re.search(r'\"t\":\"([^\"]+)\"', unesc)
desc_m = re.search(r'\"desc\":\"([^\"]+)\"', unesc)
if murl_m:
clean_u = murl_m.group(1).replace(r'\/', '/')
t = turl_m.group(1).replace(r'\/', '/') if turl_m else clean_u
title = title_m.group(1) if title_m else ''
desc = desc_m.group(1) if desc_m else ''
if clean_u not in seen and is_valid_adult_candidate(clean_u, title, desc, clean_name):
seen.add(clean_u)
results.append({'url': clean_u, 'thumbnail': t, 'source': 'Adult Web'})
if len(results) >= limit:
break
except Exception:
pass
if len(results) >= limit:
break
return results
def fetch_image_as_base64(image_url):
req = urllib.request.Request(
image_url,
headers={
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/122.0.0.0 Safari/537.36',
'Referer': image_url
}
)
with urllib.request.urlopen(req, timeout=12) as resp:
content_type = resp.headers.get('Content-Type', 'image/jpeg')
if not content_type or not content_type.startswith('image/'):
content_type = 'image/jpeg'
raw_data = resp.read()
b64_str = base64.b64encode(raw_data).decode('utf-8')
return f"data:{content_type};base64,{b64_str}"
def normalize_host(host_str):
if not host_str:
return DEFAULT_QB_HOST
host_str = host_str.strip().rstrip('/')
if not host_str.startswith('http://') and not host_str.startswith('https://'):
host_str = 'http://' + host_str
if 'truenas.local' in host_str:
try:
ip = socket.gethostbyname('truenas.local')
host_str = host_str.replace('truenas.local', ip)
except Exception:
pass
return host_str
def normalize_url_base(url_str):
if not url_str:
return 'http://localhost:1234/v1'
url_str = url_str.strip().rstrip('/')
if not url_str.startswith('http://') and not url_str.startswith('https://'):
url_str = 'http://' + url_str
# Try resolving mDNS hostnames like odysseus.local or truenas.local
for mdns in ['odysseus.local', 'truenas.local']:
if mdns in url_str:
try:
ip = socket.gethostbyname(mdns)
url_str = url_str.replace(mdns, ip)
except Exception:
pass
return url_str
def ensure_smb_mounted():
if os.path.exists(DEFAULT_SMB_SOURCE):
return True, "Mounted and accessible"
mount_script = "/home/david/mount_isolation.sh"
if os.path.exists(mount_script):
try:
subprocess.run(["bash", mount_script], capture_output=True, text=True, timeout=15)
if os.path.exists(DEFAULT_SMB_SOURCE):
return True, "Successfully mounted via mount_isolation.sh"
except Exception as e:
return False, f"Auto-mount error: {str(e)}"
return os.path.exists(DEFAULT_SMB_MOUNT), "Mount status checked"
# ==================== Video Streaming Helper (Range 206) ====================
def resolve_stream_video_path(requested_path):
if not requested_path:
return None
ensure_smb_mounted()
# 1. Direct path check
if os.path.isfile(requested_path):
return requested_path
basename = os.path.basename(requested_path)
stem = os.path.splitext(basename)[0]
candidate_paths = [
requested_path.replace('/data/', '/mnt/isolation/videos/'),
requested_path.replace('/data', '/mnt/isolation/videos'),
requested_path.replace('/media/stashapp/generated/', '/mnt/isolation/stashapp/generated/'),
requested_path.replace('/generated/', '/mnt/isolation/stashapp/generated/'),
requested_path.replace('/isolation/videos/', '/mnt/isolation/videos/'),
requested_path.replace('/isolation/', '/mnt/isolation/'),
requested_path.replace('/downloads/', '/mnt/isolation/videos/unsorted/'),
requested_path.replace('/downloads/', '/mnt/isolation/videos/videodownloader/'),
requested_path.replace('/downloads/', '/mnt/isolation/videos/'),
requested_path.replace('/downloads/', '/mnt/isolation/incomplete/'),
requested_path.replace('/downloads/', '/mnt/isolation/downloads/'),
os.path.join('/mnt/isolation/videos/unsorted', basename),
os.path.join('/mnt/isolation/videos/videodownloader', basename),
os.path.join('/mnt/isolation/stashapp/generated/transcodes', basename),
os.path.join('/mnt/isolation/incomplete', basename),
os.path.join('/mnt/isolation/downloads', basename)
]
video_exts = ('.mp4', '.mkv', '.avi', '.mov', '.webm', '.ts', '.m4v')
for p in candidate_paths:
if os.path.isfile(p):
return p
elif os.path.isdir(p):
try:
best_file = None
best_size = -1
for root, _, files in os.walk(p):
for f in files:
if f.lower().endswith(video_exts):
full = os.path.join(root, f)
sz = os.path.getsize(full)
if sz > best_size:
best_size = sz
best_file = full
if best_file:
return best_file
except Exception:
pass
# 2. Check top-level directories
try:
search_dirs = [
'/mnt/isolation/videos/unsorted',
'/mnt/isolation/videos/videodownloader',
'/mnt/isolation/incomplete',
'/mnt/isolation/downloads'
]
for s_dir in search_dirs:
if os.path.exists(s_dir) and os.path.isdir(s_dir):
target = os.path.join(s_dir, basename)
if os.path.isfile(target):
return target
for entry in os.scandir(s_dir):
if entry.is_file() and (entry.name == basename or os.path.splitext(entry.name)[0] == stem):
return entry.path
elif entry.is_dir():
sub_target = os.path.join(entry.path, basename)
if os.path.isfile(sub_target):
return sub_target
except Exception as e:
print(f"Error scanning directories in streaming resolver: {e}")
# 3. Fallback search inside performer folders in /mnt/isolation/videos
try:
if os.path.exists('/mnt/isolation/videos'):
for p_entry in os.scandir('/mnt/isolation/videos'):
if p_entry.is_dir():
candidate = os.path.join(p_entry.path, basename)
if os.path.isfile(candidate):
return candidate
except Exception:
pass
return None
class RobustProxyHandler(http.server.SimpleHTTPRequestHandler):
def send_cors_headers(self):
origin = self.headers.get('Origin')
if origin:
self.send_header('Access-Control-Allow-Origin', origin)
self.send_header('Access-Control-Allow-Credentials', 'true')
else:
self.send_header('Access-Control-Allow-Origin', '*')
self.send_header('Access-Control-Allow-Methods', 'GET, POST, OPTIONS, PUT, DELETE')
self.send_header('Access-Control-Allow-Headers', 'Content-Type, Authorization, X-Requested-With, Cookie, X-QB-Target-Host, X-Custom-LLM-Base, Accept-Encoding, x-api-key, anthropic-version, anthropic-beta')
self.send_header('Access-Control-Max-Age', '86400')
def send_json_response(self, code, data):
resp_bytes = json.dumps(data).encode('utf-8')
self.send_response(code)
self.send_cors_headers()
self.send_header('Content-Type', 'application/json')
self.send_header('Content-Length', str(len(resp_bytes)))
self.end_headers()
self.wfile.write(resp_bytes)
def do_OPTIONS(self):
self.send_response(204)
self.send_cors_headers()
self.end_headers()
def do_HEAD(self):
if self.path.startswith('/api/stream'):
self.handle_stream_video(head_only=True)
else:
super().do_HEAD()
def do_GET(self):
# Database Cache API
if self.path.startswith('/api/db/cache'):
self.handle_db_get_cache()
# Relocation History API
elif self.path.startswith('/api/history'):
self.handle_history_list()
# Storage Analytics & Heatmap API
elif self.path.startswith('/api/analytics/storage'):
self.handle_storage_analytics()
# SMB Maintenance Cleanup Scan API
elif self.path.startswith('/api/smb/scan-cleanup'):
self.handle_smb_scan_cleanup()
# Video Stream API (Range 206 Partial Content)
elif self.path.startswith('/api/stream'):
self.handle_stream_video()
# Existing Performer Directories Intellisense API
elif self.path.startswith('/api/performers'):
self.handle_list_performers()
# Performer Image Proxy API
elif self.path.startswith('/api/performer/image'):
self.handle_performer_image()
elif self.path.startswith('/api/image/proxy'):
self.handle_image_proxy()
# Stash Status, Tasks, Duplicates & Performer Review APIs
elif self.path.startswith('/api/stash/status'):
self.handle_stash_status()
elif self.path.startswith('/api/stash/tasks/jobs'):
self.handle_stash_tasks_jobs()
elif self.path.startswith('/api/stash/duplicates'):
self.handle_stash_duplicates()
elif self.path.startswith('/api/stash/plugins'):
self.handle_stash_plugins_list()
elif self.path.startswith('/api/stash/performers/duplicates'):
self.handle_stash_performer_duplicates()
elif self.path.startswith('/api/stash/performers/missing-images'):
self.handle_stash_missing_images()
elif self.path.startswith('/api/stash/scenes/missing-metadata'):
self.handle_stash_missing_metadata()
elif self.path.startswith('/api/stash/tags/duplicates'):
self.handle_stash_tag_duplicates()
elif self.path.startswith('/api/stash/transcodes/scan'):
self.handle_stash_transcode_scan()
# SMB API
elif self.path.startswith('/api/smb/status'):
self.handle_smb_status()
elif self.path.startswith('/api/smb/files'):
self.handle_smb_list_files()
elif self.path.startswith('/qb-proxy/'):
self.proxy_qb_request('GET')
elif self.path.startswith('/custom-llm-proxy/'):
raw_target = self.headers.get('X-Custom-LLM-Base') or 'http://localhost:1234/v1'
target_base = normalize_url_base(raw_target)
self.proxy_external_api(target_base, '/custom-llm-proxy/', method='GET')
elif self.path == '/server-status':
status_data = {
"status": "online",
"target": DEFAULT_QB_HOST,
"has_sid": bool(latest_sid),
"smb_mounted": os.path.exists(DEFAULT_SMB_SOURCE),
"stash_url": DEFAULT_STASH_URL,
"db": os.path.exists(DB_PATH)
}
self.send_json_response(200, status_data)
elif self.path == '/' or self.path == '':
self.path = '/qbittorrent_performer_sorter.html'
super().do_GET()
else:
super().do_GET()
def do_POST(self):
# Database Cache Save/Delete APIs
if self.path == '/api/db/save':
self.handle_db_save()
elif self.path == '/api/db/batch-save':
self.handle_db_batch_save()
elif self.path == '/api/db/delete':
self.handle_db_delete()
# Relocation History & Undo APIs
elif self.path == '/api/history/record':
self.handle_history_record()
elif self.path == '/api/history/undo':
self.handle_history_undo()
elif self.path == '/api/history/clear':
self.handle_history_clear()
# Stash Lookup, Tasks, Scan & Image Review APIs
elif self.path == '/api/stash/lookup':
self.handle_stash_lookup()
elif self.path == '/api/stash/batch-lookup':
self.handle_stash_batch_lookup()
elif self.path == '/api/stash/scan':
self.handle_stash_metadata_scan()
elif self.path == '/api/stash/tasks/scan':
self.handle_stash_task_scan()
elif self.path == '/api/stash/tasks/identify':
self.handle_stash_task_identify()
elif self.path == '/api/stash/tasks/clean':
self.handle_stash_task_clean()
elif self.path == '/api/stash/tasks/generate':
self.handle_stash_task_generate()
elif self.path == '/api/stash/tasks/stop':
self.handle_stash_task_stop()
elif self.path == '/api/stash/plugins/run-task':
self.handle_stash_run_plugin_task()
elif self.path == '/api/stash/scenes/destroy':
self.handle_stash_scene_destroy()
elif self.path == '/api/stash/scenes/batch-autotag':
self.handle_stash_batch_autotag()
elif self.path == '/api/stash/performers/find-images':
self.handle_stash_find_images()
elif self.path == '/api/stash/performers/commit-image':
self.handle_stash_commit_image()
elif self.path == '/api/stash/performers/batch-commit-images':
self.handle_stash_batch_commit_images()
elif self.path == '/api/stash/performers/merge':
self.handle_stash_performer_merge()
elif self.path == '/api/stash/performers/batch-merge':
self.handle_stash_performer_batch_merge()
elif self.path == '/api/stash/tags/merge':
self.handle_stash_tag_merge()
elif self.path == '/api/stash/tags/batch-merge':
self.handle_stash_tag_batch_merge()
elif self.path == '/api/stash/transcodes/replace':
self.handle_stash_transcode_replace()
elif self.path == '/api/stash/transcodes/batch-replace':
self.handle_stash_transcode_batch_replace()
# SMB API Moves, Deletes & Cleanups
elif self.path == '/api/smb/move':
self.handle_smb_move()
elif self.path == '/api/smb/batch-move':
self.handle_smb_batch_move()
elif self.path == '/api/smb/delete':
self.handle_smb_delete()
elif self.path == '/api/smb/cleanup':
self.handle_smb_execute_cleanup()
# Webhook API
elif self.path == '/api/webhooks/send':
self.handle_send_webhook()
elif self.path.startswith('/qb-proxy/'):
self.proxy_qb_request('POST')
elif self.path.startswith('/anthropic-proxy/'):
self.proxy_external_api('https://api.anthropic.com', '/anthropic-proxy/', method='POST')
elif self.path.startswith('/openai-proxy/'):
self.proxy_external_api('https://api.openai.com', '/openai-proxy/', method='POST')
elif self.path.startswith('/custom-llm-proxy/'):
raw_target = self.headers.get('X-Custom-LLM-Base') or 'http://localhost:1234/v1'
target_base = normalize_url_base(raw_target)
self.proxy_external_api(target_base, '/custom-llm-proxy/', method='POST')
else:
self.send_error(404, "Endpoint not found")
# ==================== SQLite Persistence Handlers ====================
def handle_db_get_cache(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
item_type = params.get('type', ['all'])[0]
try:
conn = get_db_connection()
cursor = conn.cursor()
if item_type != 'all':
cursor.execute('SELECT * FROM performer_analysis_cache WHERE item_type = ?', (item_type,))
else:
cursor.execute('SELECT * FROM performer_analysis_cache')
rows = cursor.fetchall()
cache_map = {}
for r in rows:
all_perf = []
try:
all_perf = json.loads(r['all_performers']) if r['all_performers'] else []
except Exception:
pass
cache_map[r['item_key']] = {
"item_type": r['item_type'],
"item_key": r['item_key'],
"name": r['name'],
"file_size": r['file_size'],
"file_mtime": r['file_mtime'],
"primary_performer": r['primary_performer'],
"all_performers": all_perf,
"selected_performer": r['selected_performer'],
"target_path": r['target_path'],
"confidence": r['confidence'],
"reasoning": r['reasoning'],
"move_status": r['move_status'],
"updated_at": r['updated_at']
}
conn.close()
self.send_json_response(200, {"cache": cache_map, "count": len(cache_map)})
except Exception as e:
self.send_json_response(500, {"error": f"Failed reading SQLite cache: {str(e)}"})
def handle_db_save(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
item = json.loads(post_data.decode('utf-8'))
item_type = item.get('item_type', 'smb')
item_key = item.get('item_key')
name = item.get('name', '')
file_size = item.get('file_size', 0)
file_mtime = item.get('file_mtime', 0)
primary_performer = item.get('primary_performer', '')
all_performers = json.dumps(item.get('all_performers', []))
selected_performer = item.get('selected_performer', primary_performer)
target_path = item.get('target_path', '')
confidence = item.get('confidence', 'medium')
reasoning = item.get('reasoning', '')
move_status = item.get('move_status', 'ready')
if not item_key or not name:
self.send_response(400)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": "Missing item_key or name"}).encode('utf-8'))
return
conn = get_db_connection()
cursor = conn.cursor()
cursor.execute('''
INSERT INTO performer_analysis_cache (
item_type, item_key, name, file_size, file_mtime,
primary_performer, all_performers, selected_performer, target_path,
confidence, reasoning, move_status, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP)
ON CONFLICT(item_key) DO UPDATE SET
primary_performer = excluded.primary_performer,
all_performers = excluded.all_performers,
selected_performer = excluded.selected_performer,
target_path = excluded.target_path,
confidence = excluded.confidence,
reasoning = excluded.reasoning,
move_status = excluded.move_status,
updated_at = CURRENT_TIMESTAMP
''', (item_type, item_key, name, file_size, file_mtime, primary_performer, all_performers, selected_performer, target_path, confidence, reasoning, move_status))
conn.commit()
conn.close()
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"success": True, "item_key": item_key}).encode('utf-8'))
except Exception as e:
self.send_response(500)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": f"Failed saving to DB: {str(e)}"}).encode('utf-8'))
def handle_db_batch_save(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
items = payload.get('items', [])
conn = get_db_connection()
cursor = conn.cursor()
saved_count = 0
for item in items:
item_type = item.get('item_type', 'smb')
item_key = item.get('item_key')
name = item.get('name', '')
file_size = item.get('file_size', 0)
file_mtime = item.get('file_mtime', 0)
primary_performer = item.get('primary_performer', '')
all_performers = json.dumps(item.get('all_performers', []))
selected_performer = item.get('selected_performer', primary_performer)
target_path = item.get('target_path', '')
confidence = item.get('confidence', 'medium')
reasoning = item.get('reasoning', '')
move_status = item.get('move_status', 'ready')
if not item_key or not name:
continue
cursor.execute('''
INSERT INTO performer_analysis_cache (
item_type, item_key, name, file_size, file_mtime,
primary_performer, all_performers, selected_performer, target_path,
confidence, reasoning, move_status, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP)
ON CONFLICT(item_key) DO UPDATE SET
primary_performer = excluded.primary_performer,
all_performers = excluded.all_performers,
selected_performer = excluded.selected_performer,
target_path = excluded.target_path,
confidence = excluded.confidence,
reasoning = excluded.reasoning,
move_status = excluded.move_status,
updated_at = CURRENT_TIMESTAMP
''', (item_type, item_key, name, file_size, file_mtime, primary_performer, all_performers, selected_performer, target_path, confidence, reasoning, move_status))
saved_count += 1
conn.commit()
conn.close()
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"success": True, "saved_count": saved_count}).encode('utf-8'))
except Exception as e:
self.send_response(500)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": f"Failed batch saving to DB: {str(e)}"}).encode('utf-8'))
def handle_db_delete(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
item_key = payload.get('item_key')
item_keys = payload.get('item_keys', [])
if item_key and item_key not in item_keys:
item_keys.append(item_key)
if not item_keys:
self.send_response(400)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": "Missing item_key or item_keys"}).encode('utf-8'))
return
conn = get_db_connection()
cursor = conn.cursor()
for key in item_keys:
cursor.execute("DELETE FROM performer_analysis_cache WHERE item_key = ? OR name = ?", (key, key))
conn.commit()
conn.close()
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"success": True, "deleted_keys": item_keys}).encode('utf-8'))
except Exception as e:
self.send_response(500)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": f"Failed deleting from DB: {str(e)}"}).encode('utf-8'))
# ==================== Video Streaming Handlers (Range 206) ====================
def handle_stream_video(self):
head_only = False
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
raw_path = params.get('path', [''])[0]
file_path = resolve_stream_video_path(raw_path)
if not file_path or not os.path.exists(file_path) or not os.path.isfile(file_path):
self.send_response(404)
self.send_cors_headers()
self.end_headers()
if not head_only:
self.wfile.write(b"File not found")
return
file_size = os.path.getsize(file_path)
ext = os.path.splitext(file_path)[1].lower()
# MIME mapping
mime_map = {
'.mp4': 'video/mp4',
'.m4v': 'video/mp4',
'.webm': 'video/webm',
'.mkv': 'video/mp4',
'.mov': 'video/quicktime',
'.avi': 'video/x-msvideo',
'.ts': 'video/mp2t'
}
mime_type = mime_map.get(ext, 'video/mp4')
range_header = self.headers.get('Range', '')
if range_header and range_header.startswith('bytes='):
try:
range_match = re.match(r'bytes=(\d+)-(\d*)', range_header)
if range_match:
start = int(range_match.group(1))
end_str = range_match.group(2)
end = int(end_str) if end_str else file_size - 1
end = min(end, file_size - 1)
if start > end or start >= file_size:
self.send_response(416)
self.send_header('Content-Range', f'bytes */{file_size}')
self.send_cors_headers()
self.end_headers()
return
length = end - start + 1
self.send_response(206)
self.send_header('Content-Type', mime_type)
self.send_header('Content-Range', f'bytes {start}-{end}/{file_size}')
self.send_header('Content-Length', str(length))
self.send_header('Accept-Ranges', 'bytes')
self.send_cors_headers()
self.end_headers()
if head_only:
return
with open(file_path, 'rb') as f:
f.seek(start)
remaining = length
chunk_size = 128 * 1024
while remaining > 0:
read_bytes = min(chunk_size, remaining)
data = f.read(read_bytes)
if not data:
break
self.wfile.write(data)
remaining -= len(data)
return
except (ConnectionResetError, BrokenPipeError):
return
except Exception as e:
print(f"Video stream range error: {e}")
return
# Serve full file
self.send_response(200)
self.send_header('Content-Type', mime_type)
self.send_header('Content-Length', str(file_size))
self.send_header('Accept-Ranges', 'bytes')
self.send_cors_headers()
self.end_headers()
if head_only:
return
try:
with open(file_path, 'rb') as f:
chunk_size = 128 * 1024
while True:
data = f.read(chunk_size)
if not data:
break
self.wfile.write(data)
except (ConnectionResetError, BrokenPipeError):
pass
except Exception as e:
print(f"Video stream full error: {e}")
# ==================== Existing Performer Intellisense Handlers ====================
def handle_list_performers(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
target_dir = params.get('path', [DEFAULT_SMB_TARGET_ROOT])[0] or DEFAULT_SMB_TARGET_ROOT
if not os.path.exists(target_dir):
ensure_smb_mounted()
performers_set = set()
if os.path.exists(target_dir) and os.path.isdir(target_dir):
try:
for entry in os.scandir(target_dir):
if entry.is_dir() and entry.name not in ('videodownloader', 'unsorted') and not entry.name.startswith('.'):
performers_set.add(entry.name)
except Exception as e:
print(f"Error scanning existing performer directories: {e}")
# Also merge distinct performers from SQLite cache
try:
conn = sqlite3.connect(DB_PATH)
cursor = conn.cursor()
cursor.execute("SELECT DISTINCT selected_performer FROM performer_analysis_cache WHERE selected_performer IS NOT NULL AND selected_performer != ''")
for row in cursor.fetchall():
if row[0]:
performers_set.add(row[0])
conn.close()
except Exception as e:
print(f"Error merging performers from DB: {e}")
# Also merge female performers from Stash database if online
try:
stash_q = '{ findPerformers(performer_filter: { gender: { value: FEMALE, modifier: EQUALS } }, filter: { per_page: -1 }) { performers { name } } }'
res = execute_stash_graphql(stash_q, timeout=3)
for p in res.get('data', {}).get('findPerformers', {}).get('performers', []):
slug = format_performer_slug(p.get('name'))
if slug:
performers_set.add(slug)
except Exception as e:
pass
performers = sorted(list(performers_set), key=lambda s: s.lower())
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"performers": performers, "count": len(performers)}).encode('utf-8'))
# ==================== Performer Image Proxy Handlers ====================
def handle_performer_image(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
name = params.get('name', [''])[0]
p_id = params.get('id', [''])[0]
stash_url = params.get('url', [''])[0]
try:
content_type, img_bytes = fetch_performer_image_data(name, p_id, stash_url)
self.send_response(200)
self.send_header('Content-Type', content_type)
self.send_header('Cache-Control', 'public, max-age=86400')
self.send_cors_headers()
self.end_headers()
self.wfile.write(img_bytes)
except (BrokenPipeError, ConnectionResetError):
pass
except Exception as e:
try:
self.send_response(500)
self.end_headers()
except Exception:
pass
def handle_image_proxy(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
target_url = params.get('url', [''])[0]
if not target_url:
self.send_response(400)
self.end_headers()
return
try:
req = urllib.request.Request(
target_url,
headers={
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/122.0.0.0 Safari/537.36',
'Referer': target_url
}
)
with urllib.request.urlopen(req, timeout=10) as resp:
content_type = resp.headers.get('Content-Type', 'image/jpeg')
img_data = resp.read()
self.send_response(200)
self.send_header('Content-Type', content_type)
self.send_header('Cache-Control', 'public, max-age=86400')
self.send_cors_headers()
self.end_headers()
self.wfile.write(img_data)
except Exception:
self.send_response(404)
self.end_headers()
# ==================== Stash Server API Handlers ====================
def handle_stash_status(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
stash_url = params.get('url', [''])[0] or DEFAULT_STASH_URL
try:
q = '{ version { version } stats { scene_count performer_count } }'
res = execute_stash_graphql(q, stash_url, timeout=3)
data = res.get('data', {})
version = data.get('version', {}).get('version', 'unknown')
scene_count = data.get('stats', {}).get('scene_count', 0)
performer_count = data.get('stats', {}).get('performer_count', 0)
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({
"online": True,
"version": version,
"sceneCount": scene_count,
"performerCount": performer_count,
"url": normalize_stash_url(stash_url)
}).encode('utf-8'))
except Exception as e:
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({
"online": False,
"error": str(e),
"url": normalize_stash_url(stash_url)
}).encode('utf-8'))
def handle_stash_lookup(self):
content_length = int(self.headers.get('Content-Length', 0))
body = self.rfile.read(content_length).decode('utf-8')
try:
req_data = json.loads(body)
file_name = req_data.get('name', '')
file_path = req_data.get('path', '')
item_type = req_data.get('item_type', 'smb')
stash_url = req_data.get('stash_url', '')
match = query_stash_single(file_name, file_path, item_type, stash_url)
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"match": match, "found": bool(match)}).encode('utf-8'))
except Exception as e:
self.send_response(500)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": str(e), "match": None, "found": False}).encode('utf-8'))
def handle_stash_batch_lookup(self):
content_length = int(self.headers.get('Content-Length', 0))
body = self.rfile.read(content_length).decode('utf-8')
try:
req_data = json.loads(body)
items = req_data.get('items', [])
stash_url = req_data.get('stash_url', '')
results = query_stash_batch(items, stash_url)
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"results": results, "count": len(results)}).encode('utf-8'))
except Exception as e:
self.send_response(500)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": str(e), "results": []}).encode('utf-8'))
def handle_stash_missing_images(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
stash_url = params.get('url', [''])[0] or DEFAULT_STASH_URL
gender_filter = params.get('gender', ['FEMALE'])[0].upper()
only_scenes = params.get('has_scenes', ['true'])[0].lower() in ['true', '1']
full_names_only = params.get('full_names_only', ['true'])[0].lower() in ['true', '1']
query = '''
query {
allPerformers {
id
name
disambiguation
gender
image_path
scene_count
image_count
alias_list
country
birthdate
}
}
'''
try:
res = execute_stash_graphql(query, stash_url, timeout=15)
all_p = res.get('data', {}).get('allPerformers', [])
missing = []
for p in all_p:
img_path = p.get('image_path') or ''
is_missing = not img_path or 'default=true' in img_path or 'default' in img_path.lower()
if not is_missing:
continue
# Filter out ambiguous single-word names (only first names)
p_name = (p.get('name') or '').strip()
if full_names_only:
name_parts = [pt for pt in re.split(r'[\s\.\-_]+', p_name) if len(pt) > 1]
if len(name_parts) < 2:
continue
p_gender = (p.get('gender') or '').upper()
if gender_filter != 'ALL':
if gender_filter == 'FEMALE' and p_gender not in ['FEMALE', '', 'NONE', None]:
continue
elif gender_filter != 'FEMALE' and p_gender != gender_filter:
continue
scenes = p.get('scene_count') or 0
if only_scenes and scenes <= 0:
continue
missing.append(p)
# Sort: performers with more scenes first, then by name
missing.sort(key=lambda x: (-(x.get('scene_count') or 0), x.get('name', '').lower()))
self.send_json_response(200, {
"performers": missing,
"count": len(missing),
"total_in_stash": len(all_p)
})
except Exception as e:
self.send_json_response(500, {"error": f"Failed fetching performers from Stash: {str(e)}", "performers": []})
def handle_stash_find_images(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
name = payload.get('name', '')
aliases = payload.get('aliases', [])
candidates = search_performer_candidate_images(name, aliases)
self.send_json_response(200, {
"name": name,
"candidates": candidates,
"count": len(candidates)
})
except Exception as e:
self.send_json_response(500, {"error": str(e), "candidates": []})
def handle_stash_commit_image(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
performer_id = str(payload.get('id', ''))
image_url = payload.get('image_url', '')
base64_data = payload.get('base64', '')
stash_url = payload.get('stash_url', '') or DEFAULT_STASH_URL
if not performer_id:
self.send_json_response(400, {"error": "Missing performer ID"})
return
if not base64_data and image_url:
base64_data = fetch_image_as_base64(image_url)
if not base64_data:
self.send_json_response(400, {"error": "Missing image data or URL"})
return
mutation = f'''
mutation {{
performerUpdate(input: {{ id: {json.dumps(performer_id)}, image: {json.dumps(base64_data)} }}) {{
id
name
image_path
}}
}}
'''
res = execute_stash_graphql(mutation, stash_url, timeout=20)
if 'errors' in res and res['errors']:
err_msg = res['errors'][0].get('message', 'GraphQL error')
self.send_json_response(500, {"error": err_msg})
return
updated = res.get('data', {}).get('performerUpdate', {})
self.send_json_response(200, {
"success": True,
"id": performer_id,
"name": updated.get('name'),
"image_path": updated.get('image_path')
})
except Exception as e:
self.send_json_response(500, {"error": f"Failed committing image to Stash: {str(e)}"})
def handle_stash_batch_commit_images(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
updates = payload.get('updates', [])
stash_url = payload.get('stash_url', '') or DEFAULT_STASH_URL
results = []
for item in updates:
p_id = str(item.get('id', ''))
img_url = item.get('image_url', '')
b64 = item.get('base64', '')
try:
if not b64 and img_url:
b64 = fetch_image_as_base64(img_url)
if not b64:
results.append({"id": p_id, "success": False, "error": "No image data"})
continue
mutation = f'''
mutation {{
performerUpdate(input: {{ id: {json.dumps(p_id)}, image: {json.dumps(b64)} }}) {{
id
name
image_path
}}
}}
'''
res = execute_stash_graphql(mutation, stash_url, timeout=20)
if 'errors' in res and res['errors']:
err_msg = res['errors'][0].get('message', 'GraphQL error')
results.append({"id": p_id, "success": False, "error": err_msg})
else:
updated = res.get('data', {}).get('performerUpdate', {})
results.append({
"id": p_id,
"success": True,
"name": updated.get('name'),
"image_path": updated.get('image_path')
})
except Exception as ex:
results.append({"id": p_id, "success": False, "error": str(ex)})
self.send_json_response(200, {
"results": results,
"total": len(updates),
"successful": len([r for r in results if r.get('success')])
})
except Exception as e:
self.send_json_response(500, {"error": str(e), "results": []})
# ==================== SMB File Management Handlers ====================
def handle_smb_status(self):
mounted, msg = ensure_smb_mounted()
resp = {
"mounted": mounted,
"message": msg,
"sourcePath": DEFAULT_SMB_SOURCE,
"targetRoot": DEFAULT_SMB_TARGET_ROOT,
"exists": os.path.exists(DEFAULT_SMB_SOURCE)
}
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps(resp).encode('utf-8'))
def handle_smb_list_files(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
dir_path = params.get('path', [DEFAULT_SMB_SOURCE])[0]
if not os.path.exists(dir_path):
ensure_smb_mounted()
if not os.path.exists(dir_path):
self.send_response(404)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": f"Directory not found or SMB unmounted: {dir_path}"}).encode('utf-8'))
return
files_list = []
try:
with os.scandir(dir_path) as it:
for entry in it:
if entry.is_file():
ext = os.path.splitext(entry.name)[1].lower()
if ext in VIDEO_EXTENSIONS:
try:
stat_res = entry.stat()
files_list.append({
"name": entry.name,
"path": entry.path,
"size": stat_res.st_size,
"modified": stat_res.st_mtime,
"ext": ext
})
except Exception:
files_list.append({
"name": entry.name,
"path": entry.path,
"size": 0,
"modified": 0,
"ext": ext
})
# Sort by name
files_list.sort(key=lambda x: x["name"].lower())
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"files": files_list, "count": len(files_list), "dir": dir_path}).encode('utf-8'))
except Exception as e:
self.send_response(500)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": f"Failed reading directory: {str(e)}"}).encode('utf-8'))
def handle_smb_move(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
source_path = payload.get('source')
target_dir = payload.get('targetDir')
target_file = payload.get('targetFile')
performer = payload.get('performer', '')
if not source_path or not os.path.exists(source_path):
self.send_response(400)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": f"Source file does not exist: {source_path}"}).encode('utf-8'))
return
if not target_file and target_dir:
filename = os.path.basename(source_path)
target_file = os.path.join(target_dir, filename)
if not target_file:
self.send_response(400)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": "Missing target path"}).encode('utf-8'))
return
dest_dir = os.path.dirname(target_file)
os.makedirs(dest_dir, exist_ok=True)
f_size = 0
try:
f_size = os.path.getsize(source_path)
except Exception:
pass
# Move file
shutil.move(source_path, target_file)
# Record in relocation history
record_relocation_history('smb', source_path, os.path.basename(source_path), source_path, target_file, performer=performer, file_size=f_size)
# Update DB status if entry exists
try:
conn = get_db_connection()
cursor = conn.cursor()
cursor.execute("UPDATE performer_analysis_cache SET move_status = 'moved', target_path = ?, updated_at = CURRENT_TIMESTAMP WHERE item_key = ?", (dest_dir, source_path))
conn.commit()
conn.close()
except Exception:
pass
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"success": True, "source": source_path, "target": target_file}).encode('utf-8'))
except Exception as e:
self.send_response(500)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": f"Move failed: {str(e)}"}).encode('utf-8'))
def handle_smb_batch_move(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
items = payload.get('items', [])
results = []
for item in items:
source = item.get('source')
target_dir = item.get('targetDir')
target_file = item.get('targetFile')
performer = item.get('performer', '')
if not target_file and target_dir and source:
target_file = os.path.join(target_dir, os.path.basename(source))
if not source or not os.path.exists(source):
results.append({"source": source, "success": False, "error": "File does not exist"})
continue
try:
dest_dir = os.path.dirname(target_file)
os.makedirs(dest_dir, exist_ok=True)
f_size = 0
try:
f_size = os.path.getsize(source)
except Exception:
pass
shutil.move(source, target_file)
# Record history
record_relocation_history('smb', source, os.path.basename(source), source, target_file, performer=performer, file_size=f_size)
# Update DB
try:
conn = get_db_connection()
cursor = conn.cursor()
cursor.execute("UPDATE performer_analysis_cache SET move_status = 'moved', target_path = ?, updated_at = CURRENT_TIMESTAMP WHERE item_key = ?", (dest_dir, source))
conn.commit()
conn.close()
except Exception:
pass
results.append({"source": source, "target": target_file, "success": True})
except Exception as e:
results.append({"source": source, "success": False, "error": str(e)})
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"results": results, "total": len(items)}).encode('utf-8'))
except Exception as e:
self.send_response(500)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": f"Batch move failed: {str(e)}"}).encode('utf-8'))
def handle_smb_delete(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
path = payload.get('path')
paths = payload.get('paths', [])
if path and path not in paths:
paths.append(path)
if not paths:
self.send_response(400)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": "Missing path or paths"}).encode('utf-8'))
return
deleted = []
errors = []
conn = get_db_connection()
cursor = conn.cursor()
for p in paths:
try:
resolved = resolve_stream_video_path(p) or p
if os.path.exists(resolved):
if os.path.isdir(resolved):
shutil.rmtree(resolved)
else:
os.remove(resolved)
deleted.append(p)
cursor.execute("DELETE FROM performer_analysis_cache WHERE item_key = ? OR name = ?", (p, os.path.basename(p)))
else:
errors.append(f"File not found: {p}")
except Exception as del_err:
errors.append(f"Failed deleting {p}: {str(del_err)}")
conn.commit()
conn.close()
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({
"success": len(deleted) > 0,
"deleted": deleted,
"deleted_count": len(deleted),
"errors": errors
}).encode('utf-8'))
except Exception as e:
self.send_response(500)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"error": f"Delete failed: {str(e)}"}).encode('utf-8'))
# ==================== Relocation History & Undo Handlers ====================
def handle_history_list(self):
try:
conn = get_db_connection()
cursor = conn.cursor()
cursor.execute("SELECT * FROM relocation_history ORDER BY id DESC LIMIT 250")
rows = [dict(r) for r in cursor.fetchall()]
cursor.execute("SELECT COUNT(*) as total, SUM(CASE WHEN status='completed' THEN 1 ELSE 0 END) as active, SUM(CASE WHEN status='undone' THEN 1 ELSE 0 END) as undone, SUM(file_size) as total_bytes FROM relocation_history")
raw_stats = cursor.fetchone()
stats = dict(raw_stats) if raw_stats else {}
conn.close()
self.send_json_response(200, {"history": rows, "stats": stats, "count": len(rows)})
except Exception as e:
self.send_json_response(500, {"error": str(e), "history": []})
def handle_history_record(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
p = json.loads(post_data.decode('utf-8'))
hid = record_relocation_history(
p.get('source_type', 'qb'),
p.get('item_key', ''),
p.get('item_name', ''),
p.get('source_path', ''),
p.get('target_path', ''),
performer=p.get('performer', ''),
file_size=p.get('file_size', 0)
)
self.send_json_response(200, {"success": True, "history_id": hid})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
def handle_history_undo(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
history_id = payload.get('history_id')
if not history_id:
self.send_json_response(400, {"error": "Missing history_id"})
return
conn = get_db_connection()
cursor = conn.cursor()
cursor.execute("SELECT * FROM relocation_history WHERE id = ?", (history_id,))
row = cursor.fetchone()
if not row:
conn.close()
self.send_json_response(404, {"error": "History entry not found"})
return
item = dict(row)
source_path = item['source_path']
target_path = item['target_path']
source_type = item['source_type']
item_key = item['item_key']
if item['status'] == 'undone':
conn.close()
self.send_json_response(400, {"error": "Entry has already been undone"})
return
# Execute Undo
if source_type == 'smb':
if os.path.exists(target_path):
os.makedirs(os.path.dirname(source_path), exist_ok=True)
shutil.move(target_path, source_path)
elif os.path.exists(os.path.join(target_path, os.path.basename(source_path))):
actual_file = os.path.join(target_path, os.path.basename(source_path))
os.makedirs(os.path.dirname(source_path), exist_ok=True)
shutil.move(actual_file, source_path)
else:
conn.close()
self.send_json_response(404, {"error": f"Target file no longer exists at {target_path}"})
return
# Update analysis cache
cursor.execute("UPDATE performer_analysis_cache SET move_status = 'ready', updated_at = CURRENT_TIMESTAMP WHERE item_key = ?", (source_path,))
cursor.execute("UPDATE relocation_history SET status = 'undone', undone_at = CURRENT_TIMESTAMP WHERE id = ?", (history_id,))
conn.commit()
conn.close()
self.send_json_response(200, {"success": True, "undone_item": item})
except Exception as e:
self.send_json_response(500, {"error": f"Undo failed: {str(e)}"})
def handle_history_clear(self):
try:
conn = get_db_connection()
cursor = conn.cursor()
cursor.execute("DELETE FROM relocation_history")
conn.commit()
conn.close()
self.send_json_response(200, {"success": True})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
# ==================== Stash Task & Tool Control Handlers ====================
def handle_stash_tasks_jobs(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
stash_url = params.get('stash_url', [DEFAULT_STASH_URL])[0]
query = '''
query JobQueue {
jobQueue {
id
status
subTasks
description
progress
startTime
endTime
addTime
error
}
}
'''
try:
res = execute_stash_graphql(query, stash_url, timeout=10)
jobs = res.get('data', {}).get('jobQueue') or []
self.send_json_response(200, {"jobs": jobs, "count": len(jobs)})
except Exception as e:
self.send_json_response(500, {"error": str(e), "jobs": []})
def handle_stash_task_scan(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
paths = payload.get('paths', [])
rescan = payload.get('rescan', False)
scanGenerateCovers = payload.get('scanGenerateCovers', True)
scanGeneratePreviews = payload.get('scanGeneratePreviews', False)
scanGenerateSprites = payload.get('scanGenerateSprites', False)
scanGeneratePhashes = payload.get('scanGeneratePhashes', False)
input_obj = {
"paths": paths if paths else None,
"rescan": bool(rescan),
"scanGenerateCovers": bool(scanGenerateCovers),
"scanGeneratePreviews": bool(scanGeneratePreviews),
"scanGenerateSprites": bool(scanGenerateSprites),
"scanGeneratePhashes": bool(scanGeneratePhashes)
}
# Remove None values
input_obj = {k: v for k, v in input_obj.items() if v is not None}
mutation = build_stash_mutation("metadataScan", input_obj=input_obj)
res = execute_stash_graphql(mutation, stash_url, timeout=15)
self.send_json_response(200, {"success": True, "result": res.get('data', {}).get('metadataScan')})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
def handle_stash_task_identify(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
paths = payload.get('paths', [])
mode = payload.get('mode', 'autotag')
if mode == 'autotag':
input_obj = {
"paths": paths if paths else None,
"performers": ["*"] if bool(payload.get('performers', True)) else None,
"studios": ["*"] if bool(payload.get('studios', True)) else None,
"tags": ["*"] if bool(payload.get('tags', True)) else None
}
input_obj = {k: v for k, v in input_obj.items() if v is not None}
mutation = build_stash_mutation("metadataAutoTag", input_obj=input_obj)
res = execute_stash_graphql(mutation, stash_url, timeout=15)
self.send_json_response(200, {"success": True, "job_id": res.get('data', {}).get('metadataAutoTag')})
else:
input_obj = {
"paths": paths if paths else None,
"options": {
"setCoverImage": True,
"setOrganized": True
}
}
input_obj = {k: v for k, v in input_obj.items() if v is not None}
mutation = build_stash_mutation("metadataIdentify", input_obj=input_obj)
res = execute_stash_graphql(mutation, stash_url, timeout=15)
self.send_json_response(200, {"success": True, "job_id": res.get('data', {}).get('metadataIdentify')})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
def handle_stash_task_clean(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
clean_type = payload.get('clean_type', 'metadata')
dry_run = bool(payload.get('dry_run', False))
paths = payload.get('paths', [])
results = {}
if clean_type in ['metadata', 'both']:
input_clean = {"dryRun": dry_run, "paths": paths if paths else None}
input_clean = {k: v for k, v in input_clean.items() if v is not None}
mutation = build_stash_mutation("metadataClean", input_obj=input_clean)
res = execute_stash_graphql(mutation, stash_url, timeout=15)
results['metadataClean'] = res.get('data', {}).get('metadataClean')
if clean_type in ['generated', 'both']:
input_gen = {
"dryRun": dry_run,
"blobFiles": True,
"sprites": True,
"screenshots": True,
"transcodes": True,
"markers": True,
"imageThumbnails": True
}
mutation_gen = build_stash_mutation("metadataCleanGenerated", input_obj=input_gen)
res_gen = execute_stash_graphql(mutation_gen, stash_url, timeout=15)
results['metadataCleanGenerated'] = res_gen.get('data', {}).get('metadataCleanGenerated')
self.send_json_response(200, {"success": True, "results": results})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
def handle_stash_task_generate(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
input_gen = {
"covers": bool(payload.get('covers', True)),
"sprites": bool(payload.get('sprites', True)),
"previews": bool(payload.get('previews', True)),
"imagePreviews": bool(payload.get('imagePreviews', True)),
"markers": bool(payload.get('markers', True)),
"transcodes": bool(payload.get('transcodes', False)),
"phashes": bool(payload.get('phashes', True)),
"overwrite": bool(payload.get('overwrite', False)),
"paths": payload.get('paths') if payload.get('paths') else None
}
input_gen = {k: v for k, v in input_gen.items() if v is not None}
mutation = build_stash_mutation("metadataGenerate", input_obj=input_gen)
res = execute_stash_graphql(mutation, stash_url, timeout=15)
self.send_json_response(200, {"success": True, "job_id": res.get('data', {}).get('metadataGenerate')})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
def handle_stash_task_stop(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
job_id = payload.get('job_id')
stop_all = payload.get('stop_all', False)
if stop_all or not job_id:
mutation = "mutation { stopAllJobs }"
res = execute_stash_graphql(mutation, stash_url, timeout=10)
self.send_json_response(200, {"success": True, "result": res.get('data', {}).get('stopAllJobs')})
else:
mutation = f'mutation {{ stopJob(job_id: "{job_id}") }}'
res = execute_stash_graphql(mutation, stash_url, timeout=10)
self.send_json_response(200, {"success": True, "result": res.get('data', {}).get('stopJob')})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
def handle_stash_metadata_scan(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
paths = payload.get('paths', [])
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
mutation = '''
mutation MetadataScan($input: ScanMetadataInput!) {
metadataScan(input: $input)
}
'''
variables = {
"input": {
"paths": paths if paths else [],
"rescan": False,
"scanGenerateCovers": True
}
}
res = execute_stash_graphql(mutation, stash_url, timeout=20)
self.send_json_response(200, {"success": True, "result": res})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
# ==================== Stash Duplicate Checker & Plugin Integration Handlers ====================
def handle_stash_duplicates(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
stash_url = params.get('stash_url', [DEFAULT_STASH_URL])[0]
distance = int(params.get('distance', [0])[0])
duration_diff = float(params.get('duration_diff', [5.0])[0])
query = '''
query FindDuplicateScenes($distance: Int, $duration_diff: Float) {
findDuplicateScenes(distance: $distance, duration_diff: $duration_diff) {
id
title
date
details
rating100
organized
paths {
screenshot
preview
stream
}
files {
id
path
size
duration
video_codec
width
height
bit_rate
frame_rate
}
performers {
id
name
image_path
}
studio {
id
name
image_path
}
tags {
id
name
}
}
}
'''
try:
res = execute_stash_graphql(query, stash_url, variables={"distance": distance, "duration_diff": duration_diff}, timeout=30)
raw_groups = res.get('data', {}).get('findDuplicateScenes') or []
duplicate_groups = []
for group in raw_groups:
if isinstance(group, list) and len(group) > 1:
sorted_group = sorted(group, key=lambda s: (
(s.get('files', [{}])[0].get('height') or 0),
(s.get('files', [{}])[0].get('bit_rate') or 0),
(s.get('files', [{}])[0].get('size') or 0)
), reverse=True)
duplicate_groups.append(sorted_group)
self.send_json_response(200, {
"groups": duplicate_groups,
"group_count": len(duplicate_groups),
"total_duplicate_scenes": sum(len(g) for g in duplicate_groups),
"distance": distance,
"duration_diff": duration_diff
})
except Exception as e:
self.send_json_response(500, {"error": str(e), "groups": []})
def handle_stash_plugins_list(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
stash_url = params.get('stash_url', [DEFAULT_STASH_URL])[0]
query = '''
{
plugins {
id
name
description
version
enabled
tasks {
name
description
}
}
}
'''
try:
res = execute_stash_graphql(query, stash_url, timeout=10)
plugins = res.get('data', {}).get('plugins') or []
self.send_json_response(200, {"plugins": plugins, "count": len(plugins)})
except Exception as e:
self.send_json_response(500, {"error": str(e), "plugins": []})
def handle_stash_run_plugin_task(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
plugin_id = payload.get('plugin_id')
task_name = payload.get('task_name')
if not plugin_id or not task_name:
self.send_json_response(400, {"error": "Missing plugin_id or task_name"})
return
mutation = f'mutation {{ runPluginTask(plugin_id: "{plugin_id}", task_name: "{task_name}") }}'
res = execute_stash_graphql(mutation, stash_url, timeout=15)
self.send_json_response(200, {"success": True, "job_id": res.get('data', {}).get('runPluginTask')})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
def handle_stash_scene_destroy(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
scene_id = payload.get('scene_id')
delete_file = bool(payload.get('delete_file', True))
delete_generated = bool(payload.get('delete_generated', True))
if not scene_id:
self.send_json_response(400, {"error": "Missing scene_id"})
return
input_obj = {
"id": str(scene_id),
"delete_file": delete_file,
"delete_generated": delete_generated
}
mutation = build_stash_mutation("sceneDestroy", input_obj=input_obj)
res = execute_stash_graphql(mutation, stash_url, timeout=15)
self.send_json_response(200, {"success": True, "result": res.get('data', {}).get('sceneDestroy')})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
# ==================== Performer Deduplication & Merge Handlers ====================
def handle_stash_performer_duplicates(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
stash_url = params.get('stash_url', [DEFAULT_STASH_URL])[0]
query = '''
{
findPerformers(filter: { per_page: -1 }) {
count
performers {
id
name
gender
scene_count
image_path
alias_list
stash_ids {
endpoint
stash_id
}
}
}
}
'''
try:
res = execute_stash_graphql(query, stash_url, timeout=30)
performers = res.get('data', {}).get('findPerformers', {}).get('performers') or []
# Group by normalized clean name
clusters = {}
for p in performers:
raw_name = p.get('name') or ''
# Normalize: remove dots, underscores, dashes, trailing roman numerals/numbers
norm = re.sub(r'\s*\([ivx0-9]+\)\s*$', '', raw_name, flags=re.IGNORECASE)
norm = re.sub(r'[\._\-]+', ' ', norm).strip().lower()
norm_compact = re.sub(r'[^a-z0-9]', '', norm)
if not norm_compact:
continue
if norm_compact not in clusters:
clusters[norm_compact] = []
clusters[norm_compact].append(p)
# Filter only clusters with > 1 performer
duplicate_clusters = []
for k, group in clusters.items():
if len(group) > 1:
# Sort so performer with most scenes / image is primary
sorted_group = sorted(group, key=lambda x: (
bool(x.get('image_path')),
x.get('scene_count') or 0,
len(x.get('stash_ids') or [])
), reverse=True)
duplicate_clusters.append({
"normalized_name": k,
"primary": sorted_group[0],
"candidates": sorted_group[1:],
"all": sorted_group,
"count": len(sorted_group)
})
self.send_json_response(200, {
"clusters": duplicate_clusters,
"cluster_count": len(duplicate_clusters),
"total_performers": len(performers)
})
except Exception as e:
self.send_json_response(500, {"error": str(e), "clusters": []})
def handle_stash_performer_merge(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
source_ids = [str(sid) for sid in payload.get('source_ids', [])]
dest_id = str(payload.get('destination_id'))
if not source_ids or not dest_id:
self.send_json_response(400, {"error": "Missing source_ids or destination_id"})
return
input_obj = {
"source": source_ids,
"destination": dest_id
}
mutation = build_stash_mutation("performerMerge", input_obj=input_obj, fields="id name scene_count")
res = execute_stash_graphql(mutation, stash_url, timeout=20)
self.send_json_response(200, {"success": True, "merged_performer": res.get('data', {}).get('performerMerge')})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
def handle_stash_performer_batch_merge(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
clusters = payload.get('clusters', [])
if not clusters:
self.send_json_response(400, {"error": "No clusters provided for batch merge"})
return
merged_count = 0
errors = []
for c in clusters:
source_ids = [str(sid) for sid in c.get('source_ids', [])]
dest_id = str(c.get('destination_id'))
if not source_ids or not dest_id:
continue
try:
mutation = build_stash_mutation("performerMerge", input_obj={"source": source_ids, "destination": dest_id}, fields="id name")
res = execute_stash_graphql(mutation, stash_url, timeout=20)
if res.get('data', {}).get('performerMerge'):
merged_count += 1
elif res.get('errors'):
errors.append({"destination_id": dest_id, "error": str(res.get('errors'))})
except Exception as ex:
errors.append({"destination_id": dest_id, "error": str(ex)})
self.send_json_response(200, {
"success": True,
"merged_count": merged_count,
"total_requested": len(clusters),
"errors": errors
})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
# ==================== Missing Metadata Radar Handlers ====================
def handle_stash_missing_metadata(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
stash_url = params.get('stash_url', [DEFAULT_STASH_URL])[0]
missing_type = params.get('type', ['studio'])[0]
page = int(params.get('page', [1])[0])
per_page = int(params.get('per_page', [24])[0])
query = '''
query FindMissingMetadataScenes($scene_filter: SceneFilterType, $filter: FindFilterType) {
findScenes(scene_filter: $scene_filter, filter: $filter) {
count
scenes {
id
title
date
rating100
paths {
screenshot
preview
}
files {
id
path
size
duration
video_codec
width
height
}
performers {
id
name
image_path
}
studio {
id
name
}
tags {
id
name
}
}
}
}
'''
try:
variables = {
"scene_filter": {"is_missing": missing_type},
"filter": {"per_page": per_page, "page": page}
}
res = execute_stash_graphql(query, stash_url, variables=variables, timeout=20)
data = res.get('data', {}).get('findScenes') or {}
self.send_json_response(200, {
"type": missing_type,
"count": data.get('count', 0),
"scenes": data.get('scenes', []),
"page": page,
"per_page": per_page
})
except Exception as e:
self.send_json_response(500, {"error": str(e), "scenes": [], "count": 0})
def handle_stash_batch_autotag(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
paths = payload.get('paths', [])
input_obj = {
"paths": paths if paths else None,
"performers": ["*"] if bool(payload.get('performers', True)) else None,
"studios": ["*"] if bool(payload.get('studios', True)) else None,
"tags": ["*"] if bool(payload.get('tags', True)) else None
}
input_obj = {k: v for k, v in input_obj.items() if v is not None}
mutation = build_stash_mutation("metadataAutoTag", input_obj=input_obj)
res = execute_stash_graphql(mutation, stash_url, timeout=20)
self.send_json_response(200, {"success": True, "job_id": res.get('data', {}).get('metadataAutoTag')})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
# ==================== Tag Normalizer & Deduplication Handlers ====================
def handle_stash_tag_duplicates(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
stash_url = params.get('stash_url', [DEFAULT_STASH_URL])[0]
query = '''
{
findTags(filter: { per_page: -1 }) {
count
tags {
id
name
scene_count
image_path
aliases
}
}
}
'''
try:
res = execute_stash_graphql(query, stash_url, timeout=20)
tags = res.get('data', {}).get('findTags', {}).get('tags') or []
clusters = {}
for t in tags:
raw_name = t.get('name') or ''
norm = re.sub(r'[\._\-]+', ' ', raw_name).strip().lower()
norm_compact = re.sub(r'[^a-z0-9]', '', norm)
if not norm_compact:
continue
if norm_compact not in clusters:
clusters[norm_compact] = []
clusters[norm_compact].append(t)
duplicate_clusters = []
for k, group in clusters.items():
if len(group) > 1:
sorted_group = sorted(group, key=lambda x: (x.get('scene_count') or 0), reverse=True)
duplicate_clusters.append({
"normalized_name": k,
"primary": sorted_group[0],
"candidates": sorted_group[1:],
"all": sorted_group,
"count": len(sorted_group)
})
self.send_json_response(200, {
"clusters": duplicate_clusters,
"cluster_count": len(duplicate_clusters),
"total_tags": len(tags)
})
except Exception as e:
self.send_json_response(500, {"error": str(e), "clusters": []})
def handle_stash_tag_merge(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
source_ids = [str(sid) for sid in payload.get('source_ids', [])]
dest_id = str(payload.get('destination_id'))
if not source_ids or not dest_id:
self.send_json_response(400, {"error": "Missing source_ids or destination_id"})
return
input_obj = {
"source": source_ids,
"destination": dest_id
}
mutation = build_stash_mutation("tagsMerge", input_obj=input_obj, fields="id name scene_count")
res = execute_stash_graphql(mutation, stash_url, timeout=20)
self.send_json_response(200, {"success": True, "merged_tag": res.get('data', {}).get('tagsMerge')})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
def handle_stash_tag_batch_merge(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
clusters = payload.get('clusters', [])
if not clusters:
self.send_json_response(400, {"error": "No clusters provided for batch merge"})
return
merged_count = 0
errors = []
for c in clusters:
source_ids = [str(sid) for sid in c.get('source_ids', [])]
dest_id = str(c.get('destination_id'))
if not source_ids or not dest_id:
continue
try:
mutation = build_stash_mutation("tagsMerge", input_obj={"source": source_ids, "destination": dest_id}, fields="id name")
res = execute_stash_graphql(mutation, stash_url, timeout=20)
if res.get('data', {}).get('tagsMerge'):
merged_count += 1
elif res.get('errors'):
errors.append({"destination_id": dest_id, "error": str(res.get('errors'))})
except Exception as ex:
errors.append({"destination_id": dest_id, "error": str(ex)})
self.send_json_response(200, {
"success": True,
"merged_count": merged_count,
"total_requested": len(clusters),
"errors": errors
})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
# ==================== Stash Transcode Correlator & Storage Saver Handlers ====================
def handle_stash_transcode_scan(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
stash_url = params.get('stash_url', [DEFAULT_STASH_URL])[0]
transcode_dir = params.get('transcode_dir', ['/mnt/isolation/stashapp/generated/transcodes'])[0]
stash_prefix = params.get('stash_prefix', ['/data'])[0]
local_prefix = params.get('local_prefix', ['/mnt/isolation/videos'])[0]
if not os.path.exists(transcode_dir):
self.send_json_response(400, {"error": f"Transcode directory not found: {transcode_dir}", "items": []})
return
# Build dictionary of all transcode files on disk
transcode_files = {}
for f in os.listdir(transcode_dir):
if f.endswith('.mp4'):
h = f[:-4]
full_p = os.path.join(transcode_dir, f)
try:
transcode_files[h] = {
"filename": f,
"path": full_p,
"size": os.path.getsize(full_p)
}
except OSError:
pass
# Scan Stash findFiles in batches of 5000 until all matched or end
matched_items = []
page = 1
page_size = 5000
unmatched_transcodes = set(transcode_files.keys())
while unmatched_transcodes:
q = f'''
query GetFilesBatch {{
findFiles(filter: {{ per_page: {page_size}, page: {page} }}) {{
files {{
id
path
basename
size
fingerprints {{
type
value
}}
... on VideoFile {{
width
height
duration
video_codec
}}
}}
}}
}}
'''
try:
res = execute_stash_graphql(q, stash_url, timeout=30)
files = res.get('data', {}).get('findFiles', {}).get('files', [])
if not files:
break
for f in files:
fps = {fp['type']: fp['value'] for fp in f.get('fingerprints', [])}
oshash = fps.get('oshash')
if oshash and oshash in transcode_files:
tinfo = transcode_files[oshash]
orig_stash_path = f.get('path', '')
orig_local_path = orig_stash_path.replace(stash_prefix, local_prefix) if stash_prefix and orig_stash_path.startswith(stash_prefix) else orig_stash_path
orig_size = f.get('size') or 0
transcode_size = tinfo['size']
saved_bytes = orig_size - transcode_size
pct_change = ((saved_bytes / orig_size) * 100) if orig_size > 0 else 0
matched_items.append({
"file_id": f['id'],
"oshash": oshash,
"basename": f.get('basename', ''),
"orig_stash_path": orig_stash_path,
"orig_local_path": orig_local_path,
"orig_exists": os.path.exists(orig_local_path),
"orig_size": orig_size,
"orig_resolution": f"{f.get('width', '')}x{f.get('height', '')}",
"orig_codec": f.get('video_codec', ''),
"transcode_filename": tinfo['filename'],
"transcode_path": tinfo['path'],
"transcode_size": transcode_size,
"saved_bytes": saved_bytes,
"pct_change": round(pct_change, 1)
})
unmatched_transcodes.discard(oshash)
if len(files) < page_size:
break
page += 1
except Exception as e:
print(f"Error querying Stash findFiles page {page}: {e}")
break
# Sort by space saved descending
matched_items.sort(key=lambda x: x['saved_bytes'], reverse=True)
total_saved_bytes = sum(x['saved_bytes'] for x in matched_items if x['saved_bytes'] > 0)
self.send_json_response(200, {
"items": matched_items,
"total_matched": len(matched_items),
"total_transcodes_on_disk": len(transcode_files),
"total_potential_savings_bytes": total_saved_bytes
})
def handle_stash_transcode_replace(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
orig_local_path = payload.get('orig_local_path')
transcode_path = payload.get('transcode_path')
delete_transcode = bool(payload.get('delete_transcode', True))
trigger_rescan = bool(payload.get('trigger_rescan', True))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
if not orig_local_path or not transcode_path:
self.send_json_response(400, {"error": "Missing orig_local_path or transcode_path"})
return
if not os.path.exists(transcode_path):
self.send_json_response(400, {"error": f"Transcode file does not exist: {transcode_path}"})
return
if not os.path.exists(orig_local_path):
self.send_json_response(400, {"error": f"Original file does not exist at local path: {orig_local_path}"})
return
orig_size = os.path.getsize(orig_local_path)
transcode_size = os.path.getsize(transcode_path)
if transcode_size == 0:
self.send_json_response(400, {"error": "Transcode file is 0 bytes (corrupted), aborting replacement"})
return
# Step 1: Copy transcode to temporary staging file next to original
tmp_staging_path = f"{orig_local_path}.transcode.tmp"
if os.path.exists(tmp_staging_path):
try:
os.remove(tmp_staging_path)
except OSError:
pass
shutil.copy2(transcode_path, tmp_staging_path)
# Step 2: Verify size of staging file
staged_size = os.path.getsize(tmp_staging_path)
if staged_size != transcode_size:
if os.path.exists(tmp_staging_path):
os.remove(tmp_staging_path)
self.send_json_response(500, {"error": f"Size mismatch during staging: expected {transcode_size}, got {staged_size}"})
return
# Step 3: Atomically overwrite original file with staged transcode
os.replace(tmp_staging_path, orig_local_path)
# Step 4: If delete_transcode is requested, remove original transcode cache
if delete_transcode:
try:
os.remove(transcode_path)
except OSError as e:
print(f"Warning: Could not remove transcode cache file {transcode_path}: {e}")
saved_bytes = orig_size - transcode_size
# Step 5: Record to Relocation History audit log
try:
record_relocation_history(
source_type="stash_transcode",
item_key=orig_local_path,
item_name=os.path.basename(orig_local_path),
source_path=transcode_path,
target_path=orig_local_path,
performer="Transcode Optimization",
file_size=transcode_size
)
except Exception as e:
print(f"Warning: History record error: {e}")
# Step 6: Trigger Stash metadata scan on parent folder if requested
if trigger_rescan:
try:
parent_dir = os.path.dirname(orig_local_path)
stash_parent = parent_dir.replace('/mnt/isolation/videos', '/data')
mutation = build_stash_mutation("metadataScan", input_obj={"paths": [stash_parent], "rescan": True})
execute_stash_graphql(mutation, stash_url, timeout=10)
except Exception as e:
print(f"Background Stash rescan trigger error: {e}")
self.send_json_response(200, {
"success": True,
"orig_local_path": orig_local_path,
"orig_size": orig_size,
"new_size": transcode_size,
"saved_bytes": saved_bytes
})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
def handle_stash_transcode_batch_replace(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
items = payload.get('items', [])
delete_transcode = bool(payload.get('delete_transcode', True))
trigger_rescan = bool(payload.get('trigger_rescan', True))
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
replaced_count = 0
total_saved_bytes = 0
errors = []
rescanned_dirs = set()
for item in items:
orig_local_path = item.get('orig_local_path')
transcode_path = item.get('transcode_path')
if not orig_local_path or not transcode_path or not os.path.exists(transcode_path) or not os.path.exists(orig_local_path):
continue
try:
orig_size = os.path.getsize(orig_local_path)
transcode_size = os.path.getsize(transcode_path)
if transcode_size == 0:
continue
tmp_staging_path = f"{orig_local_path}.transcode.tmp"
shutil.copy2(transcode_path, tmp_staging_path)
if os.path.getsize(tmp_staging_path) != transcode_size:
if os.path.exists(tmp_staging_path):
os.remove(tmp_staging_path)
continue
os.replace(tmp_staging_path, orig_local_path)
if delete_transcode:
try:
os.remove(transcode_path)
except OSError:
pass
replaced_count += 1
saved = orig_size - transcode_size
total_saved_bytes += saved
rescanned_dirs.add(os.path.dirname(orig_local_path))
try:
record_relocation_history(
source_type="stash_transcode",
item_key=orig_local_path,
item_name=os.path.basename(orig_local_path),
source_path=transcode_path,
target_path=orig_local_path,
performer="Transcode Optimization",
file_size=transcode_size
)
except Exception:
pass
except Exception as e:
errors.append({"path": orig_local_path, "error": str(e)})
# Trigger Stash scan for all touched directories
if trigger_rescan and rescanned_dirs:
try:
stash_dirs = [d.replace('/mnt/isolation/videos', '/data') for d in rescanned_dirs]
mutation = build_stash_mutation("metadataScan", input_obj={"paths": stash_dirs, "rescan": True})
execute_stash_graphql(mutation, stash_url, timeout=15)
except Exception as e:
print(f"Batch Stash rescan trigger error: {e}")
self.send_json_response(200, {
"success": True,
"replaced_count": replaced_count,
"total_saved_bytes": total_saved_bytes,
"errors": errors
})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
# ==================== SMB Maintenance & Cleanup Handlers ====================
def handle_smb_scan_cleanup(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
scan_dir = params.get('path', [DEFAULT_SMB_SOURCE])[0]
ensure_smb_mounted()
if not os.path.exists(scan_dir):
self.send_json_response(404, {"error": f"Directory does not exist: {scan_dir}", "zero_byte_files": [], "empty_dirs": []})
return
zero_byte = []
empty_dirs = []
try:
for root, dirs, files in os.walk(scan_dir, topdown=False):
for f in files:
fp = os.path.join(root, f)
try:
if os.path.isfile(fp) and os.path.getsize(fp) == 0:
zero_byte.append({"name": f, "path": fp, "size": 0})
except Exception:
pass
# Check empty directories
if root != scan_dir and not os.listdir(root):
empty_dirs.append({"name": os.path.basename(root), "path": root})
self.send_json_response(200, {
"target_dir": scan_dir,
"zero_byte_files": zero_byte,
"empty_dirs": empty_dirs,
"zero_byte_count": len(zero_byte),
"empty_dir_count": len(empty_dirs)
})
except Exception as e:
self.send_json_response(500, {"error": str(e), "zero_byte_files": [], "empty_dirs": []})
def handle_smb_execute_cleanup(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
files_to_delete = payload.get('files', [])
dirs_to_prune = payload.get('dirs', [])
deleted_files = 0
pruned_dirs = 0
for fp in files_to_delete:
try:
if os.path.isfile(fp):
os.remove(fp)
deleted_files += 1
except Exception:
pass
for dp in dirs_to_prune:
try:
if os.path.isdir(dp) and not os.listdir(dp):
os.rmdir(dp)
pruned_dirs += 1
except Exception:
pass
self.send_json_response(200, {
"success": True,
"deleted_files": deleted_files,
"pruned_dirs": pruned_dirs
})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
# ==================== Storage Analytics & Heatmap ====================
def handle_storage_analytics(self):
parsed = urllib.parse.urlparse(self.path)
params = urllib.parse.parse_qs(parsed.query)
target_root = params.get('path', [DEFAULT_SMB_TARGET_ROOT])[0]
ensure_smb_mounted()
if not os.path.exists(target_root):
self.send_json_response(404, {"error": f"Target root not accessible: {target_root}", "performers": []})
return
performer_stats = []
total_library_size = 0
total_video_count = 0
try:
entries = os.listdir(target_root)
for entry in entries:
full_p = os.path.join(target_root, entry)
if os.path.isdir(full_p) and not entry.startswith('.'):
f_count = 0
f_bytes = 0
for root, _, files in os.walk(full_p):
for f in files:
ext = os.path.splitext(f)[1].lower()
if ext in VIDEO_EXTENSIONS:
f_count += 1
try:
f_bytes += os.path.getsize(os.path.join(root, f))
except Exception:
pass
if f_count > 0:
total_video_count += f_count
total_library_size += f_bytes
performer_stats.append({
"name": entry,
"folder": entry,
"path": full_p,
"video_count": f_count,
"total_size": f_bytes,
"formatted_size": f"{f_bytes / (1024**3):.2f} GB" if f_bytes >= 1024**3 else f"{f_bytes / (1024**2):.1f} MB"
})
# Sort largest folders first
performer_stats.sort(key=lambda x: -x['total_size'])
# Free disk space check
disk_info = {}
try:
st = os.statvfs(target_root)
disk_info = {
"total": st.f_blocks * st.f_frsize,
"free": st.f_bavail * st.f_frsize,
"used": (st.f_blocks - st.f_bavail) * st.f_frsize,
"free_gb": f"{(st.f_bavail * st.f_frsize) / (1024**3):.1f} GB",
"total_gb": f"{(st.f_blocks * st.f_frsize) / (1024**3):.1f} GB"
}
except Exception:
pass
self.send_json_response(200, {
"target_root": target_root,
"total_performers": len(performer_stats),
"total_videos": total_video_count,
"total_size_bytes": total_library_size,
"total_size_formatted": f"{total_library_size / (1024**3):.2f} GB",
"performers": performer_stats[:100],
"disk": disk_info
})
except Exception as e:
self.send_json_response(500, {"error": str(e), "performers": []})
# ==================== Webhook Dispatcher ====================
def handle_send_webhook(self):
content_length = int(self.headers.get('Content-Length', 0))
post_data = self.rfile.read(content_length)
try:
payload = json.loads(post_data.decode('utf-8'))
webhook_url = payload.get('webhook_url')
title = payload.get('title', 'Performer Organizer Notification')
message = payload.get('message', '')
items_count = payload.get('items_count', 0)
if not webhook_url:
self.send_json_response(400, {"error": "Missing webhook_url"})
return
body = {
"content": f"🚀 **{title}**\n{message}\n*Organized {items_count} items into library.*",
"username": "qBittorrent Organizer"
}
req = urllib.request.Request(
webhook_url,
data=json.dumps(body).encode('utf-8'),
headers={'Content-Type': 'application/json', 'User-Agent': 'Mozilla/5.0'}
)
with urllib.request.urlopen(req, timeout=10) as res:
self.send_json_response(200, {"success": True, "status": res.status})
except Exception as e:
self.send_json_response(500, {"error": str(e)})
self.end_headers()
self.wfile.write(json.dumps({"error": f"Delete request failed: {str(e)}"}).encode('utf-8'))
# ==================== External API Proxies ====================
def proxy_external_api(self, upstream_base, proxy_prefix, method='POST'):
upstream_base = upstream_base.rstrip('/')
target_subpath = self.path[len(proxy_prefix):]
target_url = f"{upstream_base}/{target_subpath.lstrip('/')}"
body = None
content_length = int(self.headers.get('Content-Length', 0))
if content_length > 0:
body = self.rfile.read(content_length)
req = urllib.request.Request(target_url, data=body, method=method)
for header, val in self.headers.items():
h_lower = header.lower()
if h_lower in ['content-type', 'authorization', 'x-api-key', 'anthropic-version', 'anthropic-beta']:
req.add_header(header, val)
try:
with urllib.request.urlopen(req, timeout=90) as resp:
resp_data = resp.read()
self.send_response(resp.status)
self.send_cors_headers()
self.send_header('Content-Type', resp.headers.get('Content-Type', 'application/json'))
self.end_headers()
self.wfile.write(resp_data)
except urllib.error.HTTPError as e:
err_data = e.read()
self.send_response(e.code)
self.send_cors_headers()
self.send_header('Content-Type', 'application/json')
self.end_headers()
self.wfile.write(err_data)
except Exception as e:
self.send_response(502)
self.send_cors_headers()
self.send_header('Content-Type', 'application/json')
self.end_headers()
self.wfile.write(json.dumps({"error": f"Error connecting to LLM endpoint at {target_url}: {str(e)}"}).encode('utf-8'))
def proxy_qb_request(self, method):
global latest_sid
raw_target = self.headers.get('X-QB-Target-Host') or DEFAULT_QB_HOST
target_base = normalize_host(raw_target)
# Target Endpoint
target_subpath = self.path[len('/qb-proxy/'):]
target_url = f"{target_base}/{target_subpath.lstrip('/')}"
# Read POST body
body = None
content_length = int(self.headers.get('Content-Length', 0))
if content_length > 0:
body = self.rfile.read(content_length)
# Build Request
req = urllib.request.Request(target_url, data=body, method=method)
# Forward headers from browser
browser_cookie = self.headers.get('Cookie', '')
for header, val in self.headers.items():
h_lower = header.lower()
if h_lower not in ['host', 'origin', 'referer', 'content-length', 'x-qb-target-host', 'cookie', 'accept-encoding']:
req.add_header(header, val)
req.add_header('Accept-Encoding', 'identity')
# Attach SID cookie
if 'SID=' in browser_cookie:
req.add_header('Cookie', browser_cookie)
elif latest_sid:
req.add_header('Cookie', f"SID={latest_sid}")
req.add_header('Referer', target_base)
req.add_header('Origin', target_base)
try:
with urllib.request.urlopen(req, timeout=25) as resp:
status_code = resp.status
resp_data = resp.read()
resp_headers = resp.getheaders()
# Handle gzip decompression if returned
encoding = resp.headers.get('Content-Encoding', '').lower()
if encoding == 'gzip' or resp_data.startswith(b'\x1f\x8b'):
try:
resp_data = gzip.decompress(resp_data)
except Exception as gz_err:
print(f"Gzip decompress error: {gz_err}")
# Capture SID cookie from login response
for header, val in resp_headers:
if header.lower() == 'set-cookie':
match = re.search(r'SID=([^;]+)', val)
if match:
latest_sid = match.group(1)
self.send_response(status_code)
self.send_cors_headers()
content_type = resp.headers.get('Content-Type', 'application/json')
self.send_header('Content-Type', content_type)
self.send_header('Content-Length', str(len(resp_data)))
for header, val in resp_headers:
if header.lower() == 'set-cookie':
self.send_header(header, val)
self.end_headers()
self.wfile.write(resp_data)
except urllib.error.HTTPError as e:
err_body = e.read()
if e.headers.get('Content-Encoding', '').lower() == 'gzip' or err_body.startswith(b'\x1f\x8b'):
try:
err_body = gzip.decompress(err_body)
except Exception:
pass
err_text = err_body.decode('utf-8', errors='ignore')
self.send_response(e.code)
self.send_cors_headers()
self.send_header('Content-Type', 'application/json')
self.end_headers()
err_json = json.dumps({"error": f"qBittorrent error ({e.code}): {err_text}", "code": e.code})
self.wfile.write(err_json.encode('utf-8'))
except Exception as e:
print(f"Proxy connection error: {e}")
self.send_response(502)
self.send_cors_headers()
self.send_header('Content-Type', 'application/json')
self.end_headers()
err_json = json.dumps({"error": f"Failed to connect to qBittorrent at {target_url}: {str(e)}"})
self.wfile.write(err_json.encode('utf-8'))
def main():
os.chdir(os.path.dirname(os.path.abspath(__file__)))
server_address = ('0.0.0.0', PORT)
httpd = http.server.ThreadingHTTPServer(server_address, RobustProxyHandler)
print(f"===========================================================")
print(f"🚀 qBittorrent + SMB Direct Sorter Server running!")
print(f"👉 Open in browser: http://localhost:{PORT}")
print(f"📁 SMB Target: {DEFAULT_SMB_SOURCE}")
print(f"💾 SQLite Database: {DB_PATH}")
print(f"===========================================================")
try:
httpd.serve_forever()
except KeyboardInterrupt:
print("\nStopping server...")
httpd.server_close()
if __name__ == '__main__':
main()