From cac887c8fd2d97fefefecbfb0446c02b7c536307 Mon Sep 17 00:00:00 2001 From: johnmalek312 Date: Mon, 9 Mar 2026 18:41:14 +1100 Subject: [PATCH] fix: replace workflow step with direct method for message injection WorkflowValidationError: llama-index validates that all consumed events are produced by some step. ExternalUserMessageEvent is external-only, so the @step approach fails validation. Replace with DroidAgent.send_user_message(text) -> QueuedUserMessage. Directly queues on shared state, no event routing needed. Drain points in FastAgent and Manager still emit Applied/Dropped stream events. --- droidrun/agent/droid/droid_agent.py | 40 +++++++++++----------- tests/test_external_message.py | 53 +++++++++++++++++++++++++++++ 2 files changed, 73 insertions(+), 20 deletions(-) create mode 100644 tests/test_external_message.py diff --git a/droidrun/agent/droid/droid_agent.py b/droidrun/agent/droid/droid_agent.py index 24909fd..25aa894 100644 --- a/droidrun/agent/droid/droid_agent.py +++ b/droidrun/agent/droid/droid_agent.py @@ -28,8 +28,6 @@ from droidrun.agent.droid.events import ( ExecutorInputEvent, ExecutorResultEvent, ExternalUserMessageDroppedEvent, - ExternalUserMessageEvent, - ExternalUserMessageQueuedEvent, FastAgentExecuteEvent, FastAgentResultEvent, FinalizeEvent, @@ -595,33 +593,35 @@ class DroidAgent(Workflow): return event # ======================================================================== - # External user message ingestion + # External user message injection # ======================================================================== - @step - async def ingest_external_user_message( - self, ctx: Context, ev: ExternalUserMessageEvent - ) -> None: - """Accept an external user message and queue it in shared state. + def send_user_message(self, message: str) -> "QueuedUserMessage": + """Inject a user message into the running agent loop. - This step runs any time during the workflow. It does NOT touch - message_history directly — the active agent loop drains the queue - at its next safe checkpoint. + Thread-safe to call from any context while the agent is running. + The message is queued in shared state and drained at the next safe + checkpoint (after tool results in direct mode, or at Manager's + prepare_context in reasoning mode). + + Args: + message: The user's message text. + + Returns: + QueuedUserMessage with a unique ID for tracking. + + Raises: + RuntimeError: If the workflow has already completed. """ - queued = self.shared_state.queue_user_message(ev.message) + from droidrun.agent.droid.state import QueuedUserMessage + + queued = self.shared_state.queue_user_message(message) logger.info( f"📩 External user message queued [id={queued.id}] " f"(queue length: {len(self.shared_state.pending_user_messages)})", extra={"color": "cyan"}, ) - ctx.write_event_to_stream( - ExternalUserMessageQueuedEvent( - message_id=queued.id, - message=ev.message, - queue_length=len(self.shared_state.pending_user_messages), - step_number=self.shared_state.step_number, - ) - ) + return queued # ======================================================================== # execute_task — FastAgent / CodeActAgent diff --git a/tests/test_external_message.py b/tests/test_external_message.py new file mode 100644 index 0000000..ece0247 --- /dev/null +++ b/tests/test_external_message.py @@ -0,0 +1,53 @@ +"""Test mid-run external user message injection. + +Run with: python tests/test_external_message.py +Requires a connected Android device and configured LLM in ~/.config/droidrun/config.yaml +""" + +import asyncio + +from droidrun.agent.droid import DroidAgent +from droidrun.agent.droid.events import ( + ExternalUserMessageAppliedEvent, + ExternalUserMessageDroppedEvent, +) +from droidrun.config_manager.loader import ConfigLoader + + +async def main(): + config = ConfigLoader.load() + agent = DroidAgent( + goal="Open the settings app and check the android version", config=config + ) + handler = agent.run() + + async def inject_after_delay(): + await asyncio.sleep(10) + print( + "\n>>> Sending: 'Actually open Chrome and go to google.com'\n" + ) + queued = agent.send_user_message( + "Actually open Chrome and go to google.com" + ) + print(f"[QUEUED] id={queued.id}") + + task = asyncio.create_task(inject_after_delay()) + + async for ev in handler.stream_events(): + if isinstance(ev, ExternalUserMessageAppliedEvent): + print( + f"[APPLIED] ids={ev.message_ids} consumer={ev.consumer} step={ev.step_number}" + ) + elif isinstance(ev, ExternalUserMessageDroppedEvent): + print(f"[DROPPED] ids={ev.message_ids} reason={ev.reason}") + + result = await handler + task.cancel() + + print( + f"\nResult: success={result.success} reason={result.reason} steps={result.steps}" + ) + + +if __name__ == "__main__": + asyncio.run(main())