Browse Source

feat(arch): resilient local audio ring-buffer, server-to-mac fast-path direct dispatch (<10ms), and HTTP dictate fallback

main
Ali Alavi 20 hours ago
parent
commit
e110017be7
  1. 14
      android/app/src/main/java/com/soniox/remotemic/MainActivity.kt
  2. 122
      android/app/src/main/java/com/soniox/remotemic/StreamDictationClient.kt
  3. 5
      mac/src/FocusedInputSync.swift
  4. 82
      server/relay_server.py

14
android/app/src/main/java/com/soniox/remotemic/MainActivity.kt

@ -89,7 +89,7 @@ class MainActivity : AppCompatActivity() {
setupUI() setupUI()
checkPermissions() checkPermissions()
AppLogger.log("Main", "اپلیکیشن همگام‌سازی صوتی سانی‌اوکس راه‌اندازی شد (نسخه بهینه‌شده فوق‌سریع v5.8)")
AppLogger.log("Main", "اپلیکیشن همگام‌سازی صوتی سانی‌اوکس راه‌اندازی شد (نسخه بهینه‌شده فوق‌سریع v5.9)")
// Initialize Collaborative WebSocket Client // Initialize Collaborative WebSocket Client
streamDictationClient = StreamDictationClient( streamDictationClient = StreamDictationClient(
@ -105,8 +105,8 @@ class MainActivity : AppCompatActivity() {
if (state.source != "android") { if (state.source != "android") {
val timeSinceLocalEdit = System.currentTimeMillis() - lastLocalUserEditTime 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 return@StreamDictationClient
} }
@ -241,12 +241,11 @@ class MainActivity : AppCompatActivity() {
lastLocalUserEditTime = System.currentTimeMillis() lastLocalUserEditTime = System.currentTimeMillis()
binding.etTranscript.setText(mergedText) binding.etTranscript.setText(mergedText)
binding.etTranscript.setSelection(newCursor) binding.etTranscript.setSelection(newCursor)
val wordCount = if (mergedText.trim().isEmpty()) 0 else mergedText.trim().split("\\s+".toRegex()).size
binding.tvCharCount.text = "$wordCount کلمه"
isApplyingRemoteUpdate = false isApplyingRemoteUpdate = false
AppLogger.log("Main", "تزریق گفتار در نشانگر: '$trimmedSpeech' (موقعیت جدید: $newCursor)") 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() { private fun setupUI() {
@ -287,7 +286,10 @@ class MainActivity : AppCompatActivity() {
binding.etTranscript.setText("") binding.etTranscript.setText("")
lastLocalText = "" lastLocalText = ""
binding.tvCharCount.text = "0 کلمه" binding.tvCharCount.text = "0 کلمه"
voiceInsertionCursorStart = 0
voiceInsertionCursorEnd = 0
lastLocalUserEditTime = System.currentTimeMillis() lastLocalUserEditTime = System.currentTimeMillis()
currentRevision++
isApplyingRemoteUpdate = false isApplyingRemoteUpdate = false
pendingSyncRunnable?.let { debounceHandler.removeCallbacks(it) } pendingSyncRunnable?.let { debounceHandler.removeCallbacks(it) }

122
android/app/src/main/java/com/soniox/remotemic/StreamDictationClient.kt

@ -8,9 +8,13 @@ import android.os.Handler
import android.os.Looper import android.os.Looper
import android.util.Log import android.util.Log
import okhttp3.* import okhttp3.*
import okhttp3.MediaType.Companion.toMediaType
import okhttp3.RequestBody.Companion.toRequestBody
import okio.ByteString.Companion.toByteString import okio.ByteString.Companion.toByteString
import org.json.JSONObject import org.json.JSONObject
import java.io.ByteArrayOutputStream
import java.util.UUID import java.util.UUID
import java.util.concurrent.ConcurrentLinkedQueue
import java.util.concurrent.TimeUnit import java.util.concurrent.TimeUnit
import java.util.concurrent.atomic.AtomicBoolean import java.util.concurrent.atomic.AtomicBoolean
import kotlin.math.sqrt import kotlin.math.sqrt
@ -44,6 +48,10 @@ class StreamDictationClient(
private val isRecording = AtomicBoolean(false) private val isRecording = AtomicBoolean(false)
private val mainHandler = Handler(Looper.getMainLooper()) private val mainHandler = Handler(Looper.getMainLooper())
// Resilient Audio Ring Buffer (Holds chunks during temporary disconnects)
private val audioChunkQueue = ConcurrentLinkedQueue<ByteArray>()
private val fullAudioSessionBuffer = ByteArrayOutputStream()
private val candidateHosts: List<String> = listOf(host).filter { it.isNotBlank() }.distinct() private val candidateHosts: List<String> = listOf(host).filter { it.isNotBlank() }.distinct()
private var currentHostIndex = 0 private var currentHostIndex = 0
@ -60,6 +68,7 @@ class StreamDictationClient(
private val isConnecting = AtomicBoolean(false) private val isConnecting = AtomicBoolean(false)
private var currentSessionId: String = "" private var currentSessionId: String = ""
private var isSessionActive = AtomicBoolean(false) private var isSessionActive = AtomicBoolean(false)
private var startFrameSentForSession = AtomicBoolean(false)
private var retryAttempt = 0 private var retryAttempt = 0
init { init {
@ -78,7 +87,7 @@ class StreamDictationClient(
val req = Request.Builder() val req = Request.Builder()
.url(wsUrl) .url(wsUrl)
.header("User-Agent", "SonioxAndroidRemote/5.6")
.header("User-Agent", "SonioxAndroidRemote/5.9")
.build() .build()
activeWebSocket = okHttpClient.newWebSocket(req, object : WebSocketListener() { activeWebSocket = okHttpClient.newWebSocket(req, object : WebSocketListener() {
@ -89,6 +98,9 @@ class StreamDictationClient(
isConnecting.set(false) isConnecting.set(false)
retryAttempt = 0 retryAttempt = 0
mainHandler.post { onConnectionStateChanged(true) } mainHandler.post { onConnectionStateChanged(true) }
// Flush buffered audio chunks if recording was in progress during reconnect
drainAudioQueue(webSocket)
} }
override fun onMessage(webSocket: WebSocket, text: String) { override fun onMessage(webSocket: WebSocket, text: String) {
@ -121,7 +133,7 @@ class StreamDictationClient(
"final" -> { "final" -> {
val finalText = json.optString("text") val finalText = json.optString("text")
isSessionActive.set(false) isSessionActive.set(false)
AppLogger.log(tag, "⚡ صوت پردازش شد: '$finalText'")
AppLogger.log(tag, "⚡ صوت پردازش و در مک درج شد: '$finalText'")
mainHandler.post { onSpeechCompleted(finalText) } mainHandler.post { onSpeechCompleted(finalText) }
} }
"error" -> { "error" -> {
@ -151,7 +163,7 @@ class StreamDictationClient(
activeWebSocket = null activeWebSocket = null
mainHandler.post { onConnectionStateChanged(false) } mainHandler.post { onConnectionStateChanged(false) }
// Fast failover between LAN and WAN
// Fast failover retry
currentHostIndex++ currentHostIndex++
retryAttempt++ retryAttempt++
val delayMs = if (retryAttempt <= 2) 350L else minOf(500L * (1L shl minOf(retryAttempt, 2)), 2000L) 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) * 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)}" currentSessionId = "sess_${System.currentTimeMillis()}_${UUID.randomUUID().toString().take(6)}"
isRecording.set(true) isRecording.set(true)
isSessionActive.set(true) isSessionActive.set(true)
startFrameSentForSession.set(false)
audioChunkQueue.clear()
synchronized(fullAudioSessionBuffer) { fullAudioSessionBuffer.reset() }
if (!isConnected.get() || activeWebSocket == null) { if (!isConnected.get() || activeWebSocket == null) {
connectWebSocket() 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...") AppLogger.log(tag, "🎙️ شروع استریم گفتار در موقعیت نشانگر $cursorPos...")
val minBufferSize = AudioRecord.getMinBufferSize(sampleRate, channelConfig, audioFormat) val minBufferSize = AudioRecord.getMinBufferSize(sampleRate, channelConfig, audioFormat)
@ -266,7 +298,16 @@ class StreamDictationClient(
val bytesRead = audioRecord?.read(chunk, 0, chunk.size) ?: -1 val bytesRead = audioRecord?.read(chunk, 0, chunk.size) ?: -1
if (bytesRead > 0) { if (bytesRead > 0) {
val slice = if (bytesRead == chunk.size) chunk.clone() else chunk.copyOf(bytesRead) 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 var sum = 0.0
val samplesCount = bytesRead / 2 val samplesCount = bytesRead / 2
@ -309,13 +350,58 @@ class StreamDictationClient(
Log.e(tag, "Error releasing audio hardware", e) 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() { fun release() {

5
mac/src/FocusedInputSync.swift

@ -244,10 +244,11 @@ public final class FocusedInputSync {
@discardableResult @discardableResult
public func applyRemoteUpdate(text: String, cursor: Int? = nil, isFullReplace: Bool = true) -> Bool { public func applyRemoteUpdate(text: String, cursor: Int? = nil, isFullReplace: Bool = true) -> Bool {
isApplyingRemoteChange = true isApplyingRemoteChange = true
remoteChangeExpiryTime = Date().timeIntervalSince1970 + 0.25
let expiryDelay = text.isEmpty ? 0.60 : 0.25
remoteChangeExpiryTime = Date().timeIntervalSince1970 + expiryDelay
defer { defer {
DispatchQueue.main.asyncAfter(deadline: .now() + 0.25) {
DispatchQueue.main.asyncAfter(deadline: .now() + expiryDelay) {
self.isApplyingRemoteChange = false self.isApplyingRemoteChange = false
} }
} }

82
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): 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:
await asyncio.wait_for(stop_event.wait(), timeout=0.35)
await asyncio.wait_for(stop_event.wait(), timeout=0.45)
except asyncio.TimeoutError: except asyncio.TimeoutError:
pass pass
raw_final = "".join(full_final_tokens) + current_non_final 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) 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): if is_ws_open(ws):
await ws.send_str(json.dumps({ await ws.send_str(json.dumps({
"type": "final", "type": "final",
@ -382,6 +395,68 @@ async def handle_health(request):
"current_text_len": len(state_copy.get("text", "")) "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 "<fin>" in txt:
clean = txt.replace("<fin>", "")
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("<fin>" 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): async def handle_paste(request):
try: try:
data = await request.json() data = await request.json()
@ -402,6 +477,7 @@ def create_app():
app.on_startup.append(start_background_tasks) app.on_startup.append(start_background_tasks)
app.router.add_get("/health", handle_health) app.router.add_get("/health", handle_health)
app.router.add_get("/status", 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_post("/paste", handle_paste)
app.router.add_get("/ws/stream", handle_phone_stream_ws) app.router.add_get("/ws/stream", handle_phone_stream_ws)
app.router.add_get("/ws/mac", handle_mac_ws) app.router.add_get("/ws/mac", handle_mac_ws)

Loading…
Cancel
Save