One owner per spool directory; adopt a crashed process's spool; refuse shared filesystems - #163
Merged
Merged
Conversation
….lock
Recover() deletes every .open file its own Spool object is not writing.
Nothing kept a second process from running it against a live writer: the
old test_native_live_spool.py only showed that an upload SCAN (ListPending)
spares the writer's temp file, and one process per spool was left to the
caller (storage_service.h said so). A second process pointed at the same
directory -- two jobs on one node, a restart racing its predecessor --
swept the writer's in-flight temp file and failed its stage with "cannot
link ready file".
The spool now has an owner lock, flock(LOCK_EX|LOCK_NB) on
<dir>/.owner.lock, recording "<host> <pid>" in the file:
* SpoolConfig.owner_lock = take (the default) takes it in Open(), before
the accounting walk and so before any Recover(), and holds it for the
Spool's life. A second process is refused with kOwned: "spool directory
X is owned by pid N on host H (it holds X/.owner.lock)". The kernel
drops the lock with its holder, SIGKILL included.
* held_by_caller takes none, for a process that runs a sink and a service
on one directory: two takes in ONE process refuse each other (flock
binds to an open file description), so such a process holds one
SpoolOwnerLock and opens both Spools with held_by_caller. Open()
refuses it when nothing holds the lock.
* A directory nested under, or containing, an owned one (a .owner.lock,
held or not) is refused: Scan walks recursively.
* NFS and Lustre are refused by statfs f_type (allow_shared_filesystem
overrides; SetFilesystemTypeForTesting stands in for statfs), since
flock there does not keep out a process on another node.
* Scan, Open's accounting and the capacity reconcile skip <dir>/_refs/,
where the upload handoff's ref files will live (plan section 2.4).
SpoolOwnerLock is the lock on its own: Acquire creates a new directory
beside its already-held lock file and renames it into place, so no scan of
the parent meets it unowned; TryAdopt locks an existing directory whose
owner is gone; ReleaseAndRemoveIfEmpty removes a drained one. spool.h also
gains the section 2.3 layout, <base>/<catalog_key>/r<rank>-<incarnation>/,
for the service's adoption and the engine, next.
The drivers accept "owner_lock". Every Spool takes by default, so two in
one process refuse each other: test_spool_reservations.cpp's second Spool
object on a root now opens held_by_caller, and the storage service's and
the pack sink's Spools -- which test_native_capture_chain_live.py and the
storage live suite run in one process -- refuse each other until the next
commit gives both an owner-lock mode.
The Python DurablePackSpool (spool.py:208-214) stays unlocked and its
recover() still deletes every .open file: the C++ spool is deliberately
stricter, and that is not ported to the reference.
Tests, each new one red before the change for the reason it names:
* test_native_spool_owner_lock.py (drivers, across processes): a second
process was admitted ({"ok": true}) where it is now refused naming the
holder's pid and host; no lock file was written; nested and containing
roots were admitted; held_by_caller with nothing held, and an unknown
mode, were admitted; recovery deleted <root>/_refs/*.open; Open
counted 150 bytes under _refs where it now counts 0.
* test_native_spool_owner_lock_unit.py compiles
tests/native/test_spool_owner_lock.cpp (red: the API did not exist):
two takes in one process refuse each other; held_by_caller beside a
SpoolOwnerLock stages and lists; a forked holder is named and its lock
freed by SIGKILL; nesting both ways; NFS/Lustre refused through the
test seam, admitted with the override; adoption's try-lock; the
layout.
* test_native_live_spool.py is rewritten: the upload scan beside a paused
live writer is now refused, naming the writer (red: it ran); and a
writer SIGKILLed mid-stage is recovered by the next process, its stale
.open swept and its sealed pack uploaded (green before too: it guards
that the lock does not outlive its holder).
* tests/test_native_spool.py (all 20) and test_native_spool_reservations.py
pass; so does every other cpu test under tests/test_native_*.py.
… siblings
With the spool's owner lock (previous commit) every Spool takes its
directory by default, and the storage service's and the pack sink's Spools
became two takes: in ONE process they refuse each other, since flock binds
to an open file description. Run at the previous commit,
test_native_capture_chain_live.py -- the real NativePackSink and the real
CaptureStorageService on one root in one process, as the engine runs them
-- fails every case with exactly that:
NativePackSink: sink start failed: cannot open spool: spool directory
.../spool is owned by pid <this process> on host ...
Both now take the directory's owner-lock mode: StorageServiceConfig
.spool_owner_lock and SinkConfig.spool_owner_lock, "take" or
"held_by_caller" (with the shared-filesystem override beside it), bound as
the service dict's spool_owner_lock and NativePackSink(owner_lock=...).
_dmi_native_store binds the lock itself -- SpoolOwnerLock (a context
manager; refused with SpoolOwnedError, a RuntimeError naming the holder,
or ValueError for a shared filesystem or a nested directory), spool_owner
(who holds a directory), spool_rank_directory / spool_catalog_key (the
section 2.3 layout) and a statfs test seam -- so a process holds one lock
and opens both Spools held_by_caller. The chain suite now does exactly
that, which is the engine's arrangement (next commit).
adopt_sibling_spools: the service's directory is one rank directory
<base>/<catalog_key>/r<rank>-<incarnation>/ of the layout -- refused at
construction unless it is, and unless the key is THIS catalog's, since
adopted packs are indexed here -- and start(), after the lease and the
sweep of its own directory and before the reconcile, adopts every sibling
whose owner is dead (its lock can be taken; the kernel dropped it with the
process): Recover sweeps its .open files, the uploader sends its .ready
packs under the keys their paths give, they are indexed like the service's
own (or owed in pending_index_), and the directory is removed once
nothing but its lock file is left. A live sibling is skipped. Adoption
follows the cycle's upload rules -- no upload without the lease or while a
pack is owed -- so while the service cannot index, a dead spool's packs stay
where they are durable, and a sibling left undrained is retried by the
loop and keeps flush() from reporting drained. The snapshot says
adopted_spools, adopted_packs and adoption_owed. E2b extends this with ref
replay.
The storage live suite stages and uploads through driver processes while
a service has the same directory open, and builds second services on
directories the first still holds. Its harness now holds each directory's
lock for the test and opens every service and driver held_by_caller --
the engine's arrangement across processes. No assertion changed.
Tests:
* test_native_spool_ownership.py (cpu, new): the real sink and the real
service share one rank directory under one SpoolOwnerLock and stage a
pack; two takes refuse each other; a second process is refused
naming this pid and host; spool_owner; nesting; NFS refused through
the seam and admitted with the override; the layout; adoption refused
for a directory outside the layout or under another catalog's key;
a drained directory removed with its lock, a non-empty one kept. Red at
the previous commit (the bindings did not exist), but for the case
showing two takes refusing each other, which the previous commit
introduced.
* test_native_spool_adoption_live.py (ClickHouse + fake S3, new): a
child running the engine's composition gets part of its capture
indexed, stages the rest and is SIGKILLed holding the lease; a new
incarnation waits the lease out, adopts the dead directory (its stale
.open swept), and every capture of the dead process reads back
byte-equal, while a live sibling is left alone. With the object store
cut the dead directory stays as it was, flush() times out, and the
loop adopts it once the store is back. Red with adopt_siblings()
stubbed out: adopted_spools 0 != 1, and adoption_owed False.
* test_native_capture_chain_live.py, test_native_capture_storage_live.py
(48) and the new live suite pass; the lease byte-identity suite is
untouched.
Under storage_backend="persistent" with capture_storage_config, the engine
pointed its storage service and its native pack sink at
capture_sink_config.spool_root itself. Now it claims a directory of its
own there, the plan's section 2.3 layout:
<spool_root>/<catalog_key>/r<rank>-<incarnation>/
catalog_key being the first 12 hex digits of sha256 of
database/table_prefix/store_id, rank torchrun's RANK (0 when unset or not a
rank: it only labels the directory), and the incarnation fresh for every
create_record_runtime, so two jobs on a node, or a restart racing its
predecessor, never share a directory. claim_spool_directory creates it
with its owner lock already held, BEFORE the service is built; the service
starts (sweeping the directory, then adopting the dead siblings under the
same catalog key) and only then does the sink open it, both held_by_caller.
The claim is let go last -- after the sink is sealed, the ring stopped and
the service drained and stopped -- on close(), on enable_ring_transport
replacing the record ring, and on a create_record_runtime that fails after
the service started (or whose service fails to start). A drained directory
is removed then; one still holding packs stays for the next process on the
node to adopt.
Without capture_storage_config the sink takes spool_root itself, as the
one owner of packs something else will drain. With an explicit
record_sink the service drains spool_root as that sink writes it --
unswept, as before, and adopting nothing -- opening it held_by_caller when
this process already holds its lock (a NativePackSink the caller built
took it; a second take would refuse this very process), and taking it
otherwise, which another process's lock refuses by name.
NativeSinkConfig.spool_allow_shared_filesystem carries the NFS/Lustre
override to the lock and both Spools. flush_and_wait's TimeoutError now
says when a dead process's spool is what is still owed. The v1 contract
describes the directory, its owner, adoption and the node-local rule.
Tests (test_native_capture_storage_wiring.py, the native modules faked),
each red against the previous engine for the reason it names:
* the order is now layout, lock, service construct, service start, sink
open -- on the rank directory, both held_by_caller, adoption on (was:
service first, on spool_root, no lock);
* RANK names the directory, and a RANK that is not a rank names rank 0;
* the shared-filesystem override reaches the lock and both Spools (the
field did not exist);
* a service that fails to start leaves no lock and opens no sink;
* close(), a failed drain, a failed attach and replacing the record ring
each release the lock after the service stops;
* an explicit sink: spool_root, unswept, no adoption, take -- or
held_by_caller when this process holds the lock;
* with no storage config the sink takes spool_root itself;
* a flush that times out on an owed adoption says so (the message named
only the unindexed packs).
The full cpu tier passes.
Spool::Open with owner_lock=held_by_caller only probed that SOMEONE held
<root>/.owner.lock. That passes in exactly the dangerous case: another live
process owns the directory. So any process could open a live writer's
spool held_by_caller and Recover it, deleting the writer's in-flight .open
file -- the "cannot link ready file" race B6 exists to make impossible.
Reproduced with tests/native/live_spool_stage.cpp paused mid-stage (take):
conformance_spool recover with take was refused naming the writer, the same
op with held_by_caller returned ok, the .open file was gone and the writer
exited 1 with "cannot link ready file".
held_by_caller now also requires that one of THIS process's descriptors
holds the lock (SpoolOwnedByThisProcess): /proc/self/fdinfo lists the
flocks each open file description holds, so the kernel answers, not the
owner record, which is written after the lock is taken and names a pid
that means nothing across pid namespaces. The record (host, pid) is the
fallback only where /proc cannot be read. Beside another process's lock
the open is refused with kOwned (SpoolOwnedError from the store module),
naming the holder; with nothing holding the lock it stays kBadArgument.
The tests had encoded the race as intended behaviour, and change with it:
* test_native_spool_owner_lock.py: held_by_caller beside a conformance_sink
process's lock is now refused, naming its pid and host (it asserted the
recover succeeded);
* test_native_live_spool.py: the finding's repro -- held_by_caller beside
the paused writer is refused, its .open file survives, and the writer's
stage completes (was: ok, file gone, writer rc=1);
* test_native_spool_ownership.py: the service and the sink opened
held_by_caller beside another process's SpoolOwnerLock are refused
naming it (both opened);
* test_spool_owner_lock.cpp: the same across a fork, in process (kOk).
Each was red against the previous spool.
The storage live harness ran its conformance_sink and conformance_store
drivers held_by_caller beside the test process's lock, which is what this
refuses. The drivers now take, as the plan has standalone callers do: the
sink driver stages into a scratch directory it owns and _stage renames the
sealed packs into the spool (one rename per ready file, so a scanning
service meets whole files), and the store driver uploads from a spool no
service of the test has opened yet, which is already true at all four call
sites. Services keep holding the harness's lock held_by_caller, in process.
…locks above
Two defects in SpoolOwnerLock::Acquire's nesting check, one ordering and
one rule.
Ordering: CheckNotNested ran before the lock was taken, so an outer
directory and one nested in it, taken at the same moment by two processes,
each checked before the other's lock file existed, and both won. Two
owners of overlapping trees is what the check exists to prevent: the
outer one's recursive Scan, UploadPending or Recover reaches the inner live
directory, sweeping its .open files and uploading its packs under the wrong
keys. The engine really produces the pair -- a sink-only or explicit
record_sink run takes spool_root, a default persistent run claims
spool_root/<key>/r0-<inc>. Acquire now takes the lock (LockInPlace or
CreateLocked, which renames the new directory into place first) and checks
after, so each side publishes before it looks and at least one sees the
other. A refused take lets go of its lock and leaves nothing behind: the
directory it created, while still empty, is removed, and a lock file it
added to an existing directory is unlinked.
Rule: an ANCESTOR with a .owner.lock file refused, held or not. A take
never deletes its lock file, so after one sink-only run on spool_root, or
one explicit-record_sink run whose service took it (the documented
rollback), every later default run was refused as "nested under the owned
spool directory <spool_root>" until someone deleted the file by hand. An
ancestor now refuses only while its lock is held. Safety does not need
more: the next take of that ancestor walks its descendants, meets this
directory's lock file, and is refused. Descendants still refuse on the
lock file alone, held or not -- a dead directory's packs are its
successor's to adopt, not an outer spool's to sweep.
Tests, each red before:
* test_spool_owner_lock.cpp: 200 trials of two forked processes released
at once on an outer directory and a rank directory in it -- both won in
195 of 200 before, never now; a stale lock file above no longer refuses
a nested take, and the lock file it leaves still refuses the next take
of the outer directory; a refused take leaves no directory or lock file;
* test_native_spool_owner_lock.py: the stale-lock test asserted the
nested refusal, and now asserts the admission and the outer refusal;
* test_native_spool_ownership.py: claim_spool_directory under a
spool_root that a SpoolOwnerLock once held succeeds (it raised
ValueError "nested under the owned spool directory").
…aves
SpoolOwnerLock::Acquire creates a new spool directory as a staging copy,
<parent>/.<name>.<8 hex>.creating, with its lock file created and flocked
inside, then renames it into place. A claim killed anywhere between the
lock file's open and the rename (or a power loss before the parent's
fsync) leaves that copy and its lock file behind, and nothing removed it:
adoption skipped it, since the name is not a rank directory, and the
nesting check's descendant walk counted its lock file. So a later take of
<spool_root> or <spool_root>/<key> -- a sink-only run, or the explicit
record_sink rollback, both of which take spool_root -- was refused for
good as "contains the owned spool directory .../.r0-....creating", by a
directory nobody owns that holds no pack.
* The nesting check skips a staging copy whose lock nobody holds. One
that is held is a claim in progress and still refuses; one not yet
flocked publishes after the walk, and its own ancestor check (which
runs after its rename) then sees this lock.
* Adoption clears the dead ones it meets under the catalog key: a staging
copy older than 60 s whose lock it can take is removed with it. A
younger one is left alone, since a claim between its mkdir and its
flock holds no lock yet and clearing its copy would fail it. Nothing is
owed either way.
* IsSpoolClaimStagingName names the pattern, and CreateLocked builds it.
Tests: test_spool_owner_lock.cpp -- with an unheld staging copy under
root/<key>, taking root and then root/<key> succeed (refused before), and
a held one still refuses the outer take; test_native_spool_adoption_live.py
-- the successor of a SIGKILLed process removes an hour-old staging copy
beside the dead directory and leaves a fresh one in place (the old one
survived before).
flock binds to an open file description, and a child created with fork() and no exec shares its parent's. O_CLOEXEC only acts at exec, so a fork-started worker (a multiprocessing or DataLoader worker under the fork start method, the default on Linux for the venv's Python 3.10) kept the spool's owner lock after its parent was SIGKILLed. The dead parent's directory then read as owned, naming the dead pid; its successor's adopt_sibling took kOwned for "it lives" and never came back to it, so its ready packs waited for the next restart after the worker exited. Reproduced from Python: the owner took SpoolOwnerLock, os.fork()ed a sleeping child and was SIGKILLed; spool_owner still named the dead pid and a new SpoolOwnerLock raised SpoolOwnedError until the child exited. Every held SpoolOwnerLock is now registered (per binary: each extension and driver that compiles spool.cpp keeps its own registry), and a pthread_atfork child handler closes the child's descriptor of each and marks it not held. It closes, never LOCK_UN: an unlock on the shared description would drop the parent's hold as well, and closing one of its descriptors does not. The registry mutex is held across the fork by the prepare handler, so no thread is mid-update when the child copies it. A kTake Spool's lock is a SpoolOwnerLock too, so the sink-only case is covered. posix_spawn and vfork run no atfork handler, and exec closes the descriptor there anyway. Tests, red before: test_spool_owner_lock.cpp forks an owner that takes the lock and forks a worker; the worker sees its own lock as not held while the parent's still is, and once the owner is SIGKILLed the lock is free with the worker alive (it stayed held, naming the dead owner). test_native_spool_ownership.py does the same through Python's os.fork.
adopt_sibling took a sibling whose lock was held for "it lives", which is not owed, so the loop never came back to it: a sibling owned at start() was left for the process's whole lifetime, possibly days for a serving process. The previous commit removes the fork-child case, where the owner was dead but the lock still read as held. Two remain. A sibling can be briefly held at start() -- a predecessor still inside its close() in a rolling restart, which leaves the directory with packs if it cannot drain them -- and a rank on the node can die while this one runs. Either way its packs waited for the next restart on the node. While the last pass found a live sibling, run_cycle now passes over the siblings again every adoption_recheck_interval_ns (30 s by default, 0 for never; a native config knob, not on NativeCaptureStorageConfig), under the same rules as an owed pass: with the lease, nothing owed, no upload of its own failing. A live sibling is still not owed, so flush() is not held up by another process's directory. adopt_sibling probes the lock without blocking first, so a live sibling costs one probe per pass rather than TryAdopt's few milliseconds of retries. The snapshot's live_siblings says how many the last pass left alone. Test (test_native_spool_adoption_live.py, red before): a sibling this process holds, with three packs staged by the real sink, is left alone at start (live_siblings 1, nothing adopted or owed); once its lock is released the service adopts it within the recheck interval, removes the directory, and every capture reads back from the catalog.
…ll stage
close() and enable_ring_transport's replacement released the engine's
spool claim in a finally once the service stopped, whether or not the sink
had sealed: _seal_capture_sink swallows a flush timeout (close_flush_timeout_s
is shared with the storage drain) and only logs it. The NativePackSink
outlives that release. RingEngine keeps record_sink_ after stop(),
NativePackSink::on_engine_release does nothing, and the user's RecordRuntime
keeps the RingEngine through its RingTransport. So its stagers carry on
and ~PackSink stages whatever it still holds, into a directory nobody owns
any more, or into a removed one that Spool::Stage recreates with no
.owner.lock. A successor's AdoptSiblings (another process on the node, or
this process's next engine) then takes the directory and its Recover
deletes the old sink's in-flight .open file, losing those records; or it
adopts a directory still being written. The review reproduced it with the
real sink and lock: one record submitted and not flushed, claim released and
the directory removed; destroying the sink recreated it with a ready pack and
no lock file. With a pack flushed and one record still queued, a
conformance_spool recover (take) from another process succeeded on the
directory while the sink was alive.
Dropping an engine without close() had the same effect: the SpoolClaim went
with it, and ~SpoolOwnerLock unlocked while the globally activated ring and
its sink kept capturing.
* _seal_capture_sink returns whether the sink sealed. The claim is let go
only then; otherwise the engine keeps the directory owned until the
process exits, logging that it does and which directory, and the next
process on the node adopts it.
* A SpoolClaim is registered in native_capture._HELD_SPOOL_CLAIMS until
release(), so garbage collection never unlocks it; only release() or
the process's exit does.
* A create_record_runtime that fails in attach still releases at once:
no record reached the sink.
* The v1 contract says so.
Tests, red before (test_native_capture_storage_wiring.py, fakes): a close
whose sink flush times out stops the service but leaves the lock held and
registered, and warns naming the directory (the lock was released); the
same for replacing the record ring. test_native_spool_ownership.py: a real
claim dropped and garbage-collected still reads as owned by this process
(it read as unowned).
…nly the names The adoption domain, <spool_root>/<catalog_key>/, was keyed by sha256(database/table_prefix/store_id)[:12] -- the plan's section 2.3 definition -- which names no server. With the defaults (default, dmi, s3) every deployment on a node that shares a spool_root got the same key, so the first to start after another one's crash adopted its dead directories: their packs went up through its own S3 client to its own bucket and were indexed into its own ClickHouse, and never reached their own. Nothing on the adoption path checks a pack's destination to stop that. A site-wide spool path with staging and production under one service account is enough. The key now hashes where the packs go: <database>/<table_prefix>/<store_id>\n clickhouse <clickhouse_host>:<clickhouse_port>\n s3 <s3_endpoint>/<s3_bucket> SpoolCatalogKey and SpoolRankDirectory take a SpoolDestination; the store module's spool_catalog_key / spool_rank_directory take its fields as a dict (every one required); NativeCaptureStorageConfig._spool_destination builds it for claim_spool_directory; and the service's adopt_sibling_spools check computes the same key from its ClickHouse connection, writer, S3 config and store id. The servers are hashed as spelled, so a successor that reaches the same store through a different spelling does not see the dead directory as its sibling: it waits for a process that spells it the same way. That is the safe side of the trade (packs delayed, not misdirected), and spool.h and the v1 contract say so. This departs from the plan's key definition, which needs amending to match. Tests: test_spool_owner_lock.cpp -- the key is that sha256, and changing any one of the seven fields changes it; test_native_spool_ownership.py -- the layout through the binding, a destination missing a field is a KeyError, a staging and a production config differing only in their servers claim directories under different keys, and a staging service refuses to adopt from under production's (they shared one before). The adoption live test's store-down case now has the dead process reach the store through the same switch URL as its successor.
The node-local check refused two statfs f_type values, NFS (0x6969) and
Lustre (0x0BD00BD0). Others give a flock that does not keep out a process
on another node just as well, and the owner lock is then no lock at all
across the nodes that mount the directory:
* BeeGFS (0x19830326), common for HPC scratch, keeps flock client-local
unless tuneUseGlobalFileLocks is set, and that is off by default;
* CIFS (0xFF534D42) and SMB2 (0xFE534D42);
* FUSE (0x65735546), which covers sshfs, s3fs, gcsfuse and the GlusterFS
client -- and local filesystems too (fuse-overlayfs, ntfs-3g), which
f_type cannot tell apart. Refusing a local one is loud and has the
override; admitting a network one is silent, so FUSE is refused, and
its refusal names the local case the override is for.
allow_shared_filesystem admits any of them, as before. The list still
needs maintenance, as the plan's risk section says.
Tests, red before: test_spool_owner_lock.cpp names all six and refuses a
lock and a Spool on each through the statfs seam, with the override still
admitting them, while ext4, xfs, overlayfs and tmpfs pass;
test_native_spool_ownership.py refuses each through the binding, naming
it.
The one service-level check that adoption leaves a live sibling alone wrote a file named "marker" into it and asserted it still read "live", plus live.held -- the test's own Python object, always true. Recover() unlinks only *.open files and validates only *.ready ones, so a service that swept a live sibling passed: the review's mutation, making the live branch open the sibling held_by_caller and Recover() it, left every CPU and live adoption test green. The live sibling now holds what a live writer's directory does: ready packs its real sink staged, and a stage in flight (a .open temp file). After the successor has adopted the dead directory and flushed, both are exactly where they were, none of the live packs' ids is in the object store, live_siblings is 1, and the catalog holds only the dead process's captures. The sibling's lock is held by this process, which is the case that matters: an adopter that took it for dead would get past the held_by_caller check, and only the liveness probe keeps it out. Checked: with the review's mutation applied to adopt_sibling's live branch, the test fails on the swept .open file (FileNotFoundError); as the code stands it passes.
adopt_sibling rechecks, for each sibling, that the lease is held and nothing is owed to the catalog (pending_index_ empty) before it uploads: the owner's decision that packs stay in the durable spool while the service cannot index them, applied to adoption. Nothing tested it. Deleting the line passed every CPU test and both live adoption tests; the only outage test cut the object store, not the catalog, and had one dead sibling. The new live test has two dead siblings, each staged by the real sink into its own rank directory. A switch in front of the fake S3 cuts a second switch, in front of ClickHouse, at the first upload: the first sibling's packs reach the object store, their index fails and they are owed in memory. The test asserts that the second sibling's ready packs are all still in its spool, adoption is owed, and flush() times out; then restores the catalog and asserts that both siblings are adopted and removed and every capture of both reads back from the catalog. Checked: with `|| !pending_index_.empty()` removed from adopt_sibling, the test fails on the second sibling's packs having left its spool; as the code stands it passes.
…tested
Three mutations passed every test, and one test checked less than its
name says:
* SpoolClaim.release() replaced by a plain release(): every drained
engine directory and its lock file kept. The wiring fake's lock always
reports removal, and nothing else called release() on a real claim.
test_releasing_a_claim_removes_its_directory_once_drained claims
through claim_spool_directory, and asserts that a drained directory is
removed (release() is True) and one holding a pack is kept, unowned,
with the claim out of the held registry either way.
* Spool::Open(held_by_caller) without its node-local check.
test_held_by_caller_still_refuses_a_shared_filesystem holds a real lock
and opens the service held_by_caller with the statfs seam reporting
NFS: refused, and admitted with the override.
* LockInPlace without its re-check that the locked file is still the one
at the path (a remover unlinks the lock file before it removes a
drained directory). This one is race-only, so the store gains a test
seam, SetLockOpenHookForTesting, called between the open and the
flock. test_spool_owner_lock.cpp unlinks the file there once: the take
must open and lock again (two calls), and what it holds must be the
file at the path -- ReadSpoolOwner sees it held and a rival TryAdopt is
refused.
* test_held_by_caller_with_nothing_held_is_refused passed a directory
that does not exist, so it exercised realpath() failing rather than
"exists, but nothing holds its lock". It now creates the directory,
and checks again with a lock file nobody holds.
Each new check was run against its mutation and failed (release() is
False; DID NOT RAISE; calls == 1 and the lock not seen at the path).
The previous commit brought the node-local refusal to NFS, Lustre,
BeeGFS, CIFS/SMB2 and FUSE. The review named more network filesystems
that clusters mount where a spool root could land, each of which lets a
successor on another node take a live rank directory for a dead one --
adoption then sweeps its .open files, the race B6 removes:
* GPFS / IBM Storage Scale (0x47504653): its fcntl locks are
cluster-wide, but flock is node-local (NCAR's experience on its GPFS
home filesystem: a flock on one login node does not exclude another);
* 9p (0x01021997), whose flock the client only mirrors to a server
that may not enforce it (QEMU's does not);
* AFS, OpenAFS's (0x5346414F) and kAFS's (0x6B414653);
* OrangeFS (0x20030528), which has no cross-client flock.
The rule the plan states is that a spool root is node-local, which none
of these is; allow_shared_filesystem still admits any of them. CephFS,
whose flock is coherent across clients, is left alone, as are the local
ones (ext4, xfs, btrfs, zfs, tmpfs, overlayfs).
Tests, red before: test_spool_owner_lock.cpp names and refuses each new
magic through the statfs seam (the override admitting it), and btrfs and
zfs pass; test_native_spool_ownership.py refuses each through the
binding, by name.
spool.cpp was portable POSIX at d3fe7c2. The owner lock added <sys/vfs.h>, which only Linux has (macOS declares statfs in <sys/mount.h>), read statfs's Linux-only f_type magic, and fell back to plain rename() where RENAME_NOREPLACE is missing -- which on macOS silently replaces an empty target directory. Three CPU tests compile spool.cpp directly and carry Homebrew OpenSSL paths for macOS (test_native_spool_reservations.py, test_native_pack_sink_timeout.py, test_native_spool_owner_lock_unit.py); on a Mac they would now fail to compile, and CI, Ubuntu-only, would not notice. * <sys/vfs.h> only on Linux; elsewhere <sys/mount.h> and <sys/param.h>. * The node-local check reads f_type on Linux, as before, and elsewhere f_fstypename, refusing nfs, smbfs, afpfs, webdav, lustre, afs and the FUSE mounts (macfuse, osxfuse, FreeBSD's fusefs[.<name>]) under the same override. The statfs test seam still takes a Linux magic. * The new directory's rename uses renameat2(RENAME_NOREPLACE) where it exists (Linux, unchanged), renamex_np(RENAME_EXCL) on macOS, and otherwise refuses a target that exists before a plain rename(), so an existing directory still goes to the lock-in-place path. /proc/self/fdinfo, which held_by_caller reads, is already optional: where /proc cannot be read it falls back to the owner record. Not compiled on macOS here. Each non-Linux branch was syntax-checked on Linux with __linux__ undefined against a stand-in <sys/mount.h> (and the plain-rename branch forced), and the Linux build and the spool CPU tests (66) pass unchanged.
Before the section 2.3 layout every restart reused one spool_root, and
Spool::Open counted what was already there against spool_max_bytes, so
the node's spool stayed within one budget however often a job
restarted. Now every create_record_runtime claims a fresh
r<rank>-<incarnation> directory, and its Spool counts only that
directory. So while uploads are blocked (the object store down, the
catalog up, so the successor still starts and its adoption is merely
owed) each crash-restart adds a whole spool_max_bytes: the sink fills its
directory, fails the worker ("spool stage failed"), HF's on_failure=raise
ends the job, torchelastic / Slurm --requeue / a k8s restartPolicy starts
the next incarnation, and so on until the local disk is full. The review
measured four incarnations under one base with spool_max_bytes=65536
holding 63748 bytes each, 254992 in all. The plan's per-node budget line
was not implemented alongside the layout that causes this.
SpoolConfig.charge_dead_siblings: the root is a rank directory, and what
its sibling rank directories hold (ready packs and temp files, _refs/
aside) counts against max_bytes too. A sibling another live process
holds is left out -- that is its own budget -- while one nobody holds (a
dead incarnation waiting to be adopted) or one THIS process holds (its
service adopting it) counts. The charge is refreshed wherever the
committed account is: at Open, and in the reconcile a stage runs before
it refuses, so capacity comes back as adoption drains them. The kFull
message says how much of the total is the dead directories', and the
snapshot carries sibling_bytes. The engine's sink sets it
(SinkConfig.spool_charge_dead_siblings, NativePackSink's
charge_dead_siblings, create_native_pack_sink's) for the directory it
claims; the sink-only and explicit-sink modes, which have no layout,
do not.
A deployment that runs several live ranks on a node still gives each
its own spool_max_bytes; dividing a node budget among them is the
plan's per-node line, still to do.
Tests, red before (the fields did not exist):
* test_spool_owner_lock.cpp: beside a dead sibling holding 300 bytes,
a 450-byte spool admits one 100-byte pack and refuses the next,
naming the dead bytes; once the sibling's packs are gone the stage
goes through; without the charge the same budget admits 400; a
sibling a live child process holds is not charged; one this process
holds is;
* test_native_spool_ownership.py: the real sink in a claimed directory
beside a dead one holding 1 MiB, with 1 MiB + 256 bytes of budget,
refuses its first pack with charge_dead_siblings and stages it
without;
* test_native_capture_storage_wiring.py: the engine passes
charge_dead_siblings=True to its sink with a claim, False without.
…lush()
Adoption uploaded a dead sibling's whole backlog synchronously, with no
deadline: at start(), so create_record_runtime waited for it, and in
step 3a of every flush() cycle, so flush(timeout) did too, even with
nothing of this process's own pending. With the object store refusing
connections every dead pack ran its whole retry chain (four uploader
attempts of four client attempts each) first: the review measured start()
at 16.4 s and flush(0.5) at 14.7 s for eight packs, linear in the
backlog (about 33 minutes for 1000), and far worse against a store that
accepts and never answers. Ctrl-C cannot interrupt the native call.
Before B6 start() uploaded nothing; the loop did. And every retry swept
and re-hashed the whole dead directory twice (Recover, then the
uploader's ListPending).
* start() only owes a look at the siblings -- still after the lease
and the sweep of its own directory, the plan's order -- and kicks
the loop, whose first cycle runs at once.
* The loop's cycles adopt; flush()'s never do (run_cycle(adopt)). A
cycle probes the siblings' locks, queues the dead ones, and works
through them: it takes one's lock, sweeps and validates it once
(Recover), then uploads and indexes its packs a round of
uploader.max_workers at a time (SpoolUploader::UploadEntries, the
uploader's upload of entries already listed), checking the lease and
that nothing is owed to the catalog before every round, until
adoption_slice_ns (1 s) has passed or stop() is asking. The lock, the
Spool and the remaining packs are kept between cycles, so a backlog
is hashed once, and neither a flush() nor the service's own uploads
wait behind all of it; stop() lets go of a half-adopted sibling,
which keeps what is left for the next process.
* flush() is about this process's records (plan section 2.6): a
sibling still to adopt no longer keeps it from reporting drained,
though adopted packs uploaded and not yet indexed are owed in
pending_index_ like its own. Only an adoption upload that failed
counts towards the loop's backoff. NativeCaptureStorage.flush's
TimeoutError no longer mentions adoption, and the storage part of
capture_status() is where it shows (adopted_spools, adoption_owed).
* A sibling that cannot be locked or opened is recorded, and looked
at again after the backoff.
Tests (test_native_spool_adoption_live.py): the store-down case, red
before -- with the dead lease lapsed and the store cut, start() returned
after 11.8 s; it must return within 3 s, and flush(0.5) within 2 s,
while the dead spool stays whole and adoption owed; once the store is
back the loop adopts it. The other adoption tests wait for the loop's
adoption rather than expect it done when start() returns, and the
owed-pack test's mutation (no pending-index check before a round) still
fails it. The wiring test for the old flush message goes. The storage live
suite, the chain suite and the wiring tests pass.
When a dead sibling held a pack that can never be uploaded, adoption
failed on it every pass and never finished: adoption_owed stayed set,
every cycle reported failed, the loop sat at its maximum backoff (the
review saw 1 cycle in 20 s where the poll interval gives about 400), and
each retry re-hashed the dead pack. With the adoption of the previous
commit it also kept every sibling queued after it waiting. A pack larger
than this process's uploader_max_in_flight_bytes (a crashed run with a
larger bound left it), a different object already at its key, and
staged bytes that no longer match are all like that; so is a dead
directory whose lock cannot be taken at all (EACCES on another user's
lock file).
* dmi_store::UploadFailure gains `retryable`: false for the oversized
pack, the conflict and the corrupt staged bytes -- the three exits
the uploader already treats as final -- and true otherwise
(transport, a 5xx, a failed HEAD). The store driver reports it.
* An adoption round re-queues only retryable failures, and only they
back the loop off. A non-retryable one blocks its directory: the
rest of its packs still go up, then the directory is left in place
with its lock let go, for a process that can adopt it (or a
person), recorded once in last_error and listed in the snapshot's
blocked_siblings. It is never looked at again by this service, and
is not owed. A sibling that cannot be locked or opened, and one
drained of packs that still holds other files, are blocked the same
way.
Tests, red before: test_native_spool_adoption_live.py -- two dead
siblings, one holding packs over a 4096-byte in-flight bound: the other
is adopted and removed, the first is listed in blocked_siblings with its
packs where they were and the in-flight limit in last_error, no more
uploads fail after that, the loop keeps its poll interval (at least 5
cycles a second), adoption is not owed, a flush returns at once, and the
catalog holds exactly the adoptable sibling's captures (before: the
adoption never settled). test_native_uploader.py: the oversized pack and
the conflict are reported not retryable, a TLS failure retryable.
…and them over The engine before the layout spooled into <spool_root>/v1/... and started its service there with the sweep on, so the next start uploaded what a crashed run had left. The engine now claims <spool_root>/<key>/r<rank>-<inc>/, and adoption looks only at rank directories beside it: packs a pre-layout crash left under <spool_root>/v1 are never swept, uploaded or reported, and the next run starts cleanly without a word. The same holds for what a sink-only or explicit-record_sink run leaves in spool_root itself. Adopting them in place is not open to the engine -- spool_root contains owned rank directories, which the nesting rule refuses to let anyone take -- but moving them into a rank directory nobody owns is: the next start adopts it like a dead process's. claim_spool_directory now walks spool_root by name only, skipping the rank directories (adoption's), and when it finds ready packs outside the layout logs a warning with their count, one of them, and the directory to move them into, <spool_root>/<catalog_key>/r0-00000000, paths below spool_root kept, if they are bound for this catalog and store. The v1 contract describes the upgrade. Tests: test_native_spool_ownership.py, red before (no warning) -- a pack inside the layout warns of nothing; two flat packs under spool_root/v1 (and a temp file, not counted) give one warning naming 2, their directory and the suggested rank directory. test_native_spool_adoption_live.py -- the real sink, sink-only, stages into spool_root; the claim's warning names the directory; moved there, the packs are adopted, the directory removed, and every capture reads back from the catalog.
A child forked without exec closes its copies of the owner locks, so a fork-started worker never keeps a dead parent's directory looking live. But the handler only knew locks already held: a take registered its descriptor after LockInPlace or CreateLocked had opened, locked and recorded it, and a release unregistered it before the close. A fork from another thread inside either window -- the storage service's loop takes and lets go of dead siblings' locks without the GIL, while the main thread forks DataLoader workers -- gave the child a descriptor the handler did not close. One copied before the flock shares the description the flock then locks, so the child kept the lock after its owner died. The handler now tracks descriptors, not held objects. Every descriptor on a lock file is added to the set by the open() that makes it and removed by the close() that ends it, each under the mutex the fork's prepare handler takes; the child closes every one in the set. A SpoolOwnerLock records the fork generation it was taken in, which a child bumps, so the child's copies of the objects read as not held and never close a descriptor number reused since. Tests: test_spool_owner_lock.cpp, red before -- a fork inside a take, between the lock file's open() and its flock (the lock-open test seam), left the worker holding the lock after its owner was SIGKILLed; now the lock goes with the owner and the worker lives on.
Whether a spool directory's owner lives was judged by whichever inode sat at <dir>/.owner.lock at that moment. The owner locks a file it never touches again, and nothing checked that it stayed in place: once something other than its holder removed it -- systemd-tmpfiles aging /tmp (`q /tmp 1777 root root 10d` on RHEL and Fedora) takes a stale file and leaves fresh packs, or a person -- the live directory read as dead. TryAdopt then created a new lock file, took it, and Recover swept the owner's in-flight .open files: the race B6 is there to remove. A take now also flocks the directory itself (O_RDONLY|O_DIRECTORY), after its lock file, and CreateLocked locks the staging copy before the rename, so a new directory still appears with both held. A directory cannot be unlinked while it holds anything, so a lock file replaced behind a live owner's back no longer lets anyone in: the next take meets the directory's lock and is refused (kOwned, saying the lock file was replaced and so cannot name the holder). ReadSpoolOwner probes the directory's lock when the file's is free, and SpoolOwnedByThisProcess accepts this process's lock on either inode, so held_by_caller keeps working for the owner. systemd-tmpfiles, for its part, skips a flocked directory and everything below it, which keeps it from aging a live spool at all. The v1 contract says so, and to keep spool_root out of what other cleaners age. Tests: test_spool_owner_lock.cpp, red before -- with a live owner's lock file removed, the directory reads as held, TryAdopt and a take are refused, and once the owner dies it is adopted; in the owner's own process, SpoolOwnedByThisProcess and a held_by_caller Spool still see it held. test_native_spool_ownership.py -- the same from another process through the bound spool_owner and SpoolOwnerLock.
held_by_caller -- the engine's service, its sink and every adoption -- asks the kernel whether one of this process's descriptors holds the directory's lock, through the "lock:" lines of /proc/self/fdinfo, and fell back to the owner record only when /proc/self/fd could not be opened. A procfs that lists descriptors but no locks answered "no" for the process that really held the lock: gVisor's fdinfo prints only pos, flags and mnt_id (Modal, GKE Sandbox), and WSL1's lists none either. There every held_by_caller open was refused kOwned, naming the caller's own pid, so the default persistent path could not start at all, and every adoption would have blocked its sibling. Before B6 that configuration worked. Whether this kernel lists flocks is now probed once per process, with a flock on a temporary file; where it does not (or no temporary file can be made), the record decides -- this host and this pid -- as the Python spool_owner_lock_beside already does. SetFdinfoHidesLocksForTesting makes this binary's fdinfo reads see no lock lines. Tests: test_spool_owner_lock.cpp, red before -- with the seam hiding lock lines, the process holding a lock is told so and opens its own held_by_caller Spool; beside another process's lock the open is still refused, naming that holder.
charge_dead_siblings exempted only siblings another process held, and charged every one this process held on the grounds that its service was adopting it: SpoolOwnedByThisProcess cannot tell an adoption's lock from a claim's. But two kinds of this-process-held directory are no adoption's -- a claim kept owned because its sink did not seal (by design, until exit), and the claim of an engine dropped without close() -- and the service reads them as live and never adopts them. So after a seal that timed out while the store was down, the same process's next record runtime (a notebook, a server, an adapter re-attach, a ring replacement) had its budget cut by the kept directory's bytes for the rest of the process's life, possibly to nothing, and was refused with "still to be adopted". Siblings the service blocked were charged by every sink on the node for good. The lock file's record now says why its holder took it: TryAdopt, the adopter's take, adds an "adopting" line, and an adopter leaving a directory blocked writes "blocked: <why>" into it (MarkBlocked) before it lets go; the next take rewrites the record. ReadSpoolOwner fills the record whether or not the lock is held. A sink then charges a sibling only while an adoption can drain it: one nobody holds that is not marked blocked and that this process could take and empty (another user's is not), and one this process holds as an adoption. A kept claim, another process's directory, and a blocked one are not charged, and the refusal says the charged bytes are in dead directories this process's storage service adopts. A blocked directory's lock file also says why for a person. The v1 contract says what is charged. Tests: test_spool_owner_lock.cpp, red before -- a sibling this process holds through a take is not charged and the sink gets its whole budget; taken by TryAdopt it is charged; marked blocked and let go it is not (and the record says why), taken again the mark is gone and it is charged; one this process cannot write is not charged. test_native_spool_ownership.py, red before -- a kept claim beside a new one leaves the new sink its budget. test_native_spool_adoption_live.py -- a blocked sibling's lock file carries the reason.
The C++ spool skips <root>/_refs/ in Open's accounting walk, the sibling charge, the committed reconcile and Scan (so Recover and ListPending). The Python reference DurablePackSpool does not: its constructor counts, and its recover() sweeps and quarantines, .open and .ready files there. spool.h's parity paragraph and the tests named only the lock as deliberately stricter, so a parity comparison meeting the difference once ref files live there would not know it was intended. It is, and it is not ported to the reference; spool.h and the driver tests now say so.
Two holes in the nesting rule, one in each direction. The descendant half of CheckNotNested refused any lock file below the directory, held or not. A default-mode run that leaves its rank directory under spool_root -- a SIGKILL, a close whose drain left packs, an unsealed close once the process exits, a blocked sibling -- then refused both other modes on that spool_root: sink-only (the sink takes spool_root) and the explicit-record_sink rollback (the service takes it), with "contains the owned spool directory" although nobody held it, until a default-mode start adopted the directory, and for good for a blocked one. The rollback was blocked exactly when it is most likely needed, and nothing said so. The other way, TryAdopt runs no nesting check, although the rule relied on the next take of a directory meeting the lock file below it. An adopter of a dead rank directory holding a live nested spool (a root put there, which an unheld ancestor admits) had its Recover sweep that spool's in-flight .open files, and would have uploaded its packs under the outer directory's keys. The refusal was standing in for what the walks should do. Every walk of a spool -- Open's accounting, the committed reconcile, the sibling charge and Scan (Recover and ListPending) -- now passes over a subdirectory with a lock file of its own: another spool directory, live or dead, is never this one's to count, sweep, list or upload. So a flat spool_root leaves the rank directories under it to adoption, and an adopter leaves a spool nested in the dead directory it drains alone (then keeps the directory, which still holds it, and says so). CheckNotNested refuses a descendant, like an ancestor, only while its lock is held -- the outer walk could meet a live inner directory before its lock file is there -- and names the holder. The two layouts on one spool_root take turns and never run at once. The v1 contract says what switching modes does; skipping nested spools is C++-only, not ported to the Python reference. Tests: test_spool_owner_lock.cpp, red before -- a held rank directory under a root refuses its take naming the holder; once dead, with a pack and a stage in flight, the root is taken, counts and recovers nothing of it, stages its own, and the rank directory cannot be taken while the root is held; an adopter of a dead directory holding a live nested spool recovers only its own pack, leaves the nested .open in place, and keeps the directory. test_native_spool_owner_lock.py, red before -- a dead spool inside another is left alone by the outer recover. test_native_spool_ownership.py, red before -- after a crashed default run, the explicit-sink service and a sink-only sink take spool_root and leave the dead directory as it was; a live claim still refuses both.
When a foreign lease outlives 2 x TTL, the service latches and its loop returns -- without letting go of adopting_, which holds the dead sibling's owner lock between cycles. The sibling stayed locked by this process, with nobody working on it, until the engine called stop() at close, while the engine kept capturing, possibly for hours. Meanwhile every process on the node, the one now holding the catalog included, read it as live and could only recheck it every 30 s. A latched loop now lets go of the sibling being adopted and of those queued (what stop() does, now one helper), and looks no more; the snapshot's adoption_owed says so. latch_failure wakes the loop, so this happens at once rather than after its wait, which an outage stretches to max_backoff_ns. Tests: test_native_spool_adoption_live.py, red before -- a service half-way through adopting a dead sibling (the object store cut) loses the catalog to a rival for 2 x TTL and latches; within 5 s the sibling's lock is free, with every pack still in it.
The flush half of the fix that moved adoption into the loop -- flush()'s cycles run with adopt=false -- had no test that could fail. The one assertion aimed at it, flush(0.5) returning within 2 s while the store is down, ran while the loop's first cycle held the cycle mutex through its own upload round, so the flush timed out on the lock without running a cycle at all (its own comment says as much). A flush that adopted again survived the whole adoption live suite, and would reopen the overrun for flush() and close()'s drain: a dead pack's retry chain past the deadline, minutes against a stalled store. The new test idles the loop (an hour's poll interval) once its first adoption round has failed with the object store cut, so the flush holds the cycle itself. It must return within 2 s, with upload_failures and adopted_packs unchanged and the dead packs where they were. Tests: test_native_spool_adoption_live.py -- passes; with flush() running run_cycle(true) (the review's surviving mutation) it fails, the flush uploading dead packs and timing out.
TestANewDirectoryAppearsWithItsLockHeld promised that a new directory appears with its lock held -- CreateLocked's staging copy and rename -- but checked only the end state: one entry, its lock file present. Creating the directory and then locking it in place passed every CPU and live suite, although that is exactly the window in which an adopter's scan meets a brand-new sibling unlocked, takes it for dead, and removes it from under its claimer (whose create_record_runtime then fails). The test now races a watcher against the takes: a forked child lists the parent over and over and probes every rank directory it has not yet seen held, while this process creates 100 and holds every one, so their creation is the only moment one could read as unheld. Each take is slowed at the lock-open seam, where a lock file is open and not yet locked, so a directory there to be seen before its lock is seen. Tests: test_spool_owner_lock.cpp -- passes; against a spool.cpp that creates the directory and then locks it (the review's surviving mutation) the watcher meets all 100 unheld, in every run.
The set of lock descriptors the fork handler closes in a child holds the directory's descriptor as well as the lock file's since the directory is locked too; its comment named only the lock file.
… on time) into the spool owner lock
Conflicts, and how each was resolved so both milestones' behaviour stays:
- store/spool.h: B5's Cancellation forward declaration beside B6's
OwnerLock enum. Scan's walk keeps both B6's skipped directories (_refs/,
nested spools) and B5's cut between packs.
- store/uploader.{h,cpp}: B5's UploadStaged replaces B6's UploadEntries,
which did the same thing. UploadFailure carries both `cancelled` (B5)
and `retryable` (B6); UploadOne reports both, and the aggregate
initialisers now name them.
- store/conformance_store.cpp: upload_pending failures report both
`cancelled` and `retryable`.
- sink/bindings_sink.cpp: NativePackSink takes release_flush_timeout_s and
owner_lock, allow_shared_filesystem, charge_dead_siblings.
- catalog/storage_service.{h,cpp}: run_cycle(deadline_ns, allow_reconcile,
adopt) -- the loop's cycles reconcile and adopt, flush()'s do neither.
B6's member index_or_owe takes B5's deadline and returns what a cancel
left owed. The constructor checks B6's layout and wires B5's two
clients and Cancellations. The loop lets go of an adoption once
latched, then follows B5's wait-after-a-flush rule. B5's chunked upload
path stands; adoption runs after it, and a failed adoption upload still
fails the cycle. Adoption is otherwise as B6 left it here; it is moved
onto B5's clients, chunks and cancels in the commits that follow.
- engine.py: _seal_capture_sink keeps B5's docstring and B6's result.
- storage/native_capture.py: flush()'s docstring says both what B5 bounds
and what B6's adoption leaves out.
The merge left adoption as B6 wrote it, beside B5's upload path rather than on it. Its rounds uploaded through the index-read client (s3_) with an uploader that had no Cancellation, so stop() cut an adopted transfer only by accident and the uploader went on retrying; a cut pack was booked as an upload failure, logged as "adopting dead spool ... upload failed" and pushed to the back of the queue. Its listing of a dead backlog -- Recover, which hashes every pack -- could not be cut at all, so stop() waited for the whole hash (9 s for 32 sparse 256 MiB packs). - A round is now one chunk of the service's own upload path: upload_chunk (UploadStaged, then the same booking of uploaded, failed and cancelled packs), then index_or_owe before the next chunk. It is at most uploader.max_workers and indexer.max_packs packs, so at most one chunk of a dead spool is out of it and not yet in the catalog. - The adoption uploader writes through upload_s3_ and has upload_cancel_, as the service's own does: stop() cuts its transfers, retries and backoff, and its index reads go through s3_, which read_cancel_ cuts. - A pack a cancel cut goes back to the front of the sibling's remaining packs, is counted in cancelled_uploads, not upload_failures, and makes the cycle cut short rather than failed; what a cancel left owed counts in the cycle's deferred, as the service's own does. - Spool::Recover gains B5's (cancel, cut) form, and begin_adoption lists a dead spool with upload_cancel_: a cut listing lets the directory go, whole and unlocked, and it is looked at again from the start. - A cycle adopts only once every staged pack of its own went up, nothing is owed and no cancel came. stop() still never removes a dead directory: only finish_adoption does, once none of its packs is left.
…lice flush() never adopts, and B5 has it return on time; but it takes the cycle lock, and a loop cycle adopting held that for adoption_slice_ns (a second by default) plus a round. Against a slow store a flush of a process with nothing of its own timed out behind a dead backlog: with a second a PUT and a one-minute slice, flush(4.0) returned False. flush() now counts itself in flushes_in_progress_ for its whole call, and a cycle adopting lets go of the cycle at its next step while one is -- after the round, or the sibling's listing, it is in. Each cycle still takes one step at least, so flushes that keep coming slow adoption down but never stop it.
B6 lets go of the engine's spool directory once the sink and the service are done, and judged the sink by close()'s flush before the ring stops. B5 added a second flush after it: stopping the ring drains what it still queued into the sink, and the release backstop (on_engine_release) stages that, bounded by its own timeout. A backstop that timed out left a stage on its way into a directory whose lock close() then gave up -- another process's adoption could sweep it -- and a backstop that sealed a sink whose first flush ran out of budget kept the directory owned until exit. NativePackSink::sealed_on_release() says whether the release's flush went through; a released sink admits nothing, so then no stage is left to come. It is false while attached, after a release flush that failed or timed out, and with the backstop off. close() and a record ring's replacement judge the sink by it once the ring has stopped; a sink that does not report it (not a native pack sink) is still judged by the flush before the stop.
…r of adoption The integration guide judged the sink by close()'s own flush and said a flush does not wait for adoption. Since the merge the directory stays owned only when the release backstop did not seal the sink, a flush waits for the adoption step in flight and no more, and close() cuts an adoption as it cuts the service's own uploads.
B5's live tests build their services from a raw native dict, which opens the spool with owner_lock take. Under the owner lock that service keeps the directory until it is destroyed, stop() or not, so the successor each test then starts on the same spool (_service, held_by_caller beside the harness's lock) was refused: five tests failed with SpoolOwnedError. They now open it held_by_caller under the harness's lock, as B6 did for every raw service already in the file.
The wiring tests fake the sink, so only here does a real ring release a real NativePackSink whose backstop seals it: close() without a flush then drains the tail into the catalog and, the sink sealed on release, lets go of the engine's rank directory and removes it, leaving nothing to adopt.
A child forked without exec closes its copy of the spool lock only in the fork handler, which runs once the child is first scheduled. The two fork tests that killed the owner without hearing from the worker then read the lock at once, and on a loaded host the worker had not run yet: the lock outlived its owner for tens of milliseconds and the check failed (the C++ case every time at load 60-80 on 32 cores). The worker now says it has run, past fork() and its handler, before the owner is killed, as the sibling case already did.
… of its length The switch forwarded a request whose client went away mid-body -- a flush deadline's cancel landing between libcurl's head and its body, after the switch had answered 100 Continue itself -- and the fake S3, whose signature check takes the payload hash from x-amz-content-sha256, stored the short body. A PUT held back by a delay could so land after the successor had uploaded the same pack, and overwrite it with 0 bytes: test_a_close_whose_budget_ends_mid_upload_leaves_nothing_owed then failed reading the pack's trailer. The switch now drops a request it did not receive whole, and one a delay still held back when close() ran; the fake S3, like S3, stores no body short of its Content-Length or not the one x-amz-content-sha256 hashes.
The multipart-abort flush had 1 s to list and hash its 65 MiB pack, send the preflight HEAD, hash it again and create the upload before its part was stalled; at load 70-80 on 32 cores the deadline came first, cut the upload before the part and left nothing to abort, on main as here. Its budget is 4 s now, the bound moved with it. The listing-cut uploader test ended within one 256 MiB pack's hash of its cancel, which on such a runner took most of its 2 s bound; its backlog is 128 packs of 64 MiB, so one pack is a small part of the bound and the whole listing still far over it.
…ff shows stop() cancels the upload client and the uploader alike. With only the client's cancel an attempt in flight is still aborted and every later one refused at once, but the uploader's own backoff sleeps on: at the default four attempts that is about 1.75 s, which fit under both tests' bounds, so neither the adoption's uploader nor the service's losing its Cancellation turned anything red. At eight attempts it is about 26 s; each mutation now fails its test (27 s, and past the 15 s join).
…e pack An adoption listed a dead spool in one step: Spool::Recover hashed every byte of every ready pack under the cycle lock, and only stop() cut it. A flush, which yields to the adoption only between steps, waited out the whole listing -- as long as the backlog is big, up to the dead spool's 1 TiB budget -- and timed out with its own pack still staged; close() then reported capture storage undrained and left the tail pack for the next process. With 512 sparse packs of 64 MiB, a flush(8.0) of one pack staged during the listing timed out; it now drains in about a second. Spool::BeginRecovery sweeps a dead spool's .open files and lists its ready packs without hashing them, and each ContinueRecovery validates the next, quarantining one that fails; the last rebuilds the account as Recover() does. Adoption validates one pack per step, keeps the listing across cycles, and checks the slice between packs as between rounds, so a flush waits for at most one pack's hash (or one round), the service's own uploads for at most a slice, and stop() still cuts between packs. The cancellable Recover overload, whose only caller this was, is gone. The flush docstring, close_flush_timeout_s and the integration guide say what a flush waits for of an adoption.
… go first Two scheduling rules of adoption had no test: every adopting cycle takes one step at least, however many flushes wait for it (f8ae2d9), and a cycle adopts only once every pack of its own went up (10b4fe7). Dropping either left the adoption suite green. The first is now pinned by two threads flushing back to back while a sibling dies; each cycle lists a 32 MiB own pack that fails validation and cannot be quarantined, so a flush is always waiting by the time the loop's cycle comes to adopt (Python flushes alone leave gaps a cycle slips through). The second by an own pack the store always refuses: the dead sibling is not even taken while it is there, and is adopted once it is gone. Each rule's mutant fails its test.
The uploader's note still had adoption listing a dead spool through Recover, and the flush docstring named a knob the Python config does not have: a round is at most four packs, one per upload worker.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Milestone B6 of the native capture production plan. It does three things:
What changes
One owner per spool directory, via
native/csrc/store/spool.{h,cpp}.SpoolConfig.owner_lockis eithertake(the default) orheld_by_caller.takeflocks<dir>/.owner.lockinOpen, beforeRecover, and records the holder's host and pid. A second process on the same directory is refused fast, and the error names the holder. That makes the old Recover race ("cannot link ready file") impossible.held_by_calleris accepted only when this process really holds the lock./proc/self/fdinfois checked, falling back to the owner record.statfsf_type, unless an explicit override (spool_allow_shared_filesystem) is set.Scan,Recoverand the accounting walk skip_refs/and nested spools.The engine holds the lock. A bound
SpoolOwnerLockin_dmi_native_storeis taken before the storage service starts. The in-process sink and service both open withheld_by_caller; a per-Spool lock would have made them refuse each other, and a test pins that.The lock is released only when nothing can still write the spool. That's after the sink and service close, and after B5's
on_engine_releasebackstop has sealed. The newNativePackSink.sealed_on_releasereports it. If a stage may still be coming, the lock is kept until the process exits.Adoption of dead siblings (
storage_service.cpp)..open, uploads.ready, indexes them, then removes the directory. The catalog directory is keyed by the servers its packs go to, not only by names.Cancellations, sostop()cuts it promptly. A pack a cancel cut stays with its sibling for the successor, and a directory is removed only once none of its packs is left.stop()waits behind a whole backlog being hashed..open/.readynaming is unchanged.Deliberately stricter than the Python reference (
spool.py): the owner lock, the shared-filesystem refusal, the nesting refusal and adoption are not ported to it. The parity suites pass unchanged.Integration with B5 (#162)
main gained B5 after this branch was written, so main is merged in as one merge commit (ac868a4). All 9 conflicted files keep both sides;
git show --remerge-diff ac868a4shows the resolutions. Examples:run_cycle(deadline, allow_reconcile, adopt);UploadFailurecarries bothcancelledandretryable;release_flush_timeout_splus the owner-lock and shared-filesystem options.Integration commits after it:
Evidence
CPU tier:
pytest -m cpu, 2660 passed. The 1 skip is a case that needs CUDA.Live suites:
test_native_spool_adoption_live.pyThe new adoption tests cover:
stop()cutting an adoption stalled on its uploads, and the successor finishing it;stop()cutting a dead-backlog listing (32 × 256 MiB: 9.2 s before, under 2.5 s after);GPU, on one RTX 4090 through an idle-GPU gate:
test_native_capture_storage_gpu_e2e.pypassed. That coversclose()withoutflush_and_waitleaving every capture queryable, andclose()removing the engine's rank directory.Mutations, each turning tests red:
takemode;AdoptSiblingsas a no-op;Flushbackstop;read_cancel_fromstop();stop()ignoring cancellation during adoption.Under heavy machine load (load average 55–78), some older timing tests failed: three in
test_native_lease_request_boundand one S3 multipart-cancel test. They fail the same way on an origin/main build under that load, and they passed on rerun.Review
tests/test_native_catalog_lease_live.py's_catalog_drop_onlynever drops its tables. It's already on main, and the file is the lease SQL byte-identity gate, so it needs its own change.Known gap:
release_flush_timeout_sexists only as aNativePackSinkbinding keyword, not as aNativeSinkConfigfield, so the engine always uses the 30 s default.