17 KiB
Streaming Implementation Specification
1. Назначение документа
Этот документ переводит streaming architecture в execution-oriented plan, пригодный для агентной разработки.
Он отвечает на вопросы:
- в каком порядке реализовывать фичи;
- какие файлы менять в каждом срезе;
- какие новые модули создавать;
- какие тесты обязательны;
- какие риски и зависимости есть между срезами;
- какой Definition of Done нужен для каждого шага.
Документ дополняет:
- streaming-mcp-plan.md
- streaming-admin-api.md
- streaming-runtime-design.md
- streaming-ui-contract.md
- 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
Порядок реализации:
feat/mcp-streamable-http-alignmentfeat/streaming-core-modelfeat/stream-session-storefeat/runtime-window-modefeat/rest-sse-adapterfeat/grpc-server-streaming-adapterfeat/session-and-job-toolsfeat/streaming-ui-configfeat/streaming-e2efeat/websocket-upstream-adapterfeat/soap-architecture-and-core-modelfeat/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
- apps/mcp-server/src/jsonrpc.rs
- apps/mcp-server/src/session.rs
- apps/mcp-server/src/main.rs
Добавить при необходимости:
apps/mcp-server/src/transport.rsapps/mcp-server/src/sse.rs
4.3. Что должно быть сделано
- parse
Acceptкорректно дляapplication/jsonиtext/event-stream; GETendpoint перестает быть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;
GETSSE 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
- crates/crank-core/src/operation.rs
- crates/crank-core/src/ids.rs
- crates/crank-core/src/protocol.rs
Добавить:
crates/crank-core/src/streaming.rscrates/crank-core/src/stream_session.rs
5.3. Что должно быть сделано
- добавить
ExecutionMode; - добавить
TransportBehavior; - добавить
AggregationMode; - добавить
StreamingConfig; - добавить
StreamSessionId,AsyncJobId; - добавить
StreamSession,AsyncJobHandle,StreamStatus,JobStatus; - вшить
streaming: Option<StreamingConfig>в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
- crates/crank-registry/src/model.rs
- crates/crank-registry/src/migrations.rs
- 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
- crates/crank-runtime/src/model.rs
- crates/crank-runtime/src/executor.rs
- crates/crank-runtime/src/error.rs
Добавить:
crates/crank-runtime/src/streaming.rscrates/crank-runtime/src/aggregation.rscrates/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. Файлы
Изменить:
Ожидаемые файлы:
src/lib.rssrc/error.rssrc/client.rssrc/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. Файлы
Изменить:
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
- apps/mcp-server/src/app.rs
- apps/mcp-server/src/catalog.rs
- apps/mcp-server/src/session.rs
- 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/*.htmlapps/ui/js/wizard.jsapps/ui/js/api.jsapps/ui/js/i18n.js- CSS файлы wizard/settings/pages по необходимости
Добавить:
apps/ui/js/streaming-form.jsapps/ui/js/stream-test-run.jsapps/ui/html/stream-sessions.htmlapps/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.tomlcrates/crank-adapter-websocket/src/lib.rscrates/crank-adapter-websocket/src/error.rscrates/crank-adapter-websocket/src/client.rscrates/crank-adapter-websocket/src/session.rs
Изменить:
- crates/crank-core/src/protocol.rs
- crates/crank-core/src/operation.rs
- 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/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.tomlcrates/crank-adapter-soap/src/lib.rscrates/crank-adapter-soap/src/error.rscrates/crank-adapter-soap/src/wsdl.rscrates/crank-adapter-soap/src/xml.rscrates/crank-adapter-soap/src/client.rs
Изменить:
apps/admin-apiapps/uicrates/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 -> runningrunning -> runningrunning -> stoppedrunning -> expiredrunning -> failedstopped -> deletedexpired -> deletedfailed -> deleted
Forbidden:
stopped -> runningexpired -> runningdeleted -> *
16.2. AsyncJob
Allowed:
created -> runningrunning -> runningrunning -> completedrunning -> failedrunning -> cancelledcompleted -> expiredfailed -> expiredcancelled -> expired
Forbidden:
completed -> runningcancelled -> running
17. Sequence outlines
17.1. Window
- MCP client calls tool.
mcp-serverresolves tool binding.runtime.execute_window_operation.- adapter collects bounded upstream data.
- runtime aggregates and truncates.
mcp-serverreturns final JSON-RPC result.
17.2. Session
- MCP client calls
{tool}_start. - runtime starts upstream session and persists
StreamSession. - client calls
{tool}_poll. - runtime loads session and collects next bounded chunk.
- client calls
{tool}_stop. - runtime closes upstream and marks session stopped.
17.3. Async Job
- MCP client calls
{tool}_start. - runtime starts long-running job and persists
AsyncJob. - client calls
{tool}_status. - runtime returns current progress.
- client calls
{tool}_result. - runtime returns final normalized result if ready.
18. Acceptance checklist per slice
Перед merge каждого slice агент обязан проверить:
- docs updated if contract changed;
just fmt-checkjust checkjust clippyjust 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-alignmentfeat/streaming-core-modelfeat/stream-session-storefeat/runtime-window-modefeat/rest-sse-adapterfeat/grpc-server-streaming-adapterfeat/session-and-job-toolsfeat/streaming-ui-configfeat/streaming-e2efeat/websocket-upstream-adapterfeat/soap-architecture-and-core-modelfeat/soap-adapter-foundation
Ожидаемые commit classes:
docs: ...feat: ...test: ...refactor: ...fix: ...
20. Practical rule for agents
Агент не должен брать следующий slice, пока не выполнены DoD и тесты предыдущего slice.
Если slice требует новый contract:
- обновить docs;
- обновить tests;
- обновить code;
- только потом переходить дальше.