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
8 changes: 8 additions & 0 deletions sdk/agentserver/azure-ai-agentserver-core/CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,13 @@
# Release History

## 2.2.1 (Unreleased)

### Bugs Fixed

- Failed steering queue appends no longer retain an unreturned acknowledgment
future. Per-slot acknowledgment IDs prevent a committed append with a lost
response from routing the next accepted input's result incorrectly.

## 2.2.0 (2026-09-23)

### Other Changes
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,4 +2,4 @@
# Copyright (c) Microsoft Corporation. All rights reserved.
# ---------------------------------------------------------

VERSION = "2.2.0"
VERSION = "2.2.1"
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,29 @@
_STEERING_QUEUE_CAP = 9


def _steering_pending_ack_ids(steering: dict[str, Any], pending_count: int) -> list[str | None]:
"""Align internal acknowledgment IDs with queued inputs from older records.

:param steering: Persisted steering state.
:type steering: dict[str, Any]
:param pending_count: Number of queued inputs.
:type pending_count: int
:return: IDs aligned with pending inputs, including untagged legacy slots.
:rtype: list[str | None]
:raises ValueError: If stored IDs cannot align with the pending queue.
"""
if pending_count == 0:
return []
ids = steering.get("pending_ack_ids")
if ids is None:
return [None] * pending_count
if not isinstance(ids, list) or len(ids) > pending_count or any(
value is not None and not isinstance(value, str) for value in ids
):
raise ValueError("Invalid steering pending_ack_ids for pending queue")
return ids + [None] * (pending_count - len(ids))


# --------------------------------------------------------------------------- #
# Hash helper
# --------------------------------------------------------------------------- #
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -798,6 +798,7 @@ async def _append_steering_input( # pylint: disable=protected-access,too-many-l
task_id: str,
input_val: Any,
existing: Any,
ack_id: str,
input_id: str | None = None,
if_last_input_id: str | None = None,
) -> None:
Expand All @@ -812,6 +813,8 @@ async def _append_steering_input( # pylint: disable=protected-access,too-many-l
:keyword existing: The previously-fetched task record (used for the
first etag attempt; later attempts re-fetch internally).
:paramtype existing: Any
:keyword ack_id: Internal identifier binding this queue slot to its acknowledgment.
:paramtype ack_id: str
:keyword input_id: When set, the new input's identity.
Used to advance ``payload["last_input_id"]``
atomically with the queue append.
Expand Down Expand Up @@ -867,8 +870,10 @@ async def _append_steering_input( # pylint: disable=protected-access,too-many-l
_STEERING_INPUT_KEY_PREFIX,
_STEERING_THRESHOLD_BYTES,
_resolve_input_storage,
_steering_pending_ack_ids,
)

pending_ack_ids = _steering_pending_ack_ids(steering, len(pending))
next_seq = int(steering.get("next_input_seq", 0))
steering_key = f"{_STEERING_INPUT_KEY_PREFIX}{next_seq}"
store_mode, queue_entry = _resolve_input_storage(
Expand All @@ -883,7 +888,9 @@ async def _append_steering_input( # pylint: disable=protected-access,too-many-l
steering["next_input_seq"] = next_seq + 1

pending.append(queue_entry)
pending_ack_ids.append(ack_id)
steering["pending_inputs"] = pending
steering["pending_ack_ids"] = pending_ack_ids
steering["cancel_requested"] = True
# SOT: the
# internal _steering["generation"] payload field is removed
Expand Down Expand Up @@ -949,8 +956,8 @@ def _create_steering_ack_run(
manager: Any,
task_id: str,
future: Any,
ack_id: str,
input_id: str | None = None,
input_val: Any = None,
) -> TaskRun[Output]:
"""Create a TaskRun for a queued steering input.

Expand All @@ -960,11 +967,10 @@ def _create_steering_ack_run(
:type task_id: str
:param future: Future that will resolve with the next-turn outcome.
:type future: Any
:param ack_id: Internal identifier of the queued steering slot.
:type ack_id: str
:param input_id: The input_id stamped on the queued input (if any).
:type input_id: str | None
:param input_val: The raw queued input value (used to identify the
slot when ``cancel()`` is invoked on the returned handle).
:type input_val: Any
:return: A :class:`TaskRun` whose result resolves with the queued turn.
:rtype: TaskRun[Output]
"""
Expand All @@ -973,8 +979,7 @@ async def _queued_cancel_cb() -> None:
await manager._cancel_queued_steering_input( # pylint: disable=protected-access
task_id=task_id,
future=future,
input_id=input_id,
input_val=input_val,
ack_id=ack_id,
)

return TaskRun(
Expand Down Expand Up @@ -1306,21 +1311,31 @@ async def _lifecycle_start_inner( # pylint: disable=too-many-locals,too-many-st
if self._opts.steerable:
# Steering path: append input to queue, signal cancel, return ack
# pylint: disable=protected-access
ack_future = manager._register_steering_future(task_id)
await self._append_steering_input(
manager,
task_id=task_id,
input_val=input,
existing=existing,
input_id=input_id,
if_last_input_id=if_last_input_id,
)
# Keep drain from binding a future before its append is accepted.
async with manager._get_task_write_lock(task_id):
ack_id = _generate_input_id()
ack_future = manager._register_steering_future(task_id, ack_id)
appended = False
try:
await self._append_steering_input(
manager,
task_id=task_id,
input_val=input,
existing=existing,
ack_id=ack_id,
input_id=input_id,
if_last_input_id=if_last_input_id,
)
appended = True
finally:
if not appended:
manager._unregister_steering_future(task_id, ack_id, ack_future)
# Set cancel on in-memory context if task runs in this process
active = manager._active_tasks.get(task_id)
# pylint: enable=protected-access
if active:
active.context.cancel.set()
return self._create_steering_ack_run(manager, task_id, ack_future, input_id=input_id, input_val=input)
return self._create_steering_ack_run(manager, task_id, ack_future, ack_id, input_id=input_id)
raise TaskConflictError(task_id, "in_progress")

# completed (or any other terminal status)
Expand Down Expand Up @@ -1748,8 +1763,8 @@ async def delete(self, task_id: str) -> None:
exec_task.cancel()

# 2. Resolve all queued steerer futures with TaskCancelled.
pending = getattr(mgr, "_pending_steering_futures", {}).pop(task_id, [])
for queued_fut in pending:
pending = getattr(mgr, "_pending_steering_futures", {}).pop(task_id, {})
for queued_fut in pending.values():
if not queued_fut.done():
queued_fut.set_exception(TaskCancelled())

Expand Down
Loading
Loading