# Streaming MCP Plan ## 1. Назначение документа Этот документ фиксирует целевую архитектуру потоковой обработки в Crank. Он описывает: - какие transport- и upstream-протоколы поддерживаются; - какие streaming-сценарии считаются допустимыми; - какие ограничения обязательны для безопасности и управляемости; - какие сущности, API, UI-поля и runtime-механизмы нужно добавить; - в каком порядке это реализовывать. Документ синхронизирован с MCP transport specification `2025-06-18`, где `Streamable HTTP` определен как основной HTTP transport, а SSE допускается как часть `POST`-response и как отдельный `GET`-stream. ## 2. Базовое архитектурное решение Crank поддерживает streaming не как бесконечный текстовый поток в чат, а как управляемую инструментальную модель поверх MCP tools. Принцип: - transport может быть long-lived; - tool contract обязан оставаться ограниченным и управляемым; - upstream stream всегда преобразуется в bounded result, session poll или async job; - оператор в UI настраивает не "трубу", а режим сбора, агрегации и завершения. ## 3. Поддерживаемые transport- и upstream-протоколы ### 3.1. Downstream MCP transport Поддерживается: - `Streamable HTTP` как канонический network transport; - HTTP `POST` с ответом `application/json` или `text/event-stream`; - HTTP `GET` для server-to-client SSE stream; - несколько SSE streams одновременно в рамках одного MCP session; - `Mcp-Session-Id` и `MCP-Protocol-Version` headers; - resumable SSE streams как опциональная возможность. Не является текущим приоритетом: - `stdio` как обязательная часть продукта; - собственные нестандартные transport-режимы поверх MCP. ### 3.2. Upstream protocols Поддерживаются: - REST unary; - REST SSE в bounded режимах; - GraphQL `query` и `mutation`; - gRPC unary; - gRPC server-streaming в bounded режимах; - WebSocket в bounded режимах. Отложено: - GraphQL `subscription`; - gRPC client-streaming; - gRPC bidirectional streaming; - arbitrary websocket passthrough; - raw infinite stream forwarding в MCP client. ### 3.3. Protocol capability matrix | Protocol | Unary | Window | Session | Async Job | Notes | | --- | --- | --- | --- | --- | --- | | REST | Yes | Yes | Limited | Yes | SSE and long-poll sources are supported in controlled form | | GraphQL | Yes | No | No | Limited | `query` and `mutation` only; `subscription` is future scope | | gRPC | Yes | Yes | Yes | Yes | server-streaming only; client/bidi deferred | | WebSocket | No | Yes | Yes | Yes | upstream adapter only; not downstream MCP transport | | SOAP | Yes | Limited | Limited | Yes | primarily request-response enterprise workflows | ## 4. Поддерживаемые streaming modes Crank поддерживает четыре режима выполнения operation. ### 4.1. `unary` Обычный request-response вызов. Подходит для: - REST; - GraphQL `query` и `mutation`; - gRPC unary; - SOAP. ### 4.2. `window` Runtime открывает upstream stream или repeatedly polls upstream source, собирает данные в пределах окна и возвращает один bounded ответ. Параметры: - `window_duration_ms` - `max_items` - `max_bytes` - `upstream_timeout_ms` - `aggregation_mode` Подходит для: - логи за период; - метрики за период; - event window; - SSE stream snapshot; - gRPC server-stream window; - WebSocket event window. ### 4.3. `session` Runtime создает stream session, после чего данные читаются по шагам через session-oriented tool family. Обязательные операции: - `start` - `poll` - `stop` Подходит для: - follow logs; - telemetry follow; - alert/event feed; - контроль длительных stream-подписок; - WebSocket subscriptions. ### 4.4. `async_job` Runtime запускает long-running upstream operation и возвращает `job_id`. Обязательные операции: - `start` - `status` - `result` - `cancel` Подходит для: - import/export; - deploy/reindex; - batch processing; - инфраструктурные control-plane действия; - SOAP workflows with deferred status polling. ## 5. Бизнес-кейсы ### 5.1. Logs Window LLM запрашивает: - ошибки сервиса за последние `30s`; - top errors за последние `100` записей; - логи по конкретному `correlation_id`. Runtime: - собирает bounded окно; - агрегирует counts, уровни, sample lines; - возвращает summary плюс ограниченный список items. ### 5.2. Metrics Window LLM запрашивает: - latency/error summary по сервису; - CPU/memory snapshot; - anomaly summary за окно. Runtime: - читает поток метрик; - агрегирует min/max/avg/p95 или anomaly set; - возвращает компактный JSON. ### 5.3. Event Feed LLM запрашивает: - audit events за период; - queue events по фильтру; - security alerts за окно. Runtime: - фильтрует events; - ограничивает количество; - возвращает items plus cursor. ### 5.4. Long-running Operation Status LLM запускает: - reindex; - import job; - rollout; - repair task. Runtime: - создает job handle; - возвращает `job_id`; - позволяет дальше получать status/result/cancel. ### 5.5. Control Plane Follow LLM инициирует: - reboot; - rollout; - node drain; - workflow transition. Runtime: - стартует действие; - пишет progress в session/job state; - возвращает snapshots по `poll`. ### 5.6. WebSocket Realtime Feeds LLM запрашивает: - realtime alert snapshot; - device telemetry window; - market data slice; - status feed по subscription channel. Runtime: - открывает upstream WebSocket; - подписывается на channel; - собирает bounded окно или session step; - возвращает summary и limited items. ### 5.7. SOAP Enterprise Operations LLM запрашивает: - создание/поиск сущности в ERP; - запуск enterprise workflow; - получение статуса batch operation; - B2B request через SOAP gateway. Runtime: - строит SOAP envelope из MCP input; - вызывает enterprise endpoint; - нормализует XML response или SOAP Fault; - возвращает JSON-oriented output. ## 6. Функциональные требования ### 6.1. Общие - operation должна явно указывать `execution_mode`; - streaming operation обязана быть bounded; - runtime обязан поддерживать timeout, max items и max bytes; - tool output обязан иметь предсказуемую схему; - stream/session/job state должен быть наблюдаемым и логируемым; - cancel/stop должен быть явной операцией, а не побочным эффектом disconnect. ### 6.2. Для `window` - задать размер окна; - задать лимит items; - задать лимит bytes; - выбрать aggregation mode; - вернуть `truncated` и `window_complete` flags; - поддерживать optional cursor для следующего окна. ### 6.3. Для `session` - создать `session_id`; - поддерживать `poll`; - поддерживать `stop`; - хранить курсор и session status; - иметь `idle_timeout`; - иметь `max_session_lifetime`; - удалять expired sessions. ### 6.4. Для `async_job` - создать `job_id`; - хранить progress, status и final result metadata; - поддерживать `cancel`; - поддерживать retrieval последнего готового результата. ## 7. Нефункциональные требования ### 7.1. Безопасность - никакого неограниченного passthrough потока; - обязательные лимиты по времени, items и bytes; - обязательный redact layer для secret-bearing полей; - audit trail на `start`, `poll`, `stop`, `cancel`; - session и job identifiers должны быть криптографически стойкими. ### 7.2. Производительность - bounded memory per session; - bounded upstream read buffer; - ограничение числа параллельных sessions и jobs на workspace и agent; - backpressure при медленных клиентах; - возможность early cut-off после достижения лимита. ### 7.3. Надежность - TTL и cleanup для sessions/jobs; - устойчивость к disconnect downstream client; - poll должен быть идемпотентным; - long-running upstream action не должен считаться отмененным из-за SSE disconnect; - resumability для downstream SSE допускается как следующая волна, но не является обязательной в первой реализации. ### 7.4. UX - оператор должен видеть, что operation является `unary`, `window`, `session` или `async_job`; - UI должен явно показывать все лимиты и режим агрегации; - результат тестового вызова должен показывать `truncated`, `window_complete`, `has_more`, `status`. ## 8. UI contract ### 8.1. Новый блок `Execution mode` Поля: - `mode`: `unary | window | session | async_job` - `transport_behavior`: `request_response | server_stream` ### 8.2. Блок `Collection limits` Поля: - `window_duration_ms` - `poll_interval_ms` - `upstream_timeout_ms` - `idle_timeout_ms` - `max_session_lifetime_ms` - `max_items` - `max_bytes` ### 8.3. Блок `Aggregation` Поля: - `aggregation_mode`: `raw_items | summary_only | summary_plus_samples | stats | latest_state` - `summary_path` - `items_path` - `cursor_path` - `status_path` - `done_path` ### 8.4. Блок `Safety` Поля: - `truncate_item_fields` - `max_field_length` - `redacted_paths` - `drop_duplicates` - `sampling_rate` ### 8.5. Block `Tool family` Для `session`: - `start_tool_name` - `poll_tool_name` - `stop_tool_name` Для `async_job`: - `start_tool_name` - `status_tool_name` - `result_tool_name` - `cancel_tool_name` ## 9. Domain model changes ### 9.1. `ExecutionMode` Новый enum: - `Unary` - `Window` - `Session` - `AsyncJob` ### 9.2. `StreamingConfig` Новая часть `execution_config`: - `mode` - `window_duration_ms` - `poll_interval_ms` - `upstream_timeout_ms` - `idle_timeout_ms` - `max_session_lifetime_ms` - `max_items` - `max_bytes` - `aggregation_mode` - `items_path` - `summary_path` - `cursor_path` - `status_path` - `done_path` - `redacted_paths` ### 9.3. `StreamSession` Новая runtime/store сущность: - `id` - `workspace_id` - `agent_id` - `operation_id` - `protocol` - `mode` - `status` - `cursor` - `state_json` - `expires_at` - `last_poll_at` - `created_at` - `closed_at` ### 9.4. `AsyncJobHandle` Новая runtime/store сущность: - `id` - `workspace_id` - `agent_id` - `operation_id` - `status` - `progress_json` - `result_json` - `error_json` - `expires_at` - `created_at` - `updated_at` - `finished_at` ## 10. MCP publishing model for streaming tools ### 10.1. Unary and Window `unary` и `window` публикуются как один MCP tool: - один input contract; - один bounded result; - transport может использовать `application/json` или SSE response stream до финального JSON-RPC response. ### 10.2. Session `session` публикуется как tool family: - `{tool}_start` - `{tool}_poll` - `{tool}_stop` Причина: - lifecycle становится явным; - LLM получает контролируемую state machine; - runtime не скрывает долговременное состояние за одним "магическим" вызовом. ### 10.3. Async Job `async_job` публикуется как tool family: - `{tool}_start` - `{tool}_status` - `{tool}_result` - `{tool}_cancel` ## 11. Module decomposition and responsibilities ### 11.1. `crank-core` Новые модули: - `streaming` - `stream_session` Новые типы: - `ExecutionMode` - `StreamingConfig` - `AggregationMode` - `StreamSession` - `AsyncJobHandle` - `StreamStatus` - `JobStatus` ### 11.2. `crank-registry` Новые обязанности: - хранение `stream_sessions`; - хранение `async_jobs`; - cleanup expired rows; - optimistic updates on poll/stop/cancel. Ожидаемые функции: - `create_stream_session` - `get_stream_session` - `advance_stream_session` - `close_stream_session` - `create_async_job` - `get_async_job` - `update_async_job_status` - `cancel_async_job` - `delete_expired_stream_sessions` ### 11.3. `crank-runtime` Новые orchestration функции: - `execute_unary_operation` - `execute_window_operation` - `start_stream_session` - `poll_stream_session` - `stop_stream_session` - `start_async_job` - `get_async_job_status` - `get_async_job_result` - `cancel_async_job` ### 11.4. Protocol adapters REST: - unary HTTP; - bounded SSE collection; - bounded long-poll collection. GraphQL: - `query` и `mutation`; - `subscription` отложен на отдельную protocol wave. gRPC: - unary; - bounded server-stream collection; - client/bidi отложены на отдельную protocol wave. WebSocket: - bounded event collection; - subscribe/poll/stop orchestration; - heartbeat and reconnect policy. SOAP: - WSDL-driven request/response adapter; - XML normalization; - SOAP Fault normalization; - future WS-Security expansion. ### 11.5. `apps/mcp-server` Новые обязанности: - корректно вести `Streamable HTTP` lifecycle; - принимать `POST` с `Accept: application/json, text/event-stream`; - отдавать `application/json` или `text/event-stream`; - поддерживать `GET` SSE stream для server-to-client messages и notifications; - вести `Mcp-Session-Id`; - публиковать tool families для `session` и `async_job`. ### 11.6. `apps/admin-api` Новые обязанности: - CRUD и versioning для streaming config; - тестовые window/session/job runs; - UI-oriented validation ошибок для streaming fields. ### 11.7. `apps/ui` - конфиг execution mode; - конфиг limits/aggregation/safety; - test-run screen для bounded window/session/job behavior; - отдельные предупреждения про truncation и timeouts. ## 12. Границы текущей продуктовой волны В первой продуктовой волне входит: - `Streamable HTTP` и SSE на MCP transport; - `window` mode; - `async_job` mode; - REST SSE; - gRPC server streaming; - tool family generation; - bounded session/job state. Во второй продуктовой волне: - WebSocket upstream adapter; - SOAP adapter foundation; - richer session tooling; - expanded protocol smoke suite. Отложено: - GraphQL subscriptions; - gRPC client streaming; - gRPC bidirectional streaming; - raw infinite stream passthrough; - guaranteed resumability across all stream types; - generic websocket proxy mode. ## 13. Порядок реализации ### 13.1. `feat/streaming-mcp-architecture` - зафиксировать docs; - обновить protocol support matrix; - синхронизировать `TASKS.md`. ### 13.2. `feat/mcp-streamable-http-alignment` - довести `mcp-server` до полного соответствия `Streamable HTTP`; - session headers; - GET SSE; - protocol version header validation; - explicit cancel behavior. ### 13.3. `feat/streaming-core-model` - ввести `ExecutionMode`, `StreamingConfig`, `StreamSession`, `AsyncJobHandle`. ### 13.4. `feat/stream-session-store` - таблицы `stream_sessions` и `async_jobs`; - cleanup; - optimistic state transitions. ### 13.5. `feat/runtime-window-mode` - bounded collection для `window`; - `truncated`, `window_complete`, `has_more`. ### 13.6. `feat/rest-sse-adapter` - поддержка REST SSE upstream. ### 13.7. `feat/grpc-server-streaming-adapter` - поддержка bounded gRPC server-streaming. ### 13.8. `feat/session-and-job-tools` - генерация tool families; - `start/poll/stop`; - `start/status/result/cancel`. ### 13.9. `feat/websocket-upstream-adapter` - bounded WebSocket collection; - subscribe/unsubscribe templates; - heartbeat/reconnect policy; - session integration. ### 13.10. `feat/soap-architecture-and-core-model` - WSDL/XSD-driven domain model; - SOAP execution config; - XML normalization strategy. ### 13.11. `feat/soap-adapter-foundation` - runtime SOAP adapter; - envelope builder; - fault normalization; - test-run support. ### 13.12. `feat/streaming-ui-config` - новый execution mode selector; - limits/aggregation/safety blocks; - test-run UX. ### 13.13. `feat/streaming-e2e` - публичные smoke targets; - e2e сценарии; - manual regression plan. ## 14. Практический итог Crank должен поддерживать streaming как полнофункциональный MCP proxy, но в управляемой форме: - transport-level SSE и `Streamable HTTP` поддерживаются; - upstream streaming поддерживается там, где его можно bounded-ить; - tool contract остается контролируемым; - UI настраивает лимиты, aggregation и lifecycle; - платформа не превращается в бесконечную data pipe.