feat(targets): add an opt-in NATS JetStream publish path for the notify and audit targets
The NATS notify and audit targets publish through NATS Core, which returns
before the server has durably accepted the message. A broker restart or a
connection drop between the publish and the flush loses the event, even though
the send queue has already cleared it, and no acknowledgement gates that clear.
An opt-in JetStream publish path clears a queued event only after the server
returns a durable PublishAck, so delivery is at-least-once across a broker
restart or a reconnect. It applies to both the notify and audit NATS targets, is
off by default, and is byte-identical to the NATS Core path when disabled.
The path includes durable store-and-forward, a stable dedup id sent as the
Nats-Msg-Id header so a replayed event is collapsed by the stream duplicate
window, pre-flight stream validation, and a bounded failed-events store for
terminally-failed and retry-exhausted events. Three configuration keys per
target select it: JETSTREAM_ENABLE, JETSTREAM_STREAM_NAME, and
JETSTREAM_ACK_TIMEOUT_SECS, under the RUSTFS_NOTIFY_NATS_ and RUSTFS_AUDIT_NATS_
prefixes.
The on-disk batch filename separator changes from colon to underscore so
batch names are valid on Windows filesystems, with transparent read-back
of files written under the previous separator. The migration affects the
shared queue store for every target type and lands with this feature
because the store gains its first Windows-exercised paths here.
Co-authored-by: houseme <housemecn@gmail.com>
* fix(targets): make persistent queue store crash-safe and replay lifecycle correct
Harden the target notification persistent queue (store.rs) and the replay
worker lifecycle (runtime) against data loss, silent truncation, ordering
drift, orphaned tasks, and a few low-risk robustness gaps.
store.rs
- Atomic, durable writes: write to a per-key temp file, fsync (sync_all),
then rename into place; best-effort parent-dir fsync. A crash mid-write
can no longer lose an acknowledged event or leave a half-written payload
that reads as a valid entry.
- open() now removes leftover .tmp residue and zero-byte files, and only
indexes files matching the queue extension, so ghosts/foreign files are
never replayed.
- FIFO ordering is derived from time-ordered UUIDv7 entry names instead of
coarse, clock-dependent file mtimes, so replay order is stable and
identical after a restart.
- Clamp HashMap/Vec pre-allocation derived from untrusted inputs
(entry_limit, batch item_count) to avoid capacity-overflow panics / giant
allocations.
target/mod.rs
- QueuedPayload::decode validates body length against the recorded
payload_len, rejecting torn/truncated writes instead of delivering a
silently truncated body.
- send_from_store purges a NotFound/empty entry (index + file) instead of
skipping it, so it cannot occupy a queue slot and be replayed forever.
- sanitize_queue_dir_component appends a stable hash suffix when the id was
lossy, so distinct target ids can no longer collapse onto the same queue
directory; path-safe ids are unchanged (no migration).
runtime
- Replay backoff, idle waits, and inter-scan pauses are now cancel-aware, so
reload/shutdown is not blocked for the full retry delay.
- ReplayWorkerManager::stop_all signals cancellation and then awaits each
worker's exit (bounded, with abort fallback), preventing orphaned tasks
and overlapping drain of the same store.
- Fix the always-true replay flush condition so batching is real
(size/timeout based, one semaphore permit per batch) rather than one
permit per entry; dedup keys already pending in the batch.
- clear_and_close aggregates and reports per-target close failures instead
of swallowing them; explicit shutdown surfaces them.
Relates to rustfs/backlog#966
Relates to rustfs/backlog#967
Relates to rustfs/backlog#975
Relates to rustfs/backlog#970
Relates to rustfs/backlog#983
Co-Authored-By: heihutu <heihutu@gmail.com>
* style(targets): apply rustfmt to replay batch dedup guard
Fixes the Quick Checks rustfmt failure on the multi-line `.iter().any(...)`
closure in the replay batch dedup guard.
Co-Authored-By: heihutu <heihutu@gmail.com>
---------
Co-authored-by: heihutu <heihutu@gmail.com>
Addresses security/correctness audit findings in the target-plugin control
plane, TLS reload coordinator, and runtime extension registries.
control_plane (backlog#977):
- Install now preserves the currently installed revision as previous_revision
so Rollback actually restores the prior version instead of a no-op.
- Split the gate: circuit-breaker and runtime-activation checks apply only to
Install/Enable; Disable and Rollback stay available as break-glass
remediation while the breaker is open.
- Enforce the sidecar runtime protocol version for every external transport
(previously skippable by declaring a non-gRPC transport) and validate the
plugin api_compatibility_version at planning time.
- Download/signature/provenance host allowlisting now matches the full
host authority, so an allowlisted host never implicitly authorizes a
different host:port.
TLS reload coordinator/validate (backlog#981, #970-coordinator):
- Compute the fingerprint before building material in register() to remove the
TOCTOU that could permanently pin an old certificate; rotation self-heals.
- Always start a detection loop when reload is enabled; Watch mode no longer
returns success without any loop.
- validate_cert_key_pairing now verifies the private key matches the
certificate's public key instead of only parsing both files.
- Normalize a zero poll interval to a positive minimum to avoid a panic that
silently killed the poll loop.
- Serialize reload cycles with a per-target mutex; a first-step fingerprint
read failure now records last_error and a failure metric.
- Duplicate label registration stops and joins the previous loop before
publishing the replacement (stop-before-start), preventing orphaned loops.
runtime extension points (backlog#983 runtime subset):
- ops_profiler/ops_diagnostics authorize before probing the registry so an
unauthorized caller cannot learn whether a backend/surface exists.
- s3_hooks dispatch_post_auth actually traverses the registered hooks for the
point instead of unconditionally returning Continue.
- sidecar send_with_timeout redacts errors and uses the configurable failure
threshold via policy, and a successful send resets the breaker accounting.
Relates to rustfs/backlog#977
Relates to rustfs/backlog#981
Relates to rustfs/backlog#970
Relates to rustfs/backlog#983
Co-authored-by: heihutu <heihutu@gmail.com>