Files
qbittorrent-performer-sorter/server.py
T

1245 lines
50 KiB
Python

#!/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 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)')
conn.commit()
conn.close()
init_db()
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'''<svg xmlns="http://www.w3.org/2000/svg" viewBox="0 0 160 200" width="160" height="200">
<defs>
<linearGradient id="g" x1="0%" y1="0%" x2="100%" y2="100%">
<stop offset="0%" stop-color="#1e1b4b"/>
<stop offset="50%" stop-color="#312e81"/>
<stop offset="100%" stop-color="#4338ca"/>
</linearGradient>
</defs>
<rect width="100%" height="100%" fill="url(#g)" rx="12"/>
<circle cx="80" cy="72" r="32" fill="#6366f1" opacity="0.85"/>
<circle cx="80" cy="62" r="18" fill="#e0e7ff"/>
<path d="M48 122 Q80 94 112 122 Z" fill="#e0e7ff"/>
<text x="80" y="162" font-family="system-ui, -apple-system, sans-serif" font-size="13" font-weight="600" fill="#ffffff" text-anchor="middle">{clean}</text>
</svg>'''
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, timeout=4):
endpoint = normalize_stash_url(stash_url)
req = urllib.request.Request(
endpoint,
data=json.dumps({'query': query_str}).encode('utf-8'),
headers={'Content-Type': 'application/json'}
)
with urllib.request.urlopen(req, timeout=timeout) as res:
return json.loads(res.read().decode())
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 query_stash_single(file_name, file_path=None, item_type='smb', stash_url=None):
basename = os.path.basename(file_path or file_name)
stash_paths = [
f'/data/videodownloader/{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 {{
name
gender
}}
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. Basename search in Stash
stem = os.path.splitext(basename)[0]
q = f'''
{{
findScenes(filter: {{q: {json.dumps(stem[:35])}, per_page: 3}}) {{
scenes {{
id
title
files {{
basename
}}
performers {{
name
gender
}}
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, 'fuzzy_match')
if parsed:
return parsed
except Exception:
pass
return None
def query_stash_batch(items, stash_url=None):
if not items:
return []
queries = []
for idx, item in enumerate(items):
file_path = item.get('path') or item.get('name')
basename = os.path.basename(file_path)
p = f'/data/videodownloader/{basename}'
queries.append(f'item_{idx}: findScenes(scene_filter: {{path: {{value: {json.dumps(p)}, modifier: EQUALS}}}}) {{ scenes {{ id title performers {{ name gender }} studio {{ name }} }} }}')
batch_q = '{\n' + '\n'.join(queries) + '\n}'
results = [None] * len(items)
unmatched_indices = []
try:
res = execute_stash_graphql(batch_q, stash_url, timeout=6)
data = res.get('data', {})
for idx in range(len(items)):
scenes = data.get(f'item_{idx}', {}).get('scenes', [])
if scenes and scenes[0].get('performers'):
parsed = parse_stash_scene_data(scenes[0], 'exact_path')
if parsed:
results[idx] = parsed
if not results[idx]:
unmatched_indices.append(idx)
except Exception as e:
print(f"Batch exact Stash query failed: {e}")
unmatched_indices = list(range(len(items)))
# Second pass for unmatched items
for idx in unmatched_indices:
item = items[idx]
file_name = item.get('name', '')
file_path = item.get('path', '')
item_type = item.get('item_type', 'smb')
match = query_stash_single(file_name, file_path, item_type, stash_url)
results[idx] = match
return results
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"
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 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()
# 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()
# Stash Status API
elif self.path.startswith('/api/stash/status'):
self.handle_stash_status()
# 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':
self.send_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
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.wfile.write(json.dumps(status_data).encode('utf-8'))
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 APIs
if self.path == '/api/db/save':
self.handle_db_save()
elif self.path == '/api/db/batch-save':
self.handle_db_batch_save()
# Stash Lookup APIs
elif self.path == '/api/stash/lookup':
self.handle_stash_lookup()
elif self.path == '/api/stash/batch-lookup':
self.handle_stash_batch_lookup()
# SMB API Moves
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.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_response(200)
self.send_header('Content-Type', 'application/json')
self.send_cors_headers()
self.end_headers()
self.wfile.write(json.dumps({"cache": cache_map, "count": len(cache_map)}).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 SQLite cache: {str(e)}"}).encode('utf-8'))
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'))
# ==================== 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)
file_path = params.get('path', [''])[0]
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 != 'videodownloader' 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
# ==================== 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'))
# ==================== 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')
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)
# Move file
shutil.move(source_path, target_file)
# 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')
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)
shutil.move(source, target_file)
# 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'))
# ==================== 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()