1471 lines
60 KiB
Python
1471 lines
60 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 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 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('/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/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()
|
|
# 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':
|
|
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()
|
|
# 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 & Deletes
|
|
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.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 != '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'))
|
|
|
|
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 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()
|