368 lines
18 KiB
Bash
Executable File
368 lines
18 KiB
Bash
Executable File
#!/usr/bin/env bash
|
|
# _wake-common.sh — shared state resolution + atomic-write primitive for the
|
|
# wake/heartbeat durable queue (W2 of the wake canon, EPIC #892).
|
|
#
|
|
# CONTRACT ANCHORS (docs/scratchpads/heartbeat-planning/CONVERGED-DESIGN.md):
|
|
# §1.2 three-cursor durable queue; ALL state XDG, atomic write-tmp+rename.
|
|
# §2.3 durability is NEVER bypassed.
|
|
#
|
|
# This file is sourced by store.sh and ack.sh. It defines NO top-level actions;
|
|
# sourcing it is side-effect-free except for setting readonly path vars.
|
|
#
|
|
# Operator-agnostic (framework firewall): state location is derived purely from
|
|
# XDG / env. No operator paths, names, or secrets appear here.
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# State layout (XDG). §1.2: "all state XDG".
|
|
# base = ${WAKE_STATE_HOME:-${XDG_STATE_HOME:-$HOME/.local/state}/mosaic/wake}
|
|
# agent = ${WAKE_AGENT:-default} (per-agent queue namespace)
|
|
# STATE_DIR = <base>/<agent>
|
|
# ---------------------------------------------------------------------------
|
|
wake_state_dir() {
|
|
local base agent
|
|
base="${WAKE_STATE_HOME:-${XDG_STATE_HOME:-$HOME/.local/state}/mosaic/wake}"
|
|
agent="${WAKE_AGENT:-default}"
|
|
printf '%s/%s' "$base" "$agent"
|
|
}
|
|
|
|
# File names within STATE_DIR.
|
|
# observed_seq — highest observed_seq ever recorded (max monotonic int).
|
|
# consumed_seq — consumer cursor: top of the contiguous consumed prefix.
|
|
# observed.set — observed seqs in the live window (> consumed_seq), one int
|
|
# per line; the gap-detector's source of truth. Immune to
|
|
# coalescing (coalescing removes a pending ENTRY, never the
|
|
# fact that its seq was observed).
|
|
# pending.jsonl — the durable pending-inbox: one entry JSON object per line,
|
|
# {observed_seq, locators, class, emit_ts, hmac}.
|
|
# ack-ledger.jsonl — local-write-only ack ledger (RECEIVED / CONSUMED).
|
|
# ack-sync.state — background-sync bookkeeping (last shipped / outage flag).
|
|
# consumed-hashes.jsonl — #932 store-owned last-consumed record: one object per
|
|
# (kind,id), {kind,id,observed_hash,observed_seq}, written at
|
|
# consume-truncation. The reconciler's THIRD accounting source
|
|
# (a consumed state matches neither the truncated inbox nor the
|
|
# reconciler's own seen-ledger). ADDITIVE / lazily created;
|
|
# older code ignores it (on-disk read-compat preserved).
|
|
|
|
# _wake_tmp_glob DIR — the glob used for atomic-write temp files, so readers can
|
|
# ignore in-flight/crashed writes. A crash leaves one of these; it is NEVER the
|
|
# live file (only rename promotes content), so it can never corrupt a read.
|
|
_wake_tmp_prefix='.wake.tmp.'
|
|
|
|
# _atomic_write TARGET (content on stdin)
|
|
# §1.2 / task: EVERY state mutation is atomic write-tmp+rename. Write a temp
|
|
# file in the SAME directory (so mv is a same-filesystem atomic rename), fsync
|
|
# is best-effort, then rename over the target. A crash before the rename leaves
|
|
# the old target fully intact and a stale .wake.tmp.* that readers ignore.
|
|
_atomic_write() {
|
|
local target="$1" dir tmp
|
|
# --- TEST-ONLY FAULT SEAM (issue #934) — PROD-INERT. ------------------------
|
|
# Forces the ALREADY-EXISTING atomic-write failure PATH (the fail-loud +
|
|
# rollback handling #908/#917 built) to be taken for ONE named write target, so
|
|
# the seq-integrity failure assertions (T9 arrow-1 no-burn, T11 cursor-write
|
|
# gate) RUN UNPRIVILEGED in the real non-privileged CI runner instead of being
|
|
# skipped behind an unshare+bind-mount injection. It is honored ONLY when the
|
|
# test-only env var WAKE_TEST_FAULT is explicitly set to name a write point; it
|
|
# writes nothing, adds no new behavior, and changes no on-disk format. No
|
|
# production input (CLI args, locators JSON, on-disk state, watch-list) can set
|
|
# a process env var, so with WAKE_TEST_FAULT unset this is a no-op and the write
|
|
# proceeds exactly as before. Map: pending->pending.jsonl, cursor->observed_seq.
|
|
if [ -n "${WAKE_TEST_FAULT:-}" ]; then
|
|
case "${WAKE_TEST_FAULT}:$(basename -- "$target")" in
|
|
pending:pending.jsonl | cursor:observed_seq)
|
|
cat >/dev/null 2>&1 || true # drain the producer, then report the commit as failed
|
|
return 1
|
|
;;
|
|
esac
|
|
fi
|
|
dir="$(dirname "$target")"
|
|
[ -d "$dir" ] || mkdir -p "$dir"
|
|
tmp="$(mktemp "$dir/${_wake_tmp_prefix}XXXXXX")" || return 1
|
|
if ! cat >"$tmp"; then
|
|
rm -f "$tmp"
|
|
return 1
|
|
fi
|
|
# Best-effort durability of the temp file before the rename. Not fatal if the
|
|
# platform lacks it — the rename atomicity is the load-bearing guarantee.
|
|
sync "$tmp" 2>/dev/null || true
|
|
if ! mv -f "$tmp" "$target"; then
|
|
rm -f "$tmp"
|
|
return 1
|
|
fi
|
|
}
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Enqueue serialization (single store-side observed_seq allocator, #908).
|
|
#
|
|
# store.sh is the SOLE allocator of observed_seq: it reads the observed_seq
|
|
# cursor, computes next = observed_seq + 1, and writes the pending record +
|
|
# observed.set + the cursor as ONE transaction. That read-modify-write MUST be
|
|
# serialized so two concurrent enqueues cannot read the same cursor and collide
|
|
# on a seq. These helpers take an exclusive lock for the duration of the
|
|
# transaction. flock (Linux/CI) is the load-bearing mechanism; where flock is
|
|
# absent the open still succeeds so a single-threaded enqueue is unaffected (the
|
|
# concurrency guarantee then degrades — the concurrency test SKIPs without flock,
|
|
# exactly as the detector's single-instance test already does).
|
|
#
|
|
# The lock uses a FIXED fd (8) within one store.sh process. Each store.sh
|
|
# invocation is its own process (the detector/reconciler call it as a
|
|
# subprocess), so fd 8 is always free here and never clashes with the detector
|
|
# run-loop lock (fd 9, a DIFFERENT process).
|
|
# ---------------------------------------------------------------------------
|
|
_wake_lock_acquire() {
|
|
# _wake_lock_acquire LOCKFILE — open fd 8 on LOCKFILE and take an exclusive
|
|
# (blocking) lock. Returns non-zero if the lock cannot be taken.
|
|
local lf="$1" dir
|
|
dir="$(dirname "$lf")"
|
|
[ -d "$dir" ] || mkdir -p "$dir" 2>/dev/null || true
|
|
# Open the lock fd. If the file does not yet exist and the dir is writable it
|
|
# is created; if the dir is read-only but the file exists, opening it O_WRONLY
|
|
# still succeeds (write perm on the file, not the dir).
|
|
# NB: a command-less `exec` redirection persists for the WHOLE shell, so we must
|
|
# NOT append `2>/dev/null` here (it would permanently silence the caller's
|
|
# stderr and swallow every later fail-loud diagnostic). Redirect only fd 8.
|
|
exec 8>"$lf" || return 1
|
|
if command -v flock >/dev/null 2>&1; then
|
|
flock 8 || return 1
|
|
fi
|
|
return 0
|
|
}
|
|
|
|
_wake_lock_release() {
|
|
# _wake_lock_release — drop the enqueue lock (closing fd 8 releases the flock).
|
|
# The `2>/dev/null` is SCOPED to the brace group (suppressing a "bad fd" close
|
|
# error) — it must NOT sit on a bare `exec`, where the redirection would persist
|
|
# for the whole shell and silence every later fail-loud diagnostic.
|
|
{ exec 8>&-; } 2>/dev/null || true
|
|
}
|
|
|
|
# _wake_clean_stale_tmp DIR — reap ORPHANED atomic-write temp files left by a
|
|
# crash mid-write (a crash before the rename leaves a .wake.tmp.* that no reader
|
|
# ever promotes). These are never the live store.
|
|
#
|
|
# AGE-SCOPED (#927): only tmp files whose mtime is older than
|
|
# ${WAKE_TMP_STALE_MIN:-5} minutes are removed. A LIVE in-flight atomic write's
|
|
# tmp is at most milliseconds old (mktemp -> cat -> sync -> rename all complete
|
|
# well under a second), so it can NEVER match this age filter. That is what makes
|
|
# the cleanup safe even if it ever overlaps a concurrent enqueue's atomic write:
|
|
# it deletes only DEMONSTRABLY-orphaned tmps, never another process's live write.
|
|
#
|
|
# This is why #927 is fixed: an UNCONDITIONAL delete of every .wake.tmp.* (the
|
|
# old behaviour) clobbered a concurrent enqueue's in-flight tmp -> spurious
|
|
# "durable pending write FAILED". It is ALSO no longer invoked from the per-
|
|
# enqueue hot path (see _wake_init_dir); it runs only at maintenance / daemon-
|
|
# start (store.sh init) and the detector poll tick, where accumulation is bounded
|
|
# once per pass rather than raced on every enqueue.
|
|
_wake_clean_stale_tmp() {
|
|
local dir="$1" min="${WAKE_TMP_STALE_MIN:-5}"
|
|
[ -d "$dir" ] || return 0
|
|
case "$min" in '' | *[!0-9]*) min=5 ;; esac
|
|
find "$dir" -maxdepth 1 -name "${_wake_tmp_prefix}*" -type f -mmin "+$min" -delete 2>/dev/null || true
|
|
}
|
|
|
|
# _wake_read_int FILE DEFAULT — read a single integer from FILE, or DEFAULT.
|
|
_wake_read_int() {
|
|
local file="$1" def="$2" val
|
|
if [ -f "$file" ]; then
|
|
val="$(tr -d '[:space:]' <"$file")"
|
|
case "$val" in
|
|
'' | *[!0-9]*) printf '%s' "$def" ;;
|
|
*) printf '%s' "$val" ;;
|
|
esac
|
|
else
|
|
printf '%s' "$def"
|
|
fi
|
|
}
|
|
|
|
# _wake_init_dir STATE_DIR — ensure the state layout exists; idempotent.
|
|
#
|
|
# #927: this runs on the HOT enqueue path (cmd_enqueue calls it BEFORE taking the
|
|
# enqueue lock) as well as on consume/cursors/ack. It must therefore NEVER touch
|
|
# another process's tmp files: the old _wake_clean_stale_tmp call here deleted a
|
|
# concurrent enqueue's LIVE in-flight tmp mid-write -> spurious durable-write
|
|
# abort. Stale-tmp reaping is now an explicit maintenance action (store.sh init /
|
|
# detector poll tick), NOT a side effect of ensuring the layout. Keep this
|
|
# function limited to creating the dir + seeding the cursor files.
|
|
_wake_init_dir() {
|
|
local dir="$1"
|
|
mkdir -p "$dir" || return 1
|
|
[ -f "$dir/observed_seq" ] || printf '0' | _atomic_write "$dir/observed_seq"
|
|
[ -f "$dir/consumed_seq" ] || printf '0' | _atomic_write "$dir/consumed_seq"
|
|
[ -f "$dir/observed.set" ] || printf '' | _atomic_write "$dir/observed.set"
|
|
[ -f "$dir/pending.jsonl" ] || printf '' | _atomic_write "$dir/pending.jsonl"
|
|
[ -f "$dir/ack-ledger.jsonl" ] || printf '' | _atomic_write "$dir/ack-ledger.jsonl"
|
|
}
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# #973 — three-valued grep assertion helpers for the wake test suites.
|
|
#
|
|
# grep's exit contract is three-valued: 0 = match, 1 = no match, >1 = ERROR
|
|
# (bad file, bad pattern, resource failure). Every wake-suite assertion used
|
|
# to read all non-zero as "absent", so a grep that COULD NOT LOOK wore the
|
|
# colour of a verdict: OR-polarity sites (`|| fail`) went falsely red,
|
|
# AND-polarity sites (`&& fail` — including the credential canaries) went
|
|
# falsely green. The repair is to refuse to answer: rc 0 -> match, rc 1 -> no
|
|
# match, anything else -> loud abort naming the call site, the raw exit code,
|
|
# and the arguments. An error NEVER becomes a verdict.
|
|
#
|
|
# Production tools (store.sh, ack.sh) source this file but call none of the
|
|
# helpers below; they are inert outside the suites.
|
|
#
|
|
# Suite integration contract:
|
|
# - Call `wake_assert_init` ONCE at suite top level, right after sourcing.
|
|
# It dups the suite's real stderr to a saved fd BEFORE any call-site
|
|
# redirect exists, so an abort stays loud even at sites that append
|
|
# `2>/dev/null` (the preimage credential canaries pre-swallow stderr —
|
|
# exactly where a silent abort would recreate the defect being fixed).
|
|
# - Assertion sites live inside `( ... ) && ok` subshell blocks, pipelines,
|
|
# and `$(...)` substitutions, where a plain `exit` dies one layer deep and
|
|
# the suite would carry on to emit a verdict. The abort therefore signals
|
|
# the suite's MAIN shell ($$ is the main PID in every subshell) and then
|
|
# exits the current context: the suite dies by signal, non-zero, with NO
|
|
# verdict line emitted.
|
|
#
|
|
# Validation instrumentation (#973 evidence, not part of the assertion fix):
|
|
# - WAKE_ASSERT_LEDGER=<file>: every helper call appends
|
|
# "<helper> <caller-file>:<caller-line>" to <file>. That is the ONLY
|
|
# divergence from production behaviour — the suite otherwise runs its
|
|
# normal arms, so a validate run exercises exactly the shipped paths.
|
|
# - WAKE_ASSERT_FORCE_GREP_ERROR_AT=<caller-file>:<caller-line>: at exactly
|
|
# that call site, the invocation is routed through a REAL grep driven onto
|
|
# its real error path (unknown option -> rc 2) — a genuinely executed
|
|
# failing process, not a stubbed return — to prove per-site that the abort
|
|
# fires. Unset in production; matching no site is a no-op.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# wake_assert_init — dup the suite's real stderr once, for abort loudness.
|
|
# MUST be called at suite TOP LEVEL, immediately after sourcing and before any
|
|
# test block: a lazy (first-call) dup could capture an already-redirected
|
|
# stderr if the first executed helper call sat under a call-site 2>/dev/null,
|
|
# silencing every abort thereafter. The fd is allocated dynamically (>= 10),
|
|
# so it cannot collide with the wake lock fds (8) or the detector run-loop
|
|
# lock (9).
|
|
#
|
|
# Init also PINS the BASH_LINENO convention the site coordinates depend on:
|
|
# a helper call written across a backslash continuation must report at its
|
|
# FIRST physical line (the denominator artifact's convention). That was
|
|
# measured on a developer bash (5.3.x); CI runs whatever bash its base image
|
|
# baked in, and that version floats silently between image rebuilds. A bash
|
|
# that disagrees would shift every continuation-site coordinate by one line
|
|
# UNDER the validation instead of in front of it — so the convention is
|
|
# asserted at runtime, in the same bash binary that runs the suite, and a
|
|
# disagreeing bash aborts the suite loudly instead of skewing coordinates.
|
|
_wake_assert_lineno_pin() {
|
|
local _wa_pin_tmp _wa_pin_got
|
|
_wa_pin_tmp="$(mktemp)" || {
|
|
_wake_assert_err_note "WAKE-ASSERT INIT ABORT: mktemp failed; cannot pin the BASH_LINENO convention — a pin that silently does not run is not a pin (#973)"
|
|
exit 97
|
|
}
|
|
cat >"$_wa_pin_tmp" <<'WAKE_ASSERT_PIN'
|
|
_wap() { printf '%s\n' "${BASH_LINENO[0]}"; }
|
|
(
|
|
_wap simple
|
|
_wap \
|
|
continuation
|
|
)
|
|
WAKE_ASSERT_PIN
|
|
# WAKE_ASSERT_PIN_BASH: test-only interpreter override so the pin's abort
|
|
# arm can be PROVEN to fire (microtest C10) — bash resets $BASH at startup,
|
|
# so the real probe interpreter cannot be spoofed from the environment.
|
|
_wa_pin_got="$("${WAKE_ASSERT_PIN_BASH:-${BASH:-bash}}" "$_wa_pin_tmp" 2>/dev/null)"
|
|
rm -f "$_wa_pin_tmp"
|
|
if [ "$_wa_pin_got" != "$(printf '3\n4')" ]; then
|
|
_wake_assert_err_note "WAKE-ASSERT INIT ABORT: BASH_LINENO convention violated on bash ${BASH_VERSION}: probe reported [${_wa_pin_got:-<no output>}], expected [3 4] (simple call at own line, continuation call at FIRST physical line) — site coordinates are untrustworthy on this bash (#973)"
|
|
exit 97
|
|
fi
|
|
}
|
|
|
|
wake_assert_init() {
|
|
if [ -z "${_wake_assert_err_fd:-}" ]; then
|
|
exec {_wake_assert_err_fd}>&2
|
|
_wake_assert_lineno_pin
|
|
fi
|
|
}
|
|
|
|
# _wake_assert_err_note MSG — write MSG to the saved real-stderr fd, falling
|
|
# back to the current stderr if init was never called.
|
|
_wake_assert_err_note() {
|
|
if [ -n "${_wake_assert_err_fd:-}" ]; then
|
|
printf '%s\n' "$1" >&"$_wake_assert_err_fd" 2>/dev/null ||
|
|
printf '%s\n' "$1" >&2
|
|
else
|
|
printf '%s\n' "$1" >&2
|
|
fi
|
|
}
|
|
|
|
# _wake_assert_abort HELPER SITE RC ARGS... — refuse to answer, loudly.
|
|
# Writes the named reason to the saved real-stderr fd (falling back to the
|
|
# current stderr), signals the suite's main shell, and exits this context.
|
|
_wake_assert_abort() {
|
|
local _wa_helper="$1" _wa_where="$2" _wa_code="$3"
|
|
shift 3
|
|
_wake_assert_err_note "WAKE-ASSERT ABORT: ${_wa_helper} at ${_wa_where}: grep exit ${_wa_code} is an error, not a verdict (args: $*) — refusing to answer (#973)"
|
|
if [ -n "${BASHPID:-}" ] && [ "$BASHPID" != "$$" ]; then
|
|
kill -TERM "$$" 2>/dev/null || true
|
|
fi
|
|
exit 97
|
|
}
|
|
|
|
# _wake_assert_armed SITE — true iff the forced-error arm targets SITE; on a
|
|
# match it emits a positive confirmation FIRST, so "site did not abort" can
|
|
# never conflate SITE NOT CONVERTED with ARM NEVER REACHED IT: an armed run
|
|
# with no ARMED line means the arm matched nothing (typo/renumber/drift), and
|
|
# an ARMED line with no abort means the site's error path is broken. The two
|
|
# defects are separable on stderr alone.
|
|
_wake_assert_armed() {
|
|
[ "${WAKE_ASSERT_FORCE_GREP_ERROR_AT:-}" = "$1" ] || return 1
|
|
_wake_assert_err_note "WAKE-ASSERT ARMED: forcing real grep error at $1 (#973)"
|
|
return 0
|
|
}
|
|
|
|
# has_match GREP_ARGS... — three-valued grep verdict.
|
|
# Drop-in for verdict-bearing `grep` calls (flags, files, stdin all pass
|
|
# through; stdout is not captured, so extract-form call sites may use it
|
|
# inside a substitution). Returns 0 on match, 1 on no-match; any other grep
|
|
# exit aborts the suite via _wake_assert_abort.
|
|
has_match() {
|
|
local _wa_site="${BASH_SOURCE[1]##*/}:${BASH_LINENO[0]}" _wa_rc=0
|
|
if [ -n "${WAKE_ASSERT_LEDGER:-}" ]; then
|
|
printf 'has_match %s\n' "$_wa_site" >>"$WAKE_ASSERT_LEDGER"
|
|
fi
|
|
if _wake_assert_armed "$_wa_site"; then
|
|
command grep --wake-assert-forced-error -- /dev/null
|
|
_wa_rc=$?
|
|
else
|
|
command grep "$@"
|
|
_wa_rc=$?
|
|
fi
|
|
case "$_wa_rc" in
|
|
0) return 0 ;;
|
|
1) return 1 ;;
|
|
*) _wake_assert_abort has_match "$_wa_site" "$_wa_rc" "$@" ;;
|
|
esac
|
|
}
|
|
|
|
# count_lines GREP_ARGS... — `grep -c` with the same three-way discipline.
|
|
# Call sites drop their `-c` (the helper supplies it) and keep every other
|
|
# argument. Prints the count on rc 0 AND rc 1 (rc 1 is grep's "count is 0" —
|
|
# a valid measurement, not an error); any other exit aborts. The abort still
|
|
# kills the suite from inside a `$(...)` capture: the substitution subshell
|
|
# cannot exit the suite, but the signal to the main shell can — a count from
|
|
# a failed measurement is never printed.
|
|
count_lines() {
|
|
local _wa_site="${BASH_SOURCE[1]##*/}:${BASH_LINENO[0]}" _wa_rc=0 _wa_out=""
|
|
if [ -n "${WAKE_ASSERT_LEDGER:-}" ]; then
|
|
printf 'count_lines %s\n' "$_wa_site" >>"$WAKE_ASSERT_LEDGER"
|
|
fi
|
|
if _wake_assert_armed "$_wa_site"; then
|
|
_wa_out="$(command grep --wake-assert-forced-error -c -- /dev/null)"
|
|
_wa_rc=$?
|
|
else
|
|
_wa_out="$(command grep -c "$@")"
|
|
_wa_rc=$?
|
|
fi
|
|
case "$_wa_rc" in
|
|
0 | 1) printf '%s\n' "$_wa_out" ;;
|
|
*) _wake_assert_abort count_lines "$_wa_site" "$_wa_rc" "$@" ;;
|
|
esac
|
|
}
|