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
42 changes: 37 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -203,9 +203,14 @@ The Agent module also provides:
- Conditional request and representation caching with injectable `Cache` and `Transport` protocols.

Default fallback cache lifetimes are four hours for Service documents, one hour for Collections,
and five minutes for Offerings. HTTP cache directives take precedence. Provide distinct `transport`
and `supporting_transport` instances when protocol resources and linked schemas require different
credentials or network policy.
five minutes for Offerings, zero for searches, one hour for Filter and Sort Definitions, and
24 hours for Attribute Schemas. Set each independently using `CacheFallbacks` (`service_document`,
`collection`, `offering`, `search`, `filters`, `sorts`, and `attribute_schema`). HTTP cache directives
take precedence. Continuations retain their originating operation's fallback; an unrecognized
continuation uses the search fallback. `ServiceClient` uses a
separate anonymous transport for linked schemas and OpenAPI documents, even when its primary
`transport` has authentication configured. An explicit `supporting_transport` override must also
send these requests anonymously; it must not share the primary transport's credentials or cookies.

### Search across Services

Expand Down Expand Up @@ -264,11 +269,22 @@ invoke the resolved target.

### Caching and HTTP transport

Clients with a supplied `transport` use separate cache partitions by default. To share cached
responses between these clients, explicitly supply the same `cache_partition` only when they use
the same authentication context. Create a new client or select a new partition when changing
credentials. SDK-owned anonymous transports can share their anonymous partition.

`MemoryCache` is the default process-local cache. Implement the `Cache` protocol when representations
must survive process restarts or share storage across workers. A custom `Transport` implements
asynchronous `send()` and `aclose()` methods. Caller-provided caches and transports remain owned by
the caller.

`HttpRequest.maximum_response_bytes` gives a custom transport the response budget. Enforce it while
reading, rather than buffering the complete response first. The built-in transport closes responses
on overflow, read failure, and cancellation. A successful response exceeding its budget raises
`TransportError` with `code="RESPONSE_LIMIT_EXCEEDED"`; `ServiceClient` preserves that code on
`AgentError`. Oversized error bodies are discarded while retaining the HTTP status and headers.

The built-in HTTP transport resolves and validates every destination before connecting, pins the
connection to a validated public address, does not inherit proxy settings from the environment, and
sends supporting-document requests without credentials. A custom transport must preserve those ODP
Expand All @@ -278,7 +294,8 @@ network and credential-isolation requirements. Local HTTP development is disable

Attribute Schema resolution accepts JSON Schema Draft 2020-12 and is limited to 256 KiB per
document, 16 documents, eight reference levels, and one MiB for the complete graph. OpenAPI
documents are limited to one MiB. These are fixed SDK safety ceilings. Linked schema documents must
documents are limited to one MiB and 32 levels of JSON nesting. Every other ODP response is limited
to 16 levels of nesting, and the Service Document to eight. These are fixed SDK safety ceilings. Linked schema documents must
use HTTPS. Cross-document schema composition uses `$ref`; `$dynamicRef` accepts only a fragment
reference such as `#node`.

Expand Down Expand Up @@ -340,10 +357,21 @@ its Actions can advertise enrollment, payment, and trust protocols, but ODP does
credentials, invoke Actions, submit payments, or implement trust protocols. Applications compose
the appropriate protocol clients around an Action resolved through ODP.

`parse_service_document` is the strict current-version Service parser. Agent inspection and
`parse_service_document` validates Service metadata against the supported ODP major version.
Compatible minor versions such as `1.7` are accepted without rewriting the received version;
SDK-generated documents use `1.0`. Agent inspection and
Directory results filter unrecognized enrollment, payment, and trust descriptors while retaining
strict validation for recognized descriptors.

Individual Offering and Collection GETs default to full representations; list and search operations
default to terse items. The Service handler writes `odp_version` on standalone resources and page
envelopes, omitting it from embedded page items without changing the Catalog's models.
A Catalog receives the requested language in `CatalogRequest.language` and
declares the language it actually returns on each resource. Static catalogs do not translate content.

Refinement parsing detects duplicate JSON values without guessing the type of a string. Comparing
decimal or date-time strings by their meaning requires the referenced Filter Definition.

## Errors and validation

Each role exposes typed errors:
Expand All @@ -357,6 +385,10 @@ Protocol models preserve additive members in `model.additional` and round-trip t
`model.to_dict()`. Parsing remains strict for normative constraints and fields that prohibit unknown
members.

Directory records retain additional metadata such as branding and MCP endpoints. These are discovery
hints, not authorization or authoritative routing data. The default Agent factory uses the record's
`service_origin` and retrieves that Service's own document before making catalog requests.

Handle the narrowest error that the application can act upon and use the role's base error for the
remaining failures:

Expand Down
164 changes: 160 additions & 4 deletions scripts/conformance_adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@

from offering_protocol.agent import ServiceClient
from offering_protocol.core import (
Collection,
CollectionSearchRequest,
Offering,
OfferingPage,
Expand Down Expand Up @@ -38,7 +39,14 @@
)
from offering_protocol.core.validation import _normalize_agent_response
from offering_protocol.directory.transport import HttpRequest, HttpResponse
from offering_protocol.service import CatalogRequest, Request, ServiceBuilder
from offering_protocol.service import (
CatalogError,
CatalogRequest,
Request,
ServiceBuilder,
StaticCatalog,
StaticCatalogOptions,
)


class MapTransport:
Expand Down Expand Up @@ -241,7 +249,149 @@ async def evaluate_errors_limits(case: dict[str, Any]) -> bool | None:
return (result.status == 200) == case["valid"]


def capability_definition(value: dict[str, Any], kind: str) -> dict[str, Any]:
defaults: dict[str, Any] = {"title": value["id"], "description": value["id"]}
if kind == "filters":
defaults.update(type="string", operators=["eq"])
else:
defaults["keys"] = [{"filter_id": "region", "direction": "ascending", "missing": "last"}]
result = {**defaults, **value}
if kind == "sorts":
result["keys"] = [
{"direction": "ascending", "missing": "last", **key} for key in result["keys"]
]
return result


def capability_advertisement(value: dict[str, Any]) -> dict[str, Any]:
result = {}
for kind, source in value.items():
if isinstance(source, list):
if not source:
continue
source = {"inline": source}
result[kind] = {
**source,
**(
{"inline": [capability_definition(item, kind) for item in source["inline"]]}
if "inline" in source
else {}
),
}
return result


async def evaluate_capabilities(case: dict[str, Any]) -> bool | None:
operation = case["operation"]
origin = case.get("service_origin", "https://service.example")
document = service_document()
document["operations"] = [
{"name": name, "authentication": "not-required"}
for name in {
"get-offering",
"list-offerings",
*case.get("operations", ["search-offerings"]),
}
]
if operation == "validate-advertisement":
advertisement = capability_advertisement(case["advertisement"])
document["search_capabilities"] = advertisement
try:
parse_service_document(json.dumps(document))
for source in advertisement.values():
if "linked" in source:
resolve_continuation(source["linked"]["href"], origin)
valid = True
except ValueError:
valid = False
return valid == case["valid"]

documents = {}
if operation in {"validate-linked-source", "validate-linked-page-count"}:
reference = case.get("href", "/filters/0")
document["search_capabilities"] = {"filters": {"linked": {"href": reference}}}
if operation == "validate-linked-source":
pages = case["pages"]
else:
pages = [
{
"odp_version": "1.0",
"items": [{"id": f"f{i}"}],
**({"next": f"/filters/{i + 1}"} if i + 1 < case["page_count"] else {}),
}
for i in range(case["page_count"])
]
for page in pages:
documents[resolve_continuation(reference, origin)] = response(
{
**page,
"items": [capability_definition(item, "filters") for item in page["items"]],
}
)
reference = page.get("next", "")
elif operation == "merge-capabilities":
document["search_capabilities"] = capability_advertisement(case["service"])
document["operations"].append({"name": "get-collection", "authentication": "not-required"})
if "collection_id" in case:
documents[f"{origin}/odp/collections/{case['collection_id']}?representation=full"] = (
response(
{
"odp_version": "1.0",
"id": case["collection_id"],
"name": "Collection",
"search_capabilities": capability_advertisement(
case["selected_collection"]
),
}
)
)
else:
return None
documents[f"{origin}/.well-known/odp"] = response(document)
async with ServiceClient(origin, transport=MapTransport(documents)) as client:
if "collection_id" in case:
result = await client.get_collection_search_capabilities(case["collection_id"])
else:
result = await client.get_offering_search_capabilities()
if operation == "merge-capabilities":
expected = case["expected"]
return (
sorted(result.filters) == sorted(expected["filter_ids"])
and sorted(result.sorts) == sorted(expected["sort_ids"])
and len(result.issues) == len(expected["issues"])
)
return (not result.issues) == case["valid"]


async def evaluate_case(subject: str, case: dict[str, Any], role: str) -> bool | None:
if subject == "search-capability-contract" and role == "agent":
return await evaluate_capabilities(case)
if subject == "collection-hierarchy" and role == "service":
if "chain_length" in case:
collections = [
{"id": f"c{i}", **({"parent_ids": [f"c{i - 1}"]} if i else {})}
for i in range(case["chain_length"] + 1)
]
else:
collections = case["collections"]
try:
StaticCatalog(
StaticCatalogOptions(
collections=tuple(
Collection.model_validate(
{"odp_version": "1.0", "name": value["id"], **value}
)
for value in collections
)
)
)
valid = True
except (CatalogError, ValueError):
valid = False
return valid == case["valid"]
if subject == "protocol-version":
document = {"odp_version": case["received"], "id": "item", "name": "Item"}
return succeeds(lambda: parse_offering(json.dumps(document))) == case["compatible"]
if subject == "local-identifier":
return is_local_resource_identifier(case["value"]) == case["valid"]
if subject == "identity-comparison":
Expand Down Expand Up @@ -315,8 +465,13 @@ async def evaluate_case(subject: str, case: dict[str, Any], role: str) -> bool |
)
return valid == case["valid"]
if case.get("operation") == "validate-limit":
limit = case["limit"]
valid = isinstance(limit, int) and not isinstance(limit, bool) and 1 <= limit <= 100
service = ServiceBuilder("Conformance", "Conformance", "en", "/odp").build(
EmptyCatalog()
)
result = await service.handle(
Request("GET", "/odp/offerings", query=f"limit={case['limit']}")
)
valid = result.status == 200
return valid == case["valid"]
if case.get("operation") == "validate-next":
valid = succeeds(lambda: resolve_continuation(case["next"], case["service_origin"]))
Expand Down Expand Up @@ -378,7 +533,8 @@ async def evaluate(request: dict[str, Any]) -> dict[str, object]:
return {
"status": "skipped",
"message": (
f"No public Python operation maps {request['vector']['subject']}/{operation}"
f"Adapter does not exercise {request['vector']['subject']}/{operation}; "
"this is not a passing conformance result"
),
}
return (
Expand Down
22 changes: 18 additions & 4 deletions scripts/node_interoperability.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
import sys

from offering_protocol.agent import ServiceClient
from offering_protocol.core import PriceType
from offering_protocol.core import OfferingSearchRequest, PriceType


async def run(service_url: str) -> None:
Expand All @@ -31,10 +31,24 @@ async def run(service_url: str) -> None:
raise RuntimeError("download Action did not resolve to the Node.js Service")


async def search(service_url: str) -> None:
async with ServiceClient(service_url, allow_local_network=True) as client:
page = await client.search_offerings(OfferingSearchRequest(query="gpu", limit=2))
if [item.id for item in page.items] != ["gpu-00000000", "gpu-00000001"]:
raise RuntimeError("Offering search did not match the Node.js marketplace")
if not page.next:
raise RuntimeError("Offering search did not return a continuation")
next_page = await client.continue_offerings(page.next)
if [item.id for item in next_page.items] != ["gpu-00000002", "gpu-00000003"]:
raise RuntimeError("Offering search continuation did not retain its query and limit")
if (await client.search_offerings(OfferingSearchRequest(query="plants"))).items:
raise RuntimeError("Offering search ignored the query")


def main() -> None:
if len(sys.argv) != 2:
raise SystemExit("usage: node_interoperability.py SERVICE_URL")
asyncio.run(run(sys.argv[1]))
if len(sys.argv) != 3 or sys.argv[2] not in {"catalog", "search"}:
raise SystemExit("usage: node_interoperability.py SERVICE_URL catalog|search")
asyncio.run((run if sys.argv[2] == "catalog" else search)(sys.argv[1]))
print("Python Agent interoperates with the Node.js example Service")


Expand Down
34 changes: 21 additions & 13 deletions scripts/run-node-interoperability.sh
Original file line number Diff line number Diff line change
Expand Up @@ -8,18 +8,26 @@ log_file=${TMPDIR:-/tmp}/odp-node-interop.log

package_manager=$(node -p 'require(require("node:path").resolve(process.argv[1])).packageManager' "$node_dir/package.json")
corepack "$package_manager" --dir "$node_dir" build
HOST=127.0.0.1 PORT="$port" node "$node_dir/examples/odp-service-small/dist/index.js" >"$log_file" 2>&1 &
service_pid=$!
trap 'kill "$service_pid" 2>/dev/null || true' EXIT INT TERM
run_example() {
HOST=127.0.0.1 PORT="$port" node "$node_dir/examples/$1/dist/index.js" >"$log_file" 2>&1 &
service_pid=$!
trap 'kill "$service_pid" 2>/dev/null || true' EXIT INT TERM

attempt=0
until curl --fail --silent --output /dev/null "$service_url/.well-known/odp"; do
attempt=$((attempt + 1))
if [ "$attempt" -ge 50 ]; then
sed -n '1,120p' "$log_file" >&2
exit 1
fi
sleep 0.1
done
attempt=0
until curl --fail --silent --output /dev/null "$service_url/.well-known/odp"; do
attempt=$((attempt + 1))
if [ "$attempt" -ge 50 ]; then
sed -n '1,120p' "$log_file" >&2
exit 1
fi
sleep 0.1
done

uv run python scripts/node_interoperability.py "$service_url"
uv run python scripts/node_interoperability.py "$service_url" "$2"
kill "$service_pid"
wait "$service_pid" 2>/dev/null || true
trap - EXIT INT TERM
}

run_example odp-service-small catalog
run_example odp-service-marketplace search
Loading
Loading