Skip to main content
Glama
J-X0
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.*