diff --git a/.gitignore b/.gitignore index ab2e70d..5f441bf 100644 --- a/.gitignore +++ b/.gitignore @@ -28,6 +28,8 @@ static/arthur/js/ # the Django half's own state: the document database, the content-addressed blob # store (tiers 2 and 3), and collectstatic's output db.sqlite3 +db.sqlite3-shm +db.sqlite3-wal /var/ # vim swap files diff --git a/clips/blobs.py b/clips/blobs.py index 7469e66..2216fcb 100644 --- a/clips/blobs.py +++ b/clips/blobs.py @@ -16,11 +16,13 @@ addressing that answers questions about work not yet done. import hashlib import os import tempfile +import zlib from pathlib import Path from django.conf import settings CHUNK = 1 << 20 +CROP_MEDIA_TYPE = "application/zlib" def digest_bytes(data: bytes) -> str: @@ -86,6 +88,20 @@ def write_stream(chunks) -> tuple[str, int]: return digest.hexdigest(), size +def write_compressed_stream(chunks) -> tuple[str, int]: + """Store a losslessly compressed stream; the digest names stored bytes.""" + compressor = zlib.compressobj() + + def compressed(): + for chunk in chunks: + if part := compressor.compress(chunk): + yield part + if part := compressor.flush(): + yield part + + return write_stream(compressed()) + + def adopt(source: Path) -> tuple[str, int]: """Store a file already on disk, by hard link where the filesystem allows it. diff --git a/clips/management/commands/compress_crop_blocks.py b/clips/management/commands/compress_crop_blocks.py new file mode 100644 index 0000000..6367b81 --- /dev/null +++ b/clips/management/commands/compress_crop_blocks.py @@ -0,0 +1,57 @@ +"""Compress existing raw mouth crop blocks without changing their public bytes.""" + +import hashlib +import zlib + +from django.core.management.base import BaseCommand, CommandError +from django.db import transaction +from django.db.models.deletion import ProtectedError + +from clips import blobs +from clips.models import Blob, Block + + +class Command(BaseCommand): + help = "Compress existing source/crops blobs and remove unreferenced raw copies" + + def handle(self, *args, **options): + converted = 0 + before = after = 0 + for block in Block.objects.filter(role="source/crops").select_related("data"): + old = block.data + if old.media_type == blobs.CROP_MEDIA_TYPE: + continue + old_digest = old.digest + with open(blobs.path_for(old_digest), "rb") as source: + digest, size = blobs.write_compressed_stream( + iter(lambda: source.read(blobs.CHUNK), b"") + ) + check = hashlib.sha256() + decompressor = zlib.decompressobj() + with open(blobs.path_for(digest), "rb") as compressed: + while chunk := compressed.read(blobs.CHUNK): + check.update(decompressor.decompress(chunk)) + check.update(decompressor.flush()) + if not decompressor.eof or check.hexdigest() != old_digest: + raise CommandError(f"crop compression failed verification: {block.key}") + with transaction.atomic(): + new, _ = Blob.objects.get_or_create( + digest=digest, + defaults={"size": size, "media_type": blobs.CROP_MEDIA_TYPE}, + ) + changed = Block.objects.filter(key=block.key, data=old).update(data=new) + if not changed: + continue + converted += 1 + before += old.size + after += size + if old_digest != new.digest: + try: + old.delete() + except ProtectedError: + pass + else: + blobs.path_for(old_digest).unlink(missing_ok=True) + self.stdout.write( + f"Compressed {converted} crop blocks: {before:,} -> {after:,} bytes" + ) diff --git a/clips/tests/test_api.py b/clips/tests/test_api.py index 411c0d6..da3b534 100644 --- a/clips/tests/test_api.py +++ b/clips/tests/test_api.py @@ -16,6 +16,7 @@ of the system rather than a convention in ClojureScript: The rest is the load/save round trip, the conditional write, and the footage manifest that makes the frames the backend's to serve. """ +import base64 import hashlib import json import shutil @@ -23,11 +24,13 @@ import struct import subprocess import tempfile import zlib +from io import StringIO from pathlib import Path from unittest import skipUnless from unittest.mock import Mock, patch from django.core.files.uploadedfile import SimpleUploadedFile +from django.core.management import call_command from django.test import TestCase, override_settings from clips import blobs, extraction @@ -215,6 +218,47 @@ class Tier2Tests(TestCase): self.assertEqual("AAE=", fetched["state"]) self.assertEqual(descriptor, fetched["descriptor"]) + @override_settings(DATA_UPLOAD_MAX_MEMORY_SIZE=1024, FILE_UPLOAD_MAX_MEMORY_SIZE=1024) + def test_large_block_upload_streams_past_json_body_limit(self): + analysis = self.register_analysis() + descriptor = block_descriptor(analysis, role="source/crops") + key = key_for(descriptor) + payload = bytes(range(256)) * 16 + response = self.client.post("/api/blocks", { + "key": key, + "descriptor": descriptor, + "data": SimpleUploadedFile("block.bin", payload), + "state": SimpleUploadedFile("state.bin", b"\x00\x01"), + }) + self.assertEqual(201, response.status_code, response.content) + row = Block.objects.get(key=key) + self.assertEqual(blobs.CROP_MEDIA_TYPE, row.data.media_type) + self.assertLess(row.data.size, len(payload)) + self.assertEqual(payload, zlib.decompress(blobs.read(row.data_id))) + self.assertEqual(base64.b64encode(payload).decode(), + self.client.get(f"/api/blocks/{key}").json()["data"]) + self.assertEqual(b"\x00\x01", blobs.read(row.state_id)) + + def test_existing_raw_crop_block_is_compressed_without_changing_its_key_or_read(self): + analysis = self.register_analysis() + descriptor = block_descriptor(analysis, role="source/crops") + key = key_for(descriptor) + payload = b"raw crop pixels" * 100 + digest, size = blobs.write(payload) + old = Blob.objects.create(digest=digest, size=size) + Block.objects.create(key=key, descriptor=descriptor, role="source/crops", + analysis_id=analysis, data=old) + + call_command("compress_crop_blocks", stdout=StringIO()) + row = Block.objects.select_related("data").get(key=key) + self.assertEqual(blobs.CROP_MEDIA_TYPE, row.data.media_type) + self.assertEqual(base64.b64encode(payload).decode(), + self.client.get(f"/api/blocks/{key}").json()["data"]) + self.assertFalse(blobs.path_for(digest).exists()) + compressed_digest = row.data_id + call_command("compress_crop_blocks", stdout=StringIO()) + self.assertEqual(compressed_digest, Block.objects.get(key=key).data_id) + def test_a_block_whose_analysis_is_unknown_is_refused(self): descriptor = block_descriptor("sha256:" + "f" * 64) response = self.post("/api/blocks", { diff --git a/clips/views.py b/clips/views.py index b8a1c59..6203fe7 100644 --- a/clips/views.py +++ b/clips/views.py @@ -26,6 +26,7 @@ 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 @@ -98,6 +99,22 @@ def _blob(b64, media_type="application/octet-stream"): 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 @@ -456,7 +473,10 @@ def blocks(request): """Store one dense block: its bytes, its optional absence mask, and the descriptor its key is the hash of.""" try: - data = _body(request) + 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) @@ -478,17 +498,27 @@ def blocks(request): "version that produced it", analysis=analysis_key, ) - if not data.get("data"): + 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": _blob(data["data"]), - "state": _blob(data["state"]) if data.get("state") else None, + "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) @@ -504,10 +534,13 @@ def block_detail(request, key): 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(blobs.read(row.data.digest)).decode("ascii"), + "data": base64.b64encode(data).decode("ascii"), } if row.state_id: out["state"] = base64.b64encode(blobs.read(row.state.digest)).decode("ascii") diff --git a/docs/architecture.md b/docs/architecture.md index 86f3f6f..90d6536 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -444,9 +444,12 @@ must never run while the transport is moving. The mouth crops are the non-obvious entry, and they are what makes remote work possible at all. `extractTeeth` reads source pixels, so without them a -collaborator holding the analysis but not the 600 source PNGs cannot touch a -single teeth knob. A 40×30 crop is about 1.2KB; a 600-frame take is under a -megabyte against hundreds for the footage. +collaborator holding the analysis but not the source video cannot touch a +single teeth knob without decoding video again. These are RGBA crops: a 40×30 +crop is 4.8KB raw, and current 200-pixel-wide crops can total over 12MB for a +take. The blob store compresses `source/crops` losslessly with zlib; block reads +return the original pixels. Run `python manage.py compress_crop_blocks` once to +convert existing raw crop blobs and remove their unreferenced copies. ### Bake B — resolved geometry. For scale. @@ -573,7 +576,7 @@ POST /api/analyses {key, descriptor} idempotent GET /api/analyses/ metadata + source block keys PUT /api/analyses/ link dense landmarks, mask, crops POST /api/blocks/missing {keys} -> {missing} -POST /api/blocks {key, descriptor, data, state} +POST /api/blocks multipart: key, descriptor, data file, optional state file (JSON also accepted) GET /api/blocks/ GET /api/footage/ the manifest: video and stream URLs, audio, a URL per tracing still GET /blob/ immutable bytes, with byte ranges for video playback diff --git a/frontend/src/arthur/domain/wire.cljs b/frontend/src/arthur/domain/wire.cljs index 272c673..f16e453 100644 --- a/frontend/src/arthur/domain/wire.cljs +++ b/frontend/src/arthur/domain/wire.cljs @@ -15,8 +15,8 @@ a codec that returned a sorted map would work locally and stop working after one round trip, which is the failure the plain-map rule already prevents. - The bytes are separate and base64, because tier 2 is typed arrays and transit - has nothing to say about them. `channel/dense-at` reads a block as + The bytes are separate from transit: JSON block reads carry base64, while + uploads send binary file parts. `channel/dense-at` reads a block as `{:data :state }` and both sides of the wire must hold byte-for-byte the same array — a handle that names a sha256 has to name the bytes you actually hold." diff --git a/frontend/src/arthur/events/project.cljs b/frontend/src/arthur/events/project.cljs index 7f776f7..d9e1439 100644 --- a/frontend/src/arthur/events/project.cljs +++ b/frontend/src/arthur/events/project.cljs @@ -23,6 +23,7 @@ fx, which is the only thing in this namespace that is not pure." (:require [arthur.domain.clip :as clip] [arthur.domain.project :as project] + [arthur.domain.wire :as wire] [arthur.events.playback :as pb] [arthur.footage.store :as store] [arthur.flow.address :as address] @@ -38,6 +39,20 @@ (defn- block-keys [^js doc] (into-array (map #(.-key %) (array-seq (.-blocks doc))))) +(defn- block-bytes [value] + (if (string? value) + (wire/bytes-of value) + (js/Uint8Array. (.-buffer value) (.-byteOffset value) (.-byteLength value)))) + +(defn- block-form [^js block] + (let [form (js/FormData.)] + (.append form "key" (.-key block)) + (.append form "descriptor" (.-descriptor block)) + (.append form "data" (js/Blob. #js [(block-bytes (.-data block))]) "block.bin") + (when-let [state (.-state block)] + (.append form "state" (js/Blob. #js [(block-bytes state)]) "state.bin")) + form)) + (defn- upload-missing! "POST the blocks the server said it does not have, and nothing else. @@ -45,9 +60,9 @@ this and it made sqlite answer \"database is locked\" on a save — which reaches the page as a 500 with nothing wrong with the request. The backend was fixed too (WAL, and a busy timeout, in server/settings.py), and this stays sequential - anyway: the uploads are a few kilobytes each, nothing is waiting on them, and a - burst of parallel writes to buy nothing is how the same bug comes back the first - time a take has sixty blocks instead of eleven." + anyway: most uploads are small, and a burst of parallel writes to buy nothing + is how the same bug comes back the first time a take has sixty blocks instead + of eleven." [^js doc] (-> (http/POST "/api/blocks/missing" #js {:keys (block-keys doc)}) (.then (fn [^js answer] @@ -55,7 +70,9 @@ todo (filterv #(contains? missing (.-key ^js %)) (array-seq (.-blocks doc)))] (-> (reduce (fn [chain block] - (.then chain (fn [_] (http/POST "/api/blocks" block)))) + (.then chain + (fn [_] + (http/POST-form "/api/blocks" (block-form block))))) (js/Promise.resolve nil) todo) (.then (fn [_] (count todo))))))))) @@ -104,7 +121,7 @@ (.then (fn [_] (when (seq source-blocks) (-> (upload-missing! - #js {:blocks (source/wire-blocks source-blocks)}) + #js {:blocks (source/upload-blocks source-blocks)}) (.then (fn [_] (http/PUT (str "/api/analyses/" (:id analysis)) diff --git a/frontend/src/arthur/flow/source.cljs b/frontend/src/arthur/flow/source.cljs index 91a7190..c196577 100644 --- a/frontend/src/arthur/flow/source.cljs +++ b/frontend/src/arthur/flow/source.cljs @@ -75,6 +75,13 @@ #js {:key key :descriptor descriptor :data (wire/base64 data)})) roles))) +(defn upload-blocks [blocks] + (into-array + (map (fn [role] + (let [{:keys [key descriptor data]} (get blocks role)] + #js {:key key :descriptor descriptor :data data})) + roles))) + (defn unpack "The three block-detail responses -> inputs for measurement and freeze." [^js responses [width height]]