diff --git a/Cargo.lock b/Cargo.lock index a334665..0a53347 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1079,6 +1079,7 @@ dependencies = [ "serde_json", "thiserror", "url", + "yoke-derive", ] [[package]] diff --git a/README.md b/README.md index d2ef2b8..df91d9b 100644 --- a/README.md +++ b/README.md @@ -12,6 +12,12 @@ ODP separates Service discovery from catalog discovery. An Agent searches the ca for candidate Services, inspects each Service's live ODP document, and then navigates or searches that Service's Collections and Offerings. +`DirectoryClient::search` discovers indexed Services and submitted Collections. Use +`search_services` for a Service-only response or `collect_services` for bounded Service-only +aggregation. `suggest` returns mixed target names; `suggest_services` returns Service-only +keyword suggestions. See the [Directory guide](./crates/odp-directory/README.md) for result types, +the 100-result cap, and migration from the earlier method names. + ## Workspace | Goal | Crate | Guide | diff --git a/crates/odp-agent/README.md b/crates/odp-agent/README.md index 801e774..6d37f12 100644 --- a/crates/odp-agent/README.md +++ b/crates/odp-agent/README.md @@ -57,6 +57,11 @@ MiB. These are fixed SDK safety ceilings. Cross-document schema composition uses ## Search across Services +Federated discovery uses `DirectoryClient::collect_services` and remains Service-only. For mixed +discovery, use `DirectoryClient::search`, inspect each Collection result's owning Service, then +call `ServiceClient::get_collection` with its Collection ID. Collection results are not separate +Services. See the [Directory guide](../odp-directory/README.md#search-services-and-collections). + ```rust,no_run use odp_agent::{Agent, FederatedSearchRequest}; use odp_core::{OfferingSearchRequest, VERSION}; diff --git a/crates/odp-agent/src/agent.rs b/crates/odp-agent/src/agent.rs index 23179c7..4ab28d0 100644 --- a/crates/odp-agent/src/agent.rs +++ b/crates/odp-agent/src/agent.rs @@ -76,7 +76,7 @@ impl Agent { let concurrency = bounded(request.concurrency, 4, 16, "concurrency")?; let services = self .directory - .search_services( + .collect_services( &request.services, IterationOptions { max_items: maximum_services, diff --git a/crates/odp-core/Cargo.toml b/crates/odp-core/Cargo.toml index 7289137..d43883f 100644 --- a/crates/odp-core/Cargo.toml +++ b/crates/odp-core/Cargo.toml @@ -21,6 +21,8 @@ serde.workspace = true serde_json.workspace = true thiserror.workspace = true url.workspace = true +# yoke-derive 0.8.3 does not compile on Rust 1.85; constrain consumer resolution too. +yoke-derive = "=0.8.2" [lints] workspace = true diff --git a/crates/odp-directory/README.md b/crates/odp-directory/README.md index 672b47c..4c3561a 100644 --- a/crates/odp-directory/README.md +++ b/crates/odp-directory/README.md @@ -20,7 +20,7 @@ use odp_directory::{ # async fn main() -> Result<(), Box> { let directory = DirectoryClient::new(Environment::Production)?; let services = directory - .search_services( + .collect_services( &SearchRequest { filters: Some(ServiceFilters { payments: vec![PaymentFilter { @@ -45,17 +45,19 @@ for service in services { } let suggestions = directory - .suggest(&SuggestionRequest { + .suggest_services(&SuggestionRequest { limit: 5, prefix: "pla".to_owned(), + ..Default::default() }) .await?; # Ok(()) # } ``` -`search` returns one page. `continue_search` follows one opaque `next` reference. -`search_pages` and `search_services` perform bounded traversal for callers that want aggregation. +`search_services` returns one Service-only response. `continue_search_services` follows one opaque +`next` reference. `collect_services` performs bounded Service-only traversal, stopping at the +response or item limit without fetching another response. Search filters cover keywords, ODP operations, enrollment protocols, payment protocols, payment options, trust protocols, and the authentication requirements attached to operations and payments. Search responses can also carry facets for building data-driven filters without packaging the @@ -66,3 +68,75 @@ validated. See the [workspace guide](../../README.md) and the [ODP specification](https://www.offeringprotocol.org/). + +## Search Services and Collections + +```rust,no_run +use odp_directory::{DirectoryClient, DirectoryResult, Environment, ResourceSearchRequest}; + +# #[tokio::main(flavor = "current_thread")] +# async fn main() -> Result<(), Box> { +let directory = DirectoryClient::new(Environment::Production)?; +let response = directory.search(&ResourceSearchRequest { + query: "weather forecast".to_owned(), + limit: 25, + ..ResourceSearchRequest::default() +}).await?; +for item in response.items { + match item { + DirectoryResult::Service(item) => println!("Service: {}", item.service.name), + DirectoryResult::Collection(item) => println!("Collection: {} ({}, through {})", + item.collection.name, item.collection.id, item.service.service_origin), + DirectoryResult::Unknown { kind, .. } => println!("Unsupported result type: {kind}"), + } +} +for issue in response.issues { + eprintln!("Skipped result {}: {}", issue.index, issue.message); +} +# Ok(()) +# } +``` + +`ResourceSearchRequest.types` can restrict results to `ResultType::Service` or +`ResultType::Collection`. `None` selects both; explicit lists must be nonempty and distinct. +Filters apply to the owning Service. An empty query is omitted, allowing browsing. + +Collection identity is its owning Service origin plus its case-sensitive Collection ID. +The result's `indexed_at` describes the Collection's freshness; `service.indexed_at` describes +the parent's freshness. Both are timestamp strings. `service.service_id()` identifies the +Directory's Service record. Service results may include `available_through` platform attribution; +Collection attribution is the owning `service` itself. + +Inspect the owning Service's live ODP document, then use the Agent client's `get_collection` to +retrieve current details. Directory metadata is not authority to execute an Action or send +credentials. Unknown future result types retain their full raw JSON and are not interpreted as +Services. Malformed known results become indexed `issues` without discarding valid results. +Additional fields are retained in `additional` maps. + +Mixed search returns at most 100 results. `limit: 0` omits the limit, using the server's default +of 100. The server does not currently offer continuation: absent `next` does not mean every match +was returned. `continue_search` accepts an opaque same-origin continuation if one is supplied. +Each call returns one response. Facets count all matching targets, not just returned items; +a Service and two Collections count as three. Collection search does not depend on permission +to display its card on the Directory landing page. + +`suggest` sends POST `/v1/directory/suggestions`. `SuggestionRequest.filters` accepts the same +`ServiceFilters` as search, including AEP, keywords, ODP operations, payments and trust. +Collection filters apply to their owning Service. Matching spans names, descriptions and keywords, +but output contains deduplicated **names of matching Services and Collections**. Despite the +argument name `prefix`, matching uses substrings and whitespace-separated alternative terms. +`suggest_services` uses GET `/v1/services/suggestions` for Service-only keyword-prefix suggestions +and does not accept filters. +Both return strings from the server's `items` array; the default and maximum limit are 25. + +See the [canonical Directory example](../../examples/README.md#canonical-directory-discovery). + +## Migration + +- Service-only `search` calls become `search_services`; `continue_search` calls become + `continue_search_services`. +- Aggregating `search_services(request, options)` calls become `collect_services(request, options)`. +- `search_pages` is removed. To retain individual responses, call `search_services` followed by + `continue_search_services` with an explicit application limit. +- Service-only `suggest` calls become `suggest_services`. +- `search`, `continue_search`, and `suggest` select mixed discovery. diff --git a/crates/odp-directory/src/client.rs b/crates/odp-directory/src/client.rs index 0c8080e..f2ea851 100644 --- a/crates/odp-directory/src/client.rs +++ b/crates/odp-directory/src/client.rs @@ -6,8 +6,9 @@ use thiserror::Error; use url::Url; use crate::{ - DirectoryService, Environment, HttpRequest, HttpResponse, IterationOptions, SearchPage, - SearchRequest, SuggestionRequest, Transport, TransportError, default_transport, + DirectoryService, Environment, HttpRequest, HttpResponse, IterationOptions, + ResourceSearchRequest, SearchPage, SearchRequest, SearchResponse, SuggestionRequest, Transport, + TransportError, default_transport, }; const MAXIMUM_REDIRECTS: usize = 5; @@ -54,7 +55,35 @@ impl DirectoryClient { self.environment } - pub async fn search(&self, request: &SearchRequest) -> Result { + pub async fn search( + &self, + request: &ResourceSearchRequest, + ) -> Result { + validate_search(&request.query, request.limit, request.filters.as_ref())?; + if let Some(types) = &request.types { + if types.is_empty() || types.len() > 2 || (types.len() == 2 && types[0] == types[1]) { + return Err(DirectoryError::InvalidRequest( + "types must contain distinct service or collection values".to_owned(), + )); + } + } + let body = serde_json::to_vec(request) + .map_err(|error| DirectoryError::InvalidRequest(error.to_string()))?; + let target = self.continuation_url("/v1/directory/search")?; + let response = self.request("POST", target, body).await?; + crate::results::decode(&response.body) + } + + pub async fn continue_search(&self, next: &str) -> Result { + let target = self.continuation_url(next)?; + let response = self.request("GET", target, Vec::new()).await?; + crate::results::decode(&response.body) + } + + pub async fn search_services( + &self, + request: &SearchRequest, + ) -> Result { validate_search_request(request)?; let body = serde_json::to_vec(request) .map_err(|error| DirectoryError::InvalidRequest(error.to_string()))?; @@ -66,47 +95,54 @@ impl DirectoryClient { .await } - pub async fn continue_search(&self, next: &str) -> Result { + pub async fn continue_search_services(&self, next: &str) -> Result { let target = self.continuation_url(next)?; self.request_page("GET", target.as_str(), Vec::new()).await } - pub async fn search_pages( + pub async fn collect_services( &self, request: &SearchRequest, options: IterationOptions, - ) -> Result, DirectoryError> { - let maximum_pages = bounded(options.max_pages, 16, 16, "max_pages")?; - let mut pages = Vec::new(); - let mut page = self.search(request).await?; - for _ in 0..maximum_pages { - let next = page.next.clone(); - pages.push(page); - if next.is_empty() { - return Ok(pages); + ) -> Result, DirectoryError> { + let maximum_items = bounded(options.max_items, 10_000, 10_000, "max_items")?; + let maximum_responses = bounded(options.max_pages, 16, 16, "max_pages")?; + let mut services = Vec::new(); + let mut page = self.search_services(request).await?; + for index in 0..maximum_responses { + services.extend(page.items.into_iter().take(maximum_items - services.len())); + if page.next.is_empty() + || services.len() == maximum_items + || index + 1 == maximum_responses + { + break; } - page = self.continue_search(&next).await?; + page = self.continue_search_services(&page.next).await?; } - Ok(pages) + Ok(services) } - pub async fn search_services( + pub async fn suggest( &self, - request: &SearchRequest, - options: IterationOptions, - ) -> Result, DirectoryError> { - let maximum_items = bounded(options.max_items, 10_000, 10_000, "max_items")?; - let pages = self.search_pages(request, options).await?; - Ok(pages - .into_iter() - .flat_map(|page| page.items) - .take(maximum_items) - .collect()) + request: &SuggestionRequest, + ) -> Result, DirectoryError> { + self.suggestions("/v1/directory/suggestions", request, true) + .await } - pub async fn suggest( + pub async fn suggest_services( + &self, + request: &SuggestionRequest, + ) -> Result, DirectoryError> { + self.suggestions("/v1/services/suggestions", request, false) + .await + } + + async fn suggestions( &self, + path: &str, request: &SuggestionRequest, + mixed: bool, ) -> Result, DirectoryError> { let prefix = request.prefix.trim(); if prefix.is_empty() || prefix.chars().count() > 128 { @@ -119,20 +155,38 @@ impl DirectoryClient { "limit must be from 1 through 25".to_owned(), )); } - let mut target = Url::parse(&format!( - "{}/v1/services/suggestions", - self.environment.origin() - )) - .map_err(|error| DirectoryError::InvalidRequest(error.to_string()))?; - target.query_pairs_mut().append_pair("prefix", prefix); - if request.limit != 0 { - target - .query_pairs_mut() - .append_pair("limit", &request.limit.to_string()); + let mut target = Url::parse(&format!("{}{}", self.environment.origin(), path)) + .map_err(|error| DirectoryError::InvalidRequest(error.to_string()))?; + let response = if mixed { + validate_search("", 0, request.filters.as_ref())?; + let payload = SuggestionRequest { + prefix: prefix.to_owned(), + ..request.clone() + }; + let body = serde_json::to_vec(&payload) + .map_err(|error| DirectoryError::InvalidRequest(error.to_string()))?; + self.request("POST", target, body).await? + } else { + if request.filters.is_some() { + return Err(DirectoryError::InvalidRequest( + "Service-only suggestions do not support filters".to_owned(), + )); + } + target.query_pairs_mut().append_pair("prefix", prefix); + if request.limit != 0 { + target + .query_pairs_mut() + .append_pair("limit", &request.limit.to_string()); + } + self.request("GET", target, Vec::new()).await? + }; + #[derive(serde::Deserialize)] + struct Suggestions { + items: Vec, } - let response = self.request("GET", target, Vec::new()).await?; - let suggestions = serde_json::from_slice::>(&response.body) - .map_err(|error| DirectoryError::InvalidResponse(error.to_string()))?; + let suggestions = serde_json::from_slice::(&response.body) + .map_err(|error| DirectoryError::InvalidResponse(error.to_string()))? + .items; if suggestions.len() > 25 || suggestions.iter().any(|value| { value.trim() != value || value.is_empty() || value.chars().count() > 128 @@ -247,6 +301,11 @@ impl DirectoryClient { } fn continuation_url(&self, next: &str) -> Result { + if next.trim().is_empty() { + return Err(DirectoryError::InvalidRequest( + "Directory continuation is empty".to_owned(), + )); + } let origin = Url::parse(self.environment.origin()) .map_err(|error| DirectoryError::InvalidResponse(error.to_string()))?; let target = origin @@ -316,17 +375,25 @@ fn bounded( } fn validate_search_request(request: &SearchRequest) -> Result<(), DirectoryError> { - if request.limit > 100 { + validate_search(&request.query, request.limit, request.filters.as_ref()) +} + +fn validate_search( + query: &str, + limit: usize, + filters: Option<&crate::ServiceFilters>, +) -> Result<(), DirectoryError> { + if limit > 100 { return Err(DirectoryError::InvalidRequest( "limit must be from 1 through 100".to_owned(), )); } - if request.query.trim() != request.query || request.query.chars().count() > 512 { + if query.trim() != query || query.chars().count() > 512 { return Err(DirectoryError::InvalidRequest( "query must contain at most 512 characters without surrounding whitespace".to_owned(), )); } - if let Some(filters) = &request.filters { + if let Some(filters) = filters { if filters.keywords.len() > 32 || filters .keywords @@ -423,7 +490,7 @@ mod tests { }); let client = DirectoryClient::with_transport(Environment::Sandbox, transport.clone()); let page = client - .search(&SearchRequest { + .search_services(&SearchRequest { filters: Some(ServiceFilters { trust: vec![odp_core::TrustProtocol { name: odp_core::Protocol::Tap, @@ -462,7 +529,10 @@ mod tests { Arc::new(ResponseTransport(body.to_vec())), ); - let page = client.search(&SearchRequest::default()).await.unwrap(); + let page = client + .search_services(&SearchRequest::default()) + .await + .unwrap(); assert_eq!( page.facets.unwrap().trust, @@ -479,7 +549,12 @@ mod tests { Environment::Production, Arc::new(ResponseTransport(invalid.to_vec())), ); - assert!(client.search(&SearchRequest::default()).await.is_err()); + assert!( + client + .search_services(&SearchRequest::default()) + .await + .is_err() + ); } #[tokio::test] @@ -491,7 +566,7 @@ mod tests { }), ); let error = client - .search_services( + .collect_services( &SearchRequest::default(), IterationOptions { max_items: 10_001, @@ -512,7 +587,7 @@ mod tests { }), ); let error = client - .search(&SearchRequest { + .search_services(&SearchRequest { filters: Some(ServiceFilters { trust: vec![odp_core::TrustProtocol { name: odp_core::Protocol::Mpp, @@ -533,7 +608,10 @@ mod tests { Environment::Production, Arc::new(ResponseTransport(body.to_vec())), ); - let page = client.search(&SearchRequest::default()).await.unwrap(); + let page = client + .search_services(&SearchRequest::default()) + .await + .unwrap(); let protocols = page.items[0].protocols.as_ref().unwrap(); assert_eq!(protocols.payments.len(), 1); assert_eq!(protocols.trust.len(), 1); @@ -544,6 +622,11 @@ mod tests { Environment::Production, Arc::new(ResponseTransport(malformed.into_bytes())), ); - assert!(client.search(&SearchRequest::default()).await.is_err()); + assert!( + client + .search_services(&SearchRequest::default()) + .await + .is_err() + ); } } diff --git a/crates/odp-directory/src/lib.rs b/crates/odp-directory/src/lib.rs index d35ecbc..9aba635 100644 --- a/crates/odp-directory/src/lib.rs +++ b/crates/odp-directory/src/lib.rs @@ -2,6 +2,7 @@ mod client; mod models; +mod results; mod transport; pub use client::*; diff --git a/crates/odp-directory/src/models.rs b/crates/odp-directory/src/models.rs index 588b572..db06fc2 100644 --- a/crates/odp-directory/src/models.rs +++ b/crates/odp-directory/src/models.rs @@ -67,6 +67,83 @@ fn is_zero(value: &usize) -> bool { *value == 0 } +#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(rename_all = "lowercase")] +pub enum ResultType { + Service, + Collection, +} + +#[derive(Clone, Debug, Default, Deserialize, PartialEq, Serialize)] +pub struct ResourceSearchRequest { + #[serde(skip_serializing_if = "Option::is_none")] + pub filters: Option, + #[serde(default, skip_serializing_if = "is_zero")] + pub limit: usize, + #[serde(default, skip_serializing_if = "String::is_empty")] + pub query: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub types: Option>, +} + +#[derive(Clone, Debug, PartialEq)] +pub enum DirectoryResult { + Service(Box), + Collection(Box), + Unknown { kind: String, raw: Value }, +} + +#[derive(Clone, Debug, Deserialize, PartialEq)] +pub struct ServiceResult { + pub service: DirectoryService, + pub indexed_at: String, + pub available_through: Option, + #[serde(flatten)] + pub additional: AdditionalMembers, +} + +#[derive(Clone, Debug, Deserialize, PartialEq)] +pub struct CollectionResult { + pub service: DirectoryService, + pub indexed_at: String, + pub collection: CollectionSummary, + #[serde(flatten)] + pub additional: AdditionalMembers, +} + +#[derive(Clone, Debug, Deserialize, PartialEq)] +pub struct ServiceReference { + pub service_id: String, + pub service_origin: String, + pub name: Option, + #[serde(flatten)] + pub additional: AdditionalMembers, +} + +#[derive(Clone, Debug, Deserialize, PartialEq)] +pub struct CollectionSummary { + pub id: String, + pub name: String, + pub description: Option, + #[serde(flatten)] + pub additional: AdditionalMembers, +} + +#[derive(Clone, Debug, PartialEq)] +pub struct DirectoryIssue { + pub index: usize, + pub message: String, +} + +#[derive(Clone, Debug, PartialEq)] +pub struct SearchResponse { + pub facets: Option, + pub items: Vec, + pub next: Option, + pub issues: Vec, + pub additional: AdditionalMembers, +} + #[derive(Clone, Debug, Deserialize, PartialEq)] pub struct DirectoryService { pub description: String, @@ -92,6 +169,12 @@ pub struct DirectoryService { pub additional: AdditionalMembers, } +impl DirectoryService { + pub fn service_id(&self) -> Option<&str> { + self.additional.get("service_id").and_then(Value::as_str) + } +} + #[derive(Clone, Debug, Deserialize, PartialEq)] pub struct Facet { pub count: u64, @@ -131,8 +214,11 @@ pub struct SearchPage { pub additional: BTreeMap, } -#[derive(Clone, Debug, Default, PartialEq)] +#[derive(Clone, Debug, Default, PartialEq, Serialize)] pub struct SuggestionRequest { + #[serde(skip_serializing_if = "Option::is_none")] + pub filters: Option, + #[serde(skip_serializing_if = "is_zero")] pub limit: usize, pub prefix: String, } diff --git a/crates/odp-directory/src/results.rs b/crates/odp-directory/src/results.rs new file mode 100644 index 0000000..6af7007 --- /dev/null +++ b/crates/odp-directory/src/results.rs @@ -0,0 +1,164 @@ +use odp_core::{derive_service_origin, is_local_resource_identifier, parse_agent_service_document}; +use serde::Deserialize; +use serde_json::{Value, json}; + +use crate::{ + CollectionResult, DirectoryError, DirectoryIssue, DirectoryResult, Facets, SearchResponse, + ServiceResult, +}; + +pub(crate) fn decode(body: &[u8]) -> Result { + #[derive(Deserialize)] + struct Envelope { + items: Vec, + facets: Option, + next: Option, + #[serde(flatten)] + additional: odp_core::AdditionalMembers, + } + let envelope: Envelope = serde_json::from_slice(body).map_err(invalid)?; + if envelope.items.len() > 100 { + return Err(invalid("Directory response exceeds 100 results")); + } + if envelope + .next + .as_ref() + .is_some_and(|next| next.trim().is_empty()) + { + return Err(invalid("Directory continuation is empty")); + } + if envelope.facets.as_ref().is_some_and(|facets| { + facets + .trust + .iter() + .any(|facet| facet.value.name != odp_core::Protocol::Tap) + }) { + return Err(invalid("Directory trust facets are invalid")); + } + let mut items = Vec::new(); + let mut issues = Vec::new(); + for (index, raw) in envelope.items.into_iter().enumerate() { + match result(raw) { + Ok(item) => items.push(item), + Err(error) => issues.push(DirectoryIssue { + index, + message: error.to_string(), + }), + } + } + Ok(SearchResponse { + items, + issues, + facets: envelope.facets, + next: envelope.next, + additional: envelope.additional, + }) +} + +fn result(mut raw: Value) -> Result { + let kind = text(&raw, "type", 128)?.to_owned(); + if kind != "service" && kind != "collection" { + return Ok(DirectoryResult::Unknown { kind, raw }); + } + text(&raw, "indexed_at", 64)?; + let service = raw + .get_mut("service") + .ok_or_else(|| invalid("Missing service"))?; + text(service, "service_id", 128)?; + origin(service)?; + text(service, "indexed_at", 64)?; + let object = service + .as_object_mut() + .ok_or_else(|| invalid("service must be an object"))?; + for name in [ + "branding", + "http", + "mcp", + "odp_version", + "payment_origins", + "search_capabilities", + ] { + object.remove(name); + } + let mut document = Value::Object(object.clone()); + document["odp_version"] = json!("1.0"); + document["http"] = json!({"endpoint_base":"/"}); + let parsed = parse_agent_service_document(&serde_json::to_vec(&document).map_err(invalid)?) + .map_err(invalid)?; + object.insert( + "operations".to_owned(), + serde_json::to_value(parsed.operations).map_err(invalid)?, + ); + if let Some(protocols) = parsed.protocols { + object.insert( + "protocols".to_owned(), + serde_json::to_value(protocols).map_err(invalid)?, + ); + } else { + object.remove("protocols"); + } + if kind == "service" { + if let Some(reference) = raw.get("available_through") { + text(reference, "service_id", 128)?; + origin(reference)?; + if reference.get("name").is_some() { + text(reference, "name", 128)?; + } + } + raw.as_object_mut() + .ok_or_else(|| invalid("Invalid result"))? + .remove("type"); + Ok(DirectoryResult::Service(Box::new( + serde_json::from_value::(raw).map_err(invalid)?, + ))) + } else { + let collection = raw + .get("collection") + .ok_or_else(|| invalid("Missing collection"))?; + if !is_local_resource_identifier(text(collection, "id", 128)?) { + return Err(invalid("collection.id must be a local resource identifier")); + } + text(collection, "name", 128)?; + if let Some(description) = collection.get("description") { + if description + .as_str() + .is_none_or(|value| value.chars().count() > 1024) + { + return Err(invalid( + "collection.description must be a string of at most 1024 characters", + )); + } + } + raw.as_object_mut() + .ok_or_else(|| invalid("Invalid result"))? + .remove("type"); + Ok(DirectoryResult::Collection(Box::new( + serde_json::from_value::(raw).map_err(invalid)?, + ))) + } +} + +fn origin(value: &Value) -> Result<(), DirectoryError> { + let origin = text(value, "service_origin", 2048)?; + if !origin.starts_with("https://") || derive_service_origin(origin).map_err(invalid)? != origin + { + return Err(invalid("service_origin must be a canonical HTTPS origin")); + } + Ok(()) +} + +fn text<'a>(value: &'a Value, field: &str, maximum: usize) -> Result<&'a str, DirectoryError> { + value + .get(field) + .and_then(Value::as_str) + .filter(|text| !text.trim().is_empty() && text.chars().count() <= maximum) + .ok_or_else(|| { + invalid(format!( + "{field} must be a nonempty string of at most {maximum} characters" + )) + }) +} + +fn invalid(error: impl std::fmt::Display) -> DirectoryError { + DirectoryError::InvalidResponse(error.to_string()) +} diff --git a/crates/odp-directory/tests/mixed.rs b/crates/odp-directory/tests/mixed.rs new file mode 100644 index 0000000..982189c --- /dev/null +++ b/crates/odp-directory/tests/mixed.rs @@ -0,0 +1,404 @@ +use std::{ + collections::{BTreeMap, VecDeque}, + sync::{Arc, Mutex}, +}; + +use async_trait::async_trait; +use odp_directory::*; +use serde_json::{Value, json}; + +#[derive(Default)] +struct Stub { + requests: Mutex>, + replies: Mutex>, +} + +impl Stub { + fn reply(&self, status: u16, body: Value, headers: BTreeMap) { + self.replies.lock().unwrap().push_back(HttpResponse { + status, + body: serde_json::to_vec(&body).unwrap(), + headers: BTreeMap::from([("content-type".to_owned(), "application/json".to_owned())]) + .into_iter() + .chain(headers) + .collect(), + }); + } + fn ok(&self, body: Value) { + self.reply(200, body, BTreeMap::new()); + } +} + +#[async_trait] +impl Transport for Stub { + async fn send(&self, request: HttpRequest) -> Result { + self.requests.lock().unwrap().push(request); + Ok(self + .replies + .lock() + .unwrap() + .pop_front() + .expect("unexpected request")) + } +} + +fn setup() -> (DirectoryClient, Arc) { + let stub = Arc::new(Stub::default()); + ( + DirectoryClient::with_transport(Environment::Sandbox, stub.clone()), + stub, + ) +} + +fn service() -> Value { + json!({"service_id":"parent", "service_origin":"https://api.example.com", + "indexed_at":"2026-09-18T11:00:00Z", "name":"Data", "description":"Data services.", + "language":"en", "localizations":["en"], "operations":[ + {"name":"get-offering","authentication":"not-required"}, + {"name":"list-offerings","authentication":"not-required"}], + "protocols":{"trust":[{"name":"tap"},{"name":"future"}]}}) +} + +#[tokio::test] +async fn sends_suggestion_filters_without_mutating_the_request() { + let (client, stub) = setup(); + stub.ok(json!({"items":["Weather"]})); + let mut request = SuggestionRequest { + prefix: " we ".to_owned(), + filters: Some(ServiceFilters { + keywords: vec!["weather".to_owned()], + ..Default::default() + }), + ..Default::default() + }; + assert_eq!(client.suggest(&request).await.unwrap(), ["Weather"]); + assert_eq!(request.prefix, " we "); + assert!(client.suggest_services(&request).await.is_err()); + request.filters.as_mut().unwrap().keywords = vec![String::new()]; + assert!(client.suggest(&request).await.is_err()); + let requests = stub.requests.lock().unwrap(); + assert_eq!(requests.len(), 1); + assert_eq!(requests[0].method, "POST"); + assert_eq!( + serde_json::from_slice::(&requests[0].body).unwrap(), + json!({"prefix":"we", "filters":{"keywords":["weather"]}}) + ); +} + +fn item(kind: &str) -> Value { + let mut item = json!({"type":kind, "service":service(), "indexed_at":"2026-09-18T12:00:00Z"}); + if kind == "collection" { + item["collection"] = + json!({"id":"Weather","name":"Weather forecasts","description":"Forecasts."}); + } + item +} + +#[tokio::test] +async fn decodes_mixed_results_and_preserves_unknown_types() { + let (client, stub) = setup(); + let mut first = item("service"); + first["available_through"] = json!({"service_id":"platform", "service_origin":"https://platform.example", "name":"Platform"}); + first["extra"] = json!(true); + let future = json!({"type":"future","nested":{"data":42}}); + stub.ok(json!({"items":[first,item("collection"),future],"extra":42, + "facets":{"keywords":[{"value":"weather","count":12}]}})); + let response = client + .search(&ResourceSearchRequest::default()) + .await + .unwrap(); + assert!(response.issues.is_empty()); + assert_eq!(response.additional["extra"], 42); + assert_eq!(response.facets.unwrap().keywords[0].count, 12); + let DirectoryResult::Service(service) = &response.items[0] else { + panic!("service") + }; + assert_eq!(service.service.service_id(), Some("parent")); + assert_eq!( + service.available_through.as_ref().unwrap().name.as_deref(), + Some("Platform") + ); + assert_eq!(service.additional["extra"], true); + assert_eq!(service.service.protocols.as_ref().unwrap().trust.len(), 1); + let DirectoryResult::Collection(collection) = &response.items[1] else { + panic!("collection") + }; + assert_eq!(collection.collection.id, "Weather"); + assert_ne!(collection.indexed_at, collection.service.indexed_at); + let DirectoryResult::Unknown { kind, raw } = &response.items[2] else { + panic!("unknown") + }; + assert_eq!(kind, "future"); + assert_eq!(raw, &future); +} + +#[tokio::test] +async fn isolates_malformed_known_items_and_normalizes_future_operations() { + let (client, stub) = setup(); + let mut invalid = Vec::new(); + for (pointer, value) in [ + ("/type", json!(null)), + ("/service/service_id", json!("")), + ( + "/service/service_origin", + json!("https://user@api.example.com"), + ), + ( + "/service/service_origin", + json!("https://api.example.com/path"), + ), + ("/service/operations", json!([])), + ("/indexed_at", json!(false)), + ("/collection/id", json!("../bad")), + ("/collection/name", json!("")), + ("/collection/description", json!(null)), + ] { + let mut candidate = item("collection"); + *candidate.pointer_mut(pointer).unwrap() = value; + invalid.push(candidate); + } + let count = invalid.len(); + let mut valid = item("collection"); + valid["collection"]["description"] = json!(""); + valid["service"]["operations"] + .as_array_mut() + .unwrap() + .push(json!({"name":"future-operation","authentication":"not-required"})); + valid["service"]["http"] = json!({"endpoint_base":"https://untrusted.example"}); + invalid.push(valid); + stub.ok(json!({"items":invalid})); + let response = client + .search(&ResourceSearchRequest::default()) + .await + .unwrap(); + assert_eq!(response.issues.len(), count); + assert_eq!( + response + .issues + .iter() + .map(|issue| issue.index) + .collect::>(), + (0..count).collect::>() + ); + assert_eq!(response.items.len(), 1); + let DirectoryResult::Collection(item) = &response.items[0] else { + panic!("collection") + }; + assert_eq!(item.service.operations.len(), 2); + assert!(!item.service.additional.contains_key("http")); +} + +#[tokio::test] +async fn keeps_routes_bodies_and_suggestions_separate() { + let (client, stub) = setup(); + stub.ok(json!({"items":[]})); + client + .search(&ResourceSearchRequest { + query: "weather".to_owned(), + types: Some(vec![ResultType::Collection]), + ..Default::default() + }) + .await + .unwrap(); + stub.ok(json!({"items":[],"next":"/v1/directory/search?cursor=opaque"})); + assert_eq!( + client + .continue_search("/v1/directory/search?cursor=opaque") + .await + .unwrap() + .next + .as_deref(), + Some("/v1/directory/search?cursor=opaque") + ); + stub.ok(json!({"items":[service()]})); + assert_eq!( + client + .search_services(&SearchRequest::default()) + .await + .unwrap() + .items + .len(), + 1 + ); + stub.ok(json!({"items":[]})); + client + .continue_search_services("/v1/services/search?cursor=old") + .await + .unwrap(); + for mixed in [true, false] { + stub.ok(json!({"items":["Weather forecasts"]})); + let request = SuggestionRequest { + prefix: "we".to_owned(), + limit: 10, + ..Default::default() + }; + let names = if mixed { + client.suggest(&request).await + } else { + client.suggest_services(&request).await + } + .unwrap(); + assert_eq!(names, ["Weather forecasts"]); + } + let requests = stub.requests.lock().unwrap(); + assert_eq!( + requests + .iter() + .map(|request| request.method.as_str()) + .collect::>(), + ["POST", "GET", "POST", "GET", "POST", "GET"] + ); + assert_eq!( + serde_json::from_slice::(&requests[0].body).unwrap(), + json!({"query":"weather","types":["collection"]}) + ); + assert!(requests[1].body.is_empty()); + for (request, path) in requests.iter().zip([ + "/v1/directory/search", + "/v1/directory/search?cursor=opaque", + "/v1/services/search", + "/v1/services/search?cursor=old", + "/v1/directory/suggestions", + "/v1/services/suggestions?prefix=we&limit=10", + ]) { + assert_eq!(request.url, format!("https://sandbox.inflowpay.ai{path}")); + } +} + +#[tokio::test] +async fn rejects_invalid_requests_before_transport() { + let (client, stub) = setup(); + for request in [ + ResourceSearchRequest { + types: Some(vec![]), + ..Default::default() + }, + ResourceSearchRequest { + types: Some(vec![ResultType::Service, ResultType::Service]), + ..Default::default() + }, + ResourceSearchRequest { + limit: 101, + ..Default::default() + }, + ResourceSearchRequest { + query: " ".to_owned(), + ..Default::default() + }, + ] { + assert!(matches!( + client.search(&request).await, + Err(DirectoryError::InvalidRequest(_)) + )); + } + for next in [ + "", + " ", + "https://other.example/", + "https://user@sandbox.inflowpay.ai/", + ] { + assert!(client.continue_search(next).await.is_err()); + } + assert!(client.suggest(&SuggestionRequest::default()).await.is_err()); + assert!(stub.requests.lock().unwrap().is_empty()); + assert_eq!( + serde_json::to_value(ResourceSearchRequest::default()).unwrap(), + json!({}) + ); +} + +#[tokio::test] +async fn bounds_service_aggregation_without_extra_requests() { + for options in [ + IterationOptions { + max_items: 1, + max_pages: 10, + }, + IterationOptions { + max_items: 10, + max_pages: 1, + }, + ] { + let (client, stub) = setup(); + stub.ok(json!({"items":[service()],"next":"/v1/services/search?cursor=more"})); + assert_eq!( + client + .collect_services(&SearchRequest::default(), options) + .await + .unwrap() + .len(), + 1 + ); + assert_eq!(stub.requests.lock().unwrap().len(), 1); + } + let (client, stub) = setup(); + stub.ok(json!({"items":[service()],"next":"/v1/services/search?cursor=more"})); + stub.ok(json!({"items":[service()]})); + assert_eq!( + client + .collect_services(&SearchRequest::default(), IterationOptions::default()) + .await + .unwrap() + .len(), + 2 + ); + assert_eq!(stub.requests.lock().unwrap().len(), 2); +} + +#[tokio::test] +async fn rejects_invalid_envelopes_and_transport_failures() { + let (client, stub) = setup(); + for body in [ + json!(null), + json!({}), + json!({"items":null}), + json!({"items":vec![item("service");101]}), + json!({"items":[],"next":false}), + json!({"items":[],"facets":{"trust":[{"value":{"name":"mpp"},"count":1}]}}), + ] { + stub.ok(body); + assert!( + client + .search(&ResourceSearchRequest::default()) + .await + .is_err() + ); + } + stub.reply( + 429, + json!({"detail":"rate limited"}), + BTreeMap::from([("retry-after".to_owned(), "5".to_owned())]), + ); + assert!(matches!( + client.search(&ResourceSearchRequest::default()).await, + Err(DirectoryError::Request { status: 429, .. }) + )); + stub.reply( + 307, + json!(null), + BTreeMap::from([("location".to_owned(), "https://other.example/".to_owned())]), + ); + assert!( + client + .search(&ResourceSearchRequest::default()) + .await + .is_err() + ); + stub.reply( + 303, + json!(null), + BTreeMap::from([("location".to_owned(), "/redirected".to_owned())]), + ); + stub.ok(json!({"items":[]})); + client + .search(&ResourceSearchRequest { + query: "weather".to_owned(), + ..Default::default() + }) + .await + .unwrap(); + let requests = stub.requests.lock().unwrap(); + let last = requests.last().unwrap(); + assert_eq!(last.method, "GET"); + assert!(last.body.is_empty()); +} diff --git a/examples/Cargo.toml b/examples/Cargo.toml index d452a14..7bf3b16 100644 --- a/examples/Cargo.toml +++ b/examples/Cargo.toml @@ -24,5 +24,9 @@ path = "odp-agent-discovery/main.rs" name = "odp-service-small" path = "odp-service-small/main.rs" +[[bin]] +name = "odp-directory-discovery" +path = "odp-directory-discovery/main.rs" + [lints] workspace = true diff --git a/examples/README.md b/examples/README.md index 298f93a..5683f54 100644 --- a/examples/README.md +++ b/examples/README.md @@ -6,6 +6,7 @@ The examples demonstrate both sides of a minimal ODP integration: | --- | --- | | [`odp-service-small`](odp-service-small/) | Publishes a validated in-memory catalog through the framework-neutral Service runtime. | | [`odp-agent-discovery`](odp-agent-discovery/) | Injects a mock Directory into the top-level Agent, performs federated discovery, and navigates Collections, Offerings, and Actions. | +| [`odp-directory-discovery`](odp-directory-discovery/main.rs) | Searches the canonical Directory for Services and Collections, then retrieves anonymous Collection details. | Run the small Service in one terminal, then the Agent in another: @@ -17,3 +18,16 @@ cargo run -p odp-examples --bin odp-agent-discovery The Agent also accepts Service origins as positional arguments. It converts the compatible origins into an in-memory Directory response, while all Service requests continue to use the real HTTP endpoints. + +## Canonical Directory discovery + +```sh +cargo run -p odp-examples --bin odp-directory-discovery -- sandbox weather +``` + +Use `production` instead of `sandbox` for production. Omit `weather` to browse rather than search. +This example requires a deployment with `/v1/directory/search`. It requests at most five results, +prints Service and Collection names, reports malformed or unknown results, and retrieves full +Collection details only after inspecting the owning Service's advertised anonymous operation. +It does not enroll, pay, or invoke Actions. Unlike `odp-agent-discovery`, it uses the real Directory, +not a mock. The server's bounded result list is not an exhaustive catalog listing. diff --git a/examples/odp-directory-discovery/main.rs b/examples/odp-directory-discovery/main.rs new file mode 100644 index 0000000..97931c3 --- /dev/null +++ b/examples/odp-directory-discovery/main.rs @@ -0,0 +1,51 @@ +use odp_agent::ServiceClient; +use odp_core::{AuthenticationRequirement, Operation}; +use odp_directory::{DirectoryClient, DirectoryResult, Environment, ResourceSearchRequest}; + +#[tokio::main(flavor = "current_thread")] +async fn main() -> Result<(), Box> { + let mut args = std::env::args().skip(1); + let environment = match args.next().as_deref().unwrap_or("production") { + "production" => Environment::Production, + "sandbox" => Environment::Sandbox, + _ => return Err("Environment must be production or sandbox".into()), + }; + let directory = DirectoryClient::new(environment)?; + let response = directory + .search(&ResourceSearchRequest { + query: args.collect::>().join(" "), + limit: 5, + ..Default::default() + }) + .await?; + for issue in response.issues { + eprintln!("Skipped result {}: {}", issue.index, issue.message); + } + for result in response.items { + match result { + DirectoryResult::Service(item) => println!( + "Service: {} ({})", + item.service.name, item.service.service_origin + ), + DirectoryResult::Collection(item) => { + println!( + "Collection: {} ({}, through {})", + item.collection.name, item.collection.id, item.service.service_origin + ); + let client = ServiceClient::new(&item.service.service_origin)?; + let inspection = client.inspect().await?; + if inspection.document.operations.iter().any(|operation| { + operation.name == Operation::GetCollection + && operation.authentication != AuthenticationRequirement::Required + }) { + let collection = client.get_collection(&item.collection.id).await?; + println!("{}", serde_json::to_string_pretty(&collection)?); + } else { + println!("The Service does not advertise anonymous Collection retrieval."); + } + } + DirectoryResult::Unknown { kind, .. } => println!("Unsupported result type: {kind}"), + } + } + Ok(()) +} diff --git a/scripts/verify-consumer.sh b/scripts/verify-consumer.sh index 9b6d545..409da83 100755 --- a/scripts/verify-consumer.sh +++ b/scripts/verify-consumer.sh @@ -42,7 +42,7 @@ EOF cat > "$consumer/src/main.rs" <<'EOF' use odp_agent::ServiceClient; use odp_core::{Representation, ResourceIdentity, ResourceType}; -use odp_directory::{DirectoryClient, Environment}; +use odp_directory::{DirectoryClient, Environment, ResourceSearchRequest, ResultType}; use odp_service::ServiceBuilder; fn main() { @@ -53,6 +53,10 @@ fn main() { "rubber-plant", ); let _ = DirectoryClient::new(Environment::Production); + let _ = ResourceSearchRequest { + types: Some(vec![ResultType::Collection]), + ..Default::default() + }; let _ = ServiceBuilder::new("Example", "Example Service", "en", "/odp"); let _ = Representation::Terse; }