arthur/clips/views.py
Olive Vaughn 6a53adb5e0 Projects live at URLs, have owners, and are edited together live
A project is only ever at /p/<id>/<slug>; / is the index of the projects
you own or edit. Every project has an owner, who can name editors;
anyone with the link can view. Every edit saves itself, one request in
flight at a time, as a patch of the leaves that changed, and a websocket
(channels + daphne) carries presence and each committed write to
everyone else in the project. The first write to a leaf wins, and the
loser is told.

Undo is per person: a step undoes only if the leaves it touched still
hold what it left, so it never takes a collaborator's work with it.
Named snapshots replace saving, and restore as an ordinary write.

An empty symbol now survives the leaf round trip with `:nodes {}`.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-29 22:04:03 -04:00

1056 lines
42 KiB
Python

"""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, 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 _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/<cid>/symbol/<sid>`, 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"])
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)
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 `<video>` element seeks this URL, and a media element
that is handed 200OK with no `Accept-Ranges` cannot seek: it reports an empty
`seekable` range, every `currentTime` write is a no-op, and detection then runs
ninety times over frame one without anything raising. Django's `FileResponse`
does not do this for us — there is no Range handling anywhere in it — so the
absence of these thirty lines presents as "MediaPipe's video mode is broken".
"""
try:
row = Blob.objects.get(digest=digest)
path = blobs.path_for(digest)
except (Blob.DoesNotExist, ValueError):
return JsonResponse({"error": "no such blob"}, status=404)
size = path.stat().st_size
span = _byte_range(request.headers.get("Range"), size)
if span is False:
response = HttpResponse(status=416)
response["Content-Range"] = f"bytes */{size}"
elif span is None:
response = FileResponse(open(path, "rb"), content_type=row.media_type)
else:
start, end = span
handle = open(path, "rb")
handle.seek(start)
response = FileResponse(_Slice(handle, end - start + 1),
status=206, content_type=row.media_type)
response["Content-Range"] = f"bytes {start}-{end}/{size}"
response["Content-Length"] = str(end - start + 1)
response["Accept-Ranges"] = "bytes"
response["Cache-Control"] = "public, max-age=31536000, immutable"
response["ETag"] = f'"{digest}"'
return response
# ---------------------------------------------------------------------------
# tier 2: analyses and blocks
@require_http_methods(["POST"])
def analyses(request):
"""Register an analysis artifact's identity. Idempotent: the same inputs are
the same key are the same row."""
try:
data = _body(request)
key = data.get("key")
descriptor = data.get("descriptor")
parsed = _check_key(key, descriptor)
for field in ("detector", "version"):
if not parsed.get(field):
raise Bad(
f"the descriptor does not declare a {field}: a cache key that "
"omits the detector version lets a model upgrade silently reuse "
"old landmarks",
missing=field,
)
footage = None
if data.get("footage"):
digest = str(data["footage"]).removeprefix("sha256:")
footage = Footage.objects.filter(digest=digest).first()
if footage is None:
raise Bad("the analysis names footage this server does not have",
footage=data["footage"])
row, created = Analysis.objects.get_or_create(
key=key,
defaults={
"descriptor": descriptor,
"detector": parsed["detector"],
"version": str(parsed["version"]),
"footage": footage,
},
)
return JsonResponse({"key": row.key, "created": created}, status=201 if created else 200)
except Bad as exc:
return _error(exc)
@require_http_methods(["GET", "PUT"])
def analysis_detail(request, key):
try:
row = Analysis.objects.get(key=key)
except Analysis.DoesNotExist:
return JsonResponse({"error": "no such analysis"}, status=404)
if request.method == "GET":
return JsonResponse({
"key": row.key, "descriptor": row.descriptor,
"detector": row.detector, "version": row.version,
"footage": str(row.footage_id) if row.footage_id else None,
"source_blocks": sorted(row.source_blocks.values_list("key", flat=True)),
})
try:
keys = _body(request).get("source_blocks")
roles = {"source/dense", "source/detected", "source/crops"}
if (not isinstance(keys, list) or not keys
or not all(isinstance(k, str) for k in keys) or len(set(keys)) != len(keys)):
raise Bad("an analysis needs distinct source block keys")
blocks = list(Block.objects.filter(key__in=keys))
if len(blocks) != len(keys) or any(b.analysis_id != key for b in blocks):
raise Bad("source blocks must exist and name this analysis")
by_subject = {}
for block in blocks:
subjects = json.loads(block.descriptor).get("features", [])
if (not isinstance(subjects, list) or len(subjects) > 1
or any(not isinstance(s, str) or not s for s in subjects)):
raise Bad("a source block must name one subject")
# Older single-face analyses used an empty feature list.
by_subject.setdefault(tuple(subjects), []).append(block.role)
if any(len(found) != len(roles) or set(found) != roles
for found in by_subject.values()):
raise Bad("each subject needs one block for each source role")
with transaction.atomic():
row = Analysis.objects.select_for_update().get(key=key)
current = set(row.source_blocks.values_list("key", flat=True))
if current and current != set(keys):
raise Bad("the source blocks of an analysis are immutable", status=409)
row.source_blocks.set(blocks)
return JsonResponse({"key": key, "source_blocks": sorted(keys)})
except Bad as exc:
return _error(exc)
@require_http_methods(["POST"])
def blocks_missing(request):
"""Which of these keys the server does not have.
The return on content addressing, as one request: a save uploads the blocks
that are new and nothing else, so re-saving a document after a knob-free edit
moves kilobytes.
"""
try:
keys = _body(request).get("keys") or []
if not isinstance(keys, list):
raise Bad("keys must be a list")
have = set(Block.objects.filter(key__in=keys).values_list("key", flat=True))
return JsonResponse({"missing": [k for k in keys if k not in have]})
except Bad as exc:
return _error(exc)
@require_http_methods(["POST"])
def blocks(request):
"""Store one dense block: its bytes, its optional absence mask, and the
descriptor its key is the hash of."""
try:
multipart = request.content_type == "multipart/form-data"
data = request.POST if multipart else _body(request)
upload = request.FILES.get("data") if multipart else None
state_upload = request.FILES.get("state") if multipart else None
key = data.get("key")
descriptor = data.get("descriptor")
parsed = _check_key(key, descriptor)
role = parsed.get("role")
if not role:
raise Bad("a block's descriptor names its role")
if not parsed.get("layout", {}).get("type"):
raise Bad(
"a block's descriptor must say what its elements are: an Int16Array "
"and a Float32Array over the same bytes are both valid readings and "
"only one of them is the block"
)
analysis_key = parsed.get("analysis")
analysis = Analysis.objects.filter(key=analysis_key).first()
if analysis is None:
raise Bad(
"this block names an analysis the server does not know; register the "
"analysis first, so that every stored block can name the detector "
"version that produced it",
analysis=analysis_key,
)
if not (upload and upload.size) and not data.get("data"):
raise Bad("a block with no bytes")
with transaction.atomic():
if role == "source/crops":
if upload:
data_blob = _crop_blob(upload.chunks())
else:
import base64
data_blob = _crop_blob([base64.b64decode(data["data"])])
else:
data_blob = _uploaded_blob(upload) if upload else _blob(data["data"])
row, created = Block.objects.get_or_create(
key=key,
defaults={
"descriptor": descriptor,
"role": role,
"analysis": analysis,
"data": data_blob,
"state": (_uploaded_blob(state_upload) if state_upload else
_blob(data["state"]) if data.get("state") else None),
},
)
return JsonResponse({"key": row.key, "created": created}, status=201 if created else 200)
except Bad as exc:
return _error(exc)
@require_http_methods(["GET"])
def block_detail(request, key):
import base64
try:
row = Block.objects.select_related("data", "state").get(key=key)
except Block.DoesNotExist:
return JsonResponse({"error": "no such block"}, status=404)
data = blobs.read(row.data.digest)
if row.data.media_type == blobs.CROP_MEDIA_TYPE:
data = zlib.decompress(data)
out = {
"key": row.key,
"descriptor": row.descriptor,
"data": base64.b64encode(data).decode("ascii"),
}
if row.state_id:
out["state"] = base64.b64encode(blobs.read(row.state.digest)).decode("ascii")
response = JsonResponse(out)
response["Cache-Control"] = "public, max-age=31536000, immutable"
return response
# ---------------------------------------------------------------------------
# tier 1: projects, clips, leaves
def _who(user):
return {"username": user.get_username() if user.is_authenticated else None}
def _project_json(project: Project, user):
leaves = list(project.leaves.all())
clips = []
for clip in project.clips.all():
prefix = f"clip/{clip.cid}/"
clips.append(
{
"cid": clip.cid,
"name": clip.name,
"footage": str(clip.footage_id) if clip.footage_id else None,
"analysis": clip.analysis_id,
"blocks": sorted(clip.blocks.values_list("key", flat=True)),
"leaves": {leaf.path: leaf.value for leaf in leaves if leaf.path.startswith(prefix)},
}
)
return {
"id": str(project.id),
"name": project.name,
"schema_version": project.schema_version,
"seq": project.seq,
"palette": project.palette,
"owner": project.owner.get_username(),
"editors": sorted(project.editors.values_list("username", flat=True)),
"can_edit": project.can_edit(user),
"clips": clips,
}
def _project(project_id):
try:
return Project.objects.select_related("owner").get(id=project_id)
except Project.DoesNotExist:
raise Bad("no such project", status=404)
def _writable(request, project_id):
project = _project(project_id)
if not project.can_edit(request.user):
raise Bad("only the owner and the editors can change this project; "
"save a copy instead", status=403)
return project
# ---------------------------------------------------------------------------
# who you are
#
# Django's session cookie, and the page's CSRF cookie on every write. Nothing
# here that a signed-in admin does not already have; the API gains a way in that
# is not the admin's login page.
@require_http_methods(["GET"])
def me(request):
return JsonResponse(_who(request.user))
@require_http_methods(["POST"])
def login(request):
data = json.loads(request.body or b"{}")
user = authenticate(request, username=data.get("username"), password=data.get("password"))
if user is None:
return JsonResponse({"error": "wrong username or password"}, status=400)
auth_login(request, user)
return JsonResponse(_who(user))
@require_http_methods(["POST"])
def signup(request):
data = json.loads(request.body or b"{}")
username = (data.get("username") or "").strip()
password = data.get("password") or ""
if not username or len(password) < 8:
return JsonResponse({"error": "a username, and a password of 8 or more"}, status=400)
User = get_user_model()
if User.objects.filter(username__iexact=username).exists():
return JsonResponse({"error": "that username is taken"}, status=409)
user = User.objects.create_user(username=username, password=password)
auth_login(request, user)
return JsonResponse(_who(user), status=201)
@require_http_methods(["POST"])
def logout(request):
auth_logout(request)
return JsonResponse(_who(request.user))
# ---------------------------------------------------------------------------
# tier 1: projects, clips, leaves
@require_http_methods(["GET", "POST"])
def projects(request):
"""GET lists what you own and are an editor of — nothing, signed out; POST
makes one, owned by you. Every project has an owner, so making one needs you
signed in."""
if request.method == "GET":
if not request.user.is_authenticated:
return JsonResponse({"projects": []})
visible = Q(owner=request.user) | Q(editors=request.user)
return JsonResponse(
{
"projects": [
{"id": str(p.id), "name": p.name,
"schema_version": p.schema_version, "seq": p.seq,
"owner": p.owner.get_username(),
"updated": p.updated.isoformat()}
for p in Project.objects.filter(visible).distinct()
.select_related("owner")[:100]
]
}
)
if not request.user.is_authenticated:
return JsonResponse({"error": "sign in to make a project"}, status=403)
try:
data = _body(request)
project = Project.objects.create(name=data.get("name") or "untitled", owner=request.user)
return JsonResponse(_project_json(project, request.user), status=201)
except Bad as exc:
return _error(exc)
@require_http_methods(["GET", "PUT"])
def project_detail(request, project_id):
try:
if request.method == "GET":
return JsonResponse(_project_json(_project(project_id), request.user))
project = _writable(request, project_id)
return _save(project, _body(request), request.user)
except Bad as exc:
return _error(exc)
@require_http_methods(["POST", "DELETE"])
def editors(request, project_id, username=None):
"""The owner names who else can write. POST {username} adds; DELETE
`editors/<username>` removes."""
try:
project = _project(project_id)
if not (request.user.is_authenticated and request.user.id == project.owner_id):
raise Bad("only the owner can change who edits", status=403)
if request.method == "POST":
username = _body(request).get("username")
user = get_user_model().objects.filter(username__iexact=username or "").first()
if user is None:
raise Bad(f"nobody is called {username!r}", status=404)
if request.method == "POST":
project.editors.add(user)
else:
project.editors.remove(user)
broadcast(project.id, {}, kind="access")
return JsonResponse({"editors": sorted(project.editors.values_list("username", flat=True))})
except Bad as exc:
return _error(exc)
@transaction.atomic
def _save(project: Project, data, user):
"""A save: one clip's leaves, written.
SCOPED BY CLIP, not by project. A payload that carries clip `a` does not
disturb clip `b`'s leaves.
Two shapes. Without `base`, a clip's leaves REPLACE that clip's leaves — the
whole-document save. With `base`, the seq the client last caught up to, the
save is a PATCH: `leaves` are the ones it changed, `removed` the ones it
deleted, and nothing it did not mention is touched. A leaf it names that
somebody else changed after `base`, to something else, is a conflict, and the
whole save answers 409 with their values — last-writer-wins per leaf, with the
loser told rather than silently clobbered. docs/architecture.md, "Make the
merge unit small instead of clever".
A leaf whose value is unchanged keeps its VERSION. That is what makes the
entity tag mean something: a save of a document where one channel moved
invalidates one leaf's etag, not all four hundred.
"""
base = data.get("base")
seq = project.bump()
if data.get("name"):
project.name = data["name"]
if data.get("palette"):
project.palette = data["palette"]
project.save(update_fields=["name", "palette"])
written, removed, unchanged, conflicts, deltas = [], [], [], {}, []
for spec in data.get("clips") or []:
cid = spec.get("cid")
if not cid:
raise Bad("every clip in a save names its cid")
leaves = spec.get("leaves") or {}
gone = spec.get("removed") or [] if base is not None else []
prefix = f"clip/{cid}/"
for path in [*leaves, *gone]:
if not path.startswith(prefix):
raise Bad(
f"leaf {path!r} is not addressed to clip {cid!r}",
clip=cid, path=path,
)
keys = spec.get("blocks") or []
have = set(Block.objects.filter(key__in=keys).values_list("key", flat=True))
if missing := [k for k in keys if k not in have]:
# Referential integrity across the tiers, enforced where it can be:
# a document that names blocks the server does not hold would load
# into a blank stage on any other machine.
raise Bad(
"this clip names tier-2 blocks the server does not have; upload them "
"before saving the document that points at them",
status=409, missing=missing,
)
existing = {leaf.path: leaf for leaf in project.leaves.filter(path__startswith=prefix)}
if base is not None:
for path in [*leaves, *gone]:
theirs = existing.get(path)
if theirs and theirs.seq > base and (
path not in leaves or theirs.value != leaves[path]):
conflicts[path] = theirs.value
if conflicts:
continue
analysis = Analysis.objects.filter(key=spec.get("analysis")).first()
footage = None
if spec.get("footage"):
footage = Footage.objects.filter(id=spec["footage"]).first()
clip, _ = Clip.objects.update_or_create(
project=project,
cid=cid,
defaults={"name": spec.get("name") or "", "analysis": analysis, "footage": footage},
)
blocks = Block.objects.filter(key__in=keys)
if base is None:
clip.blocks.set(blocks)
gone = [path for path in existing if path not in leaves]
else:
clip.blocks.add(*blocks)
changed = {}
for path, value in leaves.items():
leaf = existing.get(path)
if leaf is None:
Leaf.objects.create(project=project, path=path, value=value, seq=seq)
elif leaf.value != value:
leaf.value, leaf.seq = value, seq
leaf.version += 1
leaf.save(update_fields=["value", "version", "seq", "updated"])
else:
unchanged.append(path)
continue
changed[path] = value
dropped = [path for path in gone if path in existing]
project.leaves.filter(path__in=dropped).delete()
written += changed
removed += dropped
deltas.append({"cid": cid, "leaves": changed, "removed": dropped, "blocks": keys})
if conflicts:
raise Bad(
"somebody else changed these since you last caught up",
status=409, seq=seq - 1, conflicts=conflicts,
)
by = user.get_username() if user.is_authenticated else None
transaction.on_commit(lambda: broadcast(project.id, {
"seq": seq, "by": by, "name": project.name, "clips": deltas,
}))
return JsonResponse(
{
"id": str(project.id),
"schema_version": project.schema_version,
"seq": seq,
"written": sorted(written),
"removed": sorted(removed),
"unchanged": len(unchanged),
}
)
@require_http_methods(["GET", "PUT"])
def leaf_detail(request, project_id, leaf_path):
"""One leaf, conditionally.
`If-Match` and a 409 whose body carries the CURRENT value, so the client can
offer keep-mine / take-theirs. A PUT that replaced unconditionally is the bug
docs/architecture.md calls out in tl: the loser's work disappears silently, and
for a painted cel that is the class of bug that ends trust in a tool.
"""
try:
project = _project(project_id) if request.method == "GET" else _writable(request, project_id)
except Bad as exc:
return _error(exc)
leaf = project.leaves.filter(path=leaf_path).first()
if request.method == "GET":
if leaf is None:
return JsonResponse({"error": "no such leaf"}, status=404)
response = JsonResponse({"path": leaf.path, "value": leaf.value, "version": leaf.version})
response["ETag"] = leaf.etag
return response
try:
data = _body(request)
except Bad as exc:
return _error(exc)
if "value" not in data:
return _error(Bad("a leaf write carries a value"))
match = request.headers.get("If-Match")
with transaction.atomic():
seq = project.bump()
leaf = project.leaves.filter(path=leaf_path).first()
if leaf is None:
# ANY `If-Match` on a leaf that does not exist is a failed precondition,
# `*` included: RFC 7232 gives `*` the meaning "the resource must already
# exist", which is exactly the write a client makes when it believes it
# is editing something. Creating it instead would turn "somebody deleted
# this node" into a silent resurrection.
if match:
transaction.set_rollback(True)
return JsonResponse(
{"error": "no such leaf", "path": leaf_path}, status=409
)
leaf = Leaf.objects.create(project=project, path=leaf_path, value=data["value"], seq=seq)
else:
if match and match not in ("*", leaf.etag):
transaction.set_rollback(True)
response = JsonResponse(
{
"error": "stale write",
"path": leaf.path,
"version": leaf.version,
"value": leaf.value,
},
status=409,
)
response["ETag"] = leaf.etag
return response
leaf.value, leaf.seq = data["value"], seq
leaf.version += 1
leaf.save(update_fields=["value", "version", "seq", "updated"])
cid = leaf_path.split("/")[1] if leaf_path.startswith("clip/") else None
by = request.user.get_username() if request.user.is_authenticated else None
transaction.on_commit(lambda: broadcast(project.id, {
"seq": seq, "by": by, "name": project.name,
"clips": [{"cid": cid, "leaves": {leaf.path: leaf.value}, "removed": [], "blocks": []}],
}))
response = JsonResponse({"path": leaf.path, "version": leaf.version, "seq": seq})
response["ETag"] = leaf.etag
return response
@require_http_methods(["GET", "POST"])
def revisions(request, project_id):
"""Named snapshots: GET lists them, POST {summary} takes one of the document
as it is now."""
try:
project = _project(project_id) if request.method == "GET" else _writable(request, project_id)
except Bad as exc:
return _error(exc)
if request.method == "GET":
return JsonResponse(
{
"revisions": [
{"id": r.id, "seq": r.seq, "author": r.author, "summary": r.summary,
"created": r.created.isoformat(), "leaves": len(r.document)}
for r in project.revisions.all()[:100]
]
}
)
data = json.loads(request.body or b"{}")
revision = Revision.objects.create(
project=project,
seq=project.seq,
author=request.user.get_username(),
summary=(data.get("summary") or "").strip()[:500],
document={leaf.path: leaf.value for leaf in project.leaves.all()},
blocks={clip.cid: sorted(clip.blocks.values_list("key", flat=True))
for clip in project.clips.all()},
)
return JsonResponse({"id": revision.id, "seq": revision.seq,
"leaves": len(revision.document)}, status=201)
@require_http_methods(["POST"])
def restore(request, project_id, revision_id):
"""Put a snapshot back: an ordinary write of every leaf that differs, so
everybody in the room receives it the way they receive any other."""
try:
project = _writable(request, project_id)
revision = project.revisions.filter(id=revision_id).first()
if revision is None:
raise Bad("no such snapshot", status=404)
except Bad as exc:
return _error(exc)
with transaction.atomic():
seq = project.bump()
existing = {leaf.path: leaf for leaf in project.leaves.all()}
deltas = {}
def delta(path):
cid = path.split("/")[1]
return deltas.setdefault(cid, {"cid": cid, "leaves": {}, "removed": [],
"blocks": revision.blocks.get(cid, [])})
for path, value in revision.document.items():
leaf = existing.get(path)
if leaf is None:
Leaf.objects.create(project=project, path=path, value=value, seq=seq)
elif leaf.value != value:
leaf.value, leaf.seq = value, seq
leaf.version += 1
leaf.save(update_fields=["value", "version", "seq", "updated"])
else:
continue
delta(path)["leaves"][path] = value
gone = [path for path in existing if path not in revision.document]
project.leaves.filter(path__in=gone).delete()
for path in gone:
delta(path)["removed"].append(path)
for cid, keys in revision.blocks.items():
clip = project.clips.filter(cid=cid).first()
if clip:
clip.blocks.add(*Block.objects.filter(key__in=keys))
by = request.user.get_username()
transaction.on_commit(lambda: broadcast(project.id, {
"seq": seq, "by": by, "name": project.name, "clips": list(deltas.values()),
}))
return JsonResponse({"seq": seq, "changed": sum(len(d["leaves"]) + len(d["removed"])
for d in deltas.values())})