diff --git a/README.md b/README.md index c1db141..2b9af4e 100644 --- a/README.md +++ b/README.md @@ -22,7 +22,7 @@ This Server keeps track of how many images were displayed with the [Neko Fans](h ### Gallery query API -`POST /gallery/query` runs a constrained `gallery-dl` JSON query through the internal worker and caches the worker result in Redis. +`POST /gallery/query` runs a constrained `gallery-dl` JSON query through the internal worker, normalizes the result shape, and caches the normalized result in Redis. Request body: @@ -39,26 +39,51 @@ Fields: | Field | Required | Description | | -------- | -------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | -| `source` | yes | One of `danbooru`, `gelbooru`, `safebooru`, `konachan`, or `yandere`. | +| `source` | yes, unless `url` is set | One of `danbooru`, `gelbooru`, `safebooru`, `konachan`, or `yandere`. | +| `url` | yes, unless `source` is set | Direct `https://` URL for any `gallery-dl` extractor. This cannot be combined with `source`, `rating`, or `tags`. | | `rating` | no | Convenience field for a booru rating tag such as `safe`, `questionable`, or `explicit`. `rating:safe` is also accepted and normalized to `safe`. The worker receives this as a normal `rating:` search tag. | | `tags` | no | Up to 20 tag strings. Tags may contain ASCII letters, numbers, `_`, `-`, and `:`. Existing clients may still pass rating filters as tags, for example `rating:safe`. | | `limit` | no | Number of gallery-dl items to request. Defaults to `50`; maximum is `100`. | -At least one query term is required: either `rating` or one `tags` entry. +For `source` queries, at least one query term is required: either `rating` or one `tags` entry. Response body: ```json -[] +{ + "items": [ + { + "url": "https://cdn.example/image.jpg", + "source": "danbooru", + "category": "danbooru", + "subcategory": "post", + "id": 123, + "title": null, + "filename": "image", + "extension": "jpg", + "file_url": "https://cdn.example/image.jpg", + "preview_url": "https://cdn.example/preview.jpg", + "sample_url": null, + "width": 1200, + "height": 900, + "rating": "s", + "score": 42, + "tags": ["cat_girl"], + "created_at": "2026-05-24T10:00:00Z", + "metadata": {} + } + ], + "errors": [] +} ``` -The response body is exactly the JSON value returned by `gallery-dl`. +The server consumes `gallery-dl -J` event rows such as `[2, metadata]`, `[3, url, metadata]`, and `[-1, error]`, then returns a stable object. Source-specific fields are preserved under `metadata`. Caching behavior: | Setting | Default | Description | | ------------------------------ | ------- | ---------------------------------------------------------------------------------------------------------------------------- | -| `GALLERY_DL_CACHE_TTL_SECONDS` | `900` | Redis result cache TTL for normalized `source` + query terms + `limit`. | +| `GALLERY_DL_CACHE_TTL_SECONDS` | `900` | Redis result cache TTL for normalized query target + `limit`. | | worker lock TTL | `30` | Short Redis lock used to avoid duplicate worker calls for the same cache key. Concurrent identical misses return HTTP `202`. | Successful responses are always backed by Redis: either `X-Server-Cache: HIT` or the worker result is written to Redis before returning `X-Server-Cache: MISS`. `X-Server-Cache-Ttl-Seconds` reports the remaining server-side TTL. HTTP responses include `Cache-Control: no-store` because this is a POST endpoint whose response varies by request body; caching is server-side in Redis, not browser/proxy caching. diff --git a/gallery-dl-worker/worker.py b/gallery-dl-worker/worker.py index 2277922..0855c39 100755 --- a/gallery-dl-worker/worker.py +++ b/gallery-dl-worker/worker.py @@ -1,6 +1,8 @@ import asyncio import json +import logging import os +import shlex from typing import Annotated from fastapi import FastAPI, HTTPException, Response @@ -12,6 +14,9 @@ MAX_OUTPUT_BYTES = 25 * 1024 * 1024 TIMEOUT_SECONDS = 20 MAX_CONCURRENT_JOBS = int(os.environ.get("MAX_CONCURRENT_JOBS", "8")) +logging.basicConfig(level=os.environ.get("LOG_LEVEL", "INFO")) +logger = logging.getLogger("gallery-dl-worker") + app = FastAPI(docs_url=None, redoc_url=None, openapi_url=None) job_semaphore = asyncio.Semaphore(MAX_CONCURRENT_JOBS) @@ -40,11 +45,13 @@ async def query(request: QueryRequest): output = await run_gallery_dl(request) if len(output) > MAX_OUTPUT_BYTES: + logger.error("gallery-dl output too large: %s bytes", len(output)) raise HTTPException(status_code=502, detail="gallery-dl output too large") try: json.loads(output) except json.JSONDecodeError as exc: + logger.error("invalid gallery-dl JSON output: %s", exc) raise HTTPException( status_code=502, detail="invalid gallery-dl output" ) from exc @@ -61,22 +68,35 @@ async def run_gallery_dl(request: QueryRequest) -> bytes: *request.args, request.url, ] + logger.info("running command: %s", shlex.join(command)) process = await asyncio.create_subprocess_exec( *command, stdout=asyncio.subprocess.PIPE, - stderr=asyncio.subprocess.DEVNULL, + stderr=asyncio.subprocess.PIPE, ) try: - stdout, _ = await asyncio.wait_for( + stdout, stderr = await asyncio.wait_for( process.communicate(), timeout=TIMEOUT_SECONDS ) except asyncio.TimeoutError as exc: process.kill() await process.wait() + logger.error("gallery-dl timed out after %s seconds", TIMEOUT_SECONDS) raise HTTPException(status_code=504, detail="gallery-dl timed out") from exc if process.returncode != 0: + logger.error( + "gallery-dl failed with exit code %s: %s", + process.returncode, + stderr.decode("utf-8", errors="replace").strip(), + ) raise HTTPException(status_code=502, detail="gallery-dl failed") + if stderr: + logger.warning( + "gallery-dl wrote to stderr: %s", + stderr.decode("utf-8", errors="replace").strip(), + ) + return stdout diff --git a/renovate.json b/renovate.json index f03e6e2..a98728b 100644 --- a/renovate.json +++ b/renovate.json @@ -1,18 +1,16 @@ { "$schema": "https://docs.renovatebot.com/renovate-schema.json", - "extends": [ - "config:recommended" - ], + "extends": ["config:recommended"], + "minimumReleaseAge": "2 days", + "internalChecksFilter": "strict", + "schedule": ["before 6am"], + "prConcurrentLimit": 1, + "prHourlyLimit": 1, "packageRules": [ { - "groupName": "Rust Dependencies", - "groupSlug": "rust-dependencies", - "matchManagers": [ - "cargo" - ], - "matchPackageNames": [ - "/.*/" - ] + "description": "Keep all dependency updates together so Renovate opens one PR.", + "groupName": "dependencies", + "matchPackageNames": ["*"] } ] } diff --git a/src/gallery_dl.rs b/src/gallery_dl.rs index 0d6a818..014420b 100644 --- a/src/gallery_dl.rs +++ b/src/gallery_dl.rs @@ -19,7 +19,8 @@ const COMMAND_TIMEOUT: Duration = Duration::from_secs(20); #[derive(Deserialize)] pub struct GalleryQuery { - source: String, + source: Option, + url: Option, tags: Option>, rating: Option, limit: Option, @@ -36,35 +37,62 @@ struct WorkerRequest<'a> { args: Vec, } +#[derive(Serialize)] +struct GalleryResponse { + items: Vec, + errors: Vec, +} + +#[derive(Serialize)] +struct GalleryItem { + url: Option, + source: Option, + category: Option, + subcategory: Option, + id: Option, + title: Option, + filename: Option, + extension: Option, + file_url: Option, + preview_url: Option, + sample_url: Option, + width: Option, + height: Option, + rating: Option, + score: Option, + tags: Option, + created_at: Option, + metadata: serde_json::Value, +} + +#[derive(Serialize)] +struct GalleryExtractorError { + error: Option, + message: Option, + metadata: serde_json::Value, +} + +struct GalleryTarget { + cache_source: String, + cache_terms: Vec, + url: String, +} + pub async fn query( - auth: Option, request: GalleryQuery, mut redis: ConnectionManager, ) -> Result { - if !is_authorized(auth.as_deref()) { - return Ok(json_error(StatusCode::UNAUTHORIZED, "unauthorized")); - } - let cache_ttl_seconds = cache_ttl_seconds(); let limit = request.limit.unwrap_or(50); - let source = match normalize_source(request.source) { - Ok(source) => source, - Err(message) => return Ok(json_error(StatusCode::BAD_REQUEST, message)), - }; - let query = match normalize_query(request.tags, request.rating) { - Ok(query) => query, - Err(message) => return Ok(json_error(StatusCode::BAD_REQUEST, message)), - }; if limit == 0 || limit > MAX_LIMIT { return Ok(json_error(StatusCode::BAD_REQUEST, "invalid limit")); } - - let url = match build_url(&source, &query.terms, limit) { - Some(url) => url, - None => return Ok(json_error(StatusCode::BAD_REQUEST, "unsupported source")), + let target = match gallery_target(request, limit) { + Ok(target) => target, + Err(message) => return Ok(json_error(StatusCode::BAD_REQUEST, message)), }; - let cache_key = cache_key(&source, &query.terms, limit); + let cache_key = cache_key(&target.cache_source, &target.cache_terms, limit); let cached: Result, _> = redis.get(&cache_key).await; let cached = match cached { Ok(value) => value, @@ -107,13 +135,14 @@ pub async fn query( } } - let output = match run_gallery_dl(&url, limit).await { + let output = match run_gallery_dl(&target.url, limit).await { Ok(value) => value, Err(message) => { let _: Result<(), _> = redis.del(lock_key).await; return Ok(json_error(StatusCode::BAD_GATEWAY, message)); } }; + let output = normalize_gallery_output(output); let serialized = match serde_json::to_string(&output) { Ok(serialized) => serialized, @@ -145,17 +174,46 @@ pub async fn query( )) } -fn is_authorized(auth: Option<&str>) -> bool { - let Ok(token) = env::var("GALLERY_DL_TOKEN") else { - return true; +struct NormalizedQuery { + terms: Vec, +} + +fn gallery_target(request: GalleryQuery, limit: u16) -> Result { + if let Some(url) = request.url { + if request.source.is_some() || request.tags.is_some() || request.rating.is_some() { + return Err("url cannot be combined with source, tags, or rating"); + } + let url = normalize_url(url)?; + return Ok(GalleryTarget { + cache_source: "url".to_string(), + cache_terms: vec![url.clone()], + url, + }); + } + + let source = match request.source { + Some(source) => normalize_source(source)?, + None => return Err("missing source or url"), + }; + let query = normalize_query(request.tags, request.rating)?; + let url = match build_url(&source, &query.terms, limit) { + Some(url) => url, + None => return Err("unsupported source"), }; - let expected = format!("Bearer {}", token); - auth == Some(expected.as_str()) + Ok(GalleryTarget { + cache_source: source, + cache_terms: query.terms, + url, + }) } -struct NormalizedQuery { - terms: Vec, +fn normalize_url(url: String) -> Result { + let url = url.trim().to_string(); + if url.len() > 2048 || !url.starts_with("https://") || url.contains('\0') { + return Err("invalid url"); + } + Ok(url) } fn normalize_source(source: String) -> Result { @@ -291,6 +349,111 @@ fn cache_key(source: &str, tags: &[String], limit: u16) -> String { format!("gallery:v1:result:{}", hex) } +fn normalize_gallery_output(value: serde_json::Value) -> GalleryResponse { + let mut items = Vec::new(); + let mut errors = Vec::new(); + let mut pending_metadata: Option = None; + + let serde_json::Value::Array(events) = value else { + return GalleryResponse { items, errors }; + }; + + for event in events { + let serde_json::Value::Array(fields) = event else { + continue; + }; + let Some(code) = fields.first().and_then(serde_json::Value::as_i64) else { + continue; + }; + + match code { + -1 => { + if let Some(metadata) = fields.get(1).cloned() { + errors.push(gallery_error(metadata)); + } + } + 2 => { + pending_metadata = fields.get(1).cloned(); + } + 3 => { + let url = fields + .get(1) + .and_then(serde_json::Value::as_str) + .map(str::to_string); + let metadata = fields + .get(2) + .cloned() + .or_else(|| pending_metadata.clone()) + .unwrap_or(serde_json::Value::Null); + items.push(gallery_item(url, metadata)); + pending_metadata = None; + } + _ => {} + } + } + + if items.is_empty() { + if let Some(metadata) = pending_metadata { + items.push(gallery_item(None, metadata)); + } + } + + GalleryResponse { items, errors } +} + +fn gallery_item(url: Option, metadata: serde_json::Value) -> GalleryItem { + GalleryItem { + url, + source: metadata_string(&metadata, "category"), + category: metadata_string(&metadata, "category"), + subcategory: metadata_string(&metadata, "subcategory"), + id: metadata.get("id").cloned(), + title: metadata_string(&metadata, "title"), + filename: metadata_string(&metadata, "filename"), + extension: metadata_string(&metadata, "extension") + .or_else(|| metadata_string(&metadata, "file_ext")), + file_url: metadata_string(&metadata, "file_url"), + preview_url: metadata_string(&metadata, "preview_url") + .or_else(|| metadata_string(&metadata, "preview_file_url")), + sample_url: metadata_string(&metadata, "sample_url"), + width: metadata_u64(&metadata, "width").or_else(|| metadata_u64(&metadata, "image_width")), + height: metadata_u64(&metadata, "height") + .or_else(|| metadata_u64(&metadata, "image_height")), + rating: metadata_string(&metadata, "rating"), + score: metadata_i64(&metadata, "score"), + tags: metadata.get("tags").cloned(), + created_at: metadata.get("created_at").cloned(), + metadata, + } +} + +fn gallery_error(metadata: serde_json::Value) -> GalleryExtractorError { + GalleryExtractorError { + error: metadata_string(&metadata, "error"), + message: metadata_string(&metadata, "message"), + metadata, + } +} + +fn metadata_string(metadata: &serde_json::Value, key: &str) -> Option { + metadata + .get(key) + .and_then(serde_json::Value::as_str) + .map(str::to_string) +} + +fn metadata_u64(metadata: &serde_json::Value, key: &str) -> Option { + metadata + .get(key) + .and_then(|value| value.as_u64().or_else(|| value.as_str()?.parse().ok())) +} + +fn metadata_i64(metadata: &serde_json::Value, key: &str) -> Option { + metadata + .get(key) + .and_then(|value| value.as_i64().or_else(|| value.as_str()?.parse().ok())) +} + async fn run_gallery_dl(url: &str, limit: u16) -> Result { let worker_url = env::var("GALLERY_DL_WORKER_URL").map_err(|_| "GALLERY_DL_WORKER_URL is not configured")?; @@ -353,3 +516,42 @@ fn json_error(status: StatusCode, message: &str) -> Response { None, ) } + +#[cfg(test)] +mod tests { + use super::*; + use serde_json::json; + + #[test] + fn normalizes_download_events() { + let response = normalize_gallery_output(json!([ + [2, {"category": "danbooru", "id": 123, "image_width": 800, "image_height": 600}], + [3, "https://example.test/file.jpg", {"category": "danbooru", "id": 123, "file_ext": "jpg"}] + ])); + + assert_eq!(response.items.len(), 1); + assert_eq!(response.errors.len(), 0); + assert_eq!( + response.items[0].url.as_deref(), + Some("https://example.test/file.jpg") + ); + assert_eq!(response.items[0].category.as_deref(), Some("danbooru")); + assert_eq!(response.items[0].extension.as_deref(), Some("jpg")); + } + + #[test] + fn normalizes_error_events() { + let response = normalize_gallery_output(json!([[ + -1, + {"error": "AuthRequired", "message": "credentials missing"} + ]])); + + assert_eq!(response.items.len(), 0); + assert_eq!(response.errors.len(), 1); + assert_eq!(response.errors[0].error.as_deref(), Some("AuthRequired")); + assert_eq!( + response.errors[0].message.as_deref(), + Some("credentials missing") + ); + } +} diff --git a/src/lib.rs b/src/lib.rs index 4f96358..38299da 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -190,7 +190,6 @@ pub async fn init(port: u16) { let gallery_query = warp::path("gallery") .and(warp::path("query")) .and(warp::post()) - .and(warp::header::optional::("authorization")) .and(warp::body::content_length_limit(16 * 1024)) .and(warp::body::json()) .and(with_redis(redis.clone())) diff --git a/tests/test_gallery_dl.py b/tests/test_gallery_dl.py index 0ba0ce9..9467f61 100755 --- a/tests/test_gallery_dl.py +++ b/tests/test_gallery_dl.py @@ -68,9 +68,9 @@ def main() -> int: response, headers = run_query(args.server_url, query) if response is None: return 1 - if not isinstance(response, list): + if not is_normalized_gallery_response(response): print( - f"failed {query['name']}: expected gallery-dl JSON list", + f"failed {query['name']}: expected normalized gallery response", file=sys.stderr, ) return 1 @@ -85,7 +85,10 @@ def main() -> int: f"failed {query['name']}: expected server cache header", file=sys.stderr ) return 1 - print(f"ok {query['name']}: server-cache={headers.get('X-Server-Cache')}") + items = response.get("items", []) if isinstance(response, dict) else [] + print( + f"ok {query['name']}: server-cache={headers.get('X-Server-Cache')} items={len(items)}" + ) return 0 @@ -121,5 +124,20 @@ def post_json(url: str, value: dict): return json.loads(response.read()), response.headers +def is_normalized_gallery_response(response) -> bool: + if not isinstance(response, dict): + return False + if not isinstance(response.get("items"), list): + return False + if not isinstance(response.get("errors"), list): + return False + for item in response["items"]: + if not isinstance(item, dict): + return False + if "url" not in item or "metadata" not in item: + return False + return True + + if __name__ == "__main__": raise SystemExit(main())