ADR-0011: Surfaces: ingest, REST, WebSocket, schedules and MCP¶
Status: Superseded by ADR-0044 Date: 2026-09-28 Deciders: Alex Nodeland
Context¶
Events come from other services (over HTTP, in batches, from webhooks), from people and dashboards (live, over WebSocket), from timetables, and from external agents (over MCP). Operators need to see and fix what rules and runs are doing. Reflex offered HTTP and WebSocket ingest and a dead-letter endpoint, with rate limiting and health checks, each with its own handling code.
Decision¶
v0.1 ships four surfaces, all thin adapters over one command handler on a Stream handle, so they behave identically:
- HTTP ingest and REST (
reflexr.fastapi): publish one event or a batch, idempotent by id; read the log; inspect and administer rules (cursor, lag, replay), runs (retry, skip, cancel) and dead letters. - WebSocket (
reflexr.fastapi): a resumable subscription with artifactr's handshake (hello, replay,replay_complete), filters by event type, the same command frames, and the same close codes. - Schedules (
reflexr.stream): interval and cron schedules that publish into streams. Tick ids derive from the schedule and the tick's time, so racing replicas publish each tick once; a lease keeps them from racing. - MCP (
reflexr.mcp): publishing, reading and administration as tools, and runs as resources with tenant-qualified URIs and update notifications.
Authentication is the host's: each surface takes a resolver that returns the tenant and actor. Rate limiting and health checks are the application's (FastAPI middleware and its own routes), not the library's.
Amendment (2026-09-29): every surface takes the authorize hook¶
ReflexrMcp(authorize=...)takes the router's hook,authorize(tenant_id, workspace_id, actor) -> bool, and asks it on every tool call and resource read that names a workspace, before the workspace is opened. A refusal is a tool error, or a failed resource read, carrying theforbiddenrejection's message, as every other rejection over MCP carries its message.list_rulesnames no workspace and is not asked. Without the hook, an MCP client could use any workspace of its tenant however the application restricted REST and the WebSocket.- The hook's type,
Authorize, lives inreflexr.workspace, besideWorkspaces.open, since adapters may not import each other (ADR-0025).reflexr.fastapi.Authorizeis the same alias, re-exported.
Amendment (2026-09-29): MCP reports what REST reports, and checks subscriptions¶
MCP's rule_status computed its own listing from rule_progress(), so it disagreed with REST's; MCP had no schedule status; and a client could listen for any tenant's run notifications.
- What a surface reports, the workspace computes.
RuleStatusandScheduleStatusmove fromreflexr.fastapitoreflexr.workspace, andWorkspace.rule_statuses()andWorkspace.schedule_statuses()compute them from the registered rules and schedules. REST returns them unchanged, with the same JSON, and MCP formats them as text, so neither adapter computes a status itself (ADR-0025).reflexr.fastapistill exports both, as the same classes. rule_statuslists every registered rule, asGET /workspaces/{id}/rulesdoes: whether it is enabled, and its cursor, lag, generation and dead letters. It had listed every rule with a cursor, so it showed rules no longer registered and left out registered rules that had never evaluated the workspace, such as a disabled rule that had never run.schedule_status, a new tool, lists each schedule that targets the workspace with its last and next tick, asGET /workspaces/{id}/schedulesdoes. It asksauthorize, as every tool that names a workspace does, and a test fails if a new tool is not asked.list_runstakesscope_key, asGET /workspaces/{id}/runsdoes. It holds the scope's values, since the MCP SDK parses a JSON array given as a string before validation, andreflexr.core.scope_keybuilds the key from them.- A subscription is checked when it opens. The MCP SDK serves
subscriptions/listenitself, so a middleware, passed throughMCPServer(middleware=...), makes a resource read's checks for each run URI a listen request names: that the run is the client's tenant's, and thenauthorize. A refusal fails the request withINVALID_PARAMSand the read's message; a request that names no run is not checked, andresolveis not called for it. Until now any client could listen for notifications on any tenant's runs, learning which moved and when. URIs are matched asrun_uriwrites them, which is what notifications name. The SDK calls its middleware provisional; its version is pinned inuv.lock, and the tests exercise the check through the SDK's client. artifactr made the same change in its ADR-0012. - Deliberate differences remain. MCP's
read_eventsandlist_runsreturn 50 and 20 by default, for a model's context, where REST returns everything unless given alimit. Batch publishing and command ids for deduplication are REST's and the WebSocket's; over MCP an event'sidmakes a publish idempotent.GET /scheduleshas no MCP tool yet: its narrowing of targets to the caller's tenant would move to the workspace layer first.
Amendment (2026-09-29): windows, type filters and tails of the log, and joining at its head¶
Workspace.read had no type filter and no tail. So REST's type=, MCP's read_events, EventContext.read_events and oncall's runbook all read the whole log and filtered it in Python (#46). A WebSocket client that wanted only what happens next still received the whole replay. artifactr shares the protocol's shape, and made the same additions with the same names in its ADR-0005.
- A read is a window, a filter and a page.
Storage.read,Workspace.read,GET /workspaces/{id}/eventsand MCP'sread_eventstake:after_seqand a newbefore_seq, for the windowafter_seq < seq < before_seq;types, which REST spells as a repeatedtype;limit(the first so many) or a newlast(the last so many, still oldest first).
- Using the tail.
lastreads the tail, andlastwithbefore_seqset to the oldestseqa client has pages backwards. Every read returns the log in order, whichever end the page comes from, so a client handles every page the same way. - Checking arguments. A read with both
limitandlast, or with a negative number, isvalidation_failed.Workspace.readchecks this once, so storage adapters don't have to. - Storage filters, not the handle. The port takes the new parameters, so each adapter pushes them down. SQL storage keeps each event's type in a new
event_typecolumn, besideevent_id. Migration 0003 fills the column from the stored envelopes and indexes it with the workspace andseq. SQL storage reads the last so many backwards from the end.Transaction.read, which evaluation uses, is unchanged. - Every reader passes its filter down: REST, MCP,
EventContext.read_events(which also takesbefore_seqandlast, capped atread_limit) and the runbook'sdiagnose. - Defaults. MCP still returns 50 when given neither
limitnorlast, andEventContext.read_eventsreturns 20. hellotakesfrom_head. Withfrom_head: true, a connection starts athead_seq: nothing is replayed, andreplay_completefollowswelcomeat once.core.resumedecides this. It refusesfrom_headwith a nonzeroresume_after_seq, and the stream closes with 4400.- A client that reconnects resumes from its position, so it cannot skip what it missed by asking for the head again.
- It is a new optional field, so the protocol stays
reflexr.v1.
- The stream has no tail of its own. A tail read over REST, followed by
hellowithresume_after_seqset to the tail's lastseq, shows recent history and then everything after it, with no gap. Sohelloneeds only a place to start, not a count.
Options considered¶
| Option | Consistency across surfaces | Scope for v0.1 |
|---|---|---|
| All four over one command handler (chosen) | By construction | Largest |
| HTTP ingest and REST only | Trivial | Smallest; operators and agents get no live view or tools |
| Per-surface handlers, as in Reflex | Drifts | Medium |
Consequences¶
- Easier: a new surface is authentication plus translation into commands.
- Easier: one JSON Schema describes every frame, for clients in any language.
- Harder: four surfaces to test to 100% coverage, and a protocol to version.
Action items¶
- Implement the schedule runner (RFC-0001 phase 2), and
reflexr.fastapiandreflexr.mcpwith generated schemas (phase 5).