diff --git a/docs/weather-starter.md b/docs/weather-starter.md index c94a95a..09c0f7f 100644 --- a/docs/weather-starter.md +++ b/docs/weather-starter.md @@ -20,6 +20,11 @@ npm init -y npm install ./smart-data-engines-sde-0.1.0-dev.0.tgz pg ``` +`setup`, `doctor` and `run` refuse before writing anything when a binding's driver cannot be imported, +naming what to install (`[postgres]` or `[clickhouse]` in Python, `pg` in Node). `setup` also says +why the bootstrap's map does not load - for example that verifying its signature needs the `signed` +extra - instead of reporting an incomplete operation. + The version strings above identify development artifacts, not a promise that every package with that version contains this feature. Retain their SHA-256 checksums and the supplying source commit. The project owner can build a wheel with `python -m build --wheel --outdir ARTIFACT_DIRECTORY python` @@ -89,6 +94,9 @@ 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. +One unfinished run is known exactly and does not stop it: a run that failed before its first write +(`incomplete`, no `pending` range, and no acknowledged or verified row), because a run records each +batch before writing it. `--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` diff --git a/python/src/sde_demo/diagnostics.py b/python/src/sde_demo/diagnostics.py index c6db089..2338f6a 100644 --- a/python/src/sde_demo/diagnostics.py +++ b/python/src/sde_demo/diagnostics.py @@ -12,12 +12,13 @@ from . import resources from .model import model -from .project import config, connections, public_keys +from .project import config, connections, public_keys, require_drivers def doctor(root: Path) -> dict[str, Any]: with transaction(root): settings = config(root) + require_drivers(binding["dialect"] for binding in settings["engines"].values()) resources.verify(root) return _diagnose(root, settings) diff --git a/python/src/sde_demo/project.py b/python/src/sde_demo/project.py index 4f68405..4577f0e 100644 --- a/python/src/sde_demo/project.py +++ b/python/src/sde_demo/project.py @@ -3,10 +3,12 @@ from __future__ import annotations import base64 +import binascii +import importlib import json import re import stat -from collections.abc import Iterator, Mapping +from collections.abc import Iterable, Iterator, Mapping from contextlib import contextmanager from pathlib import Path from typing import Any @@ -72,7 +74,10 @@ def public_keys(value: Any) -> dict[str, bytes]: for name, encoded in value.items(): if not isinstance(name, str) or not name or not isinstance(encoded, str): raise DemoRefused("Invalid public key configuration.") - decoded = base64.b64decode(encoded, validate=True) + try: + decoded = base64.b64decode(encoded, validate=True) + except binascii.Error: + raise DemoRefused("Invalid public key configuration.") from None if len(decoded) != 32: raise DemoRefused("Ed25519 public keys must contain 32 bytes.") result[name] = decoded @@ -105,9 +110,14 @@ def bootstrap(value: dict[str, Any]) -> tuple[dict[str, Any], sde.PlacementMap]: "Weather demo needs one or two uniquely named PostgreSQL/ClickHouse bindings." ) keys = public_keys(value["public_keys"]) - placement = sde.load_map( - value["current_map"], model=logical, public_key=keys, require_signature=True - ) + try: + placement = sde.load_map( + value["current_map"], model=logical, public_key=keys, require_signature=True + ) + except sde.MapError as exc: + # Refused before setup writes anything, and the reason is about the map - a missing + # 'signed' extra, a bad signature - never about a credential, so it can be said in full. + raise DemoRefused(f"The bootstrap's map does not load: {exc}") from None check_map_project(placement, value["project_id"]) if placement.contract < GENERATIONS_SINCE or placement.map_version != 1: raise DemoRefused( @@ -170,6 +180,32 @@ def credentials(root: Path, purpose: str, bindings: Mapping[str, Any]) -> dict[s return value +DRIVERS = {"postgres": ("psycopg", "postgres"), "clickhouse": ("clickhouse_connect", "clickhouse")} +"""The driver module each binding's adapter imports, and the SDK extra that installs it.""" + + +def require_drivers(dialects: Iterable[str]) -> None: + """Refuse before anything is written when a binding's driver cannot be imported. + + Without this the first session open failed inside a run's retry loop - ten seconds of retries, + a message that could only say "incomplete", and a report left `incomplete` for later fleet runs + to stop at. The import is the adapter's own, so a module that is present but cannot load counts + as missing too. + """ + missing = [] + for dialect in sorted(set(dialects)): + module, extra = DRIVERS[dialect] + try: + importlib.import_module(module) + except ImportError: + missing.append(extra) + if missing: + raise DemoRefused( + "Install the drivers for this demo's engines: " + f"pip install 'smart-data-engine-sdk[{','.join(missing)}]'" + ) + + def engine(dialect: str, dsn: str) -> Any: if dialect == "postgres": from sde.engines.postgres import PostgresEngine @@ -214,6 +250,7 @@ def setup(root: Path, supplied: dict[str, Any], admin_dsns: Mapping[str, str]) - from .resources import allocate, grant_tables settings, placement = bootstrap(supplied) + require_drivers(supplied["engines"].values()) logical = model() with transaction(root): if (root / "reset-request.json").exists(): diff --git a/python/src/sde_demo/runtime.py b/python/src/sde_demo/runtime.py index 8b8b55f..7edc86a 100644 --- a/python/src/sde_demo/runtime.py +++ b/python/src/sde_demo/runtime.py @@ -24,7 +24,16 @@ model, reading, ) -from .project import DemoRefused, config, credentials, engine, public_keys, read, write +from .project import ( + DemoRefused, + config, + credentials, + engine, + public_keys, + read, + require_drivers, + write, +) T = TypeVar("T") FLEET_PAGE = 100 @@ -36,6 +45,19 @@ ALERT_PAGE = 100 +def _failed_before_writing(report: Mapping[str, Any]) -> bool: + """A run that stopped, failed, before it recorded its first write: it wrote no row.""" + return ( + report.get("status") == "incomplete" + and report.get("pending", False) is None + and report.get("generator_id") == GENERATOR_ID + and all( + type(report.get(field)) is int and report[field] == 0 + for field in ("acknowledged_rows", "verified_after_uncertain_rows", "verified_rows") + ) + ) + + def fleet_runs(root: Path, *, project_id: str) -> dict[str, int]: """Every earlier run of this project, by run ID, with the rows it verified. @@ -44,6 +66,12 @@ def fleet_runs(root: Path, *, project_id: str) -> dict[str, int]: generated values depend only on the sequence number, and each run's local report says how many rows it verified. A run that did not complete - interrupted, or with an uncertain batch - makes that sum unknowable, so the workload refuses instead of comparing a read with a guess. + + One unfinished run is known exactly: one that failed before its first write. A run records the + batch it is about to write before writing it, so a failed run with no pending batch and no + acknowledged or verified row wrote nothing, and it adds nothing to the sum. Without this, a + first run that could not open its session - a missing driver - stopped every later fleet run + in that directory. """ runs: dict[str, int] = {} for path in sorted((root / "runs").glob("*/report.json")): @@ -51,6 +79,10 @@ def fleet_runs(root: Path, *, project_id: str) -> dict[str, int]: if report.get("project_id") != project_id: continue identity = report.get("run_id") + if _failed_before_writing(report) and isinstance(identity, str) and ( + path.parent.name == identity + ): + continue if ( not isinstance(identity, str) or re.fullmatch(r"[0-9a-f]{32}", identity) is None @@ -137,6 +169,7 @@ def run( if workload not in WORKLOADS: raise DemoRefused("Choose the mixed, point, analytics, fleet or alerts workload.") settings = config(root) + require_drivers(binding["dialect"] for binding in settings["engines"].values()) # 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 {} logical, keys = model(), public_keys(settings["public_keys"]) diff --git a/python/tests/test_demo_starter.py b/python/tests/test_demo_starter.py index c4a8330..bd4c4c2 100644 --- a/python/tests/test_demo_starter.py +++ b/python/tests/test_demo_starter.py @@ -589,3 +589,119 @@ def test_the_cli_and_the_runtime_offer_the_same_workloads(tmp_path: Path) -> Non cli.main(["--directory", str(tmp_path), "run", "--workload", "bogus"]) with pytest.raises(project.DemoRefused, match="fleet"): runtime.run(tmp_path, workload="bogus") + + +def test_fleet_skips_a_run_that_failed_before_its_first_write(tmp_path: Path) -> None: + """A run that could not open its first session wrote nothing, so it adds nothing to the sum. + + It was found on the installed Weather acceptance of 26 September. A TypeScript run without the + `pg` package failed before its first write, and its `incomplete` report then stopped every + later fleet run in that directory, with no command to settle it. + """ + _report(tmp_path, "a" * 32) + _report( + tmp_path, "b" * 32, status="incomplete", failure="EngineError", acknowledged_rows=0, + verified_after_uncertain_rows=0, verified_rows=0, + ) + assert runtime.fleet_runs(tmp_path, project_id="p" * 32) == {"a" * 32: 4} + + +@pytest.mark.parametrize( + "fields", + [ + {"acknowledged_rows": 10}, + {"verified_after_uncertain_rows": 10}, + {"verified_rows": 3}, + {"pending": {"first": 1, "count": 10}}, + {"status": "running"}, + {"acknowledged_rows": None}, + {"generator_id": "weather-v0:other"}, + ], + ids=["acknowledged", "resolved", "verified", "pending", "running", "unsaid", "generator"], +) +def test_fleet_still_refuses_an_unfinished_run_that_may_have_written( + tmp_path: Path, fields: dict[str, Any] +) -> None: + before_writing = { + "status": "incomplete", "acknowledged_rows": 0, "verified_after_uncertain_rows": 0, + "verified_rows": 0, + } + _report(tmp_path, "a" * 32, **{**before_writing, **fields}) + with pytest.raises(project.DemoRefused, match="every earlier run"): + runtime.fleet_runs(tmp_path, project_id="p" * 32) + + +def test_setup_names_why_the_bootstrap_map_does_not_load_before_writing( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + """Without the 'signed' extra the map cannot be verified: said in full, and nothing written. + + It was found on the installed acceptance of 26 September, where the starter answered + "Operation incomplete ... uncertain writes must not be replayed" for a setup that had written + nothing, and hid the one sentence that said what to install. + """ + import sys + + bundle, _ = supplied() + monkeypatch.setitem(sys.modules, "cryptography.hazmat.primitives.asymmetric.ed25519", None) + with pytest.raises(project.DemoRefused, match=r"map does not load: .*'signed' extra"): + project.setup(tmp_path / "refused", bundle, {}) + assert not (tmp_path / "refused").exists() + + +def test_setup_refuses_a_public_key_that_is_not_base64(tmp_path: Path) -> None: + bundle, _ = supplied() + bundle["public_keys"] = {name: "not base64!" for name in bundle["public_keys"]} + with pytest.raises(project.DemoRefused, match="Invalid public key configuration"): + project.setup(tmp_path / "refused", bundle, {}) + assert not (tmp_path / "refused").exists() + + +def _without(monkeypatch: pytest.MonkeyPatch, module: str) -> None: + """As if `module` were not installed, for the driver check's import (and nothing else).""" + import importlib + + real = importlib.import_module + + def imported(name: str, *args: Any) -> Any: + if name == module: + raise ImportError(f"No module named {module!r}") + return real(name, *args) + + monkeypatch.setattr(importlib, "import_module", imported) + + +def test_setup_refuses_a_missing_driver_before_writing( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + bundle, _ = supplied() + _without(monkeypatch, "psycopg") + with pytest.raises(project.DemoRefused, match=r"smart-data-engine-sdk\[postgres\]"): + project.setup(tmp_path / "refused", bundle, {}) + assert not (tmp_path / "refused").exists() + + +def test_a_run_refuses_a_missing_driver_before_its_report_exists( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + """The refusal comes before the run's report, so no `incomplete` run is left to stop a fleet.""" + local(tmp_path) + _without(monkeypatch, "psycopg") + monkeypatch.setattr( + sde.Session, "connect", lambda *a, **k: pytest.fail("a session opened without a driver") + ) + with pytest.raises(project.DemoRefused, match=r"smart-data-engine-sdk\[postgres\]"): + runtime.run(tmp_path, iterations=1, batch_size=1, interval_ms=0) + assert not (tmp_path / "runs").exists() + + +def test_the_driver_check_imports_each_dialects_module(monkeypatch: pytest.MonkeyPatch) -> None: + import importlib + + imported: list[str] = [] + monkeypatch.setattr(importlib, "import_module", imported.append) + project.require_drivers(["clickhouse", "postgres", "postgres"]) + assert imported == ["clickhouse_connect", "psycopg"] + _without(monkeypatch, "clickhouse_connect") + with pytest.raises(project.DemoRefused, match=r"\[clickhouse\]"): + project.require_drivers(["clickhouse", "postgres"]) diff --git a/typescript/src/demo/weather.ts b/typescript/src/demo/weather.ts index 409de29..9b3a96a 100644 --- a/typescript/src/demo/weather.ts +++ b/typescript/src/demo/weather.ts @@ -43,6 +43,17 @@ function recoverable(error: unknown) { function same(actual: Row | null, expected: Row) { if (!isDeepStrictEqual(actual, expected)) throw new DemoRefused('A logical read did not match this run.') } +/** + * A run that stopped, failed, before it recorded its first write: it wrote no row. A run records the + * batch it is about to write before writing it, so no pending batch and no acknowledged or verified + * row means nothing reached the engine. Without this, a first run that could not open its session - + * a missing driver - stopped every later fleet run in that directory. + */ +function failedBeforeWriting(report: Record): boolean { + return report.status === 'incomplete' && report.pending === null && report.generator_id === generatorId && + (['acknowledged_rows', 'verified_after_uncertain_rows', 'verified_rows'] as const) + .every(field => report[field] === 0) +} /** * Every earlier run of this project, by run ID, with the rows it verified. * @@ -50,7 +61,8 @@ function same(actual: Row | null, expected: Row) { * the same timeline, so the exact answer is a sum over the runs - known exactly, because generated * values depend only on the sequence number and each run's local report says how many rows it * verified. A run that did not complete makes that sum unknowable, so the workload refuses instead - * of comparing a read with a guess. The same rules as the Python starter's fleet_runs. + * of comparing a read with a guess - except a run that failed before its first write, which wrote no + * row (failedBeforeWriting). The same rules as the Python starter's fleet_runs. */ export function fleetRuns(root: string, projectId: string): Map { const runs = new Map() @@ -62,6 +74,7 @@ export function fleetRuns(root: string, projectId: string): Map const report = read(path).value if (report.project_id !== projectId) continue const identity = report.run_id, rows = report.verified_rows + if (failedBeforeWriting(report) && identity === name) continue if (typeof identity !== 'string' || !/^[0-9a-f]{32}$/.test(identity) || name !== identity || report.status !== 'complete' || report.pending !== null || report.generator_id !== generatorId || typeof rows !== 'number' || !Number.isSafeInteger(rows) || rows < 0 || rows > maxRunRows) { @@ -111,6 +124,20 @@ export function alertExpected(runId: string, through: number, limit: number): [R } return [page, total, total ? `${cents / 100n}.${String(cents % 100n).padStart(2, '0')}` : null] } +/** + * Refuse before a run's report exists when a binding's driver cannot be imported. Without this the + * first session open failed inside the run's retry loop: ten seconds of retries, a message that could + * only say the operation did not complete, and a report left `incomplete`. ClickHouse is spoken over + * fetch and needs no package. The import is the adapter's own, resolved from this package. + */ +export async function requireDrivers( + dialects: Iterable, load: (name: string) => Promise = name => import(name), +): Promise { + if (![...dialects].includes('postgres')) return + try { await load('pg') } catch { + throw new DemoRefused("The PostgreSQL binding needs the 'pg' package: npm install pg") + } +} 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) || @@ -118,6 +145,7 @@ export async function runWeather(root: string, options: RunOptions = {}): Promis intervalMs < 0 || intervalMs > 1000 || recoveryMs < 0 || recoveryMs > 30000 || !(workloads as readonly string[]).includes(workload)) throw new DemoRefused('Invalid bounded workload options.') const { model, bindings, projectId, activeMap } = project(root) + await requireDrivers(Object.values(bindings).map(binding => binding.dialect)) // Read before this run's own report exists; a run that starts later is not in the window. const earlier = workload === 'fleet' ? fleetRuns(root, projectId) : new Map() const factories: Record ManagedEngine> = {} diff --git a/typescript/tests/weather.test.ts b/typescript/tests/weather.test.ts index cd2cb3b..6905630 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 { alertExpected, alertHumidity, fleetExpected, fleetRuns, runWeather, workloads } from '../src/demo/weather.js' +import { alertExpected, alertHumidity, fleetExpected, fleetRuns, requireDrivers, 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 }) }) @@ -244,3 +244,40 @@ it('offers the fleet and alerts workloads and refuses anything else before readi expect(workloads).toEqual(['mixed', 'point', 'analytics', 'fleet', 'alerts']) await expect(runWeather('/nonexistent', { workload: 'bogus' as never })).rejects.toThrow(/Invalid bounded workload/) }) + +it('skips a run that failed before its first write, and still refuses one that may have written', () => { + // Found on the installed Weather acceptance of 26 September: a run without the `pg` package failed + // before writing, and its `incomplete` report stopped every later fleet run in that directory. + const beforeWriting = { status: 'incomplete', failure: 'EngineError', acknowledged_rows: 0, + verified_after_uncertain_rows: 0, verified_rows: 0 } + const { root } = fixture() + report(root, 'a'.repeat(32)); report(root, 'b'.repeat(32), beforeWriting) + expect(fleetRuns(root, 'p'.repeat(32))).toEqual(new Map([['a'.repeat(32), 4]])) + for (const change of [{ acknowledged_rows: 10 }, { verified_after_uncertain_rows: 10 }, { verified_rows: 3 }, + { pending: { first: 1, count: 10 } }, { status: 'running' }, { acknowledged_rows: null }, + { generator_id: 'weather-v0:other' }]) { + const { root: other } = fixture() + report(other, 'a'.repeat(32), { ...beforeWriting, ...change }) + expect(() => fleetRuns(other, 'p'.repeat(32)), JSON.stringify(change)).toThrow(/every earlier run/) + } +}) +it('checks the pg package only for a PostgreSQL binding, and names the install when it is missing', async () => { + const loaded: string[] = [] + await requireDrivers(['clickhouse'], async name => { loaded.push(name) }) + expect(loaded).toEqual([]) + await requireDrivers(['clickhouse', 'postgres'], async name => { loaded.push(name) }) + expect(loaded).toEqual(['pg']) + await expect(requireDrivers(['postgres'], async () => { throw new Error("Cannot find package 'pg'") })) + .rejects.toThrow("needs the 'pg' package: npm install pg") +}) +it('refuses a run without the pg package before its report exists', async () => { + const { root } = fixture() + vi.doMock('pg', () => { throw new Error("Cannot find package 'pg'") }) + const connect = vi.spyOn(Session, 'connect') + try { + await expect(runWeather(root, { iterations: 1, batchSize: 1, intervalMs: 0 })) + .rejects.toThrow("needs the 'pg' package") + } finally { vi.doUnmock('pg') } + expect(connect).not.toHaveBeenCalled() + expect(() => readdirSync(join(root, 'runs'))).toThrow() +})