Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion docs/weather-starter.md
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
2 changes: 1 addition & 1 deletion python/src/sde_demo/model.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
53 changes: 52 additions & 1 deletion python/src/sde_demo/runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,8 @@
CELSIUS_BASE_CENTS,
CELSIUS_MODULUS,
GENERATOR_ID,
HUMIDITY_BASE,
HUMIDITY_MODULUS,
WORKLOADS,
model,
reading,
Expand All @@ -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]:
Expand Down Expand Up @@ -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))
Expand All @@ -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 {}
Expand Down Expand Up @@ -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(
Expand Down
115 changes: 114 additions & 1 deletion python/tests/test_demo_starter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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"):
Expand Down
10 changes: 10 additions & 0 deletions python/tests/test_demo_starter_live.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
3 changes: 2 additions & 1 deletion typescript/src/demo/model.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
42 changes: 40 additions & 2 deletions typescript/src/demo/weather.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -88,6 +94,23 @@ export function fleetExpected(runs: Map<string, number>, 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<RunReport> {
const { iterations = 10, batchSize = 10, intervalMs = 100, recoveryMs = 10000, workload = 'mixed' } = options
if (![iterations, batchSize, intervalMs, recoveryMs].every(Number.isSafeInteger) ||
Expand Down Expand Up @@ -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),
Expand Down
8 changes: 8 additions & 0 deletions typescript/tests/weather.live.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Loading
Loading