diff --git a/python/packages/orchestrations/agent_framework_orchestrations/_group_chat.py b/python/packages/orchestrations/agent_framework_orchestrations/_group_chat.py index eb12cfbac8a..0ac34212c1c 100644 --- a/python/packages/orchestrations/agent_framework_orchestrations/_group_chat.py +++ b/python/packages/orchestrations/agent_framework_orchestrations/_group_chat.py @@ -52,7 +52,7 @@ ) from ._feature_usage import FeatureIndex from ._orchestration_request_info import AgentApprovalExecutor -from ._orchestrator_helpers import clean_conversation_for_handoff +from ._orchestrator_helpers import clean_conversation_for_handoff, extract_markdown_fence_bodies from ._participant_output_config import ( UNSET, _coalesce_output_from, # pyright: ignore[reportPrivateUsage] @@ -450,7 +450,9 @@ def _parse_agent_output(cls, agent_response: Any) -> AgentOrchestrationOutput: Preferred path is structured output (`agent_response.value`) when available. If only text is available, first attempt strict JSON parsing and then apply a - temporary concatenated-JSON fallback as a stop-gap. + temporary concatenated-JSON fallback as a stop-gap. Providers that ignore + `response_format` often wrap the JSON in a Markdown code fence, so the body of + the last fenced block is tried as well. """ try: structured_value = agent_response.value @@ -470,6 +472,11 @@ def _parse_agent_output(cls, agent_response: Any) -> AgentOrchestrationOutput: if response_text and response_text not in text_candidates: text_candidates.append(response_text) + for candidate in list(text_candidates): + fence_bodies = extract_markdown_fence_bodies(candidate) + if fence_bodies and fence_bodies[-1] not in text_candidates: + text_candidates.append(fence_bodies[-1]) + last_error: Exception | None = None for candidate in text_candidates: try: diff --git a/python/packages/orchestrations/agent_framework_orchestrations/_magentic.py b/python/packages/orchestrations/agent_framework_orchestrations/_magentic.py index fc828b0e095..6206022f043 100644 --- a/python/packages/orchestrations/agent_framework_orchestrations/_magentic.py +++ b/python/packages/orchestrations/agent_framework_orchestrations/_magentic.py @@ -4,7 +4,6 @@ import contextlib import json import logging -import re import sys from abc import ABC, abstractmethod from collections.abc import Callable, Sequence @@ -39,6 +38,7 @@ ParticipantRegistry, ) from ._feature_usage import FeatureIndex +from ._orchestrator_helpers import extract_markdown_fence_bodies from ._participant_output_config import ( UNSET, _coalesce_output_from, # pyright: ignore[reportPrivateUsage] @@ -414,9 +414,12 @@ def _extract_json(text: str) -> dict[str, Any]: The `text` method is concatenating multiple text contents from diff msgs into a single string. """ - fence = re.search(r"```(?:json)?\s*(\{[\s\S]*?\})\s*```", text, flags=re.IGNORECASE) - if fence: - candidate = fence.group(1) + fenced_object = next( + (body for body in extract_markdown_fence_bodies(text) if body.startswith("{") and body.endswith("}")), + None, + ) + if fenced_object is not None: + candidate = fenced_object else: # Find first balanced JSON object start = text.find("{") diff --git a/python/packages/orchestrations/agent_framework_orchestrations/_orchestrator_helpers.py b/python/packages/orchestrations/agent_framework_orchestrations/_orchestrator_helpers.py index b4da5b7e282..776110cf178 100644 --- a/python/packages/orchestrations/agent_framework_orchestrations/_orchestrator_helpers.py +++ b/python/packages/orchestrations/agent_framework_orchestrations/_orchestrator_helpers.py @@ -92,3 +92,72 @@ def create_completion_message( contents=[message_text], author_name=author_name, ) + + +_FENCE_MARKER = "```" + + +def _backtick_run_end(text: str, start: int) -> int: + """Return the index just past the run of backticks that begins at ``start``.""" + end = start + while end < len(text) and text[end] == "`": + end += 1 + return end + + +def _fence_info_string_end(text: str, start: int) -> int: + """Return the index just past the info string of an opening fence such as ``json`` or ``application/json``. + + An info string has to start with a letter, so a same-line body such as ``{"a": 1}`` that + follows the fence directly is left in place. + """ + index = start + while index < len(text) and text[index] in " \t": + index += 1 + if index >= len(text) or not text[index].isalpha(): + return start + while index < len(text) and (text[index].isalnum() or text[index] in "_-+./"): + index += 1 + return index + + +def extract_markdown_fence_bodies(text: str) -> list[str]: + """Return the stripped bodies of the Markdown code fences in ``text``, in source order. + + Model output that ignores ``response_format`` often wraps JSON in a fence. An opening fence + is a run of three or more backticks followed by an optional info string. The closing fence + is the next run that is at least as long and ends its line, so backticks inside a JSON + string value (which always end in a quote on the same line) never close the block, and an + outer fence can use more backticks than the ones inside it. A fence that is never closed + yields nothing. + + The scan advances through ``text`` without backtracking, so it stays linear in the input + length for malformed output such as an unterminated fence followed by whitespace. + """ + bodies: list[str] = [] + length = len(text) + open_start = text.find(_FENCE_MARKER) + while open_start != -1: + open_end = _backtick_run_end(text, open_start) + closing_marker = text[open_start:open_end] + body_start = _fence_info_string_end(text, open_end) + + search_from = body_start + close_start = close_end = -1 + while (run_start := text.find(closing_marker, search_from)) != -1: + run_end = _backtick_run_end(text, run_start) + line_end = run_end + while line_end < length and text[line_end] in " \t\r": + line_end += 1 + if line_end == length or text[line_end] == "\n": + close_start, close_end = run_start, run_end + break + search_from = line_end + if close_start == -1: + break + + body = text[body_start:close_start].strip() + if body: + bodies.append(body) + open_start = text.find(_FENCE_MARKER, close_end) + return bodies diff --git a/python/packages/orchestrations/tests/test_group_chat.py b/python/packages/orchestrations/tests/test_group_chat.py index 961ba2ed2c7..eb2b7fa7369 100644 --- a/python/packages/orchestrations/tests/test_group_chat.py +++ b/python/packages/orchestrations/tests/test_group_chat.py @@ -1,5 +1,7 @@ # Copyright (c) Microsoft. All rights reserved. +import json +import time from collections.abc import AsyncIterable, Callable, Sequence from typing import Any, cast @@ -30,7 +32,8 @@ MagenticProgressLedgerItem, ) -from agent_framework_orchestrations import BaseGroupChatOrchestrator +from agent_framework_orchestrations import AgentBasedGroupChatOrchestrator, BaseGroupChatOrchestrator +from agent_framework_orchestrations._orchestrator_helpers import extract_markdown_fence_bodies class StubAgent(BaseAgent): @@ -178,6 +181,40 @@ async def run( # type: ignore[override] # ty: ignore[invalid-method-override] ) +class FencedJsonManagerAgent(Agent): + """Manager agent that ignores response_format and wraps its JSON in a Markdown code fence.""" + + def __init__(self) -> None: + super().__init__(client=cast(Any, MockChatClient()), name="fenced_manager", description="Fenced JSON manager") + self._call_count = 0 + + async def run( # type: ignore[override] # ty: ignore[invalid-method-override] + self, + messages: str | Content | Message | Sequence[str | Content | Message] | None = None, + *, + session: AgentSession | None = None, + **kwargs: Any, + ) -> AgentResponse[Any]: + if self._call_count == 0: + self._call_count += 1 + text = ( + "```json\n" + '{"terminate": false, "reason": "delegate", "next_speaker": "agent", "final_message": null}\n' + "```" + ) + else: + text = ( + "```json\n" + "{\n" + ' "terminate": true,\n' + ' "reason": "Task complete",\n' + ' "final_message": "fenced manager final"\n' + "}\n" + "```" + ) + return AgentResponse(messages=[Message(role="assistant", contents=[text], author_name=self.name)]) + + def make_sequence_selector() -> Callable[[GroupChatState], str]: state_counter = {"value": 0} @@ -411,6 +448,100 @@ async def test_agent_manager_handles_concatenated_json_output() -> None: assert final_update.text == "concatenated manager final" +async def test_agent_manager_handles_fenced_json_output() -> None: + manager = FencedJsonManagerAgent() + worker = StubAgent("agent", "worker response") + + workflow = GroupChatBuilder( + participants=[worker], + orchestrator_agent=manager, + ).build() + + updates: list[AgentResponseUpdate] = [] + async for event in workflow.run("coordinate task", stream=True): + if event.type == "output" and isinstance(event.data, AgentResponseUpdate): + updates.append(event.data) + + assert updates + final_update = updates[-1] + # terminate=true inside the fence must end the chat, not fall through to max_rounds. + assert final_update.author_name == manager.name + assert final_update.text == "fenced manager final" + + +@pytest.mark.parametrize( + "text", + [ + '```json\n{"terminate": true, "reason": "done", "final_message": "bye"}\n```', + '```\n{"terminate": true, "reason": "done", "final_message": "bye"}\n```', + '```JSON\r\n{"terminate": true, "reason": "done", "final_message": "bye"}\r\n```', + '```json {"terminate": true, "reason": "done", "final_message": "bye"} ```', + 'Here is my decision:\n```json\n{"terminate": true, "reason": "done", "final_message": "bye"}\n```', + '```application/json\n{"terminate": true, "reason": "done", "final_message": "bye"}\n```', + '``` json\n{"terminate": true, "reason": "done", "final_message": "bye"}\n```', + '```{"terminate": true, "reason": "done", "final_message": "bye"}```', + '````json\n{"terminate": true, "reason": "done", "final_message": "bye"}\n````', + ( + 'Draft:\n```json\n{"terminate": false, "reason": "draft"}\n```\n' + 'Final:\n```json\n{"terminate": true, "reason": "done", "final_message": "bye"}\n```' + ), + ], +) +def test_agent_orchestrator_parses_fenced_json(text: str) -> None: + response = AgentResponse(messages=[Message(role="assistant", contents=[text])]) + + output = AgentBasedGroupChatOrchestrator._parse_agent_output(response) + + assert output.terminate is True + assert output.final_message == "bye" + + +def test_agent_orchestrator_parses_fenced_json_with_nested_fence_in_final_message() -> None: + final_message = "Run this:\n```python\nprint('hi')\n```\nDone." + payload = json.dumps({"terminate": True, "reason": "done", "final_message": final_message}, indent=2) + response = AgentResponse(messages=[Message(role="assistant", contents=[f"```json\n{payload}\n```"])]) + + output = AgentBasedGroupChatOrchestrator._parse_agent_output(response) + + assert output.terminate is True + assert output.final_message == final_message + + +def test_agent_orchestrator_rejects_unterminated_fence_in_linear_time() -> None: + # A backtracking extractor stalled for seconds on an unterminated fence followed by + # whitespace, and this parser runs on the event loop. + text = "```json\n" + " " * 200_000 + "\n" + "\t" * 200_000 + response = AgentResponse(messages=[Message(role="assistant", contents=[text])]) + + started = time.perf_counter() + with pytest.raises(ValueError, match="Failed to parse agent orchestration output"): + AgentBasedGroupChatOrchestrator._parse_agent_output(response) + assert time.perf_counter() - started < 1.0 + + +@pytest.mark.parametrize( + ("text", "expected"), + [ + ("no fences here", []), + ("```json\n{}\n```", ["{}"]), + ("```\nfirst\n```\ntext\n```py\nsecond\n```", ["first", "second"]), + ("````md\nouter\n```py\ninner\n```\n````", ["outer\n```py\ninner\n```"]), + ('```json\n{"a": "x ``` y"}\n```', ['{"a": "x ``` y"}']), + ("```json\n{}", []), + ("```\n```", []), + ], +) +def test_extract_markdown_fence_bodies(text: str, expected: list[str]) -> None: + assert extract_markdown_fence_bodies(text) == expected + + +def test_agent_orchestrator_rejects_fenced_non_json() -> None: + response = AgentResponse(messages=[Message(role="assistant", contents=["```json\nnot json\n```"])]) + + with pytest.raises(ValueError, match="Failed to parse agent orchestration output"): + AgentBasedGroupChatOrchestrator._parse_agent_output(response) + + # Comprehensive tests for group chat functionality diff --git a/python/packages/orchestrations/tests/test_magentic.py b/python/packages/orchestrations/tests/test_magentic.py index e9f8df1fea1..38a3dcc718c 100644 --- a/python/packages/orchestrations/tests/test_magentic.py +++ b/python/packages/orchestrations/tests/test_magentic.py @@ -1440,4 +1440,20 @@ def test_standard_manager_checkpoint_restore_empty_state(): assert mgr.task_ledger is None +@pytest.mark.parametrize( + "text", + [ + 'Ledger:\n```json\n{"is_request_satisfied": true}\n```', + '```application/json\n{"is_request_satisfied": true}\n```', + '```py\nprint(1)\n```\n```json\n{"is_request_satisfied": true}\n```', + '```json\n{"is_request_satisfied": true, "note": "use ```py``` here"}\n```', + 'No fence: {"is_request_satisfied": true} trailing', + ], +) +def test_extract_json_reads_first_fenced_object(text: str) -> None: + from agent_framework_orchestrations._magentic import _extract_json # type: ignore + + assert _extract_json(text)["is_request_satisfied"] is True + + # endregion