diff --git a/agent_core/core/event_stream/event.py b/agent_core/core/event_stream/event.py index 9cb1f050..0d022288 100644 --- a/agent_core/core/event_stream/event.py +++ b/agent_core/core/event_stream/event.py @@ -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 @@ -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]: """ @@ -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 @@ -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 diff --git a/agent_core/core/impl/event_stream/event_stream.py b/agent_core/core/impl/event_stream/event_stream.py index a596cb00..5aa93647 100644 --- a/agent_core/core/impl/event_stream/event_stream.py +++ b/agent_core/core/impl/event_stream/event_stream.py @@ -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. @@ -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``. @@ -270,6 +274,7 @@ def log( action_output=action_output, platform=platform, continue_work=continue_work, + question=question, ) rec = EventRecord(event=ev) diff --git a/agent_core/core/impl/event_stream/manager.py b/agent_core/core/impl/event_stream/manager.py index 79d562bb..44882448 100644 --- a/agent_core/core/impl/event_stream/manager.py +++ b/agent_core/core/impl/event_stream/manager.py @@ -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: """ @@ -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 diff --git a/agent_core/core/prompts/action.py b/agent_core/core/prompts/action.py index d35958a9..814ec83e 100644 --- a/agent_core/core/prompts/action.py +++ b/agent_core/core/prompts/action.py @@ -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). diff --git a/app/data/action/send_message.py b/app/data/action/send_message.py index 4a6d2f56..24679b6d 100644 --- a/app/data/action/send_message.py +++ b/app/data/action/send_message.py @@ -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"], @@ -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": { @@ -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 diff --git a/app/internal_action_interface.py b/app/internal_action_interface.py index 92f1c28a..62fe6bbb 100644 --- a/app/internal_action_interface.py +++ b/app/internal_action_interface.py @@ -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. @@ -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( @@ -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 diff --git a/app/state/state_manager.py b/app/state/state_manager.py index ad533490..e11669f1 100644 --- a/app/state/state_manager.py +++ b/app/state/state_manager.py @@ -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 @@ -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. @@ -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 @@ -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, ) diff --git a/app/ui_layer/adapters/base.py b/app/ui_layer/adapters/base.py index c3119ed1..07e4b3b2 100644 --- a/app/ui_layer/adapters/base.py +++ b/app/ui_layer/adapters/base.py @@ -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, @@ -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, ) ) @@ -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. @@ -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 @@ -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, ) ) diff --git a/app/ui_layer/adapters/browser_adapter.py b/app/ui_layer/adapters/browser_adapter.py index 7a7c8834..e0f2f08d 100644 --- a/app/ui_layer/adapters/browser_adapter.py +++ b/app/ui_layer/adapters/browser_adapter.py @@ -104,7 +104,12 @@ StatusBarProtocol, FootageComponentProtocol, ) -from app.ui_layer.components.types import ChatMessage, ActionItem, Attachment +from app.ui_layer.components.types import ( + ChatMessage, + ActionItem, + Attachment, + QUESTION_DISMISSED_VALUE, +) from app.ui_layer.events import UIEvent, UIEventType from app.ui_layer.onboarding import OnboardingFlowController from app.ui_layer.metrics import MetricsCollector @@ -203,48 +208,54 @@ def _init_storage(self) -> None: # Load recent messages from storage (initial page) stored_messages = self._storage.get_recent_messages(limit=50) for stored in stored_messages: - attachments = None - if stored.attachments: - attachments = [ - Attachment( - name=att.get("name", ""), - path=att.get("path", ""), - type=att.get("type", ""), - size=att.get("size", 0), - url=att.get("url", ""), - ) - for att in stored.attachments - ] - options = None - if stored.options: - from app.ui_layer.components.types import ChatMessageOption - - options = [ - ChatMessageOption( - label=o.get("label", ""), - value=o.get("value", ""), - style=o.get("style", "default"), - ) - for o in stored.options - ] - self._messages.append( - ChatMessage( - sender=stored.sender, - content=stored.content, - style=stored.style, - timestamp=stored.timestamp, - message_id=stored.message_id, - attachments=attachments, - session_id=stored.session_id, - options=options, - option_selected=stored.option_selected, - continue_work=stored.continue_work, - ) - ) + self._messages.append(self._stored_to_chat_message(stored)) except Exception: # Storage may not be available, continue without persistence pass + @staticmethod + def _stored_to_chat_message(stored) -> ChatMessage: + """Rehydrate a StoredChatMessage row into the live ChatMessage shape.""" + from app.ui_layer.components.types import ChatMessageOption + + attachments = None + if stored.attachments: + attachments = [ + Attachment( + name=att.get("name", ""), + path=att.get("path", ""), + type=att.get("type", ""), + size=att.get("size", 0), + url=att.get("url", ""), + ) + for att in stored.attachments + ] + options = None + if stored.options: + options = [ + ChatMessageOption( + label=o.get("label", ""), + value=o.get("value", ""), + style=o.get("style", "default"), + ) + for o in stored.options + ] + return ChatMessage( + sender=stored.sender, + content=stored.content, + style=stored.style, + timestamp=stored.timestamp, + message_id=stored.message_id, + attachments=attachments, + session_id=stored.session_id, + options=options, + option_selected=stored.option_selected, + continue_work=stored.continue_work, + is_question=stored.is_question, + allow_free_text=stored.allow_free_text, + requires_choice=not stored.is_question, + ) + async def append_message(self, message: ChatMessage) -> None: """Append message and broadcast to clients.""" self._messages.append(message) @@ -283,6 +294,8 @@ async def append_message(self, message: ChatMessage) -> None: session_id=message.session_id, options=options_data, continue_work=message.continue_work, + is_question=message.is_question, + allow_free_text=message.allow_free_text, ) self._storage.insert_message(stored) except Exception: @@ -343,47 +356,7 @@ def get_messages_before( stored = self._storage.get_messages_before( before_timestamp, session_id=session_id, limit=limit ) - messages = [] - for s in stored: - attachments = None - if s.attachments: - attachments = [ - Attachment( - name=att.get("name", ""), - path=att.get("path", ""), - type=att.get("type", ""), - size=att.get("size", 0), - url=att.get("url", ""), - ) - for att in s.attachments - ] - options = None - if s.options: - from app.ui_layer.components.types import ChatMessageOption - - options = [ - ChatMessageOption( - label=o.get("label", ""), - value=o.get("value", ""), - style=o.get("style", "default"), - ) - for o in s.options - ] - messages.append( - ChatMessage( - sender=s.sender, - content=s.content, - style=s.style, - timestamp=s.timestamp, - message_id=s.message_id, - attachments=attachments, - session_id=s.session_id, - options=options, - option_selected=s.option_selected, - continue_work=s.continue_work, - ) - ) - return messages + return [self._stored_to_chat_message(s) for s in stored] except Exception: return [] @@ -1353,6 +1326,15 @@ async def _handle_ws_message(self, data: Dict[str, Any], ws=None) -> None: message_id = data.get("messageId", "") await self._handle_option_click(value, session_id, message_id) + elif msg_type == "question_response": + value = data.get("value", "") + session_id = data.get("sessionId", "") + message_id = data.get("messageId", "") + dismissed = bool(data.get("dismissed", False)) + await self._handle_question_response( + value, session_id, message_id, dismissed + ) + # Settings operations elif msg_type == "settings_get": await self._handle_settings_get() @@ -3821,6 +3803,65 @@ async def _handle_option_click( f"[OPTION_CLICK] Error handling option click: {e}", exc_info=True ) + async def _handle_question_response( + self, value: str, session_id: str, message_id: str, dismissed: bool + ) -> None: + """Handle the user answering (or dismissing) a pinned agent question. + + Marks the question message answered (which un-pins it everywhere), + then feeds the answer back into the agent as a regular user message + so the normal trigger queue/merge behavior applies. + """ + try: + question_text = "" + recorded = QUESTION_DISMISSED_VALUE if dismissed else value + pending_questions: list = [] + if self._chat and message_id: + for m in self._chat._messages: + if m.message_id == message_id: + question_text = m.content + m.option_selected = recorded + break + if self._chat._storage: + try: + self._chat._storage.update_option_selected( + message_id, recorded + ) + # After marking this one, whatever question messages + # remain unanswered are still pinned in the user's UI. + pending_questions = self._chat._storage.get_pending_questions( + session_id + ) + except Exception: + pass + + # Un-pin on every connected client (the answering client already + # marked the selection optimistically). + await self._broadcast( + { + "type": "question_answered", + "data": { + "sessionId": session_id, + "messageId": message_id, + "value": recorded, + }, + } + ) + + await self._controller.submit_question_answer( + value, + question_text, + session_id, + dismissed, + adapter_id=self._adapter_id, + pending_questions=pending_questions, + ) + except Exception as e: + logger.error( + f"[QUESTION_RESPONSE] Error handling question response: {e}", + exc_info=True, + ) + # ───────────────────────────────────────────────────────────────────── # Session Handlers (sidebar surface) # ───────────────────────────────────────────────────────────────────── @@ -7637,46 +7678,9 @@ async def _reply(payload: Dict[str, Any]) -> None: if storage else [] ) - messages = [] - for s in stored: - attachments = None - if s.attachments: - attachments = [ - Attachment( - name=att.get("name", ""), - path=att.get("path", ""), - type=att.get("type", ""), - size=att.get("size", 0), - url=att.get("url", ""), - ) - for att in s.attachments - ] - options = None - if s.options: - from app.ui_layer.components.types import ChatMessageOption - - options = [ - ChatMessageOption( - label=o.get("label", ""), - value=o.get("value", ""), - style=o.get("style", "default"), - ) - for o in s.options - ] - messages.append( - ChatMessage( - sender=s.sender, - content=s.content, - style=s.style, - timestamp=s.timestamp, - message_id=s.message_id, - attachments=attachments, - session_id=s.session_id, - options=options, - option_selected=s.option_selected, - continue_work=s.continue_work, - ) - ) + messages = [ + BrowserChatComponent._stored_to_chat_message(s) for s in stored + ] await _reply( { diff --git a/app/ui_layer/browser/frontend/src/components/Chat/Chat.tsx b/app/ui_layer/browser/frontend/src/components/Chat/Chat.tsx index 4e729b72..ba03f9d9 100644 --- a/app/ui_layer/browser/frontend/src/components/Chat/Chat.tsx +++ b/app/ui_layer/browser/frontend/src/components/Chat/Chat.tsx @@ -26,7 +26,9 @@ import { selectSessionHasMoreMessages, selectSessionLoadingOlderMessages, selectSessionOldestMessageTimestamp, + selectPendingQuestions, } from '../../store/selectors/messages' +import { QuestionBox } from './QuestionBox' import { selectSessionActivity } from '../../store/selectors/activity' import { selectSessionBusy, selectSessionRunState } from '../../store/selectors/agent' import type { ActionItem, ChatMessage } from '../../types' @@ -160,6 +162,7 @@ export function Chat({ sessionId, placeholder }: ChatProps) { sendCommand, stopSession, sendOptionClick, + sendQuestionAnswer, openFile, openFolder, lastSeenBySession, @@ -178,6 +181,9 @@ export function Chat({ sessionId, placeholder }: ChatProps) { const messages = useAppSelector(state => selectSessionMessages(state, sessionId)) const activity = useAppSelector(state => selectSessionActivity(state, sessionId)) + // Unanswered agent questions (oldest first). The first one is pinned in a + // QuestionBox above the composer; answering/dismissing advances the queue. + const pendingQuestions = useAppSelector(state => selectPendingQuestions(state, sessionId)) const hasMoreMessages = useAppSelector(state => selectSessionHasMoreMessages(state, sessionId)) const loadingOlderMessages = useAppSelector(state => selectSessionLoadingOlderMessages(state, sessionId)) const oldestMessageTimestamp = useAppSelector(state => selectSessionOldestMessageTimestamp(state, sessionId)) @@ -852,6 +858,16 @@ export function Chat({ sessionId, placeholder }: ChatProps) { sendOptionClick(value, messageId, sessionId) }, [navigate, sendOptionClick, sessionId]) + const handleQuestionAnswer = useCallback((value: string) => { + const q = pendingQuestions[0] + if (q) sendQuestionAnswer(q.messageId, value, sessionId) + }, [pendingQuestions, sendQuestionAnswer, sessionId]) + + const handleQuestionDismiss = useCallback(() => { + const q = pendingQuestions[0] + if (q) sendQuestionAnswer(q.messageId, '', sessionId, true) + }, [pendingQuestions, sendQuestionAnswer, sessionId]) + // Reply action from an agent bubble — arm the reply bar and focus the // input so the user can type straight away. const handleChatReply = useCallback((displayName: string, originalContent: string) => { @@ -1354,6 +1370,19 @@ export function Chat({ sessionId, placeholder }: ChatProps) { row inside it ("+" menu on the left, mic/lang + send on the right). The area is width-capped and centered like the timeline. */}