Stop the lyrics worker from filling the disk: skip long audio, sweep stale temp files, back off after a crash
This commit is contained in:
@@ -27,6 +27,7 @@ import tempfile
|
|||||||
import time
|
import time
|
||||||
import urllib.error
|
import urllib.error
|
||||||
import urllib.request
|
import urllib.request
|
||||||
|
import glob
|
||||||
import http.cookiejar
|
import http.cookiejar
|
||||||
|
|
||||||
FILLER = re.compile(r"^(?:(?:oh|ooh|ohh|oh-oh|ah|ahh|hey|yeah|mm|mm-mm|mm-mm-mm|hmm|whoa|woah|la|na|uh|come on)[\s,.!?-]*)+$", re.I)
|
FILLER = re.compile(r"^(?:(?:oh|ooh|ohh|oh-oh|ah|ahh|hey|yeah|mm|mm-mm|mm-mm-mm|hmm|whoa|woah|la|na|uh|come on)[\s,.!?-]*)+$", re.I)
|
||||||
@@ -133,6 +134,24 @@ def main():
|
|||||||
run_once(args, Api(args.base, os.environ.get('YTP_TOKEN'), os.environ.get('YTP_ADMIN_PASSWORD')))
|
run_once(args, Api(args.base, os.environ.get('YTP_TOKEN'), os.environ.get('YTP_ADMIN_PASSWORD')))
|
||||||
|
|
||||||
|
|
||||||
|
# A whisper run killed by the container's memory cap skips the tempfile
|
||||||
|
# cleanup, and the restart used to retry the same song forever — 490 OOM
|
||||||
|
# kills left 61 GB of copies of one long video in /tmp. Long audio is skipped
|
||||||
|
# up front, stale temp files are swept on start, and a song that was in
|
||||||
|
# progress when the worker died is recorded as a failure (with backoff).
|
||||||
|
TMP_PREFIX = 'ytp-lyrics-'
|
||||||
|
MAX_AUDIO_BYTES = int(os.environ.get('LYRICS_MAX_AUDIO_MB', '40')) * 1048576
|
||||||
|
|
||||||
|
|
||||||
|
def sweep_stale_tmp():
|
||||||
|
for f in glob.glob(os.path.join(tempfile.gettempdir(), 'tmp*.m4a')) + \
|
||||||
|
glob.glob(os.path.join(tempfile.gettempdir(), TMP_PREFIX + '*')):
|
||||||
|
try:
|
||||||
|
os.remove(f)
|
||||||
|
except OSError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
def load_state(path):
|
def load_state(path):
|
||||||
try:
|
try:
|
||||||
with open(path) as f:
|
with open(path) as f:
|
||||||
@@ -165,6 +184,15 @@ def watch(args):
|
|||||||
time.sleep(3600)
|
time.sleep(3600)
|
||||||
args.missing, args.ids, args.overwrite, args.dry_run = True, None, False, False
|
args.missing, args.ids, args.overwrite, args.dry_run = True, None, False, False
|
||||||
model = None
|
model = None
|
||||||
|
sweep_stale_tmp()
|
||||||
|
state = load_state(args.state)
|
||||||
|
for vid, entry in state.items():
|
||||||
|
if entry.get('status') == 'in_progress': # the worker died mid-song
|
||||||
|
fails = entry.get('fails', 0) + 1
|
||||||
|
state[vid] = {'status': 'failed', 'fails': fails, 'retry_at': time.time() + min(86400, 900 * 2 ** fails),
|
||||||
|
'error': 'worker died while transcribing (likely out of memory)'}
|
||||||
|
print(f'lyrics worker: {vid} crashed the previous run — backing off', flush=True)
|
||||||
|
save_state(args.state, state)
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
api = Api(args.base, token, password)
|
api = Api(args.base, token, password)
|
||||||
@@ -182,9 +210,11 @@ def watch(args):
|
|||||||
print(f'lyrics worker: loading {args.model}', flush=True)
|
print(f'lyrics worker: loading {args.model}', flush=True)
|
||||||
model = WhisperModel(args.model, device='cpu', compute_type='int8', cpu_threads=args.threads)
|
model = WhisperModel(args.model, device='cpu', compute_type='int8', cpu_threads=args.threads)
|
||||||
m = todo[0] # one song per cycle keeps the worker's footprint small
|
m = todo[0] # one song per cycle keeps the worker's footprint small
|
||||||
result = transcribe_one(args, api, model, m['id'])
|
|
||||||
entry = state.get(m['id'], {})
|
entry = state.get(m['id'], {})
|
||||||
if result.startswith('instrumental'):
|
state[m['id']] = {**entry, 'status': 'in_progress'}
|
||||||
|
save_state(args.state, state)
|
||||||
|
result = transcribe_one(args, api, model, m['id'])
|
||||||
|
if result.startswith('instrumental') or result.startswith('too long'):
|
||||||
state[m['id']] = {'status': 'instrumental', 'at': now}
|
state[m['id']] = {'status': 'instrumental', 'at': now}
|
||||||
elif result.startswith('saved') or result.startswith('skip'):
|
elif result.startswith('saved') or result.startswith('skip'):
|
||||||
state.pop(m['id'], None)
|
state.pop(m['id'], None)
|
||||||
@@ -233,7 +263,9 @@ def transcribe_one(args, api, model, vid):
|
|||||||
st, audio = api.call('GET', f'/api/media/{vid}?a=1', raw=True)
|
st, audio = api.call('GET', f'/api/media/{vid}?a=1', raw=True)
|
||||||
if st != 200:
|
if st != 200:
|
||||||
return f'no cached audio ({st})'
|
return f'no cached audio ({st})'
|
||||||
with tempfile.NamedTemporaryFile(suffix='.m4a') as f:
|
if len(audio) > MAX_AUDIO_BYTES:
|
||||||
|
return f'too long: {len(audio) // 1048576} MB audio (cap {MAX_AUDIO_BYTES // 1048576} MB) — skipped'
|
||||||
|
with tempfile.NamedTemporaryFile(suffix='.m4a', prefix=TMP_PREFIX) as f:
|
||||||
f.write(audio)
|
f.write(audio)
|
||||||
f.flush()
|
f.flush()
|
||||||
t0 = time.time()
|
t0 = time.time()
|
||||||
|
|||||||
Reference in New Issue
Block a user