diff --git a/laguna_core.py b/laguna_core.py index 54dd0d3..a69cb99 100644 --- a/laguna_core.py +++ b/laguna_core.py @@ -77,6 +77,17 @@ OUT_SINK_WAIT_MARGIN_S = 2.0 # folga sobre a duracao da frase ao esperar o sink OUT_SINK_POLL_S = 0.2 # granularidade do stop na thread do sink +# Robustez da saida (issue #54): o sink ja se recupera sozinho de uma falha — o +# stream cai, a frase seguinte reabre — mas o relato para fora era terminal, e a +# UI marcava a direcao inteira como parada (com o Stop desabilitado) por causa de +# um soluco de driver, com o worker ainda traduzindo. Mesma doutrina do orcamento +# de retries da captura (CAPTURE_MAX_RETRIES) e das frases (SEGMENT_MAX_FAILURES): +# falha isolada e recuperavel, so falhas CONSECUTIVAS no MESMO device sao fatais. +# 3 (e nao os 5 da captura) porque aqui cada falha ja custou uma frase muda naquele +# device — o usuario precisa saber cedo. Numero de robustez, nao de latencia: o +# caminho feliz nunca toca este contador. +OUT_SINK_MAX_ERRORS = 3 # falhas consecutivas por device antes do fatal + def _open_output_stream(device: int, samplerate: int): """Abre e inicia o OutputStream mono de um device de saida. @@ -104,19 +115,24 @@ class _OutputSink: uma vez so, ficando bloqueado ~1 frase, nao N. Falha de device e local: a excecao fecha so este stream (a proxima frase - reabre) e vai para `on_error(device, exc)`; os outros sinks seguem tocando. + reabre) e vai para `on_error(device, exc, consecutivas)`; os outros sinks + seguem tocando. O contador de falhas CONSECUTIVAS vive aqui (issue #54) — + e este objeto que sabe se o `write()` deu certo, e um contador por sink ja + e um contador por device: a falha de um nunca contamina o do outro. Quem + decide se `consecutivas` ja e fatal e o worker (`_on_sink_error`). """ def __init__( self, device: int, samplerate: int, - on_error: Callable[[int, Exception], None], + on_error: Callable[[int, Exception, int], None], open_stream: Optional[Callable[[int, int], object]] = None, ) -> None: self.device = int(device) self.samplerate = int(samplerate) self._on_error = on_error + self._erros = 0 # falhas consecutivas; so a thread do sink toca self._open_stream = open_stream or _open_output_stream self._q: "queue.Queue[np.ndarray]" = queue.Queue(maxsize=1) self._done = threading.Event() @@ -184,6 +200,11 @@ def _run(self) -> None: except Exception as e: self._close_stream() # forca reabertura na proxima frase self._report(e) + else: + # Frase tocada: o device esta vivo de novo. Zera o orcamento para + # que falhas ESPACADAS ao longo de uma call (um soluco por hora) + # nunca somem ate o teto e matem a direcao (issue #54). + self._erros = 0 finally: self._done.set() @@ -195,10 +216,15 @@ def _report(self, exc: Exception) -> None: Stop (ou troca de config, que faz `laguna_server` parar o worker antigo) com uma frase tocando acendia `error.play` na UI — e `static/app.js` trata erro como terminal, deixando o painel travado. Parar nao e erro. + + O contador so anda em falha REPORTADA: o write abortado pelo Stop nao + conta, senao o shutdown envenenaria o orcamento de um sink que sequer + vai sobreviver a ele. """ if self._stop.is_set(): return - self._on_error(self.device, exc) + self._erros += 1 + self._on_error(self.device, exc, self._erros) def _ensure_stream(self): with self._lock: @@ -314,6 +340,10 @@ def __init__(self, cfg: DirectionConfig, on_event: EventCb) -> None: self._last_xrun_log = float("-inf") # throttle de log de xrun self._drops = DropCounter() # descartes de fila sob sobrecarga self._sinks: list[_OutputSink] = [] # saidas da traducao (abertas em _run) + # Devices de saida ja declarados PERDIDOS (terminal `error.play` emitido). + # So a thread de cada sink escreve, e cada uma escreve a propria chave — + # `set.add`/`in` de builtin bastam, sem lock no caminho de erro. + self._sinks_perdidos: set[int] = set() # volume control (linear factors; atualizados via setters threadsafe) self._out_gain = 10.0 ** (float(cfg.output_gain_db) / 20.0) self._pt_gain = 10.0 ** (float(cfg.passthrough_gain_db) / 20.0) @@ -921,6 +951,7 @@ def _open_sinks(self, samplerate: int) -> None: # Stop chegou durante o warmup (carga de modelo estoura o join de 2s # do stop()): nao adianta abrir device para fechar no finally. return + self._sinks_perdidos.clear() # sinks novos, orcamento e veredito zerados self._sinks = [ _OutputSink(dev, samplerate, self._on_sink_error) for dev in self.cfg.output_devices @@ -931,12 +962,57 @@ def _close_sinks(self) -> None: for sink in sinks: sink.close() - def _on_sink_error(self, dev: int, exc: Exception) -> None: + def _on_sink_error(self, dev: int, exc: Exception, consecutivas: int) -> None: + """Decide se a falha de UM device de saida e soluco ou perda definitiva. + + Antes, toda falha virava `error.play` — e `static/app.js` trata erro + nao-recuperavel como fim de sessao (status 'error', `running.delete`, + `toggleButtons(false)`). Resultado: um fone entrando em power save + desabilitava o Stop com o worker ainda capturando e mandando audio pro + Discord, e a unica saida era recarregar a pagina (issue #54). + + O degrau e o mesmo do laco de traducao (issue #45): falha isolada e + `recoverable=True` — aviso ambar transitorio, a direcao segue rodando e o + Stop segue habilitado — e so `OUT_SINK_MAX_ERRORS` falhas CONSECUTIVAS no + mesmo device viram o terminal `error.play`. Preferimos `recoverable` a um + `kind:"status"` novo porque status pinta o painel de VERDE ("rodando") + com um texto de falha e nao volta sozinho; o caminho recuperavel ja existe + na UI, entao nenhuma linha de `app.js` precisa mudar para isto funcionar. + + O veredito terminal e SEM VOLTA para aquele device (`_sinks_perdidos`). + Sem essa trava, o zerar-em-sucesso abriria um estado que nao existia antes + desta issue: depois do `error.play` a UI ja derrubou os botoes + (`app.js`: `toggleButtons(false)` + `running.delete`), e uma frase boa + seguida de nova falha voltaria como `recoverable` — o painel repintaria + VERDE "Rodando" com o Parar desabilitado. Um device que ja custou + `OUT_SINK_MAX_ERRORS` frases seguidas nao volta a ser soluco; tambem evita + re-emitir o terminal a cada frase. + """ + if dev in self._sinks_perdidos: + return + if consecutivas >= OUT_SINK_MAX_ERRORS: + self._sinks_perdidos.add(dev) + self._emit_key( + "error", + "error.play", + {"dev": dev, "max": OUT_SINK_MAX_ERRORS, "detail": str(exc)}, + msg=( + f"Saida perdida (dev={dev}): {OUT_SINK_MAX_ERRORS} frases " + f"seguidas falharam. ({exc})" + ), + ) + return self._emit_key( "error", - "error.play", - {"dev": dev, "detail": str(exc)}, - msg=f"play(dev={dev}): {exc}", + "error.output_retry", + { + "dev": dev, + "attempt": consecutivas, + "max": OUT_SINK_MAX_ERRORS, + "detail": str(exc), + }, + msg=f"Saida falhou (dev={dev}) {consecutivas}/{OUT_SINK_MAX_ERRORS}: {exc}", + recoverable=True, ) def _play(self, pcm: np.ndarray, sr: int) -> None: diff --git a/static/i18n.js b/static/i18n.js index 11e15bb..5a9fc03 100644 --- a/static/i18n.js +++ b/static/i18n.js @@ -44,7 +44,8 @@ window.LAGUNA_I18N = { "error.passthrough_invalid_devices": "passthrough: dispositivos inválidos (src={src}, dst={dst})", "error.passthrough_query": "passthrough query_devices: {detail}", "error.passthrough_generic": "passthrough: {detail}", - "error.play": "play(dev={dev}): {detail}", + "error.output_retry": "Saída falhou (dev={dev}) {attempt}/{max} — a tradução continua. ({detail})", + "error.play": "Saída perdida (dev={dev}) — {max} frases seguidas falharam. Verifique o dispositivo. ({detail})", "error.segment_failed": "Frase perdida ({attempt}/{max}) — a tradução continua. ({detail})", "error.direction_lost": "Direção encerrada: {max} frases seguidas falharam. ({detail})", "error.start_capture_device_invalid": "Dispositivo de captura [{device}] indisponível — atualize a lista de dispositivos e selecione de novo. ({detail})", @@ -193,7 +194,8 @@ window.LAGUNA_I18N = { "error.passthrough_invalid_devices": "passthrough: invalid devices (src={src}, dst={dst})", "error.passthrough_query": "passthrough query_devices: {detail}", "error.passthrough_generic": "passthrough: {detail}", - "error.play": "play(dev={dev}): {detail}", + "error.output_retry": "Output failed (dev={dev}) {attempt}/{max} — translation continues. ({detail})", + "error.play": "Output lost (dev={dev}) — {max} consecutive phrases failed. Check the device. ({detail})", "error.segment_failed": "Phrase dropped ({attempt}/{max}) — translation continues. ({detail})", "error.direction_lost": "Direction stopped: {max} consecutive phrases failed. ({detail})", "error.start_capture_device_invalid": "Capture device [{device}] unavailable — refresh the device list and pick it again. ({detail})", diff --git a/tests_unit/test_output_sink.py b/tests_unit/test_output_sink.py index 39c46c8..0895031 100644 --- a/tests_unit/test_output_sink.py +++ b/tests_unit/test_output_sink.py @@ -14,7 +14,13 @@ import numpy as np import pytest -from laguna_core import DirectionConfig, DirectionWorker, _OutputSink +import laguna_core +from laguna_core import ( + OUT_SINK_MAX_ERRORS, + DirectionConfig, + DirectionWorker, + _OutputSink, +) SR = 22050 FRASE_S = 0.15 # "duracao" simulada de uma frase @@ -65,11 +71,17 @@ def _wait_streams(store: dict, n: int, timeout: float = 2.0) -> None: time.sleep(0.005) +def _coletor(errors): + """Callback `on_error(dev, exc, consecutivas)` que so acumula (issue #54).""" + + def _cb(dev, exc, consecutivas): + errors.append((dev, exc, consecutivas)) + + return _cb + + def _sinks(devices, store, errors, dur=FRASE_S): - return [ - _OutputSink(d, SR, lambda dev, exc: errors.append((dev, exc)), _factory(store, dur)) - for d in devices - ] + return [_OutputSink(d, SR, _coletor(errors), _factory(store, dur)) for d in devices] def test_abre_stream_no_start_e_reusa_entre_frases(): @@ -113,12 +125,12 @@ def test_duas_saidas_comecam_juntas_e_tocam_em_paralelo(): def test_falha_de_um_device_nao_impede_os_outros(): store, errors = {}, [] - bom = _OutputSink(1, SR, lambda dev, exc: errors.append((dev, exc)), _factory(store)) + bom = _OutputSink(1, SR, _coletor(errors), _factory(store)) def _quebrado(device: int, samplerate: int): raise RuntimeError("device sumiu") - ruim = _OutputSink(2, SR, lambda dev, exc: errors.append((dev, exc)), _quebrado) + ruim = _OutputSink(2, SR, _coletor(errors), _quebrado) try: _wait_streams(store, 1) for s in (bom, ruim): @@ -126,7 +138,7 @@ def _quebrado(device: int, samplerate: int): for s in (bom, ruim): assert s.wait(2.0) assert len(store[1].writes) == 1 # o device bom tocou - assert errors and all(dev == 2 for dev, _ in errors) # erro so do ruim + assert errors and all(dev == 2 for dev, _, _ in errors) # erro so do ruim finally: bom.close() ruim.close() @@ -147,7 +159,7 @@ def _open(device: int, samplerate: int): abertos.append(st) return st - sink = _OutputSink(3, SR, lambda dev, exc: errors.append((dev, exc)), _open) + sink = _OutputSink(3, SR, _coletor(errors), _open) try: sink.submit(PCM) assert sink.wait(2.0) @@ -202,7 +214,7 @@ def _open(device: int, samplerate: int): criados.append(st) return st - sink = _OutputSink(11, SR, lambda dev, exc: errors.append((dev, exc)), _open) + sink = _OutputSink(11, SR, _coletor(errors), _open) sink.submit(PCM) assert criados[0].no_write.wait(2.0) # thread esta DENTRO do write sink.close() @@ -220,13 +232,11 @@ class _SempreFalha(FakeStream): def write(self, data): raise RuntimeError("device sumiu no meio da frase") - sink = _OutputSink( - 12, SR, lambda dev, exc: errors.append((dev, exc)), lambda d, sr: _SempreFalha(d, sr) - ) + sink = _OutputSink(12, SR, _coletor(errors), lambda d, sr: _SempreFalha(d, sr)) try: sink.submit(PCM) assert sink.wait(2.0) - assert [dev for dev, _ in errors] == [12] # sem stop: erro chega a UI + assert [dev for dev, _, _ in errors] == [12] # sem stop: erro chega a UI finally: sink.close() @@ -253,3 +263,139 @@ def _tempo_de_play(n_saidas: int) -> float: assert duas < uma * 1.8 # em serie seria ~2x assert not errors assert w._sinks == [] # _close_sinks nao deixa sink pendurado + + +# --- issue #54: soluco de UM device nao pode matar a direcao ----------------- + + +def _worker(): + """Worker inerte (nenhuma thread sobe no __init__) + lista de eventos.""" + eventos: list[dict] = [] + cfg = DirectionConfig(name="falar", src_lang="pt", tgt_lang="en") + return DirectionWorker(cfg, on_event=eventos.append), eventos + + +def test_falha_isolada_de_saida_e_recuperavel_e_nao_para_a_direcao(): + """O bug da #54: uma falha virava `error.play`, e app.js trata erro + nao-recuperavel como fim de sessao — Stop desabilitado com o worker vivo.""" + w, eventos = _worker() + w._on_sink_error(5, RuntimeError("power save"), 1) + + (ev,) = eventos + assert ev["kind"] == "error" and ev["key"] == "error.output_retry" + assert ev["recoverable"] is True # aviso ambar; painel segue 'rodando' + assert ev["args"] == { + "dev": 5, + "attempt": 1, + "max": OUT_SINK_MAX_ERRORS, + "detail": "power save", + } + + +def test_falhas_consecutivas_ate_o_teto_viram_erro_terminal_com_o_device(): + w, eventos = _worker() + for n in range(1, OUT_SINK_MAX_ERRORS + 1): + w._on_sink_error(7, RuntimeError("device sumiu"), n) + + chaves = [e["key"] for e in eventos] + assert chaves == ["error.output_retry"] * (OUT_SINK_MAX_ERRORS - 1) + ["error.play"] + fatal = eventos[-1] + assert not fatal.get("recoverable") # terminal: agora a saida esta perdida + assert fatal["args"]["dev"] == 7 and fatal["args"]["max"] == OUT_SINK_MAX_ERRORS + + +def test_depois_do_terminal_o_device_nao_volta_a_ser_recuperavel(): + """Sem a trava, o zerar-em-sucesso repintaria o painel de VERDE 'rodando' + com o Parar ja desabilitado pelo terminal — estado pior que o da #54.""" + w, eventos = _worker() + for n in range(1, OUT_SINK_MAX_ERRORS + 1): + w._on_sink_error(7, RuntimeError("fone sumiu"), n) + assert eventos[-1]["key"] == "error.play" + + # Frase boa zerou o contador do sink; a falha seguinte chega como 1/3. + w._on_sink_error(7, RuntimeError("fone sumiu de novo"), 1) + assert len(eventos) == OUT_SINK_MAX_ERRORS # nenhum evento novo + assert not any(e.get("recoverable") for e in eventos[OUT_SINK_MAX_ERRORS - 1 :]) + + +def test_terminal_de_um_device_nao_silencia_o_vizinho(): + w, eventos = _worker() + for n in range(1, OUT_SINK_MAX_ERRORS + 1): + w._on_sink_error(7, RuntimeError("fone sumiu"), n) + w._on_sink_error(8, RuntimeError("soluco"), 1) + + ultimo = eventos[-1] + assert ultimo["key"] == "error.output_retry" and ultimo["args"]["dev"] == 8 + assert ultimo["recoverable"] is True + + +def test_open_sinks_zera_o_veredito_dos_devices_perdidos(monkeypatch): + """Direcao reiniciada nao pode herdar device marcado como perdido.""" + w, eventos = _worker() + for n in range(1, OUT_SINK_MAX_ERRORS + 1): + w._on_sink_error(7, RuntimeError("fone sumiu"), n) + + # Stream falso: `_open_sinks` nao aceita injecao, entao trocamos o global que + # o `_OutputSink` resolve no __init__ — segue sem PortAudio/device real. + monkeypatch.setattr(laguna_core, "_open_output_stream", _factory({})) + w.cfg.output_devices = [7] + w._open_sinks(SR) + try: + w._on_sink_error(7, RuntimeError("soluco novo"), 1) + finally: + w._close_sinks() + assert eventos[-1]["key"] == "error.output_retry" + + +def test_frase_tocada_zera_o_orcamento_do_device(): + """Falhas ESPACADAS (um soluco por hora) nunca podem somar ate o teto.""" + errors: list = [] + falhar = {"v": True} + + class _Controlado(FakeStream): + def write(self, data): + if falhar["v"]: + raise RuntimeError("device engasgou") + super().write(data) + + sink = _OutputSink(21, SR, _coletor(errors), lambda d, sr: _Controlado(d, sr)) + try: + for _ in range(2): + sink.submit(PCM) + assert sink.wait(2.0) + assert [n for _, _, n in errors] == [1, 2] + + falhar["v"] = False # frase boa no meio + sink.submit(PCM) + assert sink.wait(2.0) + + falhar["v"] = True + sink.submit(PCM) + assert sink.wait(2.0) + # 1, nao 3: o sucesso zerou o contador e a direcao nao morre. + assert [n for _, _, n in errors] == [1, 2, 1] + finally: + sink.close() + + +def test_orcamento_e_por_device_e_nao_contamina_o_vizinho(): + errors: list = [] + + class _SempreFalha(FakeStream): + def write(self, data): + raise RuntimeError("device sumiu") + + a = _OutputSink(31, SR, _coletor(errors), lambda d, sr: _SempreFalha(d, sr)) + b = _OutputSink(32, SR, _coletor(errors), lambda d, sr: _SempreFalha(d, sr)) + try: + for _ in range(2): + a.submit(PCM) + assert a.wait(2.0) + b.submit(PCM) + assert b.wait(2.0) + + assert [n for dev, _, n in errors if dev == 31] == [1, 2] + assert [n for dev, _, n in errors if dev == 32] == [1] # vizinho no zero + finally: + a.close() + b.close()