Skip to content
Closed
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
The table of contents is too big for display.
Diff view
Diff view
  •  
  •  
  •  
The diff you're trying to view is too large. We only load the first 3000 changed files.
9 changes: 8 additions & 1 deletion acp_adapter/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -1534,7 +1534,14 @@ def _notify_title_update(_title: str) -> None:
)
except Exception:
logger.debug("Failed to auto-title ACP session %s", session_id, exc_info=True)
if final_response and conn and not streamed_message:
if final_response and conn:
# Always deliver the final response — plugins may have transformed
# it after streaming finished (e.g. transform_llm_output hook).
logger.warning(
"TRACE: delivering final_response to ACP, streamed=%s len=%d has_debug=%s",
streamed_message, len(final_response),
"🔧 [Debug]" in final_response,
)
update = acp.update_agent_message_text(final_response)
await conn.session_update(session_id, update)

Expand Down
5 changes: 5 additions & 0 deletions agent/conversation_loop.py
Original file line number Diff line number Diff line change
Expand Up @@ -3938,6 +3938,8 @@ def _stop_spinner():
except Exception as _ver_err:
logger.debug("file-mutation verifier footer failed: %s", _ver_err)

_response_transformed = False

# Plugin hook: transform_llm_output
# Fired once per turn after the tool-calling loop completes.
# Plugins can transform the LLM's output text before it's returned.
Expand All @@ -3952,9 +3954,11 @@ def _stop_spinner():
model=agent.model,
platform=getattr(agent, "platform", None) or "",
)
_response_transformed = False
for _hook_result in _transform_results:
if isinstance(_hook_result, str) and _hook_result:
final_response = _hook_result
_response_transformed = True
break # First non-empty string wins
except Exception as exc:
logger.warning("transform_llm_output hook failed: %s", exc)
Expand Down Expand Up @@ -4005,6 +4009,7 @@ def _stop_spinner():
"turn_exit_reason": _turn_exit_reason,
"partial": False, # True only when stopped due to invalid tool calls
"interrupted": interrupted,
"response_transformed": _response_transformed,
"response_previewed": getattr(agent, "_response_was_previewed", False),
"model": agent.model,
"provider": agent.provider,
Expand Down
34 changes: 33 additions & 1 deletion gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -16955,6 +16955,7 @@ def _title_failure_cb(task: str, exc: BaseException) -> None:
"context_length": _context_length,
"session_id": effective_session_id,
"response_previewed": result.get("response_previewed", False),
"response_transformed": result.get("response_transformed", False),
}

# Start progress message sender if enabled
Expand Down Expand Up @@ -17575,7 +17576,16 @@ async def _notify_long_running():
_content_delivered = bool(
_sc and getattr(_sc, "final_content_delivered", False)
)
if not _is_empty_sentinel and (_streamed or _previewed or _content_delivered):
# Plugin hooks (e.g. transform_llm_output) may have appended content
# after streaming finished — when the response was transformed, always
# send the final version so the appended content reaches the client.
_transformed = bool(response.get("response_transformed"))
logger.warning(
"TRACE: final-send check transformed=%s streamed=%s previewed=%s content_delivered=%s final_len=%s",
_transformed, _streamed, _previewed, _content_delivered,
len(response.get("final_response") or ""),
)
if not _is_empty_sentinel and not _transformed and (_streamed or _previewed or _content_delivered):
logger.info(
"Suppressing normal final send for session %s: final delivery already confirmed (streamed=%s previewed=%s content_delivered=%s).",
session_key or "?",
Expand All @@ -17584,6 +17594,28 @@ async def _notify_long_running():
_content_delivered,
)
response["already_sent"] = True
elif not _is_empty_sentinel and _transformed and _sc is not None:
# Plugin hooks transformed the response after streaming — edit the
# existing streamed message instead of sending a duplicate.
_sc_msg_id = _sc.message_id
if _sc_msg_id:
try:
await _sc.adapter.edit_message(
chat_id=source.chat_id,
message_id=_sc_msg_id,
content=response["final_response"],
finalize=True,
)
response["already_sent"] = True
logger.info(
"Edited streamed message %s for session %s to include plugin-transformed content.",
_sc_msg_id, session_key or "?",
)
except Exception as _edit_err:
logger.warning(
"Failed to edit streamed message for session %s: %s",
session_key or "?", _edit_err,
)

# Schedule deletion of tracked temporary progress bubbles after the
# final response lands. Failed runs skip this so bubbles remain as
Expand Down
10 changes: 10 additions & 0 deletions gateway/stream_consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -192,6 +192,16 @@ def final_response_sent(self) -> bool:
"""True when the stream consumer delivered the final assistant reply."""
return self._final_response_sent

@property
def message_id(self) -> str | None:
"""The Discord/chat message ID of the last-sent or edited message."""
return self._message_id

@property
def accumulated_text(self) -> str:
"""The accumulated streamed text (without think-block content)."""
return self._accumulated

@property
def final_content_delivered(self) -> bool:
"""True when the final response content reached the user, even if
Expand Down
Loading