Files
music_bot_v2/soloist/bridge.py
tkrmagid 8087e2e98b feat: 스포티파이 곡을 Soloist(스포티파이 공식 클라이언트)로 먼저 재생, 안 되면 유튜브로
- soloist/: soloist-bridge 중계 서버(docker). 곡마다 soloist --single-track을 띄우고
  PulseAudio 빈 싱크를 캡처해 MP3 스트림으로 Lavalink(http 소스)에 넘긴다.
  첫 소리가 나와야 성공으로 응답, 15초 안에 소리가 없거나 오류면 실패.
  한 계정 한 곡 제한 → 다른 길드가 쓰는 중이면 busy. 빌드 만료(90일) 자동 갱신, 페어링 대기 자동 진입.
- bot: SOLOIST_URL 설정 시 스포티파이 곡은 재생 직전에 Soloist를 먼저 시도, 실패하면 기존 유튜브 경로.
  곡 도중 끊김·위치 이동·음성 재접속은 같은 곡의 유튜브 음원으로 그 위치부터 이어 재생.
  오래 걸려 실패하면 5분 동안 유튜브로만 재생.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-03 20:33:46 +09:00

553 lines
23 KiB
Python

#!/usr/bin/env python3
"""
soloist-bridge — 스포티파이 공식 헤드리스 클라이언트(Soloist)로 곡을 재생하고, 그 소리를 MP3 스트림으로
Lavalink(http 소스)에 넘겨 주는 중계 서버.
흐름 (곡 하나마다)
봇 → GET /prepare?track=<스포티파이 곡 id>&owner=<봇:길드>&duration=<ms>
→ parec(빈 싱크 모니터) + ffmpeg(MP3 인코딩)을 먼저 띄우고 soloist --single-track 실행
→ 앞쪽 무음을 버리고 실제 소리가 나오기 시작하면 {"path": "/stream/<sid>.mp3"} 응답
(정해진 시간 안에 소리가 없거나 soloist가 먼저 끝나면 오류 응답 → 봇이 유튜브로 재생)
Lavalink → GET /stream/<sid>.mp3 (처음부터 끝까지 받은 MP3를 그대로 흘려 보냄)
→ soloist가 곡을 다 재생하고 종료하면 스트림도 끝남(→ Lavalink "finished")
→ 곡 도중에 끊기면(다른 기기에서 같은 계정 재생 등) 청크 종료 없이 연결을 끊어 Lavalink가 오류로 보게 함
계정 하나는 한 번에 한 곡만 재생할 수 있으므로 세션은 동시에 하나만 돈다. 다른 길드가 쓰는 중이면 409(busy)를 준다.
같은 owner의 새 요청은 이전 세션을 끊고 새로 시작한다(건너뛰기).
환경변수
SOLOIST_API_KEY 개발자 대시보드에서 만든 Soloist API 키(필수, 없으면 비활성 → /prepare 503)
SOLOIST_DEVICE_NAME 스포티파이 앱에 보일 기기 이름(기본 music_bot)
SOLOIST_DATA 세션·실행 파일 보관 폴더(기본 /data, 볼륨으로 유지해야 페어링이 남는다)
SOLOIST_BIN soloist 실행 파일 경로(기본 $SOLOIST_DATA/bin/soloist, 없으면 자동 다운로드)
SOLOIST_PORT HTTP 포트(기본 8780)
SOLOIST_PREPARE_TIMEOUT 첫 소리가 나올 때까지 기다리는 최대 초(기본 15)
PULSE_SINK 캡처할 PulseAudio 싱크 이름(기본 soloist)
"""
import json
import os
import re
import shutil
import subprocess
import sys
import tarfile
import threading
import time
import urllib.request
import uuid
from array import array
from collections import deque
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from urllib.parse import parse_qs, urlparse
PORT = int(os.environ.get("SOLOIST_PORT", "8780"))
API_KEY = os.environ.get("SOLOIST_API_KEY", "").strip()
DEVICE_NAME = os.environ.get("SOLOIST_DEVICE_NAME", "music_bot").strip() or "music_bot"
DATA = os.environ.get("SOLOIST_DATA", "/data")
BIN = os.environ.get("SOLOIST_BIN", "").strip() or os.path.join(DATA, "bin", "soloist")
STATE_DIR = os.path.join(DATA, "state")
CACHE_DIR = os.path.join(DATA, "cache")
SINK = os.environ.get("PULSE_SINK", "soloist")
PREPARE_TIMEOUT = float(os.environ.get("SOLOIST_PREPARE_TIMEOUT", "15"))
BUILD_URL = os.environ.get(
"SOLOIST_BUILD_URL", "https://soloist-builds.spotifycdn.com/soloist_release_x86_64.tar.gz")
RATE, CHANNELS = 48000, 2
PCM_BYTES_PER_SEC = RATE * CHANNELS * 2
SILENCE_LEVEL = 8 # 이 값보다 큰 샘플이 나오면 '소리 시작'으로 본다(빈 싱크의 무음은 정확히 0)
TAIL_SEC = 0.8 # soloist 종료 후 싱크에 남은 소리를 더 받는 시간
IDLE_SEC = 10 # 스트림을 받는 쪽이 없을 때 세션을 정리하기까지의 시간
MAX_BUF = 80 * 1024 * 1024 # 세션당 MP3 보관 상한(320kbps 기준 30분 이상)
TRACK_RE = re.compile(r"^[A-Za-z0-9]{22}$")
EXIT_EXPIRED = 10 # soloist: 빌드 만료
def log(msg: str) -> None:
print(f"[soloist-bridge] {msg}", flush=True)
class Soloist:
"""soloist 실행 파일·페어링 상태 관리."""
def __init__(self) -> None:
self.lock = threading.Lock()
self.version = ""
self.expires_days: int | None = None
self.pairing = False
self.last_update = 0.0
self.last_error = ""
@property
def enabled(self) -> bool:
return bool(API_KEY)
@property
def paired(self) -> bool:
# soloist는 로그인 토큰을 <data-dir>/cache/dbrts 에 저장한다(페어링 전에는 없거나 0바이트).
p = os.path.join(STATE_DIR, "cache", "dbrts")
return os.path.exists(p) and os.path.getsize(p) > 0
def base_args(self) -> list[str]:
return [BIN, "-n", DEVICE_NAME, "-k", API_KEY, "-D", STATE_DIR, "-C", CACHE_DIR, "-i", "100", "-v"]
def read_version(self) -> None:
try:
out = subprocess.run([BIN, "--version"], capture_output=True, text=True, timeout=15)
self.version = (out.stdout or out.stderr).strip()
except Exception as e: # noqa: BLE001
self.version = ""
log(f"버전 확인 실패: {e}")
def note_output(self, line: str) -> None:
m = re.search(r"client expires in (-?\d+) days", line)
if m:
self.expires_days = int(m.group(1))
if self.expires_days <= 7:
self.update_async(f"만료 {self.expires_days}일 전")
def update(self, reason: str) -> bool:
"""최신 빌드를 받아 실행 파일을 교체한다. 실행 중인 프로세스는 이전 파일을 계속 쓴다."""
with self.lock:
if time.time() - self.last_update < 3600 and os.path.exists(BIN):
return False
self.last_update = time.time()
log(f"soloist 빌드 내려받기 ({reason})")
try:
os.makedirs(os.path.dirname(BIN), exist_ok=True)
tmp_tar = BIN + ".tar.gz"
with urllib.request.urlopen(BUILD_URL, timeout=120) as r, open(tmp_tar, "wb") as f:
shutil.copyfileobj(r, f)
with tarfile.open(tmp_tar) as tar:
member = tar.getmember("soloist")
src = tar.extractfile(member)
if src is None:
raise RuntimeError("tar 안에 soloist 파일이 없습니다")
with open(BIN + ".new", "wb") as f:
shutil.copyfileobj(src, f)
os.chmod(BIN + ".new", 0o755)
check = subprocess.run([BIN + ".new", "--version"], capture_output=True, text=True, timeout=15)
if check.returncode != 0:
raise RuntimeError(f"새 빌드 실행 실패(code {check.returncode})")
os.replace(BIN + ".new", BIN)
os.remove(tmp_tar)
self.expires_days = None
self.read_version()
log(f"soloist 갱신 완료: {self.version}")
return True
except Exception as e: # noqa: BLE001
log(f"soloist 갱신 실패: {e}")
return False
def update_async(self, reason: str) -> None:
threading.Thread(target=self.update, args=(reason,), daemon=True).start()
def start_pairing(self) -> None:
"""페어링이 안 돼 있으면 soloist --pair 를 띄워 둔다. 같은 네트워크의 스포티파이 앱 기기 목록에서
이 기기를 한 번 고르면 로그인 정보가 저장되고 종료된다."""
if not self.enabled:
return
with self.lock:
if self.pairing:
return
self.pairing = True
threading.Thread(target=self._pair_loop, daemon=True).start()
def _pair_loop(self) -> None:
try:
while not self.paired:
if sessions.running():
time.sleep(5)
continue
log(f'페어링 대기: 스포티파이 앱의 기기 목록에서 "{DEVICE_NAME}"를 고르세요')
proc = subprocess.Popen(self.base_args() + ["-p"], stdout=subprocess.PIPE,
stderr=subprocess.STDOUT, text=True, errors="replace")
assert proc.stdout is not None
for line in proc.stdout:
line = strip_ansi(line.rstrip())
self.note_output(line)
if " E " in line or " W " in line or "pair" in line.lower():
log(f"[pair] {line}")
code = proc.wait()
if code == 0 and self.paired:
log("페어링 완료")
break
if code == EXIT_EXPIRED:
self.update("빌드 만료")
self.last_error = f"pair exit {code}"
time.sleep(30)
finally:
with self.lock:
self.pairing = False
ANSI_RE = re.compile(r"\x1b\[[0-9;]*m")
def strip_ansi(s: str) -> str:
return ANSI_RE.sub("", s)
def kill(proc: subprocess.Popen | None, wait: float = 3.0) -> None:
if proc is None or proc.poll() is not None:
return
try:
proc.terminate()
proc.wait(timeout=wait)
except Exception: # noqa: BLE001
try:
proc.kill()
proc.wait(timeout=wait)
except Exception: # noqa: BLE001
pass
class Session:
def __init__(self, owner: str, track: str, duration_ms: int) -> None:
self.id = uuid.uuid4().hex
self.owner = owner
self.track = track
self.duration_ms = duration_ms
self.created = time.time()
self.buf = bytearray()
self.cond = threading.Condition()
self.started = threading.Event() # 실제 소리가 나오기 시작함
self.finished = False # MP3 출력이 끝남(정상/비정상 모두)
self.truncated = False # 곡 도중에 끊김 → 클라이언트에 오류로 알림
self.error = ""
self.clients = 0
self.last_client_seen = time.time()
self.fed_bytes = 0
self.solo_exit: int | None = None
self.solo_exit_at = 0.0
self.stopping = False
self.logs: deque[str] = deque(maxlen=40)
self.parec: subprocess.Popen | None = None
self.ffmpeg: subprocess.Popen | None = None
self.solo: subprocess.Popen | None = None
# ── 상태 ──
@property
def solo_running(self) -> bool:
return self.solo is not None and self.solo.poll() is None
@property
def played_ms(self) -> int:
return self.fed_bytes * 1000 // PCM_BYTES_PER_SEC
def summary(self) -> dict:
return {"id": self.id, "owner": self.owner, "track": self.track, "started": self.started.is_set(),
"finished": self.finished, "truncated": self.truncated, "error": self.error,
"playedMs": self.played_ms, "clients": self.clients, "bytes": len(self.buf)}
# ── 실행 ──
def start(self) -> None:
self.parec = subprocess.Popen(
["parec", "-d", f"{SINK}.monitor", "--format=s16le", f"--rate={RATE}", f"--channels={CHANNELS}",
"--raw", "--latency-msec=50"], stdout=subprocess.PIPE, stderr=subprocess.DEVNULL)
self.ffmpeg = subprocess.Popen(
["ffmpeg", "-hide_banner", "-loglevel", "error",
# 입력 분석(기본 최대 5초 분량)을 끄지 않으면 첫 MP3가 몇 초 늦게 나온다
"-probesize", "32", "-analyzeduration", "0", "-fflags", "nobuffer", "-f", "s16le", "-ar", str(RATE), "-ac", str(CHANNELS),
"-i", "pipe:0", "-c:a", "libmp3lame", "-b:a", "320k", "-write_xing", "0", "-id3v2_version", "0",
"-flush_packets", "1", "-f", "mp3", "pipe:1"],
stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.DEVNULL)
self.solo = subprocess.Popen(
soloist.base_args() + ["-s", f"spotify:track:{self.track}"],
stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, errors="replace")
threading.Thread(target=self._read_solo, daemon=True).start()
threading.Thread(target=self._pump_pcm, daemon=True).start()
threading.Thread(target=self._pump_mp3, daemon=True).start()
def _read_solo(self) -> None:
assert self.solo is not None and self.solo.stdout is not None
for line in self.solo.stdout:
line = strip_ansi(line.rstrip())
if not line:
continue
self.logs.append(line)
soloist.note_output(line)
code = self.solo.wait()
self.solo_exit, self.solo_exit_at = code, time.time()
if code == EXIT_EXPIRED:
soloist.update_async("빌드 만료")
if not self.started.is_set():
self.fail(self.reason_from_logs(code))
def reason_from_logs(self, code: int) -> str:
text = "\n".join(self.logs)
if code == EXIT_EXPIRED:
return "build_expired"
if "requires stored credentials" in text or "No session to restore" in text:
soloist.start_pairing()
return "not_paired"
errs = [line for line in self.logs if " E " in line or " W " in line]
return f"soloist exit {code}: " + (errs[-1] if errs else (self.logs[-1] if self.logs else ""))
def _pump_pcm(self) -> None:
"""parec → (앞쪽 무음 버림) → ffmpeg. soloist가 끝나면 TAIL_SEC 뒤 ffmpeg 입력을 닫는다."""
assert self.parec and self.parec.stdout and self.ffmpeg and self.ffmpeg.stdin
fd = self.parec.stdout.fileno()
pending = b""
try:
while not self.stopping:
if self.solo_exit is not None and time.time() - self.solo_exit_at > TAIL_SEC:
break
chunk = os.read(fd, 9600) # 50ms
if not chunk:
break
if not self.started.is_set():
chunk = pending + chunk
usable = len(chunk) - (len(chunk) % 4)
pending = chunk[usable:]
samples = array("h", chunk[:usable])
if not samples or (max(samples) <= SILENCE_LEVEL and -min(samples) <= SILENCE_LEVEL):
continue
# 첫 소리가 있는 프레임부터 보낸다
idx = next(i for i, v in enumerate(samples) if v > SILENCE_LEVEL or -v > SILENCE_LEVEL)
chunk = chunk[(idx - idx % CHANNELS) * 2:usable]
self.started.set()
log(f"[{self.id[:8]}] 소리 시작 ({time.time() - self.created:.2f}초) track={self.track}")
self.ffmpeg.stdin.write(chunk)
self.fed_bytes += len(chunk)
except (BrokenPipeError, OSError, ValueError):
pass
finally:
kill(self.parec, 1)
try:
self.ffmpeg.stdin.close()
except Exception: # noqa: BLE001
pass
if self.started.is_set() and not self.stopping:
expected = self.duration_ms
if (self.solo_exit not in (0, None)) or (expected and self.played_ms < expected - 10000):
self.truncated = True
log(f"[{self.id[:8]}] 곡 도중 끊김: 재생 {self.played_ms}ms / 곡 {expected}ms, "
f"soloist exit={self.solo_exit} {self.logs[-1] if self.logs else ''}")
def _pump_mp3(self) -> None:
assert self.ffmpeg and self.ffmpeg.stdout
fd = self.ffmpeg.stdout.fileno()
while True:
data = os.read(fd, 65536)
if not data:
break
with self.cond:
self.buf += data
self.cond.notify_all()
if len(self.buf) > MAX_BUF:
log(f"[{self.id[:8]}] 버퍼 상한 초과, 세션 종료")
self.stop(truncate=True)
break
self.ffmpeg.wait()
with self.cond:
self.finished = True
self.cond.notify_all()
kill(self.solo)
def fail(self, reason: str) -> None:
if not self.error:
self.error = reason
self.stop(truncate=True)
self.started.set() # prepare 대기를 깨운다(error로 판별)
def stop(self, truncate: bool = False) -> None:
if truncate:
self.truncated = True
self.stopping = True
kill(self.solo)
kill(self.parec, 1)
if self.ffmpeg and self.ffmpeg.stdin:
try:
self.ffmpeg.stdin.close()
except Exception: # noqa: BLE001
pass
with self.cond:
self.cond.notify_all()
class Sessions:
def __init__(self) -> None:
self.lock = threading.Lock()
self.items: dict[str, Session] = {}
def running(self) -> Session | None:
with self.lock:
return next((s for s in self.items.values() if s.solo_running or (not s.started.is_set() and not s.error)), None)
def get(self, sid: str) -> Session | None:
with self.lock:
return self.items.get(sid)
def create(self, owner: str, track: str, duration_ms: int) -> tuple[Session | None, str]:
with self.lock:
active = [s for s in self.items.values() if s.solo_running or (not s.started.is_set() and not s.error)]
for s in active:
in_use = s.clients > 0 or time.time() - s.last_client_seen < IDLE_SEC
if s.owner != owner and in_use:
return None, "busy"
for s in active:
if s.started.is_set():
s.stop(truncate=True)
else:
s.fail("preempted") # 아직 준비 중인 요청은 바로 실패로 돌려준다
session = Session(owner, track, duration_ms)
self.items[session.id] = session
# 이전 soloist가 데이터 폴더 잠금(.lock)을 놓을 때까지 기다린다
for s in active:
if s.solo:
try:
s.solo.wait(timeout=5)
except Exception: # noqa: BLE001
pass
session.start()
return session, ""
def reap(self) -> None:
while True:
time.sleep(1)
now = time.time()
with self.lock:
items = list(self.items.values())
for s in items:
idle = s.clients == 0 and now - s.last_client_seen > IDLE_SEC
if idle and not s.finished and s.started.is_set():
log(f"[{s.id[:8]}] 받는 쪽이 없어 세션 종료")
s.stop(truncate=True)
if (s.finished or s.error) and idle:
with self.lock:
self.items.pop(s.id, None)
soloist = Soloist()
sessions = Sessions()
class Handler(BaseHTTPRequestHandler):
protocol_version = "HTTP/1.1"
def log_message(self, fmt, *args): # noqa: N802
pass
def send_json(self, code: int, obj: dict) -> None:
body = json.dumps(obj, ensure_ascii=False).encode()
self.send_response(code)
self.send_header("Content-Type", "application/json; charset=utf-8")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)
def do_GET(self): # noqa: N802
u = urlparse(self.path)
q = {k: v[0] for k, v in parse_qs(u.query).items()}
if u.path == "/health":
run = sessions.running()
return self.send_json(200, {
"enabled": soloist.enabled, "paired": soloist.paired, "pairing": soloist.pairing,
"deviceName": DEVICE_NAME, "version": soloist.version, "expiresInDays": soloist.expires_days,
"busy": run.summary() if run else None})
if u.path == "/prepare":
return self.prepare(q)
m = re.fullmatch(r"/stream/([0-9a-f]{32})\.mp3", u.path)
if m:
return self.stream(m.group(1))
self.send_json(404, {"error": "not_found"})
def prepare(self, q: dict) -> None:
track = q.get("track", "")
owner = q.get("owner", "")
try:
duration = int(q.get("duration", "0") or 0)
except ValueError:
duration = 0
if not TRACK_RE.match(track) or not owner:
return self.send_json(400, {"error": "bad_request"})
if not soloist.enabled:
return self.send_json(503, {"error": "disabled"})
if not soloist.paired:
soloist.start_pairing()
return self.send_json(503, {"error": "not_paired"})
if soloist.pairing:
return self.send_json(503, {"error": "pairing"})
t0 = time.time()
session, err = sessions.create(owner, track, duration)
if not session:
return self.send_json(409, {"error": err})
session.started.wait(PREPARE_TIMEOUT)
if not session.started.is_set():
session.fail("no_audio_timeout")
if session.error:
log(f"[{session.id[:8]}] 준비 실패 track={track}: {session.error}")
return self.send_json(502, {"error": session.error})
session.last_client_seen = time.time()
self.send_json(200, {"path": f"/stream/{session.id}.mp3", "id": session.id,
"readyMs": int((time.time() - t0) * 1000)})
def stream(self, sid: str) -> None:
s = sessions.get(sid)
if not s or s.error:
return self.send_json(404, {"error": "no_session"})
self.send_response(200)
self.send_header("Content-Type", "audio/mpeg")
self.send_header("Transfer-Encoding", "chunked")
self.send_header("Cache-Control", "no-store")
self.end_headers()
self.close_connection = True
with s.cond:
s.clients += 1
s.last_client_seen = time.time()
pos = 0
clean = False
t_open = time.time()
try:
while True:
with s.cond:
while len(s.buf) <= pos and not s.finished and not s.stopping:
s.cond.wait(1.0)
data = bytes(s.buf[pos:pos + 65536])
done = (s.finished or s.stopping) and pos + len(data) >= len(s.buf)
if data:
self.wfile.write(b"%x\r\n" % len(data) + data + b"\r\n")
pos += len(data)
if done:
if not s.truncated:
self.wfile.write(b"0\r\n\r\n")
clean = True
break
self.wfile.flush()
except (BrokenPipeError, ConnectionResetError, OSError):
pass
finally:
with s.cond:
s.clients -= 1
s.last_client_seen = time.time()
log(f"[{s.id[:8]}] 스트림 연결 종료: {pos}바이트, {time.time() - t_open:.1f}초, 정상종료={clean}")
if not clean:
# 청크 종료 없이 끊어 받는 쪽이 '비정상 종료'로 알게 한다
try:
self.connection.shutdown(2)
except OSError:
pass
def main() -> None:
if not soloist.enabled:
log("SOLOIST_API_KEY 없음 → 비활성 상태로 시작(/prepare는 503, 봇은 유튜브로 재생)")
else:
if not os.path.exists(BIN):
soloist.update("실행 파일 없음")
soloist.read_version()
log(f"soloist: {soloist.version or '(없음)'} / 기기 이름: {DEVICE_NAME} / 페어링: {soloist.paired}")
os.makedirs(STATE_DIR, exist_ok=True)
os.makedirs(CACHE_DIR, exist_ok=True)
if not soloist.paired:
soloist.start_pairing()
threading.Thread(target=sessions.reap, daemon=True).start()
server = ThreadingHTTPServer(("0.0.0.0", PORT), Handler)
server.daemon_threads = True
log(f"listening on :{PORT}")
server.serve_forever()
if __name__ == "__main__":
sys.exit(main())