1004 lines
56 KiB
Markdown
1004 lines
56 KiB
Markdown
# DarkStruct
|
||
|
||
A minimal Julia service that receives files over HTTP and hands them off to a
|
||
pool of worker threads for processing. The HTTP endpoint does no real work: it
|
||
spools each uploaded file to disk, pushes a lightweight reference onto a work
|
||
queue, and responds immediately, staying free to accept the next upload.
|
||
|
||
The per-file processing is deliberately thin. Each file runs through a small
|
||
neural-network classifier that labels it **known** (a file type resembling the
|
||
training set) or **unknown**, and logs the result. See "File classifier" below.
|
||
|
||
## Architecture
|
||
|
||
The pipeline is four stages, each with its own bounded queue and its own worker
|
||
pool (tuned independently, since classification is CPU-bound, known-file
|
||
enrichment is process-/IO-bound, content triage is cheap IO, and language
|
||
enrichment mixes CPU with a subprocess):
|
||
|
||
```
|
||
POST /upload (multipart)
|
||
│
|
||
▼
|
||
┌─────────────────┐ stream bytes to disk (never buffered)
|
||
│ HTTP handler │────────────────────────► data/spool/<uuid>-<name>
|
||
│ (streaming) │ the file's home for its
|
||
└────────┬─────────┘ enqueue reference whole time in flight
|
||
│ (non-blocking)
|
||
▼ │
|
||
202 + job IDs ▼
|
||
(503 if full) ┌────────────────────┐
|
||
│ stage-1 queue │ classification
|
||
└─────────┬──────────┘
|
||
│ dequeue
|
||
┌───────────────────┼───────────────────┐
|
||
▼ ▼ ▼
|
||
classify wkr 1 classify wkr 2 … classify wkr N
|
||
│
|
||
┌────────────┴────────────┐
|
||
:unknown :known the file itself never moves;
|
||
│ │ routing is the enqueue alone
|
||
│ enqueue (blocking) │ enqueue (blocking backpressure)
|
||
▼ ▼
|
||
┌────────────────────┐ ┌────────────────────┐
|
||
│ unknown queue │ │ known queue │ enrichment
|
||
└─────────┬──────────┘ └─────────┬──────────┘
|
||
│ dequeue │ dequeue
|
||
┌────────┼────────┐ ┌───────────┼───────────┐
|
||
▼ ▼ ▼ ▼ ▼ ▼
|
||
unk 1 unk 2 … unk K known wkr 1 known wkr 2 … known wkr M
|
||
│ binary-vs-text sniff │ exiftool → normalized sidecar
|
||
├─► data/binary/<uuid>-<name> success ──┴──► data/done/<uuid>-<name>
|
||
│ (terminal — moved) data/done/<uuid>-<name>.meta.json
|
||
│ (sidecar-first commit)
|
||
│ :text enqueue (blocking failure ───────► data/failed/<uuid>-<name>
|
||
▼ backpressure)
|
||
┌────────────────────┐
|
||
│ text queue │ language enrichment
|
||
└─────────┬──────────┘
|
||
│ dequeue
|
||
┌────────┼────────┐
|
||
▼ ▼ ▼
|
||
txt 1 txt 2 … txt P
|
||
│ Languages.jl (natural language) + github-linguist (programming language)
|
||
└─► data/text_done/<uuid>-<name> + data/text_done/<uuid>-<name>.meta.json
|
||
(sidecar-first commit)
|
||
```
|
||
|
||
A file is written once, into `data/spool/`, and stays there for its entire time
|
||
in the pipeline. Stages hand it on by enqueueing its small `Job` reference, never
|
||
by moving bytes: the queue holding the reference *is* the record of which stage
|
||
the file has reached. The only move is the last one, into a terminal sink
|
||
(`data/done/`, `data/text_done/`, `data/binary/`) or into `data/failed/` if a
|
||
worker throws. Stage 1 used to rename each file into `data/known/` or
|
||
`data/unknown/` first, and stage 3 into `data/text/`. Those three directories are
|
||
gone, and with them 11.6 µs per file, the equal of the classifier itself. Stage 1
|
||
now runs at ~85k files/s on one worker instead of ~35k.
|
||
|
||
Stages 2 (known-file enrichment) and 3 (content triage) run in parallel: stage 1
|
||
feeds both the known and unknown queues. Stage 3 in turn feeds stage 4 (language
|
||
enrichment) for every file it sorts as text.
|
||
|
||
Key properties:
|
||
|
||
- **Fast intake:** the queue only ever carries small references; file bytes live
|
||
on disk, so memory stays flat regardless of file size. This holds end to end:
|
||
intake streams each upload from the socket to the spool file a chunk at a
|
||
time (`FS_UPLOAD_CHUNK_BYTES`, default 64 KiB) rather than buffering the body,
|
||
and every worker reads only a bounded prefix. Measured: uploads of 256 MiB,
|
||
1 GiB and 2 GiB each grow resident memory by ~20 MiB, a flat line in file
|
||
size. See "Streaming intake" and "Benchmarking" below.
|
||
- **Backpressure:** each queue is bounded (default 1000). When the *intake* queue
|
||
is full, uploads get `503 Service Unavailable`. When the *known* queue is full,
|
||
the stage-1 worker blocks and retries (a classified file is never dropped).
|
||
- **Crash-resilient:** files survive on disk. On startup everything left in
|
||
`data/spool/` is re-enqueued (`recovered` in the log) and **replays from stage
|
||
1**. That is safe rather than merely tolerable: classification and the
|
||
binary/text sniff are pure functions of the file's bytes, and the terminal
|
||
commits rename with `force=true`, so a replayed file lands where it would have
|
||
landed and overwrites its own sidecar. Recovery *blocks* on a full queue rather
|
||
than dropping the excess, and runs with the worker pools already live, so a
|
||
backlog larger than one queue's capacity takes longer to re-drive but none of
|
||
it is abandoned.
|
||
|
||
The cost of replay is redoing stages a file had already cleared. That is a
|
||
property of the *queue*, not of the directory layout: the default queues are
|
||
in-process (`src/queue.jl`), so a crash destroys the only record of how far
|
||
each file got. Per-stage directories used to stand in for that record, at the
|
||
price of a rename per file per stage on the hot path: a permanent cost on
|
||
every file, to buy a cheaper restart. The seam in `src/queue.jl` is where that
|
||
is actually fixed: run with `FS_QUEUE_BACKEND=rabbitmq` and a job stays unacked
|
||
until its handler commits, so a restart resumes each file at the stage it had
|
||
reached instead of re-driving `spool/` from stage 1. See "Queue backends".
|
||
- **Graceful shutdown:** SIGINT (Ctrl-C) and SIGTERM (systemd/Docker/k8s `stop`)
|
||
both stop accepting uploads, then drain the stages *in order*: close the
|
||
stage-1 queue and wait out the classify workers (the only producer of the known
|
||
*and* unknown queues), then close those queues and wait out the enrich and
|
||
content-triage workers (content triage being the only producer of the text
|
||
queue), then close the text queue and wait out the language-enrichment workers.
|
||
(See "Shutdown" below for one cosmetic caveat on SIGTERM.)
|
||
- **Safe filenames:** client-supplied names are sanitized and prefixed with a
|
||
server-minted UUID before touching the filesystem (no path traversal).
|
||
|
||
### Streaming intake
|
||
|
||
The upload endpoint never holds a file in memory. Bytes go socket → spool file in
|
||
`FS_UPLOAD_CHUNK_BYTES` chunks, so resident memory per in-flight upload is set by
|
||
the chunk size, not the file size. A 2 GiB upload costs about what a 2 KiB one
|
||
does. Two pieces make that work, and both are deliberate:
|
||
|
||
- **`src/multipart.jl`, an incremental multipart parser.** HTTP.jl's
|
||
`parse_multipart_form` takes the *complete* body as a byte vector, so using it
|
||
means every file in the request is in memory at once (and copied again per
|
||
part). The reader here pulls fixed-size chunks and hands each part's bytes
|
||
straight to its spool file. Its interface is two calls in a loop
|
||
(`next_part!` then `write_part_body!`, or `skip_part_body!`), so the handler
|
||
keeps ordinary control flow instead of inverting into callbacks. The subtle
|
||
part is that a boundary delimiter can straddle two chunks, so the buffer always
|
||
retains the last `length(delimiter)-1` bytes; the test suite parses the same
|
||
body at chunk sizes from 1 byte upward to put that split at every offset.
|
||
- **`/upload` bypasses Oxygen's router.** Oxygen's root handler wraps
|
||
`HTTP.streamhandler`, which does `request.body = read(stream)` *before*
|
||
dispatching, even for an Oxygen `@stream` route, so no route can stream an
|
||
upload. `run` therefore passes its own `handler` to `serve`
|
||
(`root_stream_handler`), which intercepts `POST /upload` at the stream level and
|
||
delegates everything else to Oxygen unchanged. The trade-off: `/upload` is
|
||
absent from Oxygen's built-in metrics and docs.
|
||
|
||
Streaming also changes what the endpoint can promise. A buffered handler knows up
|
||
front how many files a request holds; this one discovers them as they arrive. So
|
||
when the intake queue fills mid-request it does not abandon the connection: it
|
||
stops spooling (discarding the remaining parts rather than writing files it can't
|
||
queue), drains the body, and answers `503` with the `accepted` list of whatever
|
||
got in first. Files already queued stay queued, and the client can retry the rest.
|
||
|
||
A client that hangs up mid-upload is treated as routine: the partial spool file is
|
||
removed (so restart recovery can never pick up a truncated upload as if it were
|
||
complete) and the event is logged `upload aborted by client`. One cosmetic caveat,
|
||
like the SIGTERM one below: when a request body is cut short, HTTP.jl's own
|
||
`closeread` logs an `EOFError` after the handler returns, because the connection
|
||
promised more bytes via `Content-Length` than arrived. It's harmless noise from
|
||
inside HTTP.jl: the partial file is already cleaned up and the connection closed.
|
||
|
||
### Metadata enrichment (stage 2)
|
||
|
||
Files the classifier labels **known** are handed to a second pool that extracts
|
||
metadata with [`exiftool`](https://exiftool.org/) (`exiftool -json -G -n`),
|
||
chosen because no native Julia library comes close to its multi-format coverage.
|
||
The output is normalized into a small, stable, documented schema and written as a
|
||
JSON **sidecar** next to the file in `data/done/`, e.g.
|
||
`data/done/<uuid>-<name>.meta.json`. The original bytes are never modified.
|
||
|
||
> **Prerequisite:** `exiftool` must be on `PATH` (e.g. `apt install
|
||
> libimage-exiftool-perl`). The server fails fast at startup if it's missing.
|
||
|
||
Sidecar top-level fields (all nullable, present only when available), plus the
|
||
complete raw `exiftool` object under `raw`:
|
||
|
||
| Field | Meaning |
|
||
|---|---|
|
||
| `id`, `original_name` | job id and client-supplied name |
|
||
| `file_type`, `mime_type` | e.g. `PDF` / `application/pdf` |
|
||
| `file_size` | bytes (authoritative, from intake, not exiftool) |
|
||
| `created_date`, `modified_date` | content timestamps |
|
||
| `author` | person (`Author`/`Artist`/`By-line`) |
|
||
| `created_by` | authoring app/tool (`Producer`/`CreatorTool`/`Creator`/`Software`/…) |
|
||
| `dimensions` | `{width, height}` for media |
|
||
| `duration` | seconds, for audio/video |
|
||
| `page_count` | for documents |
|
||
| `error` | set on a *degraded* sidecar (see below) |
|
||
| `raw` | full `exiftool` output |
|
||
|
||
Each normalized field is a coalesce over a priority list of exiftool tags
|
||
(`src/metadata.jl`); extend a field by appending tag names. If extraction fails
|
||
or `exiftool` times out (`FS_EXIFTOOL_TIMEOUT`, default 30s), the file still
|
||
completes to `data/done/` with a **degraded sidecar** (`file_size`/`file_type`
|
||
plus an `error` note) rather than being quarantined, because it's still a wanted
|
||
known file. Only genuine I/O errors (can't write the sidecar or move the file)
|
||
send it to `data/failed/`.
|
||
|
||
The sidecar is committed *before* the file is moved into `data/done/`, so a
|
||
file's presence there always implies its sidecar is already present; a crash in
|
||
between leaves only a harmless orphan sidecar, and recovery re-enriches
|
||
idempotently.
|
||
|
||
### Content triage (stage 3)
|
||
|
||
Files the classifier labels **unknown** are handed to a third pool that sorts
|
||
them into two coarse buckets so downstream tooling can treat them differently:
|
||
|
||
- **binary:** the file looks like binary data. This is the end of the live path,
|
||
so the file is committed to `data/binary/`, which doubles as the corpus for the
|
||
offline stage-5 discovery sweep below.
|
||
- **text:** the file looks like text. Stage 4 is still to come, so nothing
|
||
moves; the file stays in `data/spool/` and its reference goes onto the
|
||
stage-4 queue, ending up in `data/text_done/` once enriched.
|
||
|
||
The test is a **UTF-8 sniff**: read the first 8000 bytes and call the file text
|
||
when that window is valid UTF-8 and holds no control bytes outside the text-safe
|
||
set (tab, newline, CR, and friends, plus ESC for ANSI-colored logs); otherwise
|
||
binary. It's cheap (no full read) and Unicode-aware. Where the older NUL-byte and
|
||
printable-ASCII heuristics misfiled non-ASCII text, this keeps accents, CJK and
|
||
emoji in the text bucket, while binary formats, which rarely form valid UTF-8
|
||
near their start, still read as binary. A NUL byte is valid UTF-8 but not a text
|
||
control byte, so it too reads as binary. A multi-byte character split by the
|
||
8000-byte boundary is trimmed before the check so it isn't mistaken for malformed
|
||
bytes. An empty file is treated as text (`src/content.jl`).
|
||
|
||
### Language enrichment (stage 4)
|
||
|
||
Files that stage 3 sorts as **text** are handed to a fourth pool that identifies
|
||
their language and writes a `.meta.json` sidecar, mirroring the stage-2
|
||
known-file enrichment. Two detectors run per file:
|
||
|
||
- **natural language:** [`Languages.jl`](https://github.com/JuliaText/Languages.jl)'s
|
||
`LanguageDetector` (a Julia port of the `whatlang` n-gram model) reads a bounded
|
||
prefix (up to `LANG_SAMPLE_BYTES`, 64 KiB) and reports the language's English
|
||
name, ISO 639-3 code, and a confidence in `[0,1]`. Pure Julia, no subprocess.
|
||
The detector is built once at startup and shared read-only across the pool.
|
||
- **programming language:** the [`github-linguist`](https://github.com/github-linguist/linguist)
|
||
CLI recognizes source and markup by extension + content heuristics (e.g.
|
||
`Python`, `Markdown`). Plain prose reports as `Text` and unrecognized content as
|
||
`null`; both collapse to *no programming language*.
|
||
|
||
The sidecar schema:
|
||
|
||
| field | meaning |
|
||
|---|---|
|
||
| `id`, `original_name` | from intake |
|
||
| `file_size` | bytes (authoritative, from intake) |
|
||
| `content_type` | always `"text"` |
|
||
| `language` | natural-language English name (e.g. `English`), or `null` |
|
||
| `language_code` | ISO 639-3 code (e.g. `eng`), or `null` |
|
||
| `language_confidence` | detector confidence in `[0,1]`, or `null` |
|
||
| `programming_language` | e.g. `Python`, `Markdown`, or `null` |
|
||
| `error` | set if natural-language detection produced nothing |
|
||
|
||
> **`github-linguist` and the git-repo quirk:** run against a path *inside* a git
|
||
> repository, linguist reads the file's committed git blob, not the on-disk
|
||
> bytes, and an untracked file (which everything under `data/` is) has no blob,
|
||
> so it crashes. Stage 4 sidesteps this by copying each file to a fresh temp dir under
|
||
> `/tmp` (outside any repo, preserving the name so extension heuristics still
|
||
> fire) and pointing linguist there.
|
||
>
|
||
> Programming-language detection is best-effort: if `github-linguist` is
|
||
> missing (a startup warning, not a fatal error, unlike `exiftool`), fails, or
|
||
> times out (`FS_LINGUIST_TIMEOUT`, default 30s), `programming_language` is simply
|
||
> `null` and the file still completes. Natural-language detection failing produces
|
||
> a **degraded sidecar** (with an `error` note) rather than a quarantine, because
|
||
> the file is still wanted text.
|
||
|
||
Like stage 2, the sidecar is committed *before* the file is moved into
|
||
`data/text_done/`, so the file's presence there always implies its sidecar is
|
||
present; recovery re-enriches idempotently (`src/language.jl`).
|
||
|
||
### Unknown-format discovery (stage 5, offline)
|
||
|
||
The `binary/` sink from stage 3 is the pile of genuinely *unrecognized* files.
|
||
Stage 5 mines it for recurring new file formats by clustering files on their
|
||
header bytes into a growing catalog of discovered formats, each with a magic-byte
|
||
signature that can eventually be promoted into the classifier's fast path. Unlike
|
||
stages 1–4 it is not on the request hot path: it is a single-owner *batch*
|
||
process (the catalog is mutable shared state, the opposite of the stateless
|
||
classifier), and because promotion is human-gated nothing here is
|
||
latency-sensitive. The full rationale, and the assumptions we deliberately
|
||
rejected, live in [`model/DESIGN_clustering.md`](model/DESIGN_clustering.md).
|
||
|
||
The model (`src/cluster.jl`, base-Julia, no extra deps) is a Dirichlet-process
|
||
mixture of **per-position categoricals** over the first 32 header bytes, on a
|
||
257-symbol alphabet (byte `0–255` plus a `past-EOF` symbol so short fixed-length
|
||
formats are modeled honestly). Bytes are treated as categorical rather than
|
||
numeric (`0x89` and `0x88` are not "close"), so this deliberately does *not*
|
||
reuse the classifier's `[0,1]` byte scaling. A fixed uniform *background*
|
||
component absorbs structureless (compressed/encrypted) blobs so they don't mint spurious
|
||
clusters. A cluster's spiked positions become a libmagic-style signature;
|
||
clusters with enough members and enough fixed positions self-**nominate** for
|
||
promotion (a human does the one irreversible step, redefining "known").
|
||
|
||
**Status:** both phases are implemented and calibrated. Phase A (offline Gibbs)
|
||
is the science; phase B (`src/catalog.jl`) is the live catalog: a durable
|
||
single-owner state that sweeps `binary/`, folds each new file into a cluster with
|
||
the deterministic CRP-predictive rule, and writes promotion nominations.
|
||
|
||
The catalog process is run periodically (cron), single-threaded. It is the only
|
||
writer of the catalog, so it needs no locking:
|
||
|
||
```bash
|
||
julia --project=. bin/cluster_sweep.jl # incremental live sweep of new binary/ files
|
||
julia --project=. bin/cluster_sweep.jl --compact # offline Gibbs re-cluster (seed / recompact)
|
||
```
|
||
|
||
The catalog is a single durable file (`FS_CLUSTER_CATALOG`, default
|
||
`data/catalog.json`) committed with the same sidecar-first
|
||
temp→fsync→rename→fsync-dir discipline as the stage-2 sidecars, so a crash can
|
||
neither corrupt it nor lose a write. On the first run (empty catalog) the sweep
|
||
auto-promotes to a `--compact` pass to seed clusters; later runs assign
|
||
incrementally, touching only files they have not seen. A cluster that clears the
|
||
member/magic thresholds writes a nomination (a hex magic template, member count,
|
||
and example filenames) into `FS_NOMINATED_DIR` (default `data/nominated/`) for a
|
||
human to glance at and promote. Under the calibrated `bg_mass > α`, the live
|
||
sweep never mints single-file clusters; genuinely new formats surface from the
|
||
periodic `--compact` re-clustering of the background residue, not the live path.
|
||
|
||
Calibration is its own offline script, like training, and never in the request
|
||
path. It is scored against magic-collapsed ground truth (so `docx`≡`zip` and the whole
|
||
ELF family count as one format each, which is the *correct* answer, not an error):
|
||
|
||
```bash
|
||
julia --project=. bin/cluster_calibrate.jl [training_set_dir] # defaults to ../training_set
|
||
```
|
||
|
||
It grid-tunes the hyperparameters to maximize Adjusted Rand Index against known
|
||
formats and cross-checks against a model-free NCD (gzip) baseline. On the 700-file
|
||
training corpus the calibrated defaults (`n=32`, `α=1.0`, `β=0.1`) recover the
|
||
known formats at **ARI 0.77** (0.885 excluding tar), with `gzip`, `pkzip`
|
||
(`docx`+`zip` merged), and `jpeg` forming clean, promotable clusters; the NCD
|
||
baseline agrees. See `DESIGN_clustering.md` §11 for the full results, including the
|
||
one known limitation: ELF and these tarballs share a long run of header
|
||
zero-padding and merge. The v2 fix is inverse-entropy position weighting.
|
||
|
||
## Queue backends
|
||
|
||
The HTTP handler and the workers only ever call `enqueue!`, `dequeue!`, `ack!`,
|
||
`nack!` and `close!` on a `JobQueue` (`src/queue.jl`). Two implementations sit
|
||
behind those five methods, chosen at startup by `FS_QUEUE_BACKEND`:
|
||
|
||
| | `channel` (default) | `rabbitmq` |
|
||
|---|---|---|
|
||
| Where jobs live | in-process, bounded buffer | durable broker queues, persistent messages |
|
||
| Crash recovery | everything in `spool/` replays from stage 1 | each file resumes at the stage it had reached |
|
||
| External dependency | none | a RabbitMQ broker |
|
||
| Capacity | hard limit, enforced per enqueue | advisory, checked against a polled depth |
|
||
| Delivery | exactly once (nothing to redeliver) | at least once |
|
||
|
||
`ack!` is the whole difference. A job is settled only after its handler commits,
|
||
so a crash mid-enrichment leaves that job on the enrich queue and the restart
|
||
picks it up there — not at stage 1, and not lost. `ChannelQueue` implements
|
||
`ack!`/`nack!` as no-ops, which is honest rather than lazy: an in-process queue
|
||
has no delivery to settle, and its recovery story is `recover_dir!`.
|
||
|
||
### Running it
|
||
|
||
```bash
|
||
# broker + server, both in compose
|
||
docker compose -f docker-compose.yml -f docker-compose.rabbitmq.yml up
|
||
|
||
# just the broker, with the server on the host (5672 and the management UI on
|
||
# 15672 are published for exactly this)
|
||
docker compose -f docker-compose.yml -f docker-compose.rabbitmq.yml up -d rabbitmq
|
||
FS_QUEUE_BACKEND=rabbitmq \
|
||
FS_AMQP_URL=amqp://darkstruct:darkstruct@localhost:5672/ \
|
||
julia --project=. -t auto bin/server.jl
|
||
```
|
||
|
||
Queues are named `<FS_AMQP_PREFIX>.<stage>` — `darkstruct.classify`,
|
||
`.enrich`, `.triage`, `.language` — so the broker's queue list reads like the
|
||
pipeline, and `/stats` reports each one's depth from a once-a-second poll.
|
||
|
||
### What it guarantees, and what it deliberately doesn't
|
||
|
||
- **At least once, not exactly once.** Stages 1 and 3 publish downstream *then*
|
||
ack upstream, so a crash in that window redelivers a job that was already
|
||
routed. Safe for the same reason replay is: every stage is a pure function of
|
||
the file's bytes and every commit is idempotent. The one new case is that a
|
||
duplicate can run *concurrently* with the original and find the file already
|
||
committed; `worker_loop` recognises a vanished `job.path` and settles it
|
||
quietly, rather than quarantining a file that in fact succeeded.
|
||
- **Publisher confirms on intake only.** A `202` means the broker has the file,
|
||
because a client may delete its copy on the strength of it. That costs a round
|
||
trip per file and serializes intake publishes; `FS_AMQP_CONFIRMS=false` turns
|
||
it off. Inter-stage publishes are fire-and-forget by design — a lost one leaves
|
||
its source job unacked, and redelivery repairs it for free, on a handoff that
|
||
otherwise costs 0.12 µs.
|
||
- **Advisory capacity.** `enqueue!` compares `FS_QUEUE_CAPACITY` against a depth
|
||
polled once a second (and adjusted locally in between), so the existing `503`
|
||
and `blocked_ns` behaviour still works, but a burst can overshoot by up to a
|
||
poll interval. Making the limit real would mean `x-max-length` with
|
||
`overflow: reject-publish`, which needs a confirm per message to detect.
|
||
- **No reconnect.** A dropped connection invalidates every in-flight delivery
|
||
tag, so reconnecting would silently reprocess whatever the workers were
|
||
holding. Instead the server drains and exits non-zero, and the restart policy
|
||
brings it back to redelivered messages and an intact `spool/`.
|
||
- **No dead-letter queue.** A failed job is quarantined to `failed/` and then
|
||
acked: the file's story and the message's story end in the same place. A DLQ
|
||
would hold messages pointing at files that had already moved.
|
||
- **One consumer process.** `Job` is a claim check carrying a *local* path, so a
|
||
second server on another host would be handed jobs whose files it cannot see.
|
||
The broker buys durable resume across restarts, not horizontal scale; scaling
|
||
out would additionally need `spool/` on shared storage and an aggregated
|
||
`/stats`.
|
||
- **Restart recovery skips `spool/`.** The broker is the record of what is in
|
||
flight, so re-driving `spool/` would duplicate the whole backlog on every
|
||
restart. `FS_RECOVER_SPOOL=true` forces it, for the one case the broker cannot
|
||
cover: it was purged or recreated and the spooled files are all that is left.
|
||
|
||
### Tests
|
||
|
||
The pure parts (URL parsing, the wire format, backend selection, the duplicate
|
||
branch) run in the normal suite. The round trip against a real broker is gated:
|
||
|
||
```bash
|
||
docker compose -f docker-compose.yml -f docker-compose.rabbitmq.yml up -d rabbitmq
|
||
FS_TEST_AMQP_URL=amqp://darkstruct:darkstruct@localhost:5672/ \
|
||
julia --project=. -e 'using Pkg; Pkg.test()'
|
||
```
|
||
|
||
Without `FS_TEST_AMQP_URL` those tests are skipped with a notice, so the suite
|
||
still passes on a machine with no Docker.
|
||
|
||
## Running
|
||
|
||
```bash
|
||
# install deps (first time)
|
||
julia --project=. -e 'using Pkg; Pkg.instantiate()'
|
||
|
||
# external tools: exiftool (stage 2, required) and github-linguist (stage 4,
|
||
# optional, for programming-language detection). e.g. on Debian/Ubuntu:
|
||
# apt install libimage-exiftool-perl
|
||
# gem install github-linguist
|
||
|
||
# start the server; -t sets the number of OS threads available to workers
|
||
julia --project=. -t auto bin/server.jl
|
||
```
|
||
|
||
By default the pipeline's queues are in-process, and nothing external is needed.
|
||
For durable queues that survive a crash, run against RabbitMQ instead — see
|
||
"Queue backends" above.
|
||
|
||
## Shutdown
|
||
|
||
Both SIGINT and SIGTERM trigger the same idempotent graceful drain
|
||
(stop serving → close queue → wait for workers → exit):
|
||
|
||
- **SIGINT** is caught as an `InterruptException` (we call
|
||
`Base.exit_on_sigint(false)`), so shutdown is clean and quiet.
|
||
|
||
One caveat specific to the RabbitMQ backend: Julia delivers SIGINT to whichever
|
||
task happens to be running on thread 1, which is usually the main loop but is
|
||
not guaranteed to be — AMQPClient runs receiver tasks of its own, and an
|
||
interrupt that lands in one of those kills the broker connection instead of
|
||
reaching the main loop. That is not a stuck server: the depth poller notices
|
||
the closed connection within a second and asks for the same drain, so the
|
||
process still stops (within ~4s, exiting non-zero, and logging it as a lost
|
||
connection rather than an interrupt). **Under a process manager, prefer SIGTERM
|
||
for the broker backend** — it goes through `atexit`, which has no such lottery.
|
||
Shutdown on either signal also prints a few `Consumer ... task exiting` warnings
|
||
from AMQPClient, which are cosmetic: they are its consumer tasks noticing that
|
||
we cancelled them.
|
||
- **SIGTERM** can't be intercepted directly: Julia blocks it on worker threads
|
||
and handles it in its own runtime, so a user `signal()` handler never fires.
|
||
Instead we hook the drain into an `atexit` handler, which Julia's SIGTERM path
|
||
does run. Caveat: Julia prints its own `signal 15: Terminated` backtrace
|
||
*before* `atexit` runs. It's harmless noise, and the drain still completes
|
||
right after it, but for a fully quiet stop under a process manager,
|
||
configure it to send SIGINT instead (systemd: `KillSignal=SIGINT`; Docker:
|
||
`STOPSIGNAL SIGINT`). Give the stop timeout enough headroom to drain
|
||
in-flight work (systemd: `TimeoutStopSec`).
|
||
|
||
## File classifier
|
||
|
||
Each file is scored by a fixed-structure neural network (Lux.jl) that answers a
|
||
single binary question: is this file **known** (like the types in the training
|
||
set) or **unknown**? It's novelty detection, not exact file-typing: it won't
|
||
tell you "PDF", just "this looks like something I was trained on, or not".
|
||
|
||
- **Features:** the first 16 bytes + last 16 bytes of the file, each scaled
|
||
0–255 → `[0,1]`, giving a 32-dim input. Files under 32 bytes can't form that
|
||
window and are classified `unknown` without touching the model.
|
||
- **Architecture:** `Dense(32→64,relu) → Dense(64→16,relu) → Dense(16→2)`,
|
||
raw logits; decision is `argmax` (class 1 = known, class 2 = unknown).
|
||
- **Artifact:** trained weights live in `model/classifier.jld2` (committed), so
|
||
the server just loads them at startup. Missing/unreadable ⇒ the server fails
|
||
fast rather than run without classification.
|
||
- **Effect today:** *active routing*. The class is logged
|
||
(`classification=known|unknown`) and drives the pipeline split: `:known` files
|
||
go onto the stage-2 queue for metadata enrichment, `:unknown` files onto the
|
||
stage-3 queue for content triage. The class chooses the downstream stage (the
|
||
file itself stays in `data/spool/` either way); what's still unproven is the
|
||
model's *accuracy*, not whether the routing path runs.
|
||
|
||
The architecture and byte→feature mapping are defined once in `src/model.jl` and
|
||
shared by the trainer and the server, so they can't drift apart.
|
||
|
||
### Training
|
||
|
||
Training is a separate, offline script. It never runs in the request path:
|
||
|
||
```bash
|
||
julia --project=. bin/train.jl <positives_dir> [negatives_dir]
|
||
```
|
||
|
||
- **positives_dir:** every file in it (≥32 bytes) is a "known" example.
|
||
- **negatives_dir** *(optional)*: a grab-bag of *other* real file types used as
|
||
"unknown" examples. Negatives are generated ~1:1 with positives, split 50/50
|
||
between uniform-random byte vectors and grab-bag files. With no grab-bag dir,
|
||
negatives are all random (weaker: the net may just learn "high entropy =
|
||
unknown" rather than your actual types, so a grab-bag of real off-distribution
|
||
files is recommended).
|
||
|
||
The script uses an 80/20 seeded split, reports validation accuracy, and writes
|
||
`model/classifier.jld2` (path overridable via `FS_MODEL_PATH`). A fixed seed
|
||
(`FS_TRAIN_SEED`, default 42) drives negative generation, the split, and weight
|
||
init, so the artifact is exactly regenerable from the same inputs.
|
||
|
||
## Configuration (environment variables)
|
||
|
||
| Variable | Default | Meaning |
|
||
|---------------------|----------------|------------------------------------------|
|
||
| `FS_HOST` | `127.0.0.1` | Bind address |
|
||
| `FS_PORT` | `8080` | Port |
|
||
| `FS_WORKERS` | `nthreads()` | Stage-1 (classification) worker tasks |
|
||
| `FS_QUEUE_CAPACITY` | `1000` | Max pending intake jobs before `503` |
|
||
| `FS_KNOWN_WORKERS` | `nthreads()` | Stage-2 (enrichment) worker tasks |
|
||
| `FS_KNOWN_QUEUE_CAPACITY` | `1000` | Max pending enrichment jobs (then backpressure) |
|
||
| `FS_UNKNOWN_WORKERS` | `nthreads()` | Stage-3 (content triage) worker tasks |
|
||
| `FS_UNKNOWN_QUEUE_CAPACITY` | `1000` | Max pending triage jobs (then backpressure) |
|
||
| `FS_TEXT_WORKERS` | `nthreads()` | Stage-4 (language enrichment) worker tasks |
|
||
| `FS_TEXT_QUEUE_CAPACITY` | `1000` | Max pending language jobs (then backpressure) |
|
||
| `FS_SPOOL_DIR` | `data/spool` | Every in-flight file, at every stage |
|
||
| `FS_BINARY_DIR` | `data/binary` | Stage-3 sink: unknown files that look binary |
|
||
| `FS_DONE_DIR` | `data/done` | Enriched known files (+ `.meta.json`) |
|
||
| `FS_TEXT_DONE_DIR` | `data/text_done` | Enriched text files (+ `.meta.json`) |
|
||
| `FS_FAILED_DIR` | `data/failed` | Files whose processing threw |
|
||
| `FS_MODEL_PATH` | `model/classifier.jld2` | Classifier artifact loaded at startup |
|
||
| `FS_UPLOAD_CHUNK_BYTES` | `65536` | Socket read size at intake; bounds intake memory per in-flight upload |
|
||
| `FS_EXIFTOOL_TIMEOUT` | `30` | Seconds before a stuck exiftool is killed |
|
||
| `FS_LINGUIST_TIMEOUT` | `30` | Seconds before a stuck github-linguist is killed |
|
||
| `FS_CLUSTER_DIR` | `data/binary` | Stage-5 input: the unknown/binary pile to sweep |
|
||
| `FS_CLUSTER_N` | `32` | Header bytes modeled per file |
|
||
| `FS_CLUSTER_ALPHA` | `1.0` | CRP concentration (propensity to spawn new formats) |
|
||
| `FS_CLUSTER_PSEUDOCOUNT` | `0.1` | Dirichlet pseudocount β (calibrated) |
|
||
| `FS_CLUSTER_BG_MASS` | `5.0` | Mass of the uniform background component |
|
||
| `FS_PROMOTE_MIN_MEMBERS` | `20` | Cluster size threshold for promotion nomination |
|
||
| `FS_PROMOTE_MIN_MAGIC` | `3` | Required fixed signature positions to nominate |
|
||
| `FS_CLUSTER_CATALOG` | `data/catalog.json` | Durable stage-5 catalog file (single-owner) |
|
||
| `FS_NOMINATED_DIR` | `data/nominated` | One JSON per self-nominated cluster (human promote gate) |
|
||
| `FS_QUEUE_BACKEND` | `channel` | `channel` (in-process) or `rabbitmq` (durable); see "Queue backends" |
|
||
| `FS_AMQP_URL` | `amqp://guest:guest@localhost:5672/` | Broker connection, credentials included |
|
||
| `FS_AMQP_PREFIX` | `darkstruct` | Queue names are `<prefix>.<stage>` |
|
||
| `FS_AMQP_PREFETCH` | worker count | Unacked messages the broker hands one stage at a time |
|
||
| `FS_AMQP_CONFIRMS` | `true` | Wait for a publisher confirm before a `202` (intake only) |
|
||
| `FS_RECOVER_SPOOL` | `false` | Re-drive `spool/` at startup even on the broker backend (use after a purged broker) |
|
||
|
||
> To get real parallelism, start Julia with enough threads (`-t N`) to cover all
|
||
> pools. If `FS_WORKERS + FS_KNOWN_WORKERS + FS_UNKNOWN_WORKERS + FS_TEXT_WORKERS`
|
||
> exceeds available threads you'll get a warning (non-fatal) and workers will
|
||
> share threads.
|
||
|
||
## Usage
|
||
|
||
```bash
|
||
# health check
|
||
curl http://127.0.0.1:8080/health
|
||
# {"status":"ok"}
|
||
|
||
# upload one or more files (multipart/form-data)
|
||
curl -F "a=@report.pdf" -F "b=@data.csv" http://127.0.0.1:8080/upload
|
||
# 202 {"accepted":[{"id":"<uuid>","name":"report.pdf"}, ...]}
|
||
|
||
# per-stage counters
|
||
curl http://127.0.0.1:8080/stats
|
||
```
|
||
|
||
Each file in a request becomes its own job. Responses:
|
||
|
||
- `202 Accepted` — all files spooled and queued (with per-file job IDs)
|
||
- `400 Bad Request` — not multipart, or no files present
|
||
- `503 Service Unavailable` — queue full, retry later
|
||
- `500 Internal Server Error` — failed to write a file to disk
|
||
|
||
### `GET /stats` — per-stage counters
|
||
|
||
The pipeline counts its own work (`src/stats.jl`), because nothing outside it
|
||
can. Every in-flight file sits in `data/spool/` no matter which stage it has
|
||
reached. The stage is a property of the queue holding its reference, and only
|
||
the pipeline can see that; there is no directory to poll.
|
||
|
||
```jsonc
|
||
{
|
||
"now": 1785725300.5, "since": 1785725291.9, "uptime_seconds": 8.5,
|
||
"intake": { "requests": 400, "files": 400, "bytes": 6553600, "rejected": 0 },
|
||
"stages": [
|
||
{ "stage": 4, "name": "language", "workers": 16,
|
||
"queue_depth": 184, "queue_capacity": 1000,
|
||
"completed": 200, "failed": 0, "bytes": 3276800,
|
||
"busy_seconds": 170.5, // summed handler time across the pool
|
||
"blocked_seconds": 0.0, // of that, time parked on a full downstream queue
|
||
"in_flight": 3 }
|
||
]
|
||
}
|
||
```
|
||
|
||
Counters are monotonic since startup, Prometheus-style. Rates are the reader's
|
||
job, so a scrape is stateless and two readers can't disturb each other. Take two
|
||
scrapes Δt apart and subtract:
|
||
|
||
```
|
||
throughput = Δcompleted / Δt
|
||
utilization = (Δbusy_seconds − Δblocked_seconds) / (Δt × workers)
|
||
```
|
||
|
||
**Utilization is the number that names the bottleneck.** In a pipeline every
|
||
stage completes the same files, so at steady state they all report near-identical
|
||
files/s no matter which one is the constraint; what separates them is how hard
|
||
each pool worked to keep up. The bottleneck sits near 1.0 with its queue backing
|
||
up while its neighbours idle.
|
||
|
||
`blocked_seconds` is what keeps that true. Stages 1 and 3 apply *blocking*
|
||
backpressure (a full downstream queue means the handler parks and retries rather
|
||
than dropping the file), and that wait happens inside the handler. Counting it as
|
||
busy would pin stage 1 at 1.0 whenever stage 2 is the real jam, making every
|
||
stage upstream of a jam look like the jam. Subtracted out, utilization means
|
||
"doing its own work", and a high blocked share becomes its own signal: a stage
|
||
blocked 90% of the time is naming its successor.
|
||
|
||
The counters are a handful of atomic adds per file, recorded in `worker_loop`:
|
||
the one place every stage's work passes through, so a new stage is instrumented
|
||
the moment it is wired up, and never on the read path.
|
||
|
||
## Benchmarking (throughput + memory)
|
||
|
||
There are five harnesses. Only the first needs a running server:
|
||
|
||
| script | measures | server? |
|
||
|---|---|---|
|
||
| `bin/bench.jl` (below) | intake, end-to-end and per-stage throughput; server RSS | **yes** |
|
||
| [`bin/bench_stage1.jl`](#stage-1-component-benchmark-binbench_stage1jl) | stage 1 taken apart: classify vs. rename vs. enqueue vs. logging | no |
|
||
| [`bin/bench_stage2.jl`](#stage-2-component-benchmark-binbench_stage2jl) | stage 2 taken apart: exiftool spawn vs. extraction vs. commit | no |
|
||
| [`bin/bench_model.jl`](#model-microbenchmark-binbench_modeljl) | the classifier alone: inference, feature reads, thread scaling | no |
|
||
| [`bin/cluster_calibrate.jl`](#unknown-format-discovery-stage-5-offline) | stage-5 clustering quality vs. an NCD baseline | no |
|
||
|
||
Running all of them from a clean checkout:
|
||
|
||
```bash
|
||
julia --project=. -e 'using Pkg; Pkg.instantiate()' # once
|
||
|
||
# 1. the model, on its own; no server involved
|
||
julia --project=. -t auto bin/bench_model.jl
|
||
|
||
# 1b. stage 1 taken apart; also no server
|
||
julia --project=. -t auto bin/bench_stage1.jl
|
||
|
||
# 1c. stage 2 taken apart; needs a directory of real files, not generated ones
|
||
julia --project=. -t auto bin/bench_stage2.jl
|
||
|
||
# 2. the pipeline. Start the server in one terminal…
|
||
julia --project=. -t auto bin/server.jl
|
||
|
||
# …and drive it from another. Restart the server between memory runs: Julia's
|
||
# GC returns memory to the OS lazily, so a second run starts inflated.
|
||
julia --project=. -t auto bin/bench.jl --files 2000 --size 8k --concurrency 32
|
||
julia --project=. -t auto bin/bench.jl --files 2 --size 1g --concurrency 1 # memory
|
||
julia --project=. -t auto bin/bench.jl --corpus ../training_set --concurrency 16
|
||
|
||
# 3. stage-5 clustering quality (offline, needs a labelled corpus)
|
||
julia --project=. bin/cluster_calibrate.jl ../training_set
|
||
```
|
||
|
||
The test suite is `julia --project=. -t auto -e 'using Pkg; Pkg.test()'`.
|
||
|
||
`bin/bench.jl` must run on the same machine as the server (it reads the sink dirs
|
||
and `/proc`); otherwise it only speaks HTTP: `/upload`, `/health`, and `/stats`
|
||
for the per-stage numbers.
|
||
|
||
Three properties of this design dictate how it measures:
|
||
|
||
- **HTTP latency is not throughput.** `/upload` returns `202` once the bytes are
|
||
spooled and a reference is enqueued, and all four stages run *after* the
|
||
response, so `ab`/`hey`/`wrk` would only ever measure intake. The bench
|
||
instead uploads a corpus and polls the terminal sinks (`done/`, `text_done/`, `binary/`,
|
||
`failed/`) until the count stops moving, and reports both numbers separately:
|
||
intake rate *and* end-to-end completion rate. It also samples `spool/`, whose
|
||
peak depth is the high-water mark of files in flight. But *which* stage they
|
||
are waiting on comes from `/stats`, not from disk.
|
||
- **End-to-end throughput doesn't name the slow stage.** The four stages run
|
||
concurrently behind their own queues, so the pipeline's rate *is* the slowest
|
||
stage's rate and the others are invisible in it. The bench scrapes
|
||
[`/stats`](#get-stats--per-stage-counters) before and after the run and
|
||
subtracts, giving each stage its own throughput, mean service time and
|
||
utilization:
|
||
|
||
```
|
||
PER-STAGE (server counters, delta over the end-to-end window)
|
||
stage files/s MiB/s svc ms util blocked peak queue failed
|
||
1 classify 32.52 0.51 38.7 0.08 0.0% 209/1000 0
|
||
2 enrich 0.24 0.0 541.0 0.01 0.0% 0/1000 0
|
||
3 triage 32.27 0.5 33.9 0.07 0.0% 83/1000 0
|
||
4 language 16.26 0.25 852.6 0.87 0.0% 184/1000 0
|
||
bottleneck stage 4 (language) at 87.0% utilization of 16 worker(s)
|
||
```
|
||
|
||
(400 mixed files, 16 KiB each, concurrency 16. Stage 4 is the constraint,
|
||
`github-linguist` being a process spawn per file at ~850 ms, and the 200 text
|
||
files queue up behind it while stages 1 and 3 idle at under 10%.) Read *down*
|
||
the util column, not the files/s column: each stage sees a different subset of
|
||
the corpus, so a low rate can just mean little work was routed there.
|
||
- **Memory should be flat in file size, and the sweep is what proves it.** Both
|
||
halves of the pipeline are bounded: the workers read bounded prefixes (16+16
|
||
bytes to classify, 8 KB to sniff, 64 KB to language-detect), and intake streams
|
||
each upload to disk a chunk at a time. So peak RSS should track *concurrency*,
|
||
not size. Measured on this machine (16 threads, 64 KiB chunk):
|
||
|
||
| upload size | concurrency | RSS growth |
|
||
|---|---|---|
|
||
| 256 MiB × 4 | 1 | 14.0 / 20.4 MiB (two runs) |
|
||
| 1 GiB × 2 | 1 | 22.3 MiB |
|
||
| 2 GiB × 1 | 1 | 31.3 MiB |
|
||
| 256 MiB × 4 | 4 | none measurable (peak stayed under the baseline) |
|
||
|
||
**Read the shape, not the digits.** Growth is flat across an 8× range of file
|
||
sizes: 14–31 MiB whether the upload is 256 MiB or 2 GiB. That is the claim
|
||
that matters, because nothing scales with file size, so nothing is buffering.
|
||
The absolute figures are *not* precise to the megabyte. A freshly started server
|
||
settles anywhere in an ~860–985 MiB band, so run-to-run variance in the
|
||
baseline is comparable to the growth being measured; the residual is GC churn
|
||
from the chunk reads, not retained buffers.
|
||
|
||
Two consequences worth knowing before quoting these numbers:
|
||
|
||
- **Let the server settle ~30s after startup** before a memory run, or the
|
||
baseline is sampled mid-fall and the run reports less growth than it caused
|
||
(or none at all, as the concurrency-4 row above did).
|
||
- **Linear-in-concurrency is not currently demonstrated.** An earlier
|
||
measurement put 4 concurrent 256 MiB uploads at 86.9 MiB (21.7 per in-flight
|
||
upload), but that run predates a fix to how the baseline was sampled: it was
|
||
read *before* the peak counter was reset, against a different origin. The
|
||
re-measurement above cannot reproduce it, and the growth sits under the noise
|
||
floor. Expect concurrency to cost memory; don't trust a specific coefficient
|
||
without a quieter machine or many more runs.
|
||
|
||
Before intake was streamed, the same 256 MiB upload grew RSS by ~700 MiB and 4
|
||
concurrent ones pushed a 950 MiB baseline past 2 GiB. So a `--size` sweep that
|
||
slopes upward is the regression signal: something has started buffering bodies
|
||
again.
|
||
|
||
Flags: `--files`, `--size` (`8k`/`64m`/`1g`), `--concurrency`, `--kind`
|
||
(`binary`/`text`/`mixed`, which chooses the stages that get loaded), `--corpus
|
||
DIR` (real files, the only way to exercise stage 2's exiftool path), `--pid`, `--no-mem`,
|
||
`--no-stats`, `--sample-ms`, `--timeout`, `--json PATH`. Full list in the script
|
||
header. The `--json` output carries the per-stage numbers too, so a sweep can be
|
||
compared run to run.
|
||
|
||
Two caveats the script reports rather than hides: it counts sink *deltas*, so it
|
||
warns if the pipeline isn't idle at the start (in-flight leftovers would be
|
||
counted as its own throughput); and because Julia's GC returns memory to the OS
|
||
lazily, a second run in the same process starts from an inflated baseline. It
|
||
resets the kernel's peak-RSS counter (`/proc/<pid>/clear_refs`) per run and flags
|
||
a drifted baseline, but for a clean growth figure restart the server between
|
||
memory runs.
|
||
|
||
### Stage-1 component benchmark (`bin/bench_stage1.jl`)
|
||
|
||
`bin/bench.jl` reports stage 1 as one number and `bin/bench_model.jl` takes the
|
||
*classifier* apart, but stage 1 is more than the model. Per file it also
|
||
renames the file into its stage directory, pushes a reference onto the
|
||
downstream queue, and logs. `bin/bench_stage1.jl` times each of those in
|
||
isolation, then times the real `handle_classify_job` end to end so the parts can
|
||
be checked against the whole:
|
||
|
||
```bash
|
||
julia --project=. -t auto bin/bench_stage1.jl
|
||
```
|
||
|
||
Measured on this machine (Ryzen 7 2700X, 8 cores/16 threads, Julia 1.12; 2000 ×
|
||
64 KiB files, minimum of 5 trials):
|
||
|
||
| component | per file | share of the handler |
|
||
|---|---|---|
|
||
| `classify()` | 10.7 µs | 38% |
|
||
| ↳ `read_features` | 7.9 µs | 28% |
|
||
| ↳ `Lux.apply` | 2.3 µs | 8% |
|
||
| `move_to` (rename) | 11.7 µs | 41% |
|
||
| `enqueue_blocking!` | 0.12 µs | 0.4% |
|
||
| per-file logging (disabled `@debug`) | 0.29 µs | 1% |
|
||
| **`handle_classify_job`** | **28.3 µs** | 100% |
|
||
|
||
**This benchmark is why stage 1's per-file log lines are `@debug` rather than
|
||
`@info`.** As `@info` they cost ~71 µs of the handler's ~118 µs, about 6× the
|
||
classifier and 6× the rename, and nearly all of it was `ConsoleLogger`
|
||
*formatting* (~64 µs), not the `FlushLogger`'s per-message flush (~8 µs on top).
|
||
Demoting them took stage 1 from 8.5k files/s to 35.3k files/s on a single worker,
|
||
a 4.2× speedup for no algorithmic change. The script still prices a formatted
|
||
line, so the cost of turning them back on is visible: running the handler under
|
||
`JULIA_DEBUG=DarkStruct` measures 133 µs per file, a 4.7× slowdown. That is the
|
||
trade. Per-file tracing is available when you want it and off by default, with
|
||
`GET /stats` giving per-file observability that is counted, not formatted.
|
||
|
||
What's left is evenly split between the rename and the classifier, and neither
|
||
has an easy 2×. Two things worth knowing:
|
||
|
||
- **The rename, not the model, is the single largest component** (11.7 µs), and
|
||
it's a plain `mv` within one filesystem. Inside `classify`, the same pattern
|
||
holds: 7.9 µs of the 10.7 µs is `read_features` (the `open`, the two reads
|
||
and the `seek`) against 2.3 µs of actual inference. Stage 1 is now a
|
||
filesystem-bound stage with a neural network attached, not the reverse.
|
||
- **Stage 1 now peaks at ~4 workers.** With the logger removed from the hot path
|
||
the sweep reads 35.0k/s at 1 worker, 66.6k/s at 2, **73.3k/s at 4**, then
|
||
*falls back* to 60.1k/s at 8 and 53.3k/s at 16, because every worker renaming
|
||
into the same two directories contends on the same directory inode. That ceiling
|
||
coincides with the one `bin/bench_model.jl` finds for inference, so ~4 is the
|
||
number from both directions: raising `FS_WORKERS` past it costs throughput.
|
||
|
||
Reported times are the minimum over trials. Flags: `--files`, `--reps`,
|
||
`--trials`, `--size`, `--dir`, `--model`, `--threads`, `--no-threads`,
|
||
`--json PATH`.
|
||
|
||
### Stage-2 component benchmark (`bin/bench_stage2.jl`)
|
||
|
||
Stage 2 is the one stage whose cost is dominated by something outside Julia
|
||
entirely: it forks `exiftool`, a Perl program, once per file. `bin/bench.jl`
|
||
reports the stage as a single throughput number, which can't distinguish "the
|
||
extraction is slow" from "the *spawn* is slow", and those have opposite fixes.
|
||
`bin/bench_stage2.jl` times each piece in isolation, then times the real
|
||
`handle_known_job` end to end:
|
||
|
||
```bash
|
||
julia --project=. -t auto bin/bench_stage2.jl
|
||
```
|
||
|
||
Two things make this benchmark different from the stage-1 one:
|
||
|
||
- **The corpus must be real files.** exiftool's cost depends on what it finds; a
|
||
file of random bytes bails out early and understates the stage by ~10×. The
|
||
default corpus is `data/done`, files that already went through stage 2 on this
|
||
machine. `--corpus PATH` points it elsewhere.
|
||
- **It prices the alternatives to one-fork-per-file**, because if the fork
|
||
dominates then the only fixes are to stop paying it per file. `exiftool
|
||
(batched Nx)` runs the whole corpus through one process; `exiftool
|
||
(-stay_open)` keeps one process alive and feeds it one file at a time over a
|
||
pipe, which is the shape a streaming pipeline could actually adopt. Both are
|
||
measured, not assumed.
|
||
|
||
Measured on this machine (Ryzen 7 2700X, 8 cores/16 threads, Julia 1.12,
|
||
exiftool 12.40; 150 real files / 102 MiB, minimum of 2 trials):
|
||
|
||
| component | per file | share of the handler |
|
||
|---|---|---|
|
||
| `run_exiftool()` | 135.9 ms | 98% |
|
||
| ↳ bare fork + Perl boot (`exiftool -ver`) | 76.7 ms | 55% |
|
||
| ↳ `JSON3.read` + tag map | 10 µs | 0.0% |
|
||
| `normalize_metadata` | 1.7 µs | 0.0% |
|
||
| `commit_enriched!` (sidecar + fsyncs + rename) | 2.0 ms | 1.5% |
|
||
| per-file logging (`@info`, flush→file) | 47 µs | 0.0% |
|
||
| **`handle_known_job`** | **138.3 ms** | 100% |
|
||
| *alt:* `exiftool -stay_open` | 40.1 ms | 29% |
|
||
| *alt:* `exiftool` batched 150× | 37.3 ms | 27% |
|
||
|
||
**Stage 2 is exiftool and nothing else.** Everything the Julia code does
|
||
(parsing, normalizing, the durable sidecar-first commit, the log line) sums to
|
||
about 1.5% of the stage. There is no point optimizing any of it.
|
||
|
||
**More than half the stage is interpreter startup, not metadata extraction.**
|
||
The bare `exiftool -ver` (fork, Perl boot, module loads, read no file) costs
|
||
76.7 ms against a 135.9 ms full call. Both fork-free alternatives agree on what's
|
||
left: ~37–40 ms of actual work per file. So a persistent exiftool would cut the
|
||
stage by ~70%, and `-stay_open` gets there without giving up the one-file-in,
|
||
one-result-out shape the pipeline needs. That remains the single biggest
|
||
available win in this stage; it is measured here but not yet implemented.
|
||
|
||
**This benchmark is also why `run_with_timeout` no longer polls.** The original
|
||
watchdog polled with `sleep(0.1)` and then joined the polling task, so every call
|
||
paid the remainder of an in-flight sleep *after* the child had already exited:
|
||
~25 ms per file here, and a measured 101 ms on a process that exits instantly.
|
||
Replacing it with a one-shot `Timer` took the stage from 164.6 ms to 138.3 ms per
|
||
file (6→8 files/s on one worker) and cost nothing in behavior. Stage 4 shares the
|
||
wrapper and got the same fix for free.
|
||
|
||
Writing the missing tests for that wrapper turned up a second, worse problem:
|
||
**the timeout was never enforceable.** `wait(proc)` returns only once the
|
||
captured stdout pipe closes, and any grandchild inherits that pipe, so
|
||
signalling the child alone left the worker blocked until the whole process tree
|
||
finished on its own (a `sh -c "trap '' TERM; sleep 30"` child ran the full 30 s
|
||
against a 1 s timeout, under both the old and new watchdog). The child now runs
|
||
in its own process group and the timeout signals the group. The trade is that a
|
||
hard crash of the server orphans an in-flight child rather than taking it down;
|
||
these children are short-lived and timeout-bounded, which is the cheaper side of
|
||
it.
|
||
|
||
**Stage 2 scales to ~8 workers, then flattens**: 8 files/s at 1 worker, 14 at 2,
|
||
28 at 4, **49 at 8**, and 49 at 16. The machine runs out of cores to run Perl
|
||
on, which is exactly what you'd expect of a stage that is ~100% subprocess. Note
|
||
that the sweep pulls from a shared counter rather than splitting the corpus into
|
||
contiguous slices: per-file exiftool time spans two orders of magnitude on a real
|
||
corpus (one 2.1 s archive among 48 files), and a static split reports a scaling
|
||
ceiling that is really just load imbalance.
|
||
|
||
One caveat the numbers raise but don't answer: **`fsync_dir` measures 1.75 µs**,
|
||
which is far too fast to be a real disk flush. The durability that
|
||
`commit_enriched!` is written for may not survive power loss on this filesystem,
|
||
even though the code is correct. That's a correctness question, not a speed one,
|
||
and it is not yet resolved.
|
||
|
||
Reported times are the minimum over trials. Flags: `--files`, `--reps`,
|
||
`--trials`, `--corpus`, `--dir`, `--timeout`, `--threads`, `--no-threads`,
|
||
`--no-stay-open`, `--json PATH`.
|
||
|
||
### Model microbenchmark (`bin/bench_model.jl`)
|
||
|
||
`bin/bench.jl` reports stage 1 as a single number: the wall time of
|
||
`handle_classify_job`, which is a feature read, an inference, a rename, a
|
||
(disabled) debug line, and whatever contention the other three pools create.
|
||
That's the right number for capacity planning and the wrong one for "is the
|
||
model slow?".
|
||
`bin/bench_model.jl` answers that separately, with no server, queue, or HTTP
|
||
involved:
|
||
|
||
```bash
|
||
julia --project=. -t auto bin/bench_model.jl
|
||
```
|
||
|
||
Measured on this machine (Ryzen 7 2700X, 8 cores/16 threads, Julia 1.12):
|
||
|
||
| what | per file | notes |
|
||
|---|---|---|
|
||
| `Lux.apply`, batch 1 | **2.3 µs** | 768 B allocated per call |
|
||
| `read_features` | **4.2–6.0 µs** | flat across a 262,144× size range (1 KiB → 256 MiB) |
|
||
| `classify()` | **8.4 µs** | 73% feature read, 27% inference |
|
||
|
||
So the model is not the pipeline's problem, by three orders of magnitude: the
|
||
same run measured stage 1 at 38.7 ms per file, ~4,500× the 8.4 µs `classify()`
|
||
costs. Whatever stage 1 spends its time on, it isn't the network. (That 38.7 ms
|
||
predates the `@debug` demotion above and is a whole-pipeline figure. It includes
|
||
time the stage-1 worker spends *blocked* on a full downstream queue, which is why
|
||
it is three orders of magnitude above the 28.3 µs the handler costs in
|
||
isolation. For the uncontended split, see
|
||
[the stage-1 decomposition](#stage-1-component-benchmark-binbench_stage1jl).)
|
||
|
||
Two findings worth acting on if stage 1 ever *does* become the constraint:
|
||
|
||
- **Batching would buy ~13×.** A 32×1 matmul wastes most of a BLAS call:
|
||
batch 64 costs 280 ns/file and batch 512 costs 179 ns/file, against 2.34 µs
|
||
one at a time. The pipeline classifies strictly one file per job today, so it
|
||
pays the worst row in that table.
|
||
- **Inference does not scale past ~4 threads.** Concurrent `Lux.apply` on the
|
||
shared read-only `Classifier` peaks around 1.2M files/s at 4 tasks and then
|
||
*falls back* to single-thread throughput at 16. The script runs a pure-compute
|
||
control kernel through the same sweep to place the blame: the control reaches
|
||
14.3× at 16 tasks (90% efficiency) on the same box, so the machine
|
||
parallelizes and `Lux.apply` doesn't. GC is only ~1% of it, so allocation
|
||
pressure isn't the explanation either. The cause is inside Lux/BLAS and is
|
||
not diagnosed here. The practical consequence: raising `FS_WORKERS` past ~4
|
||
adds no classification throughput.
|
||
|
||
Reported times are the minimum over trials, because for a microbenchmark the
|
||
floor is the signal and everything above it is scheduler and GC noise. Every
|
||
timed loop stores its result in a sink so a pure call can't be hoisted out.
|
||
Flags: `--model`, `--reps`, `--trials`, `--batches`, `--sizes`, `--no-threads`,
|
||
`--json PATH`.
|
||
|
||
## Layout
|
||
|
||
```
|
||
src/
|
||
DarkStruct.jl module + run() (startup, recovery, workers, serve, shutdown)
|
||
config.jl Config struct + env parsing
|
||
job.jl Job (the queue reference)
|
||
queue.jl JobQueue seam + in-process ChannelQueue
|
||
rabbit.jl RabbitMQ-backed JobQueue: durable queues, acks, confirms
|
||
stats.jl per-stage counters behind GET /stats (throughput, utilization)
|
||
multipart.jl streaming multipart/form-data reader (intake never buffers a file)
|
||
spool.jl filename sanitizing, streaming spool/move, startup recovery
|
||
model.jl NN architecture + byte→feature mapping (shared with trainer)
|
||
classify.jl load artifact + classify a file at inference time
|
||
metadata.jl exiftool extraction + normalized sidecar (stage 2)
|
||
content.jl binary-vs-text sniff for unknown files (stage 3)
|
||
language.jl natural + programming language enrichment for text (stage 4)
|
||
cluster.jl header-byte clustering model + Gibbs + scoring core (stage 5, science)
|
||
catalog.jl durable single-owner format catalog + sweep + nominations (stage 5, phase B)
|
||
worker.jl parametrized worker loop + classify/enrich/triage/language handlers
|
||
server.jl HTTP routes + the streaming /upload handler
|
||
docker-compose.yml the server, in-process queues
|
||
docker-compose.rabbitmq.yml overlay: adds the broker and switches the backend
|
||
bin/
|
||
server.jl entry point
|
||
bench.jl throughput + memory harness against a running server
|
||
bench_model.jl classifier microbenchmark (inference, feature reads, scaling)
|
||
bench_stage1.jl stage-1 decomposition (classify vs. rename vs. enqueue vs. logging)
|
||
train.jl offline training script → model/classifier.jld2
|
||
cluster_calibrate.jl offline stage-5 hyperparameter calibration + NCD baseline
|
||
cluster_sweep.jl stage-5 phase-B runner: sweep binary/, update catalog, write nominations
|
||
model/
|
||
classifier.jld2 committed trained weights (loaded at startup)
|
||
DESIGN_clustering.md stage-5 design rationale + calibration results
|
||
```
|