From 45f2e017b0227200270e20d223039695bfc0a26b Mon Sep 17 00:00:00 2001 From: Krzysztof Macewicz Date: Thu, 24 Sep 2026 20:37:47 +0200 Subject: [PATCH] feat: an alerts workload in the Weather starter, a range outside the key Every read the starter could drive filtered on key columns, so traffic it generated gave a control plane no reason for the one physical change built in place: an added index. --workload alerts reads one station's readings at or above 95% humidity - a bounded page, a count and an exact celsius summary - in both languages, checked exactly against the generator (humidity is 30 + sequence mod 70, so five in seventy alert; before the first one the summary total is null, as for any summary of no values). The window's filtered_on names the station equality and the humidity range. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/weather-starter.md | 7 +- python/src/sde_demo/model.py | 2 +- python/src/sde_demo/runtime.py | 53 +++++++++++- python/tests/test_demo_starter.py | 115 ++++++++++++++++++++++++- python/tests/test_demo_starter_live.py | 10 +++ typescript/src/demo/model.ts | 3 +- typescript/src/demo/weather.ts | 42 ++++++++- typescript/tests/weather.live.test.ts | 8 ++ typescript/tests/weather.test.ts | 70 ++++++++++++++- 9 files changed, 300 insertions(+), 10 deletions(-) diff --git a/docs/weather-starter.md b/docs/weather-starter.md index 228add6..5b32ded 100644 --- a/docs/weather-starter.md +++ b/docs/weather-starter.md @@ -85,7 +85,12 @@ Every run in a directory writes the same timeline, so the window holds every run program checks them exactly against the local reports of the directory's earlier runs of this project. It therefore refuses to start while any of those runs is not complete: an interrupted run or an unresolved batch makes the expected answer unknowable, and it does not compare with a guess. -Limits are 1-1000 iterations and rows per batch, at most 10000 rows per invocation, 0-1000 ms +`--workload alerts` reads one station's readings at or above 95% humidity - a bounded page, a count +and an exact celsius summary. Humidity is not in the key, so the key narrows the read to the station +and cannot serve the range: this is the traffic an added index is for, and the window's `filtered_on` +says so (`equal: ["station"]`, `range: "humidity"`). The expected answer is exact from the generator: +humidity is 30 + sequence mod 70, so five readings in every seventy alert, and before the first one a +summary of no values has a null total. Limits are 1-1000 iterations and rows per batch, at most 10000 rows per invocation, 0-1000 ms between batches and 0-30000 ms for bounded read/recovery attempts. Each run writes `runs/RUN_ID/report.json` and a real SDK `window.json`. Reports contain counters, diff --git a/python/src/sde_demo/model.py b/python/src/sde_demo/model.py index fcde5ca..577d6ea 100644 --- a/python/src/sde_demo/model.py +++ b/python/src/sde_demo/model.py @@ -36,7 +36,7 @@ CELSIUS_BASE_CENTS, CELSIUS_MODULUS = 1525, 1000 HUMIDITY_BASE, HUMIDITY_MODULUS = 30, 70 UUID_VERSION = 4 -WORKLOADS = ("mixed", "point", "analytics", "fleet") +WORKLOADS = ("mixed", "point", "analytics", "fleet", "alerts") """What a run can drive. The CLI offers and the runtime accepts this one tuple.""" GENERATOR_SPEC = { "kind": "sde-weather-generator", diff --git a/python/src/sde_demo/runtime.py b/python/src/sde_demo/runtime.py index ec77e69..ee5b79d 100644 --- a/python/src/sde_demo/runtime.py +++ b/python/src/sde_demo/runtime.py @@ -18,6 +18,8 @@ CELSIUS_BASE_CENTS, CELSIUS_MODULUS, GENERATOR_ID, + HUMIDITY_BASE, + HUMIDITY_MODULUS, WORKLOADS, model, reading, @@ -28,6 +30,10 @@ FLEET_PAGE = 100 MAX_RUN_ROWS = 10000 """The fleet checks one cross-station page exactly; its count and sum cover the whole window.""" +ALERT_HUMIDITY = 95 +"""The alert threshold. Humidity is 30 + sequence % 70, so a reading alerts when sequence % 70 is +65 or more - five readings in seventy, known exactly from the sequence number.""" +ALERT_PAGE = 100 def fleet_runs(root: Path, *, project_id: str) -> dict[str, int]: @@ -83,6 +89,26 @@ def fleet_expected( return page, total, Decimal(f"{cents // 100}.{cents % 100:02d}") +def alert_expected( + run_id: str, through: int, limit: int +) -> tuple[list[dict[str, Any]], int, Decimal | None]: + """The first page, count and celsius total of this run's alerts among sequences 1..through. + + One station, so key order - station, then time - is sequence order. With no alert yet the total + is None, as the library reports a summary of no values on every engine. + """ + page: list[dict[str, Any]] = [] + total, cents = 0, 0 + for sequence in range(1, through + 1): + if HUMIDITY_BASE + sequence % HUMIDITY_MODULUS < ALERT_HUMIDITY: + continue + total += 1 + cents += CELSIUS_BASE_CENTS + sequence % CELSIUS_MODULUS + if len(page) < limit: + page.append(reading(run_id, 0, sequence)) + return page, total, Decimal(f"{cents // 100}.{cents % 100:02d}") if total else None + + def limits(iterations: int, batch_size: int, interval_ms: int, recovery_ms: int) -> None: if ( any(type(value) is not int for value in (iterations, batch_size, interval_ms, recovery_ms)) @@ -109,7 +135,7 @@ def run( ) -> dict[str, Any]: limits(iterations, batch_size, interval_ms, recovery_ms) if workload not in WORKLOADS: - raise DemoRefused("Choose the mixed, point, analytics or fleet workload.") + raise DemoRefused("Choose the mixed, point, analytics, fleet or alerts workload.") settings = config(root) # Read before this run's own report exists; a run that starts later is not in the window. earlier = fleet_runs(root, project_id=settings["project_id"]) if workload == "fleet" else {} @@ -281,6 +307,31 @@ def check_fleet_page( raise DemoRefused( "The exact fleet summary did not match this directory's runs." ) + elif workload == "alerts": + # One station's readings at or above the alert threshold: a range on humidity, a + # field outside the key, which the key cannot serve and an index can. The expected + # answer is exact, from the generator (``alert_expected``). + alert_page, alert_rows, alert_celsius = alert_expected(run_id, count, ALERT_PAGE) + alert_bounds = sde.Range("humidity", low=ALERT_HUMIDITY) + + def check_alert_page( + current: sde.Session, + where: dict[str, Any] = where, + bounds: sde.Range = alert_bounds, + expected: list[dict[str, Any]] = alert_page, + ) -> None: + page = current.scan( + "WeatherReading", where=where, bounds=bounds, limit=ALERT_PAGE + ) + if list(page.rows) != expected: + raise DemoRefused("The alert page did not match this run.") + + read_retry(check_alert_page) + if read_retry(partial(counted, where=where, bounds=alert_bounds)) != alert_rows: + raise DemoRefused("The alert count did not match this run.") + summary = read_retry(partial(summarized, where=where, bounds=alert_bounds)) + if summary.count != alert_rows or summary.total != alert_celsius: + raise DemoRefused("The exact alert summary did not match this run.") elif workload != "point" or iteration == iterations - 1: def check_page( diff --git a/python/tests/test_demo_starter.py b/python/tests/test_demo_starter.py index 63e555d..bc33f25 100644 --- a/python/tests/test_demo_starter.py +++ b/python/tests/test_demo_starter.py @@ -401,6 +401,116 @@ def test_fleet_expectation_is_every_run_row_in_the_window() -> None: assert runtime.fleet_expected(FLEET_RUNS, through, 100)[0] == rows +@pytest.mark.parametrize(("through", "limit"), [(150, 7), (150, 100), (64, 100), (65, 1)]) +def test_alert_expectation_is_every_alert_of_the_run(through: int, limit: int) -> None: + """Checked against brute force: every generated row at or above the threshold, in key order.""" + rows = [ + row + for row in (reading("a" * 32, 0, sequence) for sequence in range(1, through + 1)) + if row["humidity"] >= runtime.ALERT_HUMIDITY + ] + rows.sort(key=lambda row: (row["station"], row["at"])) + page, total, celsius = runtime.alert_expected("a" * 32, through, limit) + assert page == rows[:limit] + assert total == len(rows) + assert celsius == (sum((row["celsius"] for row in rows), Decimal(0)) if rows else None) + # The threshold is a humidity the generator reaches and exceeds: five in every seventy. + assert total == {150: 10, 64: 0, 65: 1}[through] + + +class Engine: + """An in-memory session answering reads as an engine does: equality, a range, key order. + + ``defect`` makes one answer wrong in the way a broken adapter or layout could, so the alerts + workload's own comparison is what must refuse the run. + """ + + def __init__(self, monkeypatch: pytest.MonkeyPatch, defect: str | None = None) -> None: + self.rows: list[dict[str, Any]] = [] + harness = self + + class Client: + def close(self) -> None: + pass + + def save_many(self, entity: str, rows: list[dict[str, Any]]) -> None: + harness.rows.extend(rows) + + def get( + self, entity: str, key: dict[str, Any], **options: Any + ) -> dict[str, Any] | None: + return next( + (row for row in harness.rows if all(row[k] == v for k, v in key.items())), None + ) + + def matching( + self, where: dict[str, Any] | None, bounds: sde.Range | None, ranged: bool = True + ) -> list[dict[str, Any]]: + found = [] + for row in sorted(harness.rows, key=lambda row: (row["station"], row["at"])): + if any(row[name] != value for name, value in (where or {}).items()): + continue + if bounds is not None and ranged: + value = row[bounds.field] + if (bounds.low is not None and value < bounds.low) or ( + bounds.high is not None and value >= bounds.high + ): + continue + found.append(row) + return found + + def scan(self, entity: str, **options: Any) -> sde.ScanPage: + rows = self.matching( + options.get("where"), options.get("bounds"), ranged=defect != "page" + ) + return sde.ScanPage(tuple(rows[: options["limit"]]), None) + + def count(self, entity: str, **options: Any) -> int: + rows = self.matching( + options.get("where"), options.get("bounds"), ranged=defect != "count" + ) + return len(rows) + + def summarize(self, entity: str, field: str, **options: Any) -> sde.NumericSummary: + rows = self.matching(options.get("where"), options.get("bounds")) + total = sum((row[field] for row in rows), Decimal(0)) if rows else None + if defect == "summary" and total is not None: + total += Decimal("0.01") + if defect == "empty_total" and total is None: + total = Decimal("0.00") + return sde.NumericSummary(len(rows), len(rows), None, None, total, None) + + def connect(_model: Any, _placement: Any, factories: Any, **options: Any) -> Client: + return Client() + + monkeypatch.setattr(sde.Session, "connect", connect) + + +@pytest.mark.parametrize( + ("defect", "refusal"), + [ + (None, None), + ("page", "alert page"), + ("count", "alert count"), + ("summary", "alert summary"), + ("empty_total", "alert summary"), + ], +) +def test_the_alerts_workload_checks_every_answer_exactly( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch, defect: str | None, refusal: str | None +) -> None: + """Two iterations of 40: none alerts in the first, five in the second.""" + local(tmp_path) + Engine(monkeypatch, defect) + options: dict[str, Any] = {"iterations": 2, "batch_size": 40, "interval_ms": 0} + if refusal is None: + report = runtime.run(tmp_path, workload="alerts", **options) + assert report["status"] == "complete" and report["verified_rows"] == 80 + return + with pytest.raises(project.DemoRefused, match=refusal): + runtime.run(tmp_path, workload="alerts", **options) + + def _report(root: Path, identity: str, **fields: Any) -> None: report = { "protocol": 2, @@ -462,7 +572,10 @@ def test_the_cli_and_the_runtime_offer_the_same_workloads(tmp_path: Path) -> Non assert cli.WORKLOADS is runtime.WORKLOADS # one tuple, imported by both # The parser accepts every workload the runtime runs (the missing setup refuses afterwards, # which returns a status) and rejects anything else before the runtime is reached. - assert isinstance(cli.main(["--directory", str(tmp_path), "run", "--workload", "fleet"]), int) + for workload in runtime.WORKLOADS: + assert isinstance( + cli.main(["--directory", str(tmp_path), "run", "--workload", workload]), int + ) with pytest.raises(SystemExit): cli.main(["--directory", str(tmp_path), "run", "--workload", "bogus"]) with pytest.raises(project.DemoRefused, match="fleet"): diff --git a/python/tests/test_demo_starter_live.py b/python/tests/test_demo_starter_live.py index 24692da..6e4b287 100644 --- a/python/tests/test_demo_starter_live.py +++ b/python/tests/test_demo_starter_live.py @@ -67,6 +67,16 @@ def active_only(dialect: str, dsn: str) -> object: filters = {entry["kind"]: entry.get("filtered_on") for entry in shapes} assert filters["range_read"] == [{"equal": [], "range": "at", "calls": 2}] assert filters["aggregate"] == [{"equal": [], "range": "at", "calls": 4}] + # Alerts read one station's readings at or above a humidity: the first iteration has no + # alert yet (an empty page and a summary of nothing), the second has five. + alerts = runtime.run(root, iterations=2, batch_size=40, interval_ms=0, workload="alerts") + assert alerts["status"] == "complete" and alerts["verified_rows"] == 80 + shapes = project.read(root / "runs" / alerts["run_id"] / "window.json")["groups"][ + "WeatherReading" + ]["shapes"] + filters = {entry["kind"]: entry.get("filtered_on") for entry in shapes} + assert filters["range_read"] == [{"equal": ["station"], "range": "humidity", "calls": 2}] + assert filters["aggregate"] == [{"equal": ["station"], "range": "humidity", "calls": 4}] before = (root / "state" / "active-map.json").read_bytes() assert project.setup(root, bundle, admin) == {"status": "ready", "map_version": 1} assert (root / "state" / "active-map.json").read_bytes() == before diff --git a/typescript/src/demo/model.ts b/typescript/src/demo/model.ts index 2b43450..aed8776 100644 --- a/typescript/src/demo/model.ts +++ b/typescript/src/demo/model.ts @@ -17,7 +17,8 @@ export function weatherModel() { export const baseTime = Timestamp.from('2026-01-01T00:00:00.000000Z') const stationTemplate = 'weather-{run_id}-{worker}', seedTemplate = 'sde-weather-v1:{run_id}:{worker}:{sequence}' export const celsiusBaseCents = 1525, celsiusModulus = 1000 -const humidityBase = 30, humidityModulus = 70, uuidVersion = 4 +export const humidityBase = 30, humidityModulus = 70 +const uuidVersion = 4 export const generatorSpec = Object.freeze({ kind: 'sde-weather-generator', version: 1, worker: 0, station_template: stationTemplate, seed_template: seedTemplate, diff --git a/typescript/src/demo/weather.ts b/typescript/src/demo/weather.ts index dbf0f2f..acaba3e 100644 --- a/typescript/src/demo/weather.ts +++ b/typescript/src/demo/weather.ts @@ -10,15 +10,21 @@ import { ClickHouseEngine } from '../engines/clickhouse.js' import { Session } from '../session.js' import type { ManagedEngine, Row } from '../session.js' import { Recorder, windowRecord } from '../telemetry.js' -import { baseTime, celsiusBaseCents, celsiusModulus, generatorId, reading } from './model.js' +import { baseTime, celsiusBaseCents, celsiusModulus, generatorId, humidityBase, humidityModulus, reading } from './model.js' import { DemoRefused, project, read, write } from './project.js' /** What a run can drive. The CLI offers and the runtime accepts this one list. */ -export const workloads = ['mixed', 'point', 'analytics', 'fleet'] as const +export const workloads = ['mixed', 'point', 'analytics', 'fleet', 'alerts'] as const export type Workload = typeof workloads[number] /** The fleet checks one cross-station page exactly; its count and sum cover the whole window. */ const fleetPage = 100 const maxRunRows = 10000 +/** + * The alert threshold. Humidity is 30 + sequence % 70, so a reading alerts when sequence % 70 is 65 + * or more - five readings in seventy, known exactly from the sequence number. + */ +export const alertHumidity = 95 +export const alertPage = 100 export interface RunOptions { iterations?: number; batchSize?: number; intervalMs?: number; recoveryMs?: number @@ -88,6 +94,23 @@ export function fleetExpected(runs: Map, through: number, limit: } return [page, total, `${cents / 100n}.${String(cents % 100n).padStart(2, '0')}`] } +/** + * The first page, count and celsius total of this run's alerts among sequences 1..through. One + * station, so key order - station, then time - is sequence order. With no alert yet the total is + * null, as the library reports a summary of no values on every engine. The same rules as the Python + * starter's alert_expected. + */ +export function alertExpected(runId: string, through: number, limit: number): [Row[], number, string | null] { + const page: Row[] = [] + let total = 0, cents = 0n + for (let sequence = 1; sequence <= through; sequence++) { + if (humidityBase + sequence % humidityModulus < alertHumidity) continue + total++ + cents += BigInt(celsiusBaseCents + sequence % celsiusModulus) + if (page.length < limit) page.push(reading(runId, 0, sequence)) + } + return [page, total, total ? `${cents / 100n}.${String(cents % 100n).padStart(2, '0')}` : null] +} export async function runWeather(root: string, options: RunOptions = {}): Promise { const { iterations = 10, batchSize = 10, intervalMs = 100, recoveryMs = 10000, workload = 'mixed' } = options if (![iterations, batchSize, intervalMs, recoveryMs].every(Number.isSafeInteger) || @@ -186,6 +209,21 @@ export async function runWeather(root: string, options: RunOptions = {}): Promis if (summary.count !== BigInt(rows) || summary.total !== celsius) { throw new DemoRefused("The exact fleet summary did not match this directory's runs.") } + } else if (workload === 'alerts') { + // One station's readings at or above the alert threshold: a range on humidity, a field + // outside the key, which the key cannot serve and an index can. Exact from the generator. + const [expected, alerts, celsius] = alertExpected(runId, count, alertPage) + const bounds = { field: 'humidity', low: BigInt(alertHumidity) } + const page = await readRetry(current => current.scan('WeatherReading', { where, bounds, limit: alertPage })) + if (page.rows.length !== expected.length) throw new DemoRefused('The alert page did not match this run.') + page.rows.forEach((row, index) => same(row, expected[index]!)) + if (await readRetry(current => current.count('WeatherReading', { where, bounds })) !== BigInt(alerts)) { + throw new DemoRefused('The alert count did not match this run.') + } + const summary = await readRetry(current => current.summarize('WeatherReading', 'celsius', { where, bounds })) + if (summary.count !== BigInt(alerts) || summary.total !== celsius) { + throw new DemoRefused('The exact alert summary did not match this run.') + } } else if (workload !== 'point' || iteration === iterations - 1) { const page = await readRetry(current => current.scan('WeatherReading', { where, bounds: { field: 'at', low: baseTime }, limit: Math.min(count, 1000), diff --git a/typescript/tests/weather.live.test.ts b/typescript/tests/weather.live.test.ts index f81f234..a1eab40 100644 --- a/typescript/tests/weather.live.test.ts +++ b/typescript/tests/weather.live.test.ts @@ -44,6 +44,14 @@ for (const source of ['postgres', 'clickhouse']) it.skipIf(!enabled)(`runs Weath const filters = Object.fromEntries(shapes.map((entry: { kind: string; filtered_on?: unknown }) => [entry.kind, entry.filtered_on])) expect(filters.range_read).toEqual([{ equal: [], range: 'at', calls: 2 }]) expect(filters.aggregate).toEqual([{ equal: [], range: 'at', calls: 4 }]) + // Alerts read one station's readings at or above a humidity: no alert in the first iteration + // (an empty page and a summary of nothing), five in the second. + const alerts = await runWeather(directory, { iterations: 2, batchSize: 40, intervalMs: 0, workload: 'alerts' }) + expect(alerts.status).toBe('complete'); expect(alerts.verified_rows).toBe(80) + const alertShapes = JSON.parse(readFileSync(join(directory, 'runs', alerts.run_id, 'window.json'), 'utf8')).groups.WeatherReading.shapes + const alertFilters = Object.fromEntries(alertShapes.map((entry: { kind: string; filtered_on?: unknown }) => [entry.kind, entry.filtered_on])) + expect(alertFilters.range_read).toEqual([{ equal: ['station'], range: 'humidity', calls: 2 }]) + expect(alertFilters.aggregate).toEqual([{ equal: ['station'], range: 'humidity', calls: 4 }]) expect(JSON.parse(call('reset')).status).toBe('reset') expect(JSON.parse(call('reset')).status).toBe('reset') } finally { diff --git a/typescript/tests/weather.test.ts b/typescript/tests/weather.test.ts index a0da23b..573fdef 100644 --- a/typescript/tests/weather.test.ts +++ b/typescript/tests/weather.test.ts @@ -12,7 +12,7 @@ import type { Row } from '../src/session.js' import { EngineError } from '../src/errors.js' import { generatorId, reading, weatherModel } from '../src/demo/model.js' import { project, read, write } from '../src/demo/project.js' -import { fleetExpected, fleetRuns, runWeather, workloads } from '../src/demo/weather.js' +import { alertExpected, alertHumidity, fleetExpected, fleetRuns, runWeather, workloads } from '../src/demo/weather.js' const roots: string[] = [] afterEach(() => { vi.restoreAllMocks(); for (const root of roots.splice(0)) rmSync(root, { recursive: true, force: true }) }) @@ -153,6 +153,70 @@ it('expects every run row in the fleet window, in key order (brute force)', () = expect(celsius).toBe(`${cents / 100n}.${String(cents % 100n).padStart(2, '0')}`) expect(fleetExpected(fleet, through, 100)[0]).toEqual(rows) }) +for (const [through, limit, count] of [[150, 7, 10], [150, 100, 10], [64, 100, 0], [65, 1, 1]] as const) { + it(`expects every alert of the run in key order, ${through} rows by ${limit} (brute force)`, () => { + const rows = Array.from({ length: through }, (_, index) => reading('a'.repeat(32), 0, index + 1)) + .filter(row => row.humidity >= BigInt(alertHumidity)) + const [page, total, celsius] = alertExpected('a'.repeat(32), through, limit) + expect(page).toEqual(rows.slice(0, limit)) + expect(total).toBe(rows.length); expect(total).toBe(count) + const cents = rows.reduce((sum, row) => sum + BigInt(row.celsius.replace('.', '')), 0n) + expect(celsius).toBe(rows.length ? `${cents / 100n}.${String(cents % 100n).padStart(2, '0')}` : null) + }) +} +/** An in-memory session answering reads as an engine does; `defect` makes one answer wrong. */ +function engine(defect?: string) { + const rows: Row[] = [] + type Bounds = { field: string; low?: bigint; high?: bigint } + const station = (row: Row) => row.station as string + const micros = (row: Row) => (row.at as { epochMicroseconds: bigint }).epochMicroseconds + function matching(where: Row = {}, bounds?: Bounds, ranged = true) { + return [...rows].sort((left, right) => station(left) < station(right) ? -1 : station(left) > station(right) ? 1 + : micros(left) < micros(right) ? -1 : 1) + .filter(row => Object.entries(where).every(([name, value]) => row[name] === value)) + .filter(row => { + if (!bounds || !ranged) return true + const value = row[bounds.field] as bigint + return (bounds.low === undefined || value >= bounds.low) && (bounds.high === undefined || value < bounds.high) + }) + } + vi.spyOn(Session, 'connect').mockImplementation(async () => ({ + async close() {}, + async saveMany(_entity: string, batch: Row[]) { rows.push(...batch) }, + async get(_entity: string, key: Row) { + return rows.find(row => row.station === key.station && isDeepStrictEqual(row.at, key.at)) ?? null + }, + async scan(_entity: string, options: { where?: Row; bounds?: Bounds; limit: number }) { + return { rows: matching(options.where, options.bounds, defect !== 'page').slice(0, options.limit), nextAfter: null } + }, + async count(_entity: string, options: { where?: Row; bounds?: Bounds }) { + return BigInt(matching(options.where, options.bounds, defect !== 'count').length) + }, + async summarize(_entity: string, _field: string, options: { where?: Row; bounds?: Bounds }) { + const found = matching(options.where, options.bounds) + let cents = found.reduce((total, row) => total + BigInt(String(row.celsius).replace('.', '')), 0n) + if (defect === 'summary' && found.length) cents += 1n + const total = found.length || defect === 'empty_total' + ? `${cents / 100n}.${String(cents % 100n).padStart(2, '0')}` : null + return { count: BigInt(found.length), total } + }, + }) as unknown as Session) + return rows +} +for (const [defect, refusal] of [[undefined, undefined], ['page', 'alert page'], ['count', 'alert count'], + ['summary', 'alert summary'], ['empty_total', 'alert summary']] as const) { + it(`checks every alerts answer exactly (${defect ?? 'faithful engine'})`, async () => { + // Two iterations of 40: none alerts in the first, five in the second. + const { root } = fixture(), rows = engine(defect) + const run = runWeather(root, { iterations: 2, batchSize: 40, intervalMs: 0, workload: 'alerts' }) + if (refusal === undefined) { + const report = await run + expect(report.status).toBe('complete'); expect(report.verified_rows).toBe(80); expect(rows).toHaveLength(80) + } else { + await expect(run).rejects.toThrow(refusal) + } + }) +} function report(root: string, identity: string, fields: Record = {}) { write(join(root, 'runs', identity, 'report.json'), { protocol: 2, run_id: identity, project_id: 'p'.repeat(32), status: 'complete', pending: null, generator_id: generatorId, verified_rows: 4, ...fields }) @@ -173,7 +237,7 @@ for (const [name, fields] of Object.entries({ report(root, 'a'.repeat(32), fields) expect(() => fleetRuns(root, 'p'.repeat(32))).toThrow(/every earlier run/) }) -it('offers the fleet workload and refuses anything else before reading setup', async () => { - expect(workloads).toEqual(['mixed', 'point', 'analytics', 'fleet']) +it('offers the fleet and alerts workloads and refuses anything else before reading setup', async () => { + expect(workloads).toEqual(['mixed', 'point', 'analytics', 'fleet', 'alerts']) await expect(runWeather('/nonexistent', { workload: 'bogus' as never })).rejects.toThrow(/Invalid bounded workload/) })