From 55e3e518b9aab7b08d26f908da4918c88ebdfc7d Mon Sep 17 00:00:00 2001 From: Cyril Poder Date: Thu, 24 Sep 2026 02:23:29 +0200 Subject: [PATCH] detect: an absence is raised when its deadline passes, and after a restart (varpulis #287) "An order not acknowledged within 4h" waited for the next event that reached the pattern: on subjects that went quiet after the order, it was never raised. And an absence the unit was waiting out when it stopped lost what its negated step waited for in the snapshot, so after the restart the deadline dropped it without an alert. Engine pinned at varpulis 21f3ba4 (#287). e2e/detect D8 (silent subjects) and D9 (kill -9 with the order in the snapshot) fail on v0.3.3 and pass now. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01N3K1TGnWTvwYzt9rKuJXES --- CHANGELOG.md | 15 ++++++ core/Cargo.lock | 20 ++++---- core/Cargo.toml | 2 +- docs/book/src/concepts/detects.md | 6 +++ e2e/detect/run.sh | 77 +++++++++++++++++++++++++++++-- 5 files changed, 106 insertions(+), 14 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0c8926d..3564173 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,21 @@ is `0`, minor versions may carry breaking changes — they are called out here. ## [Unreleased] +### Fixed — an absence is raised when its deadline passes (varpulis #287) +- **"An order not acknowledged within 4h" (`-> NOT Ack ... within 4h`) is + raised when the four hours are over**, by the first event that takes the + event time of the pattern's types past them or, when nothing more comes at + all, about a grace (`VEJAS_IDLE_CLOSE_SECS`) after them. It used to wait + for the next event that reached the pattern itself: on subjects that went + quiet after the order, never. `e2e/detect` D8 covers it: an order, an + acknowledged order, then silence; the previous runtime raises nothing. +- **An absence the unit was waiting out when it stopped is still raised + after the restart.** The snapshot kept the open sequence but not what its + negated step waited for, so after a restore the acknowledgement no longer + cancelled it and the deadline dropped it without an alert. `e2e/detect` + D9 covers it: an order, kill -9 with the order in the snapshot, restart, + silence; the previous runtime never raises it. + ## [0.3.3] — 2026-09-24 ### Fixed — a window on a quiet source closes diff --git a/core/Cargo.lock b/core/Cargo.lock index 7fc948c..be39ce5 100644 --- a/core/Cargo.lock +++ b/core/Cargo.lock @@ -1846,7 +1846,7 @@ checksum = "b6c140620e7ffbb22c2dee59cafe6084a59b5ffc27a8859a5f0d494b5d52b6be" [[package]] name = "varpulis-core" version = "0.11.0" -source = "git+https://github.com/varpulis/varpulis?rev=92f07b966709bacc042ed97a71e16e8dd625c84d#92f07b966709bacc042ed97a71e16e8dd625c84d" +source = "git+https://github.com/varpulis/varpulis?rev=21f3ba49220615a926b1d96c16075fe11c67e3b9#21f3ba49220615a926b1d96c16075fe11c67e3b9" dependencies = [ "chrono", "indexmap", @@ -1863,7 +1863,7 @@ dependencies = [ [[package]] name = "varpulis-dead-letter" version = "0.11.0" -source = "git+https://github.com/varpulis/varpulis?rev=92f07b966709bacc042ed97a71e16e8dd625c84d#92f07b966709bacc042ed97a71e16e8dd625c84d" +source = "git+https://github.com/varpulis/varpulis?rev=21f3ba49220615a926b1d96c16075fe11c67e3b9#21f3ba49220615a926b1d96c16075fe11c67e3b9" dependencies = [ "chrono", "serde", @@ -1875,7 +1875,7 @@ dependencies = [ [[package]] name = "varpulis-engine" version = "0.11.0" -source = "git+https://github.com/varpulis/varpulis?rev=92f07b966709bacc042ed97a71e16e8dd625c84d#92f07b966709bacc042ed97a71e16e8dd625c84d" +source = "git+https://github.com/varpulis/varpulis?rev=21f3ba49220615a926b1d96c16075fe11c67e3b9#21f3ba49220615a926b1d96c16075fe11c67e3b9" dependencies = [ "serde_json", "thiserror", @@ -1887,7 +1887,7 @@ dependencies = [ [[package]] name = "varpulis-hamlet" version = "0.11.0" -source = "git+https://github.com/varpulis/varpulis?rev=92f07b966709bacc042ed97a71e16e8dd625c84d#92f07b966709bacc042ed97a71e16e8dd625c84d" +source = "git+https://github.com/varpulis/varpulis?rev=21f3ba49220615a926b1d96c16075fe11c67e3b9#21f3ba49220615a926b1d96c16075fe11c67e3b9" dependencies = [ "rustc-hash", "smallvec", @@ -1898,7 +1898,7 @@ dependencies = [ [[package]] name = "varpulis-parser" version = "0.11.0" -source = "git+https://github.com/varpulis/varpulis?rev=92f07b966709bacc042ed97a71e16e8dd625c84d#92f07b966709bacc042ed97a71e16e8dd625c84d" +source = "git+https://github.com/varpulis/varpulis?rev=21f3ba49220615a926b1d96c16075fe11c67e3b9#21f3ba49220615a926b1d96c16075fe11c67e3b9" dependencies = [ "logos", "miette", @@ -1911,7 +1911,7 @@ dependencies = [ [[package]] name = "varpulis-pst" version = "0.11.0" -source = "git+https://github.com/varpulis/varpulis?rev=92f07b966709bacc042ed97a71e16e8dd625c84d#92f07b966709bacc042ed97a71e16e8dd625c84d" +source = "git+https://github.com/varpulis/varpulis?rev=21f3ba49220615a926b1d96c16075fe11c67e3b9#21f3ba49220615a926b1d96c16075fe11c67e3b9" dependencies = [ "hashlink", "rustc-hash", @@ -1921,7 +1921,7 @@ dependencies = [ [[package]] name = "varpulis-runtime" version = "0.11.0" -source = "git+https://github.com/varpulis/varpulis?rev=92f07b966709bacc042ed97a71e16e8dd625c84d#92f07b966709bacc042ed97a71e16e8dd625c84d" +source = "git+https://github.com/varpulis/varpulis?rev=21f3ba49220615a926b1d96c16075fe11c67e3b9#21f3ba49220615a926b1d96c16075fe11c67e3b9" dependencies = [ "chrono", "hashlink", @@ -1949,7 +1949,7 @@ dependencies = [ [[package]] name = "varpulis-sase" version = "0.11.0" -source = "git+https://github.com/varpulis/varpulis?rev=92f07b966709bacc042ed97a71e16e8dd625c84d#92f07b966709bacc042ed97a71e16e8dd625c84d" +source = "git+https://github.com/varpulis/varpulis?rev=21f3ba49220615a926b1d96c16075fe11c67e3b9#21f3ba49220615a926b1d96c16075fe11c67e3b9" dependencies = [ "chrono", "rustc-hash", @@ -1961,7 +1961,7 @@ dependencies = [ [[package]] name = "varpulis-simd" version = "0.11.0" -source = "git+https://github.com/varpulis/varpulis?rev=92f07b966709bacc042ed97a71e16e8dd625c84d#92f07b966709bacc042ed97a71e16e8dd625c84d" +source = "git+https://github.com/varpulis/varpulis?rev=21f3ba49220615a926b1d96c16075fe11c67e3b9#21f3ba49220615a926b1d96c16075fe11c67e3b9" dependencies = [ "varpulis-core", ] @@ -1969,7 +1969,7 @@ dependencies = [ [[package]] name = "varpulis-zdd" version = "0.11.0" -source = "git+https://github.com/varpulis/varpulis?rev=92f07b966709bacc042ed97a71e16e8dd625c84d#92f07b966709bacc042ed97a71e16e8dd625c84d" +source = "git+https://github.com/varpulis/varpulis?rev=21f3ba49220615a926b1d96c16075fe11c67e3b9#21f3ba49220615a926b1d96c16075fe11c67e3b9" dependencies = [ "rustc-hash", ] diff --git a/core/Cargo.toml b/core/Cargo.toml index 358f024..edd68c2 100644 --- a/core/Cargo.toml +++ b/core/Cargo.toml @@ -22,7 +22,7 @@ webpki-roots = "0.26" # The Varpulis CEP engine, embedded as a library (ADR-0031): compile a VPL # program, feed it events, publish its emits. No async runtime in its tree — # its own CI (scripts/check-engine-deps.py) fails if one ever appears. -varpulis-engine = { git = "https://github.com/varpulis/varpulis", rev = "92f07b966709bacc042ed97a71e16e8dd625c84d" } +varpulis-engine = { git = "https://github.com/varpulis/varpulis", rev = "21f3ba49220615a926b1d96c16075fe11c67e3b9" } [[bin]] name = "vejas-runtime" diff --git a/docs/book/src/concepts/detects.md b/docs/book/src/concepts/detects.md index 7beca2d..3221011 100644 --- a/docs/book/src/concepts/detects.md +++ b/docs/book/src/concepts/detects.md @@ -61,6 +61,12 @@ machine it connected to: one alert, on `vx.alerts.lateral`. window ends even if the source never speaks again; the unit also checks while it has nothing to read. A stream whose events arrive late declares `.watermark(out_of_order: 30s)` to hold its windows that much longer. +- **An absence is raised when its deadline passes**, on the same clocks: + "an order not acknowledged within 4h" (`-> NOT Ack ... within 4h`) is + raised by the first event that takes the time of the pattern's types past + the four hours or, when nothing more comes at all, about a grace after + them. An absence the unit was waiting out when it stopped is in its + snapshot, and is still raised after the restart. - **Types.** A string `event_type` in the payload names the event type; without it, the type is the one the `.from()` binding declares for that subject. The engine's own decoder does this, the same one every Varpulis diff --git a/e2e/detect/run.sh b/e2e/detect/run.sh index 39971a1..6b8470b 100755 --- a/e2e/detect/run.sh +++ b/e2e/detect/run.sh @@ -21,6 +21,10 @@ # D7 quiet source a count on a source that goes silent still closes: past # the idle grace (VEJAS_IDLE_CLOSE_SECS, 1 s here) the # source's event time moves on with the wall clock +# D8 absence "no acknowledgement within 5s" is raised on subjects that +# go silent after the order, and not for an order acked +# D9 absence, kill an absence the unit was waiting out when killed -9 is +# still raised after the restart (from the snapshot) # D3 vpl-check the engine's verdict on the CLI: ok / refused, with exit codes # D4 topology /topology lists the units under "detects", lang vpl, running # @@ -145,6 +149,33 @@ stream Brute = Failed .to(Bus, topic: "vxt.vpnbrute") VPL +cat > "$ROOT/detects/unacked.vpl" <<'VPL' +connector Bus = nats ( + url: "ignored-by-vejas: the bus is NATS_URL" +) + +event Order: + id: str +event Ack: + id: str + +stream Orders = Order + .from(Bus, topic: "vxt.shop.order") + +stream Acks = Ack + .from(Bus, topic: "vxt.shop.ack") + +pattern Unacked = + Order as o + -> NOT Ack where id == o.id + within 5s + partition by id + +stream Late = Unacked + .emit(rule: "unacked", order: o.id) + .to(Bus, topic: "vxt.unacked") +VPL + start_nats() { "$NATSD" -js -sd "$WORK/js" -a 127.0.0.1 -p "$NATS_P" > /dev/null 2>&1 & NATS_PID=$!; sleep 0.5; } start_rt() { VEJAS_ROOT="$ROOT" NATS_URL="$URL" VEJAS_STREAM=TTEST VEJAS_SUBJECT_ROOT=vxt \ @@ -185,14 +216,14 @@ echo "── D4 topology" deadline=$((SECONDS+10)) while [ "$SECONDS" -lt "$deadline" ]; do topo=$(curl -s "http://127.0.0.1:$HTTP_P/topology") - [ "$(printf '%s' "$topo" | grep -o '"status":"running"' | wc -l)" -ge 4 ] && break + [ "$(printf '%s' "$topo" | grep -o '"status":"running"' | wc -l)" -ge 5 ] && break sleep 0.2 done -python3 - "$topo" <<'PY' && ok "four detects listed, lang vpl, running" || bad "topology: $topo" +python3 - "$topo" <<'PY' && ok "five detects listed, lang vpl, running" || bad "topology: $topo" import json, sys t=json.loads(sys.argv[1]); d=t.get("detects", []) names=sorted(x["name"] for x in d) -assert names==["detect:bruteforce","detect:lateral","detect:quiet","detect:threshold"], names +assert names==["detect:bruteforce","detect:lateral","detect:quiet","detect:threshold","detect:unacked"], names assert all(x["lang"]=="vpl" for x in d), d assert all(x["status"]=="running" for x in d), [x["status"] for x in d] PY @@ -287,6 +318,46 @@ assert len(alerts)==1 and alerts[0]["ip"]=="10.0.0.88" and alerts[0]["n"]==3, al PY kill "$SUB_PID" 2>/dev/null; pkill -f "[s]ub vxt.vpnbrute" 2>/dev/null +echo "── D8 an absence on subjects that go silent is raised (idle grace)" +( timeout 90 "$NATS" -s "$URL" sub vxt.unacked --raw > "$WORK/unacked.txt" 2>/dev/null ) & SUB_PID=$! +sleep 0.5 +shop() { #