Repository navigation
Expand file tree
/
Copy pathmain.py
More file actions
2350 lines (2128 loc) · 107 KB
/
Copy pathmain.py
File metadata and controls
2350 lines (2128 loc) · 107 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
Speech-to-Text v4.1 - Main Script
YouTube 비디오를 다운로드하고 음성을 텍스트로 변환한 후 요약하여 마크다운으로 저장합니다.
"""
import os
import re
import sys
import subprocess
import time
import random
import logging
import json
import hashlib
from datetime import datetime
from pathlib import Path
from dataclasses import dataclass
from typing import Optional, Tuple, List
from dotenv import load_dotenv
_PROJECT_ROOT = Path(__file__).resolve().parent
# .env + local scratch (TMPDIR, XDG_CACHE_HOME) before yt-dlp import via stt
load_dotenv(_PROJECT_ROOT / ".env")
try:
os.chdir(_PROJECT_ROOT)
except OSError:
pass # launchd may lack Documents TCC; imports use __file__ dir via sys.path[0]
from config import (
APP_VERSION,
apply_work_path_scratch_env,
get_config_dict,
log_macos_deadlock_path_warnings,
resolve_data_root,
)
apply_work_path_scratch_env()
from tqdm import tqdm
import pandas as pd
from openai import OpenAI
# launchd 환경에서 한글 출력을 위한 인코딩 설정
if sys.stdout.encoding != 'utf-8':
import io
sys.stdout = io.TextIOWrapper(sys.stdout.buffer, encoding='utf-8', errors='replace')
if sys.stderr.encoding != 'utf-8':
import io
sys.stderr = io.TextIOWrapper(sys.stderr.buffer, encoding='utf-8', errors='replace')
import stt_function_v3 as stt
import channel_crawl
from filename_utils import fit_filename, validate_llm_suffix
import run_lock
from job_workspace import VideoJobWorkspace, cleanup_stale_jobs
from transcript_cache import (
TranscriptCache,
find_durable_full_transcript,
should_write_transcript_cache,
)
from subtitle_lifecycle import (
cleanup_expired_quarantine,
delete_subtitle_file,
quarantine_job_subtitles,
cleanup_legacy_yt_subs,
)
from pipeline_models import VideoProcessResult
from pipeline_context import set_pipeline_context, clear_video_context, attach_pipeline_context_filter
from claim_manager import ClaimManager
from shared_state_writer import SharedStateWriter
from admission_limiter import DownloadAdmissionLimiter, ProviderCooldown
from preprocess_backend import create_transcript_preprocessor, TranscriptPreprocessor
from llm_responses import call_responses, ResponsesCallError
from runtime_resources import device_compute_route, get_device_semaphore
import uuid
from concurrent.futures import ThreadPoolExecutor, as_completed
# Configure logging
def setup_logging(log_dir: Optional[str] = None) -> logging.Logger:
"""Setup logging configuration."""
if log_dir is None:
log_dir = str(_PROJECT_ROOT / "logs")
os.makedirs(log_dir, exist_ok=True)
log_file = os.path.join(log_dir, f"stt_{datetime.now().strftime('%Y%m%d')}.log")
# Create root logger so logs from imported modules (e.g., channel_crawl) are visible.
logger = logging.getLogger()
logger.setLevel(logging.INFO)
# Remove existing handlers
logger.handlers = []
# File handler (always used)
file_handler = logging.FileHandler(log_file, encoding='utf-8')
file_handler.setLevel(logging.INFO)
file_formatter = logging.Formatter(
'%(asctime)s - %(name)s - %(levelname)s - %(message)s'
)
file_handler.setFormatter(file_formatter)
logger.addHandler(file_handler)
# Console handler only when stdout is a real TTY (interactive). Under launchd/cron or
# when stdout is redirected to a file (e.g. iCloud log), writing/flushing to it can
# raise OSError [Errno 11]; skip StreamHandler so we only write to the log file.
if sys.stdout.isatty():
console_handler = logging.StreamHandler(sys.stdout)
console_handler.setLevel(logging.INFO)
console_formatter = logging.Formatter(
'%(asctime)s - %(levelname)s - %(message)s',
datefmt='%H:%M:%S'
)
console_handler.setFormatter(console_formatter)
logger.addHandler(console_handler)
return logger
# Initialize logger
logger = setup_logging()
attach_pipeline_context_filter(logger)
# Main LLM output length target = token_range × concise input tokens
MAIN_LLM_TOKEN_RANGE = (1.5, 2.0)
@dataclass
class MainLlmConfig:
"""Primary main summarization LLM with optional fallback route."""
primary_client: OpenAI
primary_model: str
primary_provider: str
fallback_client: Optional[OpenAI] = None
fallback_model: Optional[str] = None
fallback_provider: Optional[str] = None
@property
def has_fallback(self) -> bool:
return self.fallback_client is not None and bool(self.fallback_model)
def summarize(
self,
transcription: str,
filename: str,
prompt: str,
*,
token_range=MAIN_LLM_TOKEN_RANGE,
language: str = "Korean",
style: str = "Markdown",
) -> str:
return stt.summarize_with_chunking(
transcription=transcription,
filename=filename,
prompt=prompt,
client=self.primary_client,
token_range=token_range,
language=language,
style=style,
model=self.primary_model,
fallback_client=self.fallback_client,
fallback_model=self.fallback_model,
fallback_provider=self.fallback_provider,
primary_provider=self.primary_provider,
)
def failure_needs_long_cooldown(status: str, error_msg: Optional[str]) -> bool:
"""
When True, apply FAILURE_WAIT_MULTIPLIER between videos (network / rate-limit / transient I/O).
When False, use normal MIN/MAX wait (private/unavailable/subs quirks/mlx issues, etc.).
Transient patterns are checked first so yt-dlp messages that end with a generic suffix
still match (e.g. timeout + "Download failed after all retry attempts").
"""
err = (error_msg or "").lower()
st = (status or "").lower()
transient_markers = (
"429",
"rate limit",
"too many requests",
"quota exceeded",
"503",
"502",
"504",
"service unavailable",
"bad gateway",
"gateway timeout",
"internal server error",
"timeout",
"timed out",
"time out",
"connection reset",
"reset by peer",
"broken pipe",
"connection aborted",
"connection refused",
"network is unreachable",
"no route to host",
"temporary failure in name resolution",
"name or service not known",
"errno 11",
"resource deadlock",
"ssl",
"certificate verify",
"tls handshake",
"403",
"forbidden",
"bot detection",
"unable to connect",
"connection error",
"try again later",
"http error 500",
)
if any(m in err for m in transient_markers):
return True
if st == "file_error" and ("errno 11" in err or "resource deadlock" in err):
return True
# Content / policy / clearly non-network: short cooldown
short_markers = (
"private",
"members only",
"members-only",
"video unavailable",
"this video is private",
"video is private",
"not available in your country",
"removed by",
"copyright",
"blocked in your country",
"sign in to confirm",
"login required",
"invalid api key",
"no audio stream",
"regexmatch",
"uploader has not made",
"age-restricted",
"age restricted",
"downloaded file not found after", # post-download mismatch (not Errno 11)
"neither yt-dlp nor pytubefix",
"exhausted retries without success",
)
if any(m in err for m in short_markers):
return False
# Generic message with no captured detail → short (avoid long idle)
if err.strip() in ("", "download failed after all retry attempts"):
return False
# Unknown errors: prefer short cooldown (user preference: do not block batch for minutes)
return False
LOCAL_BASE_PATH_DEFAULT = str(_PROJECT_ROOT)
def _append_llm_usage(data_root: str, record: dict, video_id: str) -> None:
"""Append token metadata only; transcript and model reasoning are never persisted."""
os.makedirs(os.path.join(data_root, "logs"), exist_ok=True)
row = dict(record)
row["video_id"] = video_id
row.setdefault("prompt_file", DIRECT_PROMPT_FILE)
row.setdefault("prompt_sha256", DIRECT_PROMPT_SHA256)
with open(os.path.join(data_root, "logs", "llm_usage.jsonl"), "a", encoding="utf-8") as handle:
handle.write(json.dumps(row, ensure_ascii=False, sort_keys=True) + "\n")
def run_direct_summary(client, transcription: str, filename: str, video_id: str, config: dict):
"""Fail closed above the configured raw-input ceiling, then make one Responses call."""
try:
raw_tokens = len(stt.tiktoken.get_encoding("cl100k_base").encode(transcription))
except Exception as exc:
raise ResponsesCallError("token_count_unavailable", f"cannot count raw input tokens: {type(exc).__name__}", None) from exc
limit = int(config.get("DIRECT_MAX_INPUT_TOKENS", 200000))
if raw_tokens > limit:
raise ResponsesCallError("input_too_long", f"{raw_tokens} raw tokens exceeds limit {limit}; manual review required", None)
return call_responses(
client,
model=config.get("DIRECT_LLM_MODEL", "gpt-6-luna"),
instructions=DIRECT_LUNA_PROMPT,
input_text=f"# {filename}\n\n{transcription}",
source_text=transcription,
effort=config.get("DIRECT_LLM_REASONING_EFFORT", "medium"),
max_output_tokens=int(config.get("DIRECT_MAX_OUTPUT_TOKENS", 32000)),
sample_video_id=video_id,
)
def summarize_with_pipeline(mode: str, *, direct_client, main_llm, direct_transcription: str,
legacy_transcription: str, filename: str, video_id: str,
config: dict, prompt: str):
"""Keep the legacy route isolated while selecting one direct Responses call when enabled."""
if mode == "direct_luna":
return run_direct_summary(direct_client, direct_transcription, filename, video_id, config)
text = main_llm.summarize(
transcription=legacy_transcription, filename=filename, prompt=prompt,
token_range=list(MAIN_LLM_TOKEN_RANGE), language="Korean", style="Markdown",
)
return text, None
def load_config() -> dict:
"""Load configuration from .env (secrets/paths) and config.py (threshold, rate limiting, channel crawl)."""
config = {
'BASE_PATH': os.getenv('BASE_PATH', LOCAL_BASE_PATH_DEFAULT),
'WORK_PATH': os.getenv('WORK_PATH', '').strip() or None,
'HF_HOME': os.getenv('HF_HOME', str(Path.home() / '.cache' / 'whisper')),
'OUTPUT_MD_PATH': os.getenv('OUTPUT_MD_PATH', ''),
'OUTPUT_MD_GIT': os.getenv('OUTPUT_MD_GIT', ''),
'OPENAI_API_KEY': os.getenv('OPENAI_API_KEY'),
'XAI_API_KEY': os.getenv('XAI_API_KEY'),
'OPENROUTER_API_KEY': os.getenv('OPENROUTER_API_KEY', '').strip(),
'MAIN_LLM_PROVIDER': os.getenv('MAIN_LLM_PROVIDER', 'openai').strip().lower(),
'MAIN_LLM_MODEL': os.getenv('MAIN_LLM_MODEL', 'gpt-5-mini-2025-08-07').strip(),
'DIRECT_LLM_MODEL': os.getenv('DIRECT_LLM_MODEL', 'gpt-6-luna').strip(),
'LLM_PIPELINE_MODE': os.getenv('LLM_PIPELINE_MODE', 'legacy_two_stage').strip().lower(),
'DIRECT_LLM_REASONING_EFFORT': os.getenv('DIRECT_LLM_REASONING_EFFORT', 'medium').strip().lower(),
'DIRECT_MAX_INPUT_TOKENS': int(os.getenv('DIRECT_MAX_INPUT_TOKENS', '200000')),
'DIRECT_MAX_OUTPUT_TOKENS': int(os.getenv('DIRECT_MAX_OUTPUT_TOKENS', '32000')),
'MAIN_LLM_FALLBACK_PROVIDER': os.getenv('MAIN_LLM_FALLBACK_PROVIDER', '').strip().lower(),
'MAIN_LLM_FALLBACK_MODEL': os.getenv('MAIN_LLM_FALLBACK_MODEL', '').strip(),
'MAIN_LLM_OUTPUT_SUFFIX': os.getenv('MAIN_LLM_OUTPUT_SUFFIX', '5-mini').strip(), # 5-mini=gpt-5-mini, dS4f=deepseek
'PREPROCESS_LLM_MODEL': os.getenv('PREPROCESS_LLM_MODEL', 'gpt-5-nano-2025-08-07').strip(),
'PROXY_ADDRESS': os.getenv('PROXY_ADDRESS', ''),
'YOUTUBE_COOKIES_FILE': os.getenv('YOUTUBE_COOKIES_FILE', ''),
'YOUTUBE_API_KEY': os.getenv('YOUTUBE_API_KEY', ''),
}
try:
config.update(get_config_dict())
except Exception as e:
logger.warning(f"config.py load failed ({e}). Using defaults for rate limiting and AUDIO_SIZE_THRESHOLD_MB.")
config.update({
'AUDIO_SIZE_THRESHOLD_MB': 1024,
'MIN_WAIT_BETWEEN_VIDEOS': 30,
'MAX_WAIT_BETWEEN_VIDEOS': 60,
'EXTENDED_WAIT_INTERVAL': 10,
'EXTENDED_WAIT_DURATION': 300,
'MAX_CONSECUTIVE_FAILURES': 5,
'FAILURE_WAIT_MULTIPLIER': 2.0,
'CHANNEL_CRAWL': False,
'CHANNEL_BACKFILL': False,
'CHANNEL_START_DATE': '',
'CHANNEL_END_DATE': '',
'FILTERING_SHORTS_MINUTES': 3,
'CRAWL_QUEUE_MAX_RETRIES': 3,
'YOUTUBE_AUTO_SCRIPT': True,
'YOUTUBE_SUBS_LANGS': 'en,ko,jp,en-US,en-GB',
})
# Env overrides for YouTube subs ( .env takes precedence )
_auto = os.getenv('YOUTUBE_AUTO_SCRIPT', '').strip().lower()
if _auto:
config['YOUTUBE_AUTO_SCRIPT'] = _auto in ('true', '1', 'yes')
_langs = os.getenv('YOUTUBE_SUBS_LANGS', '').strip()
if _langs:
config['YOUTUBE_SUBS_LANGS'] = _langs
if config['LLM_PIPELINE_MODE'] not in {'direct_luna', 'legacy_two_stage'}:
raise ValueError('LLM_PIPELINE_MODE must be direct_luna or legacy_two_stage')
if config['DIRECT_LLM_REASONING_EFFORT'] not in {'low', 'medium', 'high'}:
raise ValueError('DIRECT_LLM_REASONING_EFFORT must be low, medium, or high')
if os.getenv('MAIN_LLM_OUTPUT_SUFFIX') is None and config['LLM_PIPELINE_MODE'] == 'direct_luna':
config['MAIN_LLM_OUTPUT_SUFFIX'] = 'luna-' + config['DIRECT_LLM_REASONING_EFFORT']
validate_llm_suffix(config['MAIN_LLM_OUTPUT_SUFFIX']) # keeps the filename tail short
_save_full = os.getenv('SAVE_FULL_WHEN_AUTO_SUBS', '').strip().lower()
if _save_full:
config['SAVE_FULL_WHEN_AUTO_SUBS'] = _save_full in ('true', '1', 'yes')
_use_job = os.getenv('USE_JOB_WORKSPACE', '').strip().lower()
if _use_job:
config['USE_JOB_WORKSPACE'] = _use_job in ('true', '1', 'yes')
_cache_en = os.getenv('TRANSCRIPT_CACHE_ENABLED', '').strip().lower()
if _cache_en:
config['TRANSCRIPT_CACHE_ENABLED'] = _cache_en in ('true', '1', 'yes')
_cache_ttl = os.getenv('TRANSCRIPT_CACHE_TTL_HOURS', '').strip()
if _cache_ttl.isdigit():
config['TRANSCRIPT_CACHE_TTL_HOURS'] = int(_cache_ttl)
_vw = os.getenv('VIDEO_WORKERS', '').strip()
if _vw.isdigit():
config['VIDEO_WORKERS'] = int(_vw)
_pb = os.getenv('PREPROCESS_BACKEND', '').strip()
if _pb:
config['PREPROCESS_BACKEND'] = _pb
config['VIDEO_WORKERS'] = min(2, max(1, int(config.get('VIDEO_WORKERS', 2))))
get_device_semaphore(int(config.get('DEVICE_COMPUTE_CONCURRENCY', 1)))
# Hot CSV / append jsonl: local DATA_ROOT when WORK_PATH or DATA_ROOT env set
config['DATA_ROOT'] = resolve_data_root(
config['BASE_PATH'],
config.get('WORK_PATH'),
)
# Validate required keys
if not config['OPENAI_API_KEY']:
raise ValueError("OPENAI_API_KEY가 .env 파일에 설정되지 않았습니다.")
if config.get('MAIN_LLM_PROVIDER') == 'openrouter' and not config.get('OPENROUTER_API_KEY'):
raise ValueError("MAIN_LLM_PROVIDER=openrouter일 때 OPENROUTER_API_KEY가 .env에 필요합니다.")
if config.get('MAIN_LLM_FALLBACK_PROVIDER') == 'openrouter' and not config.get('OPENROUTER_API_KEY'):
raise ValueError("MAIN_LLM_FALLBACK_PROVIDER=openrouter일 때 OPENROUTER_API_KEY가 .env에 필요합니다.")
if config.get('MAIN_LLM_FALLBACK_MODEL') and not config.get('MAIN_LLM_FALLBACK_PROVIDER'):
raise ValueError("MAIN_LLM_FALLBACK_MODEL이 설정되면 MAIN_LLM_FALLBACK_PROVIDER도 필요합니다.")
if config.get('CHANNEL_CRAWL') and not (config.get('YOUTUBE_API_KEY') or "").strip():
raise ValueError("CHANNEL_CRAWL=true일 때 .env에 YOUTUBE_API_KEY가 필요합니다. docs/YOUTUBE_API_SETUP.md 참고.")
fallback_bits = ""
if config.get('MAIN_LLM_FALLBACK_MODEL'):
fallback_bits = " fallback=%s (%s)" % (
config.get('MAIN_LLM_FALLBACK_MODEL'),
config.get('MAIN_LLM_FALLBACK_PROVIDER'),
)
if config.get('LLM_PIPELINE_MODE') == 'direct_luna':
logger.info(
"LLM config: mode=direct_luna model=%s effort=%s prompt=%s (%s) (openai responses, single call) output_suffix=_%s",
config.get('DIRECT_LLM_MODEL'),
config.get('DIRECT_LLM_REASONING_EFFORT'),
DIRECT_PROMPT_FILE,
DIRECT_PROMPT_SHA256[:12],
config.get('MAIN_LLM_OUTPUT_SUFFIX'),
)
else:
logger.info(
"LLM config: mode=legacy_two_stage preprocess=%s main=%s (%s)%s output_suffix=_%s",
config.get('PREPROCESS_LLM_MODEL'),
config.get('MAIN_LLM_MODEL'),
config.get('MAIN_LLM_PROVIDER'),
fallback_bits,
config.get('MAIN_LLM_OUTPUT_SUFFIX'),
)
# Validate paths
paths_to_check = {
'BASE_PATH': config['BASE_PATH'],
'HF_HOME': config['HF_HOME'],
'OUTPUT_MD_PATH': config['OUTPUT_MD_PATH'],
}
for key, path in paths_to_check.items():
if not os.path.exists(path):
logger.warning(f"{key} 경로가 존재하지 않습니다: {path}")
logger.warning(f"경로를 생성하거나 .env 파일에서 올바른 경로로 수정해주세요.")
else:
logger.info(f"{key} 경로 확인됨: {path}")
try:
os.makedirs(config['DATA_ROOT'], exist_ok=True)
logger.info("DATA_ROOT (input/output/channel/queue CSV, video_metadata_live.jsonl): %s", config['DATA_ROOT'])
except OSError as e:
logger.warning("DATA_ROOT 생성 실패 (%s): %s", config['DATA_ROOT'], e)
log_macos_deadlock_path_warnings(
logger,
data_root=config.get("DATA_ROOT"),
work_path=config.get("WORK_PATH"),
base_path=config.get("BASE_PATH"),
tmpdir=os.environ.get("TMPDIR", "").strip() or None,
xdg_cache=os.environ.get("XDG_CACHE_HOME", "").strip() or None,
)
# Optional path check
if config.get('OUTPUT_MD_GIT') and not os.path.exists(config['OUTPUT_MD_GIT']):
logger.warning(f"OUTPUT_MD_GIT 경로가 존재하지 않습니다: {config['OUTPUT_MD_GIT']}")
logger.warning("필요시 디렉토리를 생성하거나 .env에서 주석 처리하세요.")
return config
def _create_llm_provider_client(provider: str, config: dict, openai_client: OpenAI) -> OpenAI:
if provider == 'openrouter':
return OpenAI(
api_key=config['OPENROUTER_API_KEY'],
base_url="https://openrouter.ai/api/v1",
)
if provider == 'openai':
return openai_client
raise ValueError(f"Unsupported LLM provider: {provider}")
def initialize_clients(config: dict) -> Tuple[OpenAI, Optional[OpenAI], MainLlmConfig]:
"""Initialize OpenAI (preprocess), optional XAI, and main summarization clients."""
logger.info("Initializing API clients...")
try:
logger.info(" Initializing OpenAI client...")
openai_client = OpenAI(api_key=config['OPENAI_API_KEY'])
logger.info(" ✓ OpenAI client initialized")
# Test API connection with a simple request (optional)
# This can be enabled for debugging but adds overhead
# try:
# models = openai_client.models.list()
# logger.debug(f" OpenAI API connection test successful")
# except Exception as e:
# logger.warning(f" OpenAI API connection test failed: {str(e)}")
except Exception as e:
error_type = type(e).__name__
logger.error(f"[ERROR] Failed to initialize OpenAI client")
logger.error(f" Exception Type: {error_type}")
logger.error(f" Error: {str(e)}")
logger.error(f" Possible causes: Invalid API key, network issue")
raise
xai_client = None
if config.get('XAI_API_KEY'):
try:
logger.info(" Initializing XAI (Grok) client...")
xai_client = OpenAI(api_key=config['XAI_API_KEY'], base_url="https://api.x.ai/v1")
logger.info(" ✓ XAI client initialized")
except Exception as e:
error_type = type(e).__name__
logger.warning(f"[WARNING] Failed to initialize XAI client")
logger.warning(f" Exception Type: {error_type}")
logger.warning(f" Error: {str(e)}")
logger.warning(f" Continuing without XAI client (optional)")
xai_client = None
primary_provider = config.get('MAIN_LLM_PROVIDER', 'openai')
try:
logger.info(" Initializing main LLM client (%s)...", primary_provider)
primary_client = _create_llm_provider_client(primary_provider, config, openai_client)
logger.info(" ✓ Main LLM primary client initialized (model=%s)", config.get('MAIN_LLM_MODEL'))
except Exception as e:
error_type = type(e).__name__
logger.error("[ERROR] Failed to initialize main LLM primary client")
logger.error(" Exception Type: %s", error_type)
logger.error(" Error: %s", str(e))
raise
fallback_client = None
fallback_provider = config.get('MAIN_LLM_FALLBACK_PROVIDER') or None
fallback_model = config.get('MAIN_LLM_FALLBACK_MODEL') or None
if fallback_provider and fallback_model:
try:
logger.info(" Initializing main LLM fallback client (%s)...", fallback_provider)
fallback_client = _create_llm_provider_client(fallback_provider, config, openai_client)
logger.info(" ✓ Main LLM fallback client initialized (model=%s)", fallback_model)
except Exception as e:
error_type = type(e).__name__
logger.error("[ERROR] Failed to initialize main LLM fallback client")
logger.error(" Exception Type: %s", error_type)
logger.error(" Error: %s", str(e))
raise
main_llm = MainLlmConfig(
primary_client=primary_client,
primary_model=config.get('MAIN_LLM_MODEL', 'gpt-5-mini-2025-08-07'),
primary_provider=primary_provider,
fallback_client=fallback_client,
fallback_model=fallback_model,
fallback_provider=fallback_provider,
)
return openai_client, xai_client, main_llm
def load_dataframes(data_root: str) -> Tuple[pd.DataFrame, pd.DataFrame]:
"""Load input and output dataframes from DATA_ROOT."""
input_df_path = os.path.join(data_root, 'input_df.csv')
output_df_path = os.path.join(data_root, 'output_df_new.csv')
try:
if not os.path.exists(input_df_path):
error_msg = f"Input file not found: {input_df_path}"
logger.error(f"[ERROR] FileNotFoundError - {error_msg}")
logger.error(f" Expected path: {input_df_path}")
logger.error(f" Solution: Create input_df.csv with 'url' column")
raise FileNotFoundError(error_msg)
logger.info(f"Loading input dataframe: {input_df_path}")
input_df = _read_csv_with_retry(input_df_path, encoding='cp949')
if 'url' not in input_df.columns:
error_msg = "Input dataframe must have 'url' column"
logger.error(f"[ERROR] ValueError - {error_msg}")
logger.error(f" Columns found: {list(input_df.columns)}")
raise ValueError(error_msg)
logger.info(f" Loaded {len(input_df)} URLs from input file")
except pd.errors.EmptyDataError as e:
logger.error(f"[ERROR] EmptyDataError - Input file is empty: {input_df_path}")
raise
except pd.errors.ParserError as e:
logger.error(f"[ERROR] ParserError - Invalid CSV format: {input_df_path}")
logger.error(f" Error: {str(e)}")
raise
except UnicodeDecodeError as e:
logger.error(f"[ERROR] UnicodeDecodeError - Encoding issue: {input_df_path}")
logger.error(f" Tried encoding: cp949")
logger.error(f" Solution: Check file encoding or convert to UTF-8")
raise
try:
if os.path.exists(output_df_path):
logger.info(f"Loading output dataframe: {output_df_path}")
output_df = _read_csv_with_retry(output_df_path)
logger.info(f" Loaded {len(output_df)} existing records")
else:
logger.info(f"Output file not found, creating new: {output_df_path}")
output_df = pd.DataFrame(columns=['date', 'url', 'v_id', 'status'])
_save_output_df_with_retry(output_df, output_df_path, na_rep="")
logger.info(f" Created empty output dataframe")
except Exception as e:
error_type = type(e).__name__
logger.error(f"[ERROR] {error_type} - Failed to load/create output dataframe")
logger.error(f" File path: {output_df_path}")
logger.error(f" Error: {str(e)}")
raise
return input_df, output_df
def _save_output_df_with_retry(output_df: pd.DataFrame, output_df_path: str, na_rep: str = "unknown", max_retries: int = 3) -> None:
"""Save output_df to CSV with retry on iCloud lock (Errno 11). Uses temp file + os.replace."""
tmp_path = output_df_path + ".tmp"
for attempt in range(max_retries):
try:
output_df.to_csv(tmp_path, na_rep=na_rep, index=False, encoding="utf-8-sig")
os.replace(tmp_path, output_df_path)
return
except OSError as e:
if getattr(e, "errno", None) == 11 and attempt < max_retries - 1:
wait = 2 * (attempt + 1)
logger.warning("CSV write Errno 11 (iCloud lock), retry %d/%d in %ds: %s", attempt + 1, max_retries, wait, output_df_path)
time.sleep(wait)
continue
raise
def _read_csv_with_retry(path: str, encoding: str = 'utf-8-sig', max_retries: int = 6) -> pd.DataFrame:
"""Read CSV with retry on iCloud lock (Errno 11). Fallback: copy to temp then read."""
for attempt in range(max_retries):
try:
return pd.read_csv(path, encoding=encoding)
except OSError as e:
if getattr(e, "errno", None) != 11:
raise
if attempt < max_retries - 1:
wait = 3 * (attempt + 1)
logger.warning("CSV read Errno 11 (iCloud lock), retry %d/%d in %ds: %s", attempt + 1, max_retries, wait, path)
time.sleep(wait)
continue
# Last attempt failed; try reading via temp copy (avoids holding iCloud file open)
try:
import tempfile
import shutil
fd, tmp_path = tempfile.mkstemp(suffix=".csv")
os.close(fd)
for copy_attempt in range(3):
try:
shutil.copy2(path, tmp_path)
break
except OSError as e2:
if getattr(e2, "errno", None) == 11 and copy_attempt < 2:
time.sleep(5)
continue
raise
try:
return pd.read_csv(tmp_path, encoding=encoding)
finally:
try:
os.remove(tmp_path)
except OSError:
pass
except Exception:
raise
def load_output_df_only(data_root: str) -> pd.DataFrame:
"""Load or create output_df_new.csv only (for channel_crawl mode when input_df is not used)."""
output_df_path = os.path.join(data_root, 'output_df_new.csv')
if os.path.exists(output_df_path):
output_df = _read_csv_with_retry(output_df_path)
logger.info(f"Loaded output dataframe: {output_df_path} ({len(output_df)} records)")
else:
output_df = pd.DataFrame(columns=['date', 'url', 'v_id', 'status'])
_save_output_df_with_retry(output_df, output_df_path, na_rep="")
logger.info(f"Created output dataframe: {output_df_path}")
return output_df
def get_url_list(input_df: pd.DataFrame, output_df: pd.DataFrame) -> List[str]:
"""Get list of URLs to process (excluding already processed ones)."""
url_list = input_df['url'].tolist()
done_urls = output_df['url'].tolist() if 'url' in output_df.columns else []
url_list = list(set(url_list) - set(done_urls))
url_list.reverse()
return url_list
def get_input_urls_for_channel_crawl(data_root: str, output_df: pd.DataFrame) -> List[str]:
"""
When CHANNEL_CRAWL=true: load input_df.csv if present and return URLs not yet in output_df.
Returns [] if input_df.csv missing or invalid (so channel-only mode is unchanged).
"""
input_df_path = os.path.join(data_root, 'input_df.csv')
if not os.path.exists(input_df_path):
return []
try:
input_df = _read_csv_with_retry(input_df_path, encoding='cp949')
if input_df.empty or 'url' not in input_df.columns:
return []
return get_url_list(input_df, output_df)
except Exception as e:
logger.warning("Channel crawl: could not load input_df for merge (%s). Using channel URLs only.", e)
return []
# Prompt templates
def build_token_query(retention_min: int, retention_max: int, *, auto_subs: bool = False) -> str:
"""Build nano minimization prompt with configurable retention (Phase 1c)."""
filler = (
"Remove Korean filler words (e.g. 음, 어, 그, 네, 막, 좀) and repeated phrases.\n"
if auto_subs
else ""
)
return f"""
your task is to review the given text and remove any redundant or unnecessary wording (e.g., interjections, filler words), ensuring the task is covering as possible as the full meaning and content of the original text.
{filler}
The length of the revision must not exceed the original text, but DO NOT oversimplify or excessively reduce the content.
Please, you should keep {retention_min}~{retention_max}% of the length of the original contents (excluding timestamps if exists); Do NOT miss important details in the text!
Return only the revised text in plain text format **without** adding any commentary, acknowledgments, or opinions of your own.
"""
TOKEN_QUERY = build_token_query(80, 95)
TOKEN_INPUT_ROLE = """
Please effectively corrects typos and removes unnecessary words in text
"""
PRE_TASK_TYPE = "nano_preprocess"
CONTEXT_QUERY = """
YouTube transcription → Obsidian MD. Reorganize into clear topic structure (reorder/merge OK). Mobile-scannable but substantively complete — do not skim.
"""
MOBILE_STRUCTURE_QUERY = """
Section order (exact):
1. `# {title}` — one H1.
2. `## 한눈에 보기` — 3~5 bullets; tag at start: [확정] or [정황] only (video facts; no [추정]/[외부지식]).
3. Main body — 3~6 topic `##` sections mirroring the video arc; each with 2+ concrete points (names, numbers, claims) under `##` or `###`.
4. At most one table; no mermaid/diagrams.
5. Collapsible Insights callout — every line inside must start with `>`; 2~4 bullets; tag [외부지식] or [추정] only:
> [!note]- Insights
> - [외부지식] context not stated in the video (framework, parallel, industry norm)
6. Collapsible Key Takeaways — every line inside must start with `>`; 3~5 bullets; so-what / risks / watch-items (optional [추정] if speculative):
> [!note]- Key Takeaways
> - implication for the reader — not a restatement of 한눈에 보기
7. `## Tags` — 3~5 lowercase bullets (moved to YAML).
8. Optional `## 용어` — 2~4 definitions if jargon; skip if none.
"""
CONTENT_RULES_QUERY = """
Content:
- Main body: source-only facts; prefer restructuring over deleting detail.
- No invented stats/quotes; no "부족한 점" / "개선 제안".
- Use English for transliterated terms (e.g. ChatGPT not 챗지피티).
- Anti-duplication: if a point already appears in 한눈에 보기 or 본문, do not repeat it in callouts. Callouts must add new value.
"""
INSIGHTS_GROUNDING_QUERY = """
Grounding (A4) — tag placement by section:
- 한눈에 보기 / 본문: [확정]/[정황] only; source facts.
- Insights: [외부지식] or [추정] only; general domain context OK; no invented specifics (stats, quotes, named events not in source).
- Key Takeaways: implications (risks, decisions, watch-items); synthesize across topics — do not restate 한눈에 보기.
Good Key Takeaway: "의존도가 80% 이상이면 공급망 리스크를 운영 계획에 반영해야 함."
Bad Key Takeaway: "해당 국가는 기회를 가졌으나 구조적 문제로 성장하지 못함." (한눈에 보기 복붙)
"""
TONE_QUERY = """
Tone: Korean, direct (avoid ~입니다/~합니다). Raw markdown only — no code fences.
"""
INPUT_PROMPT = f"""
{CONTEXT_QUERY}
{MOBILE_STRUCTURE_QUERY}
{CONTENT_RULES_QUERY}
{INSIGHTS_GROUNDING_QUERY}
{TONE_QUERY}
"""
# Versioned direct prompt (SUE-1265): p3 = V3.1 (tech-tuned). Rollback = DIRECT_PROMPT_FILE=direct_luna_v5_p2.md (V3) or direct_luna_v5.md (V1)
DIRECT_PROMPT_FILE = os.getenv("DIRECT_PROMPT_FILE", "direct_luna_v5_p3.md").strip()
DIRECT_LUNA_PROMPT = (_PROJECT_ROOT / "prompt" / DIRECT_PROMPT_FILE).read_text(encoding="utf-8")
DIRECT_PROMPT_SHA256 = hashlib.sha256(DIRECT_LUNA_PROMPT.encode("utf-8")).hexdigest()
# Legacy alias for prompt_log compatibility
INPUT_QUERY = INPUT_PROMPT
INPUT_ROLE = """
A smart assistant specialized in organizing content with expertise in analyzing and providing insights on the given material.
"""
MAIN_TASK_TYPE = "gpt_summarization"
def _build_transcript_cache(config: dict) -> TranscriptCache:
work_path = config.get('WORK_PATH') or config.get('BASE_PATH') or str(_PROJECT_ROOT)
cache_root = os.path.join(work_path, 'cache')
return TranscriptCache(
cache_root,
enabled=bool(config.get('TRANSCRIPT_CACHE_ENABLED', True)),
ttl_hours=int(config.get('TRANSCRIPT_CACHE_TTL_HOURS', 72)),
)
def _run_batch_cleanup(config: dict, *, dry_run_legacy: bool = False) -> None:
"""Expire cache/quarantine/stale jobs; optional legacy yt_subs dry-run inventory."""
work_path = config.get('WORK_PATH') or config.get('BASE_PATH') or str(_PROJECT_ROOT)
cache = _build_transcript_cache(config)
n_cache = cache.cleanup_expired()
if n_cache:
logger.info("Cleaned %d expired transcript cache entries", n_cache)
n_q = cleanup_expired_quarantine(
work_path,
max_age_days=int(config.get('SUBTITLE_QUARANTINE_DAYS', 7)),
)
if n_q:
logger.info("Cleaned %d expired quarantine subtitle dirs", n_q)
max_job_age = int(config.get('STALE_JOB_MAX_AGE_HOURS', 6)) * 3600
removed, skipped = cleanup_stale_jobs(work_path, max_age_sec=max_job_age)
if removed or skipped:
logger.info("Stale job cleanup: removed=%d skipped_active=%d", removed, skipped)
if dry_run_legacy and config.get('USE_JOB_WORKSPACE', True):
yt_subs = os.path.join(work_path, 'yt_subs')
report = cleanup_legacy_yt_subs(
work_path,
yt_subs,
os.path.join(config['BASE_PATH'], 'output_new', 'full'),
os.path.join(config['DATA_ROOT'], 'output_df_new.csv'),
config.get('OUTPUT_MD_PATH', ''),
dry_run=True,
)
logger.info(
"Legacy yt_subs dry-run: before=%d files (%.1f MB) deletable=%d quarantine=%d preserve=%d",
report.before_count,
report.before_bytes / (1024 * 1024),
report.deleted_count,
report.quarantined_count,
report.preserved_count,
)
def derive_output_name(txt_file_name: str, llm_suffix: str, video_id: str, ext: str = "") -> str:
"""Name for a file derived from the transcript .txt name (concise txt or final md).
Adds the LLM suffix (and optionally swaps the extension), then refits to 255 bytes
so the video ID survives; every suffix/extension change must go through here.
"""
name = stt.change_filename(txt_file_name, f"_{llm_suffix}")
if ext:
name = stt.change_extension(name, ext)
return fit_filename(name, video_id or "")
def resolve_obsidian_mirror(v_url: str, direct_input_urls, channel_marker_by_url) -> bool:
"""Obsidian mirror decision for one video, by input source.
A URL the owner put in input_df.csv is a specific video request and always mirrors.
Channel-crawl videos mirror only when channel_df marks the channel `recording`;
a missing entry (unknown channel) is off.
"""
if v_url in direct_input_urls:
return True
return bool(channel_marker_by_url.get(v_url, False))
def sanitize_channel_name(name: str) -> str:
"""Sanitize channel name for filename (filesystem-safe chars only, no length limit)."""
if not name or not str(name).strip():
return ""
return re.sub(r'[/\\:*?"<>|]', '_', str(name).strip())
def _is_transient_md_write_errno(exc: BaseException) -> bool:
"""Retry macOS EDEADLK / errno 11; do not retry ENOSPC (28)."""
n = getattr(exc, "errno", None)
if n == 28:
return False
if n == 11:
return True
s = str(exc).lower()
return "errno 11" in s or "resource deadlock" in s
def atomic_write_text_with_retry(
final_path: str,
content: str,
*,
encoding: str = "utf-8-sig",
max_attempts: int = 8,
log: Optional[logging.Logger] = None,
) -> None:
"""
Write text via same-directory temp file + os.replace (atomic on same volume).
Retries on transient I/O (e.g. Errno 11 on iCloud/sync paths).
"""
lg = log or logging.getLogger(__name__)
directory = os.path.dirname(os.path.abspath(final_path))
os.makedirs(directory, exist_ok=True)
base_name = os.path.basename(final_path)
# Short temp name: a 255-byte final name must not make the temp name overflow.
tmp_path = os.path.join(
directory,
f".{hashlib.sha1(base_name.encode('utf-8')).hexdigest()[:12]}.{os.getpid()}.tmp",
)
last_exc: Optional[Exception] = None
for attempt in range(max_attempts):
try:
with open(tmp_path, "w", encoding=encoding) as f:
f.write(content)
os.replace(tmp_path, final_path)
return
except OSError as e:
last_exc = e
if _is_transient_md_write_errno(e) and attempt < max_attempts - 1:
wait = 0.45 * (attempt + 1) + random.uniform(0, 0.35)
lg.warning(
"Markdown write transient I/O (attempt %d/%d), retry in %.2fs: %s",
attempt + 1,
max_attempts,
wait,
e,
)
time.sleep(wait)
try:
if os.path.isfile(tmp_path):
os.remove(tmp_path)
except OSError:
pass
if not _is_transient_md_write_errno(e) or attempt >= max_attempts - 1:
raise
except Exception:
try:
if os.path.isfile(tmp_path):
os.remove(tmp_path)
except OSError:
pass
raise
if last_exc:
raise last_exc
def process_single_video(
v_url: str,
config: dict,
openai_client: OpenAI,
main_llm: MainLlmConfig,
output_df: pd.DataFrame,
base_path: str,
audio_path: str,
output_full_path: str,
output_smm_path: str,
output_md_path: str,
prompt_log_path: str,
hf_path: str,
*,
run_id: str = "",
worker_id: str = "w0",
download_limiter: Optional[DownloadAdmissionLimiter] = None,
provider_cooldown: Optional[ProviderCooldown] = None,
preprocessor: Optional[TranscriptPreprocessor] = None,
) -> VideoProcessResult:
"""
Process a single video: download, transcribe, summarize, and save.
Does not mutate shared CSV/queue — returns VideoProcessResult for single writer.
"""
def _vr(
status: str,
vid: Optional[str],
err: Optional[str] = None,
*,
stage: str = "complete",
transcript_source: Optional[str] = None,
output_md_path_res: Optional[str] = None,
metadata_updates: Optional[dict] = None,
catalog_updates: Optional[dict] = None,
prompt_entries: Optional[list] = None,
cache_hit: bool = False,
) -> VideoProcessResult:
retryable = status in {
"download_failed", "api_error", "file_error", "mlx_error", "error",
}
return VideoProcessResult(
video_id=vid or "unknown",
source_url=v_url,
status=status,
stage=stage,
retryable=retryable,
error_message=err,
output_md_path=output_md_path_res,
transcript_source=transcript_source,
transcript_cache_hit=cache_hit,