Skip to content
50 changes: 22 additions & 28 deletions bin/fm-classify-lib.sh
Original file line number Diff line number Diff line change
Expand Up @@ -804,12 +804,11 @@ EOF
# because no code path in this repo ever does that to a status file.
#
# The other real failure mode is OUR OWN read failing (a stat/wc/tail I/O
# error), not a malformed writer: every such read here is checked, and on
# failure this reports the already-trusted persisted set unchanged rather than
# risking a silent invalidation that would wipe it - never a bare "empty" as if
# nothing were open. Such a fallback call succeeds only when the set it reports
# is non-empty, so a caller can never read one as a computed "nothing open";
# the fallback block inside status_open_decisions_incremental owns why.
# error), not a malformed writer: every such read here is checked, leaves the
# already-trusted persisted set untouched on disk, and returns failure without
# printing that set. Even a non-empty persisted set is stale when unread bytes
# could have resolved one decision or opened another, so a caller must never
# present it as this call's authoritative result.
#
# Not a pure status-file read: this writes/rewrites the sibling cursor file as a
# side effect (state/.<task>.open-decisions-cursor), the library's second
Expand Down Expand Up @@ -915,7 +914,7 @@ _fm_status_read_span() { # <status-file> <start-offset> <byte-length>
}

status_open_decisions_incremental() { # <status-file> [<captured-end-offset>]
local f=$1 captured_end=${2:-} cf offset ident open='' trusted_open='' cursor_data first rest offset_line ident_line
local f=$1 captured_end=${2:-} cf offset ident open='' cursor_data first rest offset_line ident_line
local version='' size actual_size cur_ident resolve held chunk_file chunk_size line cursor_dirty=0
local target_cursor
[ -f "$f" ] && [ -r "$f" ] && [ ! -L "$f" ] || return 0
Expand Down Expand Up @@ -950,7 +949,6 @@ status_open_decisions_incremental() { # <status-file> [<captured-end-offset>]
case "$rest" in
*$'\n'*) open=${rest#*$'\n'} ;;
esac
if [ -n "$version" ] && [ -n "$ident" ]; then trusted_open=$open; fi
;;
*) offset=0; version='' ;;
esac
Expand All @@ -963,26 +961,20 @@ status_open_decisions_incremental() { # <status-file> [<captured-end-offset>]
esac
fi

# A stat/size-read failure is a genuine I/O error, not "the file is empty" -
# report the already-trusted persisted set unchanged rather than risking a
# silent invalidation that would wipe it. That fallback is only honest while
# the set it reports has something IN it: an EMPTY set returned with rc 0
# asserts "computed, nothing open" about bytes this call never read - just as
# wrong when the cursor was trusted but lagging a log that has since grown as
# when it was untrusted outright. So every fallback below reports the set it
# has and succeeds only when that set is non-empty; nothing to report means
# an incomplete fold, which the drain's existing notice already surfaces.
cur_ident=$(_fm_open_decisions_file_ident "$f") || { printf '%s' "$trusted_open"; [ -n "$trusted_open" ]; return; }
[ -n "$cur_ident" ] || { printf '%s' "$trusted_open"; [ -n "$trusted_open" ]; return; }
actual_size=$(_fm_status_file_size "$f") \
|| { printf '%s' "$trusted_open"; [ -n "$trusted_open" ]; return; }
# A stat, size, or span-read failure is a genuine I/O error, not "the file is
# empty" and not permission to replay the persisted open set. Leave that set
# untouched for a later recovery call, but return failure without printing it:
# unread bytes can make any cached set stale in either direction.
cur_ident=$(_fm_open_decisions_file_ident "$f") || return 1
[ -n "$cur_ident" ] || return 1
actual_size=$(_fm_status_file_size "$f") || return 1
actual_size=${actual_size//[[:space:]]/}
case "$actual_size" in ''|*[!0-9]*) printf '%s' "$trusted_open"; [ -n "$trusted_open" ]; return ;; esac
case "$actual_size" in ''|*[!0-9]*) return 1 ;; esac
if [ -n "$captured_end" ]; then
case "$captured_end" in
''|*[!0-9]*) printf '%s' "$trusted_open"; [ -n "$trusted_open" ]; return ;;
''|*[!0-9]*) return 1 ;;
esac
[ "$captured_end" -le "$actual_size" ] || { printf '%s' "$trusted_open"; [ -n "$trusted_open" ]; return; }
[ "$captured_end" -le "$actual_size" ] || return 1
size=$captured_end
else
size=$actual_size
Expand All @@ -991,19 +983,18 @@ status_open_decisions_incremental() { # <status-file> [<captured-end-offset>]
if [ -z "$version" ] || [ -z "$ident" ] || [ "$ident" != "$cur_ident" ] || [ "$offset" -gt "$actual_size" ]; then
offset=0
open=''
trusted_open=''
cursor_dirty=1
fi

if [ "$offset" -lt "$size" ]; then
chunk_file="$cf.read.$$"
_fm_status_read_span "$f" "$offset" "$((size - offset))" > "$chunk_file" 2>/dev/null \
|| { rm -f "$chunk_file"; printf '%s' "$trusted_open"; [ -n "$trusted_open" ]; return; }
|| { rm -f "$chunk_file"; return 1; }
chunk_size=$(LC_ALL=C wc -c < "$chunk_file" 2>/dev/null) \
|| { rm -f "$chunk_file"; printf '%s' "$trusted_open"; [ -n "$trusted_open" ]; return; }
|| { rm -f "$chunk_file"; return 1; }
chunk_size=${chunk_size//[[:space:]]/}
case "$chunk_size" in
''|*[!0-9]*) rm -f "$chunk_file"; printf '%s' "$trusted_open"; [ -n "$trusted_open" ]; return ;;
''|*[!0-9]*) rm -f "$chunk_file"; return 1 ;;
esac
# Test-only observability seam (off by default, no production behavior
# change): when set, records exactly how many bytes THIS call folded, so a
Expand Down Expand Up @@ -1403,6 +1394,9 @@ EOF
return "$rc"
}

# Only a drain that computed its whole presentation reaches this pass at all, so
# every task here was presented in full: an annotation or snapshot failure
# anywhere skips the acknowledge and commit entirely and no cursor moves.
status_acknowledge_presented_snapshot() { # <state> <snapshot> [<fully-presented-task-ids>]
local state=$1 snapshot=$2 fully_presented=${3:-} task endpoint ident f offset lines line safe
while IFS=$(printf '\t') read -r task endpoint ident; do
Expand Down
89 changes: 58 additions & 31 deletions bin/fm-wake-drain.sh
Original file line number Diff line number Diff line change
Expand Up @@ -399,7 +399,9 @@ EOF
# cursor-backed unread span as the annotation path, and runs on every drain -
# including the empty-queue fast path - so a buried answer cannot be swallowed
# when the fold later advances the cursor. Prints nothing when nothing is
# unread, which is the common case.
# unread, which is the common case. The header's "not re-printed" promise holds
# because this section runs only on a drain that goes on to commit its
# presentation receipt: any earlier failure skips the sections outright.
print_unread_status_section() { # <task-and-endpoint-snapshot>
local snapshot=$1 unread task line shown=0

Expand Down Expand Up @@ -543,54 +545,70 @@ EOF
printf 'RECORD DIVERGENCE: reconcile each one - record the captain'"'"'s own words with bin/fm-captain-hold.sh answer <task> --decision-file <path>, or re-open the status decision when that resolution was not the captain'"'"'s word.\n' || return 1
}

# A read or write hiccup anywhere in the drain's fleet-wide section passes (one
# task's status log, one task's open-decisions cursor, the scratch file the
# sections are prepared into) must never be indistinguishable from "computed,
# and genuinely nothing is open or unread": that silence is exactly what let a
# captain-facing OPEN DECISIONS section vanish for a drain even though several
# tasks' needs-decision/blocked lines were still open and unresolved in their
# own durable status logs: a failure on any ONE task aborts the whole
# acknowledge pass via its `|| return 1` before a single section is prepared,
# and the preparation that follows is itself all-or-nothing.
# Print this notice on every such failure instead of returning silently, so an
# empty presentation can only ever mean the passes ran to completion and found
# nothing - never that they could not be computed. print_status_presentation
# emits it from one place for every failure it observes, so no new failure
# path inside those passes can return without it.
print_status_sections_incomplete_notice() {
printf 'STATUS PRESENTATION INCOMPLETE: unread status, outcome backstop, OPEN DECISIONS, and record divergence could not be fully computed this drain (a status log or cursor read/write failed); do not read this drain'"'"'s silence as nothing open or unread - retry on the next drain.\n'
# A computation failure before the prepared section bytes reach stdout must
# never be indistinguishable from "computed, and genuinely nothing is open or
# unread". The caller supplies the failed operation so the notice says both
# what is missing and why. Receipt persistence is deliberately separate: once
# every prepared byte reached stdout, only the receipt can be incomplete.
print_status_sections_incomplete_notice() { # <reason>
local reason=$1
printf 'STATUS PRESENTATION INCOMPLETE: unread status, outcome backstop, OPEN DECISIONS, and record divergence could not be fully computed this drain (%s); do not read this drain'"'"'s silence as nothing open or unread - retry on the next drain.\n' "$reason"
}

print_status_receipt_failure_notice() {
printf 'STATUS PRESENTATION RECEIPT FAILED: the status sections above were fully computed and printed, but their presentation receipt could not be committed; they may repeat on the next drain.\n'
}

# One notice for every failure that leaves this drain with no presentation to
# make at all - the lock, the fleet snapshot, or the annotation pass, whether or
# not the failure can name a task. Such a drain acknowledges nothing and commits
# nothing, so "nothing was marked as seen" needs no per-task bookkeeping to be
# true: every presentation cursor is still where the previous drain left it.
print_status_presentation_incomplete_notice() { # <reason>
local reason=$1
printf 'STATUS PRESENTATION INCOMPLETE: %s; no status annotations or fleet-wide status sections were computed this drain, nothing was marked as seen and no presentation cursor advanced, so every unread status line is still unread.\n' \
"$reason"
}

print_status_sections() { # <task-and-endpoint-snapshot> [<fully-presented-task-ids>]
local snapshot=$1 fully_presented=${2:-} acknowledged prepared
acknowledged=$(status_acknowledge_presented_snapshot "$STATE" "$snapshot" "$fully_presented") || return 1
prepared=$(mktemp "$STATE/.status-presentation.prepared.XXXXXX") || return 1
acknowledged=$(status_acknowledge_presented_snapshot "$STATE" "$snapshot" "$fully_presented") || {
print_status_sections_incomplete_notice 'a status log or presentation cursor could not be read'
return 1
}
prepared=$(mktemp "$STATE/.status-presentation.prepared.XXXXXX") || {
print_status_sections_incomplete_notice 'the prepared-output file could not be created'
return 1
}
if ! {
print_unread_status_section "$snapshot" \
&& print_status_outcome_backstop_section "$snapshot" \
&& print_open_decisions_section "$snapshot" \
&& print_record_divergence_section
} > "$prepared"; then
rm -f -- "$prepared"
print_status_sections_incomplete_notice 'a status log, cursor, or prepared-output write could not be completed'
return 1
fi
# Prepare every section before presentation, but do not commit its receipt
# until the prepared bytes reach stdout. If the consumer closes or fails,
# leave the receipt behind so the next drain can recover the presentation.
if ! command cat "$prepared"; then
rm -f -- "$prepared"
print_status_sections_incomplete_notice 'the prepared section bytes could not be delivered to stdout'
return 1
fi
if ! status_commit_presentation_snapshot "$STATE" "$acknowledged"; then
rm -f -- "$prepared"
print_status_receipt_failure_notice
return 1
fi
rm -f -- "$prepared"
}

print_status_presentation() { # [<deduped-raw-rows>]
local rows=${1:-} lock="$STATE/.status-presentation-lock" snapshot annotation_manifest fully_presented='' rc=0
local lock_rc holder_pid
local rows=${1:-} lock="$STATE/.status-presentation-lock" snapshot='' annotation_manifest fully_presented='' rc=0
local lock_rc holder_pid incomplete_reason=''
if fm_lock_acquire_wait_bounded "$lock" "$PRESENTATION_LOCK_TIMEOUT"; then
:
else
Expand All @@ -600,29 +618,38 @@ print_status_presentation() { # [<deduped-raw-rows>]
printf 'STATUS PRESENTATION SKIPPED: lock remains held by live pid %s after %ss; retry on the next drain.\n' \
"$holder_pid" "$PRESENTATION_LOCK_TIMEOUT"
else
printf 'wake drain: status presentation lock could not be acquired safely\n' >&2
print_status_presentation_incomplete_notice 'status presentation lock could not be acquired safely'
fi
return 1
fi
# A failed snapshot read still leaves every task it printed BEFORE the failure
# in the captured value. Acknowledging and committing that truncated fleet view
# would rewrite the shared presentation-cursor manifest without the tasks the
# read never reached, resetting their unread and outcome-backstop cursors for
# good, so rc=1 here keeps the annotation, acknowledge and commit passes below
# from running at all rather than presenting a partial fleet view.
# good, so the partial capture is discarded here and nothing below runs on it.
snapshot=$(status_presentation_snapshot "$STATE") || {
printf 'STATUS PRESENTATION INCOMPLETE: status snapshot could not be read.\n'
snapshot=
incomplete_reason='status snapshot could not be read'
rc=1
}
if [ "$rc" -eq 0 ] && [ -n "$rows" ]; then
fm_wake_print_annotations "$rows" "$snapshot" || rc=1
if [ "$rc" -eq 0 ]; then
annotation_manifest=$(fm_wake_annotation_manifest "$rows") || rc=1
fully_presented=$(printf '%s\n' "$annotation_manifest" | awk -F '\t' '$2 == "direct" { sub(/\.status$/, "", $1); print $1 }') || rc=1
if fm_wake_print_annotations "$rows" "$snapshot"; then
annotation_manifest=$(fm_wake_annotation_manifest "$rows")
fully_presented=$(printf '%s\n' "$annotation_manifest" | awk -F '\t' '$2 == "direct" { sub(/\.status$/, "", $1); print $1 }')
else
incomplete_reason='a supplemental status annotation could not be computed'
rc=1
fi
fi
if [ "$rc" -eq 0 ] && [ -n "$snapshot" ]; then
print_status_sections "$snapshot" "$fully_presented" || { rc=1; print_status_sections_incomplete_notice; }
# One rule for the snapshot and the annotation pass alike, whether or not the
# failure can name a task: a drain that could not compute either presents
# nothing as seen. It says so once and skips the acknowledge and commit passes
# entirely, so no task's presentation cursor moves and the next drain still
# owes every unread line.
if [ -n "$incomplete_reason" ]; then
print_status_presentation_incomplete_notice "$incomplete_reason"
elif [ -n "$snapshot" ]; then
print_status_sections "$snapshot" "$fully_presented" || rc=1
fi
fm_lock_release "$lock"
return "$rc"
Expand Down
Loading
Loading