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
2 changes: 1 addition & 1 deletion backend/druks/chat/bridge.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down
5 changes: 2 additions & 3 deletions backend/druks/chat/channels/slack/channel.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
4 changes: 2 additions & 2 deletions backend/druks/chat/channels/whatsapp/webhooks.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
23 changes: 22 additions & 1 deletion backend/druks/chat/service.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import asyncio
import base64
import json
import logging
from contextlib import suppress
Expand All @@ -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
Expand Down Expand Up @@ -284,14 +286,33 @@ 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"})
return delivered_messages[-1]
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:
Expand Down
25 changes: 13 additions & 12 deletions backend/tests/chat/bridge.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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++) {
Expand Down Expand Up @@ -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));",
Expand Down Expand Up @@ -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");
Expand All @@ -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, "");
Expand Down
4 changes: 2 additions & 2 deletions backend/tests/test_slack.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
19 changes: 15 additions & 4 deletions backend/tests/test_whatsapp.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import base64
import hashlib
import hmac
import json
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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}]
Expand Down
16 changes: 8 additions & 8 deletions docs/chat.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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

Expand Down
2 changes: 1 addition & 1 deletion frontend/src/pages/ChatPage.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -322,7 +322,7 @@ function ConversationThread({ id, summary, onPin, pinPending, draft, onDraft, on
return <article className="chat-turn" key={message.id} aria-label="Message and reply">
<div className={message.isInternal ? 'chat-user is-internal' : 'chat-user'}>
<div className="chat-message-meta">{message.isInternal ? 'Druks' : 'You'} <time dateTime={message.createdAt} title={format.absTime(message.createdAt)}>{format.absTimeCompact(message.createdAt)}</time></div>
<div className="chat-user-body">{message.body}</div>
{message.body && <div className="chat-user-body">{message.body}</div>}
{message.file && <a href={message.file.url} target="_blank" rel="noreferrer">{message.file.name}</a>}
{isQueued && <span className="chat-queued">Queued</span>}
</div>
Expand Down
Loading