14 KiB
Streaming Runtime Design
1. Назначение документа
Этот документ фиксирует, как потоковая модель должна быть разложена по crates, traits, services и функциям.
Цель:
- избежать god-services;
- зафиксировать boundaries между
core,registry,runtime, protocol adapters иmcp-server; - расписать ожидаемые функции и их ответственность до уровня инженерной спецификации.
2. Общая схема слоев
2.1. crank-core
Отвечает за:
- типы
ExecutionMode,StreamingConfig,AggregationMode; - типы
StreamSession,AsyncJobHandle,StreamStatus,JobStatus; - validation helpers без инфраструктуры.
2.2. crank-registry
Отвечает за:
- persistence session/job state;
- optimistic transitions;
- cleanup expired state.
2.3. crank-runtime
Отвечает за:
- orchestration execution modes;
- dispatch в protocol adapter;
- bounded collection;
- aggregation;
- state transitions;
- invocation logging.
2.4. Protocol adapters
Отвечают за:
- protocol-specific upstream behavior;
- connect/request/collect/close;
- преобразование upstream payload в нормализованный JSON.
2.5. apps/mcp-server
Отвечает за:
- downstream MCP transport;
- tool-family generation;
- JSON-RPC lifecycle;
- transport/session routing.
2.6. apps/admin-api
Отвечает за:
- config CRUD;
- test-runs;
- validation;
- UI-oriented metadata and previews.
3. Core domain types
3.1. ExecutionMode
pub enum ExecutionMode {
Unary,
Window,
Session,
AsyncJob,
}
Методы:
fn is_stateful(&self) -> boolfn requires_tool_family(&self) -> bool
3.2. AggregationMode
pub enum AggregationMode {
RawItems,
SummaryOnly,
SummaryPlusSamples,
Stats,
LatestState,
}
Методы:
fn needs_items(&self) -> boolfn needs_summary(&self) -> bool
3.3. StreamingConfig
pub struct StreamingConfig {
pub mode: ExecutionMode,
pub transport_behavior: TransportBehavior,
pub window_duration_ms: Option<u64>,
pub poll_interval_ms: Option<u64>,
pub upstream_timeout_ms: Option<u64>,
pub idle_timeout_ms: Option<u64>,
pub max_session_lifetime_ms: Option<u64>,
pub max_items: Option<u32>,
pub max_bytes: Option<u32>,
pub aggregation_mode: AggregationMode,
pub summary_path: Option<String>,
pub items_path: Option<String>,
pub cursor_path: Option<String>,
pub status_path: Option<String>,
pub done_path: Option<String>,
pub redacted_paths: Vec<String>,
pub truncate_item_fields: bool,
pub max_field_length: Option<u32>,
pub drop_duplicates: bool,
pub sampling_rate: Option<f64>,
pub tool_family: ToolFamilyConfig,
}
Методы:
fn validate_common(&self) -> Result<(), StreamingConfigError>fn validate_for_protocol(&self, protocol: Protocol) -> Result<(), StreamingConfigError>fn effective_max_items(&self, limits: &StreamingLimits) -> u32fn effective_max_bytes(&self, limits: &StreamingLimits) -> u32fn effective_timeout(&self, defaults: &StreamingDefaults) -> Duration
3.4. StreamSession
pub struct StreamSession {
pub id: StreamSessionId,
pub workspace_id: WorkspaceId,
pub agent_id: Option<AgentId>,
pub operation_id: OperationId,
pub protocol: Protocol,
pub mode: ExecutionMode,
pub status: StreamStatus,
pub cursor: Option<Value>,
pub state: Value,
pub expires_at: DateTime<Utc>,
pub last_poll_at: Option<DateTime<Utc>>,
pub created_at: DateTime<Utc>,
pub closed_at: Option<DateTime<Utc>>,
}
Методы:
fn is_expired(&self, now: DateTime<Utc>) -> boolfn can_poll(&self, now: DateTime<Utc>) -> boolfn mark_polled(&mut self, now: DateTime<Utc>)fn mark_closed(&mut self, now: DateTime<Utc>)
3.5. AsyncJobHandle
pub struct AsyncJobHandle {
pub id: AsyncJobId,
pub workspace_id: WorkspaceId,
pub agent_id: Option<AgentId>,
pub operation_id: OperationId,
pub status: JobStatus,
pub progress: Value,
pub result: Option<Value>,
pub error: Option<Value>,
pub expires_at: Option<DateTime<Utc>>,
pub created_at: DateTime<Utc>,
pub updated_at: DateTime<Utc>,
pub finished_at: Option<DateTime<Utc>>,
}
Методы:
fn is_finished(&self) -> boolfn can_cancel(&self) -> boolfn mark_finished(&mut self, now: DateTime<Utc>, result: Value)fn mark_failed(&mut self, now: DateTime<Utc>, error: Value)
4. Registry contracts
4.1. StreamSessionRepository
#[async_trait]
pub trait StreamSessionRepository {
async fn create_stream_session(&self, session: NewStreamSession) -> Result<StreamSession, RegistryError>;
async fn get_stream_session(&self, id: &StreamSessionId) -> Result<Option<StreamSession>, RegistryError>;
async fn update_stream_session_state(&self, update: StreamSessionStateUpdate) -> Result<StreamSession, RegistryError>;
async fn close_stream_session(&self, id: &StreamSessionId, now: DateTime<Utc>) -> Result<(), RegistryError>;
async fn list_stream_sessions(&self, filter: StreamSessionFilter) -> Result<Page<StreamSession>, RegistryError>;
async fn delete_expired_stream_sessions(&self, now: DateTime<Utc>) -> Result<u64, RegistryError>;
}
4.2. AsyncJobRepository
#[async_trait]
pub trait AsyncJobRepository {
async fn create_async_job(&self, job: NewAsyncJob) -> Result<AsyncJobHandle, RegistryError>;
async fn get_async_job(&self, id: &AsyncJobId) -> Result<Option<AsyncJobHandle>, RegistryError>;
async fn update_async_job_status(&self, update: AsyncJobStatusUpdate) -> Result<AsyncJobHandle, RegistryError>;
async fn cancel_async_job(&self, id: &AsyncJobId, now: DateTime<Utc>) -> Result<(), RegistryError>;
async fn list_async_jobs(&self, filter: AsyncJobFilter) -> Result<Page<AsyncJobHandle>, RegistryError>;
async fn delete_expired_async_jobs(&self, now: DateTime<Utc>) -> Result<u64, RegistryError>;
}
4.3. Storage rules
- update methods должны использовать optimistic concurrency, если state transitions конфликтуют;
pollне должен silently reopen closed session;- cleanup job должен быть отдельным service;
- session/job state не должен хранить необрезанные raw payloads без лимитов.
5. Runtime contracts
5.1. ProtocolAdapter
Базовый trait не должен пытаться описать все streaming cases в одном методе.
#[async_trait]
pub trait ProtocolAdapter {
async fn execute_unary(&self, ctx: UnaryExecutionContext) -> Result<NormalizedResponse, AdapterError>;
async fn execute_window(&self, ctx: WindowExecutionContext) -> Result<WindowExecutionResult, AdapterError>;
async fn start_session(&self, ctx: SessionStartContext) -> Result<SessionStartResult, AdapterError>;
async fn poll_session(&self, ctx: SessionPollContext) -> Result<SessionPollResult, AdapterError>;
async fn stop_session(&self, ctx: SessionStopContext) -> Result<(), AdapterError>;
async fn start_async_job(&self, ctx: AsyncJobStartContext) -> Result<AsyncJobStartResult, AdapterError>;
async fn get_async_job_status(&self, ctx: AsyncJobStatusContext) -> Result<AsyncJobStatusResult, AdapterError>;
async fn get_async_job_result(&self, ctx: AsyncJobResultContext) -> Result<NormalizedResponse, AdapterError>;
async fn cancel_async_job(&self, ctx: AsyncJobCancelContext) -> Result<(), AdapterError>;
}
Не каждый adapter обязан поддерживать все методы. Capability mismatch должен проверяться выше.
5.2. OperationExecutor
Главная orchestration точка в crank-runtime.
Ожидаемые функции:
execute_unary_operationexecute_window_operationstart_stream_sessionpoll_stream_sessionstop_stream_sessionstart_async_jobget_async_job_statusget_async_job_resultcancel_async_job
execute_unary_operation
Отвечает за:
- schema validation;
- auth resolution;
- adapter dispatch;
- output mapping;
- observability write.
execute_window_operation
Отвечает за:
- schema validation;
- auth resolution;
- adapter dispatch в
execute_window; - aggregation;
- truncation flags;
- observability write.
start_stream_session
Отвечает за:
- schema validation;
- auth resolution;
- adapter
start_session; - registry
create_stream_session; - initial preview result;
- observability write.
poll_stream_session
Отвечает за:
- load session from registry;
- expiration check;
- adapter
poll_session; - registry
update_stream_session_state; - output mapping;
- observability write.
stop_stream_session
Отвечает за:
- load session;
- adapter
stop_session; - registry close;
- observability write.
start_async_job
Отвечает за:
- schema validation;
- adapter
start_async_job; - registry
create_async_job; - observability write.
get_async_job_status
Отвечает за:
- registry load;
- adapter
get_async_job_status, если upstream status is live; - registry update;
- normalized status result.
get_async_job_result
Отвечает за:
- job readiness check;
- adapter
get_async_job_result, если result lazy-loaded; - output mapping;
- observability write.
cancel_async_job
Отвечает за:
- adapter cancel;
- registry cancel transition;
- observability write.
5.3. Aggregation services
Нужно отделить aggregation от adapters.
Отдельные services:
WindowAggregatorSummaryBuilderCursorTrackerPayloadLimiterRedactionService
WindowAggregator
Функции:
collect_itemsapply_item_limitapply_byte_limitmark_truncated
SummaryBuilder
Функции:
build_summarybuild_stats_summarybuild_latest_state_summarybuild_summary_plus_samples
PayloadLimiter
Функции:
truncate_item_fieldstruncate_bytesenforce_max_items
RedactionService
Функции:
redact_pathsredact_objectredact_item
6. Adapter responsibilities
6.1. REST adapter
Функции:
execute_rest_unaryexecute_rest_windowstart_rest_sessionpoll_rest_sessionstop_rest_sessionstart_rest_async_jobget_rest_async_job_statusget_rest_async_job_resultcancel_rest_async_job
Дополнительные helpers:
open_sse_streamcollect_sse_windowparse_sse_eventclose_sse_stream
6.2. GraphQL adapter
Функции:
execute_graphql_unary
Отдельно отложено:
execute_graphql_subscription_windowstart_graphql_subscription_session
6.3. gRPC adapter
Функции:
execute_grpc_unaryexecute_grpc_windowstart_grpc_sessionpoll_grpc_sessionstop_grpc_sessionstart_grpc_async_jobget_grpc_async_job_statusget_grpc_async_job_resultcancel_grpc_async_job
Helpers:
invoke_unary_methodopen_server_streamcollect_server_stream_windowdecode_stream_item
6.4. WebSocket adapter
Функции:
execute_websocket_windowstart_websocket_sessionpoll_websocket_sessionstop_websocket_sessionstart_websocket_async_jobget_websocket_async_job_statuscancel_websocket_async_job
Helpers:
connect_websocketsend_subscribe_messagesend_unsubscribe_messageread_next_framedecode_text_frameheartbeat_tickreconnect_if_needed
6.5. SOAP adapter
Функции:
execute_soap_unarystart_soap_async_jobget_soap_async_job_statusget_soap_async_job_resultcancel_soap_async_job
Helpers:
render_soap_enveloperender_soap_headersparse_soap_envelopeparse_soap_faultnormalize_xml_value
7. MCP server design
7.1. Route handlers
Ожидаемые handlers:
handle_initializehandle_notifications_initializedhandle_pinghandle_tools_listhandle_tools_callhandle_get_sse_streamhandle_delete_session
7.2. Tool-family generation
Функции:
build_unary_tool_definitionbuild_window_tool_definitionbuild_session_tool_familybuild_async_job_tool_family
Для session:
build_session_start_toolbuild_session_poll_toolbuild_session_stop_tool
Для async_job:
build_async_job_start_toolbuild_async_job_status_toolbuild_async_job_result_toolbuild_async_job_cancel_tool
7.3. Transport session management
Функции:
ensure_protocol_versionresolve_transport_modeensure_mcp_sessionattach_session_headersstream_sse_eventclose_transport_session
8. Admin API service design
Ожидаемые service функции:
validate_streaming_configlist_protocol_capabilitieslist_streaming_presetsstart_stream_test_runpoll_stream_test_runstop_stream_test_runget_stream_test_resultlist_stream_sessionsget_stream_sessionstop_stream_sessiondelete_stream_sessionlist_async_jobsget_async_jobcancel_async_jobget_async_job_result
9. Cleanup jobs
Отдельные jobs:
expire_stream_sessionsexpire_async_jobsreap_orphaned_transport_sessionscompact_stream_payloads, если будет нужен storage optimization layer
10. Design rules
- adapters не агрегируют product-level summaries;
- runtime не знает деталей transport framing;
mcp-serverне знает про upstream protocols;admin-apiне управляет live protocol connections напрямую;- session/job lifecycle должен быть явным в names и contracts;
- cancellation не должна зависеть от TCP disconnect;
- bounded limits должны применяться раньше, чем payload попадет в final result.