Skip to content
Merged
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
9 changes: 9 additions & 0 deletions agent_core/core/event_stream/event.py
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,12 @@ class Event:
uses it to keep the run's "Working…" indicator up across the
bubble instead of treating every agent bubble as a run-ending
reply. None/False for final replies and non-chat events.
question: For AGENT_MESSAGE events only: set when the message is a
question to the user with suggested responses (send_message with
suggested_responses). Shape:
``{"options": ["Yes", "No"], "allow_free_text": true}``. The UI
renders it as answer chips plus a pinned question box above the
chat composer. None for ordinary messages.
"""

message: str
Expand All @@ -157,6 +163,7 @@ class Event:
action_output: Optional[Dict[str, Any]] = None
platform: Optional[str] = None
continue_work: Optional[bool] = None
question: Optional[Dict[str, Any]] = None

def display_text(self) -> Optional[str]:
"""
Expand Down Expand Up @@ -189,6 +196,7 @@ def to_dict(self) -> Dict[str, Any]:
"action_output": self.action_output,
"platform": self.platform,
"continue_work": self.continue_work,
"question": self.question,
}

@classmethod
Expand Down Expand Up @@ -228,6 +236,7 @@ def from_dict(cls, data: Dict[str, Any]) -> "Event":
action_output=data.get("action_output"),
platform=data.get("platform"),
continue_work=data.get("continue_work"),
question=data.get("question"),
)

@property
Expand Down
5 changes: 5 additions & 0 deletions agent_core/core/impl/event_stream/event_stream.py
Original file line number Diff line number Diff line change
Expand Up @@ -217,6 +217,7 @@ def log(
action_output: Optional[dict] = None,
platform: Optional[str] = None,
continue_work: Optional[bool] = None,
question: Optional[dict] = None,
) -> int:
"""
Append a new event to the stream and trigger summarization if needed.
Expand Down Expand Up @@ -249,6 +250,9 @@ def log(
continue_work: For AGENT_MESSAGE events: True when this is a
mid-run progress update and the agent keeps working after
sending it (drives the UI's persistent "Working…" row).
question: For AGENT_MESSAGE events: suggested-response payload
(``{"options": [...], "allow_free_text": bool}``) when the
message is a question the UI should pin above the composer.

Returns:
The zero-based index of the event within ``tail_events``.
Expand All @@ -270,6 +274,7 @@ def log(
action_output=action_output,
platform=platform,
continue_work=continue_work,
question=question,
)
rec = EventRecord(event=ev)

Expand Down
2 changes: 2 additions & 0 deletions agent_core/core/impl/event_stream/manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -293,6 +293,7 @@ def log(
action_output: Optional[dict] = None,
platform: Optional[str] = None,
continue_work: Optional[bool] = None,
question: Optional[dict] = None,
task_id: str | None = None,
) -> int:
"""
Expand Down Expand Up @@ -343,6 +344,7 @@ def log(
action_output=action_output,
platform=platform,
continue_work=continue_work,
question=question,
)

# Also log to markdown files for persistence
Expand Down
3 changes: 2 additions & 1 deletion agent_core/core/prompts/action.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,8 @@
- When you finish the work, send your final message as the ONLY action of
that turn. If you need the user's answer before you can continue, ask the
question as your final message — the session wakes automatically when they
reply.
reply. When asking, offer suggested_responses so the user can answer with
one click.
- Use 'end_turn' to end the run silently when the input needs no reaction
(e.g. third-party platform noise).

Expand Down
33 changes: 31 additions & 2 deletions app/data/action/send_message.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,12 @@
"you need the user's answer before you can continue), and the session will wait "
"for the user's next input. Set continue_work=true ONLY for progress updates "
"sent while you still have more work to do. Do not use multiple send_message "
"actions at the same time - combine messages into one."
"actions at the same time - combine messages into one. When your message is a "
"QUESTION to the user, provide suggested_responses: the UI pins the question "
"above the chat input with one-click answer buttons, so the user can answer "
"even while you keep working or other messages scroll by. Questions stay "
"pinned until answered or dismissed — NEVER re-send a question that is still "
"awaiting the user's answer."
),
default=True,
action_sets=["core"],
Expand All @@ -32,6 +37,23 @@
"will keep working after sending it."
),
},
"suggested_responses": {
"type": "array",
"example": ["Yes, go ahead", "No, skip it"],
"description": (
"Only when the message is a question: 2-5 short suggested answers "
"(plain strings) shown as one-click buttons. Cover the likely answers; "
"keep each under ~8 words. Omit for non-questions."
),
},
"allow_free_text": {
"type": "boolean",
"example": True,
"description": (
"Only with suggested_responses. True (default): the user can also type "
"a custom answer. Set False when only the listed answers are valid."
),
},
},
output_schema={
"status": {
Expand All @@ -58,13 +80,20 @@ async def send_message(input_data: dict) -> dict:
simulated_mode = input_data.get("simulated_mode", False)
# Extract session_id injected by ActionManager for multi-session isolation
session_id = input_data.get("_session_id")
raw_suggestions = input_data.get("suggested_responses") or []
suggested_responses = [str(s).strip() for s in raw_suggestions if str(s).strip()]
allow_free_text = bool(input_data.get("allow_free_text", True))

# In simulated mode, skip the actual interface call for testing
if not simulated_mode:
import app.internal_action_interface as internal_action_interface

await internal_action_interface.InternalActionInterface.do_chat(
message, session_id=session_id, continue_work=continue_work
message,
session_id=session_id,
continue_work=continue_work,
suggested_responses=suggested_responses,
allow_free_text=allow_free_text,
)

# Return 'success' for test compatibility, but keep 'ok' in production if needed
Expand Down
9 changes: 9 additions & 0 deletions app/internal_action_interface.py
Original file line number Diff line number Diff line change
Expand Up @@ -342,6 +342,8 @@ async def do_chat(
platform: Optional[str] = None,
session_id: Optional[str] = None,
continue_work: bool = False,
suggested_responses: Optional[List[str]] = None,
allow_free_text: bool = True,
) -> None:
"""Record an agent-authored chat message to the event stream.

Expand All @@ -353,6 +355,11 @@ async def do_chat(
session_id: Optional task/session ID for multi-task isolation.
continue_work: True when this is a mid-run progress update and
the agent keeps working after sending it.
suggested_responses: When the message is a question, the one-click
answers to offer. A non-empty list marks the message as a
question the UI pins above the chat composer.
allow_free_text: Whether the pinned question also accepts a typed
custom answer (only meaningful with suggested_responses).
"""
if InternalActionInterface.state_manager is None:
raise RuntimeError(
Expand All @@ -366,6 +373,8 @@ async def do_chat(
session_id=session_id,
platform=resolved_platform,
continue_work=continue_work,
suggested_responses=suggested_responses,
allow_free_text=allow_free_text,
)

@staticmethod
Expand Down
24 changes: 22 additions & 2 deletions app/state/state_manager.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
from typing import Optional, TYPE_CHECKING
from typing import List, Optional, TYPE_CHECKING

from agent_core.core.state.types import MainState
from agent_core.core.state.session import StateSession
Expand Down Expand Up @@ -136,6 +136,8 @@ def record_agent_message(
session_id: Optional[str] = None,
platform: Optional[str] = None,
continue_work: bool = False,
suggested_responses: Optional[List[str]] = None,
allow_free_text: bool = True,
) -> None:
"""Record an agent message to a session's event stream.

Expand All @@ -147,6 +149,13 @@ def record_agent_message(
the agent keeps working after sending it (send_message's
continue_work=true). Carried on the event so the UI keeps
its "Working…" indicator up across the bubble.
suggested_responses: Non-empty when the message is a question
with one-click answers. Carried on the event as `question`
so the UI renders answer chips and pins the question above
the composer; the stream copy also lists the offered answers
so the agent remembers what it suggested.
allow_free_text: Whether the pinned question also accepts a
typed custom answer.
"""
target = session_id or MAIN_SESSION_ID

Expand All @@ -155,13 +164,24 @@ def record_agent_message(
else:
event_label = "agent message"

question = None
stream_content = content
if suggested_responses:
question = {
"options": list(suggested_responses),
"allow_free_text": allow_free_text,
}
offered = " | ".join(suggested_responses)
stream_content = f"{content}\n[Offered suggested responses: {offered}]"

self.event_stream_manager.log(
event_label,
content,
stream_content,
event_type=EventType.AGENT_MESSAGE,
display_message=content,
platform=platform,
continue_work=continue_work,
question=question,
task_id=target,
)

Expand Down
26 changes: 26 additions & 0 deletions app/ui_layer/adapters/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -289,6 +289,19 @@ def _handle_agent_message(self, event: UIEvent) -> None:
)
for o in raw_options
]
# A question with suggested responses: the suggestions become the
# message's options (label == value, i.e. the answer text itself),
# and the question flags make the frontend pin it above the composer.
question = event.data.get("question")
is_question = False
allow_free_text = True
if isinstance(question, dict) and question.get("options"):
is_question = True
allow_free_text = bool(question.get("allow_free_text", True))
options = [
ChatMessageOption(label=str(o), value=str(o))
for o in question["options"]
]
asyncio.create_task(
self._display_chat_message(
agent_name,
Expand All @@ -297,6 +310,8 @@ def _handle_agent_message(self, event: UIEvent) -> None:
session_id=event.task_id,
options=options,
continue_work=bool(event.data.get("continue_work", False)),
is_question=is_question,
allow_free_text=allow_free_text,
)
)

Expand Down Expand Up @@ -428,6 +443,8 @@ async def _display_chat_message(
options: Optional[List[ChatMessageOption]] = None,
client_id: Optional[str] = None,
continue_work: bool = False,
is_question: bool = False,
allow_free_text: bool = True,
) -> None:
"""
Display a chat message.
Expand All @@ -441,6 +458,10 @@ async def _display_chat_message(
client_id: Optional client-generated UUID for reconciling with optimistic UI
continue_work: True when this is a mid-run agent progress update
(the run keeps going after this message)
is_question: True when this is a question with suggested
responses — pinned above the composer until answered
allow_free_text: Whether the pinned question also accepts a
typed custom answer
"""
import time

Expand All @@ -454,6 +475,11 @@ async def _display_chat_message(
options=options,
client_id=client_id,
continue_work=continue_work,
is_question=is_question,
allow_free_text=allow_free_text,
# The pinned box / chips are the affordance for questions; the
# bubble must not add the "Please select a response" banner.
requires_choice=not is_question,
)
)

Expand Down
Loading