Colab run 3 (A100) surfaced three GPU/API-path bugs mocks couldn't catch: 1. faster-whisper yields namedtuple Segments, but batch/bench consume them as dicts (.get) -> AttributeError 'Segment' has no attribute 'get' on every real transcription. Engine now normalizes segments to dicts (_to_dict_segments) at the boundary. 2. API TranscribeOptions carries engine-irrelevant keys (formats, timestamps, glossary_id, post_correction, diarize); worker's TranscriptionOptions(**job.options) crashed with TypeError. Worker now filters job.options to TranscriptionOptions.__slots__. 3. in-proc server never consumed its own queue (jobs stayed queued forever). Added opt-in Settings.auto_worker (default off): lifespan starts a daemon Worker thread for inproc backend, stopped on shutdown. Notebook enables it via LUKESCRIBE_AUTO_WORKER=true so the API upload -> completed flow works end to end. Notebook: bench manifest now uses clips schema (audio_path/duration_sec/ entities); cell 22 reads error_message/error_code; upload poll window raised to 4min (first-run model download). + 5 tests (namedtuple/dict segments, API-style options, auto_worker on/off); 136 tests pass, ruff clean.
304 lines
10 KiB
Python
304 lines
10 KiB
Python
"""Worker 통합 테스트 — 큐 → 클레임 → 전사 → 결과 저장 (전부 mock)."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import time
|
|
|
|
from luke_scribe.jobqueue.broker import InProcBroker
|
|
from luke_scribe.jobqueue.jobs import Job, JobStatus
|
|
from luke_scribe.jobqueue.worker import Worker, WorkerCallbacks
|
|
from luke_scribe.results.store import ResultStore
|
|
|
|
from ..conftest import FakeEngineOwner, FakeIngestor
|
|
|
|
|
|
def _job(**kw) -> Job:
|
|
defaults = dict(
|
|
type="file",
|
|
lane="batch",
|
|
options={
|
|
"model": "large-v3-turbo",
|
|
"language": "ko",
|
|
"device": "cpu",
|
|
"compute_type": "int8",
|
|
},
|
|
source_path="/tmp/fake-source.mp3",
|
|
source_name="fake-source.mp3",
|
|
)
|
|
defaults.update(kw)
|
|
return Job(**defaults)
|
|
|
|
|
|
class TestWorkerLifecycle:
|
|
def test_api_style_options_filtered(self, settings, tmp_path):
|
|
"""API TranscribeOptions(엔진 무관 키 포함) → 워커가 엔진 필드만 골라 처리.
|
|
|
|
Colab 실전에서 잡이 계속 실패한 원인: job.options에 formats/timestamps/
|
|
diarize 등이 포함돼 TranscriptionOptions(**job.options)가 TypeError를 냄.
|
|
"""
|
|
broker = InProcBroker(settings)
|
|
store = ResultStore(str(tmp_path / "results"))
|
|
worker = Worker(
|
|
settings=settings,
|
|
broker=broker,
|
|
store=store,
|
|
owner=FakeEngineOwner(),
|
|
ingestor=FakeIngestor(),
|
|
)
|
|
job = _job(
|
|
options={
|
|
"language": "ko",
|
|
"device": "auto",
|
|
"timestamps": True,
|
|
"formats": ["json", "srt"],
|
|
"word_timestamps": False,
|
|
"vad": True,
|
|
"hotwords": [],
|
|
"glossary_id": None,
|
|
"post_correction": None,
|
|
"diarize": False,
|
|
}
|
|
)
|
|
broker.enqueue(job)
|
|
worker.drain()
|
|
assert broker.get(job.id).status == JobStatus.COMPLETED
|
|
result = store.read_result(job.id)
|
|
assert result is not None
|
|
assert result.status == "completed"
|
|
|
|
def test_complete_flow(self, settings, tmp_path):
|
|
"""enqueue → worker 처리 → completed + 결과 저장 + 콜백."""
|
|
broker = InProcBroker(settings)
|
|
store = ResultStore(str(tmp_path / "results"))
|
|
owner = FakeEngineOwner()
|
|
done: list[str] = []
|
|
|
|
def on_done(job, result):
|
|
done.append(job.id)
|
|
|
|
worker = Worker(
|
|
settings=settings,
|
|
broker=broker,
|
|
store=store,
|
|
owner=owner,
|
|
ingestor=FakeIngestor(duration_sec=4.0),
|
|
callbacks=WorkerCallbacks(on_job_done=on_done),
|
|
)
|
|
job = _job()
|
|
broker.enqueue(job)
|
|
worker.drain()
|
|
|
|
assert broker.get(job.id).status == JobStatus.COMPLETED
|
|
assert done == [job.id]
|
|
result = store.read_result(job.id)
|
|
assert result is not None
|
|
assert result.status == "completed"
|
|
assert result.text # 세그먼트 텍스트 연결
|
|
assert result.execution is not None
|
|
assert result.execution.device == "cpu"
|
|
# 강등 없음
|
|
assert result.execution.downgrade_attempts == 0
|
|
|
|
def test_meta_written(self, settings, tmp_path):
|
|
"""완료 시 job 메타(완료 시각 포함)가 결과 디렉터리에 저장."""
|
|
broker = InProcBroker(settings)
|
|
store = ResultStore(str(tmp_path / "results"))
|
|
worker = Worker(
|
|
settings=settings,
|
|
broker=broker,
|
|
store=store,
|
|
owner=FakeEngineOwner(),
|
|
ingestor=FakeIngestor(),
|
|
)
|
|
job = _job()
|
|
broker.enqueue(job)
|
|
worker.drain()
|
|
import json
|
|
|
|
meta = json.loads(
|
|
tmp_path.joinpath("results", job.id, "meta.json").read_text(encoding="utf-8")
|
|
)
|
|
assert meta["status"] == "completed"
|
|
assert meta["completed_at"] is not None
|
|
|
|
def test_transcription_failure_marks_failed(self, settings, tmp_path):
|
|
"""엔진 오류 → failed + 오류 메타 저장."""
|
|
|
|
class BoomEngine:
|
|
def transcribe(self, req):
|
|
raise RuntimeError("engine exploded")
|
|
|
|
class BoomOwner:
|
|
def __init__(self):
|
|
self.engine = BoomEngine()
|
|
|
|
def transcribe(self, req):
|
|
return self.engine.transcribe(req)
|
|
|
|
def unload_all(self):
|
|
pass
|
|
|
|
broker = InProcBroker(settings)
|
|
store = ResultStore(str(tmp_path / "results"))
|
|
worker = Worker(
|
|
settings=settings,
|
|
broker=broker,
|
|
store=store,
|
|
owner=BoomOwner(),
|
|
ingestor=FakeIngestor(),
|
|
)
|
|
job = _job()
|
|
broker.enqueue(job)
|
|
worker.drain()
|
|
assert broker.get(job.id).status == JobStatus.FAILED
|
|
|
|
def test_cancel_before_claim(self, settings, tmp_path):
|
|
"""큐에서 취소된 job은 워커가 claim하지 않고 cancelled 처리."""
|
|
broker = InProcBroker(settings)
|
|
store = ResultStore(str(tmp_path / "results"))
|
|
worker = Worker(
|
|
settings=settings,
|
|
broker=broker,
|
|
store=store,
|
|
owner=FakeEngineOwner(),
|
|
ingestor=FakeIngestor(),
|
|
)
|
|
job = _job()
|
|
broker.enqueue(job)
|
|
broker.cancel(job.id)
|
|
worker.drain()
|
|
assert broker.get(job.id).status == JobStatus.CANCELLED
|
|
|
|
def test_cancel_during_processing(self, settings, tmp_path):
|
|
"""처리 중 취소 → 세그먼트 경계에서 CancelledError → cancelled (Eng P4)."""
|
|
|
|
class CancelAfterFirstOwner:
|
|
def __init__(self):
|
|
self.n = 0
|
|
self.segments = [
|
|
{"index": 0, "start": 0.0, "end": 1.0, "text": "첫 세그먼트"},
|
|
{"index": 1, "start": 1.0, "end": 2.0, "text": "두 번째"},
|
|
]
|
|
|
|
def transcribe(self, req):
|
|
# worker는 token 미전달, 취소는 세그먼트 경계에서 검사됨
|
|
return {
|
|
"segments": iter(
|
|
[
|
|
{"index": 0, "start": 0.0, "end": 1.0, "text": "첫 세그먼트"},
|
|
{"index": 1, "start": 1.0, "end": 2.0, "text": "두 번째"},
|
|
]
|
|
),
|
|
"info": {"language": "ko"},
|
|
"device": "cpu",
|
|
"compute_type": "int8",
|
|
"attempted_profiles": ["cpu/int8"],
|
|
}
|
|
|
|
# 세그먼트를 먼저 소비한 뒤 취소하도록: 첫 next 후 should_cancel True
|
|
broker = InProcBroker(settings)
|
|
store = ResultStore(str(tmp_path / "results"))
|
|
owner = CancelAfterFirstOwner()
|
|
worker = Worker(
|
|
settings=settings,
|
|
broker=broker,
|
|
store=store,
|
|
owner=owner,
|
|
ingestor=FakeIngestor(),
|
|
)
|
|
job = _job()
|
|
broker.enqueue(job)
|
|
# 처리 시작 전에 취소 요청을 미리 걸어 두면 claim 단계에서 cancelled 처리됨.
|
|
# 대신 처리 중 취소를 시뮬레이션: claim 직후 cancel 호출
|
|
claimed = broker.claim_next("w1", 60.0)
|
|
broker.cancel(claimed.id)
|
|
# 워커는 already-processing job을 직접 처리
|
|
worker._process(claimed)
|
|
assert broker.get(job.id).status == JobStatus.CANCELLED
|
|
|
|
|
|
class TestPrivacyFirst:
|
|
def test_source_deleted_after_completion(self, settings, tmp_path):
|
|
"""plan §3.7e/§6.1: 전사 완료 후 업로드 원본 오디오 즉시 삭제."""
|
|
from luke_scribe.results.store import ResultStore
|
|
|
|
settings.delete_source = True
|
|
store = ResultStore(str(tmp_path / "results"))
|
|
broker = InProcBroker(settings)
|
|
worker = Worker(
|
|
settings=settings,
|
|
broker=broker,
|
|
store=store,
|
|
owner=FakeEngineOwner(),
|
|
ingestor=FakeIngestor(),
|
|
)
|
|
job = _job()
|
|
# source가 store에 저장된 파일을 가리키도록
|
|
src = store.source_path_for(job.id, "meeting.mp3")
|
|
src.write_bytes(b"fake-audio")
|
|
job.source_path = str(src)
|
|
broker.enqueue(job)
|
|
worker.drain()
|
|
assert not src.exists() # 원본 삭제
|
|
assert store.read_result(job.id) is not None # 결과는 보존
|
|
|
|
def test_source_kept_when_delete_source_false(self, settings, tmp_path):
|
|
from luke_scribe.results.store import ResultStore
|
|
|
|
settings.delete_source = False
|
|
store = ResultStore(str(tmp_path / "results"))
|
|
broker = InProcBroker(settings)
|
|
worker = Worker(
|
|
settings=settings,
|
|
broker=broker,
|
|
store=store,
|
|
owner=FakeEngineOwner(),
|
|
ingestor=FakeIngestor(),
|
|
)
|
|
job = _job()
|
|
src = store.source_path_for(job.id, "meeting.mp3")
|
|
src.write_bytes(b"fake-audio")
|
|
job.source_path = str(src)
|
|
broker.enqueue(job)
|
|
worker.drain()
|
|
assert src.exists()
|
|
|
|
|
|
class TestWorkerProgress:
|
|
def test_progress_emitted(self, settings, tmp_path):
|
|
broker = InProcBroker(settings)
|
|
store = ResultStore(str(tmp_path / "results"))
|
|
worker = Worker(
|
|
settings=settings,
|
|
broker=broker,
|
|
store=store,
|
|
owner=FakeEngineOwner(segments=[{"index": 0, "start": 0.0, "end": 5.0, "text": "x"}]),
|
|
ingestor=FakeIngestor(duration_sec=10.0),
|
|
)
|
|
job = _job()
|
|
broker.enqueue(job)
|
|
worker.drain()
|
|
final = broker.get(job.id)
|
|
assert final.progress is not None
|
|
assert final.progress <= 1.0
|
|
assert final.processed_sec is not None
|
|
|
|
|
|
class TestCrashRecovery:
|
|
def test_stale_processing_recovered(self, settings, tmp_path):
|
|
"""Eng P11: 워커 크래시(리스 만료) → 스타트업 reconciler가 failed 처리."""
|
|
broker = InProcBroker(settings)
|
|
job = _job()
|
|
broker.enqueue(job)
|
|
claimed = broker.claim_next("w1", 60.0)
|
|
# 리스 만료 시뮬레이션
|
|
claimed.lease_expires_at = time.time() - 10
|
|
broker.save_meta(claimed)
|
|
worker = Worker(
|
|
settings=settings, broker=broker, store=ResultStore(str(tmp_path / "results"))
|
|
)
|
|
reclaimed = worker.start_reconcile()
|
|
assert job.id in reclaimed
|
|
assert broker.get(job.id).status == JobStatus.FAILED
|
|
assert broker.get(job.id).error_code == "worker_crash"
|