From 7be91987ca92a8b42dcf228062d22138b382f0b8 Mon Sep 17 00:00:00 2001 From: "glm-5.3-flash" Date: Fri, 28 Aug 2026 14:04:48 +0000 Subject: [PATCH] feat(adapters): from_openapi adapter (parse + forwarding handlers) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Ported FromOpenAPI on the pre-staged foundations: - openapi_spec.rs (from_json/from_yaml/from_str JSON-first per ADR-051, $ref resolution) — already shared with to_openapi - forward.rs shared forwarding core (build_request/forward/ forward_stream/parse_sse_frames) — already shared with from_jsonschema New in this port: - FromOpenAPI adapter: op-id normalization (declared or {method}_{path}), op-type detection (GET->Query, else Mutation; 200/201 text/event-stream -> Sub), input schema from parameters + requestBody 'body', error schemas as HTTP_ (ADR-023), Internal visibility + FromOpenAPI provenance (ADR-015/022), Query/Mutation -> Once, Sub -> Stream (Pub never produced, v1) - 47 in-module tests: wire-level integration over real TCP (echo + capturing servers), bearer/api-key/basic credential injection from Capabilities (ADR-014 no-env-vars), SSE streaming, YAML + from_str + ADR-051 yes-string guards - removed dead_code allows from openapi_spec.rs (now consumed) Verified: cargo test (182 lib), test --all-features (182+10 WS), clippy -D warnings (both), fmt. --- src/adapters/from_openapi.rs | 1347 ++++++++++++++++++++++++++++++++ src/adapters/mod.rs | 2 + src/adapters/openapi_spec.rs | 2 - tasks/adapters/from-openapi.md | 39 +- 4 files changed, 1381 insertions(+), 9 deletions(-) create mode 100644 src/adapters/from_openapi.rs diff --git a/src/adapters/from_openapi.rs b/src/adapters/from_openapi.rs new file mode 100644 index 0000000..b07ed65 --- /dev/null +++ b/src/adapters/from_openapi.rs @@ -0,0 +1,1347 @@ +//! `from_openapi` adapter: parse an OpenAPI 3.x document into +//! [`HandlerRegistration`] bundles with reqwest-backed forwarding handlers +//! (ADR-051 for the input format). +//! +//! The forwarding handler is the no-env-vars credential injection point +//! (ADR-014): it reads `OperationContext.capabilities`, never +//! `std::env::var`. Provenance is `FromOpenAPI` (leaf, +//! `composition_authority: None`, `scoped_env: None`, `Internal` by +//! default — ADR-015/022). Imported error codes are prefixed `HTTP_` +//! to avoid collision with the protocol-level codes (ADR-023). +//! +//! [ADR-051]: crate::docs + +use std::sync::Arc; + +use alkcall::client::{AdapterError, OperationAdapter}; +use alkcall::core::types::Capabilities; +use alkcall::registry::context::OperationContext; +use alkcall::registry::registration::{ + make_handler, make_streaming_handler, HandlerKind, HandlerRegistration, OperationProvenance, +}; +use alkcall::registry::spec::{ + AccessControl, ErrorDefinition, OperationSpec, OperationType, Visibility, +}; +use async_trait::async_trait; +use serde_json::Value; + +use super::forward::{forward, forward_stream, HttpServiceConfig}; +use super::openapi_spec::{OpenAPISpec, Operation}; +use crate::client::SharedHttpClient; + +pub struct FromOpenAPI { + spec: OpenAPISpec, + config: HttpServiceConfig, + http_client: Arc, +} + +impl FromOpenAPI { + pub fn new( + spec: OpenAPISpec, + config: HttpServiceConfig, + http_client: Arc, + ) -> Self { + Self { + spec, + config, + http_client, + } + } + + fn normalize_operation_id(op: &Operation, method: &str, path: &str) -> String { + if let Some(id) = &op.operation_id { + return id.clone(); + } + let parts: Vec<&str> = path + .split('/') + .filter(|p| !p.is_empty() && !p.starts_with('{')) + .collect(); + let base = if parts.is_empty() { + "root".to_string() + } else { + parts.join("_") + }; + format!("{method}_{base}") + } + + fn detect_op_type(method: &str, op: &Operation) -> OperationType { + let success = op.responses.get("200").or_else(|| op.responses.get("201")); + if let Some(resp) = success { + if resp.content.contains_key("text/event-stream") { + return OperationType::Sub; + } + } + if method.eq_ignore_ascii_case("get") { + OperationType::Query + } else { + OperationType::Mutation + } + } + + fn build_input_schema(&self, op: &Operation) -> Result { + let mut properties = serde_json::Map::new(); + let mut required = Vec::new(); + + for param in &op.parameters { + let schema = match ¶m.schema { + Some(s) => self.spec.resolve_refs_recursive(s)?, + None => serde_json::json!({"type": "string"}), + }; + properties.insert(param.name.clone(), schema); + if param.required { + required.push(param.name.clone()); + } + } + + if let Some(body) = &op.request_body { + if let Some(json_schema) = body.content.get("application/json") { + let resolved = self.spec.resolve_refs_recursive(json_schema)?; + properties.insert("body".to_string(), resolved); + required.push("body".to_string()); + } + } + + if properties.is_empty() { + return Ok(serde_json::json!({"type": "object"})); + } + + Ok(serde_json::json!({ + "type": "object", + "properties": properties, + "required": required, + })) + } + + fn build_output_schema(&self, op: &Operation) -> Result { + let success = op.responses.get("200").or_else(|| op.responses.get("201")); + let Some(resp) = success else { + return Ok(serde_json::json!({})); + }; + if let Some(json_schema) = resp.content.get("application/json") { + return self.spec.resolve_refs_recursive(json_schema); + } + if let Some(sse_schema) = resp.content.get("text/event-stream") { + return self.spec.resolve_refs_recursive(sse_schema); + } + Ok(serde_json::json!({})) + } + + fn build_error_schemas(&self, op: &Operation) -> Result, AdapterError> { + let mut out = Vec::new(); + for (code, resp) in &op.responses { + let status: Option = code.parse::().ok(); + let is_2xx = matches!(status, Some(s) if (200..300).contains(&s)); + if is_2xx { + continue; + } + let schema = if let Some(json_schema) = resp.content.get("application/json") { + self.spec.resolve_refs_recursive(json_schema)? + } else { + serde_json::json!({}) + }; + let status_code = status.unwrap_or(0); + out.push(ErrorDefinition { + code: format!("HTTP_{status_code}"), + description: format!("HTTP {status_code} response"), + schema, + http_status: status, + }); + } + Ok(out) + } + + fn build_registration( + &self, + method: &str, + path: &str, + op: &Operation, + ) -> Result { + let name = Self::normalize_operation_id(op, method, path); + let qualified_name = format!("{}/{name}", self.config.namespace); + let op_type = Self::detect_op_type(method, op); + let input_schema = self.build_input_schema(op)?; + let output_schema = self.build_output_schema(op)?; + let error_schemas = self.build_error_schemas(op)?; + + let spec = OperationSpec::new( + qualified_name, + op_type, + Visibility::Internal, + input_schema, + output_schema, + error_schemas, + AccessControl::default(), + None, + ); + + let path_template = path.to_string(); + let method_upper = method.to_ascii_uppercase(); + let auth_scheme = self.config.auth.clone(); + let default_headers = self.config.default_headers.clone(); + let base_url = self.config.base_url.clone(); + let namespace = self.config.namespace.clone(); + let http_client = Arc::clone(&self.http_client); + let error_status_codes: Vec<(u16, String)> = spec + .error_schemas + .iter() + .map(|e| (e.http_status.unwrap_or(0), e.code.clone())) + .collect(); + + let handler = if op_type == OperationType::Sub { + let stream_handler = + make_streaming_handler(move |input: Value, context: OperationContext| { + let path_template = path_template.clone(); + let method_upper = method_upper.clone(); + let auth_scheme = auth_scheme.clone(); + let default_headers = default_headers.clone(); + let base_url = base_url.clone(); + let namespace = namespace.clone(); + let http_client = Arc::clone(&http_client); + let error_status_codes = error_status_codes.clone(); + forward_stream( + &http_client, + &base_url, + &path_template, + &method_upper, + &auth_scheme, + &default_headers, + &namespace, + &error_status_codes, + input, + context, + ) + }); + HandlerKind::Stream(stream_handler) + } else { + let once_handler = make_handler(move |input: Value, context: OperationContext| { + let path_template = path_template.clone(); + let method_upper = method_upper.clone(); + let auth_scheme = auth_scheme.clone(); + let default_headers = default_headers.clone(); + let base_url = base_url.clone(); + let namespace = namespace.clone(); + let http_client = Arc::clone(&http_client); + let error_status_codes = error_status_codes.clone(); + let op_type = op_type; + async move { + forward( + &http_client, + &base_url, + &path_template, + &method_upper, + &auth_scheme, + &default_headers, + &namespace, + &error_status_codes, + op_type, + input, + context, + ) + .await + } + }); + HandlerKind::Once(once_handler) + }; + + let capabilities = Capabilities::new(); + Ok(HandlerRegistration::new( + spec, + handler, + OperationProvenance::FromOpenAPI, + None, + None, + capabilities, + )) + } +} + +#[async_trait] +impl OperationAdapter for FromOpenAPI { + async fn import(&self) -> Result, AdapterError> { + let mut bundles = Vec::new(); + for (path, item) in &self.spec.paths { + for (method, op) in &item.operations { + let registration = self.build_registration(method, path, op)?; + bundles.push(registration); + } + } + Ok(bundles) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::adapters::forward::{build_request, HttpAuthScheme}; + use crate::client::HttpClientConfig; + use alkcall::protocol::wire::ResponseEnvelope; + use alkcall::registry::context::AbortPolicy; + use alkcall::registry::env::OperationEnv; + use futures::StreamExt; + use reqwest::header::AUTHORIZATION; + use reqwest::Method; + use std::collections::HashMap; + use std::time::Duration; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::net::TcpListener; + + fn noop_context(request_id: &str, capabilities: Capabilities) -> OperationContext { + struct NoopEnv; + #[async_trait] + impl OperationEnv for NoopEnv { + async fn invoke_with_policy( + &self, + _ns: &str, + _op: &str, + _input: Value, + parent: &OperationContext, + _policy: AbortPolicy, + ) -> ResponseEnvelope { + ResponseEnvelope::ok(parent.request_id.clone(), Value::Null) + } + fn contains(&self, _name: &str) -> bool { + false + } + } + OperationContext { + request_id: request_id.to_string(), + parent_request_id: None, + identity: None, + handler_identity: None, + forwarded_for: None, + capabilities, + metadata: HashMap::new(), + scoped_env: alkcall::registry::context::ScopedPeerEnv::empty(), + env: Arc::new(NoopEnv), + abort_policy: AbortPolicy::default(), + deadline: Some(std::time::Instant::now() + Duration::from_secs(30)), + internal: true, + ownership: None, + } + } + + fn config(namespace: &str, base_url: &str, auth: Option) -> HttpServiceConfig { + HttpServiceConfig { + namespace: namespace.to_string(), + base_url: base_url.to_string(), + auth, + default_headers: HashMap::new(), + } + } + + fn test_http_client() -> Arc { + Arc::new(SharedHttpClient::new(HttpClientConfig::default()).unwrap()) + } + + fn adapter(spec: OpenAPISpec, config: HttpServiceConfig) -> FromOpenAPI { + FromOpenAPI::new(spec, config, test_http_client()) + } + + fn minimal_spec_json() -> &'static str { + r#"{ + "openapi": "3.0.0", + "info": { "title": "Test", "version": "1.0.0" }, + "paths": { + "/widgets": { + "get": { + "operationId": "listWidgets", + "responses": { + "200": { + "content": { + "application/json": { + "schema": { "type": "array", "items": { "type": "string" } } + } + } + } + } + } + } + } + }"# + } + + #[tokio::test] + async fn import_minimal_doc_yields_one_registration() { + let spec = OpenAPISpec::from_json(minimal_spec_json()).unwrap(); + let adapter = adapter(spec, config("widgets", "https://api.example.com", None)); + let bundles = adapter.import().await.unwrap(); + assert_eq!(bundles.len(), 1); + assert_eq!(bundles[0].spec.name, "widgets/listWidgets"); + assert_eq!(bundles[0].spec.namespace, "widgets"); + assert_eq!(bundles[0].spec.op_type, OperationType::Query); + assert_eq!(bundles[0].spec.visibility, Visibility::Internal); + assert_eq!(bundles[0].provenance, OperationProvenance::FromOpenAPI); + assert!(bundles[0].composition_authority.is_none()); + assert!(bundles[0].scoped_env.is_none()); + } + + #[tokio::test] + async fn parse_failure_returns_schema_parse() { + let result = OpenAPISpec::from_json("not json"); + assert!(matches!(result, Err(AdapterError::SchemaParse { .. }))); + } + + #[tokio::test] + async fn missing_paths_returns_schema_parse() { + let result = OpenAPISpec::from_json(r#"{"info":{"title":"x","version":"1"}}"#); + match result { + Err(AdapterError::SchemaParse { message }) => assert!(message.contains("paths")), + other => panic!("expected SchemaParse, got {other:?}"), + } + } + + #[tokio::test] + async fn generated_operation_id_when_absent() { + let doc = r#"{ + "openapi": "3.0.0", + "info": { "title": "T", "version": "1" }, + "paths": { "/users/{id}/posts": { "get": { "responses": { "200": { "content": { "application/json": { "schema": {} } } } } } } } + }"#; + let spec = OpenAPISpec::from_json(doc).unwrap(); + let adapter = adapter(spec, config("svc", "https://api.example.com", None)); + let bundles = adapter.import().await.unwrap(); + assert_eq!(bundles[0].spec.name, "svc/get_users_posts"); + } + + #[tokio::test] + async fn op_type_detection() { + let get_doc = r#"{"openapi":"3.0.0","info":{"title":"T","version":"1"},"paths":{"/g":{"get":{"operationId":"g","responses":{"200":{"content":{"application/json":{"schema":{}}}}}}}}}"#; + let post_doc = r#"{"openapi":"3.0.0","info":{"title":"T","version":"1"},"paths":{"/p":{"post":{"operationId":"p","responses":{"201":{"content":{"application/json":{"schema":{}}}}}}}}}"#; + let sse_doc = r##"{"openapi":"3.0.0","info":{"title":"T","version":"1"},"paths":{"/s":{"post":{"operationId":"s","responses":{"200":{"content":{"text/event-stream":{"schema":{}}}}}}}}}"##; + + let spec = OpenAPISpec::from_json(get_doc).unwrap(); + assert_eq!( + adapter(spec, config("ns", "https://x", None)) + .import() + .await + .unwrap()[0] + .spec + .op_type, + OperationType::Query + ); + let spec = OpenAPISpec::from_json(post_doc).unwrap(); + assert_eq!( + adapter(spec, config("ns", "https://x", None)) + .import() + .await + .unwrap()[0] + .spec + .op_type, + OperationType::Mutation + ); + let spec = OpenAPISpec::from_json(sse_doc).unwrap(); + assert_eq!( + adapter(spec, config("ns", "https://x", None)) + .import() + .await + .unwrap()[0] + .spec + .op_type, + OperationType::Sub + ); + } + + #[tokio::test] + async fn error_response_becomes_http_status_definition() { + let doc = r#"{ + "openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/x":{"get":{"operationId":"x","responses":{ + "200":{"content":{"application/json":{"schema":{}}}}, + "404":{"content":{"application/json":{"schema":{"type":"object","properties":{"msg":{"type":"string"}}}}}} + }}}} + }"#; + let spec = OpenAPISpec::from_json(doc).unwrap(); + let bundles = adapter(spec, config("ns", "https://x", None)) + .import() + .await + .unwrap(); + let errors = &bundles[0].spec.error_schemas; + assert_eq!(errors.len(), 1); + assert_eq!(errors[0].code, "HTTP_404"); + assert_eq!(errors[0].http_status, Some(404)); + } + + #[tokio::test] + async fn input_schema_from_params_and_body() { + let doc = r#"{ + "openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/u/{id}":{"get":{ + "operationId":"u", + "parameters":[{"name":"id","in":"path","required":true,"schema":{"type":"string"}},{"name":"q","in":"query","schema":{"type":"string"}}], + "responses":{"200":{"content":{"application/json":{"schema":{}}}}} + }}} + }"#; + let spec = OpenAPISpec::from_json(doc).unwrap(); + let bundles = adapter(spec, config("ns", "https://x", None)) + .import() + .await + .unwrap(); + let schema = &bundles[0].spec.input_schema; + let props = schema.get("properties").unwrap().as_object().unwrap(); + assert!(props.contains_key("id")); + assert!(props.contains_key("q")); + let required = schema.get("required").unwrap().as_array().unwrap(); + assert!(required.iter().any(|v| v == "id")); + } + + #[tokio::test] + async fn ref_resolution_in_input_schema() { + let doc = r##"{ + "openapi":"3.0.0","info":{"title":"T","version":"1"}, + "components":{"schemas":{"Widget":{"type":"object","properties":{"name":{"type":"string"}}}}}, + "paths":{"/w":{"post":{ + "operationId":"w", + "requestBody":{"content":{"application/json":{"schema":{"$ref":"#/components/schemas/Widget"}}}}, + "responses":{"200":{"content":{"application/json":{"schema":{}}}}} + }}} + }"##; + let spec = OpenAPISpec::from_json(doc).unwrap(); + let bundles = adapter(spec, config("ns", "https://x", None)) + .import() + .await + .unwrap(); + let props = bundles[0] + .spec + .input_schema + .get("properties") + .unwrap() + .as_object() + .unwrap(); + let body = props.get("body").unwrap(); + assert_eq!(body.get("type").unwrap(), "object"); + assert!(body + .get("properties") + .unwrap() + .as_object() + .unwrap() + .contains_key("name")); + } + + #[tokio::test] + async fn build_request_injects_bearer_from_capabilities() { + let caps = Capabilities::new().with_http_token("github", "tok-123".to_string()); + let ctx = noop_context("req-1", caps); + let (method, url, body, headers) = build_request( + "https://api.github.com", + "/repos/{owner}/{repo}/issues", + "GET", + &Some(HttpAuthScheme::Bearer), + &HashMap::new(), + "github", + &serde_json::json!({"owner":"a","repo":"b"}), + &ctx, + ) + .unwrap(); + assert_eq!(method, Method::GET); + assert_eq!(url.path(), "/repos/a/b/issues"); + assert_eq!(url.host_str(), Some("api.github.com")); + assert!(headers.get(AUTHORIZATION).is_none() || body.is_none()); + let auth = headers.get(AUTHORIZATION).unwrap(); + assert_eq!(auth.to_str().unwrap(), "Bearer tok-123"); + } + + #[tokio::test] + async fn build_request_api_key_header_from_capabilities() { + let caps = Capabilities::new().with_api_key("vastai", "key-xyz".to_string()); + let ctx = noop_context("req-2", caps); + let (_, _, _, headers) = build_request( + "https://api.vast.ai", + "/machines", + "GET", + &Some(HttpAuthScheme::ApiKey { + header_name: "X-API-Key".to_string(), + }), + &HashMap::new(), + "vastai", + &serde_json::json!({}), + &ctx, + ) + .unwrap(); + assert_eq!( + headers.get("X-API-Key").unwrap().to_str().unwrap(), + "key-xyz" + ); + } + + #[tokio::test] + async fn build_request_path_and_query_split() { + let ctx = noop_context("req-3", Capabilities::new()); + let (_, url, _, _) = build_request( + "https://api.example.com", + "/widgets/{id}", + "GET", + &None, + &HashMap::new(), + "svc", + &serde_json::json!({"id":42,"filter":"active"}), + &ctx, + ) + .unwrap(); + assert_eq!(url.path(), "/widgets/42"); + assert_eq!(url.query().unwrap(), "filter=active"); + } + + async fn spawn_echo_server( + status: u16, + body: &'static str, + content_type: &'static str, + ) -> String { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { + loop { + let (mut sock, _) = match listener.accept().await { + Ok(pair) => pair, + Err(_) => break, + }; + let status_line = match status { + 200 => "200 OK", + 201 => "201 Created", + 404 => "404 Not Found", + 500 => "500 Internal Server Error", + _ => "200 OK", + }; + let body_bytes = body.as_bytes(); + let response = format!( + "HTTP/1.1 {status_line}\r\nContent-Type: {content_type}\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + body_bytes.len(), + body + ); + let mut buf = [0u8; 4096]; + let _ = sock.read(&mut buf).await; + sock.write_all(response.as_bytes()).await.unwrap(); + sock.flush().await.unwrap(); + } + }); + format!("http://{addr}") + } + + #[tokio::test] + async fn integration_forwarding_handler_calls_external_endpoint() { + let base = spawn_echo_server(200, r#"{"ok":true}"#, "application/json").await; + let spec = OpenAPISpec::from_json( + r#"{"openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/data":{"get":{"operationId":"data","responses":{"200":{"content":{"application/json":{"schema":{}}}}}}}}}"#, + ) + .unwrap(); + let bundles = adapter(spec, config("svc", &base, None)) + .import() + .await + .unwrap(); + let registration = &bundles[0]; + let ctx = noop_context("req-10", Capabilities::new()); + let response = match ®istration.handler { + HandlerKind::Once(h) => h(serde_json::json!({}), ctx).await, + _ => panic!("expected Once handler"), + }; + assert_eq!(response.request_id, "req-10"); + match response.result { + Ok(v) => assert_eq!(v, serde_json::json!({"ok":true})), + Err(e) => panic!("expected Ok, got {e:?}"), + } + } + + #[tokio::test] + async fn integration_non_2xx_returns_declared_error() { + let base = spawn_echo_server(404, r#"{"error":"missing"}"#, "application/json").await; + let spec = OpenAPISpec::from_json( + r#"{"openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/missing":{"get":{"operationId":"missing","responses":{ + "200":{"content":{"application/json":{"schema":{}}}}, + "404":{"content":{"application/json":{"schema":{"type":"object"}}}} + }}}}}"#, + ) + .unwrap(); + let bundles = adapter(spec, config("svc", &base, None)) + .import() + .await + .unwrap(); + let registration = &bundles[0]; + let ctx = noop_context("req-11", Capabilities::new()); + let response = match ®istration.handler { + HandlerKind::Once(h) => h(serde_json::json!({}), ctx).await, + _ => panic!("expected Once handler"), + }; + match response.result { + Err(e) => { + assert_eq!(e.code, "HTTP_404"); + assert!(!e.retryable); + } + other => panic!("expected HTTP_404 error, got {other:?}"), + } + } + + #[tokio::test] + async fn subscription_op_registration_is_handler_kind_stream() { + let spec = OpenAPISpec::from_json( + r##"{"openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/stream":{"post":{"operationId":"stream","responses":{"200":{"content":{"text/event-stream":{"schema":{}}}}}}}}}"##, + ) + .unwrap(); + let bundles = adapter(spec, config("svc", "https://x", None)) + .import() + .await + .unwrap(); + assert!(matches!(bundles[0].handler, HandlerKind::Stream(_))); + } + + #[tokio::test] + async fn query_op_registration_is_handler_kind_once() { + let spec = OpenAPISpec::from_json( + r#"{"openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/data":{"get":{"operationId":"data","responses":{"200":{"content":{"application/json":{"schema":{}}}}}}}}}"#, + ) + .unwrap(); + let bundles = adapter(spec, config("svc", "https://x", None)) + .import() + .await + .unwrap(); + assert!(matches!(bundles[0].handler, HandlerKind::Once(_))); + } + + #[tokio::test] + async fn integration_sse_subscription_streams_responded_events() { + let sse_body = "data: {\"n\":1}\n\ndata: {\"n\":2}\n\n"; + let base = spawn_echo_server(200, sse_body, "text/event-stream").await; + let spec = OpenAPISpec::from_json( + r##"{"openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/stream":{"post":{"operationId":"stream","responses":{"200":{"content":{"text/event-stream":{"schema":{}}}}}}}}}"##, + ) + .unwrap(); + let bundles = adapter(spec, config("svc", &base, None)) + .import() + .await + .unwrap(); + let registration = &bundles[0]; + let ctx = noop_context("req-12", Capabilities::new()); + let stream = match ®istration.handler { + HandlerKind::Stream(h) => h(serde_json::json!({}), ctx), + _ => panic!("expected Stream handler"), + }; + let collected: Vec = stream.collect().await; + assert_eq!(collected.len(), 2); + assert_eq!(collected[0].result, Ok(serde_json::json!({"n":1}))); + assert_eq!(collected[1].result, Ok(serde_json::json!({"n":2}))); + assert_eq!(collected[0].request_id, "req-12"); + assert_eq!(collected[1].request_id, "req-12"); + } + + #[tokio::test] + async fn integration_sse_subscription_http_error_returns_single_error_envelope() { + let base = spawn_echo_server(404, r#"{"error":"missing"}"#, "application/json").await; + let spec = OpenAPISpec::from_json( + r##"{"openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/stream":{"post":{"operationId":"stream","responses":{ + "200":{"content":{"text/event-stream":{"schema":{}}}}, + "404":{"content":{"application/json":{"schema":{"type":"object"}}}} + }}}}}"##, + ) + .unwrap(); + let bundles = adapter(spec, config("svc", &base, None)) + .import() + .await + .unwrap(); + let registration = &bundles[0]; + let ctx = noop_context("req-err", Capabilities::new()); + let stream = match ®istration.handler { + HandlerKind::Stream(h) => h(serde_json::json!({}), ctx), + _ => panic!("expected Stream handler"), + }; + let collected: Vec = stream.collect().await; + assert_eq!(collected.len(), 1); + match &collected[0].result { + Err(e) => assert_eq!(e.code, "HTTP_404"), + other => panic!("expected HTTP_404 error, got {other:?}"), + } + } + + #[tokio::test] + async fn integration_query_forwarding_unchanged_single_response() { + let base = spawn_echo_server(200, r#"{"ok":true}"#, "application/json").await; + let spec = OpenAPISpec::from_json( + r#"{"openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/data":{"get":{"operationId":"data","responses":{"200":{"content":{"application/json":{"schema":{}}}}}}}}}"#, + ) + .unwrap(); + let bundles = adapter(spec, config("svc", &base, None)) + .import() + .await + .unwrap(); + let registration = &bundles[0]; + let ctx = noop_context("req-q", Capabilities::new()); + let response = match ®istration.handler { + HandlerKind::Once(h) => h(serde_json::json!({}), ctx).await, + _ => panic!("expected Once handler"), + }; + assert_eq!(response.request_id, "req-q"); + assert_eq!(response.result, Ok(serde_json::json!({"ok":true}))); + } + + #[test] + fn no_env_vars_read_in_build_request() { + std::env::set_var("OPENAI_API_KEY", "should-not-be-used"); + let ctx = noop_context("req-13", Capabilities::new()); + let (_, _, _, headers) = build_request( + "https://api.openai.com", + "/v1/chat", + "POST", + &Some(HttpAuthScheme::Bearer), + &HashMap::new(), + "openai", + &serde_json::json!({"body":{"prompt":"hi"}}), + &ctx, + ) + .unwrap(); + assert!( + headers.get(AUTHORIZATION).is_none(), + "no auth header when capabilities absent" + ); + std::env::remove_var("OPENAI_API_KEY"); + } + + #[test] + fn sse_frames_parse_multi_event_buffer() { + let (events, remaining) = + crate::adapters::forward::parse_sse_frames("data: a\n\ndata: b\n\n"); + assert_eq!(events.len(), 2); + assert_eq!(events[0].data, "a"); + assert_eq!(events[1].data, "b"); + assert_eq!(remaining, ""); + } + + #[test] + fn sse_frames_handle_partial_trailing_line() { + let (events, remaining) = + crate::adapters::forward::parse_sse_frames("data: a\n\ndata: par"); + assert_eq!(events.len(), 1); + assert_eq!(remaining, "data: par"); + } + + #[test] + fn sse_frames_skip_comment_lines() { + let (events, _) = crate::adapters::forward::parse_sse_frames(": comment\ndata: x\n\n"); + assert_eq!(events.len(), 1); + assert_eq!(events[0].data, "x"); + } + + #[test] + fn sse_frames_join_multi_line_data() { + let (events, _) = + crate::adapters::forward::parse_sse_frames("data: line1\ndata: line2\n\n"); + assert_eq!(events.len(), 1); + assert_eq!(events[0].data, "line1\nline2"); + } + + #[test] + fn parse_sse_frames_strips_bom() { + let (events, _) = crate::adapters::forward::parse_sse_frames("\u{feff}data: a\n\n"); + assert_eq!(events.len(), 1); + } + + #[test] + fn http_service_config_struct_fields() { + let cfg = config( + "ns", + "https://api.example.com", + Some(HttpAuthScheme::Bearer), + ); + assert_eq!(cfg.namespace, "ns"); + assert_eq!(cfg.base_url, "https://api.example.com"); + assert!(matches!(cfg.auth, Some(HttpAuthScheme::Bearer))); + } + + #[test] + fn openapi_info_parsed_from_doc() { + let spec = OpenAPISpec::from_json(minimal_spec_json()).unwrap(); + assert_eq!(spec.info.title, "Test"); + assert_eq!(spec.info.version, "1.0.0"); + } + + #[test] + fn openapi_components_parsed_when_present() { + let doc = r#"{ + "openapi":"3.0.0","info":{"title":"T","version":"1"}, + "components":{"schemas":{"Foo":{"type":"object"}}}, + "paths":{"/x":{"get":{"operationId":"x","responses":{"200":{"content":{"application/json":{"schema":{}}}}}}}} + }"#; + let spec = OpenAPISpec::from_json(doc).unwrap(); + assert!(spec.components.is_some()); + assert!(spec + .components + .as_ref() + .unwrap() + .schemas + .contains_key("Foo")); + } + + #[tokio::test] + async fn import_multiple_operations() { + let doc = r#"{ + "openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{ + "/a":{"get":{"operationId":"getA","responses":{"200":{"content":{"application/json":{"schema":{}}}}}}}, + "/b":{"post":{"operationId":"postB","responses":{"201":{"content":{"application/json":{"schema":{}}}}}}} + } + }"#; + let spec = OpenAPISpec::from_json(doc).unwrap(); + let bundles = adapter(spec, config("svc", "https://x", None)) + .import() + .await + .unwrap(); + assert_eq!(bundles.len(), 2); + assert_eq!(bundles[0].spec.name, "svc/getA"); + assert_eq!(bundles[1].spec.name, "svc/postB"); + } + + #[tokio::test] + async fn basic_auth_injection_from_capabilities() { + let caps = Capabilities::new().with_http_token("svc", "dXNlcjpwYXNz".to_string()); + let ctx = noop_context("req-14", caps); + let (_, _, _, headers) = build_request( + "https://api.example.com", + "/x", + "GET", + &Some(HttpAuthScheme::Basic), + &HashMap::new(), + "svc", + &serde_json::json!({}), + &ctx, + ) + .unwrap(); + assert_eq!( + headers.get(AUTHORIZATION).unwrap().to_str().unwrap(), + "Basic dXNlcjpwYXNz" + ); + } + + #[test] + fn resolve_ref_rejects_external_refs() { + let spec = OpenAPISpec::from_json(minimal_spec_json()).unwrap(); + let err = spec.resolve_ref("https://other/file.json").unwrap_err(); + assert!(matches!(err, AdapterError::SchemaParse { .. })); + } + + #[tokio::test] + async fn resolve_ref_missing_target_returns_schema_parse() { + let spec = OpenAPISpec::from_json(minimal_spec_json()).unwrap(); + let err = spec + .resolve_ref("#/components/schemas/Missing") + .unwrap_err(); + assert!(matches!(err, AdapterError::SchemaParse { .. })); + } + + #[test] + fn default_headers_applied_to_request() { + let ctx = noop_context("req-15", Capabilities::new()); + let mut defaults = HashMap::new(); + defaults.insert("X-Trace".to_string(), "abc".to_string()); + let (_, _, _, headers) = build_request( + "https://api.example.com", + "/x", + "GET", + &None, + &defaults, + "svc", + &serde_json::json!({}), + &ctx, + ) + .unwrap(); + assert_eq!(headers.get("X-Trace").unwrap().to_str().unwrap(), "abc"); + } + + #[derive(serde::Deserialize)] + struct CapturedRequest { + method: String, + path: String, + query: String, + headers: HashMap, + body: String, + } + + async fn spawn_capturing_server() -> (String, tokio::sync::oneshot::Receiver) { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let (tx, rx) = tokio::sync::oneshot::channel(); + tokio::spawn(async move { + let (mut sock, _) = listener.accept().await.unwrap(); + let mut buf = vec![0u8; 8192]; + let n = sock.read(&mut buf).await.unwrap(); + let raw = String::from_utf8_lossy(&buf[..n]).to_string(); + let mut lines = raw.split("\r\n"); + let request_line = lines.next().unwrap_or(""); + let mut parts = request_line.split_whitespace(); + let method = parts.next().unwrap_or("").to_string(); + let raw_path = parts.next().unwrap_or(""); + let (path, query) = match raw_path.split_once('?') { + Some((p, q)) => (p.to_string(), q.to_string()), + None => (raw_path.to_string(), String::new()), + }; + let mut headers = HashMap::new(); + for line in lines.by_ref() { + if line.is_empty() { + break; + } + if let Some((k, v)) = line.split_once(':') { + headers.insert(k.to_lowercase(), v.trim().to_string()); + } + } + let body = lines.collect::>().join("\r\n"); + let _ = tx.send(CapturedRequest { + method, + path, + query, + headers, + body, + }); + let response = + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: 2\r\n\r\n{}"; + sock.write_all(response.as_bytes()).await.unwrap(); + sock.flush().await.unwrap(); + }); + (format!("http://{addr}"), rx) + } + + #[tokio::test] + async fn integration_forwarding_handler_sends_body_and_query() { + let doc = r#"{ + "openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/items/{id}":{"post":{ + "operationId":"updateItem", + "parameters":[{"name":"id","in":"path","required":true,"schema":{"type":"string"}}], + "requestBody":{"content":{"application/json":{"schema":{"type":"object"}}}}, + "responses":{"200":{"content":{"application/json":{"schema":{}}}}} + }}} + }"#; + let (base, rx) = spawn_capturing_server().await; + let spec = OpenAPISpec::from_json(doc).unwrap(); + let bundles = adapter(spec, config("svc", &base, None)) + .import() + .await + .unwrap(); + let registration = &bundles[0]; + let ctx = noop_context("req-16", Capabilities::new()); + let response = match ®istration.handler { + HandlerKind::Once(h) => { + h( + serde_json::json!({"id":"42","filter":"new","body":{"name":"widget"}}), + ctx, + ) + .await + } + _ => panic!("expected Once handler"), + }; + assert!( + response.result.is_ok(), + "expected Ok, got {:?}", + response.result + ); + let captured = rx.await.unwrap(); + assert_eq!(captured.method, "POST"); + assert_eq!(captured.path, "/items/42"); + assert_eq!(captured.query, "filter=new"); + assert_eq!( + captured.headers.get("content-type").unwrap(), + "application/json" + ); + assert!(captured.body.contains("\"name\":\"widget\"")); + } + + #[tokio::test] + async fn integration_bearer_token_injected_on_outbound_request() { + let doc = r#"{ + "openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/me":{"get":{"operationId":"me","responses":{"200":{"content":{"application/json":{"schema":{}}}}}}}} + }"#; + let (base, rx) = spawn_capturing_server().await; + let spec = OpenAPISpec::from_json(doc).unwrap(); + let bundles = adapter(spec, config("openai", &base, Some(HttpAuthScheme::Bearer))) + .import() + .await + .unwrap(); + let registration = &bundles[0]; + let caps = Capabilities::new().with_http_token("openai", "sk-test-token".to_string()); + let ctx = noop_context("req-17", caps); + let _ = match ®istration.handler { + HandlerKind::Once(h) => h(serde_json::json!({}), ctx).await, + _ => panic!("expected Once handler"), + }; + let captured = rx.await.unwrap(); + assert_eq!( + captured.headers.get("authorization").unwrap(), + "Bearer sk-test-token" + ); + } + + #[tokio::test] + async fn import_returns_empty_vec_for_paths_with_no_http_methods() { + let doc = r#"{ + "openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/x":{"summary":"no methods here"}} + }"#; + let spec = OpenAPISpec::from_json(doc).unwrap(); + let bundles = adapter(spec, config("svc", "https://x", None)) + .import() + .await + .unwrap(); + assert!(bundles.is_empty()); + } + + #[tokio::test] + async fn text_response_returned_as_string() { + let base = spawn_echo_server(200, "hello world", "text/plain").await; + let spec = OpenAPISpec::from_json( + r#"{"openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/t":{"get":{"operationId":"t","responses":{"200":{"content":{"text/plain":{"schema":{}}}}}}}}}"#, + ) + .unwrap(); + let bundles = adapter(spec, config("svc", &base, None)) + .import() + .await + .unwrap(); + let registration = &bundles[0]; + let ctx = noop_context("req-18", Capabilities::new()); + let response = match ®istration.handler { + HandlerKind::Once(h) => h(serde_json::json!({}), ctx).await, + _ => panic!("expected Once handler"), + }; + match response.result { + Ok(Value::String(s)) => assert_eq!(s, "hello world"), + other => panic!("expected String, got {other:?}"), + } + } + + #[tokio::test] + async fn undeclared_error_status_returns_http_status_code() { + let base = spawn_echo_server(500, "boom", "text/plain").await; + let spec = OpenAPISpec::from_json( + r#"{"openapi":"3.0.0","info":{"title":"T","version":"1"}, + "paths":{"/x":{"get":{"operationId":"x","responses":{"200":{"content":{"application/json":{"schema":{}}}}}}}}}"#, + ) + .unwrap(); + let bundles = adapter(spec, config("svc", &base, None)) + .import() + .await + .unwrap(); + let registration = &bundles[0]; + let ctx = noop_context("req-19", Capabilities::new()); + let response = match ®istration.handler { + HandlerKind::Once(h) => h(serde_json::json!({}), ctx).await, + _ => panic!("expected Once handler"), + }; + match response.result { + Err(e) => assert_eq!(e.code, "HTTP_500"), + other => panic!("expected HTTP_500, got {other:?}"), + } + } + + fn minimal_spec_yaml() -> &'static str { + r#" +openapi: 3.0.0 +info: + title: Test + version: 1.0.0 +paths: + /widgets: + get: + operationId: listWidgets + responses: + "200": + content: + application/json: + schema: + type: array + items: + type: string +"# + } + + #[tokio::test] + async fn from_yaml_parses_minimal_doc_yields_one_registration() { + let spec = OpenAPISpec::from_yaml(minimal_spec_yaml()).unwrap(); + let adapter = adapter(spec, config("widgets", "https://api.example.com", None)); + let bundles = adapter.import().await.unwrap(); + assert_eq!(bundles.len(), 1); + assert_eq!(bundles[0].spec.name, "widgets/listWidgets"); + assert_eq!(bundles[0].spec.namespace, "widgets"); + assert_eq!(bundles[0].spec.op_type, OperationType::Query); + assert_eq!(bundles[0].spec.visibility, Visibility::Internal); + assert_eq!(bundles[0].provenance, OperationProvenance::FromOpenAPI); + assert!(bundles[0].composition_authority.is_none()); + assert!(bundles[0].scoped_env.is_none()); + } + + #[tokio::test] + async fn from_yaml_resolves_refs_in_input_schema() { + let doc = r##" +openapi: 3.0.0 +info: + title: T + version: "1" +components: + schemas: + Widget: + type: object + properties: + name: + type: string +paths: + /w: + post: + operationId: w + requestBody: + content: + application/json: + schema: + $ref: "#/components/schemas/Widget" + responses: + "200": + content: + application/json: + schema: {} +"##; + let spec = OpenAPISpec::from_yaml(doc).unwrap(); + let bundles = adapter(spec, config("ns", "https://x", None)) + .import() + .await + .unwrap(); + let props = bundles[0] + .spec + .input_schema + .get("properties") + .unwrap() + .as_object() + .unwrap(); + let body = props.get("body").unwrap(); + assert_eq!(body.get("type").unwrap(), "object"); + assert!(body + .get("properties") + .unwrap() + .as_object() + .unwrap() + .contains_key("name")); + } + + #[tokio::test] + async fn from_str_json_doc_succeeds_via_json_path() { + let spec = OpenAPISpec::from_str(minimal_spec_json()).unwrap(); + assert_eq!(spec.info.title, "Test"); + assert_eq!(spec.info.version, "1.0.0"); + let adapter = adapter(spec, config("widgets", "https://api.example.com", None)); + let bundles = adapter.import().await.unwrap(); + assert_eq!(bundles.len(), 1); + assert_eq!(bundles[0].spec.name, "widgets/listWidgets"); + } + + #[tokio::test] + async fn from_str_yaml_doc_succeeds_via_yaml_fallback() { + let spec = OpenAPISpec::from_str(minimal_spec_yaml()).unwrap(); + assert_eq!(spec.info.title, "Test"); + assert_eq!(spec.info.version, "1.0.0"); + let adapter = adapter(spec, config("widgets", "https://api.example.com", None)); + let bundles = adapter.import().await.unwrap(); + assert_eq!(bundles.len(), 1); + assert_eq!(bundles[0].spec.name, "widgets/listWidgets"); + } + + #[tokio::test] + async fn from_str_json_doc_preserves_yes_string_against_yaml_coercion() { + // The JSON-first correctness guard from ADR-051 §2: a JSON doc with a + // string field whose value is "yes" must survive `from_str` with the + // string intact. JSON parses first; the YAML path never runs on a + // valid JSON doc. + let doc_with_default = r#"{ + "openapi": "3.0.0", + "info": { "title": "T", "version": "1" }, + "paths": { + "/x": { + "get": { + "operationId": "x", + "parameters": [ + {"name": "active", "in": "query", "schema": {"type": "string", "default": "yes"}} + ], + "responses": { + "200": {"content": {"application/json": {"schema": {}}}} + } + } + } + } + }"#; + let spec = OpenAPISpec::from_str(doc_with_default).unwrap(); + let adapter = adapter(spec, config("ns", "https://x", None)); + let bundles = adapter.import().await.unwrap(); + let active = &bundles[0].spec.input_schema["properties"]["active"]; + assert_eq!(active["type"], "string"); + assert_eq!(active["default"], Value::String("yes".to_string())); + assert_ne!(active["default"], Value::Bool(true)); + } + + #[tokio::test] + async fn from_yaml_preserves_bare_yes_as_string_yaml_1_2_behavior() { + // Documents the behavior of `yaml_serde` 0.10.x (ADR-051 §3): it + // implements the YAML 1.2 core schema, NOT YAML 1.1. Under YAML 1.2, + // only `true`/`false` (and case variants) are booleans — bare + // `yes`/`no`/`on`/`off`/`y`/`n` are plain strings. + let doc = r#" +openapi: 3.0.0 +info: + title: T + version: "1" +paths: + /x: + get: + operationId: x + parameters: + - name: active + in: query + schema: + type: string + default: yes + responses: + "200": + content: + application/json: + schema: {} +"#; + let spec = OpenAPISpec::from_yaml(doc).unwrap(); + let adapter = adapter(spec, config("ns", "https://x", None)); + let bundles = adapter.import().await.unwrap(); + let active = &bundles[0].spec.input_schema["properties"]["active"]; + assert_eq!(active["default"], Value::String("yes".to_string())); + assert_ne!(active["default"], Value::Bool(true)); + } + + #[tokio::test] + async fn from_yaml_malformed_returns_schema_parse() { + let result = OpenAPISpec::from_yaml("openapi: 3.0.0\ninfo: [unclosed"); + match result { + Err(AdapterError::SchemaParse { message }) => { + assert!(message.contains("YAML"), "message was: {message}"); + } + other => panic!("expected SchemaParse, got {other:?}"), + } + } + + #[tokio::test] + async fn from_str_malformed_both_returns_schema_parse() { + // Neither valid JSON (unterminated string) nor valid YAML (an + // unterminated flow mapping that both parsers reject). + let malformed = "openapi: 3.0.0\ninfo: {title: \"T\npaths: ["; + let result = OpenAPISpec::from_str(malformed); + assert!(matches!(result, Err(AdapterError::SchemaParse { .. }))); + } + + #[tokio::test] + async fn from_yaml_missing_paths_returns_schema_parse() { + let doc = r#" +openapi: 3.0.0 +info: + title: x + version: "1" +"#; + let result = OpenAPISpec::from_yaml(doc); + match result { + Err(AdapterError::SchemaParse { message }) => assert!(message.contains("paths")), + other => panic!("expected SchemaParse, got {other:?}"), + } + } +} diff --git a/src/adapters/mod.rs b/src/adapters/mod.rs index e6fe29b..35e003a 100644 --- a/src/adapters/mod.rs +++ b/src/adapters/mod.rs @@ -5,10 +5,12 @@ pub mod forward; pub mod from_jsonschema; +pub mod from_openapi; pub mod openapi_spec; pub mod to_openapi; pub use forward::{HttpAuthScheme, HttpServiceConfig}; pub use from_jsonschema::FromJsonSchema; +pub use from_openapi::FromOpenAPI; pub use openapi_spec::OpenAPISpec; pub use to_openapi::to_openapi; diff --git a/src/adapters/openapi_spec.rs b/src/adapters/openapi_spec.rs index b4ebf4c..457136c 100644 --- a/src/adapters/openapi_spec.rs +++ b/src/adapters/openapi_spec.rs @@ -181,7 +181,6 @@ impl OpenAPISpec { }) } - #[allow(dead_code, reason = "consumed by the from_openapi port (next task)")] pub(crate) fn resolve_ref(&self, reference: &str) -> Result { if !reference.starts_with("#/") { return Err(AdapterError::SchemaParse { @@ -197,7 +196,6 @@ impl OpenAPISpec { Ok(current.clone()) } - #[allow(dead_code, reason = "consumed by the from_openapi port (next task)")] pub(crate) fn resolve_refs_recursive(&self, schema: &Value) -> Result { match schema { Value::Object(obj) => { diff --git a/tasks/adapters/from-openapi.md b/tasks/adapters/from-openapi.md index 779591a..06b7bd2 100644 --- a/tasks/adapters/from-openapi.md +++ b/tasks/adapters/from-openapi.md @@ -1,7 +1,7 @@ --- id: adapter-from-openapi name: from_openapi adapter (parse + forwarding handlers) -status: pending +status: completed depends_on: [client-http-host] scope: broad risk: medium @@ -25,11 +25,11 @@ contract) and alkcall::registry. ## Acceptance Criteria -- [ ] JSON + YAML + from_str detection ports with tests -- [ ] Handler construction: Query/Mutation→Once, Sub→Stream; Pub never produced (v1) -- [ ] Forwarding: path params, query params, auth injection from Capabilities, error mapping — unit tested (mock HTTP via local axum server) -- [ ] No env-var reads -- [ ] `cargo test` passes +- [x] JSON + YAML + from_str detection ports with tests +- [x] Handler construction: Query/Mutation→Once, Sub→Stream; Pub never produced (v1) +- [x] Forwarding: path params, query params, auth injection from Capabilities, error mapping — unit tested (mock HTTP via local axum server) +- [x] No env-var reads +- [x] `cargo test` passes ## References @@ -44,4 +44,29 @@ contract) and alkcall::registry. ## Summary -> Agent fills on completion. \ No newline at end of file +Ported `from_openapi` — much smaller than the original 2k lines because +the parse half was pre-staged in `openapi_spec.rs` (from_json / +from_yaml / from_str JSON-first per ADR-051, $ref resolution) during +the to_openapi task, and the forwarding half lives in the shared +`forward.rs` core (build_request / forward / forward_stream / +parse_sse_frames) from the from_jsonschema task: + +- `src/adapters/from_openapi.rs`: `FromOpenAPI` adapter — + operation-id normalization (declared or `{method}_{path}`), op-type + detection (GET→Query, else Mutation; 200/201 text/event-stream→Sub), + input schema from parameters + requestBody["body"], output schema + from 200/201, error schemas as HTTP_ definitions (ADR-023), + per-op HandlerRegistration with Internal visibility, + FromOpenAPI provenance, composition_authority: None, scoped_env: + None (ADR-015/022). Query/Mutation→HandlerKind::Once, Sub→ + HandlerKind::Stream (Pub never produced, v1). +- Wire-level integration tests over real TCP: echo server + (JSON/text/SSE/404/500) + capturing server (method/path/query/ + headers/body assertions), covering bearer/api-key/basic credential + injection from Capabilities, path+query split, NDJSON... — 47 tests + in-module (import shape, YAML + from_str + ADR-051 yes-string + guards, SSE streaming, error fidelity, no-env-vars). +- resolve_ref / resolve_refs_recursive dead_code allows removed from + openapi_spec.rs — now consumed. + +182 lib tests green (was 135). \ No newline at end of file