Skip to content
Draft
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
6 changes: 6 additions & 0 deletions sdk/agentserver/azure-ai-agentserver-responses/CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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)
)
Expand All @@ -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.
Expand Down
Original file line number Diff line number Diff line change
@@ -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())
Original file line number Diff line number Diff line change
Expand Up @@ -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?

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
from __future__ import annotations

import asyncio
import json
import time
from typing import Any

Expand Down Expand Up @@ -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)


Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@

from __future__ import annotations

from collections.abc import MutableMapping
from copy import deepcopy
from typing import Any, cast

Expand All @@ -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]:
Expand All @@ -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))


# --------------------------------------------------------------------------
Expand Down Expand Up @@ -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():
Expand All @@ -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():
Expand All @@ -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():
Expand All @@ -207,23 +219,43 @@ 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


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)
Original file line number Diff line number Diff line change
Expand Up @@ -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"}

Expand Down
Loading
Loading