Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down Expand Up @@ -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
Expand All @@ -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:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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]
Expand Down Expand Up @@ -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("{")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
133 changes: 132 additions & 1 deletion python/packages/orchestrations/tests/test_group_chat.py
Original file line number Diff line number Diff line change
@@ -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

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

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


Expand Down
16 changes: 16 additions & 0 deletions python/packages/orchestrations/tests/test_magentic.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Loading