"""The API's implementation. Two things in here are load-bearing and neither is Django. THE SERVER VERIFIES EVERY TIER-2 KEY IT IS HANDED. A key is the sha256 of a canonical descriptor, and this recomputes it and refuses a mismatch. That is what makes content addressing a property of the system rather than a convention in the client: nothing can store bytes under a name that does not describe them. It hashes THE TEXT IT WAS SENT rather than re-rendering the descriptor from parsed values, and that is the honest arrangement rather than a shortcut. JS prints an integral double as `1` and Python prints `1.0`, so a scheme where both sides re-render the numbers would disagree on the first parameter whose value happens to be whole — and the failure would be an upload that 409s with nothing wrong. The bytes are the contract; the schema on top of them is a convention, and the two fields this file actually reads out of that schema are checked separately. AND IT REFUSES A BLOCK WHOSE ANALYSIS IT DOES NOT KNOW. Every block descriptor names an analysis, and every analysis declares a detector and a VERSION. So the chain from a stored block to the model version that produced it cannot be broken by a client that forgot a step — which is the whole point of docs/architecture.md's insistence that the cache key include the detector version. A model upgrade that silently reused old landmarks would otherwise present as "the tool got worse", with no event to attach it to. """ import hashlib import json import re import zlib from functools import lru_cache from pathlib import Path from uuid import UUID from django.conf import settings from django.contrib.auth import authenticate, get_user_model from django.contrib.auth import login as auth_login, logout as auth_logout from django.core.exceptions import ValidationError from django.db import transaction from django.db.models import Q from django.http import FileResponse, HttpResponse, JsonResponse from django.shortcuts import render from django.views.decorators.http import require_http_methods from . import blobs, extraction from .consumers import broadcast from .models import Analysis, Block, Blob, Clip, Extraction, Footage, Leaf, Project, Revision, Sound, Source KEY_LENGTH = 71 # "sha256:" + 64 hex # --------------------------------------------------------------------------- # helpers def _body(request): try: return json.loads(request.body or b"{}") except json.JSONDecodeError as exc: raise Bad(f"the request body is not JSON: {exc}") from exc class Bad(Exception): """A 400 with a message, raised where the problem is noticed.""" def __init__(self, message, status=400, **detail): super().__init__(message) self.message = message self.status = status self.detail = detail def _error(exc: Bad): return JsonResponse({"error": exc.message, **exc.detail}, status=exc.status) def _check_key(key, descriptor): """The verification. A key is the sha256 of the descriptor stored beside it.""" if not isinstance(key, str) or len(key) != KEY_LENGTH or not key.startswith("sha256:"): raise Bad(f"not a content address: {key!r}") if not isinstance(descriptor, str) or not descriptor: raise Bad("a key without its descriptor addresses nothing") actual = hashlib.sha256(descriptor.encode("utf-8")).hexdigest() if actual != key[7:]: raise Bad( "the key is not the hash of its descriptor", status=409, expected=f"sha256:{actual}", given=key, ) try: return json.loads(descriptor) except json.JSONDecodeError as exc: raise Bad(f"the descriptor is not canonical JSON: {exc}") from exc def _blob(b64, media_type="application/octet-stream"): import base64 digest, size = blobs.write(base64.b64decode(b64)) blob, _ = Blob.objects.get_or_create( digest=digest, defaults={"size": size, "media_type": media_type} ) return blob def _uploaded_blob(upload, media_type="application/octet-stream"): digest, size = blobs.write_stream(upload.chunks()) blob, _ = Blob.objects.get_or_create( digest=digest, defaults={"size": size, "media_type": media_type} ) return blob def _crop_blob(chunks): digest, size = blobs.write_compressed_stream(chunks) blob, _ = Blob.objects.get_or_create( digest=digest, defaults={"size": size, "media_type": blobs.CROP_MEDIA_TYPE} ) return blob # --------------------------------------------------------------------------- # the page def _asset_version(relative): """A static file's modification time, for its URL. The stylesheet and the bundle are served with `Last-Modified` and nothing else, so a browser is free to keep a stale copy on heuristic freshness — and new JavaScript over an old stylesheet renders a pane the stylesheet has never heard of as bare elements. A version in the URL makes each edit a new URL.""" for root in settings.STATICFILES_DIRS: path = Path(root) / relative if path.exists(): return str(int(path.stat().st_mtime)) return "0" def page(request, project_id=None, slug=None): """The host page. This replaced `frontend/public/index.html` at step 9, and `:dev-http` in shadow-cljs.edn went away with it.""" return render(request, "clips/index.html", { "css_version": _asset_version("arthur/app.css"), "js_version": _asset_version("arthur/js/main.js"), }) # --------------------------------------------------------------------------- # the detector # # WHY THE SERVER ANSWERS THIS. The analysis key has to include the detector # version, and a version string in the client is a string somebody has to remember # to bump. The server serves the model, so it can hash the model — and then the # version is a fact about the bytes that produced the landmarks rather than a # claim about them. @lru_cache(maxsize=4) def _model_digest(path: str, mtime: float) -> str: return blobs.digest_file(Path(path)) def _package_version() -> str: pkg = Path(settings.BASE_DIR) / "frontend" / "package.json" try: deps = json.loads(pkg.read_text())["dependencies"] return deps["@mediapipe/tasks-vision"].lstrip("^~") except Exception: return "unknown" @require_http_methods(["GET"]) def detector(request): model = Path(settings.BASE_DIR) / "frontend" / "public" / "mediapipe" / "face_landmarker.task" if not model.exists(): # Honest rather than fatal: detection will fail at the MediaPipe boundary # with a better message than this one could give, and an analysis stamped # "unknown" is a take somebody can still look at and re-freeze later. return JsonResponse({"detector": "mediapipe", "version": "unknown", "model": None}) digest = _model_digest(str(model), model.stat().st_mtime) return JsonResponse( { "detector": "mediapipe", # The package version AND the model's own hash. Either alone can change # while the other does not, and both change the landmarks. "version": f"{_package_version()}+{digest[:16]}", "model": f"sha256:{digest}", } ) # --------------------------------------------------------------------------- # tier 3: footage # # THE MANIFEST NOW CARRIES URLS. It used to carry a directory and the loader built # `frames/0001.png` itself, which quietly made the frame layout a shared secret # between a shell script and a ClojureScript namespace. The server names every # frame instead, so uploaded video and command-line bundles produce the same # footage response without the client knowing where either stored its frames. @require_http_methods(["GET", "POST"]) def sources(request): if request.method == "GET": return JsonResponse({"sources": [ {"id": str(row.id), "filename": row.filename, "probe": row.probe} for row in Source.objects.order_by("-created")[:100] ]}) upload = request.FILES.get("file") if upload is None: return JsonResponse({"error": "upload a video as the file field"}, status=400) try: digest, size = blobs.write_stream(upload.chunks()) facts = extraction.probe(blobs.path_for(digest)) blob, _ = Blob.objects.get_or_create( digest=digest, defaults={"size": size, "media_type": upload.content_type or "video/mp4"}) row, created = Source.objects.get_or_create( blob=blob, defaults={"filename": Path(upload.name).name[:255], "probe": facts}) return JsonResponse({"id": str(row.id), "digest": digest, "filename": row.filename, "probe": row.probe, "created": created}, status=201 if created else 200) except (ValueError, OSError) as exc: return JsonResponse({"error": str(exc)}, status=400) def _sound_json(row): return {"id": str(row.id), "label": row.label or row.filename, "filename": row.filename, "duration": row.duration, "audio": f"/blob/{row.blob_id}"} def _relabel(request, row): """PATCH one asset's display name. A LABEL IS THE ONLY FIELD EITHER ROW LETS A CLIENT WRITE, and the body is read for that key alone. Footage is content-addressed and its frame count, rate and digest are facts about the bytes; an endpoint that merged whatever it was sent would let a rename quietly contradict them. Blank clears it, which puts the row back to the name it was uploaded under rather than leaving it nameless. """ data = _body(request) if "label" not in data: raise Bad("a rename needs a label") label = str(data["label"] or "").strip()[:200] if label != row.label: row.label = label row.save(update_fields=["label"]) return row @require_http_methods(["GET", "POST"]) def sounds(request): if request.method == "GET": return JsonResponse({"sounds": [_sound_json(row) for row in Sound.objects.order_by("-created")]}) upload = request.FILES.get("file") if upload is None: return JsonResponse({"error": "upload a sound as the file field"}, status=400) try: digest, size = blobs.write_stream(upload.chunks()) duration = extraction.probe_audio(blobs.path_for(digest)) blob, _ = Blob.objects.get_or_create( digest=digest, defaults={"size": size, "media_type": upload.content_type or "audio/mpeg"}) row, created = Sound.objects.get_or_create( blob=blob, defaults={"filename": Path(upload.name).name[:255], "duration": duration}) return JsonResponse(_sound_json(row), status=201 if created else 200) except (ValueError, OSError) as exc: return JsonResponse({"error": str(exc)}, status=400) @require_http_methods(["GET", "PATCH"]) def sound_detail(request, sound_id): try: row = Sound.objects.get(id=sound_id) except Sound.DoesNotExist: return JsonResponse({"error": "no such sound"}, status=404) try: if request.method == "PATCH": row = _relabel(request, row) except Bad as exc: return _error(exc) return JsonResponse(_sound_json(row)) def _extraction_json(row): return {"key": row.key, "source": str(row.source_id), "state": row.state, "progress": row.progress, "error": row.error, "footage": str(row.footage_id) if row.footage_id else None} @require_http_methods(["POST"]) def extractions(request): try: data = _body(request) source_id = data.get("source") if not source_id: raise Bad("an extraction needs a source id") try: source = Source.objects.get(id=UUID(str(source_id))) except (ValueError, ValidationError, Source.DoesNotExist): raise Bad("no such source", status=404) settings = data.get("settings") or {} if settings != {}: raise Bad("extraction currently keeps the source frame rate; settings must be empty") key = extraction.extraction_key(source, settings) row, _ = Extraction.objects.get_or_create( key=key, defaults={"source": source, "settings": settings}) if row.state != "done": extraction.enqueue(key) return JsonResponse(_extraction_json(row), status=202 if row.state != "done" else 200) except Bad as exc: return _error(exc) @require_http_methods(["GET"]) def extraction_detail(request, key): try: return JsonResponse(_extraction_json(Extraction.objects.get(key=key))) except Extraction.DoesNotExist: return JsonResponse({"error": "no such extraction"}, status=404) def _footage_json(footage: Footage, urls=True): out = { "id": str(footage.id), "label": footage.label or footage.source, "source": footage.source, "fps": footage.fps, "frames": footage.frames, "width": footage.width, "height": footage.height, "footage": f"sha256:{footage.digest}", "audio": f"/blob/{footage.audio.digest}", # The analysis source. `null` on footage ingested before the proxy # existed, which the loader reports as "re-extract this" rather than # failing somewhere inside MediaPipe. "video": f"/blob/{footage.video.digest}" if footage.video_id else None, # What the page actually decodes: one access unit per frame, no container. "stream": f"/blob/{footage.stream.digest}" if footage.stream_id else None, "feature-absence": footage.feature_absence or {}, } if urls: out["urls"] = [f"/blob/{f.blob.digest}" for f in footage.frame_set.select_related("blob")] return out @require_http_methods(["GET"]) def footage_list(request): return JsonResponse( {"footage": [_footage_json(f, urls=False) for f in Footage.objects.all()]} ) _SYMBOL_LEAF = re.compile(r"^clip/([^/]+)/symbol/([^/]+)$") def _transit_fields(value, *keys): """Top-level fields of a transit map leaf, by keyword name. A leaf's own facts are a small flat map, so no key repeats and transit's cache never stands in for one; anything else reads as absent.""" if not (isinstance(value, list) and value[:1] == ["^ "]): return {} pairs = dict(zip(value[1::2], value[2::2])) return {k: pairs.get(f"~:{k}") for k in keys} @require_http_methods(["GET"]) def symbols(request): """Every symbol in every saved project, for the pool's all-assets folder. Read off the leaf PATHS rather than by loading documents: a symbol's own leaf is `clip//symbol/`, so listing them is one query and no decoding beyond the name and length its value carries.""" rows = [] for leaf in Leaf.objects.filter(path__contains="/symbol/").select_related("project"): m = _SYMBOL_LEAF.match(leaf.path) if not m: continue fields = _transit_fields(leaf.value, "name", "frames") rows.append({ "project": str(leaf.project_id), "project_name": leaf.project.name, "cid": m.group(1), "symbol": m.group(2), "name": fields.get("name") or m.group(2).replace("~", "/"), "frames": fields.get("frames"), }) rows.sort(key=lambda r: (r["project_name"], r["project"], r["name"])) return JsonResponse({"symbols": rows}) @require_http_methods(["GET", "PATCH"]) def footage_detail(request, footage_id): try: footage = Footage.objects.select_related("audio", "video", "stream").get(id=footage_id) except Footage.DoesNotExist: return JsonResponse({"error": "no such footage"}, status=404) try: if request.method == "PATCH": footage = _relabel(request, footage) except Bad as exc: return _error(exc) return JsonResponse(_footage_json(footage)) _RANGE = re.compile(r"^bytes=(\d*)-(\d*)$") class _Slice: """A file, readable only up to `remaining` bytes from where it was seeked.""" def __init__(self, handle, remaining): self.handle, self.remaining = handle, remaining def read(self, size=-1): if self.remaining <= 0: return b"" if size < 0 or size > self.remaining: size = self.remaining data = self.handle.read(size) self.remaining -= len(data) return data def close(self): self.handle.close() def _byte_range(header, size): """One `Range` header -> (start, end) inclusive, or None for the whole blob. A syntactically broken header is NOT an error: RFC 9110 says an unparsable Range is ignored and the whole representation is sent, which is what a client that meant nothing by it wants. `False` is the third answer — a range that parses and cannot be satisfied — because that one is a 416. """ if not header: return None match = _RANGE.match(header.strip()) if not match or match.group(1) == "" and match.group(2) == "": return None first, last = match.group(1), match.group(2) if first == "": # `bytes=-500`: the LAST 500 bytes, which is a different question. length = int(last) if length == 0: return False return (max(0, size - length), size - 1) start = int(first) end = int(last) if last else size - 1 end = min(end, size - 1) if start >= size or start > end: return False return (start, end) @require_http_methods(["GET"]) def blob(request, digest): """Raw bytes, immutable, and serveable a slice at a time. `immutable` is not optimism here, it is the definition: the name IS the hash of the content, so a cached copy cannot be stale. That is what makes serving a take's frames out of this cheap enough to do on every load. RANGE IS NOT AN OPTIMISATION HERE, IT IS THE FEATURE. Since the analysis source became a video file, a `