diff --git a/agentex/database/migrations/alembic/versions/2026_09_03_1200_add_agents_registration_metadata_gin_index_b7d3e1f4a2c6.py b/agentex/database/migrations/alembic/versions/2026_09_03_1200_add_agents_registration_metadata_gin_index_b7d3e1f4a2c6.py new file mode 100644 index 00000000..cea26efa --- /dev/null +++ b/agentex/database/migrations/alembic/versions/2026_09_03_1200_add_agents_registration_metadata_gin_index_b7d3e1f4a2c6.py @@ -0,0 +1,48 @@ +"""add_agents_registration_metadata_gin_index + +Revision ID: b7d3e1f4a2c6 +Revises: c4e8b2a7f91d +Create Date: 2026-09-03 12:00:00.000000 + +Supports the ``agent_card_metadata`` filter on ``GET /agents``, which applies a +JSONB containment predicate (``registration_metadata @> ...``). Discovery +clients poll that filtered endpoint, while agent registration writes are +comparatively rare. A ``jsonb_path_ops`` GIN index on the full +``registration_metadata`` column serves the containment operator directly, so +the polling read path does not degrade into a sequential scan as the agent +registry grows. + +Safety: +- Index built with CREATE INDEX CONCURRENTLY inside an autocommit_block, so no + long write lock is taken on ``agents``. +- IF NOT EXISTS on both upgrade and downgrade makes re-runs a no-op. +- ``jsonb_path_ops`` matches the query's ``@>`` operator and is smaller and + faster for it than the default GIN opclass; the index intentionally covers + the whole column rather than an expression, mirroring the existing + ``ix_tasks_metadata_gin`` index and the exact predicate the repository emits. +""" + +from collections.abc import Sequence + +from alembic import op + +# revision identifiers, used by Alembic. +revision: str = "b7d3e1f4a2c6" +down_revision: str | None = "c4e8b2a7f91d" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + +_INDEX = "ix_agents_registration_metadata_gin" + + +def upgrade() -> None: + with op.get_context().autocommit_block(): + op.execute( + f"CREATE INDEX CONCURRENTLY IF NOT EXISTS {_INDEX} " + "ON agents USING GIN (registration_metadata jsonb_path_ops)" + ) + + +def downgrade() -> None: + with op.get_context().autocommit_block(): + op.execute(f"DROP INDEX CONCURRENTLY IF EXISTS {_INDEX}") diff --git a/agentex/openapi.yaml b/agentex/openapi.yaml index c73d2544..9611bb7d 100644 --- a/agentex/openapi.yaml +++ b/agentex/openapi.yaml @@ -166,6 +166,18 @@ paths: default: desc title: Order Direction description: Order direction (asc or desc) + - name: agent_card_metadata + in: query + required: false + schema: + anyOf: + - type: string + - type: 'null' + description: 'JSON-encoded object used to filter agents on `registration_metadata.agent_card.metadata` + via JSONB containment. Example: {"permits_capable": true}.' + title: Agent Card Metadata + description: 'JSON-encoded object used to filter agents on `registration_metadata.agent_card.metadata` + via JSONB containment. Example: {"permits_capable": true}.' responses: '200': description: Successful Response diff --git a/agentex/src/api/routes/agents.py b/agentex/src/api/routes/agents.py index 8d4c3f97..bcca6b1d 100644 --- a/agentex/src/api/routes/agents.py +++ b/agentex/src/api/routes/agents.py @@ -1,5 +1,8 @@ +import json +import math import secrets from collections.abc import AsyncIterator +from typing import Annotated from fastapi import APIRouter, HTTPException, Query, Request from fastapi.responses import StreamingResponse @@ -107,6 +110,57 @@ async def get_agent_by_name( return Agent.model_validate(agent_entity) +_AGENT_CARD_METADATA_DESCRIPTION = ( + "JSON-encoded object used to filter agents on " + "`registration_metadata.agent_card.metadata` via JSONB containment. " + 'Example: {"permits_capable": true}.' +) + + +def _reject_json_constant(name: str) -> float: + """Reject the non-standard ``NaN``/``Infinity`` literals ``json`` accepts.""" + raise ValueError(f"{name} is not a valid JSON number") + + +def _parse_finite_float(raw: str) -> float: + """Reject float literals whose magnitude overflows to infinity (e.g. ``1e1000000``).""" + value = float(raw) + if not math.isfinite(value): + raise ValueError(f"{raw} is out of range for a JSON number") + return value + + +def _parse_agent_card_metadata(raw: str) -> dict: + """Decode the JSON-encoded ``agent_card_metadata`` query value into a dict. + + Python's ``json`` module accepts values that aren't interoperable JSON -- + the bare ``NaN``/``Infinity`` constants, float literals that overflow to + infinity, and integers too large to render. Those all satisfy an + ``isinstance(..., dict)`` check but blow up further down at the JSONB bind + parameter, turning caller error into an uncontrolled 500. Reject them here + so every malformed input surfaces as a 400. + """ + try: + parsed = json.loads( + raw, + parse_constant=_reject_json_constant, + parse_float=_parse_finite_float, + ) + except ValueError as exc: + # json.JSONDecodeError subclasses ValueError, as do the hook rejections + # above and CPython's integer-string conversion limit. + raise HTTPException( + status_code=400, + detail=f"agent_card_metadata is not valid JSON: {exc}", + ) from exc + if not isinstance(parsed, dict): + raise HTTPException( + status_code=400, + detail="agent_card_metadata must be a JSON object", + ) + return parsed + + @router.get( "", response_model=list[Agent], @@ -121,14 +175,23 @@ async def list_agents( page_number: int = Query(1, description="Page number", ge=1), order_by: str | None = Query(None, description="Field to order by"), order_direction: str = Query("desc", description="Order direction (asc or desc)"), + agent_card_metadata: Annotated[ + str | None, + Query(description=_AGENT_CARD_METADATA_DESCRIPTION), + ] = None, ): """List all registered agents.""" + agent_card_metadata_filter: dict | None = None + if agent_card_metadata is not None: + agent_card_metadata_filter = _parse_agent_card_metadata(agent_card_metadata) + agent_entities = await agents_use_case.list( task_id=task_id, limit=limit, page_number=page_number, order_by=order_by, order_direction=order_direction, + agent_card_metadata=agent_card_metadata_filter, **{"id": _authorized_ids} if _authorized_ids is not None else {}, ) return [Agent.model_validate(agent_entity) for agent_entity in agent_entities] diff --git a/agentex/src/domain/repositories/agent_repository.py b/agentex/src/domain/repositories/agent_repository.py index 04c44726..6749ec5a 100644 --- a/agentex/src/domain/repositories/agent_repository.py +++ b/agentex/src/domain/repositories/agent_repository.py @@ -46,14 +46,36 @@ async def list( Args: filters: Dictionary of filters to apply. Currently supports: - task_id: Filter agents by task ID using the join table + - agent_card_metadata: Dict applied as an exact JSONB + containment filter (``@>``) against + ``registration_metadata['agent_card']['metadata']``. order_by: Field to order by order_direction: Direction to order by (asc or desc) """ query = select(AgentORM) - if filters and "task_id" in filters: + # Pop out non-column filters that the base repository can't map to a + # single equality column, so its create_where_clauses_from_filters call + # doesn't see them. + filters = dict(filters) if filters else {} + task_id = filters.pop("task_id", None) + agent_card_metadata = filters.pop("agent_card_metadata", None) + + if task_id is not None: query = query.join( TaskAgentORM, AgentORM.id == TaskAgentORM.agent_id - ).where(TaskAgentORM.task_id == filters["task_id"]) + ).where(TaskAgentORM.task_id == task_id) + if agent_card_metadata is not None: + # Top-level JSONB `@>` with the caller's dict wrapped under the same + # nested shape it will occupy in the stored registration_metadata. + # `@>` matches when every key/value in the right operand exists at + # the same path in the left, so agents whose registration_metadata + # is NULL, missing `agent_card`, or missing `agent_card.metadata` + # are naturally excluded. + query = query.where( + AgentORM.registration_metadata.contains( + {"agent_card": {"metadata": agent_card_metadata}} + ) + ) query = query.where(AgentORM.status != AgentStatus.DELETED) return await super().list( filters=filters, diff --git a/agentex/src/domain/use_cases/agents_use_case.py b/agentex/src/domain/use_cases/agents_use_case.py index a4602cbe..fefb6424 100644 --- a/agentex/src/domain/use_cases/agents_use_case.py +++ b/agentex/src/domain/use_cases/agents_use_case.py @@ -450,10 +450,16 @@ async def list( task_id: str | None = None, order_by: str | None = None, order_direction: str = "desc", + agent_card_metadata: dict[str, Any] | None = None, **filters, ) -> list[AgentEntity]: if task_id is not None: filters["task_id"] = task_id + if agent_card_metadata is not None: + # Reserved key consumed by the repository to apply a JSONB containment + # filter on `registration_metadata.agent_card.metadata`. Kept out of the + # generic column-equality path in `create_where_clauses_from_filters`. + filters["agent_card_metadata"] = agent_card_metadata return await self.agent_repo.list( filters=filters, diff --git a/agentex/tests/integration/api/agents/test_agents_api.py b/agentex/tests/integration/api/agents/test_agents_api.py index dd9699b2..d00cce6d 100644 --- a/agentex/tests/integration/api/agents/test_agents_api.py +++ b/agentex/tests/integration/api/agents/test_agents_api.py @@ -652,3 +652,332 @@ async def test_register_agent_duplicate_name_behavior(self, isolated_client): agents_with_name = [a for a in agents if a["name"] == "duplicate-name-test"] assert len(agents_with_name) == 1 assert agents_with_name[0]["description"] == "Second agent" + + @pytest.mark.asyncio + async def test_list_agents_filters_by_agent_card_metadata(self, isolated_client): + """`agent_card_metadata` returns only agents whose card metadata contains + the requested key/value pairs; agents without card metadata are excluded.""" + # Given - three agents: one Permits-capable, one with different metadata, + # and one with no agent_card at all. + await isolated_client.post( + "/agents/register", + json={ + "name": "card-metadata-permits", + "description": "opts into Permits", + "acp_url": "http://permits-agent:8000", + "acp_type": "sync", + "registration_metadata": { + "agent_card": { + "metadata": {"permits_capable": True, "region": "us"} + } + }, + }, + ) + await isolated_client.post( + "/agents/register", + json={ + "name": "card-metadata-other", + "description": "different capability", + "acp_url": "http://other-agent:8000", + "acp_type": "sync", + "registration_metadata": { + "agent_card": {"metadata": {"other_feature": True}} + }, + }, + ) + await isolated_client.post( + "/agents/register", + json={ + "name": "card-metadata-none", + "description": "no card metadata", + "acp_url": "http://plain-agent:8000", + "acp_type": "sync", + }, + ) + + # When - filter by exact key/value present on only one agent + response = await isolated_client.get( + '/agents?agent_card_metadata={"permits_capable":true}' + ) + assert response.status_code == 200 + agents = response.json() + names = {a["name"] for a in agents} + assert names == {"card-metadata-permits"} + + # And - non-matching value returns no agents (agent exists but with a + # different value for the same key does not match) + response = await isolated_client.get( + '/agents?agent_card_metadata={"permits_capable":false}' + ) + assert response.status_code == 200 + assert response.json() == [] + + # And - omitting the filter returns all non-deleted agents, including + # those without any agent_card metadata + response = await isolated_client.get("/agents") + assert response.status_code == 200 + names_unfiltered = {a["name"] for a in response.json()} + assert { + "card-metadata-permits", + "card-metadata-other", + "card-metadata-none", + } <= names_unfiltered + + @pytest.mark.asyncio + async def test_list_agents_agent_card_metadata_multi_key_containment( + self, isolated_client + ): + """Multi-key filter requires every key/value to be present (JSONB `@>`).""" + await isolated_client.post( + "/agents/register", + json={ + "name": "card-metadata-multi", + "description": "multi", + "acp_url": "http://multi-agent:8000", + "acp_type": "sync", + "registration_metadata": { + "agent_card": { + "metadata": { + "permits_capable": True, + "region": "us", + "extra": "value", + } + } + }, + }, + ) + + # All requested keys match -> included + response = await isolated_client.get( + '/agents?agent_card_metadata={"permits_capable":true,"region":"us"}' + ) + assert response.status_code == 200 + assert {a["name"] for a in response.json()} == {"card-metadata-multi"} + + # One requested key doesn't match -> excluded + response = await isolated_client.get( + '/agents?agent_card_metadata={"permits_capable":true,"region":"eu"}' + ) + assert response.status_code == 200 + assert response.json() == [] + + @pytest.mark.asyncio + async def test_list_agents_agent_card_metadata_combined_with_pagination( + self, isolated_client + ): + """The filter composes with existing pagination and ordering behavior.""" + for i in range(3): + await isolated_client.post( + "/agents/register", + json={ + "name": f"card-metadata-page-{i}", + "description": f"agent {i}", + "acp_url": f"http://page-agent-{i}:8000", + "acp_type": "sync", + "registration_metadata": { + "agent_card": {"metadata": {"permits_capable": True}} + }, + }, + ) + # Unrelated agent that must not leak into filtered results + await isolated_client.post( + "/agents/register", + json={ + "name": "card-metadata-page-noise", + "description": "noise", + "acp_url": "http://noise:8000", + "acp_type": "sync", + }, + ) + + response = await isolated_client.get( + '/agents?agent_card_metadata={"permits_capable":true}&limit=2&page_number=1' + ) + assert response.status_code == 200 + page_one = response.json() + assert len(page_one) == 2 + assert all(a["name"].startswith("card-metadata-page-") for a in page_one) + assert not any(a["name"] == "card-metadata-page-noise" for a in page_one) + + @pytest.mark.asyncio + async def test_list_agents_agent_card_metadata_invalid_json_returns_400( + self, isolated_client + ): + """Malformed JSON in `agent_card_metadata` is rejected up front.""" + response = await isolated_client.get("/agents?agent_card_metadata=not-json") + assert response.status_code == 400 + + response = await isolated_client.get("/agents?agent_card_metadata=[1,2,3]") + assert response.status_code == 400 + + @pytest.mark.parametrize( + "raw_filter", + [ + '{"x": NaN}', + '{"x": Infinity}', + '{"x": -Infinity}', + '{"x": 1e1000000}', + '{"x": [1, NaN]}', + '{"x": {"nested": Infinity}}', + '{"x": ' + "9" * 5000 + "}", + ], + ) + @pytest.mark.asyncio + async def test_list_agents_agent_card_metadata_non_finite_numbers_return_400( + self, isolated_client, raw_filter + ): + """Values Python's json accepts but JSON doesn't are rejected as 400, not + passed through to the JSONB bind parameter where they'd surface as a 500.""" + response = await isolated_client.get( + "/agents", params={"agent_card_metadata": raw_filter} + ) + assert response.status_code == 400 + assert "agent_card_metadata" in response.json()["message"] + + @pytest.mark.asyncio + async def test_list_agents_agent_card_metadata_empty_object_requires_metadata( + self, isolated_client + ): + """An explicit `{}` filter still applies the containment predicate: agents + must have a card metadata object, but any contents match.""" + await isolated_client.post( + "/agents/register", + json={ + "name": "card-metadata-empty-with", + "description": "has card metadata", + "acp_url": "http://with-agent:8000", + "acp_type": "sync", + "registration_metadata": { + "agent_card": {"metadata": {"anything": "at-all"}} + }, + }, + ) + await isolated_client.post( + "/agents/register", + json={ + "name": "card-metadata-empty-without", + "description": "no card metadata", + "acp_url": "http://without-agent:8000", + "acp_type": "sync", + }, + ) + + response = await isolated_client.get("/agents?agent_card_metadata={}") + assert response.status_code == 200 + names = {a["name"] for a in response.json()} + assert "card-metadata-empty-with" in names + assert "card-metadata-empty-without" not in names + + @pytest.mark.asyncio + async def test_reregistration_replaces_agent_card_and_filter_results( + self, isolated_client + ): + """Re-registering the same agent identity replaces the stored top-level + `agent_card`, and the metadata filter reflects the new card immediately. + This locks the rolling-update contract a polling discovery client relies + on: after a release re-registers with card B, filter A stops matching + and filter B starts matching.""" + response = await isolated_client.post( + "/agents/register", + json={ + "name": "card-metadata-rollout", + "description": "release 1", + "acp_url": "http://rollout-agent:8000", + "acp_type": "sync", + "registration_metadata": { + "agent_card": {"metadata": {"permits_capable": True, "rev": "a"}} + }, + }, + ) + assert response.status_code == 200 + + response = await isolated_client.get('/agents?agent_card_metadata={"rev":"a"}') + assert response.status_code == 200 + assert {a["name"] for a in response.json()} == {"card-metadata-rollout"} + + response = await isolated_client.post( + "/agents/register", + json={ + "name": "card-metadata-rollout", + "description": "release 2", + "acp_url": "http://rollout-agent:8000", + "acp_type": "sync", + "registration_metadata": { + "agent_card": {"metadata": {"permits_capable": True, "rev": "b"}} + }, + }, + ) + assert response.status_code == 200 + + response = await isolated_client.get('/agents?agent_card_metadata={"rev":"a"}') + assert response.status_code == 200 + assert response.json() == [] + + response = await isolated_client.get('/agents?agent_card_metadata={"rev":"b"}') + assert response.status_code == 200 + assert {a["name"] for a in response.json()} == {"card-metadata-rollout"} + + @pytest.mark.asyncio + async def test_reregistration_with_null_agent_card_withdraws_from_discovery( + self, isolated_client + ): + """An explicit `{"agent_card": null}` registration clears a previously + published card, so rolling back to a release that publishes no card can + withdraw the stale descriptor without deleting the agent. Omitting + `registration_metadata` entirely preserves the existing card.""" + response = await isolated_client.post( + "/agents/register", + json={ + "name": "card-metadata-withdraw", + "description": "publishes a card", + "acp_url": "http://withdraw-agent:8000", + "acp_type": "sync", + "registration_metadata": { + "agent_card": {"metadata": {"permits_capable": True}} + }, + }, + ) + assert response.status_code == 200 + + response = await isolated_client.post( + "/agents/register", + json={ + "name": "card-metadata-withdraw", + "description": "re-registers without touching metadata", + "acp_url": "http://withdraw-agent:8000", + "acp_type": "sync", + }, + ) + assert response.status_code == 200 + + response = await isolated_client.get( + '/agents?agent_card_metadata={"permits_capable":true}' + ) + assert response.status_code == 200 + assert {a["name"] for a in response.json()} == {"card-metadata-withdraw"} + + response = await isolated_client.post( + "/agents/register", + json={ + "name": "card-metadata-withdraw", + "description": "rolled back, withdraws the card", + "acp_url": "http://withdraw-agent:8000", + "acp_type": "sync", + "registration_metadata": {"agent_card": None}, + }, + ) + assert response.status_code == 200 + + response = await isolated_client.get( + '/agents?agent_card_metadata={"permits_capable":true}' + ) + assert response.status_code == 200 + assert response.json() == [] + + response = await isolated_client.get("/agents?agent_card_metadata={}") + assert response.status_code == 200 + assert "card-metadata-withdraw" not in {a["name"] for a in response.json()} + + response = await isolated_client.get("/agents") + assert response.status_code == 200 + assert "card-metadata-withdraw" in {a["name"] for a in response.json()} diff --git a/agentex/tests/unit/api/test_agents_authz.py b/agentex/tests/unit/api/test_agents_authz.py index 581d79d5..6a497a7b 100644 --- a/agentex/tests/unit/api/test_agents_authz.py +++ b/agentex/tests/unit/api/test_agents_authz.py @@ -325,6 +325,7 @@ async def test_authorized_ids_pushed_into_use_case(self): page_number=1, order_by=None, order_direction="desc", + agent_card_metadata=None, id=["agent-a", "agent-c"], ) @@ -351,6 +352,7 @@ async def test_none_authorized_ids_passes_through_unfiltered(self): page_number=1, order_by=None, order_direction="desc", + agent_card_metadata=None, ) diff --git a/agentex/tests/unit/use_cases/test_agents_use_case.py b/agentex/tests/unit/use_cases/test_agents_use_case.py index 129624e9..7a063d54 100644 --- a/agentex/tests/unit/use_cases/test_agents_use_case.py +++ b/agentex/tests/unit/use_cases/test_agents_use_case.py @@ -249,3 +249,115 @@ async def test_deployment_aware_register_uses_deployment_acp_url_for_healthcheck assert temporal_adapter.start_workflow.await_args.kwargs["args"] == [ {"agent_id": build_only_agent.id, "acp_url": deployment_acp_url} ] + + +async def _seed_agent_with_card_metadata( + agent_repository: AgentRepository, + name: str, + card_metadata: dict | None, +) -> AgentEntity: + registration_metadata: dict = {} + if card_metadata is not None: + registration_metadata["agent_card"] = {"metadata": card_metadata} + return await agent_repository.create( + AgentEntity( + id=str(uuid4()), + name=name, + description="seed", + status=AgentStatus.READY, + acp_type=ACPType.ASYNC, + acp_url="http://seed.example.com", + registration_metadata=registration_metadata or None, + ) + ) + + +@pytest.mark.asyncio +@pytest.mark.unit +async def test_list_filters_by_agent_card_metadata_exact_containment( + agents_use_case, agent_repository +): + """Filter matches only agents whose card metadata contains every requested pair. + + Uses a per-test-unique tag so assertions are robust to data seeded by other + tests sharing the same session-scoped Postgres container. + """ + suffix = uuid4().hex[:6] + tag = f"test-{suffix}" + permits = await _seed_agent_with_card_metadata( + agent_repository, + f"card-metadata-permits-{suffix}", + {"permits_capable": True, "region": "us", "test_tag": tag}, + ) + other = await _seed_agent_with_card_metadata( + agent_repository, + f"card-metadata-other-{suffix}", + {"other_feature": True, "test_tag": tag}, + ) + plain = await _seed_agent_with_card_metadata( + agent_repository, + f"card-metadata-none-{suffix}", + None, + ) + scoped_ids = {permits.id, other.id, plain.id} + + # Single-key exact filter, scoped to this test's tag → only permits agent. + matches = await agents_use_case.list( + limit=50, + page_number=1, + agent_card_metadata={"permits_capable": True, "test_tag": tag}, + ) + assert {a.id for a in matches if a.id in scoped_ids} == {permits.id} + + # Non-matching value under our tag returns no agents seeded by this test. + no_matches = await agents_use_case.list( + limit=50, + page_number=1, + agent_card_metadata={"permits_capable": False, "test_tag": tag}, + ) + assert {a.id for a in no_matches if a.id in scoped_ids} == set() + + # Multi-key containment requires every key/value to be present. + multi = await agents_use_case.list( + limit=50, + page_number=1, + agent_card_metadata={ + "permits_capable": True, + "region": "us", + "test_tag": tag, + }, + ) + assert {a.id for a in multi if a.id in scoped_ids} == {permits.id} + missing_key = await agents_use_case.list( + limit=50, + page_number=1, + agent_card_metadata={ + "permits_capable": True, + "region": "eu", + "test_tag": tag, + }, + ) + assert {a.id for a in missing_key if a.id in scoped_ids} == set() + + +@pytest.mark.asyncio +@pytest.mark.unit +async def test_list_without_agent_card_metadata_returns_all_non_deleted( + agents_use_case, agent_repository +): + """Omitting the filter preserves existing behavior (agents without card metadata still listed).""" + suffix = uuid4().hex[:6] + with_card = await _seed_agent_with_card_metadata( + agent_repository, + f"card-metadata-with-{suffix}", + {"permits_capable": True}, + ) + without_card = await _seed_agent_with_card_metadata( + agent_repository, + f"card-metadata-without-{suffix}", + None, + ) + + all_agents = await agents_use_case.list(limit=50, page_number=1) + all_ids = {a.id for a in all_agents} + assert {with_card.id, without_card.id} <= all_ids