Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
442 changes: 399 additions & 43 deletions .pi/extensions/fm-primary-pi-watch.ts

Large diffs are not rendered by default.

3 changes: 2 additions & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,7 @@ state/ runtime records and signals; gitignored
branch-outcomes.jsonl .branch-outcomes-cursor Pi supervision-branch durable outcome store and its read cursor; bin/fm-branch-outcome.sh owns the format
branch-session/ .branch-session .branch-mirror-cursor the branch's persistent conversation, its pointer, and the dialog-mirror cursor; extension-owned (docs/pi-supervision-branch.md)
.branch-eligible-rows .branch-eligible-owner .main-eligible-rows per-actor wake-row claims and branch-owner evidence; docs/watcher-continuity.md owns the acknowledgement contract
.branch-presented-rows .main-presented-rows post-output per-actor wake-row receipts used only by Pi delayed-delivery revalidation; docs/watcher-continuity.md owns the receipt contract
.lease-<task> per-task supervision lease naming which actor (main or branch) may change that task; bin/fm-lease-lib.sh owns the contract the guarded scripts enforce
x-watch.check.sh generated Relay poll shim; present only when opted in (section 14)
tool-updates.check.sh generated watched-tool update poll shim and its .check-trust binding; present only after bin/fm-tool-update-check.sh arm; its report record .tool-updates is what keeps one pending update from being reported on every poll
Expand All @@ -134,7 +135,7 @@ state/ runtime records and signals; gitignored
.watch.lock .wake-queue.lock watcher singleton and queue serialization locks
.claude-autoarm.lock .claude-autoarm-epoch .claude-autoarm-failure-notified .claude-autoarm-failure-alarmed .turnend-claude-blocks .turnend-claude-blocks.lock Claude Stop auto-arm single-flight, epoch, failure-episode, attended-alarm, guard-budget, and budget-lock records; never touch
.cursor-park-owner .cursor-park-owner.lock .turnend-cursor-blocks Cursor stop-hook owner record, publication and commit lock, and bounded repair-nag budget; never touch
.hash-* .count-* .stale-* .stale-since-* .churn-since-* .paused-* .wedge-escalations-* .writing-* .seen-* .hb-surfaced-* .last-* .heartbeat-streak watcher internals; never touch
.hash-* .count-* .stale-* .stale-since-* .churn-since-* .paused-* .wedge-escalations-* .writing-* .validation-* .seen-* .hb-surfaced-* .last-* .heartbeat-streak watcher internals; never touch
.watch-triage.log watcher's absorbed-wake debug log (size-capped); never relied on, safe to delete
.last-watcher-beat watcher liveness beacon, touched every poll (including while absorbing benign wakes); guard scripts read it
.subsuper-* .supervise-daemon.* sub-supervisor internals; never touch
Expand Down
132 changes: 110 additions & 22 deletions bin/fm-crew-state.sh
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,13 @@
#
# state: <working|parked|done|blocked|paused|failed|unknown> · source: <run-step|pane|status-log|remote-endpoint|none> · <detail>
#
# `fm-crew-state.sh <id> --activity-token` reuses the same attribution and emits
# `activity: <opaque-token>` only for a full, actively-working structured run
# whose current step log is readable. The watcher snapshots that token at
# wedge-timer start and compares it only at the escalation boundary; it is
# progress evidence, never a current-state source. Every other state prints
# nothing in this mode.
#
# Logic, in order:
# 1. Resolve worktree + backend target + kind from state/<id>.meta. A meta
# recording remote_host= is a remote secondmate: its worktree and endpoint
Expand All @@ -29,7 +36,8 @@
# 2. Attribute an active or terminal no-mistakes run under the branch, head,
# pipeline-custody, and newest-first rules owned by bin/fm-nm-run-lib.sh.
# The run-step is AUTHORITATIVE: running/fixing -> working, ci -> working,
# awaiting_approval/fix_review -> parked (with gate findings), terminal
# awaiting_approval/fix_review -> parked (with gate findings), including
# when a nonterminal run retains a stale terminal outcome; terminal
# passed/checks-passed -> done, failed/cancelled -> failed. EXCEPT: while
# the active step is ci, `axi status` alone cannot tell "still waiting on
# checks" from "checks green, waiting on merge" (see nm_ci_checks_state) -
Expand Down Expand Up @@ -68,7 +76,12 @@ STATE="${FM_STATE_OVERRIDE:-$FM_HOME/state}"
. "$SCRIPT_DIR/fm-nm-run-lib.sh"

ID=${1:-}
[ -n "$ID" ] || { echo "usage: fm-crew-state.sh <id>" >&2; exit 2; }
MODE=${2:-}
[ -n "$ID" ] || { echo "usage: fm-crew-state.sh <id> [--activity-token]" >&2; exit 2; }
case "$MODE" in ''|--activity-token) ;; *) echo "usage: fm-crew-state.sh <id> [--activity-token]" >&2; exit 2 ;; esac
[ "$#" -le 2 ] || { echo "usage: fm-crew-state.sh <id> [--activity-token]" >&2; exit 2; }
ACTIVITY_MODE=false
[ "$MODE" != --activity-token ] || ACTIVITY_MODE=true

META="$STATE/$ID.meta"
LOG="$STATE/$ID.status"
Expand All @@ -85,6 +98,7 @@ SEP=' · '
# Emit the one canonical line and exit 0. Detail is optional.
emit() { # <state> <source> [detail]
local line="state: $1${SEP}source: $2"
[ "$ACTIVITY_MODE" = false ] || exit 0
[ -n "${3:-}" ] && line="$line${SEP}$3"
printf '%s\n' "$line"
exit 0
Expand All @@ -103,6 +117,9 @@ KIND=$(meta_value kind)
HARNESS=$(meta_value harness)
REMOTE_HOST=$(meta_value remote_host)
[ -n "$KIND" ] || KIND=ship
if [ "$ACTIVITY_MODE" = true ] && { [ "$KIND" != ship ] || [ -n "$REMOTE_HOST" ]; }; then
exit 0
fi

# A torn-down (or never-created) worktree has no current state to read. A
# remote secondmate's recorded worktree is a path on ITS host, so the local
Expand Down Expand Up @@ -331,21 +348,65 @@ nm_effective_ci_step_status() {
# for the MOST RECENT recognized marker (the log is append-only/chronological,
# so the last match is current): green with nothing red after it means CI is
# green right now, still only waiting on merge/close.
NM_CI_LOG_TAIL=
NM_CI_CHECKS_STATE=unknown
nm_ci_checks_state() {
local run_id log_tail marker
local run_id marker
NM_CI_LOG_TAIL=
NM_CI_CHECKS_STATE=unknown
run_id=$(strip_quotes "$(nm_field id)")
[ -n "$run_id" ] || { printf 'unknown'; return; }
log_tail=$(nm_run axi logs --step ci --run "$run_id") || true
[ -n "$log_tail" ] || { printf 'unknown'; return; }
marker=$(printf '%s\n' "$log_tail" \
[ -n "$run_id" ] || return 0
NM_CI_LOG_TAIL=$(nm_run axi logs --step ci --run "$run_id") || true
[ -n "$NM_CI_LOG_TAIL" ] || return 0
marker=$(printf '%s\n' "$NM_CI_LOG_TAIL" \
| grep -E 'CI checks passed|no CI checks reported - still monitoring|no CI checks reported yet|checks failed|issues detected|CI checks running|base branch advanced.*re-arming CI monitor timeout' \
| tail -1)
case "$marker" in
*"checks passed"*|*"no CI checks reported - still monitoring"*) printf 'green' ;;
*"no CI checks reported yet"*|*"checks failed"*|*"issues detected"*|*"CI checks running"*|*"base branch advanced"*"re-arming CI monitor timeout"*) printf 'not-ready' ;;
*) printf 'unknown' ;;
*"checks passed"*|*"no CI checks reported - still monitoring"*) NM_CI_CHECKS_STATE=green ;;
*"no CI checks reported yet"*|*"checks failed"*|*"issues detected"*|*"CI checks running"*|*"base branch advanced"*"re-arming CI monitor timeout"*) NM_CI_CHECKS_STATE=not-ready ;;
esac
}

validation_activity_hash() {
if command -v shasum >/dev/null 2>&1; then
shasum -a 256 | awk '{print "sha256-" $1}'
elif command -v sha256sum >/dev/null 2>&1; then
sha256sum | awk '{print "sha256-" $1}'
else
cksum | awk '{print "cksum-" $1 "-" $2}'
fi
}

nm_active_step() {
local row step
row=$(printf '%s\n' "$RUN_OUT" \
| grep -E '^[[:space:]]*(intent|rebase|review|test|document|lint|push|pr|ci),[[:space:]]*"?(running|fixing)"?[[:space:]]*,' \
| tail -1)
if [ -n "$row" ]; then
row=$(trim "$row")
step=$(strip_quotes "${row%%,*}")
printf '%s' "$step"
return 0
fi
[ "${RUN_STATUS:-}" = ci ] && printf 'ci'
}

nm_validation_activity_token() {
local run_id step log_tail digest
run_id=$(strip_quotes "$(nm_field id)")
step=$(nm_active_step)
case "$run_id" in ''|.*|*[!A-Za-z0-9._-]*) return 0 ;; esac
case "$step" in intent|rebase|review|test|document|lint|push|pr|ci) ;; *) return 0 ;; esac
if [ "$step" = ci ] && [ -n "$NM_CI_LOG_TAIL" ]; then
log_tail=$NM_CI_LOG_TAIL
else
log_tail=$(nm_run axi logs --step "$step" --run "$run_id") || true
fi
[ -n "$log_tail" ] || return 0
digest=$(printf '%s' "$log_tail" | validation_activity_hash) || return 0
[ -n "$digest" ] || return 0
printf '%s-%s-%s' "$run_id" "$step" "$digest"
}
# Coarse fallback for cross-branch attribution. `no-mistakes axi status` (bare)
# reports the active-or-most-recent run for the CURRENT branch when one
# exists, else falls back to some other branch's run purely as informational
Expand Down Expand Up @@ -496,16 +557,21 @@ if [ "$HAVE_RUN" = 1 ]; then
gate_status=$(nm_gate_status)
has_gate=0
nm_has_gate && has_gate=1

if [ -n "$outcome" ]; then
case "$outcome" in
passed) RUN_STATE="done"; RUN_DETAIL="run passed: PR merged/closed" ;;
checks-passed) RUN_STATE="done"; RUN_DETAIL="checks green: PR ready for review" ;;
failed) RUN_STATE=failed; RUN_DETAIL="run failed" ;;
cancelled) RUN_STATE=failed; RUN_DETAIL="run cancelled" ;;
*) RUN_STATE=unknown; RUN_DETAIL="outcome: $outcome" ;;
esac
elif [ -n "$awaiting" ] || [ "$status" = awaiting_approval ] || [ "$status" = fix_review ] || [ -n "$gate_status" ] || [ "$has_gate" = 1 ]; then
agent_wait=0
fm_nm_run_has_agent_wait "$RUN_OUT" && agent_wait=1
gate_wait=0
if [ -n "$awaiting" ] || [ "$status" = awaiting_approval ] || [ "$status" = fix_review ] \
|| [ -n "$gate_status" ] || [ "$has_gate" = 1 ]; then
gate_wait=1
fi
current_gate=0
case "$status" in completed|failed|cancelled) ;; *) [ "$gate_wait" -eq 0 ] || current_gate=1 ;; esac

# A direct fix-review/approval wait on a nonterminal structured run is the
# current actionable state even if a stale terminal outcome field remains.
# A bare scalar gate keeps its historical behavior only when no outcome
# contradicts it; terminal top-level status always remains terminal.
if [ "$current_gate" -eq 1 ] && { [ "$agent_wait" -eq 1 ] || [ -z "$outcome" ]; }; then
if [ "$has_gate" = 1 ]; then
gate=$(nm_gate_line_name)
else
Expand All @@ -515,11 +581,22 @@ if [ "$HAVE_RUN" = 1 ]; then
[ -n "$gate" ] || gate=gate
RUN_STATE=parked
RUN_DETAIL="parked at $gate"
if [ -n "$gate_status" ] && [ "$gate_status" != "$gate" ]; then
RUN_DETAIL="$RUN_DETAIL ($gate_status)"
fi
fcount=$(nm_gate_findings_count)
[ -n "$fcount" ] && RUN_DETAIL="$RUN_DETAIL: $fcount finding(s)"
if printf '%s\n' "$RUN_OUT" | grep -q 'ask-user'; then
RUN_DETAIL="$RUN_DETAIL (ask-user: authority decision)"
fi
elif [ -n "$outcome" ]; then
case "$outcome" in
passed) RUN_STATE="done"; RUN_DETAIL="run passed: PR merged/closed" ;;
checks-passed) RUN_STATE="done"; RUN_DETAIL="checks green: PR ready for review" ;;
failed) RUN_STATE=failed; RUN_DETAIL="run failed" ;;
cancelled) RUN_STATE=failed; RUN_DETAIL="run cancelled" ;;
*) RUN_STATE=unknown; RUN_DETAIL="outcome: $outcome" ;;
esac
else
case "$status" in
ci) RUN_STATE=working; RUN_DETAIL="ci running" ;;
Expand All @@ -534,7 +611,8 @@ if [ "$HAVE_RUN" = 1 ]; then
CI_STEP_STATUS=$(nm_effective_ci_step_status)
case "$CI_STEP_STATUS" in
running)
CI_LOG_STATE=$(nm_ci_checks_state)
nm_ci_checks_state
CI_LOG_STATE=$NM_CI_CHECKS_STATE
if [ "$CI_LOG_STATE" = green ]; then
RUN_STATE="done"
RUN_DETAIL="checks green: PR ready for review (still monitoring for merge/close)"
Expand All @@ -556,7 +634,8 @@ if [ "$HAVE_RUN" = 1 ]; then
if [ "$RUN_STATUS" = fixing ]; then
CI_LOG_STATE=not-ready
elif [ "$CI_STEP_STATUS" = running ] && [ -z "$CI_LOG_STATE" ]; then
CI_LOG_STATE=$(nm_ci_checks_state)
nm_ci_checks_state
CI_LOG_STATE=$NM_CI_CHECKS_STATE
elif [ "$CI_STEP_STATUS" = fixing ]; then
CI_LOG_STATE=not-ready
fi
Expand All @@ -580,9 +659,18 @@ if [ "$HAVE_RUN" = 1 ]; then
;;
esac

if [ "$ACTIVITY_MODE" = true ]; then
if [ "$RUN_STATE" = working ] && [ "$RUN_SOURCE" = full ]; then
ACTIVITY_TOKEN=$(nm_validation_activity_token)
[ -z "$ACTIVITY_TOKEN" ] || printf 'activity: %s\n' "$ACTIVITY_TOKEN"
fi
exit 0
fi
emit "$RUN_STATE" run-step "$RUN_DETAIL"
fi

[ "$ACTIVITY_MODE" = false ] || exit 0

# --- fallback: no run attributed to this crew ------------------------------
# The run-step path above already handled any crew with a run, regardless of pane
# liveness, so a finished-but-pane-closed crew never reaches here. Down here there
Expand Down
25 changes: 20 additions & 5 deletions bin/fm-nm-run-lib.sh
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,8 @@
# fm-crew-state.sh (read-only current-state reporting) and fm-teardown.sh
# (pre-teardown run abort, see its "Fix 1" header comment). Teardown uses only
# strict branch-and-head identity; crew-state additionally permits the active
# pipeline-owned exemption defined below. Getting this wrong in either
# pipeline-owned exemption and current agent-wait evidence defined below. Getting
# this wrong in either
# direction is unsafe: a false negative hides a genuinely parked run, and a
# false positive lets teardown act on a run it does not own.
#
Expand Down Expand Up @@ -56,6 +57,18 @@ fm_nm_field() { # <toon-output> <key>
printf '%s\n' "$1" | sed -n "s/^[[:space:]]*$2:[[:space:]]*\(.*\)/\1/p" | head -1
}

# 0 when an active top-level run status carries direct evidence that the agent
# must answer an approval/fix-review wait. This stronger current-gate evidence
# outranks a stale terminal outcome field, but never a terminal top-level status.
fm_nm_run_has_agent_wait() { # <toon-output>
local out=$1 status
status=$(fm_nm_strip_quotes "$(fm_nm_field "$out" status)")
case "$status" in completed|failed|cancelled) return 1 ;; esac
case "$status" in awaiting_approval|fix_review) return 0 ;; esac
printf '%s\n' "$out" | grep -Eq \
'^[[:space:]]*awaiting_agent:|^[[:space:]]*[^,]+,[[:space:]]*"?(awaiting_approval|fix_review)"?[[:space:]]*,|^[[:space:]]*(status|state):[[:space:]]*"?(awaiting_approval|fix_review)"?[[:space:]]*$'
}

# 0 if run head $2 matches worktree $1's code identity, per the same rule
# everywhere this attribution is needed:
# - missing/empty head: cannot bind; reject
Expand Down Expand Up @@ -99,14 +112,16 @@ fm_nm_branch_sync_state() { # <toon-output>
fm_nm_strip_quotes "$s"
}

# 0 if the run in captured `axi status` TOON $1 is still in flight: no
# terminal outcome and no terminal status.
# 0 if the run in captured `axi status` TOON $1 is still in flight: its
# top-level status is nonterminal and it either has no terminal outcome or has
# direct current agent-wait evidence that resolves a contradictory stale outcome
# toward the actionable gate.
fm_nm_run_is_active() { # <toon-output>
local status outcome
status=$(fm_nm_strip_quotes "$(fm_nm_field "$1" status)")
outcome=$(fm_nm_strip_quotes "$(fm_nm_field "$1" outcome)")
[ -z "$outcome" ] || return 1
case "$status" in completed|failed|cancelled) return 1 ;; esac
outcome=$(fm_nm_strip_quotes "$(fm_nm_field "$1" outcome)")
[ -z "$outcome" ] || fm_nm_run_has_agent_wait "$1"
}

# The one exemption to the head rule above: while the pipeline OWNS the branch
Expand Down
20 changes: 16 additions & 4 deletions bin/fm-push-transition-lib.sh
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,8 @@ FM_WAKE_POST_OUTPUT_ACTION=
FM_WATCH_DELIVERED_REASON=
FM_WATCH_DELIVERY_PID=
FM_WATCH_DELIVERY_IDENTITY=
# Ledger rows are pid<TAB>process-identity<TAB>reason<TAB>Pi-ticket-sideband.
# The fourth field is empty outside Pi mode, and readers accept legacy three-field rows.
WATCH_DELIVERY_LOG="$STATE/.watch-deliveries.log"
WATCH_DELIVERY_LOCK="$STATE/.watch-deliveries.lock"
WATCH_DELIVERY_MAX_BYTES=${FM_WATCH_DELIVERY_MAX_BYTES:-65536}
Expand All @@ -40,8 +42,12 @@ watch_delivery_clean_reason() {
printf '%s' "$1" | tr '\t\r\n' ' ' | cut -c1-4096
}

watch_delivery_clean_tickets() {
printf '%s' "$1" | tr '\t\r\n' ' ;' | cut -c1-8192
}

watch_delivery_publish() {
local reason=$1 i size tmp raw
local reason=$1 tickets=${2:-} i size tmp raw
[ -n "$FM_WATCH_DELIVERY_PID" ] || return 0
[ -n "$FM_WATCH_DELIVERY_IDENTITY" ] || return 0
i=0
Expand All @@ -50,10 +56,11 @@ watch_delivery_publish() {
sleep 0.02
i=$((i + 1))
done
printf '%s\t%s\t%s\n' \
printf '%s\t%s\t%s\t%s\n' \
"$FM_WATCH_DELIVERY_PID" \
"$(watch_delivery_clean_identity "$FM_WATCH_DELIVERY_IDENTITY")" \
"$(watch_delivery_clean_reason "$reason")" >> "$WATCH_DELIVERY_LOG" 2>/dev/null || true
"$(watch_delivery_clean_reason "$reason")" \
"$(watch_delivery_clean_tickets "$tickets")" >> "$WATCH_DELIVERY_LOG" 2>/dev/null || true
size=$(wc -c < "$WATCH_DELIVERY_LOG" 2>/dev/null | tr -d '[:space:]')
case "$size" in
''|*[!0-9]*) ;;
Expand Down Expand Up @@ -95,7 +102,12 @@ wake() {
[ -z "$FM_WAKE_POST_OUTPUT_ACTION" ] || trap '' PIPE
if echo "$1"; then
output_status=0
watch_delivery_publish "$1" || true
# Pi alone consumes the machine-only identity lines. Every other primary
# keeps the historical one-line arm completion surface byte-for-byte.
if [ "${FM_PI_WATCH_DELIVERY:-0}" = 1 ] && [ -n "$FM_WAKE_APPENDED_TICKETS" ]; then
printf '%s\n' "$FM_WAKE_APPENDED_TICKETS" || true
fi
watch_delivery_publish "$1" "$FM_WAKE_APPENDED_TICKETS" || true
# shellcheck disable=SC2034 # Read by bin/fm-watch.sh's EXIT cleanup.
FM_WATCH_DELIVERED_REASON=$1
else
Expand Down
3 changes: 2 additions & 1 deletion bin/fm-supervise-daemon.sh
Original file line number Diff line number Diff line change
Expand Up @@ -512,7 +512,8 @@ clear_pause_tracking() { # <window> <state>
rm -f "$state/.subsuper-paused-$key" "$state/.subsuper-stale-$key" \
"$state/.paused-$watcher_key" "$state/.paused-rechecked-$watcher_key" "$state/.paused-resurfaced-$watcher_key" \
"$state/.stale-$watcher_key" "$state/.stale-since-$watcher_key" "$state/.wedge-escalations-$watcher_key" \
"$state/.writing-since-$watcher_key" "$state/.writing-resurfaced-$watcher_key"
"$state/.writing-since-$watcher_key" "$state/.writing-resurfaced-$watcher_key" \
"$state/.validation-since-$watcher_key" "$state/.validation-resurfaced-$watcher_key"
}

reconcile_pause_tracking() { # <window> <state> <last-status-line>
Expand Down
20 changes: 20 additions & 0 deletions bin/fm-wake-drain.sh
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ ACTOR=$(fm_lease_actor) || exit 2
ELIGIBLE_ROWS_FILE="$STATE/.branch-eligible-rows"
ELIGIBLE_OWNER_FILE="$STATE/.branch-eligible-owner"
MAIN_ROWS_FILE="$STATE/.main-eligible-rows"
PRESENTED_ROWS_FILE="$STATE/.$ACTOR-presented-rows"

rows_file_valid() {
[ -s "$1" ] && awk 'BEGIN { ok=1 } !/^[0-9]+$/ || seen[$0]++ { ok=0 } END { exit !ok }' "$1"
Expand Down Expand Up @@ -102,6 +103,24 @@ write_rows_file_locked() { # <target> <source>
_fm_atomic_replace "$source" "$target"
}

# Publish only rows whose raw presentation reached this actor's stdout
# successfully. Pi's delayed-follow-up validator uses the receipt to distinguish
# a row merely claimed before output from one already delivered to a handling
# turn. Queue acknowledgement remains the consumption authority.
publish_presented_rows_locked() { # <deduped-raw-rows>
local raw=$1 tmp
if [ -e "$PRESENTED_ROWS_FILE" ] || [ -L "$PRESENTED_ROWS_FILE" ]; then
[ -f "$PRESENTED_ROWS_FILE" ] && [ ! -L "$PRESENTED_ROWS_FILE" ] || return 1
fi
tmp=$(mktemp "$STATE/.$ACTOR-presented-rows.tmp.XXXXXX") || return 1
if ! printf '%s\n' "$raw" | awk -F '\t' '$2 ~ /^[0-9]+$/ { print $2 }' > "$tmp" \
|| ! chmod 0600 "$tmp" || ! _fm_atomic_replace "$tmp" "$PRESENTED_ROWS_FILE" \
|| [ ! -f "$PRESENTED_ROWS_FILE" ] || [ -L "$PRESENTED_ROWS_FILE" ]; then
rm -f -- "$tmp"
return 1
fi
}

claim_main_rows_locked() {
DRAIN_TMP=$(mktemp "$STATE/.main-eligible-rows.tmp.XXXXXX") || return 1
awk -F '\t' -v branch="$ELIGIBLE_ROWS_FILE" -v main="$MAIN_ROWS_FILE" '
Expand Down Expand Up @@ -605,6 +624,7 @@ case "${FM_WAKE_DRAIN_TEST_DELAY_BEFORE_COMMIT:-0}" in
esac
if [ -n "$RAW_ROWS" ]; then
printf '%s\n' "$RAW_ROWS" || exit "$?"
publish_presented_rows_locked "$RAW_ROWS" || exit 1
fi
fm_recovery_marker_snapshot "$RECOVERY_MARKER" || exit 1
RECOVERY_MARKER_TOKEN=$FM_RECOVERY_MARKER_TOKEN
Expand Down
Loading