16 KiB
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-Versionheaders;- resumable SSE streams как опциональная возможность.
Не является текущим приоритетом:
stdioкак обязательная часть продукта;- собственные нестандартные transport-режимы поверх MCP.
3.2. Upstream protocols
Поддерживаются:
- REST unary;
- REST SSE в bounded режимах;
- GraphQL
queryиmutation; - gRPC unary;
- gRPC server-streaming в bounded режимах.
Отложено:
- GraphQL
subscription; - gRPC client-streaming;
- gRPC bidirectional streaming;
- arbitrary websocket passthrough;
- raw infinite stream forwarding в MCP client.
4. Поддерживаемые streaming modes
Crank поддерживает четыре режима выполнения operation.
4.1. unary
Обычный request-response вызов.
Подходит для:
- REST;
- GraphQL
queryиmutation; - gRPC unary.
4.2. window
Runtime открывает upstream stream или repeatedly polls upstream source, собирает данные в пределах окна и возвращает один bounded ответ.
Параметры:
window_duration_msmax_itemsmax_bytesupstream_timeout_msaggregation_mode
Подходит для:
- логи за период;
- метрики за период;
- event window;
- SSE stream snapshot;
- gRPC server-stream window.
4.3. session
Runtime создает stream session, после чего данные читаются по шагам через session-oriented tool family.
Обязательные операции:
startpollstop
Подходит для:
- follow logs;
- telemetry follow;
- alert/event feed;
- контроль длительных stream-подписок.
4.4. async_job
Runtime запускает long-running upstream operation и возвращает job_id.
Обязательные операции:
startstatusresultcancel
Подходит для:
- import/export;
- deploy/reindex;
- batch processing;
- инфраструктурные control-plane действия.
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.
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_completeflags; - поддерживать 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 допускается, но не является обязательной в MVP.
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_jobtransport_behavior:request_response | server_stream
8.2. Блок Collection limits
Поля:
window_duration_mspoll_interval_msupstream_timeout_msidle_timeout_msmax_session_lifetime_msmax_itemsmax_bytes
8.3. Блок Aggregation
Поля:
aggregation_mode:raw_items | summary_only | summary_plus_samples | stats | latest_statesummary_pathitems_pathcursor_pathstatus_pathdone_path
8.4. Блок Safety
Поля:
truncate_item_fieldsmax_field_lengthredacted_pathsdrop_duplicatessampling_rate
8.5. Block Tool family
Для session:
start_tool_namepoll_tool_namestop_tool_name
Для async_job:
start_tool_namestatus_tool_nameresult_tool_namecancel_tool_name
9. Domain model changes
9.1. ExecutionMode
Новый enum:
UnaryWindowSessionAsyncJob
9.2. StreamingConfig
Новая часть execution_config:
modewindow_duration_mspoll_interval_msupstream_timeout_msidle_timeout_msmax_session_lifetime_msmax_itemsmax_bytesaggregation_modeitems_pathsummary_pathcursor_pathstatus_pathdone_pathredacted_paths
9.3. StreamSession
Новая runtime/store сущность:
idworkspace_idagent_idoperation_idprotocolmodestatuscursorstate_jsonexpires_atlast_poll_atcreated_atclosed_at
9.4. AsyncJobHandle
Новая runtime/store сущность:
idworkspace_idagent_idoperation_idstatusprogress_jsonresult_jsonerror_jsonexpires_atcreated_atupdated_atfinished_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
Новые модули:
streamingstream_session
Новые типы:
ExecutionModeStreamingConfigAggregationModeStreamSessionAsyncJobHandleStreamStatusJobStatus
11.2. crank-registry
Новые обязанности:
- хранение
stream_sessions; - хранение
async_jobs; - cleanup expired rows;
- optimistic updates on poll/stop/cancel.
Ожидаемые функции:
create_stream_sessionget_stream_sessionadvance_stream_sessionclose_stream_sessioncreate_async_jobget_async_jobupdate_async_job_statuscancel_async_jobdelete_expired_stream_sessions
11.3. crank-runtime
Новые orchestration функции:
execute_unary_operationexecute_window_operationstart_stream_sessionpoll_stream_sessionstop_stream_sessionstart_async_jobget_async_job_statusget_async_job_resultcancel_async_job
11.4. Protocol adapters
REST:
- unary HTTP;
- bounded SSE collection;
- bounded long-poll collection.
GraphQL:
queryиmutation;subscriptionне входит в MVP.
gRPC:
- unary;
- bounded server-stream collection;
- client/bidi не входят в MVP.
11.5. apps/mcp-server
Новые обязанности:
- корректно вести
Streamable HTTPlifecycle; - принимать
POSTсAccept: application/json, text/event-stream; - отдавать
application/jsonилиtext/event-stream; - поддерживать
GETSSE 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. Ограничения MVP
В MVP входит:
Streamable HTTPи SSE на MCP transport;windowmode;async_jobmode;- REST SSE;
- gRPC server streaming;
- tool family generation;
- bounded session/job state.
В MVP не входит:
- 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/streaming-ui-config
- новый execution mode selector;
- limits/aggregation/safety blocks;
- test-run UX.
13.10. 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.