refactor: extract subprocess cleanup into shared subproc helper (#210)
bird_x.py and youtube_yt.py had four near-identical copies of the same subprocess cleanup dance (Popen + os.setsid + communicate(timeout) + SIGTERM via killpg + proc.kill() fallback + wait(5)). Extract to lib.subproc.run_with_timeout(), which: - runs the child in its own process group via os.setsid where available - raises SubprocTimeout on timeout - on timeout: SIGTERM the group, fall back to proc.kill(), wait up to 5s - accepts an on_pid callback so bird_x can still register child PIDs with last30days.register_child_pid for whole-process cleanup - captures stdout/stderr as strings in a SubprocResult dataclass Migrated call sites: _run_bird_search, search_handles inner worker, search_youtube, fetch_transcript. With the helper in place, the signal and subprocess imports became dead in both files (plus os in youtube_yt) and went with them. Tests: 9 new subproc tests cover success, non-zero exit, stderr capture, timeout-raises, timeout-kills-group, missing-command, env passthrough, PID callback, and callback-exception suppression. test_env_v3 and test_youtube_yt patch subproc.run_with_timeout instead of the removed bird_x.subprocess and yt-dlp subprocess.
This commit is contained in:
@@ -7,13 +7,11 @@ See scripts/lib/vendor/bird-search/package.json for authoritative version.
|
||||
|
||||
import json
|
||||
import os
|
||||
import signal
|
||||
import shutil
|
||||
import subprocess
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
from . import http, log
|
||||
from . import http, log, subproc
|
||||
from datetime import datetime
|
||||
from typing import Any, Dict, List, Optional, Tuple
|
||||
|
||||
@@ -168,62 +166,51 @@ def _run_bird_search(query: str, count: int, timeout: int) -> Dict[str, Any]:
|
||||
"--json",
|
||||
]
|
||||
|
||||
# Use process groups for clean cleanup on timeout/kill
|
||||
preexec = os.setsid if hasattr(os, 'setsid') else None
|
||||
pid_holder: list[int] = []
|
||||
|
||||
try:
|
||||
proc = subprocess.Popen(
|
||||
cmd,
|
||||
stdout=subprocess.PIPE,
|
||||
stderr=subprocess.PIPE,
|
||||
text=True,
|
||||
encoding="utf-8",
|
||||
errors="replace",
|
||||
preexec_fn=preexec,
|
||||
env=_subprocess_env(),
|
||||
)
|
||||
|
||||
# Register for cleanup tracking (if available)
|
||||
def _register(pid: int) -> None:
|
||||
pid_holder.append(pid)
|
||||
try:
|
||||
from last30days import register_child_pid, unregister_child_pid
|
||||
register_child_pid(proc.pid)
|
||||
from last30days import register_child_pid
|
||||
register_child_pid(pid)
|
||||
except ImportError:
|
||||
pass
|
||||
|
||||
try:
|
||||
stdout, stderr = proc.communicate(timeout=timeout)
|
||||
except subprocess.TimeoutExpired:
|
||||
# Kill the entire process group
|
||||
try:
|
||||
os.killpg(os.getpgid(proc.pid), signal.SIGTERM)
|
||||
except (ProcessLookupError, PermissionError, OSError):
|
||||
proc.kill()
|
||||
proc.wait(timeout=5)
|
||||
return {"error": f"Search timed out after {timeout}s", "items": []}
|
||||
finally:
|
||||
try:
|
||||
result = subproc.run_with_timeout(
|
||||
cmd,
|
||||
timeout=timeout,
|
||||
env=_subprocess_env(),
|
||||
on_pid=_register,
|
||||
)
|
||||
except subproc.SubprocTimeout:
|
||||
return {"error": f"Search timed out after {timeout}s", "items": []}
|
||||
except Exception as e:
|
||||
return {"error": str(e), "items": []}
|
||||
finally:
|
||||
if pid_holder:
|
||||
try:
|
||||
from last30days import unregister_child_pid
|
||||
unregister_child_pid(proc.pid)
|
||||
unregister_child_pid(pid_holder[0])
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
if proc.returncode != 0:
|
||||
error = stderr.strip() if stderr else "Bird search failed"
|
||||
return {"error": error, "items": []}
|
||||
if result.returncode != 0:
|
||||
error = result.stderr.strip() or "Bird search failed"
|
||||
return {"error": error, "items": []}
|
||||
|
||||
output = stdout.strip() if stdout else ""
|
||||
if not output:
|
||||
return {"items": []}
|
||||
output = result.stdout.strip()
|
||||
if not output:
|
||||
return {"items": []}
|
||||
|
||||
try:
|
||||
parsed = json.loads(output)
|
||||
if isinstance(parsed, list):
|
||||
return {"items": parsed}
|
||||
return parsed
|
||||
|
||||
except json.JSONDecodeError as e:
|
||||
return {"error": f"Invalid JSON response: {e}", "items": []}
|
||||
except Exception as e:
|
||||
return {"error": str(e), "items": []}
|
||||
|
||||
if isinstance(parsed, list):
|
||||
return {"items": parsed}
|
||||
return parsed
|
||||
|
||||
|
||||
def search_x(
|
||||
@@ -330,47 +317,29 @@ def search_handles(
|
||||
"--json",
|
||||
]
|
||||
|
||||
preexec = os.setsid if hasattr(os, 'setsid') else None
|
||||
try:
|
||||
result = subproc.run_with_timeout(cmd, timeout=15, env=_subprocess_env())
|
||||
except subproc.SubprocTimeout:
|
||||
_log(f"Handle search timed out for @{handle}")
|
||||
return []
|
||||
except OSError as e:
|
||||
_log(f"Handle search error for @{handle}: {e}")
|
||||
return []
|
||||
|
||||
if result.returncode != 0:
|
||||
_log(f"Handle search failed for @{handle}: {result.stderr.strip()}")
|
||||
return []
|
||||
|
||||
output = result.stdout.strip()
|
||||
if not output:
|
||||
return []
|
||||
|
||||
try:
|
||||
proc = subprocess.Popen(
|
||||
cmd,
|
||||
stdout=subprocess.PIPE,
|
||||
stderr=subprocess.PIPE,
|
||||
text=True,
|
||||
encoding="utf-8",
|
||||
errors="replace",
|
||||
preexec_fn=preexec,
|
||||
env=_subprocess_env(),
|
||||
)
|
||||
|
||||
try:
|
||||
stdout, stderr = proc.communicate(timeout=15)
|
||||
except subprocess.TimeoutExpired:
|
||||
try:
|
||||
os.killpg(os.getpgid(proc.pid), signal.SIGTERM)
|
||||
except (ProcessLookupError, PermissionError, OSError):
|
||||
proc.kill()
|
||||
proc.wait(timeout=5)
|
||||
_log(f"Handle search timed out for @{handle}")
|
||||
return []
|
||||
|
||||
if proc.returncode != 0:
|
||||
_log(f"Handle search failed for @{handle}: {(stderr or '').strip()}")
|
||||
return []
|
||||
|
||||
output = (stdout or "").strip()
|
||||
if not output:
|
||||
return []
|
||||
|
||||
response = json.loads(output)
|
||||
return parse_bird_response(response, query=core_topic)
|
||||
|
||||
except json.JSONDecodeError:
|
||||
_log(f"Invalid JSON from handle search for @{handle}")
|
||||
except (OSError, subprocess.SubprocessError) as e:
|
||||
_log(f"Handle search error for @{handle}: {e}")
|
||||
return []
|
||||
return []
|
||||
return parse_bird_response(response, query=core_topic)
|
||||
|
||||
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||
|
||||
|
||||
@@ -0,0 +1,94 @@
|
||||
"""Subprocess helpers: safe timeout + process-group cleanup.
|
||||
|
||||
Used by bird_x.py (Node.js Bird search) and youtube_yt.py (yt-dlp search
|
||||
and transcript download). Both need the same os.setsid/killpg cleanup
|
||||
dance on timeout to avoid orphaning child processes.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import signal
|
||||
import subprocess
|
||||
from dataclasses import dataclass
|
||||
from typing import Optional, Sequence
|
||||
|
||||
|
||||
class SubprocTimeout(Exception):
|
||||
"""Raised when a subprocess exceeds its timeout and is killed."""
|
||||
|
||||
|
||||
@dataclass
|
||||
class SubprocResult:
|
||||
"""Result of a subprocess run that captured stdout and stderr."""
|
||||
|
||||
returncode: int
|
||||
stdout: str
|
||||
stderr: str
|
||||
|
||||
|
||||
def run_with_timeout(
|
||||
cmd: Sequence[str],
|
||||
*,
|
||||
timeout: int,
|
||||
env: Optional[dict] = None,
|
||||
on_pid: Optional[callable] = None,
|
||||
) -> SubprocResult:
|
||||
"""Run a subprocess with process-group cleanup on timeout.
|
||||
|
||||
Spawns ``cmd`` inside its own process group via ``os.setsid`` where
|
||||
available. If ``communicate(timeout=...)`` raises ``TimeoutExpired``,
|
||||
signals ``SIGTERM`` to the entire group, falls back to ``proc.kill()``
|
||||
if the signal fails, then waits up to 5 seconds for cleanup, and
|
||||
raises ``SubprocTimeout``.
|
||||
|
||||
Args:
|
||||
cmd: Command and arguments to spawn.
|
||||
timeout: Timeout in seconds passed to ``communicate()``.
|
||||
env: Optional environment dict. If None, inherits parent env.
|
||||
on_pid: Optional callable invoked with the child PID right after
|
||||
spawn. Used by bird_x.py to register child PIDs for cleanup
|
||||
tracking. Exceptions raised by the callback are suppressed.
|
||||
|
||||
Returns:
|
||||
SubprocResult with returncode, stdout, and stderr as strings.
|
||||
|
||||
Raises:
|
||||
SubprocTimeout: If the process exceeded ``timeout``.
|
||||
FileNotFoundError: If the executable is not found.
|
||||
OSError: For other spawn failures.
|
||||
"""
|
||||
preexec = os.setsid if hasattr(os, "setsid") else None
|
||||
|
||||
proc = subprocess.Popen(
|
||||
list(cmd),
|
||||
stdout=subprocess.PIPE,
|
||||
stderr=subprocess.PIPE,
|
||||
text=True,
|
||||
encoding="utf-8",
|
||||
errors="replace",
|
||||
preexec_fn=preexec,
|
||||
env=env,
|
||||
)
|
||||
|
||||
if on_pid is not None:
|
||||
try:
|
||||
on_pid(proc.pid)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
try:
|
||||
stdout, stderr = proc.communicate(timeout=timeout)
|
||||
except subprocess.TimeoutExpired:
|
||||
try:
|
||||
os.killpg(os.getpgid(proc.pid), signal.SIGTERM)
|
||||
except (ProcessLookupError, PermissionError, OSError):
|
||||
proc.kill()
|
||||
proc.wait(timeout=5)
|
||||
raise SubprocTimeout(f"Command {cmd[0]} timed out after {timeout}s")
|
||||
|
||||
return SubprocResult(
|
||||
returncode=proc.returncode,
|
||||
stdout=stdout or "",
|
||||
stderr=stderr or "",
|
||||
)
|
||||
@@ -8,11 +8,8 @@ Inspired by Peter Steinberger's toolchain approach (yt-dlp + summarize CLI).
|
||||
|
||||
import json
|
||||
import math
|
||||
import os
|
||||
import re
|
||||
import signal
|
||||
import shutil
|
||||
import subprocess
|
||||
import sys
|
||||
import tempfile
|
||||
import urllib.error
|
||||
@@ -37,7 +34,7 @@ TRANSCRIPT_LIMITS = {
|
||||
# Max words to keep from each transcript
|
||||
TRANSCRIPT_MAX_WORDS = 5000
|
||||
|
||||
from . import http, log
|
||||
from . import http, log, subproc
|
||||
from .relevance import token_overlap_relevance as _compute_relevance
|
||||
|
||||
|
||||
@@ -227,30 +224,16 @@ def search_youtube(
|
||||
"--no-download",
|
||||
]
|
||||
|
||||
preexec = os.setsid if hasattr(os, 'setsid') else None
|
||||
|
||||
try:
|
||||
proc = subprocess.Popen(
|
||||
cmd,
|
||||
stdout=subprocess.PIPE,
|
||||
stderr=subprocess.PIPE,
|
||||
text=True,
|
||||
preexec_fn=preexec,
|
||||
)
|
||||
try:
|
||||
stdout, stderr = proc.communicate(timeout=120)
|
||||
except subprocess.TimeoutExpired:
|
||||
try:
|
||||
os.killpg(os.getpgid(proc.pid), signal.SIGTERM)
|
||||
except (ProcessLookupError, PermissionError, OSError):
|
||||
proc.kill()
|
||||
proc.wait(timeout=5)
|
||||
_log("YouTube search timed out (120s)")
|
||||
return {"items": [], "error": "Search timed out"}
|
||||
result = subproc.run_with_timeout(cmd, timeout=120)
|
||||
except subproc.SubprocTimeout:
|
||||
_log("YouTube search timed out (120s)")
|
||||
return {"items": [], "error": "Search timed out"}
|
||||
except FileNotFoundError:
|
||||
return {"items": [], "error": "yt-dlp not found"}
|
||||
|
||||
if not (stdout or "").strip():
|
||||
stdout = result.stdout
|
||||
if not stdout.strip():
|
||||
_log("YouTube search returned 0 results")
|
||||
return {"items": []}
|
||||
|
||||
@@ -452,25 +435,10 @@ def _fetch_transcript_ytdlp(video_id: str, temp_dir: str) -> Optional[str]:
|
||||
f"https://www.youtube.com/watch?v={video_id}",
|
||||
]
|
||||
|
||||
preexec = os.setsid if hasattr(os, 'setsid') else None
|
||||
|
||||
try:
|
||||
proc = subprocess.Popen(
|
||||
cmd,
|
||||
stdout=subprocess.PIPE,
|
||||
stderr=subprocess.PIPE,
|
||||
text=True,
|
||||
preexec_fn=preexec,
|
||||
)
|
||||
try:
|
||||
proc.communicate(timeout=30)
|
||||
except subprocess.TimeoutExpired:
|
||||
try:
|
||||
os.killpg(os.getpgid(proc.pid), signal.SIGTERM)
|
||||
except (ProcessLookupError, PermissionError, OSError):
|
||||
proc.kill()
|
||||
proc.wait(timeout=5)
|
||||
return None
|
||||
subproc.run_with_timeout(cmd, timeout=30)
|
||||
except subproc.SubprocTimeout:
|
||||
return None
|
||||
except FileNotFoundError:
|
||||
return None
|
||||
|
||||
@@ -556,7 +524,7 @@ def fetch_transcripts_parallel(
|
||||
vid = futures[future]
|
||||
try:
|
||||
results[vid] = future.result()
|
||||
except (OSError, subprocess.SubprocessError) as exc:
|
||||
except OSError as exc:
|
||||
_log(f"Transcript fetch error for {vid}: {exc}")
|
||||
results[vid] = None
|
||||
except Exception as exc:
|
||||
|
||||
Reference in New Issue
Block a user