Files
box/bin/fleet-alert-relay.sh

337 lines
12 KiB
Bash
Executable File

#!/bin/bash
# bin/fleet-alert-relay.sh — idempotent relay: fleet-alert outbox.jsonl -> #lobby.
#
# PROPOSAL ONLY (board ticket ae28ac8b735f). No bl infra is touched by this
# branch; deployment needs opm review + sign-off. See
# docs/FLEET-ALERT-DUP-POST-GATE.md.
#
# NOTE (2026-10-06): auto-invoke from fleet-alert-check.sh is disabled by
# default (FLEET_BL_RELAY=1 re-enables). The container-side hook pages
# #lobby today; do not re-enable without retiring it first.
#
# The 2026-10-05 11:28Z incident: one RECOVERY record in the outbox became two
# identical verified #lobby posts (seq 642/643, 3.35s apart) because the relay
# leg had no idempotency: append-only outbox, no consume tracking, no content
# dedup, and the chat POST API has no idempotency key.
#
# Gates implemented here:
# 1. posted-watermark — each posted outbox record id is appended to
# posted.log; records already posted are skipped.
# 2. content-hash + 10-min TTL — sha256 of the formatted alert text in a
# file-backed seen-set (seen-hashes.log); identical text re-posts inside
# the window are dropped. Catches same-record re-posts AND identical
# text from different records.
# 3. retry discipline — never blind-retry a chat POST after a timeout /
# phantom-000. Read the #lobby tail first; re-post only if absent.
#
# Usage:
# fleet-alert-relay.sh [--dry-run] [--self-test]
#
# Env: FLEET_ALERT_DIR (default ~/.local/share/fleet-alert),
# CHAT_BASE (default https://chat.muse-dev.online),
# CHAT_IDENTITY (default operator-646), CHAT_KEYFILE (default ~/.ssh/id_frontdoor)
set -euo pipefail
ALERT_DIR="${FLEET_ALERT_DIR:-$HOME/.local/share/fleet-alert}"
OUTBOX="$ALERT_DIR/outbox.jsonl"
POSTED="$ALERT_DIR/posted.log"
SEEN="$ALERT_DIR/seen-hashes.log"
LOCKF="$ALERT_DIR/relay.lock"
CHAT_BASE="${CHAT_BASE:-https://chat.muse-dev.online}"
API="$CHAT_BASE/api/chat"
IDENTITY="${CHAT_IDENTITY:-operator-646}"
KEYFILE="${CHAT_KEYFILE:-$HOME/.ssh/id_frontdoor}"
CHANNEL="#lobby"
TTL_SECS=600
UA='Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0 Safari/537.36'
DRY_RUN=0
log() { echo "[fleet-alert-relay] $*" >&2; }
# ---------------------------------------------------------------- transport
# Factored so --self-test can stub them. Return contracts:
# transport_tail -> prints "<ts> <identity> <message>" lines for recent #lobby (one per line, tab-safe: message may contain spaces, fields split on first two spaces)
# transport_post "<text>" -> prints "OK <msg-id>" or "UNKNOWN <reason>"; exit 0 always (callers decide)
transport_tail() {
curl -s -m 15 -A "$UA" "$API/history?channel=%23lobby&limit=50" \
| python3 -c '
import json,sys
try:
d=json.load(sys.stdin)
except Exception:
sys.exit(0)
for m in d.get("messages",[]):
msg=(m.get("message") or "").replace(chr(10)," ")
print(m.get("ts"), m.get("identity"), msg)
'
}
transport_post() { # $1 = text
local text="$1" ts sig payload resp http
ts="$(date +%s)"
if ! sig="$(sign_payload "$(printf '%s\n%s\n%s' "$ts" "$CHANNEL" "$text")")"; then
echo "UNKNOWN sign-failed($KEYFILE)"; return 0
fi
payload="$(MSG="$text" TS="$ts" SIG="$sig" python3 -c '
import json,os
print(json.dumps({"channel":os.environ["CHANNEL"],"identity":os.environ["IDENTITY"],
"message":os.environ["MSG"],"ts":int(os.environ["TS"]),"signature":os.environ["SIG"]}))')"
resp="$(curl -s -m 20 -A "$UA" -X POST -H 'Content-Type: application/json' \
-d "$payload" -w '\n%{http_code}' "$API/post" 2>/dev/null || echo -e '\n000')"
http="$(printf '%s' "$resp" | tail -n1)"
resp="$(printf '%s' "$resp" | sed '$d')"
case "$http" in
2*)
local mid
mid="$(printf '%s' "$resp" | python3 -c '
import json,sys
try:
d=json.load(sys.stdin); print(d.get("id","") if d.get("ok") else "")
except Exception:
print("")')"
if [ -n "$mid" ]; then echo "OK $mid"; else echo "UNKNOWN bad-body-$http"; fi
;;
*) echo "UNKNOWN http-$http" ;;
esac
return 0
}
# file-based signing (never pipe): stdin-piped signatures intermittently
# fail verification (2026-10-02 incident).
sign_payload() { # $1 = payload -> armored signature
local tmpd pf
tmpd="$(mktemp -d)" || return 1
pf="$tmpd/payload"
printf '%s' "$1" > "$pf"
rm -f "$pf.sig"
ssh-keygen -Y sign -f "$KEYFILE" -n chat "$pf" </dev/null >/dev/null 2>&1 \
|| { rm -rf "$tmpd"; return 1; }
cat "$pf.sig"
rm -rf "$tmpd"
}
# ----------------------------------------------------------------- formatting
# Canonical text for one outbox record (JSON on stdin):
# ALERT -> [fleet-alert] CRITICAL: <node> warp PARTITION x<n>
# RECOVERY -> [fleet-alert] RECOVERED: <node> warp PARTITION
# condition is "<type>:<node>" (partition|cdp).
format_text() {
python3 -c '
import json,sys
r=json.load(sys.stdin)
cond=r.get("condition","")
typ,_,node=cond.partition(":")
label={"partition":"PARTITION","cdp":"CDP"}.get(typ,typ.upper())
kind=r.get("kind","")
if kind=="ALERT":
n=r.get("consecutive",r.get("count","?"))
print("[fleet-alert] CRITICAL: %s warp %s x%s" % (node,label,n))
elif kind=="RECOVERY":
print("[fleet-alert] RECOVERED: %s warp %s" % (node,label))
else:
print("[fleet-alert] %s: %s warp %s" % (kind,node,label))
'
}
content_hash() { printf '%s' "$1" | sha256sum | cut -d' ' -f1; }
# ------------------------------------------------------------- gate 1: watermark
already_posted() { # $1 = record id
[ -f "$POSTED" ] && grep -qxF "$1" "$POSTED"
}
mark_posted() { # $1 = record id
printf '%s\n' "$1" >> "$POSTED"
}
# ------------------------------------------------------- gate 2: content-hash TTL
seen_recently() { # $1 = hash -> 0 if seen within TTL
local h="$1" now ts
[ -f "$SEEN" ] || return 1
now="$(date +%s)"
while read -r h2 ts _rest; do
[ "$h2" = "$h" ] || continue
if [ "$(( now - ts ))" -lt "$TTL_SECS" ]; then return 0; fi
done < "$SEEN"
return 1
}
mark_seen() { # $1 = hash
printf '%s %s\n' "$1" "$(date +%s)" >> "$SEEN"
}
prune_seen() {
local now tmp
[ -f "$SEEN" ] || return 0
now="$(date +%s)"; tmp="$(mktemp)"
while read -r h ts _rest; do
[ -n "$h" ] && [ "$(( now - ts ))" -lt "$TTL_SECS" ] && printf '%s %s\n' "$h" "$ts" >> "$tmp"
done < "$SEEN"
mv "$tmp" "$SEEN"
}
# ------------------------------------------------- gate 3: check-before-(re)post
# 0 if the exact text appears in the recent #lobby tail.
lobby_has_text() { # $1 = text
local want="$1"
# transport_tail prints "<ts> <identity> <message>" per line; the message
# itself may contain spaces, so strip only the first two fields.
transport_tail | python3 -c '
import sys
want=sys.argv[1]
for line in sys.stdin:
parts=line.rstrip("\n").split(" ",2)
if len(parts)==3 and parts[2]==want:
sys.exit(0)
sys.exit(1)' "$want"
}
relay_record() { # $1 = record id, $2 = formatted text
local rid="$1" text="$2" h out
h="$(content_hash "$text")"
# Gate 1: posted watermark
if already_posted "$rid"; then
log "skip $rid: already in posted.log"
return 0
fi
# Gate 2: content hash within TTL
if seen_recently "$h"; then
log "skip $rid: identical text posted <10min ago (hash $h)"
mark_posted "$rid" # consume it so it never retries later
return 0
fi
# Gate 3a: somebody else already posted it (covers agent-driven relay overlap)
if lobby_has_text "$text"; then
log "skip $rid: text already in #lobby tail"
mark_posted "$rid"; mark_seen "$h"
return 0
fi
if [ "$DRY_RUN" = 1 ]; then
log "dry-run: would post $rid: $text"
return 0
fi
# Attempt one POST. Any ambiguous outcome -> verify via tail, never blind-retry.
out="$(transport_post "$text")"
case "$out" in
OK*)
log "posted $rid -> ${out#OK }"
mark_posted "$rid"; mark_seen "$h"
return 0
;;
UNKNOWN*)
log "post $rid outcome UNKNOWN (${out#UNKNOWN }); checking #lobby tail before any retry"
sleep 2
if lobby_has_text "$text"; then
log "post $rid landed despite ${out#UNKNOWN } — marking posted, no retry"
mark_posted "$rid"; mark_seen "$h"
return 0
fi
log "post $rid absent from tail — ONE retry only"
out="$(transport_post "$text")"
case "$out" in
OK*)
log "retry posted $rid -> ${out#OK }"
mark_posted "$rid"; mark_seen "$h"
return 0
;;
*)
log "ERROR: post $rid still ${out%% *} after one retry; leaving unposted for next run (tail-check will catch it if it landed)"
return 1
;;
esac
;;
esac
}
main() {
local rid rec text fails=0
[ -f "$OUTBOX" ] || { log "no outbox ($OUTBOX); nothing to do"; return 0; }
mkdir -p "$ALERT_DIR"
touch "$POSTED" "$SEEN"
prune_seen
while IFS= read -r rec; do
[ -n "$rec" ] || continue
rid="$(printf '%s' "$rec" | python3 -c 'import json,sys; print(json.load(sys.stdin).get("id",""))')"
[ -n "$rid" ] || { log "skip record with no id"; continue; }
text="$(printf '%s' "$rec" | format_text)"
relay_record "$rid" "$text" || fails=$((fails+1))
done < "$OUTBOX"
return "$fails"
}
# ------------------------------------------------------------------ self-test
# Acceptance: two identical alert submissions <10 min apart -> exactly one
# #lobby post. Stubs the transport; exercises all three gates.
# NOTE: transport_post is invoked via command substitution (subshell), so the
# stub counts calls with a file, not a variable.
self_test() {
local td calls lobby ok=1 n sk sig_out old_key
td="$(mktemp -d)"; export FLEET_ALERT_DIR="$td"
ALERT_DIR="$td"; OUTBOX="$td/outbox.jsonl"; POSTED="$td/posted.log"
SEEN="$td/seen-hashes.log"; LOCKF="$td/relay.lock"
calls="$td/calls.log"; lobby="$td/lobby.log"
touch "$calls" "$lobby"
# sign_payload must round-trip with a valid key and fail cleanly without
# one (2026-10-06: missing ~/.ssh/id_frontdoor broke every #lobby post
# with an undiagnosable bare "sign-failed").
old_key="$KEYFILE"
sk="$td/signkey"
ssh-keygen -t ed25519 -f "$sk" -N '' -q >/dev/null 2>&1 \
|| { echo "FAIL: cannot generate ephemeral test key"; ok=0; }
if KEYFILE="$sk" sig_out="$(sign_payload "self-test")"; then
case "$sig_out" in
*"BEGIN SSH SIGNATURE"*) : ;;
*) echo "FAIL: sign_payload output not armored"; ok=0 ;;
esac
else
echo "FAIL: sign_payload failed with a valid key"; ok=0
fi
if KEYFILE="$td/no-such-key" sign_payload "self-test" >/dev/null 2>&1; then
echo "FAIL: sign_payload succeeded with a missing key"; ok=0
fi
KEYFILE="$old_key"
# Two identical submissions: same text, different record ids (the 11:28Z shape)
printf '%s\n' \
'{"id":"rec-A","ts":1791199616,"kind":"RECOVERY","condition":"partition:def"}' \
'{"id":"rec-B","ts":1791199617,"kind":"RECOVERY","condition":"partition:def"}' \
> "$OUTBOX"
transport_post() { # stub: 1st call "times out" but the server DID accept it
echo "call" >> "$calls"
printf '%s\n' "$1" >> "$lobby"
n=$(wc -l < "$calls")
if [ "$n" = 1 ]; then echo "UNKNOWN phantom-000"; else echo "OK stub-$n"; fi
return 0
}
transport_tail() { # stub: the lobby shows whatever "landed"
while IFS= read -r m; do printf '%s opm %s\n' "$(date +%s)" "$m"; done < "$lobby"
}
main
n=$(wc -l < "$calls")
[ "$n" -eq 1 ] || { echo "FAIL: expected exactly 1 post call, got $n"; ok=0; }
grep -qx "rec-A" "$POSTED" || { echo "FAIL: rec-A not watermarked"; ok=0; }
grep -qx "rec-B" "$POSTED" || { echo "FAIL: rec-B not watermarked (gate 2 should consume it)"; ok=0; }
# Identical re-submission (new record id, same text) must not post again
printf '%s\n' '{"id":"rec-C","ts":1791199700,"kind":"RECOVERY","condition":"partition:def"}' >> "$OUTBOX"
main
n=$(wc -l < "$calls")
[ "$n" -eq 1 ] || { echo "FAIL: identical re-submission posted again (calls=$n)"; ok=0; }
rm -rf "$td"
[ "$ok" = 1 ] && echo "SELF-TEST PASS: 2 identical submissions -> 1 post call; re-submission -> 0 new calls" && return 0
return 1
}
case "${1:-}" in
--dry-run) DRY_RUN=1; shift ;;
--self-test) self_test; exit $? ;;
esac
# Single-flight the relay (defense in depth; the detector's own flock is separate).
exec 9>"$LOCKF"
if ! flock -n 9; then
log "another relay run holds the lock; exiting"
exit 0
fi
main