Stream block uploads and compress cached mouth crops

This commit is contained in:
Olive Vaughn 2026-09-28 14:37:09 -04:00
parent 65ad67c129
commit a611b86c0d
9 changed files with 195 additions and 16 deletions

2
.gitignore vendored
View file

@ -28,6 +28,8 @@ static/arthur/js/
# the Django half's own state: the document database, the content-addressed blob # the Django half's own state: the document database, the content-addressed blob
# store (tiers 2 and 3), and collectstatic's output # store (tiers 2 and 3), and collectstatic's output
db.sqlite3 db.sqlite3
db.sqlite3-shm
db.sqlite3-wal
/var/ /var/
# vim swap files # vim swap files

View file

@ -16,11 +16,13 @@ addressing that answers questions about work not yet done.
import hashlib import hashlib
import os import os
import tempfile import tempfile
import zlib
from pathlib import Path from pathlib import Path
from django.conf import settings from django.conf import settings
CHUNK = 1 << 20 CHUNK = 1 << 20
CROP_MEDIA_TYPE = "application/zlib"
def digest_bytes(data: bytes) -> str: def digest_bytes(data: bytes) -> str:
@ -86,6 +88,20 @@ def write_stream(chunks) -> tuple[str, int]:
return digest.hexdigest(), size 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]: def adopt(source: Path) -> tuple[str, int]:
"""Store a file already on disk, by hard link where the filesystem allows it. """Store a file already on disk, by hard link where the filesystem allows it.

View file

@ -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"
)

View file

@ -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 The rest is the load/save round trip, the conditional write, and the footage
manifest that makes the frames the backend's to serve. manifest that makes the frames the backend's to serve.
""" """
import base64
import hashlib import hashlib
import json import json
import shutil import shutil
@ -23,11 +24,13 @@ import struct
import subprocess import subprocess
import tempfile import tempfile
import zlib import zlib
from io import StringIO
from pathlib import Path from pathlib import Path
from unittest import skipUnless from unittest import skipUnless
from unittest.mock import Mock, patch from unittest.mock import Mock, patch
from django.core.files.uploadedfile import SimpleUploadedFile from django.core.files.uploadedfile import SimpleUploadedFile
from django.core.management import call_command
from django.test import TestCase, override_settings from django.test import TestCase, override_settings
from clips import blobs, extraction from clips import blobs, extraction
@ -215,6 +218,47 @@ class Tier2Tests(TestCase):
self.assertEqual("AAE=", fetched["state"]) self.assertEqual("AAE=", fetched["state"])
self.assertEqual(descriptor, fetched["descriptor"]) 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): def test_a_block_whose_analysis_is_unknown_is_refused(self):
descriptor = block_descriptor("sha256:" + "f" * 64) descriptor = block_descriptor("sha256:" + "f" * 64)
response = self.post("/api/blocks", { response = self.post("/api/blocks", {

View file

@ -26,6 +26,7 @@ tool got worse", with no event to attach it to.
import hashlib import hashlib
import json import json
import re import re
import zlib
from functools import lru_cache from functools import lru_cache
from pathlib import Path from pathlib import Path
from uuid import UUID from uuid import UUID
@ -98,6 +99,22 @@ def _blob(b64, media_type="application/octet-stream"):
return blob 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 # the page
@ -456,7 +473,10 @@ def blocks(request):
"""Store one dense block: its bytes, its optional absence mask, and the """Store one dense block: its bytes, its optional absence mask, and the
descriptor its key is the hash of.""" descriptor its key is the hash of."""
try: 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") key = data.get("key")
descriptor = data.get("descriptor") descriptor = data.get("descriptor")
parsed = _check_key(key, descriptor) parsed = _check_key(key, descriptor)
@ -478,17 +498,27 @@ def blocks(request):
"version that produced it", "version that produced it",
analysis=analysis_key, 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") raise Bad("a block with no bytes")
with transaction.atomic(): 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( row, created = Block.objects.get_or_create(
key=key, key=key,
defaults={ defaults={
"descriptor": descriptor, "descriptor": descriptor,
"role": role, "role": role,
"analysis": analysis, "analysis": analysis,
"data": _blob(data["data"]), "data": data_blob,
"state": _blob(data["state"]) if data.get("state") else None, "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) 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) row = Block.objects.select_related("data", "state").get(key=key)
except Block.DoesNotExist: except Block.DoesNotExist:
return JsonResponse({"error": "no such block"}, status=404) 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 = { out = {
"key": row.key, "key": row.key,
"descriptor": row.descriptor, "descriptor": row.descriptor,
"data": base64.b64encode(blobs.read(row.data.digest)).decode("ascii"), "data": base64.b64encode(data).decode("ascii"),
} }
if row.state_id: if row.state_id:
out["state"] = base64.b64encode(blobs.read(row.state.digest)).decode("ascii") out["state"] = base64.b64encode(blobs.read(row.state.digest)).decode("ascii")

View file

@ -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 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 possible at all. `extractTeeth` reads source pixels, so without them a
collaborator holding the analysis but not the 600 source PNGs cannot touch a collaborator holding the analysis but not the source video cannot touch a
single teeth knob. A 40×30 crop is about 1.2KB; a 600-frame take is under a single teeth knob without decoding video again. These are RGBA crops: a 40×30
megabyte against hundreds for the footage. 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. ### Bake B — resolved geometry. For scale.
@ -573,7 +576,7 @@ POST /api/analyses {key, descriptor} idempotent
GET /api/analyses/<key> metadata + source block keys GET /api/analyses/<key> metadata + source block keys
PUT /api/analyses/<key> link dense landmarks, mask, crops PUT /api/analyses/<key> link dense landmarks, mask, crops
POST /api/blocks/missing {keys} -> {missing} 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/<key> GET /api/blocks/<key>
GET /api/footage/<id> the manifest: video and stream URLs, audio, a URL per tracing still GET /api/footage/<id> the manifest: video and stream URLs, audio, a URL per tracing still
GET /blob/<digest> immutable bytes, with byte ranges for video playback GET /blob/<digest> immutable bytes, with byte ranges for video playback

View file

@ -15,8 +15,8 @@
a codec that returned a sorted map would work locally and stop working after one 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. 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 The bytes are separate from transit: JSON block reads carry base64, while
has nothing to say about them. `channel/dense-at` reads a block as uploads send binary file parts. `channel/dense-at` reads a block as
`{:data <typed array> :state <Uint8Array>}` and both sides of the wire must hold `{:data <typed array> :state <Uint8Array>}` and both sides of the wire must hold
byte-for-byte the same array — a handle that names a sha256 has to name the byte-for-byte the same array — a handle that names a sha256 has to name the
bytes you actually hold." bytes you actually hold."

View file

@ -23,6 +23,7 @@
fx, which is the only thing in this namespace that is not pure." fx, which is the only thing in this namespace that is not pure."
(:require [arthur.domain.clip :as clip] (:require [arthur.domain.clip :as clip]
[arthur.domain.project :as project] [arthur.domain.project :as project]
[arthur.domain.wire :as wire]
[arthur.events.playback :as pb] [arthur.events.playback :as pb]
[arthur.footage.store :as store] [arthur.footage.store :as store]
[arthur.flow.address :as address] [arthur.flow.address :as address]
@ -38,6 +39,20 @@
(defn- block-keys [^js doc] (defn- block-keys [^js doc]
(into-array (map #(.-key %) (array-seq (.-blocks 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! (defn- upload-missing!
"POST the blocks the server said it does not have, and nothing else. "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 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 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 (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 anyway: most uploads are small, and a burst of parallel writes to buy nothing
burst of parallel writes to buy nothing is how the same bug comes back the first is how the same bug comes back the first time a take has sixty blocks instead
time a take has sixty blocks instead of eleven." of eleven."
[^js doc] [^js doc]
(-> (http/POST "/api/blocks/missing" #js {:keys (block-keys doc)}) (-> (http/POST "/api/blocks/missing" #js {:keys (block-keys doc)})
(.then (fn [^js answer] (.then (fn [^js answer]
@ -55,7 +70,9 @@
todo (filterv #(contains? missing (.-key ^js %)) todo (filterv #(contains? missing (.-key ^js %))
(array-seq (.-blocks doc)))] (array-seq (.-blocks doc)))]
(-> (reduce (fn [chain block] (-> (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) (js/Promise.resolve nil)
todo) todo)
(.then (fn [_] (count todo))))))))) (.then (fn [_] (count todo)))))))))
@ -104,7 +121,7 @@
(.then (fn [_] (.then (fn [_]
(when (seq source-blocks) (when (seq source-blocks)
(-> (upload-missing! (-> (upload-missing!
#js {:blocks (source/wire-blocks source-blocks)}) #js {:blocks (source/upload-blocks source-blocks)})
(.then (fn [_] (.then (fn [_]
(http/PUT (http/PUT
(str "/api/analyses/" (:id analysis)) (str "/api/analyses/" (:id analysis))

View file

@ -75,6 +75,13 @@
#js {:key key :descriptor descriptor :data (wire/base64 data)})) #js {:key key :descriptor descriptor :data (wire/base64 data)}))
roles))) 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 (defn unpack
"The three block-detail responses -> inputs for measurement and freeze." "The three block-detail responses -> inputs for measurement and freeze."
[^js responses [width height]] [^js responses [width height]]