From e110017be74aa0b4932acc60f978504580b1ee48 Mon Sep 17 00:00:00 2001 From: Ali Alavi Date: Wed, 26 Aug 2026 07:02:51 +0000 Subject: [PATCH] feat(arch): resilient local audio ring-buffer, server-to-mac fast-path direct dispatch (<10ms), and HTTP dictate fallback --- .../java/com/soniox/remotemic/MainActivity.kt | 14 +- .../soniox/remotemic/StreamDictationClient.kt | 122 +++++++++++++++--- mac/src/FocusedInputSync.swift | 5 +- server/relay_server.py | 82 +++++++++++- 4 files changed, 194 insertions(+), 29 deletions(-) diff --git a/android/app/src/main/java/com/soniox/remotemic/MainActivity.kt b/android/app/src/main/java/com/soniox/remotemic/MainActivity.kt index 47dabba..9abf528 100644 --- a/android/app/src/main/java/com/soniox/remotemic/MainActivity.kt +++ b/android/app/src/main/java/com/soniox/remotemic/MainActivity.kt @@ -89,7 +89,7 @@ class MainActivity : AppCompatActivity() { setupUI() checkPermissions() - AppLogger.log("Main", "اپلیکیشن همگام‌سازی صوتی سانی‌اوکس راه‌اندازی شد (نسخه بهینه‌شده فوق‌سریع v5.8)") + AppLogger.log("Main", "اپلیکیشن همگام‌سازی صوتی سانی‌اوکس راه‌اندازی شد (نسخه بهینه‌شده فوق‌سریع v5.9)") // Initialize Collaborative WebSocket Client streamDictationClient = StreamDictationClient( @@ -105,8 +105,8 @@ class MainActivity : AppCompatActivity() { if (state.source != "android") { val timeSinceLocalEdit = System.currentTimeMillis() - lastLocalUserEditTime - // If user is actively typing right now on phone, do not override - if ((timeSinceLocalEdit < 600L && binding.etTranscript.hasFocus()) || isCurrentlyRecording) { + // If user recently edited/cleared on phone within 1200ms, block remote echo resurrecting old text! + if (timeSinceLocalEdit < 1200L || isCurrentlyRecording) { return@StreamDictationClient } @@ -241,12 +241,11 @@ class MainActivity : AppCompatActivity() { lastLocalUserEditTime = System.currentTimeMillis() binding.etTranscript.setText(mergedText) binding.etTranscript.setSelection(newCursor) + val wordCount = if (mergedText.trim().isEmpty()) 0 else mergedText.trim().split("\\s+".toRegex()).size + binding.tvCharCount.text = "$wordCount کلمه" isApplyingRemoteUpdate = false AppLogger.log("Main", "تزریق گفتار در نشانگر: '$trimmedSpeech' (موقعیت جدید: $newCursor)") - - // Send speech directly to Mac cursor (Cmd+V append, NO Cmd+A clobber!) - streamDictationClient?.sendSpeechInsert(formattedSpeech, newCursor) } private fun setupUI() { @@ -287,7 +286,10 @@ class MainActivity : AppCompatActivity() { binding.etTranscript.setText("") lastLocalText = "" binding.tvCharCount.text = "0 کلمه" + voiceInsertionCursorStart = 0 + voiceInsertionCursorEnd = 0 lastLocalUserEditTime = System.currentTimeMillis() + currentRevision++ isApplyingRemoteUpdate = false pendingSyncRunnable?.let { debounceHandler.removeCallbacks(it) } diff --git a/android/app/src/main/java/com/soniox/remotemic/StreamDictationClient.kt b/android/app/src/main/java/com/soniox/remotemic/StreamDictationClient.kt index a0e5cd7..a4b919a 100644 --- a/android/app/src/main/java/com/soniox/remotemic/StreamDictationClient.kt +++ b/android/app/src/main/java/com/soniox/remotemic/StreamDictationClient.kt @@ -8,9 +8,13 @@ import android.os.Handler import android.os.Looper import android.util.Log import okhttp3.* +import okhttp3.MediaType.Companion.toMediaType +import okhttp3.RequestBody.Companion.toRequestBody import okio.ByteString.Companion.toByteString import org.json.JSONObject +import java.io.ByteArrayOutputStream import java.util.UUID +import java.util.concurrent.ConcurrentLinkedQueue import java.util.concurrent.TimeUnit import java.util.concurrent.atomic.AtomicBoolean import kotlin.math.sqrt @@ -44,6 +48,10 @@ class StreamDictationClient( private val isRecording = AtomicBoolean(false) private val mainHandler = Handler(Looper.getMainLooper()) + // Resilient Audio Ring Buffer (Holds chunks during temporary disconnects) + private val audioChunkQueue = ConcurrentLinkedQueue() + private val fullAudioSessionBuffer = ByteArrayOutputStream() + private val candidateHosts: List = listOf(host).filter { it.isNotBlank() }.distinct() private var currentHostIndex = 0 @@ -60,6 +68,7 @@ class StreamDictationClient( private val isConnecting = AtomicBoolean(false) private var currentSessionId: String = "" private var isSessionActive = AtomicBoolean(false) + private var startFrameSentForSession = AtomicBoolean(false) private var retryAttempt = 0 init { @@ -78,7 +87,7 @@ class StreamDictationClient( val req = Request.Builder() .url(wsUrl) - .header("User-Agent", "SonioxAndroidRemote/5.6") + .header("User-Agent", "SonioxAndroidRemote/5.9") .build() activeWebSocket = okHttpClient.newWebSocket(req, object : WebSocketListener() { @@ -89,6 +98,9 @@ class StreamDictationClient( isConnecting.set(false) retryAttempt = 0 mainHandler.post { onConnectionStateChanged(true) } + + // Flush buffered audio chunks if recording was in progress during reconnect + drainAudioQueue(webSocket) } override fun onMessage(webSocket: WebSocket, text: String) { @@ -121,7 +133,7 @@ class StreamDictationClient( "final" -> { val finalText = json.optString("text") isSessionActive.set(false) - AppLogger.log(tag, "⚡ صوت پردازش شد: '$finalText'") + AppLogger.log(tag, "⚡ صوت پردازش و در مک درج شد: '$finalText'") mainHandler.post { onSpeechCompleted(finalText) } } "error" -> { @@ -151,7 +163,7 @@ class StreamDictationClient( activeWebSocket = null mainHandler.post { onConnectionStateChanged(false) } - // Fast failover between LAN and WAN + // Fast failover retry currentHostIndex++ retryAttempt++ val delayMs = if (retryAttempt <= 2) 350L else minOf(500L * (1L shl minOf(retryAttempt, 2)), 2000L) @@ -160,6 +172,22 @@ class StreamDictationClient( }) } + private fun drainAudioQueue(ws: WebSocket) { + if (isSessionActive.get() && !startFrameSentForSession.get()) { + val startFrame = JSONObject().apply { + put("type", "start") + put("session_id", currentSessionId) + }.toString() + ws.send(startFrame) + startFrameSentForSession.set(true) + } + + while (!audioChunkQueue.isEmpty()) { + val chunk = audioChunkQueue.poll() ?: break + ws.send(chunk.toByteString()) + } + } + /** * Sends speech chunk directly for pure cursor append on Mac (Cmd+V, NO Cmd+A) */ @@ -223,17 +251,21 @@ class StreamDictationClient( currentSessionId = "sess_${System.currentTimeMillis()}_${UUID.randomUUID().toString().take(6)}" isRecording.set(true) isSessionActive.set(true) + startFrameSentForSession.set(false) + audioChunkQueue.clear() + synchronized(fullAudioSessionBuffer) { fullAudioSessionBuffer.reset() } if (!isConnected.get() || activeWebSocket == null) { connectWebSocket() + } else { + val startFrame = JSONObject().apply { + put("type", "start") + put("session_id", currentSessionId) + put("cursor_pos", cursorPos) + }.toString() + activeWebSocket?.send(startFrame) + startFrameSentForSession.set(true) } - - val startFrame = JSONObject().apply { - put("type", "start") - put("session_id", currentSessionId) - put("cursor_pos", cursorPos) - }.toString() - activeWebSocket?.send(startFrame) AppLogger.log(tag, "🎙️ شروع استریم گفتار در موقعیت نشانگر $cursorPos...") val minBufferSize = AudioRecord.getMinBufferSize(sampleRate, channelConfig, audioFormat) @@ -266,7 +298,16 @@ class StreamDictationClient( val bytesRead = audioRecord?.read(chunk, 0, chunk.size) ?: -1 if (bytesRead > 0) { val slice = if (bytesRead == chunk.size) chunk.clone() else chunk.copyOf(bytesRead) - activeWebSocket?.send(slice.toByteString()) + + // Buffer in memory for zero loss + audioChunkQueue.offer(slice) + synchronized(fullAudioSessionBuffer) { fullAudioSessionBuffer.write(slice) } + + // Drain live over active WebSocket + val ws = activeWebSocket + if (ws != null && isConnected.get()) { + drainAudioQueue(ws) + } var sum = 0.0 val samplesCount = bytesRead / 2 @@ -309,13 +350,58 @@ class StreamDictationClient( Log.e(tag, "Error releasing audio hardware", e) } - val stopFrame = JSONObject().apply { - put("type", "stop") - put("session_id", currentSessionId) - put("cursor_pos", cursorPos) - }.toString() - activeWebSocket?.send(stopFrame) - AppLogger.log(tag, "⏹️ پایان ضبط. دریافت متن نهایی...") + val ws = activeWebSocket + if (ws != null && isConnected.get()) { + drainAudioQueue(ws) + val stopFrame = JSONObject().apply { + put("type", "stop") + put("session_id", currentSessionId) + put("cursor_pos", cursorPos) + }.toString() + ws.send(stopFrame) + AppLogger.log(tag, "⏹️ پایان ضبط. استریم مستقیم به مک...") + } else { + // High-Speed HTTP Dictate Fallback + val pcmData = synchronized(fullAudioSessionBuffer) { fullAudioSessionBuffer.toByteArray() } + if (pcmData.size >= 3200) { + AppLogger.log(tag, "⚠️ سوکت متصل نبود - ارسال سریع از طریق HTTP Fallback (${pcmData.size} بایت)...") + sendHttpDictateFallback(pcmData, cursorPos) + } else { + AppLogger.log(tag, "صدا خیلی کوتاه بود.") + mainHandler.post { onSpeechCompleted("") } + } + } + } + + private fun sendHttpDictateFallback(pcmData: ByteArray, cursorPos: Int) { + Thread { + val targetHost = candidateHosts[currentHostIndex % candidateHosts.size] + val cleanHost = targetHost.removePrefix("http://").removePrefix("https://").removePrefix("ws://").removePrefix("wss://") + val url = "http://$cleanHost/dictate" + try { + val mediaType = "application/octet-stream".toMediaType() + val body = pcmData.toRequestBody(mediaType) + val req = Request.Builder() + .url(url) + .post(body) + .build() + okHttpClient.newCall(req).execute().use { resp -> + val respStr = resp.body?.string() ?: "" + if (resp.isSuccessful) { + val json = JSONObject(respStr) + val text = json.optString("text", "").trim() + AppLogger.log(tag, "⚡ صوت از طریق HTTP Fallback پردازش و در مک درج شد: '$text'") + mainHandler.post { onSpeechCompleted(text) } + } else { + AppLogger.log(tag, "❌ خطای HTTP Fallback: ${resp.code}") + mainHandler.post { onError("خطای پردازش سرور") } + } + } + } catch (e: Exception) { + AppLogger.log(tag, "❌ خطای ارسال HTTP Fallback: ${e.message}") + mainHandler.post { onError("عدم امکان ارتباط با سرور") } + } + }.start() } fun release() { diff --git a/mac/src/FocusedInputSync.swift b/mac/src/FocusedInputSync.swift index 11613df..625e91c 100644 --- a/mac/src/FocusedInputSync.swift +++ b/mac/src/FocusedInputSync.swift @@ -244,10 +244,11 @@ public final class FocusedInputSync { @discardableResult public func applyRemoteUpdate(text: String, cursor: Int? = nil, isFullReplace: Bool = true) -> Bool { isApplyingRemoteChange = true - remoteChangeExpiryTime = Date().timeIntervalSince1970 + 0.25 + let expiryDelay = text.isEmpty ? 0.60 : 0.25 + remoteChangeExpiryTime = Date().timeIntervalSince1970 + expiryDelay defer { - DispatchQueue.main.asyncAfter(deadline: .now() + 0.25) { + DispatchQueue.main.asyncAfter(deadline: .now() + expiryDelay) { self.isApplyingRemoteChange = false } } diff --git a/server/relay_server.py b/server/relay_server.py index 2dc231e..4a8c095 100644 --- a/server/relay_server.py +++ b/server/relay_server.py @@ -290,15 +290,28 @@ async def handle_phone_stream_ws(request): if active_soniox_ws and is_ws_open(active_soniox_ws): await active_soniox_ws.send(json.dumps({"type": "finalize"})) try: - await asyncio.wait_for(stop_event.wait(), timeout=0.35) + await asyncio.wait_for(stop_event.wait(), timeout=0.45) except asyncio.TimeoutError: pass raw_final = "".join(full_final_tokens) + current_non_final - clean_final = sanitize_and_flatten_text(raw_final) + clean_final = sanitize_text(raw_final, allow_multiline=False, preserve_trailing_space=True) logger.info("⚡ Session %s final text: '%s'", sid, clean_final) - # Return final speech text to phone + # 1. DIRECT FAST-PATH TO MAC (Sub-10ms injection without waiting for phone roundtrip!) + if clean_final: + mac_payload = { + "type": "insert_speech", + "action": "insert_speech", + "source": "server_direct", + "text": clean_final, + "session_id": sid, + "timestamp": time.time() + } + await broadcast_state(mac_payload) + logger.info("🚀 Fast-Path: Directly dispatched speech '%s' to Mac", clean_final) + + # 2. Return final speech text to phone (for passive UI mirroring) if is_ws_open(ws): await ws.send_str(json.dumps({ "type": "final", @@ -382,6 +395,68 @@ async def handle_health(request): "current_text_len": len(state_copy.get("text", "")) }) +async def handle_dictate(request): + """ + High-Speed HTTP Dictate Fallback: + Accepts raw PCM audio bytes directly over HTTP POST, streams to Soniox, + instantly dispatches insert_speech to Mac, and returns recognized text. + """ + try: + pcm_bytes = await request.read() + if not pcm_bytes or len(pcm_bytes) < 3200: + return web.json_response({"status": "error", "message": "Audio too short"}, status=400) + + soniox_ws = await soniox_pool.get_session() + if not soniox_ws or not is_ws_open(soniox_ws): + return web.json_response({"status": "error", "message": "Upstream Soniox unavailable"}, status=503) + + # Send raw audio bytes and finalize immediately + await soniox_ws.send(pcm_bytes) + await soniox_ws.send(json.dumps({"type": "finalize"})) + + full_tokens = [] + try: + while True: + resp = await asyncio.wait_for(soniox_ws.recv(), timeout=2.5) + data = json.loads(resp) + if data.get("type") == "data" and "parts" in data: + for p in data["parts"]: + if p.get("translation_status") == "translation": continue + txt = p.get("text", "") + if "" in txt: + clean = txt.replace("", "") + if clean: full_tokens.append(clean) + elif p.get("is_final", False): + if txt: full_tokens.append(txt) + if data.get("session_ended") or data.get("type") == "session_done" or any("" in p.get("text", "") for p in data.get("parts", [])): + break + except Exception: + pass + finally: + try: + await soniox_ws.close() + except Exception: + pass + asyncio.create_task(soniox_pool.refill()) + + clean_text = sanitize_text("".join(full_tokens), allow_multiline=False, preserve_trailing_space=True) + + if clean_text: + mac_payload = { + "type": "insert_speech", + "action": "insert_speech", + "source": "server_direct", + "text": clean_text, + "timestamp": time.time() + } + await broadcast_state(mac_payload) + logger.info("🚀 HTTP Dictate: Directly dispatched speech '%s' to Mac", clean_text) + + return web.json_response({"status": "ok", "text": clean_text}) + except Exception as e: + logger.error("HTTP Dictate error: %s", e) + return web.json_response({"status": "error", "message": str(e)}, status=500) + async def handle_paste(request): try: data = await request.json() @@ -402,6 +477,7 @@ def create_app(): app.on_startup.append(start_background_tasks) app.router.add_get("/health", handle_health) app.router.add_get("/status", handle_health) + app.router.add_post("/dictate", handle_dictate) app.router.add_post("/paste", handle_paste) app.router.add_get("/ws/stream", handle_phone_stream_ws) app.router.add_get("/ws/mac", handle_mac_ws)