diff --git a/backend/druks/chat/bridge.mjs b/backend/druks/chat/bridge.mjs index 5db3133d..ec7eb769 100644 --- a/backend/druks/chat/bridge.mjs +++ b/backend/druks/chat/bridge.mjs @@ -212,7 +212,7 @@ class Conversation { connection.cancel({ sessionId: this.state.sessionId }).catch(console.error); }, request.timeout * 1000); try { - const result = await connection.prompt({ sessionId: this.state.sessionId, prompt: [{ type: "text", text: request.body }] }); + const result = await connection.prompt({ sessionId: this.state.sessionId, prompt: request.content }); await this.archive(); this.state.stopReason = result.stopReason; if (timedOut) this.state.status = "interrupted"; diff --git a/backend/druks/chat/channels/slack/channel.py b/backend/druks/chat/channels/slack/channel.py index 6fdda8a5..aee0aacf 100644 --- a/backend/druks/chat/channels/slack/channel.py +++ b/backend/druks/chat/channels/slack/channel.py @@ -109,11 +109,10 @@ async def save_message( text = message["text"].replace(f"<@{card.identity['bot_user_id']}>", "").strip() source_id = f"{message['channel']}:{message['ts']}" if files: - body = text or files[0].name - await conversation.create_message(session, body, source_id=source_id, file=files[0]) + await conversation.create_message(session, text, source_id=source_id, file=files[0]) for shared, file in zip(shares[1:], files[1:], strict=True): await conversation.create_message( - session, file.name, source_id=f"{source_id}:{shared['id']}", file=file + session, "", source_id=f"{source_id}:{shared['id']}", file=file ) elif text: await conversation.create_message(session, text, source_id=source_id) diff --git a/backend/druks/chat/channels/whatsapp/webhooks.py b/backend/druks/chat/channels/whatsapp/webhooks.py index 626028ed..247aa0b2 100644 --- a/backend/druks/chat/channels/whatsapp/webhooks.py +++ b/backend/druks/chat/channels/whatsapp/webhooks.py @@ -100,9 +100,9 @@ async def save_message( ) -> None: """Save the message and start its turn. A message with no text and no media, such as a shared location, gives the agent nothing to read.""" - message = self.data["payload"] + body = self.data["payload"]["body"] or "" file = await self.save_media(session, conversation) - if body := message["body"] or (file.name if file else ""): + if body or file: await conversation.create_message(session, body, source_id=key, file=file) await session.commit() await DBOS.start_workflow_async(deliver, conversation.id) diff --git a/backend/druks/chat/service.py b/backend/druks/chat/service.py index 60a780ce..5ebefebd 100644 --- a/backend/druks/chat/service.py +++ b/backend/druks/chat/service.py @@ -1,4 +1,5 @@ import asyncio +import base64 import json import logging from contextlib import suppress @@ -15,6 +16,7 @@ from druks.apps.registry import channels from druks.durable.engine import step_session from druks.durable.models import Run +from druks.files.constants import MAX_UPLOAD_BYTES from druks.files.datastructures import File from druks.files.storage import get_file_storage from druks.harnesses.base import Harness @@ -284,7 +286,7 @@ async def send_turn( "prompt", conversationId=conversation.id, messageId=delivered_messages[-1].id, - body="\n\n".join(pending.body for pending in delivered_messages), + content=await get_turn_content(delivered_messages), timeout=timeout, ) await publish(conversation.id, {"type": "messages"}) @@ -292,6 +294,25 @@ async def send_turn( return +async def get_turn_content(messages: list[Message]) -> list[dict]: + """The ACP content blocks the agent reads: each message's text, then its image.""" + content = [] + for message in messages: + if message.body: + content.append({"type": "text", "text": message.body}) + file = message.file + if file and file.content_type.startswith("image/") and file.size <= MAX_UPLOAD_BYTES: + image = await asyncio.to_thread(get_file_storage().open, file.id) + content.append( + { + "type": "image", + "data": base64.b64encode(image).decode(), + "mimeType": file.content_type, + } + ) + return content + + async def follow_turn( session: AsyncSession, conversation: Conversation, message: Message, bridge: Bridge ) -> None: diff --git a/backend/tests/chat/bridge.test.mjs b/backend/tests/chat/bridge.test.mjs index 6260c4f4..7706fbf2 100644 --- a/backend/tests/chat/bridge.test.mjs +++ b/backend/tests/chat/bridge.test.mjs @@ -32,6 +32,7 @@ async function request(port, values, abandon = false) { const id = number => `00000000-0000-7000-8000-${String(number).padStart(12, "0")}`; const message = number => id(100 + number); +const text = value => [{ type: "text", text: value }]; async function until(action, matches) { for (let attempt = 0; attempt < 200; attempt++) { @@ -71,7 +72,7 @@ test("The bridge streams detached turns, isolates archives, cancels, and reloads "prompt: async p => {", " const permission = await client.requestPermission({sessionId,toolCall:{toolCallId:'permission',title:'Read'},options:[{optionId:'no',name:'Deny',kind:'reject_once'},{optionId:'yes',name:'Allow',kind:'allow_once'}]});", " if (permission.outcome.optionId !== 'yes') throw Error('permission');", - " const text = p.prompt[0].text;", + " const text = p.prompt.map(block => block.text ?? block.type).join(' ');", " await client.sessionUpdate({sessionId,update:{sessionUpdate:'agent_message_chunk',content:{type:'text',text:memory+text}}});", " if(text==='wait') await new Promise(resolve => { release=resolve; });", " else await new Promise(resolve => setTimeout(resolve, 100));", @@ -131,27 +132,27 @@ test("The bridge streams detached turns, isolates archives, cancels, and reloads // The session opens on its model at session/new; the model option only switches it. assert.deepEqual(settings.trim().split("\n"), ["effort=high", "fast=on", "model=claude-sonnet-5", "effort=low", "fast=off"]); - await request(port, { method: "prompt", conversationId: id(1), messageId: message(1), body: "remember-one" }, true); + await request(port, { method: "prompt", conversationId: id(1), messageId: message(1), content: [...text("remember-one"), { type: "image", data: "aGk=", mimeType: "image/png" }] }, true); await until(() => request(port, { method: "events", conversationId: id(1), after: 0 }), result => result.events.length > 0); const saved = await until(() => status(id(1)), result => result.status === "replied"); - await request(port, { method: "prompt", conversationId: id(2), messageId: message(2), body: "private-two" }); + await request(port, { method: "prompt", conversationId: id(2), messageId: message(2), content: text("private-two") }); await until(() => status(id(2)), result => result.status === "replied"); - await request(port, { method: "prompt", conversationId: id(2), messageId: message(7), body: "wait", timeout: 0.2 }); + await request(port, { method: "prompt", conversationId: id(2), messageId: message(7), content: text("wait"), timeout: 0.2 }); await until(() => status(id(2)), result => result.status === "interrupted"); - await request(port, { method: "prompt", conversationId: id(1), messageId: message(3), body: "wait" }); + await request(port, { method: "prompt", conversationId: id(1), messageId: message(3), content: text("wait") }); await until(() => status(id(1)), result => result.status === "running"); - assert.equal((await request(port, { method: "prompt", conversationId: id(1), messageId: message(4), body: "duplicate" })).ok, false); + assert.equal((await request(port, { method: "prompt", conversationId: id(1), messageId: message(4), content: text("duplicate") })).ok, false); await until(() => request(port, { method: "events", conversationId: id(1), after: 0 }), result => result.events.some(event => event.messageId === message(3))); assert.equal((await request(port, { method: "events", conversationId: id(1), after: 0 })).events.some(event => event.messageId === message(1)), false); await request(port, { method: "cancel", conversationId: id(1), messageId: message(3) }); assert.equal((await until(() => status(id(1)), result => result.status === "cancelled")).stopReason, "cancelled"); const partial = await request(port, { method: "events", conversationId: id(1), after: 0 }); assert.match(partial.events.find(event => event.messageId === message(3)).notification.update.content.text, /wait/); - assert.equal((await request(port, { method: "prompt", conversationId: id(1), messageId: message(3), body: "again" })).status, "cancelled"); + assert.equal((await request(port, { method: "prompt", conversationId: id(1), messageId: message(3), content: text("again") })).status, "cancelled"); await request(port, { method: "cancel", conversationId: id(1), messageId: message(4) }); - assert.equal((await request(port, { method: "prompt", conversationId: id(1), messageId: message(4), body: "never send" })).status, "cancelled"); + assert.equal((await request(port, { method: "prompt", conversationId: id(1), messageId: message(4), content: text("never send") })).status, "cancelled"); assert.equal((await request(port, { method: "events", conversationId: id(1), after: 0 })).events.some(event => event.messageId === message(4)), false); const archive = path.join(temporary, "session.tar.gz"); @@ -164,20 +165,20 @@ test("The bridge streams detached turns, isolates archives, cancels, and reloads const second = await launch(home); assert.equal((await request(port, { ...start(id(1)), archivePath: archive })).ok, true); assert.equal((await request(port, { method: "events", conversationId: id(1), after: 0 })).events.length, 0); - await request(port, { method: "prompt", conversationId: id(1), messageId: message(5), body: "continue" }); + await request(port, { method: "prompt", conversationId: id(1), messageId: message(5), content: text("continue") }); await until(() => status(id(1)), result => result.status === "replied"); const events = (await request(port, { method: "events", conversationId: id(1), after: 0 })).events; - assert.match(events[0].notification.update.content.text, /remember-one/); + assert.match(events[0].notification.update.content.text, /remember-one image/); assert.doesNotMatch(events[0].notification.update.content.text, /private-two|replay/); assert.equal((await status(id(2))).status, "missing"); // An archive from another harness stays unread: the conversation opens a fresh session. assert.equal((await request(port, { ...start(id(3)), archivePath: archive, harness: "other" })).ok, true); - await request(port, { method: "prompt", conversationId: id(3), messageId: message(8), body: "fresh" }); + await request(port, { method: "prompt", conversationId: id(3), messageId: message(8), content: text("fresh") }); await until(() => status(id(3)), result => result.status === "replied"); const fresh = (await request(port, { method: "events", conversationId: id(3), after: 0 })).events; assert.equal(fresh[0].notification.update.content.text, "fresh"); - await request(port, { method: "prompt", conversationId: id(1), messageId: message(6), body: "wait" }); + await request(port, { method: "prompt", conversationId: id(1), messageId: message(6), content: text("wait") }); await request(port, { method: "cancel", conversationId: id(1), messageId: message(5) }); assert.equal((await status(id(1))).status, "running"); assert.equal((await status(id(1))).stopReason, ""); diff --git a/backend/tests/test_slack.py b/backend/tests/test_slack.py index d284cd21..81f04fa5 100644 --- a/backend/tests/test_slack.py +++ b/backend/tests/test_slack.py @@ -486,9 +486,9 @@ async def download(client, url): assert [ (message.body, message.source_id, message.file.name) for message in conversation.messages ] == [ - ("plan.pdf", "D1:1.0", "plan.pdf"), + ("", "D1:1.0", "plan.pdf"), ("both of these", "D1:2.0", "a.pdf"), - ("b.pdf", "D1:2.0:F3", "b.pdf"), + ("", "D1:2.0:F3", "b.pdf"), ] assert [url for _, url in downloads] == [ "https://files.slack.com/F1/plan.pdf", diff --git a/backend/tests/test_whatsapp.py b/backend/tests/test_whatsapp.py index 183ec0ae..4d9aa74b 100644 --- a/backend/tests/test_whatsapp.py +++ b/backend/tests/test_whatsapp.py @@ -1,3 +1,4 @@ +import base64 import hashlib import hmac import json @@ -334,11 +335,17 @@ async def test_a_message_reaches_the_bot_of_its_numbers_app(druks_db, helpdesk, async def test_one_turn_answers_every_pending_message_and_knows_its_own_reply( - druks_db, helpdesk, waha, monkeypatch + druks_db, helpdesk, waha, tmp_path, monkeypatch ): + monkeypatch.setenv("DRUKS_DATA_DIR", str(tmp_path)) + monkeypatch.setattr(WahaClient, "download", AsyncMock(return_value=b"\x89PNG")) connection = await link(druks_db, await bot_account(druks_db)) - for key, body in (("M1", "Hello"), ("M2", "Is my ticket open?"), ("M3", "And the second one?")): - await receive(connection, message_event(ANA, body, key=key)) + photo = message_event(ANA, "", key="M2") + media = {"url": "http://waha.test/api/files/photo.png", "mimetype": "image/png"} + photo["payload"].update(hasMedia=True, media=media) + await receive(connection, message_event(ANA, "Hello", key="M1")) + await receive(connection, photo) + await receive(connection, message_event(ANA, "And the second one?", key="M3")) [conversation] = await Conversation.list_for_connection(druks_db, connection.id) config = SimpleNamespace( harness_class=ClaudeHarness, @@ -386,7 +393,11 @@ async def copy_arrives_first() -> None: await service.deliver_pending(druks_db, conversation) [prompt] = [values for method, values in requests if method == "prompt"] - assert prompt["body"] == "Hello\n\nIs my ticket open?\n\nAnd the second one?" + assert prompt["content"] == [ + {"type": "text", "text": "Hello"}, + {"type": "image", "data": base64.b64encode(b"\x89PNG").decode(), "mimeType": "image/png"}, + {"type": "text", "text": "And the second one?"}, + ] assert prompt["timeout"] == 60 [start] = [values for method, values in requests if method == "start"] assert start["headers"] == [{"name": CONVERSATION_HEADER, "value": conversation.id}] diff --git a/docs/chat.md b/docs/chat.md index 8b244067..d9e371fb 100644 --- a/docs/chat.md +++ b/docs/chat.md @@ -225,12 +225,12 @@ it `interrupted`. A bot or bot admin account holds no credential of its own, so agents and runs bill the default account. All pending messages of a WhatsApp conversation go into one turn. Druks saves -media as a Druks file on its message. Druks drops a repeated message by its -WhatsApp ID. Druks writes to a person only to answer them, one reply for each -turn. It never starts a chat and never sends in bulk. Before it sends a reply, -Druks takes a new message ID from WAHA and records it. WAHA's copy of the sent -message then carries a known ID, even when the copy arrives before the send -returns. +media as a Druks file on its message. An image reaches the agent with its +message. Druks drops a repeated message by its WhatsApp ID. Druks writes to a +person only to answer them, one reply for each turn. It never starts a chat and +never sends in bulk. Before it sends a reply, Druks takes a new message ID from +WAHA and records it. WAHA's copy of the sent message then carries a known ID, +even when the copy arrives before the send returns. Druks also adds **internal messages** to a conversation. Each one comes from a fixed template. An internal message starts a turn like any message, and it @@ -315,8 +315,8 @@ people from other workspaces. A file you send to the bot, in a direct message or in a thread you joined, becomes a Druks file on its message. A message with only a file starts a turn -like any other, with the file's name as its text. Each further file in one -Slack message gets a message of its own. The agent sends no files back. +like any other. An image reaches the agent with its message. Each further file +in one Slack message gets a message of its own. The agent sends no files back. ### Rooms diff --git a/frontend/src/pages/ChatPage.tsx b/frontend/src/pages/ChatPage.tsx index 98ed51fe..97d7a2f3 100644 --- a/frontend/src/pages/ChatPage.tsx +++ b/frontend/src/pages/ChatPage.tsx @@ -322,7 +322,7 @@ function ConversationThread({ id, summary, onPin, pinPending, draft, onDraft, on return
{message.isInternal ? 'Druks' : 'You'}
-
{message.body}
+ {message.body &&
{message.body}
} {message.file && {message.file.name}} {isQueued && Queued}