reflexr.fastapi¶
The fastapi extra. See Serving over REST and WebSocket.
FastAPI adapters: HTTP ingest, REST reads and administration, and the WebSocket stream.
Mount the router under a prefix of your choosing::
app.include_router(reflexr_router(workspaces, resolve_actor=resolve_actor), prefix="/v1")
Every route speaks the workspace protocol in docs/protocol.md, and every command goes
through reflexr.workspace.Workspaces.execute, once per command_id, so REST, the
WebSocket and MCP behave identically.
The router¶
reflexr_router
¶
reflexr_router(
workspaces: Workspaces,
*,
resolve_actor: ResolveActor,
authorize: Authorize | None = None,
hello_timeout: float = 10.0,
outbox_size: int = 1000,
) -> APIRouter
Build the router for publishing, reads, administration and the WebSocket stream.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
workspaces
|
Workspaces
|
Opens tenant-scoped workspaces, holds the rules and schedules, and carries
out commands, once per |
required |
resolve_actor
|
ResolveActor
|
Authenticates each request and connection. |
required |
authorize
|
Authorize | None
|
Whether an actor may use a workspace; allows everything if omitted. |
None
|
hello_timeout
|
float
|
Seconds a new connection has to send |
10.0
|
outbox_size
|
int
|
Frames buffered for a slow connection before it is closed (4429). |
1000
|
Each REST request's span (FastAPI's own, when it is instrumented) and each connection's
reflexr.stream span are attributed to the tenant, workspace and actor.
ResolveActor
module-attribute
¶
Authenticates a request or connection: returns its tenant and actor, or raises Unauthorized.
Unauthorized
¶
Bases: Exception
Raise from resolve_actor to refuse a request (401) or connection (close 4401).
STATUS_CODES
module-attribute
¶
STATUS_CODES: dict[str, int] = {
"not_found": 404,
"invalid_state": 409,
"validation_failed": 422,
"depth_exceeded": 422,
"forbidden": 403,
"unsupported_protocol": 400,
}
The HTTP status for each rejection type.
Bodies and responses¶
The router's authorize hook, Authorize, the statuses it returns, RuleStatus and ScheduleStatus, and the rule it reads, WorkspaceRule, are the workspace layer's, shared by every surface.
PublishItem
pydantic-model
¶
PublishBatch
pydantic-model
¶
Bases: _Body
Events to publish atomically and in order, optionally joining a causal chain.
Fields:
-
events(list[PublishItem]) -
correlation_id(str | None)
The stream¶
One WebSocket connection, which the router serves at /workspaces/{workspace_id}/stream. It hands each command frame to Workspaces.execute, as POST .../commands does, so a command is carried out once per command_id whichever of the two it arrives on.
Stream
¶
Stream(
websocket: WebSocket,
workspace_id: WorkspaceId,
*,
workspaces: Workspaces,
open_workspace: OpenWorkspace,
hello_timeout: float,
outbox_size: int,
)
Serves one connection: handshake, replay then live events, and commands.
One task reads frames; commands run as their own tasks, so a slow command never delays the others. Every outgoing frame goes through one bounded outbox and one writer; a client too slow to keep up is disconnected (4429) rather than holding events back.
The connection is traced as a reflexr.stream span, and counted in
reflexr.stream.connections while it is open and reflexr.stream.disconnects, by close
code, when it ends.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
websocket
|
WebSocket
|
The connection. |
required |
workspace_id
|
WorkspaceId
|
The workspace it follows. |
required |
workspaces
|
Workspaces
|
Carries out its commands, once per |
required |
open_workspace
|
OpenWorkspace
|
Authenticates the connection and opens its workspace. |
required |
hello_timeout
|
float
|
Seconds the client has to send |
required |
outbox_size
|
int
|
Frames buffered for the client before it is closed (4429). |
required |