diff --git a/docs/phase-8-processing-design.md b/docs/phase-8-processing-design.md index d75d713..ec12a31 100644 --- a/docs/phase-8-processing-design.md +++ b/docs/phase-8-processing-design.md @@ -311,6 +311,95 @@ only thing standing between a bad rule and every host at once. That makes open question 1 below considerably less optional than it looked when it was written. +## Decision 5: what `aggregate_count` emits + +Deferred twice because it is the only action whose output is not simply +the input with edits. It emits a record that never existed. + +**Shape: the window's first record, tagged.** The codebase already has a +convention for synthetic records and this follows it rather than +inventing a second one — heartbeat and host-metrics records are ordinary +`LogRecord`s distinguished by a `cairnobs.heartbeat` / +`cairnobs.metrics` attribute, with `service` left as the agent's real +service. So an aggregate is the first record of its window, unchanged, +plus: + +| attribute | value | +|---|---| +| `cairnobs.aggregated` | `"true"` | +| `cairnobs.count` | number of records collapsed, as a string | +| `cairnobs.window_start_unix_nano` | timestamp of the first record | +| `cairnobs.window_last_unix_nano` | timestamp of the last contributing record | + +Keeping the first record intact means a human reading the line sees a +real example of what was collapsed, not a summary someone invented. + +**It tags even when the count is one.** Emitting a bare record for a +window that happened to see one event would be tidier and is wrong: it +makes `cairnobs.count` present only sometimes, so the correct way to +count aggregated data silently breaks on quiet windows. Uniformity beats +tidiness here. + +**Window bounds are observed, not nominal.** `window_last` is the last +record that actually contributed, never `window_start + window_ms`. A +window flushed early must not claim an end that never happened. + +### When it emits, and the problem that hides here + +Windows are measured on record time (see `/processing/README.md`), which +means a window can only be *closed* by a later record arriving. That is +fine for `suppress_duplicates`, which emits the first record immediately +and drops the rest — nothing is ever pending. + +`aggregate_count` holds state. If the matching stream goes quiet, the +pending aggregate has nothing to close it and sits unemitted, possibly +for hours. That is data loss dressed as latency, and it is the real +reason this action was harder to specify than the other nine. + +So emission has two triggers: + +1. **A later matching record** with a timestamp at or past + `window_start + window_ms`. The pending aggregate is emitted first, + then that record opens the next window. +2. **End of stream** — agent shutdown, and in production the batch flush + interval. Anything pending is emitted. + +**The cost, stated plainly:** trigger 2 is wall-clock in production, +which means `window_ms` is a *maximum*, not a guarantee, and one logical +burst can produce more than one aggregate record if a flush lands in the +middle. Consumers must treat aggregates as additive — which they already +must, since a burst can span windows anyway. + +The conformance corpus defines trigger 2 as an implicit flush after the +last input, which keeps the cases deterministic while describing real +behaviour. + +### The consequence nobody should discover in production + +**`stats count` undercounts aggregated data**, silently. A hundred +events become one record, so every existing dashboard panel and alert +rule that counts rows changes meaning the moment a rule starts +aggregating the data behind it. + +Nothing in this design fixes that, and pretending otherwise would be +worse than saying it. Three things follow: + +- The correct idiom over aggregated data is summing `cairnobs.count`, + not counting rows. That belongs in the query language reference before + this action ships. +- Teaching the query layer to do it automatically — making `count` + mean "sum `cairnobs.count` where present" — is a Phase 2 change to the + IR and executor, not a Phase 8 change, and it is the right long-term + answer. **Open**, and it should be decided before aggregation is + recommended for any data an alert already watches. +- It is another argument for aggregation staying opt-in per rule, which + it is. + +**Choosing between the two dedup actions:** `suppress_duplicates` is +cheaper and loses the count; `aggregate_count` preserves it and creates +a synthetic record with all of the above attached. Use suppress when the +repetition is noise, aggregate when the rate is the signal. + ## Where each rule runs Agent-side is the default and the cheaper place: data reduced before the @@ -360,7 +449,10 @@ release cannot do that, it is not finished. 3. ~~Canary rollout?~~ **Decided 2026-09-05:** no gate, ship without it; a canary is follow-up work. 4. Explicit per-rule placement, or platform-decided? -5. Does `aggregate_count` emit a synthetic record, and if so what does it - look like to a query that is not expecting one? This design does not - answer that and should before anyone builds it. The conformance - validator refuses any case using it until it is answered. +5. ~~What does `aggregate_count` emit?~~ **Decided 2026-09-05:** the + window's first record tagged with `cairnobs.aggregated`, + `cairnobs.count` and observed window bounds — see Decision 5. It + raises one new question in its place: whether `stats count` should + learn to sum `cairnobs.count` automatically, which is a Phase 2 + change to the IR rather than a Phase 8 one, and should be settled + before aggregation is pointed at data an alert already watches. diff --git a/processing/README.md b/processing/README.md index 28df35c..bdb81c7 100644 --- a/processing/README.md +++ b/processing/README.md @@ -113,6 +113,7 @@ no later action or rule runs. | `parse_regex` | `field`, `pattern` (named captures become attributes) | | `sample` | `keep_one_in` | | `suppress_duplicates` | `window_ms`, optional `key_fields` | +| `aggregate_count` | `window_ms`, optional `key_fields` | Rules are evaluated in the order given. Every matching rule's actions apply, to the record as left by the rule before it. @@ -150,14 +151,40 @@ close enough at volume. the same output regardless of how fast the test runs. This also means the behaviour is correct under backfill, which wall-clock windows are not. -### `aggregate_count` has no cases, deliberately +### `aggregate_count` emits a record that never existed -The design lists it as an action but does not answer what it emits, or -what a query that is not expecting a synthetic record sees. Writing cases -now would invent that answer by accident and freeze it. It stays -unspecified until that question is decided — -[`phase-8-processing-design.md`](../docs/phase-8-processing-design.md) -open question 5. +The only action whose output is not the input with edits, so its shape is +worth stating here too. It emits the window's **first record, unchanged**, +plus four attributes: + +| attribute | value | +|---|---| +| `cairnobs.aggregated` | `"true"` | +| `cairnobs.count` | records collapsed, as a string | +| `cairnobs.window_start_unix_nano` | first record's timestamp | +| `cairnobs.window_last_unix_nano` | last contributing record's timestamp | + +This follows the convention the agent already uses for heartbeat and +host-metrics records: an ordinary record distinguished by a +`cairnobs.*` attribute, rather than a second synthetic-record mechanism. + +Three behaviours the cases pin, each of which is easy to get wrong: + +- **It tags even when the count is one.** Otherwise `cairnobs.count` is + present only sometimes, and summing it silently breaks on quiet + windows. +- **`window_last` is observed, never `window_start + window_ms`.** A + window flushed early must not claim an end that never happened. +- **A pending window flushes at end of stream.** Record-time windows can + only be closed by a later record, so without this a burst that stops + leaves its aggregate unemitted indefinitely. In production the batch + flush provides the same trigger, which makes `window_ms` a maximum + rather than a guarantee. + +**`stats count` undercounts aggregated data.** The correct idiom is +summing `cairnobs.count`. Whether the query layer should do that +automatically is a Phase 2 question, recorded in +[`phase-8-processing-design.md`](../docs/phase-8-processing-design.md). ## Adding a case diff --git a/processing/conformance/cases/aggregate_count_closes_a_window_on_a_later_record.json b/processing/conformance/cases/aggregate_count_closes_a_window_on_a_later_record.json new file mode 100644 index 0000000..fc2c300 --- /dev/null +++ b/processing/conformance/cases/aggregate_count_closes_a_window_on_a_later_record.json @@ -0,0 +1,69 @@ +{ + "name": "aggregate_count_closes_a_window_on_a_later_record", + "description": "A record at or past window_start + window_ms flushes the pending aggregate first, then opens the next window. The second window is flushed at end of stream.", + "rules": [ + { + "match": [], + "actions": [ + { + "action": "aggregate_count", + "window_ms": 5000 + } + ] + } + ], + "inputs": [ + { + "timestamp_unix_nano": 0, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "x", + "attributes": {} + }, + { + "timestamp_unix_nano": 1000000000, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "x", + "attributes": {} + }, + { + "timestamp_unix_nano": 6000000000, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "x", + "attributes": {} + } + ], + "expect": [ + { + "timestamp_unix_nano": 0, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "x", + "attributes": { + "cairnobs.aggregated": "true", + "cairnobs.count": "2", + "cairnobs.window_start_unix_nano": "0", + "cairnobs.window_last_unix_nano": "1000000000" + } + }, + { + "timestamp_unix_nano": 6000000000, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "x", + "attributes": { + "cairnobs.aggregated": "true", + "cairnobs.count": "1", + "cairnobs.window_start_unix_nano": "6000000000", + "cairnobs.window_last_unix_nano": "6000000000" + } + } + ] +} diff --git a/processing/conformance/cases/aggregate_count_collapses_a_window_into_one_tagged_record.json b/processing/conformance/cases/aggregate_count_collapses_a_window_into_one_tagged_record.json new file mode 100644 index 0000000..4a5776a --- /dev/null +++ b/processing/conformance/cases/aggregate_count_collapses_a_window_into_one_tagged_record.json @@ -0,0 +1,56 @@ +{ + "name": "aggregate_count_collapses_a_window_into_one_tagged_record", + "description": "The emitted record is the window's first record unchanged, plus the aggregate attributes. Keeping the first record means a reader sees a real example of what was collapsed rather than an invented summary.", + "rules": [ + { + "match": [], + "actions": [ + { + "action": "aggregate_count", + "window_ms": 5000 + } + ] + } + ], + "inputs": [ + { + "timestamp_unix_nano": 0, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "repeated", + "attributes": {} + }, + { + "timestamp_unix_nano": 1000000000, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "repeated", + "attributes": {} + }, + { + "timestamp_unix_nano": 2000000000, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "repeated", + "attributes": {} + } + ], + "expect": [ + { + "timestamp_unix_nano": 0, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "repeated", + "attributes": { + "cairnobs.aggregated": "true", + "cairnobs.count": "3", + "cairnobs.window_start_unix_nano": "0", + "cairnobs.window_last_unix_nano": "2000000000" + } + } + ] +} diff --git a/processing/conformance/cases/aggregate_count_honours_key_fields.json b/processing/conformance/cases/aggregate_count_honours_key_fields.json new file mode 100644 index 0000000..e54eb86 --- /dev/null +++ b/processing/conformance/cases/aggregate_count_honours_key_fields.json @@ -0,0 +1,73 @@ +{ + "name": "aggregate_count_honours_key_fields", + "description": "With key_fields, each key gets its own window and its own aggregate. Emission order at end of stream follows the order the windows were opened.", + "rules": [ + { + "match": [], + "actions": [ + { + "action": "aggregate_count", + "window_ms": 5000, + "key_fields": [ + "host", + "message" + ] + } + ] + } + ], + "inputs": [ + { + "timestamp_unix_nano": 0, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "z", + "attributes": {} + }, + { + "timestamp_unix_nano": 100000000, + "host": "h2", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "z", + "attributes": {} + }, + { + "timestamp_unix_nano": 200000000, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "z", + "attributes": {} + } + ], + "expect": [ + { + "timestamp_unix_nano": 0, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "z", + "attributes": { + "cairnobs.aggregated": "true", + "cairnobs.count": "2", + "cairnobs.window_start_unix_nano": "0", + "cairnobs.window_last_unix_nano": "200000000" + } + }, + { + "timestamp_unix_nano": 100000000, + "host": "h2", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "z", + "attributes": { + "cairnobs.aggregated": "true", + "cairnobs.count": "1", + "cairnobs.window_start_unix_nano": "100000000", + "cairnobs.window_last_unix_nano": "100000000" + } + } + ] +} diff --git a/processing/conformance/cases/aggregate_count_leaves_non_matching_records_alone.json b/processing/conformance/cases/aggregate_count_leaves_non_matching_records_alone.json new file mode 100644 index 0000000..dfd66e6 --- /dev/null +++ b/processing/conformance/cases/aggregate_count_leaves_non_matching_records_alone.json @@ -0,0 +1,70 @@ +{ + "name": "aggregate_count_leaves_non_matching_records_alone", + "description": "Records the rule did not match pass through untouched and in order, interleaved with the aggregate the matched ones produced.", + "rules": [ + { + "match": [ + { + "field": "service", + "op": "eq", + "value": "noisy" + } + ], + "actions": [ + { + "action": "aggregate_count", + "window_ms": 5000 + } + ] + } + ], + "inputs": [ + { + "timestamp_unix_nano": 0, + "host": "h1", + "service": "noisy", + "severity": "SEVERITY_INFO", + "message": "n", + "attributes": {} + }, + { + "timestamp_unix_nano": 100000000, + "host": "h1", + "service": "quiet", + "severity": "SEVERITY_INFO", + "message": "keep me", + "attributes": {} + }, + { + "timestamp_unix_nano": 200000000, + "host": "h1", + "service": "noisy", + "severity": "SEVERITY_INFO", + "message": "n", + "attributes": {} + } + ], + "expect": [ + { + "timestamp_unix_nano": 100000000, + "host": "h1", + "service": "quiet", + "severity": "SEVERITY_INFO", + "message": "keep me", + "attributes": {} + }, + { + "timestamp_unix_nano": 0, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "n", + "attributes": { + "cairnobs.aggregated": "true", + "cairnobs.count": "2", + "cairnobs.window_start_unix_nano": "0", + "cairnobs.window_last_unix_nano": "200000000" + } + } + ] +} diff --git a/processing/conformance/cases/aggregate_count_tags_even_a_single_occurrence.json b/processing/conformance/cases/aggregate_count_tags_even_a_single_occurrence.json new file mode 100644 index 0000000..4b9a027 --- /dev/null +++ b/processing/conformance/cases/aggregate_count_tags_even_a_single_occurrence.json @@ -0,0 +1,40 @@ +{ + "name": "aggregate_count_tags_even_a_single_occurrence", + "description": "A window that saw one record is still tagged with count 1. Emitting a bare record here would make cairnobs.count present only sometimes, which silently breaks summing it on quiet windows.", + "rules": [ + { + "match": [], + "actions": [ + { + "action": "aggregate_count", + "window_ms": 5000 + } + ] + } + ], + "inputs": [ + { + "timestamp_unix_nano": 0, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "lonely", + "attributes": {} + } + ], + "expect": [ + { + "timestamp_unix_nano": 0, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "lonely", + "attributes": { + "cairnobs.aggregated": "true", + "cairnobs.count": "1", + "cairnobs.window_start_unix_nano": "0", + "cairnobs.window_last_unix_nano": "0" + } + } + ] +} diff --git a/processing/conformance/cases/aggregate_count_window_bounds_are_observed_not_nominal.json b/processing/conformance/cases/aggregate_count_window_bounds_are_observed_not_nominal.json new file mode 100644 index 0000000..5126588 --- /dev/null +++ b/processing/conformance/cases/aggregate_count_window_bounds_are_observed_not_nominal.json @@ -0,0 +1,48 @@ +{ + "name": "aggregate_count_window_bounds_are_observed_not_nominal", + "description": "window_last is the last record that actually contributed, never window_start + window_ms. A window flushed early must not claim an end that never happened.", + "rules": [ + { + "match": [], + "actions": [ + { + "action": "aggregate_count", + "window_ms": 60000 + } + ] + } + ], + "inputs": [ + { + "timestamp_unix_nano": 100000000, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "y", + "attributes": {} + }, + { + "timestamp_unix_nano": 250000000, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "y", + "attributes": {} + } + ], + "expect": [ + { + "timestamp_unix_nano": 100000000, + "host": "h1", + "service": "s1", + "severity": "SEVERITY_INFO", + "message": "y", + "attributes": { + "cairnobs.aggregated": "true", + "cairnobs.count": "2", + "cairnobs.window_start_unix_nano": "100000000", + "cairnobs.window_last_unix_nano": "250000000" + } + } + ] +} diff --git a/processing/conformance/validate.py b/processing/conformance/validate.py index a2195b2..39516b4 100755 --- a/processing/conformance/validate.py +++ b/processing/conformance/validate.py @@ -45,12 +45,15 @@ ACTIONS = { "parse_regex": ({"field", "pattern"}, set()), "sample": ({"keep_one_in"}, set()), "suppress_duplicates": ({"window_ms"}, {"key_fields"}), + "aggregate_count": ({"window_ms"}, {"key_fields"}), } -# Deliberately absent from ACTIONS. The design does not say what it emits -# or what a query that is not expecting a synthetic record sees, so a -# case using it would invent that answer and freeze it by accident. -UNSPECIFIED_ACTIONS = {"aggregate_count"} +# Actions the design names but does not yet specify the output of. A case +# using one would invent that answer and freeze it by accident, so the +# validator refuses them until the design says what they emit. +# aggregate_count was here until 2026-09-05; see Decision 5 in +# /docs/phase-8-processing-design.md. +UNSPECIFIED_ACTIONS = set() def valid_field(name): @@ -120,7 +123,7 @@ def check_action(a, where, err): n = a.get("keep_one_in") if not isinstance(n, int) or n < 1: err(f"{where}: keep_one_in must be an integer >= 1") - if name == "suppress_duplicates": + if name in ("suppress_duplicates", "aggregate_count"): w = a.get("window_ms") if not isinstance(w, int) or w < 1: err(f"{where}: window_ms must be an integer >= 1")