eastbournestreaming
by J-X0
README.md
# eastbournestreaming
Streaming inference MCP server for utility meter telemetry. It folds a stream of
meter readings into fixed (tumbling) time windows with **exactly-once**
semantics and runs an inference provider over each closed window.
Project codename: Eastbourne. Client: Quenby Metering.
## The constraint that shapes the design
Regulated meter data cannot leave the customer's boundary. Two things follow,
and both are enforced by tests:
- **No network at runtime.** The server transport is stdio (newline-delimited
JSON-RPC 2.0), not a socket. `tests/boundary.test.ts` scans `src/` and fails
the build if any module imports `node:net`, `node:http(s)`, `node:tls`,
`node:dgram`, `WebSocket`, or calls `fetch`.
- **In-process inference.** All model behaviour goes through
`InferenceProvider` (`src/providers/base.ts`). The bundled
`StubInferenceProvider` (`src/providers/stub.ts`) is deterministic and runs
offline with no API key. A production model ships on-premises behind the same
interface.
## Exactly-once semantics
For a fixed set of distinct `eventId`s the emitted aggregates are identical
regardless of how often or in what order events are delivered, as long as a
duplicate arrives before its window closes. This holds because:
- duplicates are dropped via an `eventId` dedup set scoped to open windows, and
- window closure is a deterministic function of event-time via a watermark
(`maxEventTime - allowedLateness`), not wall-clock.
`checkpoint()` / `WindowedAggregator.restore()` carry the open windows and dedup
set across a restart so the guarantee survives process recycling.
## Architecture
```
stdin (JSON-RPC lines) ──► server.ts ──► WindowedAggregator ──► InferenceProvider
│ │ │
Logger watermark + StubInferenceProvider
(stderr) dedup + windows (deterministic, offline)
```
The server is a thin envelope: it parses one request per line, routes to the
aggregator, runs the provider over any windows that closed, and writes one
response per line to stdout. State lives in `WindowedAggregator`; all model
behaviour lives behind `InferenceProvider`. Nothing opens a socket — see the ADR
on the stdio transport. Design decisions with a real alternative are recorded in
[`docs/adr/`](docs/adr/).
## Layout
- `src/types.ts` — domain types (`Reading`, `WindowAggregate`, ...)
- `src/aggregator.ts` — `WindowedAggregator`, the core algorithm
- `src/providers/base.ts` — `InferenceProvider` interface
- `src/providers/stub.ts` — deterministic offline provider (z-score anomaly)
- `src/config.ts` — env-var configuration, validated at startup
- `src/logger.ts` — structured JSON logging to stderr, with `timed()`
- `src/errors.ts` — `ConfigError`, `WindowLimitError`
- `src/server.ts` — stdio JSON-RPC entry point wiring the pieces together
- `tests/*.test.ts` — behavioural tests, `node:test`
## Configuration
All configuration is read from environment variables at startup and validated;
a malformed value aborts with exit code 78 (`EX_CONFIG`) and a logged reason
rather than failing later on a bad window.
| Variable | Default | Meaning |
| --- | --- | --- |
| `EB_WINDOW_MS` | `60000` | Tumbling window size (ms), must be > 0 |
| `EB_LATENESS_MS` | `30000` | Allowed lateness (ms), must be >= 0 |
| `EB_MAX_OPEN_WINDOWS` | `100000` | Ceiling on concurrent open windows (back-pressure) |
| `EB_BASELINE_MEAN` | `1.0` | Expected healthy mean per reading |
| `EB_BASELINE_STDDEV` | `0.25` | Baseline standard deviation, must be > 0 |
| `EB_Z_THRESHOLD` | `3.0` | Absolute z-score to flag an anomaly, must be > 0 |
| `EB_LOG_LEVEL` | `info` | `debug` \| `info` \| `warn` \| `error` |
## Operational behaviour
- **Logs** are newline-delimited JSON on **stderr** (stdout is the JSON-RPC
channel). `flush` timing is logged at `debug`.
- **Back-pressure**: a reading that would open a window past
`EB_MAX_OPEN_WINDOWS` is refused with JSON-RPC error `-32000` (busy). The
stream stays alive; the client retries once windows drain. Bounds memory
against unbounded meters or lateness.
- **Graceful degradation**: if the inference provider throws, the affected
window is still emitted with `label: "degraded"` and the error is logged. A
model fault never stalls the pipeline or drops a computed aggregate.
- **Graceful shutdown**: on `SIGINT`/`SIGTERM` the server flushes open windows
so no aggregate is lost, logs the closed windows, and exits 0.
- **Bad input**: malformed readings return JSON-RPC `-32602` (invalid params);
unparseable lines return `-32700` and are dropped without crashing.
## Requirements
Node 22+. Node runs the TypeScript sources directly; there is no build step and
no `dist/`.
## Install and test
```sh
npm ci
npm test # runs tests/ with the built-in runner
npm run typecheck # tsc --noEmit
```
## Running the server
```sh
npm run server
```
It reads one JSON-RPC request per line on stdin and writes one response per line
on stdout. Pipe a session straight in (windows are 60s, allowed lateness 30s):
```sh
printf '%s\n' \
'{"jsonrpc":"2.0","id":1,"method":"tools/list"}' \
'{"jsonrpc":"2.0","id":2,"method":"ingest","params":{"reading":{"eventId":"e1","meterId":"m1","timestampMs":0,"value":1.0}}}' \
'{"jsonrpc":"2.0","id":3,"method":"flush"}' | npm run --silent server 2>/dev/null
```
The third line closes the window opened by the second and returns its inference:
```json
{"jsonrpc":"2.0","id":3,"result":{"inferences":[{"key":{"meterId":"m1","windowStartMs":0,"windowEndMs":60000},"label":"normal","score":0.047426,"reason":"mean=1 z=0 threshold=3"}]}}
```
Logs go to stderr; the `2>/dev/null` above hides them so only JSON-RPC responses
show. Drop it (or set `EB_LOG_LEVEL=debug`) to see the structured log.
`ingest` returns `{ rejected: "duplicate" | "late" | null, inferences: [...] }`.
`flush` force-closes all open windows and returns their inferences. `stats`
returns `{ openWindows }`. Inferences are produced only for windows that close.
## Known limitations / deferred work
- The dedup set is scoped to open windows; a duplicate that arrives *after* its
window has closed is treated as a late reading and dropped rather than being
recognised as a duplicate. For the supported lateness horizon this is
equivalent, but a longer-horizon dedup would need a persistent seen-id log.
- The stub provider is a single-feature z-score detector — enough to exercise
the boundary and the window pipeline, not a production anomaly model.
---
*Quenby Metering is an illustrative client; this repository is a self-directed reference implementation built to work end to end.*
This server cannot be deployed
Maintenance
ActivitySlowing
ResponsivenessNo issues