Trying to fix the server
This commit is contained in:
@@ -68,67 +68,85 @@ async def offer(request: Request):
|
|||||||
# 1. Get the StreamingBody object from the 'response' key
|
# 1. Get the StreamingBody object from the 'response' key
|
||||||
streaming_body: StreamingBody = response["response"]
|
streaming_body: StreamingBody = response["response"]
|
||||||
|
|
||||||
# Storage for the answer you want
|
if streaming_body is None:
|
||||||
|
raise HTTPException(500, "Bedrock response stream missing")
|
||||||
|
|
||||||
answer_sdp = None
|
answer_sdp = None
|
||||||
|
buffer = "" # Storage for incomplete lines across chunks
|
||||||
|
|
||||||
# Storage for incomplete lines across chunks
|
# Flag to break the outer chunk loop once the answer is found
|
||||||
buffer = ""
|
answer_found = False
|
||||||
|
|
||||||
# 1. Iterate over the chunks from the streaming body
|
# --- Start Streaming Loop ---
|
||||||
for chunk in streaming_body.iter_chunks():
|
for chunk in streaming_body.iter_chunks():
|
||||||
decoded_chunk = chunk.decode("utf-8")
|
if answer_found:
|
||||||
|
# If we found the answer, break the chunk loop immediately.
|
||||||
|
break
|
||||||
|
|
||||||
# Add the new chunk to the buffer
|
decoded_chunk = chunk.decode("utf-8")
|
||||||
buffer += decoded_chunk
|
buffer += decoded_chunk
|
||||||
|
|
||||||
# Split the buffer into lines. -1 ensures we keep a potentially incomplete line at the end
|
# Use splitlines(keepends=True) for reliable line-by-line processing
|
||||||
lines = buffer.splitlines(keepends=True)
|
lines = buffer.splitlines(keepends=True)
|
||||||
|
buffer = "" # Reset buffer
|
||||||
|
|
||||||
# Clear the buffer
|
# Process all complete lines
|
||||||
buffer = ""
|
|
||||||
|
|
||||||
# 2. Process all complete lines
|
|
||||||
for i, line in enumerate(lines):
|
for i, line in enumerate(lines):
|
||||||
# If the last line of the list does not end with a newline, it's an incomplete line.
|
print(f"Processing line {i}: {line}")
|
||||||
# We keep it in the buffer for the next chunk.
|
|
||||||
|
# Check for incomplete line (last item in lines list without a newline)
|
||||||
if i == len(lines) - 1 and not line.endswith("\n"):
|
if i == len(lines) - 1 and not line.endswith("\n"):
|
||||||
buffer = line
|
buffer = line
|
||||||
continue
|
continue
|
||||||
|
|
||||||
# Remove trailing newline and strip whitespace
|
# 1. Clean up and strip whitespace
|
||||||
processed_line = line.strip()
|
processed_line = line.strip()
|
||||||
|
|
||||||
if not processed_line:
|
if not processed_line:
|
||||||
continue
|
continue
|
||||||
|
|
||||||
print("Received line:", processed_line)
|
print("Raw line received:", processed_line)
|
||||||
|
|
||||||
# 3. Parse the line as a JSON event
|
# 2. Check for the SSE 'data:' prefix
|
||||||
|
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
|
||||||
|
answer_found = True # Set flag
|
||||||
|
print("WebRTC answer found. Stopping stream processing.")
|
||||||
|
# Break the inner line loop
|
||||||
|
break
|
||||||
|
|
||||||
|
except json.JSONDecodeError:
|
||||||
|
print(f"Failed to parse extracted SSE payload as JSON: {sse_data}")
|
||||||
|
pass
|
||||||
|
|
||||||
|
# --- End Streaming Loop ---
|
||||||
|
|
||||||
|
# After the loop, the `StreamingBody` should be exhausted.
|
||||||
|
if answer_sdp is None:
|
||||||
|
# Check buffer only if answer wasn't found (for maximum safety)
|
||||||
|
if buffer.strip().startswith("data:"):
|
||||||
try:
|
try:
|
||||||
event = json.loads(processed_line)
|
sse_data = buffer.strip()[len("data:") :].strip()
|
||||||
print("Received event:", event)
|
event = json.loads(sse_data)
|
||||||
|
if "answer" in event and event["answer"].get("type") == "answer":
|
||||||
# 4. Check for the 'answer' key which contains the WebRTC SDP
|
answer_sdp = event["answer"]
|
||||||
if "answer" in event:
|
except Exception:
|
||||||
# The value of "answer" is the JSON payload from the agent
|
|
||||||
payload = event["answer"]
|
|
||||||
|
|
||||||
# Verify it's the WebRTC answer you are looking for
|
|
||||||
if payload.get("type") == "answer":
|
|
||||||
answer_sdp = payload
|
|
||||||
# Break out of both loops once the answer is found
|
|
||||||
raise StopIteration
|
|
||||||
|
|
||||||
except json.JSONDecodeError:
|
|
||||||
print(f"Failed to parse line as JSON: {processed_line}")
|
|
||||||
pass
|
pass
|
||||||
|
|
||||||
try:
|
|
||||||
# Check if we were stopped by the StopIteration (meaning we found the answer)
|
|
||||||
pass
|
|
||||||
except StopIteration:
|
|
||||||
pass
|
|
||||||
|
|
||||||
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