diff --git a/README.md b/README.md index dc1132d..fadbc25 100644 --- a/README.md +++ b/README.md @@ -20,7 +20,8 @@ Crank - платформа для публикации внешних API в в - `Agent` как curated MCP surface для LLM. - Поддержка REST для `GET`, `POST`, `PUT`, `PATCH` и `DELETE`. - Поддержка GraphQL для `query` и `mutation`. -- Поддержка только unary-методов gRPC. +- Поддержка unary и bounded server-streaming для gRPC. +- Поддержка controlled streaming modes поверх MCP `Streamable HTTP`. - Platform API keys и membership layer. - Observability: invocation logs, usage aggregates, latency/error metrics. - Импорт и экспорт operation-конфигураций в `YAML`. @@ -45,6 +46,7 @@ Crank - платформа для публикации внешних API в в - `docs/demo-runbook.md` - демонстрационный сценарий. - `docs/public-smoke-targets.md` - готовые публичные upstream-сервисы и payload-ы для smoke-проверки MCP. - `docs/secrets-auth-plan.md` - целевая модель upstream secrets, auth profiles и пошаговый план реализации. +- `docs/streaming-mcp-plan.md` - целевая модель MCP transport streaming, upstream streaming и поэтапный план реализации. - `docs/rust-design.md` - правила распределения поведения в Rust. - `docs/development-rules.md` - правила разработки и workflow. - `docs/rust-code-rules.md` - Rust-specific coding rules. diff --git a/TASKS.md b/TASKS.md index f3c76f0..ad6532e 100644 --- a/TASKS.md +++ b/TASKS.md @@ -2,23 +2,31 @@ ## Current -### `feat/secret-store-foundation` +### `feat/streaming-mcp-architecture` Status: completed DoD: -- Workspace-scoped secret store is persisted in PostgreSQL -- Secret values are encrypted with `CRANK_MASTER_KEY` -- Admin API exposes create/list/get/rotate/delete secret endpoints -- Integration tests cover secret CRUD and rotation +- Official MCP transport semantics are reflected in docs +- Streaming modes and protocol support matrix are documented +- Core docs and `TASKS.md` are synchronized around controlled streaming model ## Next -- `feat/auth-profile-secret-resolution` +- `feat/mcp-streamable-http-alignment` ## Backlog -- `feat/secret-store-foundation` +- `feat/streaming-mcp-architecture` +- `feat/mcp-streamable-http-alignment` +- `feat/streaming-core-model` +- `feat/stream-session-store` +- `feat/runtime-window-mode` +- `feat/rest-sse-adapter` +- `feat/grpc-server-streaming-adapter` +- `feat/session-and-job-tools` +- `feat/streaming-ui-config` +- `feat/streaming-e2e` - `feat/auth-profile-secret-resolution` - `feat/runtime-upstream-auth` - `feat/secrets-ui` diff --git a/docs/architecture.md b/docs/architecture.md index df215ac..d43a836 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -130,6 +130,7 @@ Crank - платформа для публикации внешних API в в - `Workspace` как tenant boundary. - Операции `REST`, `GraphQL`, `unary gRPC`. +- Controlled streaming operations поверх `Streamable HTTP`, REST SSE и gRPC server-streaming. - `Agent` и привязка операций к агенту. - Agent-scoped MCP endpoints. - Platform API keys. @@ -141,7 +142,8 @@ Crank - платформа для публикации внешних API в в ### Не входит -- gRPC streaming. +- GraphQL `subscription`. +- gRPC client-streaming и bidirectional streaming. - SOAP. - Оркестрация workflow. - Биллинг. @@ -201,10 +203,20 @@ GraphQL в MCP публикуется как фиксированная опер ### gRPC -- только unary RPC; +- unary RPC; +- bounded server-streaming через `window`, `session` и `async_job` execution modes; - `.proto` и `descriptor set`; - JSON-oriented schema model поверх protobuf; -- без streaming. +- без client-streaming и bidi. + +### Streaming + +Платформа поддерживает controlled streaming model: + +- downstream transport: `Streamable HTTP` с optional SSE; +- upstream streaming: REST SSE и gRPC server-streaming; +- execution modes: `unary`, `window`, `session`, `async_job`; +- никакого raw infinite stream passthrough в MCP client. ## 8. Работа с файлами и автогенерация черновика @@ -234,6 +246,8 @@ GraphQL в MCP публикуется как фиксированная опер - `Agent` - `AgentVersion` - `AgentOperationBinding` +- `StreamSession` +- `AsyncJobHandle` - `AuthProfile` - `PlatformApiKey` - `UserSession` diff --git a/docs/data-model.md b/docs/data-model.md index 670bdbf..d5dea27 100644 --- a/docs/data-model.md +++ b/docs/data-model.md @@ -17,7 +17,7 @@ Каждая `Operation` соответствует одному интеграционному контракту: - GraphQL -> один конкретный `query` или `mutation`; -- gRPC -> один unary method; +- gRPC -> один unary method или один server-streaming method в bounded execution mode; - REST -> один endpoint-сценарий. Однако MCP tool публикуется не напрямую из operation, а через `AgentOperationBinding` внутри конкретного `Agent`. @@ -54,6 +54,16 @@ Помимо канонической JSON-модели система поддерживает импорт и экспорт конфигураций в `YAML`. +### 2.7. Streaming operation обязана быть bounded + +Если операция использует upstream streaming, она должна работать в одном из execution modes: + +- `window` +- `session` +- `async_job` + +Бесконечный passthrough stream не является допустимой моделью `Operation`. + ## 3. Корневые сущности ### 3.1. `Workspace` @@ -96,6 +106,7 @@ - `samples` - `generated_draft` - `config_export` +- `streaming_config` - `created_at` - `updated_at` - `published_at` @@ -314,6 +325,50 @@ - `p95_ms` - `p99_ms` +### 3.16. `StreamSession` + +Поля: + +- `id` +- `workspace_id` +- `agent_id` +- `operation_id` +- `mode` +- `status` +- `cursor` +- `state` +- `expires_at` +- `last_poll_at` +- `created_at` +- `closed_at` + +Назначение: + +- хранение bounded session state для streaming tools; +- поддержка `start/poll/stop`; +- cleanup orphaned и expired sessions. + +### 3.17. `AsyncJobHandle` + +Поля: + +- `id` +- `workspace_id` +- `agent_id` +- `operation_id` +- `status` +- `progress` +- `result` +- `error` +- `expires_at` +- `created_at` +- `updated_at` +- `finished_at` + +Назначение: + +- поддержка `start/status/result/cancel` для long-running upstream actions. + ## 4. `Target` `Target` описывает конкретный внешний вызов. Это discriminated union по протоколу. diff --git a/docs/database-schema.md b/docs/database-schema.md index 9497258..88eccf3 100644 --- a/docs/database-schema.md +++ b/docs/database-schema.md @@ -49,6 +49,8 @@ - `agent_operation_bindings` - `published_agents` - `platform_api_keys` +- `stream_sessions` +- `async_jobs` - `invocation_logs` - `usage_rollups` - `yaml_import_jobs` @@ -180,6 +182,46 @@ - `unique (workspace_id, name)` +### `stream_sessions` + +- `id` +- `workspace_id` +- `agent_id` +- `operation_id` +- `mode` +- `status` +- `cursor_json` +- `state_json` +- `expires_at` +- `last_poll_at` +- `created_at` +- `closed_at` + +Индексы: + +- `(workspace_id, status, expires_at)` +- `(operation_id, status)` + +### `async_jobs` + +- `id` +- `workspace_id` +- `agent_id` +- `operation_id` +- `status` +- `progress_json` +- `result_json` +- `error_json` +- `expires_at` +- `created_at` +- `updated_at` +- `finished_at` + +Индексы: + +- `(workspace_id, status, expires_at)` +- `(operation_id, status)` + ## 7. Workspaces and access layer ### `workspaces` diff --git a/docs/implementation-plan.md b/docs/implementation-plan.md index 4aef4ce..84c5400 100644 --- a/docs/implementation-plan.md +++ b/docs/implementation-plan.md @@ -134,3 +134,17 @@ DoD: - end-to-end demo воспроизводим; - deployment и healthchecks стабильно зелёные; - документация и продуктовый сценарий совпадают. + +## 12. Этап 11. MCP streaming proxy support + +Цель: + +- довести Crank до controlled streaming model поверх MCP `Streamable HTTP`. + +DoD: + +- `mcp-server` соответствует transport semantics `2025-06-18`; +- execution modes `window`, `session`, `async_job` формально описаны и реализованы; +- REST SSE и gRPC server-streaming поддерживаются в bounded форме; +- UI умеет конфигурировать streaming limits, aggregation и lifecycle; +- e2e сценарии покрывают window/session/job calls. diff --git a/docs/mcp-interface.md b/docs/mcp-interface.md index 0467e26..f83e4d6 100644 --- a/docs/mcp-interface.md +++ b/docs/mcp-interface.md @@ -11,6 +11,8 @@ Решение: - основной transport: `Streamable HTTP`; +- `POST` может завершаться `application/json` или `text/event-stream`; +- `GET` SSE stream поддерживается как server-to-client канал; - отдельный `mcp-server` как сервис; - `stdio` не является обязательной частью текущего scope. @@ -40,6 +42,7 @@ - загрузить published agents и их bindings из registry; - преобразовать их в MCP tool definitions; +- вести `Mcp-Session-Id` и `MCP-Protocol-Version`; - принимать вызовы tools от MCP clients; - валидировать вход; - делегировать исполнение в runtime; @@ -116,6 +119,12 @@ - `tools/list` - `tools/call` +Дополнительно на transport уровне: + +- `POST /mcp/v1/{workspace_slug}/{agent_slug}` как основной MCP endpoint; +- `GET /mcp/v1/{workspace_slug}/{agent_slug}` для SSE stream; +- `DELETE /mcp/v1/{workspace_slug}/{agent_slug}` для explicit session termination, если сервер разрешает client-side session close. + ### Tool listing 1. клиент вызывает `tools/list`; @@ -133,6 +142,34 @@ 5. делегирует вызов в `crank-runtime`; 6. возвращает результат. +### Tool call и streaming + +Поддерживаются четыре execution modes: + +- `unary` +- `window` +- `session` +- `async_job` + +Правила публикации: + +- `unary` и `window` публикуются как один tool; +- `session` публикуется как `start/poll/stop` family; +- `async_job` публикуется как `start/status/result/cancel` family. + +Transport-level SSE не отменяет bounded tool semantics. Даже если `POST` отвечает через `text/event-stream`, итогом вызова должен оставаться управляемый JSON-RPC lifecycle. + +## 8.1. Multiple SSE connections + +Crank должен корректно работать, если MCP client держит несколько SSE streams одновременно. + +Правила: + +- одно server message отправляется только в один stream; +- disconnect не считается cancel; +- cancel выражается отдельным MCP notification или session/job control tool; +- resumability допустима как будущая возможность, но не обязательна в MVP. + ## 9. Обновление tools После публикации новой operation version или agent version: @@ -171,7 +208,9 @@ - `mcp-server` - отдельный сервис; - transport - `Streamable HTTP`; +- `POST` и `GET` transport semantics соответствуют MCP spec `2025-06-18`; - endpoint определяется парой `workspace + agent`; - одна published operation = один MCP tool внутри agent; +- streaming operations публикуются как bounded tools или tool families; - reload published tools без пересборки сервиса; - никакой draft-логики или admin CRUD в MCP слое. diff --git a/docs/module-decomposition.md b/docs/module-decomposition.md index 24bd200..41d107a 100644 --- a/docs/module-decomposition.md +++ b/docs/module-decomposition.md @@ -42,6 +42,7 @@ crank/ - secret management domain; - agent publishing domain; - observability domain. +- streaming execution domain. ## 4. Детальная декомпозиция по crate @@ -64,6 +65,7 @@ crank/ - `auth` - `secret` - `observability` +- `streaming` - `errors` ### 4.2. `crank-schema` @@ -98,6 +100,7 @@ crank/ - хранение workspace-scoped operations и version snapshots; - хранение workspace-scoped secrets и secret versions; - хранение agents и agent versions; +- хранение stream sessions и async jobs; - auth profiles; - platform API keys; - logs и usage aggregates; @@ -109,6 +112,7 @@ crank/ - исполнение published operation; - резолв `auth_profile_ref -> secret -> request auth`; +- orchestration window/session/job execution; - запись invocation events; - возврат нормализованного результата. @@ -133,6 +137,7 @@ crank/ - `platform_api_keys` - `logs` - `usage` +- `streaming` ### 4.9. `apps/mcp-server` @@ -141,6 +146,8 @@ crank/ - публикация published agent bindings как MCP tools; - transport handling; - JSON-RPC lifecycle; +- SSE lifecycle и `Mcp-Session-Id`; +- tool-family generation для `session` и `async_job`; - вызов runtime. Антипаттерн: diff --git a/docs/protocols/graphql.md b/docs/protocols/graphql.md index a6b11cc..b797cb1 100644 --- a/docs/protocols/graphql.md +++ b/docs/protocols/graphql.md @@ -30,6 +30,8 @@ GraphQL поддерживается как отдельный тип интег - обязательная зависимость от introspection - автоматическое построение любого запроса по полной GraphQL schema +`subscription` допускается только как future scope после появления отдельного websocket/subscription adapter и controlled streaming lifecycle. + ## 4. Ключевое архитектурное ограничение Платформа не должна публиковать в MCP общий GraphQL tool, который умеет получать любые поля и принимать любые параметры в зависимости от намерения LLM. @@ -92,7 +94,7 @@ GraphQL operation должна включать: - структура ответа зависит от `selection set`, значит она должна быть фиксирована заранее; - GraphQL endpoint обычно один, поэтому операция определяется не URL, а телом запроса; - variables должны быть строго ограничены, иначе один tool станет слишком широким и плохо управляемым; -- `subscription` по смыслу не подходит модели MCP tool, потому что это потоковая, а не request-response интеграция. +- `subscription` не входит в текущий scope, потому что требует отдельной lifecycle-модели, близкой к `session` mode, и отдельного transport adapter. - `JSONPath` используется для точечного извлечения вложенных данных из `data` и для управления структурой итогового ответа. ## 9. Почему GraphQL не считается "почти REST" diff --git a/docs/protocols/grpc.md b/docs/protocols/grpc.md index 6a1e873..ae8d36f 100644 --- a/docs/protocols/grpc.md +++ b/docs/protocols/grpc.md @@ -2,11 +2,12 @@ ## 1. Роль протокола в проекте -gRPC поддерживается как третий основной протокол платформы, но в самой узкой и управляемой форме. Цель состоит не в том, чтобы покрыть все возможности gRPC, а в том, чтобы представить unary RPC-методы как обычные MCP tools с формой входа и формой выхода. +gRPC поддерживается как третий основной протокол платформы в управляемой форме. Цель состоит не в том, чтобы покрыть все возможности gRPC, а в том, чтобы представить unary и bounded server-streaming методы как MCP tools с предсказуемым жизненным циклом. ## 2. Что поддерживается в MVP -- только unary RPC +- unary RPC +- bounded server-streaming через execution modes `window`, `session`, `async_job` - загрузка `.proto` - загрузка descriptor set - загрузка примеров JSON для MCP input/output при необходимости @@ -22,7 +23,6 @@ gRPC поддерживается как третий основной прот ## 3. Что не входит в MVP -- `server streaming` - `client streaming` - `bidirectional streaming` - обязательная поддержка server reflection @@ -31,17 +31,15 @@ gRPC поддерживается как третий основной прот ## 4. Ключевое архитектурное ограничение -В проекте поддерживаются только unary-методы, потому что MCP tool в этой архитектуре соответствует модели `один запрос -> один ответ`. +В проекте поддерживаются unary-методы и bounded server-streaming, потому что MCP tool в этой архитектуре должен оставаться управляемым. Это означает: - один request message; -- один response message; -- один завершенный вызов; -- отсутствие потоковых сообщений; -- отсутствие отдельного жизненного цикла stream-сессии. +- один bounded response или управляемая session/job-семантика; +- явно ограниченный lifecycle stream-сессии. -Streaming gRPC не нужен для выбранной модели взаимодействия с LLM и только усложнит runtime, UI и хранение состояния. +Streaming gRPC не публикуется как бесконечный raw stream. Он допускается только там, где runtime умеет bounded-ить, агрегировать и завершать результат. ## 5. Внутренняя модель gRPC operation @@ -64,7 +62,7 @@ gRPC operation должна включать: 1. Загружает `.proto` или descriptor set. 2. Система извлекает список services и methods. -3. Оператор выбирает конкретный unary-метод. +3. Оператор выбирает конкретный unary- или server-streaming метод. 4. UI показывает структуру request message и response message. 5. При необходимости загружает примеры JSON для MCP input/output. 6. Система строит черновую схему, стартовый mapping и runtime-ready snapshot descriptor set для выбранного метода. @@ -82,10 +80,11 @@ gRPC operation должна включать: 1. Валидировать MCP input по нормализованной схеме. 2. Применить input mapping. 3. Построить protobuf request message из JSON. -4. Выполнить unary RPC вызов. -5. Преобразовать protobuf response в нормализованный JSON. -6. Применить output mapping. -7. Вернуть итоговый результат. +4. Для unary выполнить unary RPC вызов. +5. Для server-streaming собрать bounded окно или session step. +6. Преобразовать protobuf response или stream items в нормализованный JSON. +7. Применить output mapping. +8. Вернуть итоговый результат. ## 8. Критические нюансы @@ -99,7 +98,7 @@ gRPC operation должна включать: - пользователь не должен видеть внутреннюю сложность protobuf-контракта больше, чем это нужно для настройки operation. - `JSONPath` используется как единый способ точечной адресации вложенных полей при настройке mapping поверх нормализованной JSON-модели. -## 9. Почему gRPC ограничивается unary +## 9. Почему gRPC ограничивается controlled streaming Причина не только в сложности реализации. Главное ограничение архитектурное: @@ -108,4 +107,4 @@ gRPC operation должна включать: - UI платформы построен вокруг формы входа и формы выхода; - streaming требует отдельной session-модели, buffering, cancellation и состояния. -Поэтому unary gRPC - это не "обрезанная" поддержка, а осознанно выбранная форма, которая действительно совместима с MCP-платформой. +Поэтому поддержка gRPC в Crank ограничивается unary и bounded server-streaming. Client-streaming и bidi остаются вне scope до появления полноценной interactive session model. diff --git a/docs/protocols/rest.md b/docs/protocols/rest.md index d42b01b..1fd05ae 100644 --- a/docs/protocols/rest.md +++ b/docs/protocols/rest.md @@ -21,6 +21,7 @@ REST - базовый и первый по очередности реализа - автогенерация чернового mapping - ручная донастройка через `JSONPath` - тестовый вызов перед публикацией +- optional REST SSE upstream в bounded `window` и `session` режимах ## 3. Что не входит в MVP @@ -31,6 +32,7 @@ REST - базовый и первый по очередности реализа - webhooks - long polling как специальный режим - `HEAD` и `OPTIONS` как отдельные пользовательские сценарии +- raw infinite SSE passthrough ## 4. Внутренняя модель REST operation @@ -49,6 +51,12 @@ REST operation в системе описывается следующими о На слое MCP REST operation всегда выглядит как вызов `запрос -> ответ` с фиксированной схемой входа и выхода. +Если upstream использует SSE, на слое MCP это все равно должно быть выражено как: + +- bounded window result; +- session-oriented `start/poll/stop`; +- async job semantics для длительных действий. + ## 5. Как оператор настраивает REST operation 1. Указывает `base_url`. diff --git a/docs/streaming-mcp-plan.md b/docs/streaming-mcp-plan.md new file mode 100644 index 0000000..e1fa4d4 --- /dev/null +++ b/docs/streaming-mcp-plan.md @@ -0,0 +1,615 @@ +# 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 режимах. + +Отложено: + +- 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_ms` +- `max_items` +- `max_bytes` +- `upstream_timeout_ms` +- `aggregation_mode` + +Подходит для: + +- логи за период; +- метрики за период; +- event window; +- SSE stream snapshot; +- gRPC server-stream window. + +### 4.3. `session` + +Runtime создает stream session, после чего данные читаются по шагам через session-oriented tool family. + +Обязательные операции: + +- `start` +- `poll` +- `stop` + +Подходит для: + +- follow logs; +- telemetry follow; +- alert/event feed; +- контроль длительных stream-подписок. + +### 4.4. `async_job` + +Runtime запускает long-running upstream operation и возвращает `job_id`. + +Обязательные операции: + +- `start` +- `status` +- `result` +- `cancel` + +Подходит для: + +- 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_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 допускается, но не является обязательной в 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_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` не входит в MVP. + +gRPC: + +- unary; +- bounded server-stream collection; +- client/bidi не входят в MVP. + +### 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. Ограничения MVP + +В MVP входит: + +- `Streamable HTTP` и SSE на MCP transport; +- `window` mode; +- `async_job` mode; +- 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.