diff --git a/Cargo.lock b/Cargo.lock index 4e60815..6ce1cb8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -382,6 +382,7 @@ version = "0.1.0" dependencies = [ "axum", "crank-core", + "futures-util", "reqwest", "serde", "serde_json", @@ -449,6 +450,7 @@ dependencies = [ "crank-core", "crank-mapping", "crank-schema", + "futures-util", "serde", "serde_json", "thiserror", @@ -1868,6 +1870,7 @@ dependencies = [ "cookie", "cookie_store", "futures-core", + "futures-util", "http", "http-body", "http-body-util", @@ -1887,12 +1890,14 @@ dependencies = [ "sync_wrapper", "tokio", "tokio-rustls", + "tokio-util", "tower", "tower-http", "tower-service", "url", "wasm-bindgen", "wasm-bindgen-futures", + "wasm-streams", "web-sys", "webpki-roots 1.0.6", ] @@ -2916,6 +2921,19 @@ dependencies = [ "wasmparser", ] +[[package]] +name = "wasm-streams" +version = "0.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "15053d8d85c7eccdbefef60f06769760a563c7f0a9d6902a13d35c7800b0ad65" +dependencies = [ + "futures-util", + "js-sys", + "wasm-bindgen", + "wasm-bindgen-futures", + "web-sys", +] + [[package]] name = "wasmparser" version = "0.244.0" diff --git a/TASKS.md b/TASKS.md index d7e8887..dcc08c9 100644 --- a/TASKS.md +++ b/TASKS.md @@ -2,20 +2,20 @@ ## Current -### `feat/runtime-window-mode` +### `feat/rest-sse-adapter` Status: completed DoD: -- runtime exposes bounded `window` execution -- item and byte limits are enforced -- `window_complete`, `truncated`, `has_more` are returned -- aggregation and redaction are covered by tests -- unary execution remains intact +- REST adapter supports bounded SSE collection for `window` mode +- SSE events are parsed into normalized JSON items +- adapter respects item/window limits and upstream timeout +- runtime dispatches REST `window` mode through SSE adapter path +- local SSE upstream tests cover success, timeout and malformed events ## Next -- `feat/rest-sse-adapter` +- `feat/grpc-server-streaming-adapter` ## Backlog diff --git a/crates/crank-adapter-rest/Cargo.toml b/crates/crank-adapter-rest/Cargo.toml index 090efc7..9e78916 100644 --- a/crates/crank-adapter-rest/Cargo.toml +++ b/crates/crank-adapter-rest/Cargo.toml @@ -7,10 +7,12 @@ version.workspace = true [dependencies] crank-core = { path = "../crank-core" } -reqwest.workspace = true +futures-util = "0.3" +reqwest = { workspace = true, features = ["stream"] } serde.workspace = true serde_json.workspace = true thiserror.workspace = true +tokio.workspace = true [dev-dependencies] axum.workspace = true diff --git a/crates/crank-adapter-rest/src/client.rs b/crates/crank-adapter-rest/src/client.rs index 5903dec..d5e82e3 100644 --- a/crates/crank-adapter-rest/src/client.rs +++ b/crates/crank-adapter-rest/src/client.rs @@ -7,7 +7,7 @@ use reqwest::{ }; use serde_json::Value; -use crate::{RestAdapterError, RestRequest, RestResponse}; +use crate::{RestAdapterError, RestRequest, RestResponse, RestWindowRequest, RestWindowResponse}; #[derive(Clone, Debug)] pub struct RestAdapter { @@ -62,6 +62,61 @@ impl RestAdapter { body, }) } + + pub async fn execute_window( + &self, + target: &RestTarget, + request: &RestWindowRequest, + ) -> Result { + let url = build_url(target, &request.request)?; + let mut headers = build_headers(target, &request.request)?; + headers.insert( + reqwest::header::ACCEPT, + HeaderValue::from_static("text/event-stream"), + ); + + let mut builder = self + .client + .request(to_reqwest_method(target.method), url) + .headers(headers) + .timeout(Duration::from_millis(request.request.timeout_ms)); + + if let Some(body) = &request.request.body { + builder = builder.json(body); + } + + let response = builder.send().await?; + let status = response.status(); + + if !status.is_success() { + let headers = normalize_headers(response.headers()); + let body = decode_body(response).await?; + return Err(RestAdapterError::UnexpectedStatus { + status: status.as_u16(), + body: Value::Object( + [ + ( + "headers".to_owned(), + serde_json::to_value(headers).unwrap_or(Value::Null), + ), + ("body".to_owned(), body), + ] + .into_iter() + .collect(), + ), + }); + } + + let (status_code, headers, body) = + crate::sse::collect_sse_window(response, request.window_duration_ms, request.max_items) + .await?; + + Ok(RestWindowResponse { + status_code, + headers, + body, + }) + } } fn build_url(target: &RestTarget, request: &RestRequest) -> Result { @@ -172,13 +227,15 @@ mod tests { Json, Router, extract::{Path, Query}, http::HeaderMap, + response::sse::{Event, KeepAlive, Sse}, routing::{get, post}, }; use crank_core::{HttpMethod, RestTarget}; + use futures_util::stream; use serde_json::{Value, json}; use tokio::net::TcpListener; - use crate::{RestAdapter, RestAdapterError, RestRequest}; + use crate::{RestAdapter, RestAdapterError, RestRequest, RestWindowRequest}; #[tokio::test] async fn executes_rest_request_and_normalizes_json_response() { @@ -242,10 +299,104 @@ mod tests { )); } + #[tokio::test] + async fn collects_sse_events_with_window_bounds() { + let base_url = spawn_test_server().await; + let adapter = RestAdapter::new(); + let target = RestTarget { + base_url, + method: HttpMethod::Get, + path_template: "/events".to_owned(), + static_headers: BTreeMap::new(), + }; + let request = RestWindowRequest { + request: RestRequest { + path_params: BTreeMap::new(), + query_params: BTreeMap::new(), + headers: BTreeMap::new(), + body: None, + timeout_ms: 1_000, + }, + window_duration_ms: 1_000, + max_items: Some(2), + }; + + let response = adapter.execute_window(&target, &request).await.unwrap(); + + assert_eq!(response.status_code, 200); + assert_eq!( + response.body, + json!({ + "items": [ + { "message": "one" }, + { "message": "two" } + ], + "done": false + }) + ); + } + + #[tokio::test] + async fn returns_timeout_window_when_no_events_arrive_before_deadline() { + let base_url = spawn_test_server().await; + let adapter = RestAdapter::new(); + let target = RestTarget { + base_url, + method: HttpMethod::Get, + path_template: "/events-idle".to_owned(), + static_headers: BTreeMap::new(), + }; + let request = RestWindowRequest { + request: RestRequest { + path_params: BTreeMap::new(), + query_params: BTreeMap::new(), + headers: BTreeMap::new(), + body: None, + timeout_ms: 1_000, + }, + window_duration_ms: 50, + max_items: Some(10), + }; + + let response = adapter.execute_window(&target, &request).await.unwrap(); + + assert_eq!(response.body, json!({ "items": [], "done": true })); + } + + #[tokio::test] + async fn rejects_malformed_sse_payloads() { + let base_url = spawn_test_server().await; + let adapter = RestAdapter::new(); + let target = RestTarget { + base_url, + method: HttpMethod::Get, + path_template: "/events-broken".to_owned(), + static_headers: BTreeMap::new(), + }; + let request = RestWindowRequest { + request: RestRequest { + path_params: BTreeMap::new(), + query_params: BTreeMap::new(), + headers: BTreeMap::new(), + body: None, + timeout_ms: 1_000, + }, + window_duration_ms: 1_000, + max_items: Some(10), + }; + + let error = adapter.execute_window(&target, &request).await.unwrap_err(); + + assert!(matches!(error, RestAdapterError::InvalidSseEvent)); + } + async fn spawn_test_server() -> String { let app = Router::new() .route("/users/{user_id}", post(create_user)) - .route("/fail", get(fail)); + .route("/fail", get(fail)) + .route("/events", get(sse_events)) + .route("/events-idle", get(sse_idle)) + .route("/events-broken", get(sse_broken)); let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); @@ -286,4 +437,27 @@ mod tests { Json(json!({ "error": "upstream failed" })), ) } + + async fn sse_events() + -> Sse>> { + let events = vec![ + Ok(Event::default().data("{\"message\":\"one\"}")), + Ok(Event::default().data("{\"message\":\"two\"}")), + Ok(Event::default().data("{\"message\":\"three\"}")), + ]; + + Sse::new(stream::iter(events)).keep_alive(KeepAlive::default()) + } + + async fn sse_idle() + -> Sse>> { + Sse::new(stream::pending()).keep_alive(KeepAlive::default()) + } + + async fn sse_broken() + -> Sse>> { + let events = vec![Ok(Event::default().data("{broken-json}"))]; + + Sse::new(stream::iter(events)).keep_alive(KeepAlive::default()) + } } diff --git a/crates/crank-adapter-rest/src/error.rs b/crates/crank-adapter-rest/src/error.rs index 1f18032..64f732d 100644 --- a/crates/crank-adapter-rest/src/error.rs +++ b/crates/crank-adapter-rest/src/error.rs @@ -15,6 +15,10 @@ pub enum RestAdapterError { InvalidHeaderValue { header: String }, #[error("request failed")] Transport(#[from] reqwest::Error), + #[error("sse collection window expired before stream completed")] + WindowExpired, #[error("rest endpoint returned status {status}")] UnexpectedStatus { status: u16, body: Value }, + #[error("sse stream produced malformed event payload")] + InvalidSseEvent, } diff --git a/crates/crank-adapter-rest/src/lib.rs b/crates/crank-adapter-rest/src/lib.rs index 8f25b10..94ee256 100644 --- a/crates/crank-adapter-rest/src/lib.rs +++ b/crates/crank-adapter-rest/src/lib.rs @@ -1,7 +1,8 @@ mod client; mod error; mod model; +mod sse; pub use client::RestAdapter; pub use error::RestAdapterError; -pub use model::{RestRequest, RestResponse}; +pub use model::{RestRequest, RestResponse, RestWindowRequest, RestWindowResponse}; diff --git a/crates/crank-adapter-rest/src/model.rs b/crates/crank-adapter-rest/src/model.rs index f04ecec..f2c7e36 100644 --- a/crates/crank-adapter-rest/src/model.rs +++ b/crates/crank-adapter-rest/src/model.rs @@ -23,3 +23,20 @@ pub struct RestResponse { pub headers: BTreeMap, pub body: Value, } + +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, Default)] +pub struct RestWindowRequest { + #[serde(flatten)] + pub request: RestRequest, + pub window_duration_ms: u64, + #[serde(skip_serializing_if = "Option::is_none")] + pub max_items: Option, +} + +#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] +pub struct RestWindowResponse { + pub status_code: u16, + #[serde(default, skip_serializing_if = "BTreeMap::is_empty")] + pub headers: BTreeMap, + pub body: Value, +} diff --git a/crates/crank-adapter-rest/src/sse.rs b/crates/crank-adapter-rest/src/sse.rs new file mode 100644 index 0000000..c60ed73 --- /dev/null +++ b/crates/crank-adapter-rest/src/sse.rs @@ -0,0 +1,160 @@ +use std::collections::BTreeMap; + +use futures_util::StreamExt; +use reqwest::header::HeaderMap; +use serde_json::{Value, json}; +use tokio::time::{Duration, Instant, timeout_at}; + +use crate::RestAdapterError; + +pub async fn collect_sse_window( + response: reqwest::Response, + window_duration_ms: u64, + max_items: Option, +) -> Result<(u16, BTreeMap, Value), RestAdapterError> { + let status = response.status(); + let headers = normalize_headers(response.headers()); + let deadline = Instant::now() + Duration::from_millis(window_duration_ms); + let mut stream = response.bytes_stream(); + let mut buffer = String::new(); + let mut items = Vec::new(); + let mut done = true; + + loop { + if max_items.is_some_and(|limit| items.len() >= limit as usize) { + done = false; + break; + } + + let next_chunk = match timeout_at(deadline, stream.next()).await { + Ok(next_chunk) => next_chunk, + Err(_) => break, + }; + + let Some(next_chunk) = next_chunk else { + break; + }; + + let chunk = next_chunk?; + buffer.push_str(&String::from_utf8_lossy(&chunk)); + + while let Some(event_end) = find_event_boundary(&buffer) { + let event = buffer[..event_end].to_owned(); + let boundary_len = boundary_length(&buffer[event_end..]); + buffer = buffer[event_end + boundary_len..].to_owned(); + + if let Some(item) = parse_sse_event(&event)? { + items.push(item); + + if max_items.is_some_and(|limit| items.len() >= limit as usize) { + done = false; + break; + } + } + } + + if !done { + break; + } + } + + Ok(( + status.as_u16(), + headers, + json!({ + "items": items, + "done": done + }), + )) +} + +fn parse_sse_event(raw: &str) -> Result, RestAdapterError> { + let mut data_lines = Vec::new(); + + for line in raw.lines() { + if line.is_empty() || line.starts_with(':') { + continue; + } + + if let Some(value) = line.strip_prefix("data:") { + data_lines.push(value.trim_start().to_owned()); + } + } + + if data_lines.is_empty() { + return Ok(None); + } + + let payload = data_lines.join("\n"); + if payload.is_empty() { + return Ok(None); + } + + serde_json::from_str::(&payload) + .map(Some) + .or_else(|_| { + if payload.starts_with('{') || payload.starts_with('[') { + Err(RestAdapterError::InvalidSseEvent) + } else { + Ok(Some(Value::String(payload))) + } + }) +} + +fn find_event_boundary(buffer: &str) -> Option { + buffer + .find("\r\n\r\n") + .or_else(|| buffer.find("\n\n")) + .or_else(|| buffer.find("\r\r")) +} + +fn boundary_length(boundary: &str) -> usize { + if boundary.starts_with("\r\n\r\n") { + 4 + } else { + 2 + } +} + +fn normalize_headers(headers: &HeaderMap) -> BTreeMap { + headers + .iter() + .filter_map(|(name, value)| { + value + .to_str() + .ok() + .map(|value| (name.as_str().to_owned(), value.to_owned())) + }) + .collect() +} + +#[cfg(test)] +mod tests { + use serde_json::json; + + use super::parse_sse_event; + use crate::RestAdapterError; + + #[test] + fn parses_json_event_payload() { + let event = "event: message\ndata: {\"message\":\"ok\"}\n\n"; + + let parsed = parse_sse_event(event).unwrap(); + + assert_eq!(parsed, Some(json!({ "message": "ok" }))); + } + + #[test] + fn ignores_comment_only_events() { + let parsed = parse_sse_event(": keepalive\n\n").unwrap(); + + assert_eq!(parsed, None); + } + + #[test] + fn rejects_malformed_json_like_payload() { + let error = parse_sse_event("data: {broken-json}\n\n").unwrap_err(); + + assert!(matches!(error, RestAdapterError::InvalidSseEvent)); + } +} diff --git a/crates/crank-runtime/Cargo.toml b/crates/crank-runtime/Cargo.toml index a2ef521..c3b6334 100644 --- a/crates/crank-runtime/Cargo.toml +++ b/crates/crank-runtime/Cargo.toml @@ -19,4 +19,5 @@ thiserror.workspace = true [dev-dependencies] axum.workspace = true crank-adapter-grpc = { path = "../crank-adapter-grpc", features = ["test-support"] } +futures-util = "0.3" tokio.workspace = true diff --git a/crates/crank-runtime/src/executor.rs b/crates/crank-runtime/src/executor.rs index 16d3253..1611c37 100644 --- a/crates/crank-runtime/src/executor.rs +++ b/crates/crank-runtime/src/executor.rs @@ -2,8 +2,8 @@ use std::collections::BTreeMap; use crank_adapter_graphql::{GraphqlAdapter, GraphqlRequest}; use crank_adapter_grpc::{GrpcAdapter, GrpcRequest}; -use crank_adapter_rest::{RestAdapter, RestRequest}; -use crank_core::{ExecutionMode, Target}; +use crank_adapter_rest::{RestAdapter, RestRequest, RestWindowRequest}; +use crank_core::{ExecutionMode, Target, TransportBehavior}; use serde_json::{Map, Value, json}; use crate::{ @@ -60,7 +60,16 @@ impl RuntimeExecutor { } let prepared_request = self.prepare_request(operation, input)?; - let adapter_response = self.execute_adapter(operation, prepared_request).await?; + let adapter_response = if matches!( + streaming.transport_behavior, + TransportBehavior::ServerStream + ) && matches!(operation.target, Target::Rest(_)) + { + self.execute_window_adapter(operation, prepared_request) + .await? + } else { + self.execute_adapter(operation, prepared_request).await? + }; crate::aggregation::collect_window_result(&adapter_response.body, streaming) } @@ -159,6 +168,48 @@ impl RuntimeExecutor { } } } + + async fn execute_window_adapter( + &self, + operation: &RuntimeOperation, + prepared_request: PreparedRequest, + ) -> Result { + match &operation.target { + Target::Rest(target) => { + let Some(streaming) = operation.execution_config.streaming.as_ref() else { + return Err(RuntimeError::MissingStreamingConfig { + operation_id: operation.operation_id.as_str().to_owned(), + }); + }; + let request = RestWindowRequest { + request: RestRequest { + path_params: prepared_request.path_params.clone(), + query_params: prepared_request.query_params.clone(), + headers: merge_headers( + &target.static_headers, + &operation.execution_config.headers, + &prepared_request.headers, + ), + body: prepared_request.body.clone(), + timeout_ms: streaming + .upstream_timeout_ms + .unwrap_or(operation.execution_config.timeout_ms), + }, + window_duration_ms: streaming.window_duration_ms.unwrap_or_default(), + max_items: streaming.max_items.map(|value| value.saturating_add(1)), + }; + let response = self.rest_adapter.execute_window(target, &request).await?; + + Ok(AdapterResponse { + status_code: response.status_code, + headers: response.headers, + body: response.body, + data: Value::Null, + }) + } + _ => self.execute_adapter(operation, prepared_request).await, + } + } } impl PreparedRequest { @@ -264,7 +315,11 @@ fn non_empty_payload(value: Option) -> Option { mod tests { use std::collections::BTreeMap; - use axum::{Json, Router, routing::post}; + use axum::{ + Json, Router, + response::sse::{Event, KeepAlive, Sse}, + routing::post, + }; use crank_adapter_grpc::test_support as grpc_test_support; use crank_core::{ AggregationMode, DescriptorId, ExecutionConfig, ExecutionMode, GeneratedDraft, @@ -274,6 +329,7 @@ mod tests { }; use crank_mapping::{MappingRule, MappingSet}; use crank_schema::{Schema, SchemaKind}; + use futures_util::stream; use serde_json::{Value, json}; use tokio::net::TcpListener; @@ -378,7 +434,8 @@ mod tests { async fn executes_window_mode_with_raw_items() { let base_url = spawn_runtime_server().await; let executor = RuntimeExecutor::new(); - let operation = test_window_operation(&base_url, AggregationMode::RawItems, None, None); + let operation = + test_window_snapshot_operation(&base_url, AggregationMode::RawItems, None, None); let result = executor .execute_window(&operation, &json!({})) @@ -396,7 +453,7 @@ mod tests { let base_url = spawn_runtime_server().await; let executor = RuntimeExecutor::new(); let operation = - test_window_operation(&base_url, AggregationMode::SummaryOnly, Some(10), None); + test_window_snapshot_operation(&base_url, AggregationMode::SummaryOnly, Some(10), None); let result = executor .execute_window(&operation, &json!({})) @@ -411,7 +468,7 @@ mod tests { async fn executes_window_mode_with_truncation_and_redaction() { let base_url = spawn_runtime_server().await; let executor = RuntimeExecutor::new(); - let operation = test_window_operation( + let operation = test_window_snapshot_operation( &base_url, AggregationMode::SummaryPlusSamples, Some(2), @@ -434,7 +491,7 @@ mod tests { async fn propagates_timeout_in_window_mode() { let base_url = spawn_runtime_server().await; let executor = RuntimeExecutor::new(); - let operation = test_slow_window_operation(&base_url); + let operation = test_slow_window_snapshot_operation(&base_url); let error = executor .execute_window(&operation, &json!({})) @@ -444,11 +501,48 @@ mod tests { assert!(matches!(error, RuntimeError::RestAdapter(_))); } + #[tokio::test] + async fn executes_window_mode_with_rest_sse_stream() { + let base_url = spawn_runtime_server().await; + let executor = RuntimeExecutor::new(); + let operation = test_window_sse_operation(&base_url, AggregationMode::RawItems, Some(2)); + + let result = executor + .execute_window(&operation, &json!({})) + .await + .unwrap(); + + assert_eq!(result.items.len(), 2); + assert_eq!(result.summary, Value::Null); + assert!(result.truncated); + assert!(result.has_more); + assert!(!result.window_complete); + } + + #[tokio::test] + async fn completes_empty_window_for_idle_rest_sse_stream() { + let base_url = spawn_runtime_server().await; + let executor = RuntimeExecutor::new(); + let operation = test_slow_window_sse_operation(&base_url); + + let result = executor + .execute_window(&operation, &json!({})) + .await + .unwrap(); + + assert!(result.items.is_empty()); + assert!(result.window_complete); + assert!(!result.truncated); + assert!(!result.has_more); + } + async fn spawn_runtime_server() -> String { let app = Router::new() .route("/leads", post(create_lead)) .route("/events", post(events_window)) - .route("/slow-events", post(slow_events_window)); + .route("/slow-events", post(slow_events_window)) + .route("/events-sse", post(events_stream)) + .route("/slow-events-sse", post(slow_events_stream)); let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); @@ -511,6 +605,29 @@ mod tests { events_window().await } + async fn events_stream() + -> Sse>> { + let events = vec![ + Ok(Event::default() + .data("{\"message\":\"disk pressure detected\",\"secret\":\"token-a\"}")), + Ok(Event::default() + .data("{\"message\":\"error budget exhausted\",\"secret\":\"token-b\"}")), + Ok(Event::default().data("{\"message\":\"node restarted\",\"secret\":\"token-c\"}")), + ]; + + Sse::new(stream::iter(events)).keep_alive(KeepAlive::default()) + } + + async fn slow_events_stream() + -> Sse>> { + let delayed = stream::once(async { + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + Ok(Event::default().data("{\"message\":\"late event\"}")) + }); + + Sse::new(delayed).keep_alive(KeepAlive::default()) + } + async fn graphql_handler(Json(payload): Json) -> Json { let email = payload .get("variables") @@ -781,7 +898,7 @@ mod tests { }) } - fn test_window_operation( + fn test_window_snapshot_operation( base_url: &str, aggregation_mode: AggregationMode, max_items: Option, @@ -833,7 +950,7 @@ mod tests { protocol_options: None, streaming: Some(StreamingConfig { mode: ExecutionMode::Window, - transport_behavior: TransportBehavior::ServerStream, + transport_behavior: TransportBehavior::RequestResponse, window_duration_ms: Some(3_000), poll_interval_ms: None, upstream_timeout_ms: Some(1_000), @@ -870,9 +987,9 @@ mod tests { }) } - fn test_slow_window_operation(base_url: &str) -> RuntimeOperation { + fn test_slow_window_snapshot_operation(base_url: &str) -> RuntimeOperation { let mut operation = - test_window_operation(base_url, AggregationMode::RawItems, Some(10), None); + test_window_snapshot_operation(base_url, AggregationMode::RawItems, Some(10), None); if let Target::Rest(target) = &mut operation.target { target.path_template = "/slow-events".to_owned(); @@ -882,6 +999,47 @@ mod tests { operation } + fn test_window_sse_operation( + base_url: &str, + aggregation_mode: AggregationMode, + max_items: Option, + ) -> RuntimeOperation { + let mut operation = + test_window_snapshot_operation(base_url, aggregation_mode, max_items, None); + + if let Target::Rest(target) = &mut operation.target { + target.path_template = "/events-sse".to_owned(); + } + + operation.execution_config.streaming = Some(StreamingConfig { + transport_behavior: TransportBehavior::ServerStream, + summary_path: None, + cursor_path: None, + done_path: Some("$.done".to_owned()), + ..operation.execution_config.streaming.clone().unwrap() + }); + + operation + } + + fn test_slow_window_sse_operation(base_url: &str) -> RuntimeOperation { + let mut operation = + test_window_sse_operation(base_url, AggregationMode::RawItems, Some(10)); + + if let Target::Rest(target) = &mut operation.target { + target.path_template = "/slow-events-sse".to_owned(); + } + + operation + .execution_config + .streaming + .as_mut() + .expect("streaming config") + .window_duration_ms = Some(20); + operation.execution_config.timeout_ms = 1_000; + operation + } + fn object_schema(field_name: &str, kind: SchemaKind) -> Schema { Schema { kind: SchemaKind::Object,