Initial commit: qBittorrent & SMB AI Video Performer Sorter with in-browser streaming preview, Intellisense & SQLite persistence
This commit is contained in:
@@ -0,0 +1,843 @@
|
||||
#!/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
|
||||
|
||||
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()
|
||||
# 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),
|
||||
"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()
|
||||
# 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}")
|
||||
|
||||
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'))
|
||||
|
||||
# ==================== 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()
|
||||
Reference in New Issue
Block a user