Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

5 changes: 3 additions & 2 deletions crates/odp-agent/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,9 @@ categories.workspace = true
readme = "README.md"

[dependencies]
async-trait.workspace = true
futures.workspace = true
getrandom.workspace = true
httpdate.workspace = true
jsonschema.workspace = true
odp-core = { version = "0.1.1", path = "../odp-core" }
Expand All @@ -23,11 +25,10 @@ serde.workspace = true
serde_json.workspace = true
sha2.workspace = true
thiserror.workspace = true
tokio.workspace = true
url.workspace = true

[dev-dependencies]
async-trait.workspace = true
tokio.workspace = true

[lints]
workspace = true
16 changes: 16 additions & 0 deletions crates/odp-agent/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,20 @@ is validated. `ServiceClient` uses an in-memory cache by default; callers can in
set an authentication-aware cache partition, or override the Service Document, Collection, and
Offering fallback lifetimes.

Each client has an isolated cache partition, including when sharing a `Cache`. Use
`with_cache_partition` only to share responses between clients with the same authentication context.
Manual `continue_offerings` and `continue_collections` calls use HTTP freshness headers without a
fallback lifetime: an opaque continuation URL does not identify whether it originated in a search.
The automatic traversal methods retain the originating operation's fallback policy.

`ServiceClient::new` permits public destinations only. Use `ServiceClient::for_local_development`
for localhost examples. The default transport pins validated addresses for each request, disables
proxies, and checks the connected peer. An injected transport owns its network policy; wrapping it
in `SecureTransport` requires implementing `Transport::send_to` with equivalent address pinning and
peer verification. Its default implementation refuses the request rather than silently bypassing
those checks. Implement `send_limited` to enforce response limits while streaming; its default
implementation can only check the completed response.

`get_offering_details` bundles an Offering with its validated Attribute Schema, validates the
Offering attributes, and normalizes usable Action targets. `resolve_action` resolves an Action's
request schema or unique OpenAPI 3.1 operation without invoking the target. Supporting documents
Expand All @@ -54,6 +68,8 @@ resolution accepts JSON Schema Draft 2020-12 and is limited to 256 KiB per docum
eight reference levels, and one MiB for the complete graph. OpenAPI documents are limited to one
MiB. These are fixed SDK safety ceilings. Cross-document schema composition uses `$ref`;
`$dynamicRef` accepts only a fragment reference such as `#node`.
Returned schemas include their referenced resources in `$defs` with absolute identifiers, so a
caller can use them without additional network access.

## Search across Services

Expand Down
148 changes: 147 additions & 1 deletion crates/odp-agent/src/agent.rs
Original file line number Diff line number Diff line change
Expand Up @@ -184,7 +184,7 @@ mod tests {
impl Transport for ServiceTransport {
async fn send(&self, request: HttpRequest) -> Result<HttpResponse, TransportError> {
if request.url.ends_with("/.well-known/odp") {
return Ok(odp_response(br#"{"description":"Plants","http":{"endpoint_base":"/odp"},"language":"en","localizations":["en"],"name":"Plants","odp_version":"1.0","operations":[{"authentication":"not-required","name":"get-offering"},{"authentication":"not-required","name":"list-offerings"}]}"#));
return Ok(odp_response(br#"{"description":"Plants","http":{"endpoint_base":"/odp"},"language":"en","localizations":["en"],"name":"Plants","odp_version":"1.0","operations":[{"authentication":"not-required","name":"get-offering"},{"authentication":"not-required","name":"list-offerings"},{"authentication":"not-required","name":"search-offerings"}]}"#));
}
let id = if request.url.starts_with("https://one.example") {
"one"
Expand Down Expand Up @@ -243,4 +243,150 @@ mod tests {
assert_eq!(events[0].service.name, "One");
assert_eq!(events[1].service.name, "Two");
}

/// A search request drives the search operation; a bare request just lists.
#[tokio::test]
async fn searches_a_service_only_when_the_request_asks_a_question() {
let recorder = Arc::new(RecordingFactory::default());
let directory =
DirectoryClient::with_transport(Environment::Production, Arc::new(DirectoryTransport));
let agent = Agent::with_clients(directory, recorder.clone());

agent
.search_offerings_across_services(&FederatedSearchRequest {
offerings: OfferingSearchRequest {
query: "rubber".to_owned(),
..OfferingSearchRequest::default()
},
..FederatedSearchRequest::default()
})
.await
.unwrap();

let urls = recorder.urls();
assert!(
urls.iter().any(|url| url.contains("/offerings/search")),
"{urls:?}"
);
}

/// FED-04: one Service that cannot answer is reported, and the others still report offerings.
#[tokio::test]
async fn reports_a_failing_service_without_losing_the_others() {
struct HalfBroken;

impl ServiceClientFactory for HalfBroken {
fn create(&self, service: &DirectoryService) -> Result<ServiceClient, AgentError> {
if service.service_origin.starts_with("https://one.example") {
return Err(AgentError::InvalidRequest("no client for One".to_owned()));
}
ServiceClient::with_transport(&service.service_origin, Arc::new(ServiceTransport))
}
}

let directory =
DirectoryClient::with_transport(Environment::Production, Arc::new(DirectoryTransport));
let agent = Agent::with_clients(directory, Arc::new(HalfBroken));
let events = agent
.search_offerings_across_services(&FederatedSearchRequest::default())
.await
.unwrap();

assert_eq!(events.len(), 2);
assert_eq!(events[0].service.name, "One");
assert!(events[0].offering.is_none());
assert!(events[0].issue.is_some());
assert!(events[1].offering.is_some());
assert!(events[1].issue.is_none());
}

/// Each bound has a ceiling, so a request asking for more is refused before anything is sent.
#[tokio::test]
async fn refuses_a_request_that_asks_for_more_than_the_bounds_allow() {
let directory =
DirectoryClient::with_transport(Environment::Production, Arc::new(DirectoryTransport));
let agent = Agent::with_clients(directory, Arc::new(Factory));

for (request, name) in [
(
FederatedSearchRequest {
max_services: 101,
..FederatedSearchRequest::default()
},
"max_services",
),
(
FederatedSearchRequest {
max_offerings_per_service: 101,
..FederatedSearchRequest::default()
},
"max_offerings_per_service",
),
(
FederatedSearchRequest {
concurrency: 17,
..FederatedSearchRequest::default()
},
"concurrency",
),
] {
let error = agent
.search_offerings_across_services(&request)
.await
.unwrap_err();
assert!(error.to_string().contains(name), "{name}: {error}");
}
}

#[test]
fn builds_an_agent_for_an_environment_it_keeps() {
let agent = Agent::new(Environment::Sandbox).unwrap();
assert_eq!(agent.environment(), Environment::Sandbox);
}

/// The default factory reaches the Service Origin the Directory listed.
#[test]
fn builds_a_default_client_for_a_listed_service() {
let service: DirectoryService = serde_json::from_str(
r#"{"description":"One","indexed_at":"2026-08-25T00:00:00Z","language":"en","localizations":["en"],"name":"One","operations":[],"service_origin":"https://plants.example"}"#,
)
.unwrap();
let client = DefaultServiceClientFactory.create(&service).unwrap();
assert_eq!(client.service_origin(), "https://plants.example");
}

/// A factory that records the URLs its clients are asked for.
#[derive(Default)]
struct RecordingFactory {
urls: Arc<std::sync::Mutex<Vec<String>>>,
}

impl RecordingFactory {
fn urls(&self) -> Vec<String> {
self.urls.lock().unwrap().clone()
}
}

impl ServiceClientFactory for RecordingFactory {
fn create(&self, service: &DirectoryService) -> Result<ServiceClient, AgentError> {
ServiceClient::with_transport(
&service.service_origin,
Arc::new(RecordingTransport {
urls: self.urls.clone(),
}),
)
}
}

struct RecordingTransport {
urls: Arc<std::sync::Mutex<Vec<String>>>,
}

#[async_trait]
impl Transport for RecordingTransport {
async fn send(&self, request: HttpRequest) -> Result<HttpResponse, TransportError> {
self.urls.lock().unwrap().push(request.url.clone());
ServiceTransport.send(request).await
}
}
}
69 changes: 65 additions & 4 deletions crates/odp-agent/src/cache.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,15 +21,48 @@ pub trait Cache: Send + Sync {
fn set(&self, key: String, record: CacheRecord) -> Result<(), String>;
}

#[derive(Default)]
const DEFAULT_CAPACITY: usize = 256;

/// An in-memory cache bounded by entry count, so a long-lived Agent cannot grow without limit.
/// When it is full the least recently stored record makes room for the new one.
pub struct MemoryCache {
capacity: usize,
records: RwLock<BTreeMap<String, CacheRecord>>,
}

impl Default for MemoryCache {
fn default() -> Self {
Self::with_capacity(DEFAULT_CAPACITY)
}
}

impl MemoryCache {
#[must_use]
pub fn new() -> Self {
Self::default()
}

/// A cache holding at most `capacity` records. A capacity of zero keeps nothing.
#[must_use]
pub fn with_capacity(capacity: usize) -> Self {
Self {
capacity,
records: RwLock::new(BTreeMap::new()),
}
}

#[must_use]
pub fn len(&self) -> usize {
self.records
.read()
.map(|records| records.len())
.unwrap_or_default()
}

#[must_use]
pub fn is_empty(&self) -> bool {
self.len() == 0
}
}

impl Cache for MemoryCache {
Expand All @@ -51,26 +84,54 @@ impl Cache for MemoryCache {
}

fn set(&self, key: String, record: CacheRecord) -> Result<(), String> {
self.records
let mut records = self
.records
.write()
.map_err(|_| "memory cache lock is poisoned".to_owned())?
.insert(key, record);
.map_err(|_| "memory cache lock is poisoned".to_owned())?;
if self.capacity == 0 {
return Ok(());
}
if !records.contains_key(&key) && records.len() >= self.capacity {
// Capacity is at least one here, so evicting just enough always leaves room.
let mut by_age = records
.iter()
.map(|(name, value)| (value.stored_at, name.clone()))
.collect::<Vec<_>>();
by_age.sort();
for (_, name) in by_age.into_iter().take(records.len() + 1 - self.capacity) {
records.remove(&name);
}
}
records.insert(key, record);
Ok(())
}
}

/// CCH-02: the freshness an Agent assumes when a response supplies none.
///
/// CCH-03 makes each resource class independently configurable, so every class the draft names has
/// its own field rather than borrowing a neighbour's.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct CacheFallbacks {
/// Filter and Sort Definitions.
pub capabilities: Duration,
pub collection: Duration,
pub offering: Duration,
/// Attribute Schema documents.
pub schema: Duration,
/// Search responses, which describe one request and are not reused for the next.
pub search: Duration,
pub service_document: Duration,
}

impl Default for CacheFallbacks {
fn default() -> Self {
Self {
capabilities: Duration::from_secs(60 * 60),
collection: Duration::from_secs(60 * 60),
offering: Duration::from_secs(5 * 60),
schema: Duration::from_secs(24 * 60 * 60),
search: Duration::ZERO,
service_document: Duration::from_secs(4 * 60 * 60),
}
}
Expand Down
21 changes: 18 additions & 3 deletions crates/odp-agent/src/capabilities.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ use odp_core::{
};
use url::Url;

use crate::{AgentError, CacheFallbacks, ServiceClient};
use crate::{AgentError, ServiceClient};

const MAXIMUM_CAPABILITY_PAGES: usize = 16;
const MAXIMUM_FILTERS: usize = 1_024;
Expand Down Expand Up @@ -272,7 +272,7 @@ impl ServiceClient {
let data = self
.linked_odp(
target,
CacheFallbacks::default().collection,
self.cache_fallbacks().capabilities,
validate_filter_page,
)
.await?;
Expand Down Expand Up @@ -307,7 +307,7 @@ impl ServiceClient {
let data = self
.linked_odp(
target,
CacheFallbacks::default().collection,
self.cache_fallbacks().capabilities,
validate_sort_page,
)
.await?;
Expand Down Expand Up @@ -419,4 +419,19 @@ mod tests {
assert_eq!(catalog.filters["price"].title, "Price");
assert_eq!(catalog.sorts["price-lowest"].filters[0].id, "price");
}

/// A validated document never carries one, so the helper's own guard is checked here.
#[test]
fn refuses_a_capability_reference_that_is_not_an_http_url() {
for reference in ["mailto:filters@plants.example", "file:///filters.json"] {
let error = resolve_reference(reference, "https://plants.example").unwrap_err();
assert!(error.to_string().contains("HTTP"), "{reference}: {error}");
}
}

#[test]
fn refuses_a_capability_reference_it_has_no_base_for() {
let error = resolve_reference("/filters", "not a base").unwrap_err();
assert!(matches!(error, AgentError::InvalidRequest(_)), "{error}");
}
}
Loading
Loading