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
161 changes: 161 additions & 0 deletions graph_os/fleet/access_contracts.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,161 @@
"""Register connector access contracts as unapproved virtual mappings.

A fleet connector's ``ontology://`` resources may declare access contracts.
Each contract binds one ontology class to one live source operation.
Onboarding registers one ``SourceConnection`` per source and one unapproved
``VirtualMapping`` per contract. The virtual catalog serves a mapping only
after approval, so registration exposes no data.

The SDK parser ships in ``agent_connector_sdk.access_contract``. Without that
module, onboarding skips contracts and logs at debug level.
"""

from __future__ import annotations

import logging
from collections.abc import Iterable, Mapping, Sequence
from typing import Any

logger = logging.getLogger(__name__)

#: SDK access kind to AU virtual-graph source kind.
SOURCE_KIND_BY_ACCESS = {
"mcp_tool": "mcp",
"http_endpoint": "api",
"graphql_query": "graphql",
"a2a_skill": "a2a",
}
_CATALOG: list[Any] = []


def fleet_virtual_catalog() -> Any:
"""The process virtual catalog that fleet onboarding registers into."""

if not _CATALOG:
from agent_utilities.knowledge_graph.virtual_graph import VirtualCatalog

_CATALOG.append(VirtualCatalog())
return _CATALOG[0]


def _text(body: object) -> str:
return body.decode("utf-8") if isinstance(body, bytes) else str(body)


def _entries(pack: Any) -> Iterable[Any]:
entries = getattr(pack, "entries", None)
if entries is None:
entries = getattr(getattr(pack, "archive", None), "entries", ())
found = entries() if callable(entries) else entries
return tuple(found or ())


def pack_ontologies(pack: Any) -> tuple[str, ...]:
"""The body of every ``ontology://`` entry in a captured pack."""

return tuple(
_text(getattr(entry, "body", ""))
for entry in _entries(pack)
if str(getattr(entry, "uri", "")).startswith("ontology://")
)


def _parser() -> Any:
try:
from agent_connector_sdk.access_contract import parse_access_contracts
except ImportError:
logger.debug("connector SDK has no access_contract module; skipping")
return None
return parse_access_contracts


def parse_pack_contracts(pack: Any) -> tuple[Any, ...]:
"""Every access contract that a pack's ontologies declare."""

parse = _parser()
if parse is None:
return ()
return tuple(c for text in pack_ontologies(pack) for c in parse(text))


def _by_source(contracts: Sequence[Any]) -> dict[str, list[Any]]:
grouped: dict[str, list[Any]] = {}
for contract in contracts:
grouped.setdefault(contract.source_id, []).append(contract)
return grouped


def _entity(contract: Any) -> Any:
from agent_utilities.knowledge_graph.virtual_graph import DiscoveredEntity

return DiscoveredEntity(
name=contract.entity,
key=contract.key_field,
fields=tuple(sorted({f for _, f in contract.predicates})),
operation=contract.operation,
)


def _unwired_call(operation: str, arguments: Mapping[str, Any]) -> Any:
raise LookupError(f"no live binding for {operation!r} yet")


def _register_source(
catalog: Any, connector: str, source_id: str, contracts: Sequence[Any]
) -> None:
from agent_utilities.knowledge_graph.virtual_graph import (
MetadataContract,
OperationAdapter,
SourceConnection,
)

connection = SourceConnection(
source_id=source_id,
kind=SOURCE_KIND_BY_ACCESS[contracts[0].access_kind],
endpoint_ref=f"fleet://{connector}",
)
entities = {c.entity: _entity(c) for c in contracts}
metadata = MetadataContract(
source_id=source_id,
schema_version="1",
entities=tuple(entities.values()),
)
adapter = OperationAdapter(connection, metadata, _unwired_call)
catalog.register(connection, metadata, adapter)


def _add_mappings(catalog: Any, contracts: Sequence[Any]) -> int:
from agent_utilities.knowledge_graph.virtual_graph import VirtualMapping

known = {m.mapping_id for m in catalog.mappings}
added = 0
for contract in contracts:
mapping = VirtualMapping(**contract.virtual_mapping_fields())
if mapping.mapping_id in known:
continue
catalog.add_mapping(mapping)
added += 1
return added


def register_access_contracts(pack: Any, *, connector: str, catalog: Any = None) -> int:
"""Register a pack's access contracts unapproved; return mappings added."""

contracts = parse_pack_contracts(pack)
if not contracts:
return 0
target = catalog if catalog is not None else fleet_virtual_catalog()
added = 0
for source_id, items in _by_source(contracts).items():
_register_source(target, connector, source_id, items)
added += _add_mappings(target, items)
return added


__all__ = [
"SOURCE_KIND_BY_ACCESS",
"fleet_virtual_catalog",
"pack_ontologies",
"parse_pack_contracts",
"register_access_contracts",
]
17 changes: 16 additions & 1 deletion graph_os/fleet/onboarding.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,8 @@
1. register the server, and renew any lease that lapses soon;
2. capture its served tools, prompts, skills and pack resources as a
ConnectorPack through the connector SDK;
3. attest that catalog with EG and import the pack under EG's binding.
3. attest that catalog with EG and import the pack under EG's binding;
4. register its access contracts as unapproved virtual mappings.

GraphOS observes each child catalog over MCP, as the multiplexer does.
EG issues the binding through ``attest_self_served_catalog``; GraphOS must be
Expand Down Expand Up @@ -200,6 +201,8 @@ class FleetOnboarding:
session: Any
#: Replaces :func:`capture_server_pack`, for example in tests.
capture: Any = None
#: The virtual catalog that access contracts register into.
virtual_catalog: Any = None

@property
def tenant_client(self) -> Any:
Expand Down Expand Up @@ -244,6 +247,18 @@ async def onboard(self, endpoint: FleetEndpoint, *, auth: Any) -> None:
commons=self.commons_client,
pack=pack,
)
self.register_contracts(endpoint, pack)

def register_contracts(self, endpoint: FleetEndpoint, pack: Any) -> None:
"""Register the pack's access contracts as unapproved mappings."""

from graph_os.fleet.access_contracts import register_access_contracts

added = register_access_contracts(
pack, connector=endpoint.name, catalog=self.virtual_catalog
)
if added:
logger.info("%s: %d unapproved access mappings", endpoint.name, added)

async def self_endpoints(
self, connectors: Sequence[str], served_url: str
Expand Down
1 change: 1 addition & 0 deletions specs/fleet-catalog-and-tools/requirements.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,3 +31,4 @@
| `GRAPHOS-FLEET-R026` | **The catalog admits each server on its own.** The catalog reader skips a live registration without a current `mcp_server` component. It logs a warning and records the server name. `multiplexer_status` reports the names as `unadmitted_servers`. Boot never fails for this state. Receipt, revision, pin, content, duplicate and page-bound faults stay fatal. | `tests/test_fleet_catalog_reader.py` proves the reader skips and reports a bare registration. The existing fail-closed cases still raise. A multiplexer test reads `unadmitted_servers` from status. |
| `GRAPHOS-FLEET-R027` | **GraphOS onboards every configured fleet server.** One idempotent pass reads the enabled streamable-HTTP servers from the MCP config. The pass registers each server and captures its catalog through the connector SDK. It attests the catalog with EG and imports the pack. One server's failure leaves the others unaffected. `graph-os-production-ops onboard-fleet` runs a full pass and exits non-zero on any failure. Boot runs the pass on a background thread. A pass fault never blocks or stops serving. | `tests/fleet/test_fleet_onboarding.py` covers config filtering, a fresh-store pass, a failing server, an admitted server, a registry outage, a missing child credential, the CLI and a failing boot pass. |
| `GRAPHOS-FLEET-R028` | **GraphOS renews every registration before its lease lapses.** Each pass re-registers each fleet and self-served lease that expires within six hours. Each renewal window uses its own idempotency key. A self-served registration renews at the configured served URL, else at its live URL. The boot thread repeats the pass every 30 minutes. A pass that admits a new server refreshes the multiplexer catalog when no child session runs. | `tests/fleet/test_fleet_onboarding.py` covers lapsing, fresh and absent leases, windowed keys, self-served renewal, the pass schedule and the refresh rule. |
| `GRAPHOS-FLEET-R030` | **Fleet onboarding registers connector access contracts unapproved.** After a pack import, onboarding parses each `ontology://` resource with the SDK access-contract parser. It registers one `SourceConnection` per source and one unapproved `VirtualMapping` per contract. An unapproved mapping serves no reads. A missing SDK parser skips contracts without failure. A repeated pass adds no duplicate mapping. | `tests/fleet/test_fleet_access_contracts.py` covers registration, the unapproved state, idempotency and the missing parser. |
4 changes: 2 additions & 2 deletions specs/fleet-catalog-and-tools/spec.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
**Owner:** graph-os

**State:** READY FOR IMPLEMENTATION — architecture and acceptance specified; no claim that the full surface is deployed or accepted.
**Scope IDs:** GRAPHOS-FLEET-R001, GRAPHOS-FLEET-R002, GRAPHOS-FLEET-R003, GRAPHOS-FLEET-R004, GRAPHOS-FLEET-R005, GRAPHOS-FLEET-R006, GRAPHOS-FLEET-R007, GRAPHOS-FLEET-R008, GRAPHOS-FLEET-R009, GRAPHOS-FLEET-R010, GRAPHOS-FLEET-R011, GRAPHOS-FLEET-R012, GRAPHOS-FLEET-R013, GRAPHOS-FLEET-R014, GRAPHOS-FLEET-R015, GRAPHOS-FLEET-R016, GRAPHOS-FLEET-R017, GRAPHOS-FLEET-R018, GRAPHOS-FLEET-R019, GRAPHOS-FLEET-R020, GRAPHOS-FLEET-R021, GRAPHOS-FLEET-R022, PA-12, GRAPHOS-FLEET-R026, GRAPHOS-FLEET-R027, GRAPHOS-FLEET-R028.
**Scope IDs:** GRAPHOS-FLEET-R001, GRAPHOS-FLEET-R002, GRAPHOS-FLEET-R003, GRAPHOS-FLEET-R004, GRAPHOS-FLEET-R005, GRAPHOS-FLEET-R006, GRAPHOS-FLEET-R007, GRAPHOS-FLEET-R008, GRAPHOS-FLEET-R009, GRAPHOS-FLEET-R010, GRAPHOS-FLEET-R011, GRAPHOS-FLEET-R012, GRAPHOS-FLEET-R013, GRAPHOS-FLEET-R014, GRAPHOS-FLEET-R015, GRAPHOS-FLEET-R016, GRAPHOS-FLEET-R017, GRAPHOS-FLEET-R018, GRAPHOS-FLEET-R019, GRAPHOS-FLEET-R020, GRAPHOS-FLEET-R021, GRAPHOS-FLEET-R022, PA-12, GRAPHOS-FLEET-R026, GRAPHOS-FLEET-R027, GRAPHOS-FLEET-R028, GRAPHOS-FLEET-R030.

## Outcome and state legend

Expand Down Expand Up @@ -47,7 +47,7 @@ Status is per deliverable; a source commit, a green unit test, or a prior status
| Fleet invocation and safety | GRAPHOS-FLEET-R015, GRAPHOS-FLEET-R018, GRAPHOS-FLEET-R019, GRAPHOS-FLEET-R020, GRAPHOS-FLEET-R021 | Native load/call, exact scopes, policy, no bypass, parity |
| Orchestration correctness | GRAPHOS-FLEET-R009, GRAPHOS-FLEET-R010 | Atomic admission, correct generated assembly agents |
| Reload and acceptance | GRAPHOS-FLEET-R022, PA-12, GRAPHOS-FLEET-R017 | Atomic generation swap, re-ingestion, negative cases, replica/served proof |
| Fleet onboarding | GRAPHOS-FLEET-R026, GRAPHOS-FLEET-R027, GRAPHOS-FLEET-R028 | Per-server admission, attested import of each child catalog, lease renewal |
| Fleet onboarding | GRAPHOS-FLEET-R026, GRAPHOS-FLEET-R027, GRAPHOS-FLEET-R028, GRAPHOS-FLEET-R030 | Per-server admission, attested import of each child catalog, lease renewal, unapproved access mappings |

The engine owns its durable records and method contract; the connector SDK owns pack and connector certification; graph-os owns serving composition, registry projection, fleet policy, session loading, and reload. A contributor may replace a missing external service with the fixtures and local fake adapters described in [test-spec.md](test-spec.md).

Expand Down
12 changes: 10 additions & 2 deletions specs/fleet-catalog-and-tools/status.json
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,8 @@
"GRAPHOS-FLEET-R025",
"GRAPHOS-FLEET-R026",
"GRAPHOS-FLEET-R027",
"GRAPHOS-FLEET-R028"
"GRAPHOS-FLEET-R028",
"GRAPHOS-FLEET-R030"
],
"delivery_state": "SPECIFIED",
"acceptance_state": "NOT_AUDITED",
Expand Down Expand Up @@ -94,7 +95,7 @@
"title": "Bridge upgraded to FastMCP 4",
"delivery_state": "SPECIFIED",
"evidence": [],
"remaining": "fastmcp>=4.0.0b1 is declared and locked at 4.0.5, and native .tool()/Skills-over-MCP are used and tested, but no native FastMCP4 prompt/resource/resource-template registration by graph-os's own served bridge was found (grepped for .prompt(/.resource(/resource_template( — no hits), so the 'single com"
"remaining": "fastmcp>=4.0.0b1 is declared and locked at 4.0.5, and native .tool()/Skills-over-MCP are used and tested, but no native FastMCP4 prompt/resource/resource-template registration by graph-os's own served bridge was found (grepped for .prompt(/.resource(/resource_template( \u2014 no hits), so the 'single com"
},
{
"id": "GRAPHOS-FLEET-R008",
Expand Down Expand Up @@ -309,6 +310,13 @@
"delivery_state": "BUILDING",
"evidence": [],
"remaining": "Implemented on branch feat/fleet-onboarding; not yet merged to the default branch."
},
{
"id": "GRAPHOS-FLEET-R030",
"title": "Fleet onboarding registers connector access contracts unapproved",
"delivery_state": "BUILDING",
"evidence": [],
"remaining": "Not yet merged. Live reads need an approval path and a live operation binding for each source."
}
]
}
3 changes: 3 additions & 0 deletions specs/fleet-catalog-and-tools/tasks.md
Original file line number Diff line number Diff line change
Expand Up @@ -20,3 +20,6 @@ Check a task only after its linked code and tests land. A checked source task is
- [x] T16 (GRAPHOS-FLEET-R027): Add `graph_os/fleet/onboarding.py`, the `onboard-fleet` command and the background boot pass.
- [x] T17 (GRAPHOS-FLEET-R028): Renew fleet and self-served leases each pass with windowed idempotency keys; refresh the catalog after admission.
- [ ] T18 (GRAPHOS-FLEET-R027): Bind GraphOS as the fleet importer in the deployment, roll out, and record the live `multiplexer_status` child count.
- [x] T18 (GRAPHOS-FLEET-R030): Register each connector access contract as an unapproved virtual mapping during onboarding.
- [ ] T19 (GRAPHOS-FLEET-R030): Read ontology entries from the real pack archive accessor; confirm the duck-typed `entries` read against the SDK archive.
- [ ] T20 (GRAPHOS-FLEET-R030): Bind each source to a live operation call and add an operator approval path.
8 changes: 8 additions & 0 deletions specs/fleet-catalog-and-tools/test-spec.md
Original file line number Diff line number Diff line change
Expand Up @@ -48,3 +48,11 @@ CI may run this probe with local containers/processes and synthetic credentials.
| Commit | Case/gate | Command or fixture | Result | Environment | Timestamp | Receipt or trace |
|---|---|---|---|---|---|---|
| pending | pending | pending | NOT RUN | clean checkout | pending | pending |

## GRAPHOS-FLEET-R030

`tests/fleet/test_fleet_access_contracts.py` uses a fake pack with one `ontology://` entry and a fake parser.
- Onboarding registers one source connection and one mapping per contract.
- Each registered mapping is unapproved, and the catalog serves no mapping for its class.
- A second registration adds no duplicate mapping.
- A missing SDK module skips contracts and raises nothing.
Loading
Loading