1515from druks .accounts .enums import AccountKind
1616from druks .apps .loader import get_app
1717from druks .apps .registry import channels
18+ from druks .core .exceptions import TranscriptionError
19+ from druks .core .services import SpeechToText
1820from druks .durable .engine import step_session
1921from druks .durable .models import Run
2022from druks .files .constants import MAX_UPLOAD_BYTES
3840from druks .sandbox .layout import get_remote_home , get_work_root
3941from druks .sandbox .models import SandboxIdentity , SecretRef
4042from druks .sandbox .templates import get_template_id
43+ from druks .services .exceptions import ServiceNotConnectedError
4144from druks .workspaces import Workspace
4245
4346from .bots .constants import ADMIN_PROMPT , ADMIN_TOOLS
4851 FAILURE_MESSAGE ,
4952 INTERNAL_MESSAGES_PROMPT ,
5053 RESULT_MESSAGE ,
54+ TRANSCRIPTION_FAILED_MESSAGE ,
55+ VOICE_NOTE_MARKER ,
5156)
5257from .enums import MessageRole , MessageState
5358from .exceptions import ChatBridgeError , ChatHarnessError , ChatSandboxGone
@@ -235,8 +240,8 @@ async def send_turn(
235240 config : AgentConfig ,
236241 prompt : str ,
237242) -> Message | None :
238- """Start the agent and send the pending messages.
239- Return the turn's message, unless a Stop or pause came first."""
243+ """Start the agent, transcribe the pending voice notes, and send the pending
244+ messages. Return the turn's message, unless a Stop or pause came first."""
240245 host = bridge .host
241246 status = await bridge .request ("status" , conversationId = conversation .id )
242247 if status ["status" ] == "running" :
@@ -274,13 +279,32 @@ async def send_turn(
274279 expires_at = Base .utc_now () + timedelta (seconds = SANDBOX_HOST_LEASE_SECONDS )
275280 await sandbox_client .set_expiry (host_id = host .id , expires_at = expires_at )
276281 identity .expires_at = expires_at
282+ # The account's other conversations write this row too. Release it before the
283+ # transcription calls.
284+ await session .commit ()
277285 messages = [message ]
278286 if conversation .connection :
279287 # The sandbox can take seconds to start, and a person can take the chat over meanwhile.
280288 if await conversation .is_held (session ):
281289 return
282290 messages = await conversation .list_pending_messages (session )
283- delivered_messages = [pending for pending in messages if await pending .mark_delivered (session )]
291+ notes = []
292+ for pending in messages :
293+ file = pending .file
294+ if file and file .content_type .startswith ("audio/" ):
295+ try :
296+ pending .transcript = await get_transcript (session , file )
297+ except (TranscriptionError , ServiceNotConnectedError ) as error :
298+ logger .warning ("Chat message %s has no transcript: %s" , pending .id , error )
299+ notes .append (
300+ await conversation .create_message (
301+ session , TRANSCRIPTION_FAILED_MESSAGE , is_internal = True
302+ )
303+ )
304+ # The notes are the newest messages, so they close the turn.
305+ delivered_messages = [
306+ pending for pending in (* messages , * notes ) if await pending .mark_delivered (session )
307+ ]
284308 await session .commit ()
285309 if delivered_messages :
286310 await bridge .request (
@@ -295,16 +319,33 @@ async def send_turn(
295319 return
296320
297321
322+ async def get_transcript (session : AsyncSession , file : File ) -> str :
323+ """The words in a voice note, from the Speech To Text card. Druks refuses a note
324+ over the upload cap before the call."""
325+ if file .size > MAX_UPLOAD_BYTES :
326+ raise TranscriptionError (
327+ f"The voice note is { file .size } bytes. The cap is { MAX_UPLOAD_BYTES } bytes."
328+ )
329+ content = await asyncio .to_thread (get_file_storage ().open , file .id )
330+ return await SpeechToText .transcribe (
331+ session , name = file .name , content_type = file .content_type , content = content
332+ )
333+
334+
298335async def get_turn_content (
299336 host : Host , conversation_root : str , messages : list [Message ]
300337) -> list [dict ]:
301- """The ACP content blocks the agent reads: each message's text, then its file. An
302- image travels in the prompt. Audio adds nothing. Any other file goes to the
303- conversation's folder in the sandbox, and the agent gets a link to it."""
338+ """The ACP content blocks the agent reads: each message's text, with the words of
339+ its voice note under a marker, then its file. An image travels in the prompt. Audio
340+ adds nothing more. Any other file goes to the conversation's folder in the sandbox,
341+ and the agent gets a link to it."""
304342 content = []
305343 for message in messages :
306- if message .body :
307- content .append ({"type" : "text" , "text" : message .body })
344+ parts = [message .body ]
345+ if message .transcript :
346+ parts += [VOICE_NOTE_MARKER , message .transcript ]
347+ if text := "\n " .join (part for part in parts if part ):
348+ content .append ({"type" : "text" , "text" : text })
308349 if file := message .file :
309350 if file .content_type .startswith ("image/" ) and file .size <= MAX_UPLOAD_BYTES :
310351 image = await asyncio .to_thread (get_file_storage ().open , file .id )
0 commit comments