"""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_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"