Receiving the answer
This commit is contained in:
@@ -60,27 +60,7 @@ async def run_bot(transport: BaseTransport, runner_args: RunnerArguments):
|
|||||||
# Returning the answer
|
# Returning the answer
|
||||||
if isinstance(transport, SmallWebRTCTransport):
|
if isinstance(transport, SmallWebRTCTransport):
|
||||||
yield {"status": "ANSWER:START"}
|
yield {"status": "ANSWER:START"}
|
||||||
answer_dict = transport._client._webrtc_connection.get_answer()
|
yield {"answer": transport._client._webrtc_connection.get_answer()}
|
||||||
answer_json_str = json.dumps(answer_dict)
|
|
||||||
answer_bytes = answer_json_str.encode("utf-8")
|
|
||||||
encoded_answer = base64.b64encode(answer_bytes).decode("ascii")
|
|
||||||
|
|
||||||
# Break encoded_answer into multiple messages
|
|
||||||
chunk_size = 1000 # Adjust this size based on your needs
|
|
||||||
total_chunks = (
|
|
||||||
len(encoded_answer) + chunk_size - 1
|
|
||||||
) // chunk_size # Calculate total chunks
|
|
||||||
|
|
||||||
for i in range(0, len(encoded_answer), chunk_size):
|
|
||||||
chunk = encoded_answer[i : i + chunk_size]
|
|
||||||
chunk_index = i // chunk_size
|
|
||||||
yield {
|
|
||||||
"answer_chunk": chunk,
|
|
||||||
"chunk_index": chunk_index,
|
|
||||||
"total_chunks": total_chunks,
|
|
||||||
"is_last_chunk": chunk_index == total_chunks - 1,
|
|
||||||
}
|
|
||||||
|
|
||||||
yield {"status": "ANSWER:END"}
|
yield {"status": "ANSWER:END"}
|
||||||
|
|
||||||
stt = DeepgramSTTService(api_key=os.getenv("DEEPGRAM_API_KEY"))
|
stt = DeepgramSTTService(api_key=os.getenv("DEEPGRAM_API_KEY"))
|
||||||
|
|||||||
@@ -62,61 +62,21 @@ async def offer(request: Request):
|
|||||||
runtimeSessionId="user-123456-conversation-111115555",
|
runtimeSessionId="user-123456-conversation-111115555",
|
||||||
)
|
)
|
||||||
|
|
||||||
print(f"Response {response}")
|
|
||||||
|
|
||||||
# 1. Get the StreamingBody object from the 'response' key
|
|
||||||
streaming_body: StreamingBody = response["response"]
|
|
||||||
|
|
||||||
if streaming_body is None:
|
|
||||||
raise HTTPException(500, "Bedrock response stream missing")
|
|
||||||
|
|
||||||
answer_sdp = None
|
answer_sdp = None
|
||||||
|
|
||||||
# --- Start Streaming Loop using iter_lines() ---
|
if "text/event-stream" in response.get("contentType", ""):
|
||||||
for line in streaming_body.iter_lines(chunk_size=10):
|
# Handle streaming response
|
||||||
print(f"Not decoded line: {line}")
|
content = []
|
||||||
|
streaming_body: StreamingBody = response["response"]
|
||||||
|
for line in streaming_body.iter_lines(chunk_size=1):
|
||||||
|
if line:
|
||||||
|
line = line.decode("utf-8")
|
||||||
|
if line.startswith("data: "):
|
||||||
|
line = line[6:]
|
||||||
|
print(f"Received line: {line}")
|
||||||
|
content.append(line)
|
||||||
|
|
||||||
# iter_lines yields bytes, so we must decode it to a string.
|
print("\nComplete response:", "\n".join(content))
|
||||||
try:
|
|
||||||
decoded_line = line.decode("utf-8")
|
|
||||||
except UnicodeDecodeError:
|
|
||||||
print("Warning: Could not decode line from stream.")
|
|
||||||
continue
|
|
||||||
|
|
||||||
# 1. Clean up and strip whitespace
|
|
||||||
processed_line = decoded_line.strip()
|
|
||||||
|
|
||||||
print(f"Processing line: {processed_line}")
|
|
||||||
|
|
||||||
if not processed_line:
|
|
||||||
continue
|
|
||||||
|
|
||||||
# 2. Check for the SSE 'data:' prefix (now comparing string to string)
|
|
||||||
if processed_line.startswith("data:"):
|
|
||||||
# Extract the content *after* "data: "
|
|
||||||
sse_data = processed_line[len("data:") :].strip()
|
|
||||||
print("Extracted SSE Data payload:", sse_data)
|
|
||||||
|
|
||||||
# 3. Parse the extracted payload as JSON
|
|
||||||
try:
|
|
||||||
event = json.loads(sse_data)
|
|
||||||
print("Received event:", event)
|
|
||||||
|
|
||||||
# 4. Check for the 'answer' key
|
|
||||||
if "answer" in event:
|
|
||||||
payload = event["answer"]
|
|
||||||
|
|
||||||
if payload.get("type") == "answer":
|
|
||||||
answer_sdp = payload
|
|
||||||
print("WebRTC answer found. Stopping stream processing.")
|
|
||||||
# Break the line loop immediately
|
|
||||||
break
|
|
||||||
|
|
||||||
except json.JSONDecodeError:
|
|
||||||
print(f"Failed to parse extracted SSE payload as JSON: {sse_data}")
|
|
||||||
pass
|
|
||||||
|
|
||||||
# --- End Streaming Loop ---
|
|
||||||
|
|
||||||
if answer_sdp is None:
|
if answer_sdp is None:
|
||||||
raise HTTPException(500, "Did not find WebRTC answer in agent output")
|
raise HTTPException(500, "Did not find WebRTC answer in agent output")
|
||||||
|
|||||||
Reference in New Issue
Block a user