|
|
@ -1,10 +1,8 @@ |
|
|
#!/usr/bin/env python3 |
|
|
#!/usr/bin/env python3 |
|
|
""" |
|
|
""" |
|
|
Soniox Collaborative Bi-Directional Input Synchronization Gateway (v5.2) |
|
|
|
|
|
- Real-time Google Docs / Figma style collaborative mirroring between Mac and Android. |
|
|
|
|
|
- Ultra-low latency streaming speech recognition with pre-warmed Soniox pool. |
|
|
|
|
|
- Monotonic revision counters, echo-loop suppression, and sub-15ms WebSocket routing. |
|
|
|
|
|
- Process-targeted cursor-aware voice insertion. |
|
|
|
|
|
|
|
|
Soniox Collaborative Bi-Directional Input Synchronization Gateway (v5.3) |
|
|
|
|
|
- Handles all message types: sync_state, insert_speech, speech_insert, update_input, paste. |
|
|
|
|
|
- Sub-15ms WebSocket routing between Android and Mac. |
|
|
""" |
|
|
""" |
|
|
|
|
|
|
|
|
import asyncio |
|
|
import asyncio |
|
|
@ -125,6 +123,10 @@ async def broadcast_state(payload_dict: dict, exclude_ws=None): |
|
|
current_room_state.update(payload_dict) |
|
|
current_room_state.update(payload_dict) |
|
|
payload_str = json.dumps(payload_dict, ensure_ascii=False) |
|
|
payload_str = json.dumps(payload_dict, ensure_ascii=False) |
|
|
|
|
|
|
|
|
|
|
|
logger.info("📡 Broadcasting %s (len: %d) to %d Macs, %d Phones", |
|
|
|
|
|
payload_dict.get("type"), len(payload_dict.get("text", "")), |
|
|
|
|
|
len(connected_mac_websockets), len(connected_phone_websockets)) |
|
|
|
|
|
|
|
|
# 1. Send to Phone clients |
|
|
# 1. Send to Phone clients |
|
|
dead_phones = set() |
|
|
dead_phones = set() |
|
|
for ws in list(connected_phone_websockets): |
|
|
for ws in list(connected_phone_websockets): |
|
|
@ -239,10 +241,10 @@ async def handle_phone_stream_ws(request): |
|
|
msg_type = data.get("type") or data.get("action") |
|
|
msg_type = data.get("type") or data.get("action") |
|
|
sid = data.get("session_id", f"sess_{int(time.time()*1000)}") |
|
|
sid = data.get("session_id", f"sess_{int(time.time()*1000)}") |
|
|
|
|
|
|
|
|
if msg_type == "sync_state" or msg_type == "phone_input_edit" or msg_type == "update_input": |
|
|
|
|
|
# Phone edited text: broadcast to Mac immediately! |
|
|
|
|
|
|
|
|
# ALL text / sync / insert operations must be broadcast to Mac! |
|
|
|
|
|
if msg_type in ("sync_state", "insert_speech", "speech_insert", "phone_input_edit", "update_input", "paste"): |
|
|
data["source"] = "android" |
|
|
data["source"] = "android" |
|
|
data["type"] = "sync_state" |
|
|
|
|
|
|
|
|
data["type"] = msg_type |
|
|
await broadcast_state(data, exclude_ws=ws) |
|
|
await broadcast_state(data, exclude_ws=ws) |
|
|
|
|
|
|
|
|
elif msg_type == "start": |
|
|
elif msg_type == "start": |
|
|
@ -265,7 +267,6 @@ async def handle_phone_stream_ws(request): |
|
|
if active_soniox_ws and is_ws_open(active_soniox_ws): |
|
|
if active_soniox_ws and is_ws_open(active_soniox_ws): |
|
|
await active_soniox_ws.send(json.dumps({"type": "finalize"})) |
|
|
await active_soniox_ws.send(json.dumps({"type": "finalize"})) |
|
|
try: |
|
|
try: |
|
|
# 350ms fast timeout for finalize |
|
|
|
|
|
await asyncio.wait_for(stop_event.wait(), timeout=0.35) |
|
|
await asyncio.wait_for(stop_event.wait(), timeout=0.35) |
|
|
except asyncio.TimeoutError: |
|
|
except asyncio.TimeoutError: |
|
|
pass |
|
|
pass |
|
|
@ -350,7 +351,7 @@ async def handle_health(request): |
|
|
state_copy = dict(current_room_state) |
|
|
state_copy = dict(current_room_state) |
|
|
return web.json_response({ |
|
|
return web.json_response({ |
|
|
"status": "ok", |
|
|
"status": "ok", |
|
|
"service": "Soniox Collaborative Sync Gateway v5.2", |
|
|
|
|
|
|
|
|
"service": "Soniox Collaborative Sync Gateway v5.3", |
|
|
"connected_macs": len(connected_mac_websockets), |
|
|
"connected_macs": len(connected_mac_websockets), |
|
|
"connected_phones": len(connected_phone_websockets), |
|
|
"connected_phones": len(connected_phone_websockets), |
|
|
"current_app": state_copy.get("app", ""), |
|
|
"current_app": state_copy.get("app", ""), |
|
|
@ -364,7 +365,7 @@ async def handle_paste(request): |
|
|
text = data.get("text", "") |
|
|
text = data.get("text", "") |
|
|
cursor = data.get("cursor_pos") or data.get("cursor") |
|
|
cursor = data.get("cursor_pos") or data.get("cursor") |
|
|
data["source"] = "http_post" |
|
|
data["source"] = "http_post" |
|
|
data["type"] = "sync_state" |
|
|
|
|
|
|
|
|
data["type"] = "update_input" |
|
|
await broadcast_state(data) |
|
|
await broadcast_state(data) |
|
|
return web.json_response({"status": "synced", "revision": current_room_state["revision"]}) |
|
|
return web.json_response({"status": "synced", "revision": current_room_state["revision"]}) |
|
|
except Exception as e: |
|
|
except Exception as e: |
|
|
|