diff --git a/sdk/agentserver/azure-ai-agentserver-responses/CHANGELOG.md b/sdk/agentserver/azure-ai-agentserver-responses/CHANGELOG.md index 312ac5bb7990..4094322bc039 100644 --- a/sdk/agentserver/azure-ai-agentserver-responses/CHANGELOG.md +++ b/sdk/agentserver/azure-ai-agentserver-responses/CHANGELOG.md @@ -1,5 +1,11 @@ # Release History +## 2.1.0b2 (Unreleased) + +### Bugs Fixed + +- Restored JSON-string encoding for response-level `internal_metadata` so resilient response checkpoints round-trip through Foundry storage. + ## 2.1.0b1 (2026-08-11) ### Breaking Changes diff --git a/sdk/agentserver/azure-ai-agentserver-responses/azure/ai/agentserver/responses/streaming/_event_stream.py b/sdk/agentserver/azure-ai-agentserver-responses/azure/ai/agentserver/responses/streaming/_event_stream.py index 61e3656f34aa..8ed713cb2401 100644 --- a/sdk/agentserver/azure-ai-agentserver-responses/azure/ai/agentserver/responses/streaming/_event_stream.py +++ b/sdk/agentserver/azure-ai-agentserver-responses/azure/ai/agentserver/responses/streaming/_event_stream.py @@ -30,6 +30,7 @@ ) from ._state_machine import EventStreamValidator from ._checkpoint import ResponseCheckpointEvent +from ._internal_metadata import _ResponseInternalMetadataView # Event types whose payload is a full Response snapshot. # Lifecycle events nest under a "response" key on the wire. @@ -176,6 +177,7 @@ def __init__( if agent_reference is not None: self._response["agent_reference"] = deepcopy(agent_reference) + _ResponseInternalMetadataView(self._response) self._agent_reference, self._model = _internals.extract_response_fields( cast(response_models.ResponseObject, self._response) ) @@ -198,23 +200,16 @@ def internal_metadata(self) -> "MutableMapping[str, Any]": """Live, mutable response-level framework-internal metadata. A convenience proxy backed by a reserved ``_internal_metadata`` key - inside the response's public ``metadata`` map — read / write / delete in - place (``stream.internal_metadata["phase"] = 3``). Stripped from every - client-facing payload and persisted at the next ``yield - stream.checkpoint()`` (and at terminal). Values may be any - JSON-serialisable type. + inside the response's public ``metadata`` map. The private bag is + JSON-encoded into the string-valued metadata slot required by Foundry + storage. Read / write / delete in place + (``stream.internal_metadata["phase"] = 3``). Stripped from every + client-facing payload and persisted at the next + ``yield stream.checkpoint()`` (and at terminal). :rtype: ~collections.abc.MutableMapping[str, ~typing.Any] """ - metadata = self._response.get("metadata") - if not isinstance(metadata, dict): - metadata = {} - self._response["metadata"] = metadata - bag = metadata.get("_internal_metadata") - if not isinstance(bag, dict): - bag = {} - metadata["_internal_metadata"] = bag - return bag + return _ResponseInternalMetadataView(self._response) def checkpoint(self) -> "ResponseCheckpointEvent": """Return a checkpoint event to ``yield`` for persistence. diff --git a/sdk/agentserver/azure-ai-agentserver-responses/azure/ai/agentserver/responses/streaming/_internal_metadata.py b/sdk/agentserver/azure-ai-agentserver-responses/azure/ai/agentserver/responses/streaming/_internal_metadata.py new file mode 100644 index 000000000000..c016470c1e70 --- /dev/null +++ b/sdk/agentserver/azure-ai-agentserver-responses/azure/ai/agentserver/responses/streaming/_internal_metadata.py @@ -0,0 +1,91 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT license. +"""Response-level internal metadata backed by the public metadata map.""" + +from __future__ import annotations + +import json +from collections.abc import Iterator, MutableMapping +from typing import Any + +_RESERVED_KEY = "_internal_metadata" +_MAX_METADATA_KEYS = 16 +_MAX_VALUE_LEN = 512 + + +class _ResponseInternalMetadataView(MutableMapping[str, Any]): + """Live view over JSON-encoded response-level internal metadata.""" + + __slots__ = ("_response",) + + def __init__(self, response: dict[str, Any]) -> None: + self._response = response + metadata = response.get("metadata") + if isinstance(metadata, dict) and isinstance(metadata.get(_RESERVED_KEY), dict): + self._store(dict(metadata[_RESERVED_KEY])) + + def _decode(self) -> dict[str, Any]: + metadata = self._response.get("metadata") + if not isinstance(metadata, dict): + return {} + raw = metadata.get(_RESERVED_KEY) + if isinstance(raw, dict): + return dict(raw) + if not isinstance(raw, str) or not raw: + return {} + try: + decoded = json.loads(raw) + except (TypeError, ValueError): + return {} + return decoded if isinstance(decoded, dict) else {} + + def _store(self, value: dict[str, Any]) -> None: + metadata = self._response.get("metadata") + if not value: + if isinstance(metadata, dict): + metadata.pop(_RESERVED_KEY, None) + return + + encoded = json.dumps( + value, + ensure_ascii=False, + separators=(",", ":"), + sort_keys=True, + ) + if len(encoded) > _MAX_VALUE_LEN: + raise ValueError( + f"internal_metadata encodes to {len(encoded)} chars, exceeding the " + f"{_MAX_VALUE_LEN}-char limit of the response metadata value" + ) + + if not isinstance(metadata, dict): + metadata = {} + self._response["metadata"] = metadata + projected_key_count = len(metadata) + (0 if _RESERVED_KEY in metadata else 1) + if projected_key_count > _MAX_METADATA_KEYS: + raise ValueError( + f"cannot add internal_metadata: response metadata already has " + f"{len(metadata)} keys (limit {_MAX_METADATA_KEYS})" + ) + metadata[_RESERVED_KEY] = encoded + + def __getitem__(self, key: str) -> Any: + return self._decode()[key] + + def __setitem__(self, key: str, value: Any) -> None: + if not isinstance(key, str): + raise TypeError(f"internal_metadata keys must be str, got {type(key).__name__}") + metadata = self._decode() + metadata[key] = value + self._store(metadata) + + def __delitem__(self, key: str) -> None: + metadata = self._decode() + del metadata[key] + self._store(metadata) + + def __iter__(self) -> Iterator[str]: + return iter(self._decode()) + + def __len__(self) -> int: + return len(self._decode()) diff --git a/sdk/agentserver/azure-ai-agentserver-responses/docs/handler-implementation-guide.md b/sdk/agentserver/azure-ai-agentserver-responses/docs/handler-implementation-guide.md index a7c9a3d3fb48..e639bf9fe50d 100644 --- a/sdk/agentserver/azure-ai-agentserver-responses/docs/handler-implementation-guide.md +++ b/sdk/agentserver/azure-ai-agentserver-responses/docs/handler-implementation-guide.md @@ -1529,13 +1529,21 @@ stream.internal_metadata["resume_phase"] = 3 del stream.internal_metadata["scratch"] ``` +Response-level values may be any JSON-serializable type, but the complete bag +is encoded into one string-valued reserved metadata entry. Its compact JSON +representation must fit the Foundry metadata value limit of 512 characters, +and the reserved entry consumes one of the response metadata map's 16 keys. +Item-level internal metadata is stored directly on the output item and does not +use those response metadata limits. + Use it for lightweight per-turn watermarks, id mappings (e.g. an upstream framework's message id ↔ the emitted item), or stale-message / crash-recovery detection within the turn. It is persisted whenever the response is persisted — at `response.created`, at each `yield stream.checkpoint()`, and at terminal — so -on recovery you read it back from `context.persisted_response`. It is distinct -from the *public* `ResponseObject.metadata` dict (the client's own metadata, -which is NOT stripped). +on recovery, seed `ResponseEventStream(response=context.persisted_response, ...)` +and read it through `stream.internal_metadata`. It is distinct from the *public* +`ResponseObject.metadata` dict (the client's own metadata, which is NOT +stripped). ### Which state facility? diff --git a/sdk/agentserver/azure-ai-agentserver-responses/docs/responses-resilience-spec.md b/sdk/agentserver/azure-ai-agentserver-responses/docs/responses-resilience-spec.md index 29cdca04d329..17743307cb95 100644 --- a/sdk/agentserver/azure-ai-agentserver-responses/docs/responses-resilience-spec.md +++ b/sdk/agentserver/azure-ai-agentserver-responses/docs/responses-resilience-spec.md @@ -755,7 +755,10 @@ client-facing HTTP/SSE payload** — and symmetrically stripped on ingress, so clients can neither read nor inject it. Use it for lightweight per-turn watermarks, id mappings (upstream message id ↔ emitted item), or in-turn stale-message detection; read it back on recovery via -`context.persisted_response`. It is distinct from the *public* +`ResponseEventStream(response=context.persisted_response, ...)`. Response-level +internal metadata is compact-JSON encoded into one string-valued reserved +metadata entry and must fit that entry's 512-character limit. It is distinct +from the *public* `ResponseObject.metadata` (the client's own metadata, never stripped) and from `FoundryStateStore` (cross-turn application state — §8.1). Rule of thumb: cross-turn state → `FoundryStateStore`; reconstruct diff --git a/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_checkpoint.py b/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_checkpoint.py index 7d7e9a10c333..bf4087ad24ee 100644 --- a/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_checkpoint.py +++ b/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_checkpoint.py @@ -11,6 +11,7 @@ from __future__ import annotations import asyncio +import json import time from typing import Any @@ -42,8 +43,8 @@ async def update_response(self, response, *, context=None): # noqa: ANN001 def _event(**md) -> ResponseCheckpointEvent: resp: ResponseObject = {"id": "r1", "object": "response", "status": "in_progress", "output": [], "model": "m"} - for k, v in md.items(): - resp.setdefault("metadata", {}).setdefault("_internal_metadata", {})[k] = v + if md: + resp["metadata"] = {"_internal_metadata": json.dumps(md, separators=(",", ":"), sort_keys=True)} return ResponseCheckpointEvent(resp) @@ -165,7 +166,7 @@ async def test_t21_status_as_is_in_snapshot(): ) assert p.updates[0]["status"] == "in_progress" # Reserved internal_metadata is in the persisted snapshot (storage retains it). - assert p.updates[0]["metadata"]["_internal_metadata"] == {"cp": 1} + assert p.updates[0]["metadata"]["_internal_metadata"] == '{"cp":1}' @pytest.mark.asyncio diff --git a/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_foundry_storage_provider.py b/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_foundry_storage_provider.py index 8166afd0caf8..549c7b590d8b 100644 --- a/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_foundry_storage_provider.py +++ b/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_foundry_storage_provider.py @@ -16,6 +16,7 @@ ) from azure.ai.agentserver.responses._response_context import PlatformContext +from azure.ai.agentserver.responses import ResponseEventStream from azure.ai.agentserver.responses.store._foundry_errors import ( FoundryApiError, FoundryBadRequestError, @@ -167,12 +168,16 @@ async def test_get_response__gets_correct_url(credential: Any, settings: Foundry @pytest.mark.asyncio async def test_get_response__returns_deserialized_response(credential: Any, settings: FoundryStorageSettings) -> None: - provider = _make_provider(credential, settings, _make_response(200, _RESPONSE_DICT)) + response = dict(_RESPONSE_DICT) + response["metadata"] = {"_internal_metadata": '{"checkpoint_marker":"written-before-checkpoint"}'} + provider = _make_provider(credential, settings, _make_response(200, response)) result = await provider.get_response("resp_abc123") assert result["id"] == "resp_abc123" assert result["status"] == "completed" + stream = ResponseEventStream(response_id="resp_abc123", response=result) + assert dict(stream.internal_metadata) == {"checkpoint_marker": "written-before-checkpoint"} @pytest.mark.asyncio @@ -218,11 +223,13 @@ async def test_update_response__sends_serialized_response_body( ) -> None: provider = _make_provider(credential, settings, _make_response(200, {})) response = _response_object() + response["metadata"] = {"_internal_metadata": '{"checkpoint_marker":"written-before-checkpoint"}'} await provider.update_response(response) request = provider._client.send_request.call_args[0][0] payload = json.loads(request.content.decode("utf-8")) assert payload["id"] == "resp_abc123" + assert payload["metadata"]["_internal_metadata"] == '{"checkpoint_marker":"written-before-checkpoint"}' @pytest.mark.asyncio diff --git a/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_internal_metadata.py b/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_internal_metadata.py index a4a1a7c498d9..b51d0e5bb158 100644 --- a/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_internal_metadata.py +++ b/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_internal_metadata.py @@ -9,6 +9,7 @@ from __future__ import annotations +from collections.abc import MutableMapping from copy import deepcopy from typing import Any, cast @@ -17,6 +18,7 @@ from azure.ai.agentserver.responses._egress import strip_internal_metadata from azure.ai.agentserver.responses import CreateResponse, ResponseEventStream from azure.ai.agentserver.responses.models._generated import ResponseObject +from azure.ai.agentserver.responses.streaming._internal_metadata import _ResponseInternalMetadataView def _item() -> dict[str, Any]: @@ -34,9 +36,8 @@ def _item_internal_metadata(item: dict[str, Any]) -> dict[str, Any]: return item.setdefault("internal_metadata", {}) -def _response_internal_metadata(response: ResponseObject) -> dict[str, Any]: - metadata = response.setdefault("metadata", {}) - return metadata.setdefault("_internal_metadata", {}) +def _response_internal_metadata(response: ResponseObject) -> MutableMapping[str, Any]: + return _ResponseInternalMetadataView(cast(dict[str, Any], response)) # -------------------------------------------------------------------------- @@ -150,8 +151,8 @@ def test_t1r_response_empty_view_when_unset(): def test_t2r_response_stores_under_reserved_key(): resp = _response() _response_internal_metadata(resp)["phase"] = 3 - assert resp["metadata"]["_internal_metadata"] == {"phase": 3} - assert dict(resp["metadata"]["_internal_metadata"]) == {"phase": 3} + assert resp["metadata"]["_internal_metadata"] == '{"phase":3}' + assert dict(_response_internal_metadata(resp)) == {"phase": 3} def test_t3r_in_place_mutation_writes_through(): @@ -173,15 +174,25 @@ def test_t4r_does_not_clobber_client_metadata(): def test_t5r_clear_removes_only_reserved_key(): resp = _response() resp["metadata"] = {"user": "x"} - _response_internal_metadata(resp)["phase"] = 3 - resp["metadata"].pop("_internal_metadata") + metadata = _response_internal_metadata(resp) + metadata["phase"] = 3 + metadata.clear() assert dict(resp["metadata"]) == {"user": "x"} def test_t6r_512_char_guard(): resp = _response() - _response_internal_metadata(resp)["big"] = "x" * 600 - assert resp["metadata"]["_internal_metadata"]["big"] == "x" * 600 + with pytest.raises(ValueError, match="512-char limit"): + _response_internal_metadata(resp)["big"] = "x" * 600 + assert "_internal_metadata" not in resp.get("metadata", {}) + + +def test_t6r_unicode_uses_character_limit_not_ascii_escape_length(): + resp = _response() + value = "\N{SNOWMAN}" * 100 + _response_internal_metadata(resp)["unicode"] = value + assert len(resp["metadata"]["_internal_metadata"]) < 512 + assert dict(_response_internal_metadata(resp)) == {"unicode": value} def test_t6r2_16_key_guard(): @@ -192,8 +203,9 @@ def test_t6r2_16_key_guard(): resp16 = _response() resp16["metadata"] = {f"k{i}": "v" for i in range(16)} - _response_internal_metadata(resp16)["p"] = 1 - assert resp16["metadata"]["_internal_metadata"] == {"p": 1} + with pytest.raises(ValueError, match="already has 16 keys"): + _response_internal_metadata(resp16)["p"] = 1 + assert "_internal_metadata" not in resp16["metadata"] def test_t7r_v_shaped_response_empty_view(): @@ -207,8 +219,9 @@ def test_t10r_stream_proxy_is_response_view(): req = CreateResponse({"model": "m", "input": "hi"}) stream = ResponseEventStream(response_id="resp_1", request=req) stream.internal_metadata["phase"] = 3 - assert dict(stream.response["metadata"]["_internal_metadata"]) == {"phase": 3} - stream.response["metadata"]["_internal_metadata"]["x"] = 1 + assert stream.response["metadata"]["_internal_metadata"] == '{"phase":3}' + second_view = _ResponseInternalMetadataView(stream.response) + second_view["x"] = 1 assert stream.internal_metadata["x"] == 1 @@ -216,14 +229,33 @@ def test_t28d_response_reserved_key_roundtrips(): resp = _response() _response_internal_metadata(resp)["phase"] = 3 reloaded = dict(resp) - assert dict(reloaded["metadata"]["_internal_metadata"]) == {"phase": 3} - assert reloaded["metadata"]["_internal_metadata"] == {"phase": 3} + assert reloaded["metadata"]["_internal_metadata"] == '{"phase":3}' + assert dict(ResponseEventStream(response_id="resp_1", response=reloaded).internal_metadata) == {"phase": 3} def test_t7a_compact_deterministic_encoding(): # Deterministic so checkpoint idempotency byte-compare is stable. resp = _response() _response_internal_metadata(resp)["b"] = 2 - resp["metadata"]["_internal_metadata"]["a"] = 1 + _response_internal_metadata(resp)["a"] = 1 + assert resp["metadata"]["_internal_metadata"] == '{"a":1,"b":2}' stripped = strip_internal_metadata(deepcopy(resp)) assert stripped["metadata"] is None + + +def test_nested_dict_snapshot_is_normalized_without_property_access(): + resp = _response() + resp["metadata"] = {"user": "x", "_internal_metadata": {"phase": 3}} + + stream = ResponseEventStream(response_id="resp_1", response=resp) + + assert stream.response["metadata"] == {"user": "x", "_internal_metadata": '{"phase":3}'} + + +def test_nested_dict_snapshot_respects_total_metadata_key_limit(): + resp = _response() + resp["metadata"] = {f"k{i}": "v" for i in range(16)} + resp["metadata"]["_internal_metadata"] = {"phase": 3} + + with pytest.raises(ValueError, match="already has 17 keys"): + ResponseEventStream(response_id="resp_1", response=resp) diff --git a/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_internal_metadata_egress.py b/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_internal_metadata_egress.py index 31ad10ca7a24..2f36e5e0dee9 100644 --- a/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_internal_metadata_egress.py +++ b/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_internal_metadata_egress.py @@ -116,7 +116,7 @@ def test_t17a_t17r_live_objects_untouched_after_sse_encode(): # Encode a terminal carrying the full envelope. encode_sse_event(stream.emit_completed()) # T17r: live response still carries the reserved key. - assert dict(stream.response["metadata"]["_internal_metadata"]) == {"cp": 3} + assert stream.response["metadata"]["_internal_metadata"] == '{"cp":3}' # T17a: live output item still carries its bag. assert dict(stream.response["output"][0]["internal_metadata"]) == {"phase": "gather"} diff --git a/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_internal_metadata_provider_roundtrip.py b/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_internal_metadata_provider_roundtrip.py index 72e5dc0b341d..c5003680614b 100644 --- a/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_internal_metadata_provider_roundtrip.py +++ b/sdk/agentserver/azure-ai-agentserver-responses/tests/unit/test_internal_metadata_provider_roundtrip.py @@ -14,6 +14,7 @@ import pytest +from azure.ai.agentserver.responses import ResponseEventStream from azure.ai.agentserver.responses.models._generated import ResponseObject from azure.ai.agentserver.responses.store._file import FileResponseStore from azure.ai.agentserver.responses.store._memory import InMemoryResponseProvider @@ -39,7 +40,7 @@ def _response(resp_id: str, output: list) -> ResponseObject: "status": "completed", "output": output, "model": "m", - "metadata": {"_internal_metadata": {"completed_phases": 3}}, + "metadata": {"_internal_metadata": '{"completed_phases":3}'}, }, ) @@ -63,13 +64,16 @@ async def test_t28_t28a_response_output_item_internal_metadata_preserved(name, p assert loaded["output"][0]["internal_metadata"] == {"phase": "gather", "n": 7} # T28a — update + get - resp["metadata"]["_internal_metadata"]["extra"] = "x" - await provider.update_response(resp) + stream = ResponseEventStream(response_id="resp_a", response=resp) + stream.internal_metadata["extra"] = "x" + await provider.update_response(cast(ResponseObject, stream.response)) loaded2 = await provider.get_response("resp_a") assert loaded2["output"][0]["internal_metadata"]["n"] == 7 # T28d — response-level reserved key round-trips - assert loaded2["metadata"]["_internal_metadata"] == {"completed_phases": 3, "extra": "x"} + reloaded = ResponseEventStream(response_id="resp_a", response=loaded2) + assert dict(reloaded.internal_metadata) == {"completed_phases": 3, "extra": "x"} + assert loaded2["metadata"]["_internal_metadata"] == '{"completed_phases":3,"extra":"x"}' @pytest.mark.asyncio