# Streaming Implementation Specification ## 1. Назначение документа Этот документ переводит streaming architecture в execution-oriented plan, пригодный для агентной разработки. Он отвечает на вопросы: - в каком порядке реализовывать фичи; - какие файлы менять в каждом срезе; - какие новые модули создавать; - какие тесты обязательны; - какие риски и зависимости есть между срезами; - какой Definition of Done нужен для каждого шага. Документ дополняет: - [streaming-mcp-plan.md](/home/a.tolmachev/code/rust/mcpaas/docs/streaming-mcp-plan.md) - [streaming-admin-api.md](/home/a.tolmachev/code/rust/mcpaas/docs/streaming-admin-api.md) - [streaming-runtime-design.md](/home/a.tolmachev/code/rust/mcpaas/docs/streaming-runtime-design.md) - [streaming-ui-contract.md](/home/a.tolmachev/code/rust/mcpaas/docs/streaming-ui-contract.md) - [protocol-capability-matrix.md](/home/a.tolmachev/code/rust/mcpaas/docs/protocol-capability-matrix.md) ## 2. Принципы агентной реализации - один vertical slice на ветку; - docs first, code second; - не смешивать transport alignment, core model, registry, runtime и UI в одном срезе; - каждый срез должен быть протестирован отдельно; - stateful streaming нельзя вводить через скрытые side effects; - session/job lifecycle должен появляться в коде и API явно. ## 3. Dependency graph Порядок реализации: 1. `feat/mcp-streamable-http-alignment` 2. `feat/streaming-core-model` 3. `feat/stream-session-store` 4. `feat/runtime-window-mode` 5. `feat/rest-sse-adapter` 6. `feat/grpc-server-streaming-adapter` 7. `feat/session-and-job-tools` 8. `feat/streaming-ui-config` 9. `feat/streaming-e2e` 10. `feat/websocket-upstream-adapter` 11. `feat/soap-architecture-and-core-model` 12. `feat/soap-adapter-foundation` Почему именно так: - сначала transport correctness; - потом core types; - потом persistence; - потом bounded runtime; - потом protocol adapters; - потом MCP tool publishing; - потом UI; - потом e2e; - потом расширение protocol platform. ## 4. Slice 1: `feat/mcp-streamable-http-alignment` ### 4.1. Цель Довести `apps/mcp-server` до полного и явного соответствия MCP `Streamable HTTP`. ### 4.2. Файлы Изменить: - [apps/mcp-server/src/app.rs](/home/a.tolmachev/code/rust/mcpaas/apps/mcp-server/src/app.rs) - [apps/mcp-server/src/jsonrpc.rs](/home/a.tolmachev/code/rust/mcpaas/apps/mcp-server/src/jsonrpc.rs) - [apps/mcp-server/src/session.rs](/home/a.tolmachev/code/rust/mcpaas/apps/mcp-server/src/session.rs) - [apps/mcp-server/src/main.rs](/home/a.tolmachev/code/rust/mcpaas/apps/mcp-server/src/main.rs) Добавить при необходимости: - `apps/mcp-server/src/transport.rs` - `apps/mcp-server/src/sse.rs` ### 4.3. Что должно быть сделано - parse `Accept` корректно для `application/json` и `text/event-stream`; - `GET` endpoint перестает быть `405` и становится SSE-capable transport stream; - `POST` может отвечать JSON или SSE в зависимости от negotiated transport; - `Mcp-Session-Id` создается, читается и валидируется явно; - `MCP-Protocol-Version` валидируется и возвращается явно; - `DELETE` корректно завершает transport session; - disconnect не трактуется как cancel operation. ### 4.4. Тесты Обязательно добавить: - initialize через JSON response; - initialize через SSE response; - `GET` SSE handshake; - повторное использование `Mcp-Session-Id`; - invalid protocol version; - invalid accept header; - session close through `DELETE`. ### 4.5. DoD - `mcp-server` ведет себя согласно spec; - transport tests зелёные; - docs и code names совпадают. ## 5. Slice 2: `feat/streaming-core-model` ### 5.1. Цель Ввести доменные типы streaming execution. ### 5.2. Файлы Изменить: - [crates/crank-core/src/lib.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-core/src/lib.rs) - [crates/crank-core/src/operation.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-core/src/operation.rs) - [crates/crank-core/src/ids.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-core/src/ids.rs) - [crates/crank-core/src/protocol.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-core/src/protocol.rs) Добавить: - `crates/crank-core/src/streaming.rs` - `crates/crank-core/src/stream_session.rs` ### 5.3. Что должно быть сделано - добавить `ExecutionMode`; - добавить `TransportBehavior`; - добавить `AggregationMode`; - добавить `StreamingConfig`; - добавить `StreamSessionId`, `AsyncJobId`; - добавить `StreamSession`, `AsyncJobHandle`, `StreamStatus`, `JobStatus`; - вшить `streaming: Option` в `ExecutionConfig`. ### 5.4. Тесты - serde roundtrip для `StreamingConfig`; - validation rules для `ExecutionMode`; - protocol-aware validation helpers. ### 5.5. DoD - core types стабильны; - экспортированы через `lib.rs`; - naming не спорит с docs. ## 6. Slice 3: `feat/stream-session-store` ### 6.1. Цель Добавить persistent store для `stream_sessions` и `async_jobs`. ### 6.2. Файлы Изменить: - [crates/crank-registry/src/lib.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-registry/src/lib.rs) - [crates/crank-registry/src/model.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-registry/src/model.rs) - [crates/crank-registry/src/migrations.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-registry/src/migrations.rs) - [crates/crank-registry/src/postgres.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-registry/src/postgres.rs) ### 6.3. Что должно быть сделано - таблицы `stream_sessions`, `async_jobs`; - record types; - create/get/update/close/list/delete-expired methods; - filters and paging; - optimistic transitions для state updates. ### 6.4. Тесты - migration test; - create/get/update/close session; - create/get/update/cancel async job; - cleanup expired rows; - invalid transition rejection. ### 6.5. DoD - storage API соответствует runtime design; - transitions не допускают silent corruption. ## 7. Slice 4: `feat/runtime-window-mode` ### 7.1. Цель Ввести bounded `window` mode до session/job complexity. ### 7.2. Файлы Изменить: - [crates/crank-runtime/src/lib.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-runtime/src/lib.rs) - [crates/crank-runtime/src/model.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-runtime/src/model.rs) - [crates/crank-runtime/src/executor.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-runtime/src/executor.rs) - [crates/crank-runtime/src/error.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-runtime/src/error.rs) Добавить: - `crates/crank-runtime/src/streaming.rs` - `crates/crank-runtime/src/aggregation.rs` - `crates/crank-runtime/src/redaction.rs` ### 7.3. Что должно быть сделано - `execute_window_operation`; - bounded item/byte limiting; - `truncated`, `window_complete`, `has_more`; - summary building; - redaction and truncation; - observability write. ### 7.4. Тесты - raw items mode; - summary only mode; - summary plus samples; - max item truncation; - max byte truncation; - redaction; - timeout handling. ### 7.5. DoD - unary execution не ломается; - window mode работает для synthetic adapter-level fixtures. ## 8. Slice 5: `feat/rest-sse-adapter` ### 8.1. Цель Добавить bounded REST SSE support. ### 8.2. Файлы Изменить: - текущий crate [crates/crank-adapter-rest](/home/a.tolmachev/code/rust/mcpaas/crates/crank-adapter-rest) - [crates/crank-runtime/src/executor.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-runtime/src/executor.rs) Ожидаемые файлы: - `src/lib.rs` - `src/error.rs` - `src/client.rs` - `src/sse.rs` ### 8.3. Что должно быть сделано - SSE connect; - event parsing; - bounded collect window; - session start/poll/stop hooks, если adapter slice сразу включает stateful mode; - proper close. ### 8.4. Тесты - local SSE upstream fixture; - collect N events; - timeout without events; - malformed event handling; - reconnect not enabled by default. ## 9. Slice 6: `feat/grpc-server-streaming-adapter` ### 9.1. Цель Добавить bounded gRPC server-streaming. ### 9.2. Файлы Изменить: - [crates/crank-adapter-grpc](/home/a.tolmachev/code/rust/mcpaas/crates/crank-adapter-grpc) - [crates/crank-proto](/home/a.tolmachev/code/rust/mcpaas/crates/crank-proto) - [crates/crank-runtime/src/executor.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-runtime/src/executor.rs) ### 9.3. Что должно быть сделано - identify unary vs server-streaming method from descriptors; - invoke server stream; - collect bounded window; - decode stream items to normalized JSON; - integrate with window mode. ### 9.4. Тесты - local grpc fixture with server-streaming; - descriptor-backed invocation; - bounded collection; - timeout; - malformed item handling. ## 10. Slice 7: `feat/session-and-job-tools` ### 10.1. Цель Ввести `session` и `async_job` как first-class runtime and MCP constructs. ### 10.2. Файлы Изменить: - [crates/crank-runtime/src/executor.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-runtime/src/executor.rs) - [apps/mcp-server/src/app.rs](/home/a.tolmachev/code/rust/mcpaas/apps/mcp-server/src/app.rs) - [apps/mcp-server/src/catalog.rs](/home/a.tolmachev/code/rust/mcpaas/apps/mcp-server/src/catalog.rs) - [apps/mcp-server/src/session.rs](/home/a.tolmachev/code/rust/mcpaas/apps/mcp-server/src/session.rs) - [crates/crank-registry/src/postgres.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-registry/src/postgres.rs) ### 10.3. Что должно быть сделано - runtime methods for start/poll/stop and start/status/result/cancel; - MCP tool-family generation; - binding-level tool-name derivation; - JSON-RPC call routing for tool families; - session/job observability. ### 10.4. Тесты - tool listing shows generated family; - `session start -> poll -> stop`; - `async_job start -> status -> result -> cancel`; - invalid session id; - expired session; - result before ready. ## 11. Slice 8: `feat/streaming-ui-config` ### 11.1. Цель Добавить UI для streaming configuration. ### 11.2. Файлы Изменить: - `apps/ui/html/wizard/*.html` - `apps/ui/js/wizard.js` - `apps/ui/js/api.js` - `apps/ui/js/i18n.js` - CSS файлы wizard/settings/pages по необходимости Добавить: - `apps/ui/js/streaming-form.js` - `apps/ui/js/stream-test-run.js` - `apps/ui/html/stream-sessions.html` - `apps/ui/html/async-jobs.html` ### 11.3. Что должно быть сделано - protocol capabilities fetch; - execution mode selector; - protocol-specific field visibility; - validation wiring; - test-run UX for window/session/async job; - sessions/jobs views; - localization. ### 11.4. Тесты - Playwright: - configure REST window; - configure gRPC streaming; - validation error render; - session test flow; - async job test flow. ## 12. Slice 9: `feat/streaming-e2e` ### 12.1. Цель Зафиксировать end-to-end behavior на живом тестовом стеке. ### 12.2. Что должно быть сделано - добавить local fixtures: - REST SSE server - gRPC server-streaming server - WebSocket event server - SOAP fixture, если adapter уже готов - добавить smoke scenarios; - включить их в CI. ### 12.3. Тесты - tool call end-to-end through mcp-server; - usage/logging on streaming calls; - session cleanup; - async job cleanup. ## 13. Slice 10: `feat/websocket-upstream-adapter` ### 13.1. Файлы Добавить: - `crates/crank-adapter-websocket/Cargo.toml` - `crates/crank-adapter-websocket/src/lib.rs` - `crates/crank-adapter-websocket/src/error.rs` - `crates/crank-adapter-websocket/src/client.rs` - `crates/crank-adapter-websocket/src/session.rs` Изменить: - [crates/crank-core/src/protocol.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-core/src/protocol.rs) - [crates/crank-core/src/operation.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-core/src/operation.rs) - [crates/crank-runtime/src/executor.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-runtime/src/executor.rs) ### 13.2. DoD - bounded WebSocket window; - session support; - reconnect and heartbeat policy; - capability matrix exposed in admin-api and UI. ## 14. Slice 11: `feat/soap-architecture-and-core-model` ### 14.1. Файлы Изменить: - [crates/crank-core/src/protocol.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-core/src/protocol.rs) - [crates/crank-core/src/operation.rs](/home/a.tolmachev/code/rust/mcpaas/crates/crank-core/src/operation.rs) - [crates/crank-schema](/home/a.tolmachev/code/rust/mcpaas/crates/crank-schema) Добавить: - `crates/crank-core/src/soap.rs` ### 14.2. DoD - SOAP target model; - WSDL/XSD-derived metadata model; - XML normalization strategy documented in code-level types. ## 15. Slice 12: `feat/soap-adapter-foundation` ### 15.1. Файлы Добавить: - `crates/crank-adapter-soap/Cargo.toml` - `crates/crank-adapter-soap/src/lib.rs` - `crates/crank-adapter-soap/src/error.rs` - `crates/crank-adapter-soap/src/wsdl.rs` - `crates/crank-adapter-soap/src/xml.rs` - `crates/crank-adapter-soap/src/client.rs` Изменить: - `apps/admin-api` - `apps/ui` - `crates/crank-runtime` ### 15.2. DoD - WSDL import; - service/port/operation selection; - SOAP request-response execution; - fault normalization; - UI test-run support. ## 16. State transition tables ## 16.1. `StreamSession` Allowed: - `created -> running` - `running -> running` - `running -> stopped` - `running -> expired` - `running -> failed` - `stopped -> deleted` - `expired -> deleted` - `failed -> deleted` Forbidden: - `stopped -> running` - `expired -> running` - `deleted -> *` ## 16.2. `AsyncJob` Allowed: - `created -> running` - `running -> running` - `running -> completed` - `running -> failed` - `running -> cancelled` - `completed -> expired` - `failed -> expired` - `cancelled -> expired` Forbidden: - `completed -> running` - `cancelled -> running` ## 17. Sequence outlines ## 17.1. Window 1. MCP client calls tool. 2. `mcp-server` resolves tool binding. 3. `runtime.execute_window_operation`. 4. adapter collects bounded upstream data. 5. runtime aggregates and truncates. 6. `mcp-server` returns final JSON-RPC result. ## 17.2. Session 1. MCP client calls `{tool}_start`. 2. runtime starts upstream session and persists `StreamSession`. 3. client calls `{tool}_poll`. 4. runtime loads session and collects next bounded chunk. 5. client calls `{tool}_stop`. 6. runtime closes upstream and marks session stopped. ## 17.3. Async Job 1. MCP client calls `{tool}_start`. 2. runtime starts long-running job and persists `AsyncJob`. 3. client calls `{tool}_status`. 4. runtime returns current progress. 5. client calls `{tool}_result`. 6. runtime returns final normalized result if ready. ## 18. Acceptance checklist per slice Перед merge каждого slice агент обязан проверить: - docs updated if contract changed; - `just fmt-check` - `just check` - `just clippy` - `just test` - relevant Playwright/e2e if UI touched; - no hidden feature flags without docs; - no unsupported combinations exposed in UI. ## 19. Branch and commit policy Ожидаемые branch names: - `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/websocket-upstream-adapter` - `feat/soap-architecture-and-core-model` - `feat/soap-adapter-foundation` Ожидаемые commit classes: - `docs: ...` - `feat: ...` - `test: ...` - `refactor: ...` - `fix: ...` ## 20. Practical rule for agents Агент не должен брать следующий slice, пока не выполнены DoD и тесты предыдущего slice. Если slice требует новый contract: 1. обновить docs; 2. обновить tests; 3. обновить code; 4. только потом переходить дальше.