Skip to content
Merged
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
58 changes: 55 additions & 3 deletions bin/fm-procevent-remote-reply.sh
Original file line number Diff line number Diff line change
Expand Up @@ -128,6 +128,7 @@ source_id() {

cursor_path() { printf '%s/%s.cursor\n' "$CURSOR_DIR" "$1"; }
ingest_receipt_path() { printf '%s/%s.%s.ingested\n' "$CURSOR_DIR" "$1" "$2"; }
retirement_count_path() { printf '%s/%s.retirements\n' "$CURSOR_DIR" "$1"; }
mirrored_source_path() { printf '%s/.remote-reply-mirrored-%s\n' "$STATE" "$1"; }

read_cursor() { # <id>; sets CURSOR_OFFSET and CURSOR_HASH
Expand Down Expand Up @@ -164,6 +165,43 @@ write_cursor() { # <id> <offset> <hash>
mv -f -- "$tmp" "$path"
}

# A missing file is zero and is not created. Ingest only reads this.
# Retirement is the one writer, so a crash during ingest cannot change it.
read_retirement_count() { # <id>; sets RETIREMENT_COUNT
local path count lines
path=$(retirement_count_path "$1")
RETIREMENT_COUNT=0
[ -e "$path" ] || [ -L "$path" ] || return 0
[ -f "$path" ] && [ ! -L "$path" ] || die "reply retirement count is unsafe: $path"
lines=$(grep -c '^count=' "$path" 2>/dev/null || true)
[ "$lines" = 1 ] || die "reply retirement count is invalid: $path"
count=$(sed -n 's/^count=//p' "$path")
case "$count" in ''|*[!0-9]*) die "reply retirement count is invalid: $path" ;; esac
RETIREMENT_COUNT=$count
}

write_retirement_count() { # <id> <count>
local id=$1 count=$2 path tmp
case "$count" in ''|*[!0-9]*) return 1 ;; esac
mkdir -p "$CURSOR_DIR" || return 1
chmod 700 "$CURSOR_DIR" 2>/dev/null || true
path=$(retirement_count_path "$id")
[ ! -L "$path" ] || return 1
tmp=$(umask 077; mktemp "$CURSOR_DIR/.retirements.XXXXXX") || return 1
printf 'count=%s\n' "$count" > "$tmp" || { rm -f -- "$tmp"; return 1; }
chmod 600 "$tmp" || { rm -f -- "$tmp"; return 1; }
if ! mv -f -- "$tmp" "$path"; then
rm -f -- "$tmp"
return 1
fi
}

# Twelve characters distinguish breaks in the status line. The cursor keeps
# the full digest the reader uses.
continuity_prefix() {
printf '%.12s' "$CURSOR_HASH"
}

ingest_receipt_matches() { # <id> <sequence> <result>
local path stored actual count
path=$(ingest_receipt_path "$1" "$2")
Expand Down Expand Up @@ -544,7 +582,12 @@ cmd_ingest() {
die "result does not continue the current cursor for $id"
fi
if [ "$class" = continuity-broken ]; then
line="blocked [key=remote-reply-continuity-$id]: remote reply continuity broke for $id ($reason)"
# The same offset, prefix, and retirement count build the same line, so a
# retry appends nothing. Retirement removes the cursor before it records
# the next count, so a later break is a new line even when the restored
# bytes match, and a stop between those steps leaves the count unchanged.
read_retirement_count "$id"
line="blocked [key=remote-reply-continuity-$id]: remote reply continuity broke for $id ($reason) at offset ${CURSOR_OFFSET} prefix $(continuity_prefix) retirements ${RETIREMENT_COUNT}"
append_rc=0
if status_event_recorded "$status_file" "$line"; then
append_rc=1
Expand Down Expand Up @@ -726,7 +769,7 @@ cmd_retire_quiesce_locked() {
}

cmd_retire_finalize_locked() {
local id=${1:-} force=${2:-} sid path
local id=${1:-} force=${2:-} sid path cursor
validate_id "$id"
[ -z "$force" ] || [ "$force" = --force ] || die "invalid retirement option: $force"
sid=$(source_id "$id")
Expand All @@ -742,7 +785,16 @@ cmd_retire_finalize_locked() {
done
fi
fi
rm -f -- "$(cursor_path "$id")"
# Remove the cursor first. A stop before the count write leaves that count
# unchanged, so the same break still builds the same line.
cursor=$(cursor_path "$id")
rm -f -- "$cursor" || die "cannot remove remote reply cursor"
if [ -e "$cursor" ] || [ -L "$cursor" ]; then
die "cannot remove remote reply cursor"
fi
read_retirement_count "$id"
write_retirement_count "$id" "$((RETIREMENT_COUNT + 1))" \
|| die "cannot record remote reply retirement"
rm -f -- "$CURSOR_DIR/$id".*.ingested
rm -f -- "$(fm_pending_reply_remote_channel_watermark_path "$STATE" "$id")"
}
Expand Down
6 changes: 6 additions & 0 deletions docs/remote-secondmates.md
Original file line number Diff line number Diff line change
Expand Up @@ -588,6 +588,12 @@ The [process-to-event operating contract](configuration.md#process-to-event-sour
The source log is never truncated or consumed.
A shortened or changed prefix stops the relay and surfaces a continuity failure instead of silently resetting the cursor.

The failure appends one `blocked` line to the parent status stream, which opens a decision.
The line records the reason, the reader position (the cursor offset and the first 12 characters of the prefix hash), and the retirement count (how many times the route has been retired).
Reading the same break again, with the cursor where it was and the count unchanged, appends nothing.
A later break at a different reader position, or after another retirement, appends a new `blocked` line and opens the decision again.
A line written before that position was recorded does not match, so the next break appends the new line once.

### SSH exit 255 and unavailable homes

An SSH exit status of 255 always means transport failure or unknown remote completion.
Expand Down
230 changes: 230 additions & 0 deletions tests/fm-remote-reply.test.sh
Original file line number Diff line number Diff line change
Expand Up @@ -1024,6 +1024,15 @@ assert_absent "$PARENT/state/procevent/$SID.source" "continuity break was re-arm
remote_env "$ADAPTER" ingest ios "$RESULT_TWELVE" >/dev/null 2>&1 || true
[ "$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" -eq 1 ] \
|| fail "continuity replay duplicated the escalation"
first_offset=$(sed -n 's/^offset=//p' "$PARENT/state/remote-replies/ios.cursor")
first_hash=$(sed -n 's/^prefix_sha256=//p' "$PARENT/state/remote-replies/ios.cursor" | tr 'A-F' 'a-f')
first_prefix=$(printf '%.12s' "$first_hash")
assert_grep "at offset ${first_offset} prefix ${first_prefix} retirements 0" "$PARENT/state/ios.status" \
"continuity break did not record the reader position"
assert_no_grep "prefix ${first_hash}" "$PARENT/state/ios.status" \
"continuity break recorded the full prefix hash"
assert_absent "$PARENT/state/remote-replies/ios.retirements" \
"a route that has never been retired gained a retirement count"
status_line_at_epoch "$(grep -F 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" >/dev/null \
|| fail "new continuity escalation has unknown emission time"
if [ "${FM_TEST_EVIDENCE:-0}" = 1 ]; then
Expand Down Expand Up @@ -1064,8 +1073,229 @@ remote_env "$ADAPTER" handle ios "$GEN" "$RESULT_TWELVE" >/dev/null 2>&1 || [ "$
|| fail "pending continuity result could not be acknowledged after retirement refusal"
remote_env "$ADAPTER" retire ios >/dev/null
assert_absent "$PARENT/state/remote-replies/ios.cursor" "adapter retirement left its cursor"
recorded_retirements=$(cat "$PARENT/state/remote-replies/ios.retirements" 2>/dev/null || true)
[ "$recorded_retirements" = count=1 ] \
|| fail "adapter retirement did not record its count (got: ${recorded_retirements:-absent})"
assert_absent "$PARENT/state/remote-replies/ios.caught-up" \
"adapter retirement left a caught-up watermark a later route could inherit"
pass "remote reply retirement quiesces and refuses unhandled captured results"

# Empty the remote log under the committed cursor and handle the break the
# next blocking source reports. Sets RESULT_BREAK.
break_repaired_route() { # <label> [expected-handle-status]
local label=$1 expected=${2:-3} runner handle_rc
stop_reply_listener || fail "the reply listener did not stop before the $label continuity break"
: > "$REMOTE/state/parent-replies.status"
GEN=$((GEN + 1))
remote_env "$ROOT/bin/fm-procevent.sh" start "$SID" > "$TMP_ROOT/start-$label-break.out" 2>&1 &
runner=$!
wait "$runner" || fail "the $label continuity break was not captured"
RESULT_BREAK=$(find "$PARENT/state/procevent-inbox" -name "$SID.$GEN.result" -print -quit)
[ -n "$RESULT_BREAK" ] || fail "the $label continuity break produced no durable result"
[ "$(remote_env "$ADAPTER" classify "$RESULT_BREAK")" = continuity-broken ] \
|| fail "the $label truncation was not classified as a continuity break"
set +e
remote_env "$ADAPTER" handle ios "$GEN" "$RESULT_BREAK" > "$TMP_ROOT/handle-$label-break.out" 2>&1
handle_rc=$?
set -e
[ "$handle_rc" -eq "$expected" ] || fail "the $label continuity break returned an unexpected status: $handle_rc"
}

# Close the open continuity decision, put back the log the cursor was committed
# against with one more line, and let the reader advance over that line. The
# route is not retired, so the cursor moves only because new bytes were read.
resolve_and_extend_route() { # <label> <log-content>
local label=$1 content=$2 bytes blocked_before
blocked_before=$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")
printf '%s\n' 'resolved [key=remote-reply-continuity-ios]: operator accepted the break' \
>> "$PARENT/state/ios.status"
printf '%s' "$content" > "$REMOTE/state/parent-replies.status"
bytes=$(wc -c < "$REMOTE/state/parent-replies.status" | tr -d ' ')
remote_env "$ADAPTER" arm ios >/dev/null
GEN=$((GEN + 1))
await_reply_result "$PARENT/state/procevent-inbox/$SID.$GEN.result" \
|| fail "the route extended after the $label break was not read"
assert_grep "offset=$bytes" "$PARENT/state/remote-replies/ios.cursor" \
"the route extended after the $label break did not advance the cursor"
[ "$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" -eq "$blocked_before" ] \
|| fail "extending the route after the $label break appended a continuity break"
[ -z "$(status_open_decisions "$PARENT/state/ios.status")" ] \
|| fail "extending the route after the $label break reopened the continuity decision"
}

# The resolved break above stays closed through an unchanged re-read. Repair
# the log, let the reader advance, and truncate again. The later break is at
# another reader position, so its line is new and the decision opens again.
printf 'working: route restored and readable again\n' > "$REMOTE/state/parent-replies.status"
restored_bytes=$(wc -c < "$REMOTE/state/parent-replies.status" | tr -d ' ')
remote_env "$ADAPTER" arm ios >/dev/null
GEN=$((GEN + 1))
await_reply_result "$PARENT/state/procevent-inbox/$SID.$GEN.result" \
|| fail "the repaired route was not read"
assert_grep "offset=$restored_bytes" "$PARENT/state/remote-replies/ios.cursor" \
"the repaired route did not advance the cursor"
[ "$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" -eq 1 ] \
|| fail "repairing the route appended a continuity break"
[ -z "$(status_open_decisions "$PARENT/state/ios.status")" ] \
|| fail "repairing the route reopened the continuity decision"
break_repaired_route "second" 3
RESULT_SECOND=$RESULT_BREAK
[ "$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" -eq 2 ] \
|| fail "a later continuity break after repair appended nothing"
second_offset=$(sed -n 's/^offset=//p' "$PARENT/state/remote-replies/ios.cursor")
second_hash=$(sed -n 's/^prefix_sha256=//p' "$PARENT/state/remote-replies/ios.cursor" | tr 'A-F' 'a-f')
second_prefix=$(printf '%.12s' "$second_hash")
assert_grep "at offset ${second_offset} prefix ${second_prefix} retirements 1" "$PARENT/state/ios.status" \
"a later continuity break after repair did not record its reader position"
assert_contains "$(status_open_decisions "$PARENT/state/ios.status")" \
$'remote-reply-continuity-ios\t' \
"a later continuity break after repair did not reopen the decision"
remote_env "$ADAPTER" ingest ios "$RESULT_SECOND" >/dev/null 2>&1 || true
[ "$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" -eq 2 ] \
|| fail "a repeated read of the later continuity break appended again"
assert_contains "$(status_open_decisions "$PARENT/state/ios.status")" \
$'remote-reply-continuity-ios\t' \
"a repeated read of the later continuity break closed the decision"
pass "a later continuity break after repair and re-advance opens the decision again"

# A line from before the reader position was recorded names the route and the
# reason only. It does not match the new line, so this same break appends once.
awk '
/blocked \[key=remote-reply-continuity-ios\]/ {
sub(/ at offset [0-9]+ prefix [0-9a-f]+( retirements [0-9]+)?$/, "")
}
{ print }
' "$PARENT/state/ios.status" > "$TMP_ROOT/ios-status-old-format"
mv "$TMP_ROOT/ios-status-old-format" "$PARENT/state/ios.status"
printf '%s\n' 'resolved [key=remote-reply-continuity-ios]: operator accepted the break' \
>> "$PARENT/state/ios.status"
[ -z "$(status_open_decisions "$PARENT/state/ios.status")" ] \
|| fail "operator resolve left the old-format continuity decision open"
set +e
remote_env "$ADAPTER" handle ios "$GEN" "$RESULT_SECOND" > "$TMP_ROOT/handle-old-format.out" 2>&1
handle_rc=$?
set -e
[ "$handle_rc" -eq 3 ] || fail "a continuity break after an old-format line returned an unexpected status: $handle_rc"
[ "$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" -eq 3 ] \
|| fail "a continuity break after an old-format line appended nothing"
assert_grep "at offset ${second_offset} prefix ${second_prefix} retirements 1" "$PARENT/state/ios.status" \
"a continuity break after an old-format line did not record the reader position"
assert_contains "$(status_open_decisions "$PARENT/state/ios.status")" \
$'remote-reply-continuity-ios\t' \
"a continuity break after an old-format line did not reopen the decision"
remote_env "$ADAPTER" ingest ios "$RESULT_SECOND" >/dev/null 2>&1 || true
[ "$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" -eq 3 ] \
|| fail "a repeated read after the old-format upgrade appended again"
assert_contains "$(status_open_decisions "$PARENT/state/ios.status")" \
$'remote-reply-continuity-ios\t' \
"a repeated read after the old-format upgrade closed the decision"
pass "an old-format continuity line does not swallow the next break"

# No retirement this time. The cursor leaves the escalated offset only because
# the reader consumed new bytes, and that alone makes the next break new.
LOG_RESTORED=$'working: route restored and readable again\n'
LOG_EXTENDED=$LOG_RESTORED$'working: route extended without retirement\n'
resolve_and_extend_route "second" "$LOG_EXTENDED"
break_repaired_route "third" 3
[ "$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" -eq 4 ] \
|| fail "a later continuity break after the cursor moved without retirement appended nothing"
moved_offset=$(sed -n 's/^offset=//p' "$PARENT/state/remote-replies/ios.cursor")
moved_hash=$(sed -n 's/^prefix_sha256=//p' "$PARENT/state/remote-replies/ios.cursor" | tr 'A-F' 'a-f')
moved_prefix=$(printf '%.12s' "$moved_hash")
assert_grep "at offset ${moved_offset} prefix ${moved_prefix} retirements 1" "$PARENT/state/ios.status" \
"a later continuity break after the cursor moved did not record its reader position"
assert_contains "$(status_open_decisions "$PARENT/state/ios.status")" \
$'remote-reply-continuity-ios\t' \
"a later continuity break after the cursor moved without retirement did not reopen the decision"
pass "a later continuity break after the cursor moves without retirement opens the decision again"

# Retire, put the same bytes back, and break at the same offset. The retirement
# count makes that line new, so the decision opens once. The repeat stays silent.
printf '%s\n' 'resolved [key=remote-reply-continuity-ios]: operator accepted the break' \
>> "$PARENT/state/ios.status"
[ -z "$(status_open_decisions "$PARENT/state/ios.status")" ] \
|| fail "operator resolve left the extended continuity decision open"
remote_env "$ADAPTER" retire ios >/dev/null
recorded_retirements=$(cat "$PARENT/state/remote-replies/ios.retirements" 2>/dev/null || true)
[ "$recorded_retirements" = count=2 ] \
|| fail "the second retirement did not advance the count (got: ${recorded_retirements:-absent})"
printf '%s' "$LOG_EXTENDED" > "$REMOTE/state/parent-replies.status"
remote_env "$ADAPTER" arm ios >/dev/null
GEN=$((GEN + 1))
await_reply_result "$PARENT/state/procevent-inbox/$SID.$GEN.result" \
|| fail "the restored route was not read"
assert_grep "offset=${moved_offset}" "$PARENT/state/remote-replies/ios.cursor" \
"restoring the same bytes did not reach the same offset"
restored_hash=$(sed -n 's/^prefix_sha256=//p' "$PARENT/state/remote-replies/ios.cursor" | tr 'A-F' 'a-f')
[ "$restored_hash" = "$moved_hash" ] || fail "restoring the same bytes changed the prefix hash"
[ "$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" -eq 4 ] \
|| fail "restoring the same bytes appended a continuity break"
[ -z "$(status_open_decisions "$PARENT/state/ios.status")" ] \
|| fail "restoring the same bytes reopened the continuity decision"
break_repaired_route "retired-same" 3
RESULT_RETIRED=$RESULT_BREAK
[ "$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" -eq 5 ] \
|| fail "a continuity break after an identical restore appended nothing"
assert_grep "at offset ${moved_offset} prefix ${moved_prefix} retirements 2" "$PARENT/state/ios.status" \
"a continuity break after an identical restore did not record the new retirement count"
assert_contains "$(status_open_decisions "$PARENT/state/ios.status")" \
$'remote-reply-continuity-ios\t' \
"a continuity break after an identical restore did not reopen the decision"
remote_env "$ADAPTER" ingest ios "$RESULT_RETIRED" >/dev/null 2>&1 || true
[ "$(grep -cF 'blocked [key=remote-reply-continuity-ios]' "$PARENT/state/ios.status")" -eq 5 ] \
|| fail "a repeated read of the break after an identical restore appended again"
assert_contains "$(status_open_decisions "$PARENT/state/ios.status")" \
$'remote-reply-continuity-ios\t' \
"a repeated read of the break after an identical restore closed the decision"
pass "a continuity break after retirement and an identical restore opens the decision once"

# A directory cannot be removed by rm without -r, for any user. Retirement must
# fail and leave the count where it was.
count_before=$(cat "$PARENT/state/remote-replies/ios.retirements")
rm -f "$PARENT/state/remote-replies/ios.cursor"
mkdir "$PARENT/state/remote-replies/ios.cursor"
set +e
remote_env "$ADAPTER" retire ios > "$TMP_ROOT/retire-cursor-stuck.out" 2>&1
retire_rc=$?
set -e
[ "$retire_rc" -ne 0 ] || fail "retirement reported success when the cursor could not be removed"
assert_grep 'cannot remove remote reply cursor' "$TMP_ROOT/retire-cursor-stuck.out" \
"retirement did not report the failed cursor removal"
[ "$(cat "$PARENT/state/remote-replies/ios.retirements")" = "$count_before" ] \
|| fail "a failed cursor removal increased the retirement count"
[ -d "$PARENT/state/remote-replies/ios.cursor" ] \
|| fail "a failed cursor removal removed the cursor"
pass "a failed cursor removal does not increase the retirement count"

# The count write is the step after removal. Stopping there leaves the old count.
rm -rf "$PARENT/state/remote-replies/ios.cursor"
printf 'schema=fm-remote-reply-cursor.v1\noffset=1\nprefix_sha256=%064d\n' 0 \
> "$PARENT/state/remote-replies/ios.cursor"
RETIRE_COUNT_FAIL_BIN="$TMP_ROOT/retire-count-fail-bin"
mkdir -p "$RETIRE_COUNT_FAIL_BIN"
REAL_MKTEMP=$(command -v mktemp)
{
cat <<'SH'
#!/usr/bin/env bash
case "${1:-}" in
*/.retirements.XXXXXX) exit 73 ;;
esac
SH
printf 'exec %q "$@"\n' "$REAL_MKTEMP"
} > "$RETIRE_COUNT_FAIL_BIN/mktemp"
chmod +x "$RETIRE_COUNT_FAIL_BIN/mktemp"
set +e
PATH="$RETIRE_COUNT_FAIL_BIN:$PATH" remote_env "$ADAPTER" retire ios \
> "$TMP_ROOT/retire-count-stopped.out" 2>&1
retire_rc=$?
set -e
[ "$retire_rc" -ne 0 ] || fail "retirement reported success when the count could not be recorded"
assert_grep 'cannot record remote reply retirement' "$TMP_ROOT/retire-count-stopped.out" \
"retirement did not report the failed count write"
assert_absent "$PARENT/state/remote-replies/ios.cursor" \
"a retirement that stopped after removal left the cursor"
[ "$(cat "$PARENT/state/remote-replies/ios.retirements")" = "$count_before" ] \
|| fail "a retirement that stopped after removal changed the count"
pass "a retirement that stops after removing the cursor leaves the old count"

echo "ALL TESTS PASSED"
Loading