837 lines
36 KiB
Python
837 lines
36 KiB
Python
from flask import Flask, request, jsonify, render_template, send_file
|
|
import os
|
|
import sqlite3
|
|
from collections import defaultdict
|
|
import smbclient
|
|
import tempfile
|
|
import json
|
|
import asyncio
|
|
import time
|
|
import threading
|
|
import mimetypes
|
|
import io
|
|
import zipfile
|
|
import base64
|
|
import qbittorrentapi
|
|
import requests
|
|
from bs4 import BeautifulSoup
|
|
import subprocess
|
|
|
|
# --- Threading Lock ---
|
|
REMOTE_DB_LOCK = threading.Lock()
|
|
|
|
# --- App Initialization & Path Configuration ---
|
|
# Explicitly define paths relative to this file's location to avoid ambiguity.
|
|
# .../stashtoolkit_webtop/backend/app.py
|
|
APP_ROOT = os.path.dirname(os.path.abspath(__file__))
|
|
# .../stashtoolkit_webtop/instance/
|
|
INSTANCE_PATH = os.path.join(os.path.dirname(APP_ROOT), 'instance')
|
|
|
|
app = Flask(__name__, template_folder='../static')
|
|
|
|
# --- Database Setup ---
|
|
LOCAL_DB_PATH = "local_cache.db"
|
|
CONFIG_FILE = "config.json"
|
|
|
|
def get_local_db():
|
|
db_path = os.path.join(INSTANCE_PATH, LOCAL_DB_PATH)
|
|
conn = sqlite3.connect(db_path)
|
|
conn.row_factory = sqlite3.Row
|
|
return conn
|
|
|
|
def init_local_db():
|
|
"""Initializes and resets cache tables to ensure schema is correct."""
|
|
try:
|
|
os.makedirs(INSTANCE_PATH)
|
|
except OSError:
|
|
pass # Already exists
|
|
with get_local_db() as conn:
|
|
cursor = conn.cursor()
|
|
cursor.execute("CREATE TABLE IF NOT EXISTS config (key TEXT PRIMARY KEY, value TEXT)")
|
|
cursor.execute("CREATE TABLE IF NOT EXISTS sync_metadata (key TEXT PRIMARY KEY, value TEXT)")
|
|
cursor.execute("""
|
|
CREATE TABLE IF NOT EXISTS scan_history (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
scan_type TEXT NOT NULL,
|
|
timestamp DATETIME DEFAULT CURRENT_TIMESTAMP,
|
|
status TEXT NOT NULL,
|
|
message TEXT,
|
|
log TEXT
|
|
)""")
|
|
|
|
# Check if log column exists and add it if not (for existing databases)
|
|
cursor.execute("PRAGMA table_info(scan_history)")
|
|
columns = [row['name'] for row in cursor.fetchall()]
|
|
if 'log' not in columns:
|
|
cursor.execute("ALTER TABLE scan_history ADD COLUMN log TEXT")
|
|
|
|
cursor.execute("""
|
|
CREATE TABLE IF NOT EXISTS duplicate_results (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
scan_id INTEGER NOT NULL,
|
|
file_path TEXT NOT NULL,
|
|
file_size INTEGER,
|
|
file_basename TEXT,
|
|
set_id INTEGER NOT NULL,
|
|
FOREIGN KEY (scan_id) REFERENCES scan_history (id) ON DELETE CASCADE
|
|
)""")
|
|
cursor.execute("""
|
|
CREATE TABLE IF NOT EXISTS transcode_plan_results (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
scan_id INTEGER NOT NULL,
|
|
original_path TEXT NOT NULL,
|
|
transcoded_path TEXT NOT NULL,
|
|
FOREIGN KEY (scan_id) REFERENCES scan_history (id) ON DELETE CASCADE
|
|
)""")
|
|
|
|
# Drop and recreate cache tables to apply schema changes cleanly
|
|
print("Rebuilding local cache tables to ensure schema is up-to-date.")
|
|
cursor.execute("DROP TABLE IF EXISTS scenes")
|
|
cursor.execute("DROP TABLE IF EXISTS files")
|
|
cursor.execute("DROP TABLE IF EXISTS paths")
|
|
cursor.execute("CREATE TABLE scenes (id INTEGER PRIMARY KEY, oshash TEXT)")
|
|
cursor.execute("CREATE TABLE files (scene_id INTEGER, path TEXT, basename TEXT, path_id INTEGER)")
|
|
cursor.execute("CREATE TABLE paths (id INTEGER PRIMARY KEY, path TEXT)")
|
|
conn.commit()
|
|
|
|
# --- Global Variables & Constants ---
|
|
credentials = {}
|
|
DB_PATH = "//servervm.local/main/appdata/stashapp/config/stash-go.sqlite"
|
|
VIDEO_PATH = "//truenas.local/isolation"
|
|
TRANSCODES_PATH = "//truenas.local/isolation/stashapp/generated/transcodes"
|
|
COOMER_PATH = "//truenas.local/isolation/coomer"
|
|
PATH_MAPPING_FROM = "/data/"
|
|
PATH_MAPPING_TO = "//truenas.local/isolation/videos/"
|
|
|
|
# --- Credential and Configuration ---
|
|
def load_credentials():
|
|
global credentials
|
|
with app.app_context(), get_local_db() as conn:
|
|
row = conn.execute("SELECT value FROM config WHERE key = 'credentials'").fetchone()
|
|
if row and row['value']:
|
|
credentials = json.loads(row['value'])
|
|
else:
|
|
# os.path.dirname(APP_ROOT) is .../stashtoolkit_webtop/
|
|
old_config_path = os.path.join(os.path.dirname(APP_ROOT), CONFIG_FILE)
|
|
if os.path.exists(old_config_path):
|
|
print("Migrating credentials from config.json...")
|
|
try:
|
|
with open(old_config_path, "r") as f:
|
|
credentials = json.load(f)
|
|
conn.execute("REPLACE INTO config (key, value) VALUES (?, ?)", ('credentials', json.dumps(credentials)))
|
|
conn.commit()
|
|
os.rename(old_config_path, old_config_path + ".migrated")
|
|
print("Successfully migrated credentials.")
|
|
except (json.JSONDecodeError, OSError) as e:
|
|
print(f"Error migrating from config.json: {e}")
|
|
credentials = {}
|
|
|
|
def get_smb_credentials(path):
|
|
server = path.split("/")[2].split('@')[-1]
|
|
return credentials.get(server, {})
|
|
|
|
def get_file_size_smb(path):
|
|
try:
|
|
creds = get_smb_credentials(path)
|
|
with smbclient.open_file(path, mode='rb', **creds) as f:
|
|
return f.seek(0, 2)
|
|
except Exception:
|
|
return -1
|
|
|
|
def log_to_scan(scan_id, message):
|
|
with app.app_context():
|
|
with get_local_db() as conn:
|
|
# Append message to log, creating a newline if log is not empty
|
|
conn.execute(
|
|
"UPDATE scan_history SET log = CASE WHEN log IS NULL OR log = '' THEN ? ELSE log || '\n' || ? END WHERE id = ?",
|
|
(message, message, scan_id)
|
|
)
|
|
conn.commit()
|
|
|
|
def get_qb_client():
|
|
if 'qbittorrent' in credentials:
|
|
creds = credentials['qbittorrent']
|
|
return qbittorrentapi.Client(
|
|
host=creds.get('host'),
|
|
port=creds.get('port'),
|
|
username=creds.get('username'),
|
|
password=creds.get('password')
|
|
)
|
|
return None
|
|
|
|
|
|
# --- Remote Database Sync ---
|
|
def sync_remote_db():
|
|
print("Background sync thread started.")
|
|
print("Sync thread trying to acquire DB lock...")
|
|
with REMOTE_DB_LOCK:
|
|
print("Sync thread acquired DB lock.")
|
|
try:
|
|
creds = get_smb_credentials(DB_PATH)
|
|
remote_stat = smbclient.stat(DB_PATH, **creds)
|
|
remote_mtime = remote_stat.st_mtime
|
|
except Exception as e:
|
|
print(f"Sync thread failed to check remote DB status: {e}")
|
|
return
|
|
|
|
with app.app_context(), get_local_db() as conn_local:
|
|
row = conn_local.execute("SELECT value FROM sync_metadata WHERE key = 'last_sync_mtime'").fetchone()
|
|
last_sync_mtime = float(row['value']) if (row and row['value']) else 0
|
|
|
|
if remote_mtime <= last_sync_mtime:
|
|
print("Local database is already up to date.")
|
|
return
|
|
|
|
print(f"Remote DB is newer. Syncing from {remote_mtime} > {last_sync_mtime}.")
|
|
temp_db_path = None
|
|
try:
|
|
with smbclient.open_file(DB_PATH, mode='rb', **creds) as smb_file:
|
|
with tempfile.NamedTemporaryFile(delete=False, suffix=".sqlite") as temp_db:
|
|
temp_db.write(smb_file.read())
|
|
temp_db_path = temp_db.name
|
|
|
|
with sqlite3.connect(temp_db_path) as conn_remote:
|
|
conn_local.execute("DELETE FROM scenes"); conn_local.execute("DELETE FROM files"); conn_local.execute("DELETE FROM paths")
|
|
|
|
cursor_remote_scenes = conn_remote.execute("SELECT id, oshash FROM scenes")
|
|
conn_local.executemany("INSERT INTO scenes (id, oshash) VALUES (?, ?)", cursor_remote_scenes)
|
|
|
|
cursor_remote_files = conn_remote.execute("SELECT scene_id, path, basename, path_id FROM files")
|
|
conn_local.executemany("INSERT INTO files (scene_id, path, basename, path_id) VALUES (?, ?, ?, ?)", cursor_remote_files)
|
|
|
|
cursor_remote_paths = conn_remote.execute("SELECT id, path FROM paths")
|
|
conn_local.executemany("INSERT INTO paths (id, path) VALUES (?, ?)", cursor_remote_paths)
|
|
|
|
conn_local.execute("REPLACE INTO sync_metadata (key, value) VALUES (?, ?)", ('last_sync_mtime', remote_mtime))
|
|
conn_local.commit()
|
|
print("Database sync completed successfully.")
|
|
except sqlite3.Error as e:
|
|
print(f"Error during database sync: {e}")
|
|
conn_local.rollback()
|
|
finally:
|
|
if temp_db_path and os.path.exists(temp_db_path): os.remove(temp_db_path)
|
|
print("Sync thread released DB lock.")
|
|
|
|
# --- Initialize App & Background Tasks ---
|
|
with app.app_context():
|
|
init_local_db()
|
|
load_credentials()
|
|
|
|
if os.environ.get("WERKZEUG_RUN_MAIN") != "true":
|
|
threading.Thread(target=sync_remote_db, daemon=True).start()
|
|
|
|
# --- API Endpoints ---
|
|
@app.route('/')
|
|
def index(): return render_template('index.html')
|
|
|
|
@app.route('/api/save_credentials', methods=['POST'])
|
|
def save_credentials():
|
|
data = request.get_json()
|
|
with app.app_context(), get_local_db() as conn:
|
|
conn.execute("REPLACE INTO config (key, value) VALUES (?, ?)", ('credentials', json.dumps(data)))
|
|
global credentials
|
|
credentials = data
|
|
return jsonify({"success": True})
|
|
|
|
@app.route('/api/get_credentials', methods=['GET'])
|
|
def get_credentials():
|
|
with app.app_context(), get_local_db() as conn:
|
|
row = conn.execute("SELECT value FROM config WHERE key = 'credentials'").fetchone()
|
|
if row and row['value']:
|
|
return jsonify({"credentials": json.loads(row['value'])})
|
|
return jsonify({"credentials": {}})
|
|
|
|
@app.route('/api/check_credentials', methods=['GET'])
|
|
def check_credentials():
|
|
with app.app_context(), get_local_db() as conn:
|
|
row = conn.execute("SELECT value FROM config WHERE key = 'credentials'").fetchone()
|
|
return jsonify({"has_credentials": bool(row and row['value'] and row['value'] != '{}')})
|
|
|
|
@app.route('/api/check_connections', methods=['POST'])
|
|
async def check_connections():
|
|
async def run_sync(func, *args): return await asyncio.to_thread(func, *args)
|
|
def check_db_connection():
|
|
try:
|
|
print("Health check trying to acquire DB lock...")
|
|
with REMOTE_DB_LOCK:
|
|
print("Health check acquired DB lock.")
|
|
creds = get_smb_credentials(DB_PATH)
|
|
with smbclient.open_file(DB_PATH, mode='rb', **creds) as smb_file:
|
|
header = smb_file.read(100)
|
|
print("Health check released DB lock.")
|
|
if not header.startswith(b'SQLite format 3\x00'):
|
|
return 'stash_db', {"success": False, "message": "File is not a valid SQLite3 database."}
|
|
return 'stash_db', {"success": True, "message": "Successfully connected to the database."}
|
|
except Exception as e:
|
|
print(f"Health check failed: {e}")
|
|
return 'stash_db', {"success": False, "message": str(e)}
|
|
def check_truenas_connection():
|
|
try:
|
|
creds = get_smb_credentials(VIDEO_PATH)
|
|
it = smbclient.scandir(VIDEO_PATH, **creds)
|
|
try: next(it)
|
|
except StopIteration: pass
|
|
return 'truenas', {"success": True, "message": "Successfully connected to truenas.local."}
|
|
except Exception as e:
|
|
return 'truenas', {"success": False, "message": str(e)}
|
|
|
|
def check_qb_connection():
|
|
try:
|
|
client = get_qb_client()
|
|
if not client:
|
|
return 'qbittorrent', {"success": False, "message": "qBittorrent credentials not set."}
|
|
client.auth_log_in()
|
|
return 'qbittorrent', {"success": True, "message": f"Connected to qBittorrent v{client.app.version}"}
|
|
except Exception as e:
|
|
return 'qbittorrent', {"success": False, "message": str(e)}
|
|
|
|
results = await asyncio.gather(run_sync(check_db_connection), run_sync(check_truenas_connection), run_sync(check_qb_connection))
|
|
results_dict = dict(results)
|
|
if 'stash_db' in results_dict:
|
|
results_dict['servervm'] = {"success": results_dict['stash_db']['success'], "message": results_dict['stash_db']['message']}
|
|
return jsonify(results_dict)
|
|
|
|
def do_find_duplicates_task(app_context, scan_id):
|
|
with app_context:
|
|
log_to_scan(scan_id, "Starting duplicate scan...")
|
|
try:
|
|
with get_local_db() as conn:
|
|
# Find oshash values that are associated with more than one scene
|
|
log_to_scan(scan_id, "Finding duplicate hashes...")
|
|
duplicate_oshashes_query = """
|
|
SELECT oshash FROM scenes
|
|
WHERE oshash IS NOT NULL AND oshash != ''
|
|
GROUP BY oshash
|
|
HAVING COUNT(id) > 1
|
|
"""
|
|
|
|
cursor = conn.execute(duplicate_oshashes_query)
|
|
duplicate_oshashes = [row['oshash'] for row in cursor.fetchall()]
|
|
|
|
if not duplicate_oshashes:
|
|
log_to_scan(scan_id, "No duplicate hashes found.")
|
|
conn.execute("UPDATE scan_history SET status = 'completed' WHERE id = ?", (scan_id,))
|
|
conn.commit()
|
|
log_to_scan(scan_id, "Scan completed.")
|
|
return
|
|
|
|
log_to_scan(scan_id, f"Found {len(duplicate_oshashes)} sets of duplicates. Fetching file details...")
|
|
set_id_counter = 0
|
|
for oshash in duplicate_oshashes:
|
|
set_id_counter += 1
|
|
files_query = """
|
|
SELECT p.path || '/' || f.path as full_path, f.basename
|
|
FROM files f
|
|
JOIN scenes s ON s.id = f.scene_id
|
|
JOIN paths p ON p.id = f.path_id
|
|
WHERE s.oshash = ?
|
|
"""
|
|
|
|
cursor = conn.execute(files_query, (oshash,))
|
|
files = cursor.fetchall()
|
|
|
|
log_to_scan(scan_id, f"Processing set {set_id_counter} with {len(files)} files...")
|
|
for file_row in files:
|
|
full_path_smb = file_row['full_path'].replace(PATH_MAPPING_FROM, PATH_MAPPING_TO)
|
|
size = get_file_size_smb(full_path_smb)
|
|
conn.execute(
|
|
"INSERT INTO duplicate_results (scan_id, file_path, file_size, file_basename, set_id) VALUES (?, ?, ?, ?, ?)",
|
|
(scan_id, full_path_smb, size, file_row['basename'], set_id_counter)
|
|
)
|
|
|
|
conn.execute("UPDATE scan_history SET status = 'completed' WHERE id = ?", (scan_id,))
|
|
conn.commit()
|
|
log_to_scan(scan_id, "Duplicate scan completed successfully.")
|
|
except Exception as e:
|
|
error_message = f"Error in duplicate scan: {e}"
|
|
log_to_scan(scan_id, error_message)
|
|
with app.app_context(), get_local_db() as conn:
|
|
conn.execute("UPDATE scan_history SET status = 'error', message = ? WHERE id = ?", (str(e), scan_id))
|
|
conn.commit()
|
|
|
|
@app.route('/api/find_duplicates', methods=['POST'])
|
|
def find_duplicates():
|
|
with get_local_db() as conn:
|
|
cursor = conn.execute("INSERT INTO scan_history (scan_type, status) VALUES (?, ?)", ('duplicates', 'in_progress'))
|
|
scan_id = cursor.lastrowid
|
|
threading.Thread(target=do_find_duplicates_task, args=(app.app_context(), scan_id), daemon=True).start()
|
|
return jsonify({"scan_id": scan_id})
|
|
|
|
def do_transcode_plan_task(app_context, scan_id):
|
|
with app_context:
|
|
log_to_scan(scan_id, "Starting transcode replacement plan...")
|
|
try:
|
|
log_to_scan(scan_id, "Finding transcoded files...")
|
|
creds = get_smb_credentials(TRANSCODES_PATH)
|
|
basenames = [os.path.splitext(name)[0] for name in smbclient.listdir(TRANSCODES_PATH, **creds)]
|
|
if not basenames:
|
|
log_to_scan(scan_id, "No transcoded files found.")
|
|
with get_local_db() as conn:
|
|
conn.execute("UPDATE scan_history SET status = 'completed' WHERE id = ?", (scan_id,))
|
|
conn.commit()
|
|
log_to_scan(scan_id, "Scan completed.")
|
|
return
|
|
|
|
log_to_scan(scan_id, f"Found {len(basenames)} transcoded files. Matching with original files in the database...")
|
|
with get_local_db() as conn:
|
|
placeholders = ",".join(['?'] * len(basenames))
|
|
query = f"SELECT p.path || '/' || f.path as full_path, f.path as f_path FROM scenes AS s JOIN files AS f ON s.id = f.scene_id JOIN paths AS p ON f.path_id = p.id WHERE s.oshash IN ({placeholders})"
|
|
rows = conn.execute(query, basenames).fetchall()
|
|
|
|
log_to_scan(scan_id, f"Found {len(rows)} matching original files. Generating plan...")
|
|
for row in rows:
|
|
original_path = row['full_path'].replace(PATH_MAPPING_FROM, PATH_MAPPING_TO)
|
|
transcoded_path = os.path.join(TRANSCODES_PATH, row['f_path'])
|
|
conn.execute("INSERT INTO transcode_plan_results (scan_id, original_path, transcoded_path) VALUES (?, ?, ?)",
|
|
(scan_id, original_path, transcoded_path))
|
|
conn.execute("UPDATE scan_history SET status = 'completed' WHERE id = ?", (scan_id,))
|
|
conn.commit()
|
|
log_to_scan(scan_id, "Transcode replacement plan completed successfully.")
|
|
except Exception as e:
|
|
error_message = f"Error in transcode plan: {e}"
|
|
log_to_scan(scan_id, error_message)
|
|
with app.app_context(), get_local_db() as conn:
|
|
conn.execute("UPDATE scan_history SET status = 'error', message = ? WHERE id = ?", (str(e), scan_id))
|
|
conn.commit()
|
|
|
|
@app.route('/api/transcode_replacement_plan', methods=['POST'])
|
|
def transcode_replacement_plan():
|
|
with get_local_db() as conn:
|
|
cursor = conn.execute("INSERT INTO scan_history (scan_type, status) VALUES (?, ?)", ('transcode_plan', 'in_progress'))
|
|
scan_id = cursor.lastrowid
|
|
threading.Thread(target=do_transcode_plan_task, args=(app.app_context(), scan_id), daemon=True).start()
|
|
return jsonify({"scan_id": scan_id})
|
|
|
|
@app.route('/api/scan_history', methods=['GET'])
|
|
def scan_history():
|
|
with get_local_db() as conn:
|
|
rows = conn.execute("SELECT id, scan_type, strftime('%Y-%m-%d %H:%M:%S', timestamp) as timestamp, status, message FROM scan_history ORDER BY timestamp DESC").fetchall()
|
|
return jsonify([dict(row) for row in rows])
|
|
|
|
@app.route('/api/scan_result/<int:scan_id>', methods=['GET'])
|
|
def scan_result(scan_id):
|
|
with get_local_db() as conn:
|
|
history = conn.execute("SELECT * FROM scan_history WHERE id = ?", (scan_id,)).fetchone()
|
|
if not history: return jsonify({"error": "Scan ID not found"}), 404
|
|
results = []
|
|
if history['scan_type'] == 'duplicates':
|
|
rows = conn.execute("SELECT * FROM duplicate_results WHERE scan_id = ? ORDER BY set_id", (scan_id,)).fetchall()
|
|
grouped_results = defaultdict(list)
|
|
for row in rows: grouped_results[row['set_id']].append(dict(row))
|
|
results = list(grouped_results.values())
|
|
elif history['scan_type'] == 'transcode_plan':
|
|
rows = conn.execute("SELECT * FROM transcode_plan_results WHERE scan_id = ?", (scan_id,)).fetchall()
|
|
results = [dict(row) for row in rows]
|
|
elif history['scan_type'] == 'coomer_download':
|
|
# No specific results to render for this type, just the log
|
|
pass
|
|
elif history['scan_type'] == 'video_download':
|
|
# No specific results to render for this type, just the log
|
|
pass
|
|
return jsonify({"scan_info": dict(history), "results": results})
|
|
|
|
@app.route('/api/scan_log/<int:scan_id>')
|
|
def scan_log(scan_id):
|
|
with get_local_db() as conn:
|
|
row = conn.execute("SELECT log FROM scan_history WHERE id = ?", (scan_id,)).fetchone()
|
|
if row:
|
|
return jsonify({"log": row['log'] or ""})
|
|
return jsonify({"log": ""})
|
|
|
|
|
|
@app.route('/api/delete_files', methods=['POST'])
|
|
def delete_files():
|
|
files_to_delete = request.get_json().get('files', [])
|
|
deleted_files, errors = [], []
|
|
for file_path in files_to_delete:
|
|
try:
|
|
creds = get_smb_credentials(file_path)
|
|
smbclient.remove(file_path, **creds)
|
|
deleted_files.append(file_path)
|
|
except Exception as e:
|
|
errors.append({"file": file_path, "error": str(e)})
|
|
return jsonify({"deleted": deleted_files, "errors": errors})
|
|
|
|
@app.route('/api/execute_transcode_replacement', methods=['POST'])
|
|
def execute_transcode_replacement():
|
|
replacement_plan = request.get_json().get('plan', [])
|
|
replaced_files, errors = [], []
|
|
for item in replacement_plan:
|
|
original_file, new_file = item.get('original'), item.get('transcoded')
|
|
if not original_file or not new_file: continue
|
|
try:
|
|
creds_orig = get_smb_credentials(original_file)
|
|
if smbclient.exists(original_file, **creds_orig):
|
|
smbclient.remove(original_file, **creds_orig)
|
|
creds_new = get_smb_credentials(new_file)
|
|
smbclient.rename(new_file, original_file, **creds_new)
|
|
replaced_files.append({"original": original_file, "new": new_file})
|
|
except Exception as e:
|
|
errors.append({"original": original_file, "new": new_file, "error": str(e)})
|
|
return jsonify({"replaced": replaced_files, "errors": errors})
|
|
|
|
def do_coomer_download_task(app_context, scan_id, url):
|
|
with app_context:
|
|
log_to_scan(scan_id, f"Starting download for {url}")
|
|
try:
|
|
artist_name = url.strip('/').split('/')[-1]
|
|
save_path = os.path.join(COOMER_PATH, artist_name)
|
|
log_to_scan(scan_id, f"Saving files to: {save_path}")
|
|
|
|
creds = get_smb_credentials(save_path)
|
|
if not smbclient.exists(save_path, **creds):
|
|
smbclient.makedirs(save_path, exist_ok=True, **creds)
|
|
|
|
log_to_scan(scan_id, "Fetching artist page...")
|
|
response = requests.get(url, headers={'User-Agent': 'Mozilla/5.0'})
|
|
response.raise_for_status()
|
|
soup = BeautifulSoup(response.text, 'html.parser')
|
|
|
|
video_links = []
|
|
for post in soup.find_all('article'):
|
|
for link in post.find_all('a', href=True):
|
|
if link['href'].endswith('.mp4'):
|
|
video_links.append(link['href'])
|
|
|
|
if not video_links:
|
|
log_to_scan(scan_id, "No video files found on the page.")
|
|
else:
|
|
log_to_scan(scan_id, f"Found {len(video_links)} videos. Starting download...")
|
|
for video_url in video_links:
|
|
filename = video_url.split('/')[-1]
|
|
file_path = os.path.join(save_path, filename)
|
|
|
|
if smbclient.exists(file_path, **creds):
|
|
log_to_scan(scan_id, f"Skipping existing file: {filename}")
|
|
continue
|
|
|
|
log_to_scan(scan_id, f"Downloading {filename}...")
|
|
video_response = requests.get(video_url, stream=True, headers={'User-Agent': 'Mozilla/5.0'})
|
|
video_response.raise_for_status()
|
|
|
|
with smbclient.open_file(file_path, mode='wb', **creds) as f:
|
|
for chunk in video_response.iter_content(chunk_size=8192):
|
|
f.write(chunk)
|
|
log_to_scan(scan_id, f"Finished downloading {filename}")
|
|
|
|
with get_local_db() as conn:
|
|
conn.execute("UPDATE scan_history SET status = 'completed' WHERE id = ?", (scan_id,))
|
|
conn.commit()
|
|
log_to_scan(scan_id, "Coomer download completed successfully.")
|
|
|
|
except Exception as e:
|
|
error_message = f"Error during coomer download: {e}"
|
|
log_to_scan(scan_id, error_message)
|
|
with app.app_context(), get_local_db() as conn:
|
|
conn.execute("UPDATE scan_history SET status = 'error', message = ? WHERE id = ?", (str(e), scan_id))
|
|
conn.commit()
|
|
|
|
@app.route('/api/download_coomer', methods=['POST'])
|
|
def download_coomer():
|
|
url = request.get_json().get('url')
|
|
if not url:
|
|
return jsonify({"error": "URL is required"}), 400
|
|
|
|
with get_local_db() as conn:
|
|
cursor = conn.execute("INSERT INTO scan_history (scan_type, status) VALUES (?, ?)", ('coomer_download', 'in_progress'))
|
|
scan_id = cursor.lastrowid
|
|
|
|
threading.Thread(target=do_coomer_download_task, args=(app.app_context(), scan_id, url), daemon=True).start()
|
|
|
|
return jsonify({"scan_id": scan_id})
|
|
|
|
def do_video_download_task(app_context, scan_id, url, destination_path):
|
|
with app_context:
|
|
log_to_scan(scan_id, f"Starting video download for {url}")
|
|
|
|
try:
|
|
# Security check: Ensure destination is within the main VIDEO_PATH
|
|
creds = get_smb_credentials(destination_path)
|
|
common_path = os.path.commonpath([os.path.normpath(VIDEO_PATH), os.path.normpath(destination_path)])
|
|
if common_path != os.path.normpath(VIDEO_PATH):
|
|
raise ValueError("Destination path is outside the allowed base directory.")
|
|
|
|
if not smbclient.exists(destination_path, **creds):
|
|
smbclient.makedirs(destination_path, exist_ok=True, **creds)
|
|
|
|
log_to_scan(scan_id, "Destination path is valid. Starting yt-dlp...")
|
|
|
|
# We need to construct the command carefully.
|
|
# yt-dlp needs to write to a location the app can access.
|
|
# Let's download to a temporary local file first, then move it to SMB.
|
|
|
|
with tempfile.TemporaryDirectory() as temp_dir:
|
|
temp_file_path_template = os.path.join(temp_dir, '%(title)s.%(ext)s')
|
|
|
|
command = [
|
|
'yt-dlp',
|
|
'-o', temp_file_path_template,
|
|
'--no-overwrites',
|
|
'--progress',
|
|
'--verbose',
|
|
url
|
|
]
|
|
|
|
log_to_scan(scan_id, f"Executing command: {' '.join(command)}")
|
|
|
|
process = subprocess.Popen(command, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, bufsize=1, universal_newlines=True)
|
|
|
|
downloaded_file_name = None
|
|
for line in process.stdout:
|
|
log_to_scan(scan_id, line.strip())
|
|
# Try to find the downloaded file name from yt-dlp output
|
|
if '[download] Destination:' in line:
|
|
downloaded_file_name = os.path.basename(line.split('Destination:')[1].strip())
|
|
elif '[Merger]' in line and 'Merging formats into' in line:
|
|
# Example: [Merger] Merging formats into "/tmp/tmpxxxxx/Video Title.mp4"
|
|
path_part = line.split('Merging formats into "')[1].strip().rstrip('"')
|
|
downloaded_file_name = os.path.basename(path_part)
|
|
|
|
|
|
process.wait()
|
|
|
|
if process.returncode != 0:
|
|
raise RuntimeError(f"yt-dlp exited with error code {process.returncode}")
|
|
|
|
if not downloaded_file_name:
|
|
# If we couldn't find the name, find the first file in the temp dir
|
|
files_in_temp = os.listdir(temp_dir)
|
|
if not files_in_temp:
|
|
raise FileNotFoundError("yt-dlp finished but no file was found in the temporary directory.")
|
|
downloaded_file_name = files_in_temp[0]
|
|
|
|
local_file_path = os.path.join(temp_dir, downloaded_file_name)
|
|
smb_file_path = os.path.join(destination_path, downloaded_file_name)
|
|
|
|
log_to_scan(scan_id, f"Download complete. Moving '{downloaded_file_name}' to {smb_file_path}")
|
|
|
|
if smbclient.exists(smb_file_path, **creds):
|
|
log_to_scan(scan_id, f"File '{downloaded_file_name}' already exists at destination. Deleting local file.")
|
|
else:
|
|
with open(local_file_path, 'rb') as f_in, smbclient.open_file(smb_file_path, 'wb', **creds) as f_out:
|
|
f_out.write(f_in.read())
|
|
log_to_scan(scan_id, "Move to SMB share complete.")
|
|
|
|
with get_local_db() as conn:
|
|
conn.execute("UPDATE scan_history SET status = 'completed' WHERE id = ?", (scan_id,))
|
|
conn.commit()
|
|
log_to_scan(scan_id, "Video download process finished successfully.")
|
|
|
|
except Exception as e:
|
|
error_message = f"Error during video download: {e}"
|
|
log_to_scan(scan_id, error_message)
|
|
with app.app_context(), get_local_db() as conn:
|
|
conn.execute("UPDATE scan_history SET status = 'error', message = ? WHERE id = ?", (str(e), scan_id))
|
|
conn.commit()
|
|
|
|
|
|
@app.route('/api/download_video', methods=['POST'])
|
|
def download_video():
|
|
data = request.get_json()
|
|
url = data.get('url')
|
|
destination_path = data.get('destination_path')
|
|
|
|
if not url or not destination_path:
|
|
return jsonify({"error": "URL and destination_path are required"}), 400
|
|
|
|
with get_local_db() as conn:
|
|
cursor = conn.execute("INSERT INTO scan_history (scan_type, status) VALUES (?, ?)", ('video_download', 'in_progress'))
|
|
scan_id = cursor.lastrowid
|
|
|
|
threading.Thread(target=do_video_download_task, args=(app.app_context(), scan_id, url, destination_path), daemon=True).start()
|
|
|
|
return jsonify({"scan_id": scan_id})
|
|
|
|
@app.route('/api/browse', methods=['POST'])
|
|
def browse():
|
|
path = request.get_json().get('path', VIDEO_PATH)
|
|
try:
|
|
creds = get_smb_credentials(path)
|
|
items = [dict(name=e.name, path=e.path, is_dir=e.is_dir()) for e in smbclient.scandir(path, **creds)]
|
|
return jsonify(items)
|
|
except Exception as e:
|
|
return jsonify({"error": str(e)}), 500
|
|
|
|
@app.route('/api/preview_cbz')
|
|
def preview_cbz():
|
|
path = request.args.get('path')
|
|
if not path:
|
|
return "Missing path parameter", 400
|
|
|
|
try:
|
|
creds = get_smb_credentials(path)
|
|
|
|
smb_file = smbclient.open_file(path, mode='rb', **creds)
|
|
|
|
cbz_buffer = io.BytesIO(smb_file.read())
|
|
|
|
image_urls = []
|
|
with zipfile.ZipFile(cbz_buffer) as zf:
|
|
# Sort the file list to ensure correct order
|
|
file_list = sorted(zf.namelist())
|
|
for filename in file_list:
|
|
# Check if the file is an image
|
|
mimetype = mimetypes.guess_type(filename)[0]
|
|
if mimetype and mimetype.startswith('image/'):
|
|
with zf.open(filename) as image_file:
|
|
image_data = image_file.read()
|
|
encoded_image = base64.b64encode(image_data).decode('utf-8')
|
|
data_url = f"data:{mimetype};base64,{encoded_image}"
|
|
image_urls.append(data_url)
|
|
|
|
return jsonify(image_urls)
|
|
|
|
except Exception as e:
|
|
return jsonify({"error": str(e)}), 500
|
|
|
|
@app.route('/api/preview')
|
|
def preview():
|
|
path = request.args.get('path')
|
|
if not path:
|
|
return "Missing path parameter", 400
|
|
|
|
try:
|
|
creds = get_smb_credentials(path)
|
|
|
|
# Open the file with smbclient
|
|
smb_file = smbclient.open_file(path, mode='rb', **creds)
|
|
|
|
# Read the file into a seekable buffer (important for video)
|
|
buffer = io.BytesIO(smb_file.read())
|
|
buffer.seek(0)
|
|
|
|
mimetype = mimetypes.guess_type(path)[0]
|
|
|
|
return send_file(
|
|
buffer,
|
|
mimetype=mimetype,
|
|
as_attachment=False
|
|
)
|
|
except Exception as e:
|
|
return str(e), 500
|
|
|
|
@app.route('/api/qb_torrents')
|
|
def qb_torrents():
|
|
try:
|
|
client = get_qb_client()
|
|
if not client:
|
|
return jsonify({"error": "qBittorrent not configured."}), 400
|
|
|
|
torrents = []
|
|
for torrent in client.torrents_info():
|
|
torrents.append({
|
|
"hash": torrent.hash,
|
|
"name": torrent.name,
|
|
"size": torrent.size,
|
|
"progress": torrent.progress,
|
|
"status": torrent.state,
|
|
"dlspeed": torrent.dlspeed,
|
|
"upspeed": torrent.upspeed,
|
|
"eta": torrent.eta,
|
|
"save_path": torrent.save_path
|
|
})
|
|
return jsonify(torrents)
|
|
except Exception as e:
|
|
return jsonify({"error": str(e)}), 500
|
|
|
|
@app.route('/api/qb_action', methods=['POST'])
|
|
def qb_action():
|
|
try:
|
|
client = get_qb_client()
|
|
if not client:
|
|
return jsonify({"error": "qBittorrent not configured."}), 400
|
|
|
|
data = request.get_json()
|
|
action = data.get('action')
|
|
hashes = data.get('hashes')
|
|
|
|
if not action or not hashes:
|
|
return jsonify({"error": "Missing action or hashes."}), 400
|
|
|
|
if action == 'pause':
|
|
client.torrents_pause(torrent_hashes=hashes)
|
|
elif action == 'resume':
|
|
client.torrents_resume(torrent_hashes=hashes)
|
|
elif action == 'delete':
|
|
client.torrents_delete(delete_files=True, torrent_hashes=hashes)
|
|
elif action == 'force_resume':
|
|
client.torrents_force_resume(torrent_hashes=hashes)
|
|
elif action == 'recheck':
|
|
client.torrents_recheck(torrent_hashes=hashes)
|
|
else:
|
|
return jsonify({"error": "Invalid action."}), 400
|
|
|
|
return jsonify({"success": True})
|
|
except Exception as e:
|
|
return jsonify({"error": str(e)}), 500
|
|
|
|
@app.route('/api/create_folder', methods=['POST'])
|
|
def create_folder():
|
|
data = request.get_json()
|
|
path = data.get('path')
|
|
folder_name = data.get('folder_name')
|
|
|
|
if not path or not folder_name:
|
|
return jsonify({"error": "Path and folder name are required"}), 400
|
|
|
|
try:
|
|
creds = get_smb_credentials(path)
|
|
# Security check: Ensure destination is within the main VIDEO_PATH
|
|
common_path = os.path.commonpath([os.path.normpath(VIDEO_PATH), os.path.normpath(path)])
|
|
if common_path != os.path.normpath(VIDEO_PATH):
|
|
raise ValueError("Creation path is outside the allowed base directory.")
|
|
|
|
new_folder_path = os.path.join(path, folder_name)
|
|
smbclient.makedirs(new_folder_path, exist_ok=False, **creds)
|
|
return jsonify({"success": True})
|
|
except Exception as e:
|
|
return jsonify({"error": str(e)}), 500
|
|
|
|
@app.route('/api/move_file', methods=['POST'])
|
|
def move_file():
|
|
data = request.get_json()
|
|
source_path = data.get('source_path')
|
|
dest_path = data.get('dest_path')
|
|
|
|
if not source_path or not dest_path:
|
|
return jsonify({"error": "Source and destination paths are required"}), 400
|
|
|
|
try:
|
|
creds = get_smb_credentials(source_path)
|
|
# Ensure the destination directory exists
|
|
if not smbclient.exists(dest_path, **creds):
|
|
smbclient.makedirs(dest_path, exist_ok=True, **creds)
|
|
|
|
# Construct the full destination path including the filename
|
|
filename = os.path.basename(source_path)
|
|
full_dest_path = os.path.join(dest_path, filename)
|
|
|
|
smbclient.rename(source_path, full_dest_path, **creds)
|
|
return jsonify({"success": True})
|
|
except Exception as e:
|
|
return jsonify({"error": str(e)}), 500
|
|
|
|
@app.route('/api/qb_set_location', methods=['POST'])
|
|
def qb_set_location():
|
|
try:
|
|
client = get_qb_client()
|
|
if not client:
|
|
return jsonify({"error": "qBittorrent not configured."}), 400
|
|
|
|
data = request.get_json()
|
|
hashes = data.get('hashes')
|
|
location = data.get('location')
|
|
|
|
if not hashes or not location:
|
|
return jsonify({"error": "Missing hashes or location."}), 400
|
|
|
|
client.torrents_set_location(torrent_hashes=hashes, location=location)
|
|
return jsonify({"success": True})
|
|
except Exception as e:
|
|
return jsonify({"error": str(e)}), 500
|
|
|
|
if __name__ == '__main__':
|
|
app.run(debug=True, port=5000, host='0.0.0.0')
|
|
|
|
|