feat: Add Stash Transcode Correlator & Storage Saver with atomic original replacement
This commit is contained in:
@@ -809,6 +809,8 @@ class RobustProxyHandler(http.server.SimpleHTTPRequestHandler):
|
||||
self.handle_stash_missing_metadata()
|
||||
elif self.path.startswith('/api/stash/tags/duplicates'):
|
||||
self.handle_stash_tag_duplicates()
|
||||
elif self.path.startswith('/api/stash/transcodes/scan'):
|
||||
self.handle_stash_transcode_scan()
|
||||
# SMB API
|
||||
elif self.path.startswith('/api/smb/status'):
|
||||
self.handle_smb_status()
|
||||
@@ -884,6 +886,10 @@ class RobustProxyHandler(http.server.SimpleHTTPRequestHandler):
|
||||
self.handle_stash_performer_merge()
|
||||
elif self.path == '/api/stash/tags/merge':
|
||||
self.handle_stash_tag_merge()
|
||||
elif self.path == '/api/stash/transcodes/replace':
|
||||
self.handle_stash_transcode_replace()
|
||||
elif self.path == '/api/stash/transcodes/batch-replace':
|
||||
self.handle_stash_transcode_batch_replace()
|
||||
# SMB API Moves, Deletes & Cleanups
|
||||
elif self.path == '/api/smb/move':
|
||||
self.handle_smb_move()
|
||||
@@ -2612,6 +2618,281 @@ class RobustProxyHandler(http.server.SimpleHTTPRequestHandler):
|
||||
except Exception as e:
|
||||
self.send_json_response(500, {"error": str(e)})
|
||||
|
||||
# ==================== Stash Transcode Correlator & Storage Saver Handlers ====================
|
||||
def handle_stash_transcode_scan(self):
|
||||
parsed = urllib.parse.urlparse(self.path)
|
||||
params = urllib.parse.parse_qs(parsed.query)
|
||||
stash_url = params.get('stash_url', [DEFAULT_STASH_URL])[0]
|
||||
transcode_dir = params.get('transcode_dir', ['/mnt/isolation/stashapp/generated/transcodes'])[0]
|
||||
stash_prefix = params.get('stash_prefix', ['/data'])[0]
|
||||
local_prefix = params.get('local_prefix', ['/mnt/isolation/videos'])[0]
|
||||
|
||||
if not os.path.exists(transcode_dir):
|
||||
self.send_json_response(400, {"error": f"Transcode directory not found: {transcode_dir}", "items": []})
|
||||
return
|
||||
|
||||
# Build dictionary of all transcode files on disk
|
||||
transcode_files = {}
|
||||
for f in os.listdir(transcode_dir):
|
||||
if f.endswith('.mp4'):
|
||||
h = f[:-4]
|
||||
full_p = os.path.join(transcode_dir, f)
|
||||
try:
|
||||
transcode_files[h] = {
|
||||
"filename": f,
|
||||
"path": full_p,
|
||||
"size": os.path.getsize(full_p)
|
||||
}
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
# Scan Stash findFiles in batches of 5000 until all matched or end
|
||||
matched_items = []
|
||||
page = 1
|
||||
page_size = 5000
|
||||
unmatched_transcodes = set(transcode_files.keys())
|
||||
|
||||
while unmatched_transcodes:
|
||||
q = f'''
|
||||
query GetFilesBatch {{
|
||||
findFiles(filter: {{ per_page: {page_size}, page: {page} }}) {{
|
||||
files {{
|
||||
id
|
||||
path
|
||||
basename
|
||||
size
|
||||
fingerprints {{
|
||||
type
|
||||
value
|
||||
}}
|
||||
... on VideoFile {{
|
||||
width
|
||||
height
|
||||
duration
|
||||
video_codec
|
||||
}}
|
||||
}}
|
||||
}}
|
||||
}}
|
||||
'''
|
||||
try:
|
||||
res = execute_stash_graphql(q, stash_url, timeout=30)
|
||||
files = res.get('data', {}).get('findFiles', {}).get('files', [])
|
||||
if not files:
|
||||
break
|
||||
|
||||
for f in files:
|
||||
fps = {fp['type']: fp['value'] for fp in f.get('fingerprints', [])}
|
||||
oshash = fps.get('oshash')
|
||||
if oshash and oshash in transcode_files:
|
||||
tinfo = transcode_files[oshash]
|
||||
orig_stash_path = f.get('path', '')
|
||||
orig_local_path = orig_stash_path.replace(stash_prefix, local_prefix) if stash_prefix and orig_stash_path.startswith(stash_prefix) else orig_stash_path
|
||||
|
||||
orig_size = f.get('size') or 0
|
||||
transcode_size = tinfo['size']
|
||||
saved_bytes = orig_size - transcode_size
|
||||
pct_change = ((saved_bytes / orig_size) * 100) if orig_size > 0 else 0
|
||||
|
||||
matched_items.append({
|
||||
"file_id": f['id'],
|
||||
"oshash": oshash,
|
||||
"basename": f.get('basename', ''),
|
||||
"orig_stash_path": orig_stash_path,
|
||||
"orig_local_path": orig_local_path,
|
||||
"orig_exists": os.path.exists(orig_local_path),
|
||||
"orig_size": orig_size,
|
||||
"orig_resolution": f"{f.get('width', '')}x{f.get('height', '')}",
|
||||
"orig_codec": f.get('video_codec', ''),
|
||||
"transcode_filename": tinfo['filename'],
|
||||
"transcode_path": tinfo['path'],
|
||||
"transcode_size": transcode_size,
|
||||
"saved_bytes": saved_bytes,
|
||||
"pct_change": round(pct_change, 1)
|
||||
})
|
||||
unmatched_transcodes.discard(oshash)
|
||||
|
||||
if len(files) < page_size:
|
||||
break
|
||||
page += 1
|
||||
except Exception as e:
|
||||
print(f"Error querying Stash findFiles page {page}: {e}")
|
||||
break
|
||||
|
||||
# Sort by space saved descending
|
||||
matched_items.sort(key=lambda x: x['saved_bytes'], reverse=True)
|
||||
total_saved_bytes = sum(x['saved_bytes'] for x in matched_items if x['saved_bytes'] > 0)
|
||||
|
||||
self.send_json_response(200, {
|
||||
"items": matched_items,
|
||||
"total_matched": len(matched_items),
|
||||
"total_transcodes_on_disk": len(transcode_files),
|
||||
"total_potential_savings_bytes": total_saved_bytes
|
||||
})
|
||||
|
||||
def handle_stash_transcode_replace(self):
|
||||
content_length = int(self.headers.get('Content-Length', 0))
|
||||
post_data = self.rfile.read(content_length)
|
||||
try:
|
||||
payload = json.loads(post_data.decode('utf-8'))
|
||||
orig_local_path = payload.get('orig_local_path')
|
||||
transcode_path = payload.get('transcode_path')
|
||||
delete_transcode = bool(payload.get('delete_transcode', True))
|
||||
trigger_rescan = bool(payload.get('trigger_rescan', True))
|
||||
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
|
||||
|
||||
if not orig_local_path or not transcode_path:
|
||||
self.send_json_response(400, {"error": "Missing orig_local_path or transcode_path"})
|
||||
return
|
||||
|
||||
if not os.path.exists(transcode_path):
|
||||
self.send_json_response(400, {"error": f"Transcode file does not exist: {transcode_path}"})
|
||||
return
|
||||
|
||||
if not os.path.exists(orig_local_path):
|
||||
self.send_json_response(400, {"error": f"Original file does not exist at local path: {orig_local_path}"})
|
||||
return
|
||||
|
||||
orig_size = os.path.getsize(orig_local_path)
|
||||
transcode_size = os.path.getsize(transcode_path)
|
||||
if transcode_size == 0:
|
||||
self.send_json_response(400, {"error": "Transcode file is 0 bytes (corrupted), aborting replacement"})
|
||||
return
|
||||
|
||||
# Step 1: Copy transcode to temporary staging file next to original
|
||||
tmp_staging_path = f"{orig_local_path}.transcode.tmp"
|
||||
if os.path.exists(tmp_staging_path):
|
||||
try:
|
||||
os.remove(tmp_staging_path)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
shutil.copy2(transcode_path, tmp_staging_path)
|
||||
|
||||
# Step 2: Verify size of staging file
|
||||
staged_size = os.path.getsize(tmp_staging_path)
|
||||
if staged_size != transcode_size:
|
||||
if os.path.exists(tmp_staging_path):
|
||||
os.remove(tmp_staging_path)
|
||||
self.send_json_response(500, {"error": f"Size mismatch during staging: expected {transcode_size}, got {staged_size}"})
|
||||
return
|
||||
|
||||
# Step 3: Atomically overwrite original file with staged transcode
|
||||
os.replace(tmp_staging_path, orig_local_path)
|
||||
|
||||
# Step 4: If delete_transcode is requested, remove original transcode cache
|
||||
if delete_transcode:
|
||||
try:
|
||||
os.remove(transcode_path)
|
||||
except OSError as e:
|
||||
print(f"Warning: Could not remove transcode cache file {transcode_path}: {e}")
|
||||
|
||||
saved_bytes = orig_size - transcode_size
|
||||
|
||||
# Step 5: Record to Relocation History audit log
|
||||
record_relocation_history(
|
||||
source_path=transcode_path,
|
||||
target_path=orig_local_path,
|
||||
performer_name="Transcode Optimization",
|
||||
status="completed",
|
||||
size_bytes=transcode_size,
|
||||
action_type="transcode_replace"
|
||||
)
|
||||
|
||||
# Step 6: Trigger Stash metadata scan on parent folder if requested
|
||||
if trigger_rescan:
|
||||
try:
|
||||
parent_dir = os.path.dirname(orig_local_path)
|
||||
stash_parent = parent_dir.replace('/mnt/isolation/videos', '/data')
|
||||
mutation = '''
|
||||
mutation MetadataScan($input: ScanMetadataInput!) {
|
||||
metadataScan(input: $input)
|
||||
}
|
||||
'''
|
||||
execute_stash_graphql(mutation, stash_url, variables={"input": {"paths": [stash_parent], "rescan": True}}, timeout=10)
|
||||
except Exception as e:
|
||||
print(f"Background Stash rescan trigger error: {e}")
|
||||
|
||||
self.send_json_response(200, {
|
||||
"success": True,
|
||||
"orig_local_path": orig_local_path,
|
||||
"orig_size": orig_size,
|
||||
"new_size": transcode_size,
|
||||
"saved_bytes": saved_bytes
|
||||
})
|
||||
except Exception as e:
|
||||
self.send_json_response(500, {"error": str(e)})
|
||||
|
||||
def handle_stash_transcode_batch_replace(self):
|
||||
content_length = int(self.headers.get('Content-Length', 0))
|
||||
post_data = self.rfile.read(content_length)
|
||||
try:
|
||||
payload = json.loads(post_data.decode('utf-8'))
|
||||
items = payload.get('items', [])
|
||||
delete_transcode = bool(payload.get('delete_transcode', True))
|
||||
trigger_rescan = bool(payload.get('trigger_rescan', True))
|
||||
stash_url = payload.get('stash_url') or DEFAULT_STASH_URL
|
||||
|
||||
replaced_count = 0
|
||||
total_saved_bytes = 0
|
||||
errors = []
|
||||
rescanned_dirs = set()
|
||||
|
||||
for item in items:
|
||||
orig_local_path = item.get('orig_local_path')
|
||||
transcode_path = item.get('transcode_path')
|
||||
if not orig_local_path or not transcode_path or not os.path.exists(transcode_path) or not os.path.exists(orig_local_path):
|
||||
continue
|
||||
|
||||
try:
|
||||
orig_size = os.path.getsize(orig_local_path)
|
||||
transcode_size = os.path.getsize(transcode_path)
|
||||
if transcode_size == 0:
|
||||
continue
|
||||
|
||||
tmp_staging_path = f"{orig_local_path}.transcode.tmp"
|
||||
shutil.copy2(transcode_path, tmp_staging_path)
|
||||
if os.path.getsize(tmp_staging_path) != transcode_size:
|
||||
if os.path.exists(tmp_staging_path):
|
||||
os.remove(tmp_staging_path)
|
||||
continue
|
||||
|
||||
os.replace(tmp_staging_path, orig_local_path)
|
||||
if delete_transcode:
|
||||
try:
|
||||
os.remove(transcode_path)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
replaced_count += 1
|
||||
saved = orig_size - transcode_size
|
||||
total_saved_bytes += saved
|
||||
rescanned_dirs.add(os.path.dirname(orig_local_path))
|
||||
except Exception as e:
|
||||
errors.append({"path": orig_local_path, "error": str(e)})
|
||||
|
||||
# Trigger Stash scan for all touched directories
|
||||
if trigger_rescan and rescanned_dirs:
|
||||
try:
|
||||
stash_dirs = [d.replace('/mnt/isolation/videos', '/data') for d in rescanned_dirs]
|
||||
mutation = '''
|
||||
mutation MetadataScan($input: ScanMetadataInput!) {
|
||||
metadataScan(input: $input)
|
||||
}
|
||||
'''
|
||||
execute_stash_graphql(mutation, stash_url, variables={"input": {"paths": stash_dirs, "rescan": True}}, timeout=15)
|
||||
except Exception as e:
|
||||
print(f"Batch Stash rescan trigger error: {e}")
|
||||
|
||||
self.send_json_response(200, {
|
||||
"success": True,
|
||||
"replaced_count": replaced_count,
|
||||
"total_saved_bytes": total_saved_bytes,
|
||||
"errors": errors
|
||||
})
|
||||
except Exception as e:
|
||||
self.send_json_response(500, {"error": str(e)})
|
||||
|
||||
# ==================== SMB Maintenance & Cleanup Handlers ====================
|
||||
def handle_smb_scan_cleanup(self):
|
||||
parsed = urllib.parse.urlparse(self.path)
|
||||
|
||||
Reference in New Issue
Block a user