Compare commits

...

3 Commits

Author SHA1 Message Date
e18ac45d70 Updated README 2026-07-25 09:59:04 -04:00
f4e3f5be0b Update README for stage-5 phase B (durable catalog + sweep runner) 2026-07-03 17:04:37 -04:00
584bad02a7 Add stage-5 phase B: durable single-owner format catalog
Wraps the assign_file scoring core (cluster.jl) in the live catalog the
design's phase B calls for (DESIGN §5B/§9):

- src/catalog.jl: Catalog durable state (frozen-id clusters + sufficient
  stats + processed set + examples); sparse, sidecar-first durable
  save/load; incremental catalog_sweep! (deterministic CRP-predictive
  assignment of new binary/ files); offline compact! that seeds on first
  run and recompacts later; write_nominations! emitting one JSON per
  promotable cluster with a hex magic template.
- bin/cluster_sweep.jl: cron/periodic single-owner runner (--compact
  forces a recluster; first run auto-compacts to seed).
- config.jl: cluster_catalog_path + nominated_dir knobs (FS_CLUSTER_CATALOG,
  FS_NOMINATED_DIR), wired into config_from_env and ensure_dirs.
- Tests: durable round-trip, incremental sweep growth + idempotency,
  §10.1 nothing-from-noise end-to-end (zero promotions), a recurring
  format self-nominating, seed-then-live-assign (165 pass).
- DESIGN_clustering.md: mark phase B built.
2026-07-03 17:03:33 -04:00
7 changed files with 666 additions and 16 deletions

View File

@@ -223,9 +223,30 @@ clusters. A cluster's spiked positions become a libmagic-style signature;
clusters with enough members and enough fixed positions self-**nominate** for clusters with enough members and enough fixed positions self-**nominate** for
promotion (a human does the one irreversible step, redefining "known"). promotion (a human does the one irreversible step, redefining "known").
**Status:** the offline science (phase A) is implemented and calibrated; the live **Status:** both phases are implemented and calibrated. Phase A (offline Gibbs)
catalog process (phase B) is designed and its scoring core (`assign_file`) is in is the science; phase B (`src/catalog.jl`) is the live catalog: a durable
place, but its batch-runner plumbing is not yet built. 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 — never in the request Calibration is its own offline script (like training — never in the request
path), scored against magic-collapsed ground truth (so `docx``zip` and the whole path), scored against magic-collapsed ground truth (so `docx``zip` and the whole
@@ -299,9 +320,12 @@ tell you "PDF", just "this looks like something I was trained on, or not".
- **Artifact:** trained weights live in `model/classifier.jld2` (committed), so - **Artifact:** trained weights live in `model/classifier.jld2` (committed), so
the server just loads them at startup. Missing/unreadable ⇒ the server fails the server just loads them at startup. Missing/unreadable ⇒ the server fails
fast rather than run without classification. fast rather than run without classification.
- **Effect today:** *annotate-only*. The class is logged - **Effect today:** *active routing*. The class is logged
(`classification=known|unknown`) but every file still moves to `done/`; the (`classification=known|unknown`) and drives the pipeline split: `:known` files
classifier can't misroute real files while it's unproven. go to `known/` for metadata enrichment (stage 2), `:unknown` files go to
`unknown/` for content triage (stage 3). The class chooses the downstream
stage; 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 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. shared by the trainer and the server, so they can't drift apart.
@@ -359,6 +383,8 @@ init, so the artifact is exactly regenerable from the same inputs.
| `FS_CLUSTER_BG_MASS` | `5.0` | Mass of the uniform background component | | `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_MEMBERS` | `20` | Cluster size threshold for promotion nomination |
| `FS_PROMOTE_MIN_MAGIC` | `3` | Required fixed signature positions to nominate | | `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) |
> To get real parallelism, start Julia with enough threads (`-t N`) to cover all > 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` > pools. If `FS_WORKERS + FS_KNOWN_WORKERS + FS_UNKNOWN_WORKERS + FS_TEXT_WORKERS`
@@ -398,13 +424,15 @@ src/
metadata.jl exiftool extraction + normalized sidecar (stage 2) metadata.jl exiftool extraction + normalized sidecar (stage 2)
content.jl binary-vs-text sniff for unknown files (stage 3) content.jl binary-vs-text sniff for unknown files (stage 3)
language.jl natural + programming language enrichment for text (stage 4) language.jl natural + programming language enrichment for text (stage 4)
cluster.jl header-byte clustering for unknown-format discovery (stage 5, offline) 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 worker.jl parametrized worker loop + classify/enrich/triage/language handlers
server.jl HTTP routes/handlers server.jl HTTP routes/handlers
bin/ bin/
server.jl entry point server.jl entry point
train.jl offline training script → model/classifier.jld2 train.jl offline training script → model/classifier.jld2
cluster_calibrate.jl offline stage-5 hyperparameter calibration + NCD baseline 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/ model/
classifier.jld2 committed trained weights (loaded at startup) classifier.jld2 committed trained weights (loaded at startup)
DESIGN_clustering.md stage-5 design rationale + calibration results DESIGN_clustering.md stage-5 design rationale + calibration results

19
bin/cluster_sweep.jl Normal file
View File

@@ -0,0 +1,19 @@
#!/usr/bin/env julia
#
# Stage-5 phase-B runner (model/DESIGN_clustering.md §9): the single-owner,
# periodic/cron process that sweeps `binary/`, folds new files into the durable
# format catalog by sequential CRP-predictive assignment, and (re)writes
# promotion nominations. Run it single-threaded on a schedule — it is the ONLY
# writer of the catalog, so no locking is needed.
#
# julia --project=. bin/cluster_sweep.jl # incremental live sweep
# julia --project=. bin/cluster_sweep.jl --compact # offline Gibbs (seed / recompact)
#
# On a fresh catalog (nothing processed yet) the incremental sweep would send
# every file to background — there are no clusters to match — so the first run
# auto-promotes to a compaction pass to seed the catalog. Configure via the
# FS_CLUSTER_* / FS_NOMINATED_DIR env vars (see src/config.jl).
using FileServer
FileServer.cluster_sweep_cli(ARGS)

View File

@@ -1,8 +1,10 @@
# Stage-5: Unknown-format discovery by Bayesian header clustering # Stage-5: Unknown-format discovery by Bayesian header clustering
Status: **phase A implemented and calibrated** (`src/cluster.jl`, Status: **phases A and B implemented and calibrated** (`src/cluster.jl` +
`bin/cluster_calibrate.jl`, tests in `test/runtests.jl`); phase-B core scoring `src/catalog.jl`, `bin/cluster_calibrate.jl` + `bin/cluster_sweep.jl`, tests in
implemented (`assign_file`), its live batch-process plumbing still to do. Product `test/runtests.jl`). Phase A (offline Gibbs) is calibrated; phase B's durable
single-owner catalog, incremental sweep, and nomination writer are now built on
top of the `assign_file` scoring core. Product
of a design interview; captures the decisions and — as important — the of a design interview; captures the decisions and — as important — the
assumptions we *rejected* so they don't get silently reintroduced. §11 records assumptions we *rejected* so they don't get silently reintroduced. §11 records
what building it actually taught us, including three assumptions in this document what building it actually taught us, including three assumptions in this document
@@ -291,7 +293,14 @@ model-free gzip similarity.
a MAP-style stand-in for the VI/Binder posterior summary §5 defers adequate a MAP-style stand-in for the VI/Binder posterior summary §5 defers adequate
because the formats are strongly separated; revisit if compaction 5) needs it. because the formats are strongly separated; revisit if compaction 5) needs it.
- Phase B's **live single-owner batch process** 9) and the durable catalog file - Phase B's **live single-owner batch process** 9) and the durable catalog file
are not yet built; `assign_file` is the scoring core they will wrap. are implemented in `src/catalog.jl` (the `Catalog` durable state, the
incremental `catalog_sweep!`, the offline `compact!` seed/recompaction, and
`write_nominations!`), driven by `bin/cluster_sweep.jl` (cron/periodic; the
first run auto-compacts to seed, subsequent runs sweep incrementally). The
catalog is persisted with the stage-2 sidecar-first tempfsyncrenamefsync-dir
discipline. As §5B predicts, under the calibrated `bg_mass > α` the live sweep
never mints single-file clusters new formats are discovered by the offline
`compact!` re-clustering the background residue, not by the live path.
## Open items (deferred, intentionally) ## Open items (deferred, intentionally)

View File

@@ -19,7 +19,8 @@ include("classify.jl") # Classifier + load_classifier/classify (needs model.jl)
include("metadata.jl") # exiftool extraction + sidecar enrichment (stage 2) include("metadata.jl") # exiftool extraction + sidecar enrichment (stage 2)
include("content.jl") # binary-vs-text triage for unknown files (stage 3) include("content.jl") # binary-vs-text triage for unknown files (stage 3)
include("language.jl") # natural + programming language enrichment for text (stage 4) include("language.jl") # natural + programming language enrichment for text (stage 4)
include("cluster.jl") # unknown-format discovery by header clustering (stage 5) include("cluster.jl") # unknown-format discovery by header clustering (stage 5, science)
include("catalog.jl") # durable single-owner format catalog (stage 5, phase B; needs cluster.jl + metadata.jl fsync)
include("worker.jl") include("worker.jl")
# Globals the HTTP handlers read at request time. Set once in `run`, before the # Globals the HTTP handlers read at request time. Set once in `run`, before the

403
src/catalog.jl Normal file
View File

@@ -0,0 +1,403 @@
# Stage-5 phase B: the live, single-owner format catalog (DESIGN §5B/§9).
#
# `cluster.jl` is the *science* — feature extraction, the collapsed Gibbs sampler
# (phase A, offline), and `assign_file` (the phase-B scoring core). This file is
# the *plumbing* that turns that core into a durable, growing catalog:
#
# * a `Catalog` = the surviving clusters' sufficient statistics + a record of
# which files have already been folded in + a few example filenames each;
# * `catalog_sweep!` = the phase-B loop — for every *new* file in `binary/`,
# run the deterministic CRP-predictive `assign_file` and fold it into the
# chosen cluster's stats (DESIGN §5B);
# * `compact!` = the offline Gibbs pass that *seeds* the catalog on first run
# and periodically re-clusters the pile (DESIGN §5, "seed the initial
# catalog" / "periodic compaction");
# * `write_nominations!` = surface promotable clusters to a human (DESIGN §6).
#
# Concurrency model is the deliberate opposite of the stateless classify workers
# (DESIGN §9): exactly ONE process owns the catalog, so there are no locks and no
# torn reads of sufficient stats. The catalog is a single durable file mutated by
# that one process; it is committed with the same sidecar-first temp→fsync→rename
# →fsync-dir discipline as the stage-2 sidecars (`commit_enriched!`), so a crash
# mid-write can neither corrupt it nor lose the rename. This file lives in the
# module (not dependency-flat like cluster.jl) because it needs JSON3 + the fsync
# helpers from metadata.jl.
"How many example filenames to retain per cluster (for the human promotion gate,
DESIGN §6). A handful is plenty to eyeball; the count/signature carry the weight."
const CATALOG_EXAMPLE_CAP = 8
"""
Catalog
The mutable phase-B state owned by the single sweep process:
* `n` — header window these clusters were built at (must match the
feature window used to score new files; frozen once seeded).
* `clusters` — id → `ClusterStats` (per-position 257-counts + member count).
Ids are **frozen at birth** — never renumbered — so there is no
label switching across sweeps (DESIGN §5).
* `examples` — id → up to `CATALOG_EXAMPLE_CAP` member filenames, for the
human nomination glance.
* `next_id` — the next fresh cluster id to hand out (monotone; retired ids are
never reused, keeping ids globally unique over the catalog's life).
* `processed` — basenames of every `binary/` file already folded in, so a sweep
is incremental: it touches only files it has not seen. This set
grows with the `binary/` pile it mirrors — the same population,
no faster — which is acceptable for v1 (DESIGN §9).
"""
mutable struct Catalog
n::Int
clusters::Dict{Int,ClusterStats}
examples::Dict{Int,Vector{String}}
next_id::Int
processed::Set{String}
end
"An empty catalog at header window `n` (no clusters seen yet)."
Catalog(n::Integer) = Catalog(Int(n), Dict{Int,ClusterStats}(),
Dict{Int,Vector{String}}(), 1, Set{String}())
"Record `name` as an example of cluster `id`, capped at `CATALOG_EXAMPLE_CAP`."
function record_example!(cat::Catalog, id::Integer, name::AbstractString)
ex = get!(cat.examples, id, String[])
length(ex) < CATALOG_EXAMPLE_CAP && push!(ex, String(name))
return nothing
end
# ---------------------------------------------------------------------------
# Durable persistence (reuse the stage-2 sidecar-first commit discipline)
# ---------------------------------------------------------------------------
# Count tables are stored sparsely — only the non-zero `(position, symbol, count)`
# triples — because a cluster's n×257 table is overwhelmingly zero (a peaked
# magic byte touches one of 257 symbols per position). Sparse keeps the catalog
# file small and its load O(non-zeros), not O(n·257·K).
function _sparse_counts(counts::Matrix{Int})
n, A = size(counts)
triples = Vector{Vector{Int}}()
@inbounds for i in 1:n, v in 1:A
c = counts[i, v]
c != 0 && push!(triples, [i, v, c])
end
return triples
end
function _dense_counts(triples, n::Integer)
counts = zeros(Int, n, ALPHABET)
for t in triples
counts[Int(t[1]), Int(t[2])] = Int(t[3])
end
return counts
end
"Serialize a `Catalog` to a plain `NamedTuple` ready for `JSON3.write`."
function catalog_payload(cat::Catalog)
clusters = [(
id = id,
members = c.members,
counts = _sparse_counts(c.counts),
examples = get(cat.examples, id, String[]),
) for (id, c) in sort(collect(cat.clusters); by=first)]
return (
n = cat.n,
next_id = cat.next_id,
clusters = clusters,
processed = sort(collect(cat.processed)),
)
end
"""
save_catalog!(path, cat)
Durably write `cat` to `path` with the sidecar-first ordering (DESIGN §9): write
to a temp name, fsync the bytes, atomically rename into place, then fsync the
containing directory so the rename itself survives power loss. A crash can leave
at most a stale `.tmp`, never a torn catalog.
"""
function save_catalog!(path::AbstractString, cat::Catalog)
dir = dirname(path)
isempty(dir) || mkpath(dir)
tmp = string(path, ".tmp")
open(tmp, "w") do io
write(io, JSON3.write(catalog_payload(cat)))
flush(io)
fsync_fd(fd(io)) # persist bytes before the rename
end
mv(tmp, path; force=true) # atomic replace
fsync_dir(isempty(dir) ? "." : dir) # persist the rename itself
return path
end
"""
load_catalog(path; n) -> Catalog
Load the durable catalog from `path`, or return a fresh empty `Catalog(n)` if it
does not exist yet (first run). `n` is the configured header window used only for
the empty case; a loaded catalog keeps its own frozen `n`.
"""
function load_catalog(path::AbstractString; n::Integer=HEADER_N)
isfile(path) || return Catalog(n)
doc = JSON3.read(read(path, String))
cn = Int(doc.n)
cat = Catalog(cn)
cat.next_id = Int(doc.next_id)
for entry in doc.clusters
id = Int(entry.id)
c = ClusterStats(_dense_counts(entry.counts, cn), Int(entry.members))
cat.clusters[id] = c
cat.examples[id] = String[String(e) for e in entry.examples]
end
for name in doc.processed
push!(cat.processed, String(name))
end
return cat
end
# ---------------------------------------------------------------------------
# Listing the input pile
# ---------------------------------------------------------------------------
"""
binary_files(dir) -> Vector{String}
Sorted absolute paths of the regular files in `dir` to be swept, skipping
`.meta.json` sidecars and any `.tmp` scratch. Sorted so a sweep's sequential
CRP-predictive assignment (which is order-dependent) is deterministic run to run.
"""
function binary_files(dir::AbstractString)
isdir(dir) || return String[]
paths = String[]
for name in readdir(dir; join=true)
isfile(name) || continue
(endswith(name, ".meta.json") || endswith(name, ".tmp")) && continue
push!(paths, name)
end
sort!(paths)
return paths
end
# ---------------------------------------------------------------------------
# Phase B: the incremental sweep (the deliverable)
# ---------------------------------------------------------------------------
"""
catalog_sweep!(cat, cfg) -> NamedTuple
Fold every *new* file in `cfg.cluster_dir` into `cat` using the deterministic
CRP-predictive rule (`assign_file`, DESIGN §5B), updating the chosen cluster's
sufficient statistics in place. Files already in `cat.processed` are skipped, so
repeated sweeps are incremental and idempotent over the pile.
Per file, `assign_file` returns the argmax component:
* an existing cluster id → the file joins it (`add!`);
* `0` (background) → counted, not clustered — a novel-but-unmatched file parks
here by design; genuinely new formats are discovered by the offline
`compact!` re-clustering the residue, not by single-file minting;
* `-1` (mint) → a fresh cluster is seeded with a frozen `next_id`. Under the
calibrated `bg_mass > α` this never fires on the live path (fresh and
background share the one-file likelihood, so background always wins) — the
branch exists for correctness, not as a routine outcome (DESIGN §5B).
Mutates `cat` but does NOT persist it — the caller commits once, after the sweep.
Returns a summary of what happened this sweep.
"""
function catalog_sweep!(cat::Catalog, cfg::Config)
n_seen = 0; n_bg = 0; n_joined = 0; n_minted = 0
for path in binary_files(cfg.cluster_dir)
base = basename(path)
base in cat.processed && continue
x = header_symbols(path; n=cat.n)
ids = sort!(collect(keys(cat.clusters)))
k = assign_file(x, cat.clusters, ids;
α=cfg.cluster_alpha, β=cfg.cluster_pseudocount,
bg_mass=cfg.cluster_bg_mass)
if k == -1
id = cat.next_id
cat.next_id += 1
c = ClusterStats(cat.n)
add!(c, x)
cat.clusters[id] = c
record_example!(cat, id, base)
n_minted += 1
elseif k == 0
n_bg += 1
else
add!(cat.clusters[k], x)
record_example!(cat, k, base)
n_joined += 1
end
push!(cat.processed, base)
n_seen += 1
end
return (; n_seen, n_joined, n_bg, n_minted,
n_clusters=length(cat.clusters), n_processed=length(cat.processed))
end
# ---------------------------------------------------------------------------
# Offline seed / compaction (wraps the phase-A Gibbs sampler)
# ---------------------------------------------------------------------------
"""
compact!(cat, cfg; sweeps, restarts, rng) -> NamedTuple
Re-cluster the *entire* `binary/` pile with the offline collapsed Gibbs sampler
and adopt the winning partition as the catalog's clusters (DESIGN §5). This is
both the **seed** on first run (an empty catalog has no clusters, so the live
sweep alone would send everything to background) and the **periodic compaction**
that merges drifted clusters / splits bloated ones later.
Ids are taken from the Gibbs partition and frozen; because compaction re-derives
the whole partition, this is a wholesale replace of `clusters`/`examples`, and
every file in the pile is marked processed. Callers run this on an explicit
schedule (e.g. `--compact`), never on the latency path.
"""
function compact!(cat::Catalog, cfg::Config;
sweeps::Integer=150, restarts::Integer=6,
rng::AbstractRNG=Random.default_rng())
paths = binary_files(cfg.cluster_dir)
if isempty(paths)
return (; n_files=0, n_clusters=length(cat.clusters), n_bg=0)
end
X = header_matrix(paths; n=cat.n)
result = gibbs_cluster(X; α=cfg.cluster_alpha, β=cfg.cluster_pseudocount,
bg_mass=cfg.cluster_bg_mass, sweeps=sweeps,
restarts=restarts, rng=rng)
empty!(cat.clusters)
empty!(cat.examples)
empty!(cat.processed)
for (id, c) in result.clusters
cat.clusters[id] = c
end
cat.next_id = (isempty(result.clusters) ? 0 : maximum(keys(result.clusters))) + 1
for (j, path) in enumerate(paths)
base = basename(path)
push!(cat.processed, base)
z = result.assignments[j]
z > 0 && record_example!(cat, z, base)
end
n_bg = count(==(0), result.assignments)
return (; n_files=length(paths), n_clusters=length(cat.clusters), n_bg)
end
# ---------------------------------------------------------------------------
# Nominations (closing the loop to the classifier — DESIGN §6)
# ---------------------------------------------------------------------------
"Render a signature (from `signature`) into a human-readable hex template:
two hex digits for a required byte, `EOF` for a required past-EOF, `??` for a
wildcard. This is what a human eyeballs at the promotion gate."
function signature_hex(sig::AbstractVector)
parts = map(sig) do s
s === nothing ? "??" :
s == PAST_EOF ? "EOF" :
string(s; base=16, pad=2)
end
return join(parts, " ")
end
"The required (non-wildcard) positions of a signature as `(position, byte)`
records; `byte` is the raw 0255 value, or the string `\"past_eof\"`."
function signature_magic(sig::AbstractVector)
magic = Vector{NamedTuple{(:position, :byte),Tuple{Int,Any}}}()
for (i, s) in enumerate(sig)
s === nothing && continue
push!(magic, (position=i, byte=(s == PAST_EOF ? "past_eof" : s)))
end
return magic
end
"""
write_nominations!(cat, cfg) -> Vector{String}
Write one JSON nomination per promotable cluster (`is_promotable`, DESIGN §6)
into `cfg.nominated_dir`, each carrying the cluster's signature (hex template +
required magic positions), member count, and example filenames — everything a
human needs to glance and promote. The background (id 0) is never a cluster here,
so it can never be nominated, by construction.
Nominations are rewritten every sweep (membership only grows), so each file is
durably replaced via the same temp→fsync→rename→fsync-dir commit as the catalog.
Returns the paths written. Non-promotable clusters are left alone — a cluster
that *was* nominated and later fell below threshold cannot happen (members only
grow), so there is nothing to retract.
"""
function write_nominations!(cat::Catalog, cfg::Config)
mkpath(cfg.nominated_dir)
written = String[]
for (id, c) in sort(collect(cat.clusters); by=first)
sig = signature(c) # default β — signature β is
# decoupled from clustering β (DESIGN §11.3a)
is_promotable(c, sig; min_members=cfg.promote_min_members,
min_magic=cfg.promote_min_magic) || continue
payload = (
cluster_id = id,
members = c.members,
magic_length = magic_positions(sig),
signature_hex = signature_hex(sig),
magic = signature_magic(sig),
examples = get(cat.examples, id, String[]),
)
dest = joinpath(cfg.nominated_dir, "cluster-$(id).json")
tmp = string(dest, ".tmp")
open(tmp, "w") do io
write(io, JSON3.write(payload))
flush(io)
fsync_fd(fd(io))
end
mv(tmp, dest; force=true)
push!(written, dest)
end
fsync_dir(cfg.nominated_dir)
return written
end
# ---------------------------------------------------------------------------
# Orchestration + CLI (the periodic single-owner process — DESIGN §9)
# ---------------------------------------------------------------------------
"""
run_cluster_sweep(cfg; compact=false, rng) -> NamedTuple
One end-to-end pass of the single-owner stage-5 process: load the durable
catalog, either `compact!` (offline Gibbs — used to seed on first run or to
recompact) or `catalog_sweep!` (the incremental live path), then durably persist
the catalog and (re)write nominations. This is the whole job the cron/periodic
runner performs; `bin/cluster_sweep.jl` is a thin shell around it.
`compact` is forced when the catalog is empty (no clusters AND nothing processed
yet): a first live sweep against no clusters would send every file to background,
so the catalog must be seeded by an offline Gibbs pass before it can assign.
"""
function run_cluster_sweep(cfg::Config; compact::Bool=false,
rng::AbstractRNG=Random.default_rng())
cat = load_catalog(cfg.cluster_catalog_path; n=cfg.cluster_n)
if cat.n != cfg.cluster_n
@warn "configured cluster_n differs from the catalog's frozen window; using the catalog's" catalog_n=cat.n configured_n=cfg.cluster_n
end
is_empty = isempty(cat.clusters) && isempty(cat.processed)
mode = (compact || is_empty) ? :compact : :sweep
summary = mode === :compact ? compact!(cat, cfg; rng=rng) : catalog_sweep!(cat, cfg)
save_catalog!(cfg.cluster_catalog_path, cat)
nominated = write_nominations!(cat, cfg)
return (; mode, summary, n_nominated=length(nominated), nominated,
n_clusters=length(cat.clusters), n_processed=length(cat.processed))
end
"""
cluster_sweep_cli(args)
Entry point for `bin/cluster_sweep.jl`. Builds a `Config` from the environment,
runs one `run_cluster_sweep`, and logs a one-line summary. `--compact` forces the
offline Gibbs re-cluster (seed / periodic compaction) instead of the incremental
live sweep.
"""
function cluster_sweep_cli(args::AbstractVector{<:AbstractString}=String[])
compact = "--compact" in args
cfg = config_from_env()
ensure_dirs(cfg)
@info "stage-5 sweep starting" catalog=cfg.cluster_catalog_path input=cfg.cluster_dir compact=compact
r = run_cluster_sweep(cfg; compact=compact)
@info "stage-5 sweep complete" mode=r.mode clusters=r.n_clusters processed=r.n_processed nominated=r.n_nominated summary=r.summary
return r
end

View File

@@ -44,6 +44,11 @@ Base.@kwdef struct Config
cluster_bg_mass::Float64 = 5.0 # fixed mass of the uniform background 'junk drawer' cluster_bg_mass::Float64 = 5.0 # fixed mass of the uniform background 'junk drawer'
promote_min_members::Int = 20 # cluster size threshold for promotion nomination promote_min_members::Int = 20 # cluster size threshold for promotion nomination
promote_min_magic::Int = 3 # required fixed signature positions for nomination promote_min_magic::Int = 3 # required fixed signature positions for nomination
# The stage-5 catalog is a single durable file mutated by the one sweep
# process (never a worker), and nominations are written as one file per
# promotable cluster for a human to glance at before promoting (DESIGN §9/§6).
cluster_catalog_path::String = "data/catalog.json" # durable phase-B catalog (sufficient stats + processed set)
nominated_dir::String = "data/nominated" # one JSON per self-nominated cluster, awaiting a human promote
end end
""" """
@@ -62,7 +67,8 @@ Recognised variables:
FS_DONE_DIR, FS_TEXT_DONE_DIR, FS_FAILED_DIR, FS_MODEL_PATH, FS_DONE_DIR, FS_TEXT_DONE_DIR, FS_FAILED_DIR, FS_MODEL_PATH,
FS_EXIFTOOL_TIMEOUT, FS_LINGUIST_TIMEOUT, FS_EXIFTOOL_TIMEOUT, FS_LINGUIST_TIMEOUT,
FS_CLUSTER_DIR, FS_CLUSTER_N, FS_CLUSTER_ALPHA, FS_CLUSTER_PSEUDOCOUNT, FS_CLUSTER_DIR, FS_CLUSTER_N, FS_CLUSTER_ALPHA, FS_CLUSTER_PSEUDOCOUNT,
FS_CLUSTER_BG_MASS, FS_PROMOTE_MIN_MEMBERS, FS_PROMOTE_MIN_MAGIC FS_CLUSTER_BG_MASS, FS_PROMOTE_MIN_MEMBERS, FS_PROMOTE_MIN_MAGIC,
FS_CLUSTER_CATALOG, FS_NOMINATED_DIR
""" """
function config_from_env(; host=nothing, port=nothing, worker_count=nothing, function config_from_env(; host=nothing, port=nothing, worker_count=nothing,
queue_capacity=nothing, known_worker_count=nothing, queue_capacity=nothing, known_worker_count=nothing,
@@ -75,7 +81,8 @@ function config_from_env(; host=nothing, port=nothing, worker_count=nothing,
linguist_timeout=nothing, cluster_dir=nothing, cluster_n=nothing, linguist_timeout=nothing, cluster_dir=nothing, cluster_n=nothing,
cluster_alpha=nothing, cluster_pseudocount=nothing, cluster_alpha=nothing, cluster_pseudocount=nothing,
cluster_bg_mass=nothing, promote_min_members=nothing, cluster_bg_mass=nothing, promote_min_members=nothing,
promote_min_magic=nothing) promote_min_magic=nothing, cluster_catalog_path=nothing,
nominated_dir=nothing)
Config( Config(
host = something(host, get(ENV, "FS_HOST", "127.0.0.1")), host = something(host, get(ENV, "FS_HOST", "127.0.0.1")),
port = something(port, parse(Int, get(ENV, "FS_PORT", "8080"))), port = something(port, parse(Int, get(ENV, "FS_PORT", "8080"))),
@@ -105,13 +112,16 @@ function config_from_env(; host=nothing, port=nothing, worker_count=nothing,
cluster_bg_mass = something(cluster_bg_mass, parse(Float64, get(ENV, "FS_CLUSTER_BG_MASS", "5.0"))), cluster_bg_mass = something(cluster_bg_mass, parse(Float64, get(ENV, "FS_CLUSTER_BG_MASS", "5.0"))),
promote_min_members = something(promote_min_members, parse(Int, get(ENV, "FS_PROMOTE_MIN_MEMBERS", "20"))), promote_min_members = something(promote_min_members, parse(Int, get(ENV, "FS_PROMOTE_MIN_MEMBERS", "20"))),
promote_min_magic = something(promote_min_magic, parse(Int, get(ENV, "FS_PROMOTE_MIN_MAGIC", "3"))), promote_min_magic = something(promote_min_magic, parse(Int, get(ENV, "FS_PROMOTE_MIN_MAGIC", "3"))),
cluster_catalog_path = something(cluster_catalog_path, get(ENV, "FS_CLUSTER_CATALOG", "data/catalog.json")),
nominated_dir = something(nominated_dir, get(ENV, "FS_NOMINATED_DIR", "data/nominated")),
) )
end end
"Create all the pipeline-stage directories if they don't already exist." "Create all the pipeline-stage directories if they don't already exist."
function ensure_dirs(cfg::Config) function ensure_dirs(cfg::Config)
for d in (cfg.spool_dir, cfg.known_dir, cfg.unknown_dir, cfg.binary_dir, for d in (cfg.spool_dir, cfg.known_dir, cfg.unknown_dir, cfg.binary_dir,
cfg.text_dir, cfg.done_dir, cfg.text_done_dir, cfg.failed_dir) cfg.text_dir, cfg.done_dir, cfg.text_done_dir, cfg.failed_dir,
cfg.nominated_dir)
mkpath(d) mkpath(d)
end end
return nothing return nothing

View File

@@ -15,7 +15,10 @@ using FileServer: Job, Config, ChannelQueue, enqueue!, dequeue!, length,
header_symbols, header_matrix, ClusterStats, add!, remove!, header_symbols, header_matrix, ClusterStats, add!, remove!,
log_predictive, loggamma, gibbs_cluster, assign_file, log_predictive, loggamma, gibbs_cluster, assign_file,
signature, magic_positions, is_promotable, signature, magic_positions, is_promotable,
adjusted_rand_index, v_measure, HEADER_N, ALPHABET, PAST_EOF adjusted_rand_index, v_measure, HEADER_N, ALPHABET, PAST_EOF,
Catalog, load_catalog, save_catalog!, catalog_sweep!, compact!,
write_nominations!, run_cluster_sweep, binary_files, record_example!,
signature_hex, ensure_dirs
using Random: MersenneTwister using Random: MersenneTwister
using Languages: LanguageDetector using Languages: LanguageDetector
@@ -36,6 +39,9 @@ function tmp_config(root; kwargs...)
done_dir = joinpath(root, "done"), done_dir = joinpath(root, "done"),
text_done_dir = joinpath(root, "text_done"), text_done_dir = joinpath(root, "text_done"),
failed_dir = joinpath(root, "failed"), failed_dir = joinpath(root, "failed"),
cluster_dir = joinpath(root, "binary"), # stage-5 sweeps the binary sink
cluster_catalog_path = joinpath(root, "catalog.json"),
nominated_dir = joinpath(root, "nominated"),
kwargs..., kwargs...,
) )
FileServer.ensure_dirs(cfg) FileServer.ensure_dirs(cfg)
@@ -525,4 +531,178 @@ end
end end
end end
# Helper: write a "file" of raw bytes into a dir with a UUID-ish unique name,
# returning its path. Mirrors what stage-3 deposits into binary/.
function drop_binary(dir, bytes; name=string(rand(UInt128)))
mkpath(dir)
p = joinpath(dir, name)
open(p, "w") do io; write(io, Vector{UInt8}(bytes)); end
return p
end
@testset "catalog: durable save/load round-trip" begin
mktempdir() do root
n = 8
cat = Catalog(n)
c = ClusterStats(n)
add!(c, [0xCA+1, 0xFE+1, 0xBA+1, 0xBE+1, 1, 2, 3, 4])
add!(c, [0xCA+1, 0xFE+1, 0xBA+1, 0xBE+1, 5, 6, 7, 8])
cat.clusters[7] = c
cat.next_id = 8
record_example!(cat, 7, "alpha.bin")
push!(cat.processed, "alpha.bin"); push!(cat.processed, "beta.bin")
path = joinpath(root, "catalog.json")
save_catalog!(path, cat)
@test isfile(path)
back = load_catalog(path; n=n)
@test back.n == n
@test back.next_id == 8
@test back.processed == cat.processed
@test haskey(back.clusters, 7)
@test back.clusters[7].members == 2
@test back.clusters[7].counts == c.counts # sparse round-trips exactly
@test back.examples[7] == ["alpha.bin"]
end
end
@testset "catalog: load of a missing file is a fresh catalog" begin
mktempdir() do root
cat = load_catalog(joinpath(root, "nope.json"); n=16)
@test cat.n == 16
@test isempty(cat.clusters)
@test isempty(cat.processed)
@test cat.next_id == 1
end
end
@testset "catalog: binary_files skips sidecars, tmp, dirs; sorts" begin
mktempdir() do root
drop_binary(root, "a"; name="002-file")
drop_binary(root, "b"; name="001-file")
write(joinpath(root, "003-file.meta.json"), "{}") # sidecar
write(joinpath(root, "004-file.tmp"), "x") # scratch
mkpath(joinpath(root, "subdir")) # not a file
fs = binary_files(root)
@test basename.(fs) == ["001-file", "002-file"]
end
end
@testset "catalog: incremental sweep grows an existing cluster" begin
mktempdir() do root
n = 8
cfg = tmp_config(root; cluster_n=n, cluster_alpha=1.0,
cluster_pseudocount=0.1, cluster_bg_mass=5.0)
# Seed a strong cluster (magic 0xCA 0xFE 0xBA 0xBE, random tail).
cat = Catalog(n)
c = ClusterStats(n)
rng = MersenneTwister(3)
for _ in 1:40
add!(c, vcat([0xCA,0xFE,0xBA,0xBE] .+ 1, rand(rng, 1:256, 4)))
end
cat.clusters[1] = c
cat.next_id = 2
# A brand-new file that matches the magic must JOIN cluster 1.
drop_binary(cfg.cluster_dir, vcat(UInt8[0xCA,0xFE,0xBA,0xBE], rand(rng, UInt8, 4)); name="match-01")
# A structureless random blob must park in the background.
drop_binary(cfg.cluster_dir, rand(rng, UInt8, 64); name="blob-01")
s = catalog_sweep!(cat, cfg)
@test s.n_seen == 2
@test s.n_joined == 1
@test s.n_bg == 1
@test s.n_minted == 0
@test cat.clusters[1].members == 41 # grew by the matching file
@test "match-01" in cat.processed
@test "blob-01" in cat.processed
# Re-sweeping the same pile is idempotent — nothing new is seen.
s2 = catalog_sweep!(cat, cfg)
@test s2.n_seen == 0
@test cat.clusters[1].members == 41
end
end
@testset "catalog: §10.1 nothing from noise (end-to-end, no promotion)" begin
mktempdir() do root
n = 32
cfg = tmp_config(root; cluster_n=n, cluster_alpha=1.0,
cluster_pseudocount=0.1, cluster_bg_mass=5.0,
promote_min_members=20, promote_min_magic=3)
rng = MersenneTwister(10)
# The §10.1 pile: 20 small random blobs + 1 lone structured "PDF".
for i in 1:20
drop_binary(cfg.cluster_dir, rand(rng, UInt8, 40); name="blob-$(lpad(i,2,'0'))")
end
drop_binary(cfg.cluster_dir, vcat(UInt8[0x25,0x50,0x44,0x46], rand(rng, UInt8, 60)); name="lone-pdf")
# First run auto-compacts (empty catalog) to seed, then persists + nominates.
r = run_cluster_sweep(cfg; rng=MersenneTwister(10))
@test r.mode == :compact
@test isfile(cfg.cluster_catalog_path)
# The mission-critical assertion: ZERO promoted clusters from pure noise.
@test r.n_nominated == 0
@test isempty(readdir(cfg.nominated_dir))
# Every file was accounted for (clustered-as-singleton or background).
@test r.n_processed == 21
end
end
@testset "catalog: a real recurring format self-nominates" begin
mktempdir() do root
n = 32
# β=0.1 over-splits a format into pure sub-clusters (DESIGN §11 known
# limitation) — each still carries the full magic and nominates
# independently, so a modest min_members catches those sub-clusters.
cfg = tmp_config(root; cluster_n=n, cluster_alpha=1.0,
cluster_pseudocount=0.1, cluster_bg_mass=5.0,
promote_min_members=10, promote_min_magic=3)
rng = MersenneTwister(21)
# 30 files sharing a fixed 6-byte magic then random payload — a format.
magic = UInt8[0x89, 0x46, 0x4d, 0x54, 0x21, 0x0a]
for i in 1:30
drop_binary(cfg.cluster_dir, vcat(magic, rand(rng, UInt8, 40)); name="fmt-$(lpad(i,2,'0'))")
end
r = run_cluster_sweep(cfg; rng=MersenneTwister(21))
@test r.n_nominated >= 1
files = readdir(cfg.nominated_dir; join=true)
@test !isempty(files)
payload = JSON3.read(read(first(files), String))
@test payload.members >= 10
@test payload.magic_length >= 3
# The hex template exposes the shared magic bytes for the human gate.
@test occursin("89 46 4d 54", payload.signature_hex)
end
end
@testset "catalog: seeded catalog then live-assigns a matching arrival" begin
mktempdir() do root
n = 32
cfg = tmp_config(root; cluster_n=n, cluster_alpha=1.0,
cluster_pseudocount=0.1, cluster_bg_mass=5.0,
promote_min_members=20, promote_min_magic=3)
rng = MersenneTwister(31)
magic = UInt8[0x7a, 0x7a, 0x01, 0x02, 0x03]
for i in 1:25
drop_binary(cfg.cluster_dir, vcat(magic, rand(rng, UInt8, 40)); name="seed-$(lpad(i,2,'0'))")
end
# Seed pass.
run_cluster_sweep(cfg; rng=MersenneTwister(31))
cat = load_catalog(cfg.cluster_catalog_path; n=n)
@test !isempty(cat.clusters)
members_before = sum(c.members for c in values(cat.clusters))
# A new matching file arrives; an incremental sweep must fold it in
# (mode :sweep, not compact) without re-clustering the world.
drop_binary(cfg.cluster_dir, vcat(magic, rand(rng, UInt8, 40)); name="arrival-01")
r2 = run_cluster_sweep(cfg; rng=MersenneTwister(99))
@test r2.mode == :sweep
cat2 = load_catalog(cfg.cluster_catalog_path; n=n)
members_after = sum(c.members for c in values(cat2.clusters))
@test members_after == members_before + 1 # the arrival joined a cluster
end
end
end end