feat(alert): wire idempotent fleet-alert relay into check loop and stop stale chrome scopes
This commit is contained in:
@@ -163,6 +163,9 @@ sleep 3
|
||||
rm -f "/home/super/.local/share/chrome-box/profiles/${PROFILE}/home/.config/chromium/SingletonLock" \
|
||||
"/home/super/.local/share/chrome-box/profiles/${PROFILE}/home/.config/chromium/SingletonSocket" \
|
||||
"/home/super/.local/share/chrome-box/profiles/${PROFILE}/home/.config/chromium/SingletonCookie" 2>/dev/null || true
|
||||
for _sc in $(systemctl --user list-units --type=scope --plain --no-legend 2>/dev/null | awk '{print $1}' | grep -E "^netvm-chrome-${PROFILE}-"); do
|
||||
systemctl --user stop "$_sc" 2>/dev/null || true
|
||||
done
|
||||
|
||||
# Relaunch in its own systemd scope so it survives this oneshot run.
|
||||
# nohup/setsid do NOT escape: this timer's service uses KillMode=control-group
|
||||
|
||||
@@ -263,3 +263,8 @@ rm -f "$STATE_DIR/.alerts.tmp"
|
||||
|
||||
tail -500 "$LOG" > "$LOG.tmp" 2>/dev/null && mv "$LOG.tmp" "$LOG"
|
||||
log "check complete"
|
||||
|
||||
# Relay pending outbox records to #lobby with idempotency gates (posted watermark + content hash TTL)
|
||||
if [ "$DRY_RUN" -eq 0 ] && [ -x "$BIN/fleet-alert-relay.sh" ]; then
|
||||
"$BIN/fleet-alert-relay.sh" >> "$LOG" 2>&1 || true
|
||||
fi
|
||||
|
||||
Executable
+313
@@ -0,0 +1,313 @@
|
||||
#!/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.
|
||||
#
|
||||
# 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"; 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
|
||||
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"
|
||||
# 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
|
||||
Reference in New Issue
Block a user