feat(wake): W2 — three-cursor durable store + RECEIVED/CONSUMED ack-wrapper + watch-list schema (#903)
Co-authored-by: jason.woltje <jason@diversecanvas.com> Co-committed-by: jason.woltje <jason@diversecanvas.com>
This commit was merged in pull request #903.
This commit is contained in:
279
packages/mosaic/framework/tools/wake/ack.sh
Executable file
279
packages/mosaic/framework/tools/wake/ack.sh
Executable file
@@ -0,0 +1,279 @@
|
||||
#!/usr/bin/env bash
|
||||
# ack.sh — A4 of the wake canon (EPIC #892, W2): the RECEIVED / CONSUMED
|
||||
# ack-wrapper.
|
||||
#
|
||||
# CONTRACT ANCHORS (docs/scratchpads/heartbeat-planning/CONVERGED-DESIGN.md):
|
||||
# §2.2 ack protocol: RECEIVED vs CONSUMED; local-write + async ship.
|
||||
# §1.2 three-cursor: consumed_seq advances ONLY on a consumer CONSUMED ack.
|
||||
#
|
||||
# RECEIVED — delivery happened (paste landed). `wake_id` dedups DELIVERY ONLY:
|
||||
# a duplicate delivery is a re-RECEIVE, NEVER a re-action.
|
||||
# CONSUMED N — emitted ONLY after durable capture OR no-op disposition of a
|
||||
# CONTIGUOUS prefix <=N. first-action-before-capture is FORBIDDEN
|
||||
# (caller discipline: capture BEFORE invoking this). Cannot ack N
|
||||
# while N-1 is unconsumed (the store enforces the gapless prefix).
|
||||
# Acks are CUMULATIVE: CONSUMED N implies all <=N.
|
||||
#
|
||||
# LOCAL-WRITE-ONLY: the ack path writes the XDG ledger and advances the cursor
|
||||
# with NO network call. A BACKGROUND sync ships the ack (WAKE_ACK_SYNC_CMD).
|
||||
# Ack-path outage => a single re-wake after timeout then fall back to cadence
|
||||
# (no retry spam) — the background sync runs ONCE and records an outage marker;
|
||||
# it never loops.
|
||||
#
|
||||
# The `embed` subcommand prints the one copy-run line meant to be EMBEDDED in
|
||||
# each digest (W3 renders it into the digest body).
|
||||
#
|
||||
# Operator-agnostic: ledger via XDG/env only; no operator paths/names/secrets.
|
||||
set -uo pipefail
|
||||
|
||||
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
|
||||
# shellcheck source=./_wake-common.sh disable=SC1091
|
||||
. "$SCRIPT_DIR/_wake-common.sh"
|
||||
|
||||
STATE_DIR="$(wake_state_dir)"
|
||||
STORE_SH="$SCRIPT_DIR/store.sh"
|
||||
|
||||
_need_jq() {
|
||||
command -v jq >/dev/null 2>&1 || {
|
||||
echo "ack.sh: jq is required" >&2
|
||||
exit 3
|
||||
}
|
||||
}
|
||||
|
||||
usage() {
|
||||
cat >&2 <<'EOF'
|
||||
Usage: ack.sh <command> [options]
|
||||
|
||||
Commands:
|
||||
received --wake-id ID Record a RECEIVED ack (delivery). The
|
||||
wake_id DEDUPS DELIVERY only: a repeat
|
||||
is a re-RECEIVE (prints DUP), never a
|
||||
re-action. Local-write only.
|
||||
consumed --upto N [--wake-id ID] Record a CONSUMED ack for the contiguous
|
||||
prefix <=N and advance consumed_seq.
|
||||
Local-write only; a background sync
|
||||
ships it (never blocks on network).
|
||||
--no-sync suppresses the background ship.
|
||||
embed --upto N [--wake-id ID] Print the copy-run ack line to EMBED in
|
||||
a digest (does not perform the ack).
|
||||
status Print ledger tail + cursor state.
|
||||
|
||||
Environment:
|
||||
WAKE_STATE_HOME override base state dir (XDG by default).
|
||||
WAKE_AGENT per-agent queue namespace (default: default).
|
||||
WAKE_ACK_SYNC_CMD command run in the BACKGROUND to ship an ack. Receives the
|
||||
ack JSON on stdin. Never invoked on the synchronous path.
|
||||
WAKE_ACK_CLI override the embedded copy-run invocation prefix (default:
|
||||
the resolved path to this script).
|
||||
EOF
|
||||
}
|
||||
|
||||
# _ledger_append JSON — append one ack record to the local ledger, atomically.
|
||||
# §2.2: LOCAL-WRITE-ONLY. No network here.
|
||||
_ledger_append() {
|
||||
local record="$1"
|
||||
_wake_init_dir "$STATE_DIR"
|
||||
{
|
||||
cat "$STATE_DIR/ack-ledger.jsonl" 2>/dev/null
|
||||
printf '%s\n' "$record"
|
||||
} | grep -v '^[[:space:]]*$' | _atomic_write "$STATE_DIR/ack-ledger.jsonl"
|
||||
}
|
||||
|
||||
# _wake_id_seen ID — true if a RECEIVED record for this wake_id already exists.
|
||||
# Count-based (jq -s slurp): an EMPTY ledger slurps to [] => length 0. (jq -e
|
||||
# over a zero-input file returns exit 0, which would false-positive here — so
|
||||
# never rely on -e exit semantics for presence.)
|
||||
_wake_id_seen() {
|
||||
local id="$1" n
|
||||
[ -f "$STATE_DIR/ack-ledger.jsonl" ] || return 1
|
||||
n="$(jq -s --arg id "$id" \
|
||||
'[.[] | select(.type=="RECEIVED" and .wake_id==$id)] | length' \
|
||||
"$STATE_DIR/ack-ledger.jsonl" 2>/dev/null || echo 0)"
|
||||
[ "${n:-0}" -gt 0 ]
|
||||
}
|
||||
|
||||
# _ship_async JSON — fire the background sync ONCE, fully detached, so the ack
|
||||
# path returns immediately and NEVER blocks on the network (§2.2). No retry
|
||||
# loop: a failure records an outage marker and stops (single re-wake + cadence
|
||||
# fallback is the escalation policy, handled off this path).
|
||||
_ship_async() {
|
||||
local record="$1"
|
||||
[ -n "${WAKE_ACK_SYNC_CMD:-}" ] || return 0
|
||||
local marker="$STATE_DIR/ack-sync.state"
|
||||
(
|
||||
if printf '%s\n' "$record" | sh -c "$WAKE_ACK_SYNC_CMD" >/dev/null 2>&1; then
|
||||
printf 'shipped %s\n' "$(date +%s)" >"$marker" 2>/dev/null || true
|
||||
else
|
||||
printf 'outage %s\n' "$(date +%s)" >"$marker" 2>/dev/null || true
|
||||
fi
|
||||
) </dev/null >/dev/null 2>&1 &
|
||||
# Detach so no shell job-control state ties the ack turn to the ship.
|
||||
disown 2>/dev/null || true
|
||||
}
|
||||
|
||||
cmd_received() {
|
||||
_need_jq
|
||||
local wake_id=''
|
||||
while [ $# -gt 0 ]; do
|
||||
case "$1" in
|
||||
--wake-id)
|
||||
wake_id="${2:-}"
|
||||
shift 2
|
||||
;;
|
||||
*)
|
||||
echo "ack.sh received: unknown option '$1'" >&2
|
||||
exit 2
|
||||
;;
|
||||
esac
|
||||
done
|
||||
[ -n "$wake_id" ] || {
|
||||
echo "ack.sh received: --wake-id is required" >&2
|
||||
exit 2
|
||||
}
|
||||
|
||||
_wake_init_dir "$STATE_DIR"
|
||||
|
||||
# wake_id dedups DELIVERY only. A duplicate delivery is a re-RECEIVE — logged,
|
||||
# but flagged dup so no consumer re-action is triggered. It NEVER advances
|
||||
# consumed_seq (delivery != consumption).
|
||||
local dup="false" verb="RECEIVED"
|
||||
if _wake_id_seen "$wake_id"; then
|
||||
dup="true"
|
||||
verb="DUP"
|
||||
fi
|
||||
|
||||
local record
|
||||
record="$(jq -cn \
|
||||
--arg id "$wake_id" \
|
||||
--argjson dup "$dup" \
|
||||
--argjson ts "$(date +%s)" \
|
||||
'{type:"RECEIVED", wake_id:$id, dup:$dup, ts:$ts}')"
|
||||
_ledger_append "$record"
|
||||
echo "$verb"
|
||||
}
|
||||
|
||||
cmd_consumed() {
|
||||
_need_jq
|
||||
local upto='' wake_id='' do_sync="1"
|
||||
while [ $# -gt 0 ]; do
|
||||
case "$1" in
|
||||
--upto)
|
||||
upto="${2:-}"
|
||||
shift 2
|
||||
;;
|
||||
--wake-id)
|
||||
wake_id="${2:-}"
|
||||
shift 2
|
||||
;;
|
||||
--no-sync)
|
||||
do_sync="0"
|
||||
shift
|
||||
;;
|
||||
*)
|
||||
echo "ack.sh consumed: unknown option '$1'" >&2
|
||||
exit 2
|
||||
;;
|
||||
esac
|
||||
done
|
||||
case "$upto" in
|
||||
'' | *[!0-9]*)
|
||||
echo "ack.sh consumed: --upto must be a non-negative integer" >&2
|
||||
exit 2
|
||||
;;
|
||||
esac
|
||||
|
||||
# Advance consumed_seq via the store. The store enforces the CONTIGUOUS
|
||||
# gapless-prefix rule and rejects a gap (cannot ack N while N-1 unconsumed).
|
||||
# This is a LOCAL-WRITE cursor advance — no network.
|
||||
local new_cursor
|
||||
if ! new_cursor="$("$STORE_SH" consume --upto "$upto" 2>&1)"; then
|
||||
echo "ack.sh consumed: refused — $new_cursor" >&2
|
||||
exit 1
|
||||
fi
|
||||
|
||||
# Record the CONSUMED ack in the local ledger (still no network).
|
||||
local record
|
||||
record="$(jq -cn \
|
||||
--argjson upto "$upto" \
|
||||
--arg id "$wake_id" \
|
||||
--argjson ts "$(date +%s)" \
|
||||
'{type:"CONSUMED", upto:$upto, wake_id:$id, ts:$ts}')"
|
||||
_ledger_append "$record"
|
||||
|
||||
# Only now, OFF the ack path, ship it in the background. The synchronous
|
||||
# portion is already complete and durable.
|
||||
if [ "$do_sync" = "1" ]; then
|
||||
_ship_async "$record"
|
||||
fi
|
||||
|
||||
echo "CONSUMED $new_cursor"
|
||||
}
|
||||
|
||||
cmd_embed() {
|
||||
local upto='' wake_id=''
|
||||
while [ $# -gt 0 ]; do
|
||||
case "$1" in
|
||||
--upto)
|
||||
upto="${2:-}"
|
||||
shift 2
|
||||
;;
|
||||
--wake-id)
|
||||
wake_id="${2:-}"
|
||||
shift 2
|
||||
;;
|
||||
*)
|
||||
echo "ack.sh embed: unknown option '$1'" >&2
|
||||
exit 2
|
||||
;;
|
||||
esac
|
||||
done
|
||||
case "$upto" in
|
||||
'' | *[!0-9]*)
|
||||
echo "ack.sh embed: --upto must be a non-negative integer" >&2
|
||||
exit 2
|
||||
;;
|
||||
esac
|
||||
# The one copy-run line a digest embeds. WAKE_ACK_CLI lets an installer point
|
||||
# it at the deployed path; default is this script's resolved path.
|
||||
local cli="${WAKE_ACK_CLI:-$SCRIPT_DIR/ack.sh}"
|
||||
if [ -n "$wake_id" ]; then
|
||||
printf '%s consumed --upto %s --wake-id %s\n' "$cli" "$upto" "$wake_id"
|
||||
else
|
||||
printf '%s consumed --upto %s\n' "$cli" "$upto"
|
||||
fi
|
||||
}
|
||||
|
||||
cmd_status() {
|
||||
_wake_init_dir "$STATE_DIR"
|
||||
echo "# cursors"
|
||||
"$STORE_SH" cursors
|
||||
echo "# ack-ledger (tail)"
|
||||
tail -n 10 "$STATE_DIR/ack-ledger.jsonl" 2>/dev/null || true
|
||||
if [ -f "$STATE_DIR/ack-sync.state" ]; then
|
||||
echo "# sync"
|
||||
cat "$STATE_DIR/ack-sync.state"
|
||||
fi
|
||||
}
|
||||
|
||||
main() {
|
||||
[ $# -ge 1 ] || {
|
||||
usage
|
||||
exit 2
|
||||
}
|
||||
local cmd="$1"
|
||||
shift
|
||||
case "$cmd" in
|
||||
received) cmd_received "$@" ;;
|
||||
consumed) cmd_consumed "$@" ;;
|
||||
embed) cmd_embed "$@" ;;
|
||||
status) cmd_status "$@" ;;
|
||||
-h | --help | help) usage ;;
|
||||
*)
|
||||
echo "ack.sh: unknown command '$cmd'" >&2
|
||||
usage
|
||||
exit 2
|
||||
;;
|
||||
esac
|
||||
}
|
||||
|
||||
main "$@"
|
||||
Reference in New Issue
Block a user