diff --git a/snot/alpine@3.x.x.cdn.min.js b/snot/alpine@3.x.x.cdn.min.js new file mode 100644 index 0000000..2a6849c --- /dev/null +++ b/snot/alpine@3.x.x.cdn.min.js @@ -0,0 +1,5 @@ +(()=>{var nt=!1,it=!1,G=[],ot=-1;function Ut(e){In(e)}function In(e){G.includes(e)||G.push(e),$n()}function Wt(e){let t=G.indexOf(e);t!==-1&&t>ot&&G.splice(t,1)}function $n(){!it&&!nt&&(nt=!0,queueMicrotask(Ln))}function Ln(){nt=!1,it=!0;for(let e=0;ee.effect(t,{scheduler:r=>{st?Ut(r):r()}}),at=e.raw}function ct(e){N=e}function Yt(e){let t=()=>{};return[n=>{let i=N(n);return e._x_effects||(e._x_effects=new Set,e._x_runEffects=()=>{e._x_effects.forEach(o=>o())}),e._x_effects.add(i),t=()=>{i!==void 0&&(e._x_effects.delete(i),F(i))},i},()=>{t()}]}function Oe(e,t){let r=!0,n,i=N(()=>{let o=e();JSON.stringify(o),r?n=o:queueMicrotask(()=>{t(o,n),n=o}),r=!1});return()=>F(i)}var Xt=[],Zt=[],Qt=[];function er(e){Qt.push(e)}function re(e,t){typeof t=="function"?(e._x_cleanups||(e._x_cleanups=[]),e._x_cleanups.push(t)):(t=e,Zt.push(t))}function Re(e){Xt.push(e)}function Te(e,t,r){e._x_attributeCleanups||(e._x_attributeCleanups={}),e._x_attributeCleanups[t]||(e._x_attributeCleanups[t]=[]),e._x_attributeCleanups[t].push(r)}function lt(e,t){e._x_attributeCleanups&&Object.entries(e._x_attributeCleanups).forEach(([r,n])=>{(t===void 0||t.includes(r))&&(n.forEach(i=>i()),delete e._x_attributeCleanups[r])})}function tr(e){for(e._x_effects?.forEach(Wt);e._x_cleanups?.length;)e._x_cleanups.pop()()}var ut=new MutationObserver(mt),ft=!1;function pe(){ut.observe(document,{subtree:!0,childList:!0,attributes:!0,attributeOldValue:!0}),ft=!0}function dt(){jn(),ut.disconnect(),ft=!1}var de=[];function jn(){let e=ut.takeRecords();de.push(()=>e.length>0&&mt(e));let t=de.length;queueMicrotask(()=>{if(de.length===t)for(;de.length>0;)de.shift()()})}function m(e){if(!ft)return e();dt();let t=e();return pe(),t}var pt=!1,Ce=[];function rr(){pt=!0}function nr(){pt=!1,mt(Ce),Ce=[]}function mt(e){if(pt){Ce=Ce.concat(e);return}let t=[],r=new Set,n=new Map,i=new Map;for(let o=0;o{s.nodeType===1&&s._x_marker&&r.add(s)}),e[o].addedNodes.forEach(s=>{if(s.nodeType===1){if(r.has(s)){r.delete(s);return}s._x_marker||t.push(s)}})),e[o].type==="attributes")){let s=e[o].target,a=e[o].attributeName,c=e[o].oldValue,l=()=>{n.has(s)||n.set(s,[]),n.get(s).push({name:a,value:s.getAttribute(a)})},u=()=>{i.has(s)||i.set(s,[]),i.get(s).push(a)};s.hasAttribute(a)&&c===null?l():s.hasAttribute(a)?(u(),l()):u()}i.forEach((o,s)=>{lt(s,o)}),n.forEach((o,s)=>{Xt.forEach(a=>a(s,o))});for(let o of r)t.some(s=>s.contains(o))||Zt.forEach(s=>s(o));for(let o of t)o.isConnected&&Qt.forEach(s=>s(o));t=null,r=null,n=null,i=null}function Me(e){return k(B(e))}function D(e,t,r){return e._x_dataStack=[t,...B(r||e)],()=>{e._x_dataStack=e._x_dataStack.filter(n=>n!==t)}}function B(e){return e._x_dataStack?e._x_dataStack:typeof ShadowRoot=="function"&&e instanceof ShadowRoot?B(e.host):e.parentNode?B(e.parentNode):[]}function k(e){return new Proxy({objects:e},Fn)}var Fn={ownKeys({objects:e}){return Array.from(new Set(e.flatMap(t=>Object.keys(t))))},has({objects:e},t){return t==Symbol.unscopables?!1:e.some(r=>Object.prototype.hasOwnProperty.call(r,t)||Reflect.has(r,t))},get({objects:e},t,r){return t=="toJSON"?Bn:Reflect.get(e.find(n=>Reflect.has(n,t))||{},t,r)},set({objects:e},t,r,n){let i=e.find(s=>Object.prototype.hasOwnProperty.call(s,t))||e[e.length-1],o=Object.getOwnPropertyDescriptor(i,t);return o?.set&&o?.get?o.set.call(n,r)||!0:Reflect.set(i,t,r)}};function Bn(){return Reflect.ownKeys(this).reduce((t,r)=>(t[r]=Reflect.get(this,r),t),{})}function ne(e){let t=n=>typeof n=="object"&&!Array.isArray(n)&&n!==null,r=(n,i="")=>{Object.entries(Object.getOwnPropertyDescriptors(n)).forEach(([o,{value:s,enumerable:a}])=>{if(a===!1||s===void 0||typeof s=="object"&&s!==null&&s.__v_skip)return;let c=i===""?o:`${i}.${o}`;typeof s=="object"&&s!==null&&s._x_interceptor?n[o]=s.initialize(e,c,o):t(s)&&s!==n&&!(s instanceof Element)&&r(s,c)})};return r(e)}function Ne(e,t=()=>{}){let r={initialValue:void 0,_x_interceptor:!0,initialize(n,i,o){return e(this.initialValue,()=>zn(n,i),s=>ht(n,i,s),i,o)}};return t(r),n=>{if(typeof n=="object"&&n!==null&&n._x_interceptor){let i=r.initialize.bind(r);r.initialize=(o,s,a)=>{let c=n.initialize(o,s,a);return r.initialValue=c,i(o,s,a)}}else r.initialValue=n;return r}}function zn(e,t){return t.split(".").reduce((r,n)=>r[n],e)}function ht(e,t,r){if(typeof t=="string"&&(t=t.split(".")),t.length===1)e[t[0]]=r;else{if(t.length===0)throw error;return e[t[0]]||(e[t[0]]={}),ht(e[t[0]],t.slice(1),r)}}var ir={};function y(e,t){ir[e]=t}function K(e,t){let r=Hn(t);return Object.entries(ir).forEach(([n,i])=>{Object.defineProperty(e,`$${n}`,{get(){return i(t,r)},enumerable:!1})}),e}function Hn(e){let[t,r]=_t(e),n={interceptor:Ne,...t};return re(e,r),n}function or(e,t,r,...n){try{return r(...n)}catch(i){ie(i,e,t)}}function ie(...e){return sr(...e)}var sr=Kn;function ar(e){sr=e}function Kn(e,t,r=void 0){e=Object.assign(e??{message:"No error message given."},{el:t,expression:r}),console.warn(`Alpine Expression Error: ${e.message} + +${r?'Expression: "'+r+`" + +`:""}`,t),setTimeout(()=>{throw e},0)}var oe=!0;function De(e){let t=oe;oe=!1;let r=e();return oe=t,r}function T(e,t,r={}){let n;return x(e,t)(i=>n=i,r),n}function x(...e){return cr(...e)}var cr=xt;function lr(e){cr=e}var ur;function fr(e){ur=e}function xt(e,t){let r={};K(r,e);let n=[r,...B(e)],i=typeof t=="function"?Vn(n,t):Un(n,t,e);return or.bind(null,e,t,i)}function Vn(e,t){return(r=()=>{},{scope:n={},params:i=[],context:o}={})=>{if(!oe){me(r,t,k([n,...e]),i);return}let s=t.apply(k([n,...e]),i);me(r,s)}}var gt={};function qn(e,t){if(gt[e])return gt[e];let r=Object.getPrototypeOf(async function(){}).constructor,n=/^[\n\s]*if.*\(.*\)/.test(e.trim())||/^(let|const)\s/.test(e.trim())?`(async()=>{ ${e} })()`:e,o=(()=>{try{let s=new r(["__self","scope"],`with (scope) { __self.result = ${n} }; __self.finished = true; return __self.result;`);return Object.defineProperty(s,"name",{value:`[Alpine] ${e}`}),s}catch(s){return ie(s,t,e),Promise.resolve()}})();return gt[e]=o,o}function Un(e,t,r){let n=qn(t,r);return(i=()=>{},{scope:o={},params:s=[],context:a}={})=>{n.result=void 0,n.finished=!1;let c=k([o,...e]);if(typeof n=="function"){let l=n.call(a,n,c).catch(u=>ie(u,r,t));n.finished?(me(i,n.result,c,s,r),n.result=void 0):l.then(u=>{me(i,u,c,s,r)}).catch(u=>ie(u,r,t)).finally(()=>n.result=void 0)}}}function me(e,t,r,n,i){if(oe&&typeof t=="function"){let o=t.apply(r,n);o instanceof Promise?o.then(s=>me(e,s,r,n)).catch(s=>ie(s,i,t)):e(o)}else typeof t=="object"&&t instanceof Promise?t.then(o=>e(o)):e(t)}function dr(...e){return ur(...e)}function pr(e,t,r={}){let n={};K(n,e);let i=[n,...B(e)],o=k([r.scope??{},...i]),s=r.params??[];if(t.includes("await")){let a=Object.getPrototypeOf(async function(){}).constructor,c=/^[\n\s]*if.*\(.*\)/.test(t.trim())||/^(let|const)\s/.test(t.trim())?`(async()=>{ ${t} })()`:t;return new a(["scope"],`with (scope) { let __result = ${c}; return __result }`).call(r.context,o)}else{let a=/^[\n\s]*if.*\(.*\)/.test(t.trim())||/^(let|const)\s/.test(t.trim())?`(()=>{ ${t} })()`:t,l=new Function(["scope"],`with (scope) { let __result = ${a}; return __result }`).call(r.context,o);return typeof l=="function"&&oe?l.apply(o,s):l}}var wt="x-";function C(e=""){return wt+e}function mr(e){wt=e}var ke={};function d(e,t){return ke[e]=t,{before(r){if(!ke[r]){console.warn(String.raw`Cannot find directive \`${r}\`. \`${e}\` will use the default order of execution`);return}let n=J.indexOf(r);J.splice(n>=0?n:J.indexOf("DEFAULT"),0,e)}}}function hr(e){return Object.keys(ke).includes(e)}function _e(e,t,r){if(t=Array.from(t),e._x_virtualDirectives){let o=Object.entries(e._x_virtualDirectives).map(([a,c])=>({name:a,value:c})),s=Et(o);o=o.map(a=>s.find(c=>c.name===a.name)?{name:`x-bind:${a.name}`,value:`"${a.value}"`}:a),t=t.concat(o)}let n={};return t.map(xr((o,s)=>n[o]=s)).filter(br).map(Gn(n,r)).sort(Jn).map(o=>Wn(e,o))}function Et(e){return Array.from(e).map(xr()).filter(t=>!br(t))}var yt=!1,he=new Map,_r=Symbol();function gr(e){yt=!0;let t=Symbol();_r=t,he.set(t,[]);let r=()=>{for(;he.get(t).length;)he.get(t).shift()();he.delete(t)},n=()=>{yt=!1,r()};e(r),n()}function _t(e){let t=[],r=a=>t.push(a),[n,i]=Yt(e);return t.push(i),[{Alpine:z,effect:n,cleanup:r,evaluateLater:x.bind(x,e),evaluate:T.bind(T,e)},()=>t.forEach(a=>a())]}function Wn(e,t){let r=()=>{},n=ke[t.type]||r,[i,o]=_t(e);Te(e,t.original,o);let s=()=>{e._x_ignore||e._x_ignoreSelf||(n.inline&&n.inline(e,t,i),n=n.bind(n,e,t,i),yt?he.get(_r).push(n):n())};return s.runCleanups=o,s}var Pe=(e,t)=>({name:r,value:n})=>(r.startsWith(e)&&(r=r.replace(e,t)),{name:r,value:n}),Ie=e=>e;function xr(e=()=>{}){return({name:t,value:r})=>{let{name:n,value:i}=yr.reduce((o,s)=>s(o),{name:t,value:r});return n!==t&&e(n,t),{name:n,value:i}}}var yr=[];function se(e){yr.push(e)}function br({name:e}){return wr().test(e)}var wr=()=>new RegExp(`^${wt}([^:^.]+)\\b`);function Gn(e,t){return({name:r,value:n})=>{let i=r.match(wr()),o=r.match(/:([a-zA-Z0-9\-_:]+)/),s=r.match(/\.[^.\]]+(?=[^\]]*$)/g)||[],a=t||e[r]||r;return{type:i?i[1]:null,value:o?o[1]:null,modifiers:s.map(c=>c.replace(".","")),expression:n,original:a}}}var bt="DEFAULT",J=["ignore","ref","data","id","anchor","bind","init","for","model","modelable","transition","show","if",bt,"teleport"];function Jn(e,t){let r=J.indexOf(e.type)===-1?bt:e.type,n=J.indexOf(t.type)===-1?bt:t.type;return J.indexOf(r)-J.indexOf(n)}function Y(e,t,r={}){e.dispatchEvent(new CustomEvent(t,{detail:r,bubbles:!0,composed:!0,cancelable:!0}))}function P(e,t){if(typeof ShadowRoot=="function"&&e instanceof ShadowRoot){Array.from(e.children).forEach(i=>P(i,t));return}let r=!1;if(t(e,()=>r=!0),r)return;let n=e.firstElementChild;for(;n;)P(n,t,!1),n=n.nextElementSibling}function E(e,...t){console.warn(`Alpine Warning: ${e}`,...t)}var Er=!1;function vr(){Er&&E("Alpine has already been initialized on this page. Calling Alpine.start() more than once can cause problems."),Er=!0,document.body||E("Unable to initialize. Trying to load Alpine before `` is available. Did you forget to add `defer` in Alpine's ` + diff --git a/snot/snot.py b/snot/snot.py new file mode 100755 index 0000000..0ac2c85 --- /dev/null +++ b/snot/snot.py @@ -0,0 +1,785 @@ +#!/usr/bin/env -S uv run --script +# /// script +# dependencies = ["psycopg[binary]", "click", "httpx"] +# /// + +import argparse +import contextlib +import dataclasses +import datetime +import hashlib +import itertools +import json +import logging +import os +import re +import shlex +import shutil +import socket +import subprocess +import threading +import time +import typing +from functools import partial, wraps +from os import getenv +from pathlib import Path +from typing import Any, Callable +from urllib.parse import parse_qs +from wsgiref.simple_server import WSGIRequestHandler, make_server + +import click +import httpx +import psycopg +from psycopg import sql as pgsql + +logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s: %(message)s") + +HandlerFunc = Callable[["Request"], "Response"] + + +@dataclasses.dataclass +class FfprobeResult: + duration_sec: float + duration_human: str + codec: str + fps: float + size_bytes: int + width: int + height: int + bitrate: int + container: str + ar: float + sample_ar: float + tags: dict[str, str] + + @property + def resolution(self) -> str: + return f"{self.width}x{self.height}" + + +def ffprobe(video_path: Path) -> FfprobeResult: + proc = subprocess.run( + args=[ + "ffprobe", + "-v", + "error", + "-show_streams", + "-show_format", + "-print_format", + "json", + str(video_path), + ], + text=True, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + ) + try: + proc.check_returncode() + except subprocess.CalledProcessError as e: + logging.error(f"failed to run ffprobe. stderr={e.stderr}") + raise + output: dict = json.loads(proc.stdout) + video_stream: dict = [ + s for s in output.get("streams", []) if s.get("codec_type") == "video" + ][0] + video_format: dict = output.get("format", {}) + width, height = video_stream.get("width"), video_stream.get("height") + codec = video_stream.get("codec_name") + duration = float(video_format.get("duration", 0)) + fps_long = eval(video_stream.get("r_frame_rate", "0")) + fps = float(f"{fps_long:.3f}") if fps_long else 0.0 + size = int(video_format.get("size", 0)) + duration_time = datetime.timedelta(seconds=int(duration)) + bitrate = int(video_stream.get("bit_rate", video_format.get("bit_rate", 0))) + tags = video_format.get("tags", {}) + + try: + sample_w, sample_h = list( + map(int, video_stream["sample_aspect_ratio"].split(":")) + ) + sample_ar = (sample_w / sample_h) or 1.0 + except: + sample_ar = 1.0 + + return FfprobeResult( + duration_sec=duration, + duration_human=str(duration_time), + codec=codec, + fps=fps, + size_bytes=size, + width=width, + bitrate=bitrate, + height=height, + container=video_path.suffix.lstrip(".").lower(), + ar=width / height if height else 1.0, + sample_ar=sample_ar, + tags=tags, + ) + + +def hash_partial(f: Path) -> str: + sha1 = hashlib.sha1() + chunk_size = 1024 * 1024 * 10 # 10MB chunk size + + total_read = 0 + with f.open("rb") as file: + while chunk := file.read(chunk_size): + total_read += len(chunk) + sha1.update(chunk) + break # Only reads the first 10MB + + return f"sha1:{total_read}:{sha1.hexdigest()}" + + +def chunked(iterable, n): + it = iter(iterable) + while True: + chunk = tuple(itertools.islice(it, n)) + if not chunk: + return + yield chunk + + +def retry(max_attempts: int = 3, delay: float = 1.0): + def decorator(func): + @wraps(func) + def wrapper(*args, **kwargs): + last_exception = None + for attempt in range(1, max_attempts + 1): + try: + return func(*args, **kwargs) + except Exception as e: + last_exception = e + if attempt < max_attempts: + logging.warning( + f"Attempt {attempt} failed: {e}. Retrying in {delay}s..." + ) + time.sleep(delay) + else: + logging.error(f"All {max_attempts} attempts failed.") + raise last_exception + + return wrapper + + return decorator + + +def video_exists(con: psycopg.Connection, file_path: str) -> bool: + stmt = con.execute( + "UPDATE videos SET last_seen_at = NOW() WHERE file_path = %s RETURNING true AS exists", + [file_path], + ) + row = stmt.fetchone() + if not row: + return False + return row["exists"] + + +def save_video( + con: psycopg.Connection, + category: str, + file_path: str, + ffprobe_data: dict, + video_hash: str, +): + sql = """ + INSERT INTO videos(category, file_path, ffprobe, hash_partial) + VALUES (%s, %s, %s, %s) + ON CONFLICT(file_path) DO UPDATE SET + last_seen_at = NOW();""" + con.execute( + sql, + [ + category, + file_path, + json.dumps(ffprobe_data), + video_hash, + ], + ) + + +def save_error(con: psycopg.Connection, file_path: str, stderr: str): + sql = """INSERT INTO scan_state(file_path, stderr) VALUES (%s, %s) ON conflict do nothing""" + con.execute( + sql, + [ + file_path, + stderr, + ], + ) + + +def category_for_file(file_path: Path) -> str: + s = str(file_path) + if s.startswith("/Volumes/"): + return f"downloaded.{file_path.parts[1].lower()}" + if s.startswith("/mnt/box"): + return "box" + return "local" + + +def scan_videos(con: psycopg.Connection, video_paths: list[Path]) -> None: + new_videos = [] + for i, f in enumerate(video_paths): + if "/_picks/" in str(f): + continue + progress = f"[{i + 1}/{len(video_paths)}]" + if f.is_dir(): + # Recursively find files in directory + scan_videos(con, list(f.glob("**/*"))) + continue + + if f.suffix.lower() not in [".mp4", ".mkv", ".avi", ".mov", ".wmv", ".webm"]: + continue + + file_path = str(f.absolute()) + + if video_exists(con, file_path=file_path): + logging.debug(f"{progress} video already exists {f=}") + continue + + logging.info(f"{progress} probing {f=}") + try: + probed = ffprobe(f) + except Exception as e: + logging.error(f"Failed to probe {f}: {e}") + save_error(con, file_path=file_path, stderr=str(e)) + continue + + file_hash = hash_partial(f) + + save_video( + category=category_for_file(f), + con=con, + file_path=file_path, + ffprobe_data=probed.__dict__, + video_hash=file_hash, + ) + logging.info(f"{progress} saved: {f.name}") + new_videos.append(file_path) + con.commit() + + if new_videos: + parse_releases(con, new_videos) + + +def parse_releases(con: psycopg.Connection, file_paths: list[str]) -> None: + api_key = os.getenv("OPENROUTER_API_KEY") + if not api_key: + logging.warning("OPENROUTER_API_KEY not set, skipping filename parsing") + return + + # Chunk file_paths by 20 items to avoid prompt token limits + for i, chunk in enumerate(chunked(file_paths, 20), 1): + filenames = [Path(p).name for p in chunk] + logging.info( + f"Parsing filenames for {len(filenames)} videos (chunk {i}) via OpenRouter..." + ) + + try: + parsed_data = parse_filenames_with_ai(filenames, api_key) + for item in parsed_data: + actors = item.get("actors") + if not isinstance(actors, list) or len(actors) == 0: + continue + + fname = item.get("filename") + if not fname: + continue + + # Find the full path that matches this filename in the current chunk + full_path = next((p for p in chunk if Path(p).name == fname), None) + if not full_path: + continue + + # Update the release column + release_info = {k: v for k, v in item.items() if k != "filename"} + con.execute( + "UPDATE videos SET release = %s WHERE file_path = %s", + [json.dumps(release_info), full_path], + ) + con.commit() + logging.info(f"Successfully updated release info for chunk {i}") + except Exception as e: + logging.error( + f"Failed to parse filenames or update database for chunk {i}: {e}" + ) + + +@retry(max_attempts=3, delay=0.1) +def parse_filenames_with_ai(filenames: list[str], api_key: str) -> list[dict]: + system_prompt = """you are an expert in parsing file names. + +your task is to parse the actors, studio and date, title out of filenames give me a JSON array with these fields [{filename, studio, released_at, title, actors: [...]}}]. Omit the missing / empty fields. + +Actor names are usually 2 words (name and last name) but sometimes they only contain a single word. Ignore male names. +Title is the remaining part after studio, date, actors; and don't usually contain the actor names. Ignore the quality and category indicators. + +output only valid JSON without any wrappers or quotes. + +for example: +InTheCrack.E1890.Casey.Norhman.Provence.XXX.1080p.HEVC.x265.PRT.mp4 -> studio=InTheCrack, actors=["Casey Norhman"], title=E1890 +BlackedRaw.26.05.16.Agatha.Vega.And.Ella.Hughes.Knockout.Babes.Fuck.Two.Cops.On.Duty.XXX.1080p.HEVC.x265.PRT.torrent -> {studio=BlackedRaw, actors=["Agatha Vega", "Ella Hughes"], title="Knockout Babes Fuck Two Cops On Duty", released_at=2026-05-16} +StepSiblingsCaught.26.05.14.Nata.Gold.XXX.720p.HEVC.x265.PRT.mp4 -> studio=StepSiblingsCaught, released_at=2026-05-14, actors=["Nata Gold"] +HookupHotshot.26.02.06.Episode.453.Shrooms.Q.XXX.720p.HEVC.x265.PRT.mp4 -> studio=HookupHotshot, released_at=2026-02-06, actors=["Shrooms Q"], title="Episode 453" +""" + + user_prompt = "\n".join(filenames) + + response = httpx.post( + "https://openrouter.ai/api/v1/chat/completions", + headers={ + "Authorization": f"Bearer {api_key}", + }, + json={ + "model": "mistralai/ministral-3b-2512", + "prompt_cache_key": "file_parsing", + "messages": [ + {"role": "system", "content": system_prompt}, + {"role": "user", "content": user_prompt}, + ], + "temperature": 0, + }, + timeout=20, + ) + response.raise_for_status() + result = response.json() + content = result["choices"][0]["message"]["content"].strip() + + # Remove potential markdown code blocks + if content.startswith("```"): + content = re.sub(r"^```(?:json)?\s*|\s*```$", "", content, flags=re.MULTILINE) + + parsed_data = json.loads(content) + if not isinstance(parsed_data, list): + parsed_data = [parsed_data] + + # Clean output: omit empty strings and empty arrays + cleaned_data = [] + for entry in parsed_data: + cleaned_entry = { + k: v + for k, v in entry.items() + if v != "" and not (isinstance(v, list) and not v) + } + actors = cleaned_entry.get("actors", []) + if isinstance(actors, list) and len(actors) > 0: + actors = [a for a in actors if a and "@" not in a] + cleaned_entry["actors"] = actors + cleaned_data.append(cleaned_entry) + + return cleaned_data + + +@dataclasses.dataclass +@dataclasses.dataclass +class Request: + path: str + method: str + payload: dict[str, Any] + query: dict[str, str] + environ: dict[str, Any] + + @classmethod + def from_environ(cls, environ: dict) -> "Request": + method = environ["REQUEST_METHOD"].upper() + + query = parse_qs(environ.get("QUERY_STRING", ""), keep_blank_values=True) + + payload = {} + content_type = environ.get("CONTENT_TYPE", "") + if method == "POST" and "application/json" in content_type: + try: + length = int(environ.get("CONTENT_LENGTH", 0)) + if length > 0: + payload = json.loads(environ["wsgi.input"].read(length)) + except (ValueError, json.JSONDecodeError): + pass + + return cls( + path=environ["PATH_INFO"], + method=method, + payload=payload, + query={k: v[0] for k, v in query.items()}, + environ=environ, + ) + + +@dataclasses.dataclass +class Response: + status: str + headers: list[tuple[str, str]] + body: Any + + @classmethod + def error(cls, message: str) -> "Response": + return cls( + status="400 Bad Request", + headers=[("Content-Type", "application/json")], + body=json.dumps({"error": message}) + "\n", + ) + + @classmethod + def from_html(cls, html: str) -> "Response": + return cls( + status="200 OK", + headers=[("Content-Type", "text/html; charset=utf-8")], + body=html, + ) + + @classmethod + def from_json(cls, data: Any) -> "Response": + def _default(obj): + if isinstance(obj, (datetime.date, datetime.datetime)): + return obj.isoformat() + raise TypeError(f"Object of type {type(obj)} is not JSON serializable") + + return cls( + status="200 OK", + headers=[("Content-Type", "application/json")], + body=json.dumps(data, default=_default) + "\n", + ) + + @classmethod + def from_exception(cls, e: Exception) -> "Response": + return cls( + status="500 Internal Server Error", + headers=[("Content-Type", "application/json")], + body=json.dumps({"error": str(e)}) + "\n", + ) + + +class TinyAPI: + def __init__(self, handlers: dict[str, HandlerFunc]): + self.routes = self._prepare_routes(handlers) + + def _prepare_routes(self, handlers: dict[str, HandlerFunc]): + routes = [] + for k, h in handlers.items(): + parts = k.split(maxsplit=1) + method = parts[0].upper() + path = parts[1].rstrip("/") or "/" + try: + routes.append((method, re.compile(f"^{path}$"), h)) + except re.error as e: + raise ValueError(f"Invalid regex pattern '{path}': {e}") + return routes + + @staticmethod + def find_free_port() -> int: + with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s: + s.bind(("", 0)) + return s.getsockname()[1] + + def __call__(self, environ: dict, start_response: Callable): + req = Request.from_environ(environ) + path = req.path.rstrip("/") or "/" + + handler = None + for method, pattern, h in self.routes: + if req.method == method and pattern.match(path): + handler = h + break + + if not handler: + res = Response.from_json({"error": "Not Found"}) + res.status = "404 Not Found" + else: + try: + res = handler(req) + except Exception as e: + logging.exception("Handler crash") + res = Response.from_exception(e) + + body = res.body if isinstance(res.body, bytes) else res.body.encode("utf-8") + headers = res.headers + [("Content-Length", str(len(body)))] + start_response(res.status, headers) + return [body] + + @contextlib.contextmanager + def serve(self, host: str = "localhost", port: int = 0): + if port == 0: + port = self.find_free_port() + + class LoggedRequestHandler(WSGIRequestHandler): + def log_message(self, format: str, *args: Any) -> None: + # args usually contains (request_line, status_code, size) + # We redirect to our configured logger instead of sys.stderr + logging.info("%s - %s", self.address_string(), format % args) + + server = make_server(host, port, self, handler_class=LoggedRequestHandler) + url = f"http://{host}:{port}" + + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + + logging.info(f"Serving on {url}") + try: + yield url + finally: + logging.info("Shutting down server...") + server.shutdown() + server.server_close() + thread.join(timeout=5) + + +@contextlib.contextmanager +def connect_db() -> typing.Generator[psycopg.Connection, typing.Any, typing.Any]: + with psycopg.connect( + "postgres://abdus:abdus@db.abdus.dev:5444/snot?sslmode=disable" + ) as conn: + conn.row_factory = psycopg.rows.dict_row + yield conn + + +def handle_search(req: Request, conn: psycopg.Connection) -> Response: + query: str = req.payload.get("query", "") + if not query: + return Response.error("Query is empty") + + where_sql, where_params = filter_to_where(query) + + query_stmt = pgsql.SQL( + """ + with q as (select + id, + filename_from_path(file_path) as file_name, + file_path, + actors, + size_bytes / 1048576 as size_mb, + (ffprobe->>'width') || 'x' || (ffprobe->>'height') as resolution, + duration_human, + created_at + from videos) + select * from q + where {filter} + order by created_at desc + """ + ).format(filter=pgsql.SQL(where_sql)) + with conn.cursor() as cursor: + rows = cursor.execute(query_stmt, where_params).fetchall() + return Response.from_json( + [ + { + **row, + "download_url": "https://u201686:T6672ICVoWAedECH@u201686.your-storagebox.de/files" + + row["file_path"].replace("/mnt/box/files", ""), + } + for row in rows + ] + ) + + +def tokenize_filter_groups(filter_query: str) -> list[list[str]]: + lexer = shlex.shlex(filter_query, posix=True, punctuation_chars="|") + lexer.commenters = "" + lexer.whitespace_split = True + tokens = list(lexer) + + groups: list[list[str]] = [[]] + for token in tokens: + if token == "|": + if groups[-1]: + groups.append([]) + continue + groups[-1].append(token) + + return [group for group in groups if group] + + +def token_to_condition(token: str) -> tuple[str | None, list[typing.Any]]: + if ":" in token: + key, value = token.split(":", 1) + key = key.strip().lower() + value = value.strip() + if not value: + return None, [] + + if key == "file_name": + return "file_name ILIKE %s", [f"%{value}%"] + + if key == "size_mb": + match = re.match( + r"^(>=|<=|>|<|=)?\s*(\d+(?:\.\d+)?)$", value, re.IGNORECASE + ) + if not match: + return None, [] + op = match.group(1) or "=" + num = float(match.group(2)) + return f"size_mb {op} %s", [num] + + if key == "actor": + return "actors @> %s::text[]", [[value]] + + return None, [] + + return "file_name ILIKE %s", [f"%{token}%"] + + +def filter_to_where(filter_query: str) -> tuple[str, list[typing.Any]]: + groups = tokenize_filter_groups(filter_query=filter_query) + if not groups: + return "TRUE", [] + + or_clauses: list[str] = [] + params: list[typing.Any] = [] + + for group in groups: + and_clauses: list[str] = [] + for token in group: + clause, clause_params = token_to_condition(token=token) + if not clause: + continue + and_clauses.append(clause) + params.extend(clause_params) + if and_clauses: + or_clauses.append("(" + " AND ".join(and_clauses) + ")") + + if not or_clauses: + return "TRUE", [] + + return " OR ".join(or_clauses), params + + +def handle_home(req: Request, query: str = "") -> Response: + template_file = Path(__file__).parent / "snot.html" + html = template_file.read_text() + if query: + injected_json = json.dumps({"query": query}) + html = f"\n" + html + return Response.from_html(html) + + +def handle_assets(req: Request, base_path: Path) -> Response: + asset_path = base_path / req.path.lstrip("/") + if not asset_path.exists() or not asset_path.is_file(): + return Response( + status="404 Not Found", + headers=[("Content-Type", "text/plain")], + body="Asset not found\n", + ) + + content_type = "text/plain" + if asset_path.suffix == ".js": + content_type = "application/javascript" + elif asset_path.suffix == ".css": + content_type = "text/css" + elif asset_path.suffix in [".html", ".htm"]: + content_type = "text/html" + elif asset_path.suffix == ".json": + content_type = "application/json" + + return Response( + status="200 OK", + headers=[("Content-Type", content_type)], + body=asset_path.read_text(), + ) + + +@contextlib.contextmanager +def run_htmlpopup(address: str): + exe_path = shutil.which("htmlpopup") + if not exe_path: + raise RuntimeError( + "htmlpopup executable not found in PATH. Please install it to use the HTML popup feature." + ) + proc = subprocess.Popen( + [exe_path, "--title", "s·m·u·t", address], text=True, stdout=subprocess.PIPE + ) + try: + yield + proc.wait() + finally: + proc.terminate() + + +@click.group(invoke_without_command=True) +@click.pass_context +def cli(ctx: click.Context): + """Snot - Video Browser and Scanner""" + if ctx.invoked_subcommand is None: + ctx.invoke(serve) + + +@cli.command() +@click.option( + "--port", + type=int, + default=int(getenv("PORT", "0")), + help="Port to run the server on", +) +@click.option("--query", type=str, default="", help="Initial query string") +def serve(port: int, query: str): + """Start the web UI""" + with connect_db() as conn: + handlers = { + "POST /search": partial(handle_search, conn=conn), + "GET /": partial(handle_home, query=query), + "GET /.*": partial(handle_assets, base_path=Path(__file__).parent), + } + api = TinyAPI(handlers=handlers) + + with api.serve(port=port) as server_url: + logging.info(f"Server running at {server_url}") + logging.info("Available endpoints:") + for path in handlers.keys(): + logging.info(f" {path}") + with run_htmlpopup(server_url): + logging.info("HTML popup started. Press Ctrl+C to stop.") + + +@cli.command() +@click.argument("paths", nargs=-1, type=Path) +def scan(paths: list[Path]): + """Scan video files and update the database""" + if not paths: + click.echo("No paths provided to scan.") + return + + with connect_db() as con: + scan_videos(con=con, video_paths=list(paths)) + + +@cli.command("mark-deletion") +@click.argument("paths", nargs=-1, type=Path) +def mark_deletion(paths: list[Path]): + """Mark video files for deletion by setting marked_for_deletion_at""" + if not paths: + click.echo("No paths provided.") + return + + with connect_db() as con: + for path in paths: + file_name = path.name + rows = con.execute( + "UPDATE videos SET marked_for_deletion_at = NOW() WHERE file_name = %s RETURNING file_path", + [file_name], + ).fetchall() + if rows: + for row in rows: + click.echo(f"Marked for deletion: {row['file_path']}") + else: + click.echo(f"Not found in database: {file_name}", err=True) + continue + + if not click.get_text_stream("stdin").isatty(): + continue + + stem = path.stem + siblings = con.execute( + "SELECT file_path, file_name FROM videos WHERE file_name != %s AND file_name LIKE %s AND marked_for_deletion_at IS NULL", + [file_name, f"{stem}.%"], + ).fetchall() + for sibling in siblings: + if click.confirm(f" Also mark {sibling['file_name']}?", default=False): + con.execute( + "UPDATE videos SET marked_for_deletion_at = NOW() WHERE file_path = %s", + [sibling["file_path"]], + ) + click.echo(f" Marked for deletion: {sibling['file_path']}") + + con.commit() + + +if __name__ == "__main__": + cli()