diff --git a/Cargo.lock b/Cargo.lock index d690f90..882a9a6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -391,6 +391,20 @@ dependencies = [ "tokio", ] +[[package]] +name = "crank-adapter-websocket" +version = "0.1.0" +dependencies = [ + "crank-core", + "futures-util", + "reqwest", + "serde", + "serde_json", + "thiserror", + "tokio", + "tokio-tungstenite", +] + [[package]] name = "crank-core" version = "0.1.0" @@ -448,6 +462,7 @@ dependencies = [ "crank-adapter-graphql", "crank-adapter-grpc", "crank-adapter-rest", + "crank-adapter-websocket", "crank-core", "crank-mapping", "crank-schema", @@ -456,6 +471,7 @@ dependencies = [ "serde_json", "thiserror", "tokio", + "tokio-tungstenite", ] [[package]] @@ -519,6 +535,12 @@ dependencies = [ "cipher", ] +[[package]] +name = "data-encoding" +version = "2.10.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d7a1e2f27636f116493b8b860f5546edb47c8d8f8ea73e1d2a20be88e28d1fea" + [[package]] name = "deranged" version = "0.5.8" @@ -2125,6 +2147,17 @@ dependencies = [ "syn", ] +[[package]] +name = "sha1" +version = "0.10.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e3bf829a2d51ab4a5ddf1352d8470c140cadc8301b2ae1789db023f01cedd6ba" +dependencies = [ + "cfg-if", + "cpufeatures", + "digest", +] + [[package]] name = "sha2" version = "0.10.9" @@ -2495,6 +2528,22 @@ dependencies = [ "tokio", ] +[[package]] +name = "tokio-tungstenite" +version = "0.26.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7a9daff607c6d2bf6c16fd681ccb7eecc83e4e2cdc1ca067ffaadfca5de7f084" +dependencies = [ + "futures-util", + "log", + "rustls", + "rustls-pki-types", + "tokio", + "tokio-rustls", + "tungstenite", + "webpki-roots 0.26.11", +] + [[package]] name = "tokio-util" version = "0.7.18" @@ -2693,6 +2742,25 @@ version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" +[[package]] +name = "tungstenite" +version = "0.26.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4793cb5e56680ecbb1d843515b23b6de9a75eb04b66643e256a396d43be33c13" +dependencies = [ + "bytes", + "data-encoding", + "http", + "httparse", + "log", + "rand 0.9.2", + "rustls", + "rustls-pki-types", + "sha1", + "thiserror", + "utf-8", +] + [[package]] name = "typenum" version = "1.19.0" @@ -2772,6 +2840,12 @@ dependencies = [ "serde", ] +[[package]] +name = "utf-8" +version = "0.7.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "09cc8ee72d2a9becf2f2febe0205bbed8fc6615b7cb429ad062dc7b7ddd036a9" + [[package]] name = "utf8_iter" version = "1.0.4" diff --git a/Cargo.toml b/Cargo.toml index beb18c3..d5a4d81 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -11,6 +11,7 @@ members = [ "crates/crank-adapter-rest", "crates/crank-adapter-graphql", "crates/crank-adapter-grpc", + "crates/crank-adapter-websocket", ] resolver = "3" @@ -46,4 +47,5 @@ tonic-build = "0.14" tonic-prost-build = "0.14" tracing = "0.1" tracing-subscriber = { version = "0.3", features = ["env-filter", "fmt"] } +tokio-tungstenite = { version = "0.26", features = ["rustls-tls-webpki-roots"] } uuid = { version = "1", features = ["serde", "v7"] } diff --git a/TASKS.md b/TASKS.md index 35b5bb0..7d0827b 100644 --- a/TASKS.md +++ b/TASKS.md @@ -2,25 +2,19 @@ ## Current -### `feat/streaming-e2e` +### `feat/websocket-upstream-adapter` Status: completed DoD: -- local fixture stack includes: - - REST SSE server - - gRPC server-streaming server - - WebSocket event server -- Playwright covers: - - REST window flow - - gRPC streaming flow - - session tool flow - - async job flow -- CI runs streaming e2e against the local fixture stack +- bounded WebSocket window execution is supported +- WebSocket protocol is exposed through capability APIs +- session mode can seed from WebSocket window collection +- runtime and adapter tests cover WebSocket collection and reconnects ## Next -- `feat/websocket-upstream-adapter` +- `feat/soap-architecture-and-core-model` ## Backlog diff --git a/apps/admin-api/src/error.rs b/apps/admin-api/src/error.rs index b0fe3f2..28b241b 100644 --- a/apps/admin-api/src/error.rs +++ b/apps/admin-api/src/error.rs @@ -238,6 +238,7 @@ fn runtime_test_failure_code(error: &RuntimeError) -> &'static str { RuntimeError::GraphqlAdapter(_) => "runtime_graphql_error", RuntimeError::GrpcAdapter(_) => "runtime_grpc_error", RuntimeError::RestAdapter(_) => "runtime_rest_error", + RuntimeError::WebsocketAdapter(_) => "runtime_websocket_error", RuntimeError::UnsupportedProtocol { .. } => "runtime_protocol_error", RuntimeError::InvalidPreparedRequest { .. } => "runtime_request_error", RuntimeError::MissingStreamingConfig { .. } => "runtime_streaming_config_error", diff --git a/apps/admin-api/src/service.rs b/apps/admin-api/src/service.rs index afaeab3..dab277a 100644 --- a/apps/admin-api/src/service.rs +++ b/apps/admin-api/src/service.rs @@ -1347,10 +1347,15 @@ impl AdminService { } pub async fn list_protocol_capabilities(&self) -> Vec { - [Protocol::Rest, Protocol::Graphql, Protocol::Grpc] - .into_iter() - .map(protocol_capability_view) - .collect() + [ + Protocol::Rest, + Protocol::Graphql, + Protocol::Grpc, + Protocol::Websocket, + ] + .into_iter() + .map(protocol_capability_view) + .collect() } #[instrument(skip(self))] @@ -3601,6 +3606,7 @@ fn validate_protocol_target(protocol: Protocol, target: &Target) -> Result<(), A (Protocol::Rest, Target::Rest(_)) | (Protocol::Graphql, Target::Graphql(_)) | (Protocol::Grpc, Target::Grpc(_)) + | (Protocol::Websocket, Target::Websocket(_)) ); if is_match { @@ -3827,6 +3833,7 @@ fn demo_grpc_operation_payload() -> OperationPayload { headers: BTreeMap::new(), protocol_options: Some(crank_core::ProtocolOptions { grpc: Some(crank_core::GrpcProtocolOptions { use_tls: false }), + websocket: None, }), streaming: None, }, @@ -4007,6 +4014,7 @@ fn runtime_error_code(error: &RuntimeError) -> &'static str { RuntimeError::GraphqlAdapter(_) => "graphql_error", RuntimeError::GrpcAdapter(_) => "grpc_error", RuntimeError::RestAdapter(_) => "rest_error", + RuntimeError::WebsocketAdapter(_) => "websocket_error", RuntimeError::UnsupportedProtocol { .. } => "unsupported_protocol", RuntimeError::MissingStreamingConfig { .. } => "streaming_config_error", RuntimeError::UnsupportedExecutionMode { .. } => "streaming_mode_error", @@ -4034,7 +4042,7 @@ fn protocol_capability_view(protocol: Protocol) -> ProtocolCapabilityView { .filter(|behavior| protocol.supports_transport_behavior(*behavior)) .collect(); let supports_upload_artifacts = match protocol { - Protocol::Rest | Protocol::Graphql => Vec::new(), + Protocol::Rest | Protocol::Graphql | Protocol::Websocket => Vec::new(), Protocol::Grpc => vec!["proto".to_owned(), "descriptor_set".to_owned()], }; diff --git a/apps/mcp-server/src/app.rs b/apps/mcp-server/src/app.rs index b48c1cd..8488684 100644 --- a/apps/mcp-server/src/app.rs +++ b/apps/mcp-server/src/app.rs @@ -1725,6 +1725,7 @@ fn runtime_error_code(error: &RuntimeError) -> &'static str { RuntimeError::GraphqlAdapter(_) => "adapter_execution_error", RuntimeError::GrpcAdapter(_) => "adapter_execution_error", RuntimeError::RestAdapter(_) => "adapter_execution_error", + RuntimeError::WebsocketAdapter(_) => "adapter_execution_error", RuntimeError::UnsupportedProtocol { .. } => "unsupported_protocol", RuntimeError::MissingStreamingConfig { .. } => "streaming_config_error", RuntimeError::UnsupportedExecutionMode { .. } => "streaming_mode_error", diff --git a/apps/ui/html/workspace-setup.html b/apps/ui/html/workspace-setup.html index 97b3eb1..2782632 100644 --- a/apps/ui/html/workspace-setup.html +++ b/apps/ui/html/workspace-setup.html @@ -95,6 +95,7 @@ +
Pre-selected in the operation wizard.
diff --git a/apps/ui/index.html b/apps/ui/index.html index fa03f88..8a04077 100644 --- a/apps/ui/index.html +++ b/apps/ui/index.html @@ -181,13 +181,13 @@