diff --git a/TASKS.md b/TASKS.md index 549c2f5..15ddbf7 100644 --- a/TASKS.md +++ b/TASKS.md @@ -8,6 +8,7 @@ Status: in_progress DoD: - `__REFACTORING.md` documents the full remediation plan for registry modularization, sqlx compile-time validation, KDF correction, typed timestamps, ID display support, pool tuning, and error normalization +- `__REFACTORING.md` also documents correlation IDs, rate limiting/backpressure, and a distributed `mcp-server` session store - `.gitignore` ignores local `__*.md` working notes and refactoring specs - `TASKS.md` points to the refactoring track as the active work item @@ -24,5 +25,8 @@ DoD: - `feat/id-display-support` - `feat/postgres-pool-config` - `feat/error-structure-normalization` +- `feat/correlation-id-propagation` +- `feat/runtime-rate-limiting-and-backpressure` +- `feat/distributed-mcp-session-store` - `feat/websocket-test-run-polish` - `feat/live-authenticated-staging-entry` diff --git a/__REFACTORING.md b/__REFACTORING.md index 46a9dbe..66c7739 100644 --- a/__REFACTORING.md +++ b/__REFACTORING.md @@ -13,7 +13,10 @@ The refactoring scope covers: 5. adding `Display` support to `define_id!`; 6. introducing explicit PostgreSQL pool configuration; 7. removing duplicated SQL timestamp formatting; -8. making runtime and API error payloads more structured. +8. making runtime and API error payloads more structured; +9. introducing correlation IDs across transport, runtime, and adapters; +10. adding rate limiting and backpressure controls; +11. replacing in-memory MCP session state with a multi-instance-safe store. ## Guiding Rules @@ -879,6 +882,401 @@ or dedicated variants such as: - API can map key failure classes deterministically; - logs and telemetry can group errors by category instead of by ad hoc message text. +## Track I: Correlation IDs and Cross-Layer Tracing + +### Problem + +Today, a request can enter through: + +- `admin-api` HTTP routes; +- `mcp-server` JSON-RPC over Streamable HTTP; +- internal runtime execution; +- protocol adapters (`REST`, `gRPC`, `WebSocket`, `SOAP`) + +without a single stable correlation key that survives the entire path. + +There are already partial pieces in the codebase: + +- `request_id` exists in observability models and storage; +- MCP mappings can inject `$.mcp.correlation_id`; +- some protocol fixtures mention `correlation_id`. + +But there is no guaranteed end-to-end propagation model. When an adapter fails, logs cannot always be tied back to: + +- the incoming HTTP request; +- the MCP session; +- the runtime operation execution; +- the invocation log row. + +### Goal + +Introduce a first-class correlation model that is generated once at ingress and propagated across: + +- transport; +- service/runtime orchestration; +- adapter calls; +- persistence and logs. + +### Target Model + +Introduce two related identifiers: + +- `request_id` + - unique per incoming HTTP or MCP request; + - used for ingress/egress tracing and log correlation. +- `correlation_id` + - stable business-level identifier propagated into upstream payloads when configured; + - may equal `request_id` by default, but remains distinct conceptually. + +### Flow Rules + +#### `admin-api` + +At HTTP ingress: + +- read `X-Request-Id` if supplied and valid; +- otherwise generate a UUID v7 request ID; +- store it in request extensions; +- include it in structured logs; +- pass it into runtime execution context; +- persist it into `invocation_logs.request_id` when a test run or runtime execution creates a log row. + +#### `mcp-server` + +At JSON-RPC ingress: + +- derive a request-scoped `request_id`; +- bind it to the current transport session and JSON-RPC request; +- propagate it into runtime execution and adapter context; +- return it in logs and diagnostics. + +#### Runtime + +Add an execution context struct, for example: + +```rust +pub struct ExecutionContext { + pub request_id: String, + pub correlation_id: String, + pub workspace_id: WorkspaceId, + pub agent_id: Option, + pub operation_id: OperationId, + pub mcp_session_id: Option, +} +``` + +This context must be passed through: + +- unary execution; +- window execution; +- session tool start/poll/stop; +- async job start/status/result/cancel. + +#### Adapters + +Each adapter must receive execution context and apply it consistently: + +- `REST` + - add `X-Request-Id` + - add `X-Correlation-Id` + - allow explicit mapping into headers/query/body where configured +- `gRPC` + - attach metadata keys such as `x-request-id`, `x-correlation-id` +- `WebSocket` + - inject correlation metadata into subscribe payloads when configured +- `SOAP` + - inject correlation value into configured header/body mappings + +### Files and Functions + +#### New or expanded files + +- `apps/admin-api/src/request_context.rs` +- `apps/mcp-server/src/request_context.rs` +- `crates/crank-runtime/src/context.rs` +- adapter request builders in each protocol crate + +#### Functions to add or refactor + +`admin-api` + +- `extract_or_generate_request_id(...)` +- `request_context_middleware(...)` +- `runtime_request_context(...)` + +`mcp-server` + +- `extract_mcp_request_context(...)` +- `runtime_request_context_from_message(...)` + +`crank-runtime` + +- `ExecutionContext` +- `PreparedRequest` / `RuntimeOperation` updates to carry context +- `execute_*` functions accept `&ExecutionContext` + +### Acceptance Criteria + +- every runtime invocation has a request ID; +- every adapter call can log the request ID; +- invocation logs persist the same request ID; +- operator can trace a failing adapter call back to the original HTTP or MCP request. + +## Track J: Rate Limiting and Backpressure + +### Problem + +There is currently no explicit protection layer for: + +- excessive API request rates; +- excessive tool invocation rates; +- excessive stream/session polling; +- unbounded adapter output pressure. + +This is a production-readiness gap, especially now that the platform supports: + +- unary execution; +- windowed streaming; +- session tools; +- async jobs; +- WebSocket and SSE upstream adapters. + +### Goal + +Introduce explicit controls for request rate, concurrency, and bounded output pressure. + +### Scope + +#### API-level controls + +Apply to: + +- `admin-api` routes; +- `mcp-server` transport endpoints. + +Dimensions: + +- per-IP +- per-session +- per-workspace +- per-platform-key + +#### Runtime-level controls + +Apply to: + +- concurrent unary executions; +- concurrent window executions; +- concurrent live stream sessions; +- concurrent async jobs; +- polling frequency for session/job tools. + +#### Adapter-level backpressure + +Apply to: + +- REST SSE +- gRPC server-streaming +- WebSocket windows/sessions + +Controls: + +- max buffered items +- max buffered bytes +- max window duration +- max session polls +- max per-message bytes + +### Target Design + +#### `admin-api` + +Add middleware with a keyed limiter, for example: + +- by authenticated user ID; +- by workspace ID; +- by remote IP for unauthenticated endpoints. + +#### `mcp-server` + +Add a limiter keyed by: + +- platform key ID; +- workspace slug; +- MCP session ID; +- optionally agent slug. + +#### `crank-runtime` + +Add a runtime guard layer, for example: + +```rust +pub struct RuntimeLimits { + pub max_concurrent_unary: usize, + pub max_concurrent_window: usize, + pub max_concurrent_sessions: usize, + pub max_concurrent_jobs: usize, + pub max_polls_per_session: u32, +} +``` + +And an execution gate: + +- acquire permit before execution; +- release permit on completion; +- reject with stable error code if limits are exceeded. + +### Implementation Slices + +##### Slice J1: API request throttling + +- add middleware and config for `admin-api`; +- add middleware and config for `mcp-server`. + +##### Slice J2: Runtime concurrency limits + +- add semaphores/permits around execution paths. + +##### Slice J3: Streaming/session pressure limits + +- cap poll frequency and active session count; +- cap pending buffer sizes in adapters. + +### Acceptance Criteria + +- platform can reject overload deterministically instead of degrading unpredictably; +- stream/session polling is bounded; +- slow clients cannot force unbounded buffering in streaming adapters. + +## Track K: Distributed MCP Session Store + +### Problem + +`apps/mcp-server/src/session.rs` currently keeps transport/session state in process memory through `Arc>>`. + +This is acceptable only for a single instance. It breaks when: + +- multiple `mcp-server` instances run behind a load balancer; +- a container restarts and loses in-memory session state; +- GET/POST/DELETE transport requests land on different instances. + +### Goal + +Replace the in-memory session store with a multi-instance-safe shared store. + +### Target Design + +Introduce a `SessionStore` trait: + +```rust +#[async_trait] +pub trait SessionStore { + async fn create_transport_session(&self, session: TransportSessionRecord) -> Result<...>; + async fn get_transport_session(&self, session_id: &str) -> Result<...>; + async fn touch_transport_session(&self, session_id: &str, now: OffsetDateTime) -> Result<...>; + async fn delete_transport_session(&self, session_id: &str) -> Result<...>; + async fn append_pending_message(&self, session_id: &str, message: Value) -> Result<...>; + async fn drain_pending_messages(&self, session_id: &str) -> Result, ...>; + async fn cleanup_expired_transport_sessions(&self, now: OffsetDateTime) -> Result; +} +``` + +Provide two implementations: + +- `InMemorySessionStore` + - test-only and local dev fallback +- `PostgresSessionStore` + - production default + +### Storage Model + +Add tables such as: + +- `mcp_transport_sessions` +- `mcp_transport_outbox` + +Suggested fields: + +`mcp_transport_sessions` + +- `id` +- `workspace_id` +- `agent_id` +- `platform_api_key_id` +- `transport_kind` +- `status` +- `created_at` +- `updated_at` +- `expires_at` + +`mcp_transport_outbox` + +- `id` +- `session_id` +- `sequence_no` +- `payload_json` +- `created_at` + +### Semantics + +#### POST initialize + +- create transport session record in shared store. + +#### GET SSE stream + +- attach to existing session record; +- stream messages from outbox or polling bridge; +- support reconnect without requiring process-local state. + +#### POST JSON-RPC + +- resolve session from shared store; +- append outbound messages to outbox. + +#### DELETE + +- mark session closed in shared store; +- best-effort cleanup of outbox rows. + +### Files and Modules + +Add: + +- `apps/mcp-server/src/session/mod.rs` +- `apps/mcp-server/src/session/store.rs` +- `apps/mcp-server/src/session/in_memory.rs` +- `apps/mcp-server/src/session/postgres.rs` + +Potential shared data types: + +- `crates/crank-core/src/transport_session.rs` + +### Migration Order + +##### Slice K1: trait and in-memory implementation + +- preserve current behavior but hide it behind a trait. + +##### Slice K2: persistent schema and Postgres store + +- add migrations and registry/storage support. + +##### Slice K3: transport handlers on shared store + +- wire GET/POST/DELETE through the trait-backed store. + +##### Slice K4: cleanup and failover semantics + +- cleanup job for expired transport sessions; +- reconnect tests across store-backed state. + +### Acceptance Criteria + +- transport session state survives process restart if store is persistent; +- GET/POST/DELETE requests can hit different instances safely; +- in-memory session state is no longer the production bottleneck. + ## Delivery Order Recommended order: @@ -891,7 +1289,10 @@ Recommended order: 6. `feat/secret-crypto-hkdf` 7. `feat/typed-timestamps` 8. `feat/error-structure-normalization` -9. `feat/post-refactor-cleanup` +9. `feat/correlation-id-propagation` +10. `feat/runtime-rate-limiting-and-backpressure` +11. `feat/distributed-mcp-session-store` +12. `feat/post-refactor-cleanup` Reasoning: