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() { #