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
8 changes: 8 additions & 0 deletions docs/weather-starter.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`
Expand Down Expand Up @@ -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`
Expand Down
3 changes: 2 additions & 1 deletion python/src/sde_demo/diagnostics.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
47 changes: 42 additions & 5 deletions python/src/sde_demo/project.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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():
Expand Down
35 changes: 34 additions & 1 deletion python/src/sde_demo/runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.

Expand All @@ -44,13 +66,23 @@ 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")):
report = read(path)
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
Expand Down Expand Up @@ -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"])
Expand Down
116 changes: 116 additions & 0 deletions python/tests/test_demo_starter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"])
30 changes: 29 additions & 1 deletion typescript/src/demo/weather.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,14 +43,26 @@ 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<string, unknown>): 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.
*
* Fleet analytics reads every station's rows in one time window, and every run in a directory writes
* 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<string, number> {
const runs = new Map<string, number>()
Expand All @@ -62,6 +74,7 @@ export function fleetRuns(root: string, projectId: string): Map<string, number>
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) {
Expand Down Expand Up @@ -111,13 +124,28 @@ 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<string>, load: (name: string) => Promise<unknown> = name => import(name),
): Promise<void> {
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<RunReport> {
const { iterations = 10, batchSize = 10, intervalMs = 100, recoveryMs = 10000, workload = 'mixed' } = options
if (![iterations, batchSize, intervalMs, recoveryMs].every(Number.isSafeInteger) ||
iterations < 1 || iterations > 1000 || batchSize < 1 || batchSize > 1000 || iterations * batchSize > maxRunRows ||
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<string, number>()
const factories: Record<string, () => ManagedEngine> = {}
Expand Down
Loading
Loading