Spaces:
Running
Running
Reset main to commit fad3e19 - restore static/yt_live.js version
Browse files- dvr.py +0 -452
- gitattributes +0 -35
- rebuild +0 -1
- restart_trigger +0 -1
dvr.py
DELETED
|
@@ -1,452 +0,0 @@
|
|
| 1 |
-
"""
|
| 2 |
-
DVR / Timeshift cho VTV LIVE
|
| 3 |
-
============================
|
| 4 |
-
2 tính năng:
|
| 5 |
-
1) Timeshift "xem lại LIVE lùi tới 5 phút": ffmpeg ghi luồng m3u8 thành các
|
| 6 |
-
segment HLS 2s, playlist giữ ~5 phút gần nhất (omit_endlist) -> hls.js cho
|
| 7 |
-
phép tua lùi trong cửa sổ này.
|
| 8 |
-
2) Hẹn giờ record: đặt lịch (giờ bắt đầu + thời lượng), tới giờ ffmpeg ghi ra
|
| 9 |
-
file .mp4, ghi xong (tuỳ chọn) upload lên HF Dataset để giữ vĩnh viễn
|
| 10 |
-
(vì storage của Space là ephemeral, restart sẽ mất).
|
| 11 |
-
|
| 12 |
-
Yêu cầu: ffmpeg (Dockerfile của Space đã cài sẵn).
|
| 13 |
-
Tích hợp: trong app chính -> app.include_router(dvr_router)
|
| 14 |
-
và gọi resume_schedules() lúc khởi động.
|
| 15 |
-
"""
|
| 16 |
-
import os, re, json, time, shutil, threading, subprocess, signal
|
| 17 |
-
from datetime import datetime, timedelta, timezone
|
| 18 |
-
from fastapi import APIRouter, Query, Body
|
| 19 |
-
from fastapi.responses import JSONResponse, FileResponse, Response
|
| 20 |
-
|
| 21 |
-
VN_TZ = timezone(timedelta(hours=7))
|
| 22 |
-
router = APIRouter()
|
| 23 |
-
|
| 24 |
-
# Thư mục làm việc. Trên HF Space dùng /data nếu có Persistent Storage, nếu không /tmp.
|
| 25 |
-
BASE = "/data/dvr" if os.path.isdir("/data") and os.access("/data", os.W_OK) else "/tmp/dvr"
|
| 26 |
-
BUF_DIR = os.path.join(BASE, "buffer") # timeshift buffer (HLS segments)
|
| 27 |
-
REC_DIR = os.path.join(BASE, "recordings") # file mp4 đã ghi
|
| 28 |
-
os.makedirs(BUF_DIR, exist_ok=True)
|
| 29 |
-
os.makedirs(REC_DIR, exist_ok=True)
|
| 30 |
-
SCHED_FILE = os.path.join(BASE, "schedules.json")
|
| 31 |
-
|
| 32 |
-
FFMPEG = shutil.which("ffmpeg") or os.environ.get("FFMPEG", "/usr/bin/ffmpeg")
|
| 33 |
-
YTDP_PATH = shutil.which("yt-dlp") or os.environ.get("YTDLP", "/usr/bin/yt-dlp")
|
| 34 |
-
|
| 35 |
-
# Cửa sổ timeshift: 5 phút = 150 segment * 2s
|
| 36 |
-
SEG_SEC = 2
|
| 37 |
-
TIMESHIFT_WINDOW_SEC = 5 * 60
|
| 38 |
-
LIST_SIZE = TIMESHIFT_WINDOW_SEC // SEG_SEC # 150
|
| 39 |
-
|
| 40 |
-
# ---- lấy URL m3u8 thật của kênh (tái dùng backend vtv_api nếu có) ----
|
| 41 |
-
def _resolve_stream(channel_id):
|
| 42 |
-
"""Resolve stream URL với fallback chain. Trả về list các URLs thử."""
|
| 43 |
-
try:
|
| 44 |
-
from vtv_api import fetch_vtv_stream, XEMTV_US_ENDPOINTS, FPTPLAY_URLS, VTVGO_FAILOVER
|
| 45 |
-
# Primary URL từ API
|
| 46 |
-
primary = fetch_vtv_stream(channel_id)
|
| 47 |
-
candidates = []
|
| 48 |
-
if primary:
|
| 49 |
-
candidates.append(primary)
|
| 50 |
-
# Thêm fallback URLs trực tiếp
|
| 51 |
-
ch = channel_id.lower().strip()
|
| 52 |
-
if ch in FPTPLAY_URLS and FPTPLAY_URLS[ch] not in candidates:
|
| 53 |
-
candidates.append(FPTPLAY_URLS[ch])
|
| 54 |
-
if ch in VTVGO_FAILOVER and VTVGO_FAILOVER[ch] not in candidates:
|
| 55 |
-
candidates.append(VTVGO_FAILOVER[ch])
|
| 56 |
-
return candidates if candidates else None
|
| 57 |
-
except Exception:
|
| 58 |
-
return None
|
| 59 |
-
|
| 60 |
-
# =====================================================================
|
| 61 |
-
# 1) TIMESHIFT BUFFER
|
| 62 |
-
# =====================================================================
|
| 63 |
-
_buffers = {} # channel_id -> Popen
|
| 64 |
-
_buf_lock = threading.Lock()
|
| 65 |
-
|
| 66 |
-
def _try_ffmpeg_stream(url):
|
| 67 |
-
"""Test xem ffmpeg có mở được stream URL không (timeout 10s)."""
|
| 68 |
-
try:
|
| 69 |
-
cmd = [FFMPEG, "-hide_banner", "-loglevel", "error",
|
| 70 |
-
"-reconnect", "1", "-reconnect_streamed", "1", "-reconnect_delay_max", "5",
|
| 71 |
-
"-i", url, "-t", "1", "-c", "copy", "-f", "null", "-"]
|
| 72 |
-
r = subprocess.run(cmd, capture_output=True, timeout=15)
|
| 73 |
-
return r.returncode == 0
|
| 74 |
-
except Exception:
|
| 75 |
-
return False
|
| 76 |
-
|
| 77 |
-
def _start_buffer(channel_id, stream_urls):
|
| 78 |
-
"""Khởi động buffer, retry với nhiều sources nếu cần."""
|
| 79 |
-
if isinstance(stream_urls, str):
|
| 80 |
-
stream_urls = [stream_urls]
|
| 81 |
-
ch_dir = os.path.join(BUF_DIR, channel_id)
|
| 82 |
-
os.makedirs(ch_dir, exist_ok=True)
|
| 83 |
-
playlist = os.path.join(ch_dir, "index.m3u8")
|
| 84 |
-
seg = os.path.join(ch_dir, "seg_%05d.ts")
|
| 85 |
-
|
| 86 |
-
# Xóa segments cũ
|
| 87 |
-
for f in os.listdir(ch_dir):
|
| 88 |
-
if f.endswith(".ts"):
|
| 89 |
-
try: os.remove(os.path.join(ch_dir, f))
|
| 90 |
-
except: pass
|
| 91 |
-
|
| 92 |
-
# Chọn source đầu tiên hoạt động
|
| 93 |
-
working_url = None
|
| 94 |
-
for url in stream_urls:
|
| 95 |
-
print(f"[DVR] Testing source for {channel_id}: {url[:80]}...", flush=True)
|
| 96 |
-
if _try_ffmpeg_stream(url):
|
| 97 |
-
working_url = url
|
| 98 |
-
print(f"[DVR] Source OK: {url[:80]}", flush=True)
|
| 99 |
-
break
|
| 100 |
-
else:
|
| 101 |
-
print(f"[DVR] Source FAIL: {url[:80]}", flush=True)
|
| 102 |
-
|
| 103 |
-
if not working_url:
|
| 104 |
-
print(f"[DVR] ERROR: No working source for {channel_id}", flush=True)
|
| 105 |
-
return None
|
| 106 |
-
|
| 107 |
-
cmd = [
|
| 108 |
-
FFMPEG, "-hide_banner", "-loglevel", "warning", "-y",
|
| 109 |
-
"-reconnect", "1", "-reconnect_streamed", "1", "-reconnect_delay_max", "10",
|
| 110 |
-
"-i", working_url,
|
| 111 |
-
"-c", "copy", "-copyts", "-start_at_zero",
|
| 112 |
-
"-f", "hls",
|
| 113 |
-
"-hls_time", str(SEG_SEC),
|
| 114 |
-
"-hls_list_size", str(LIST_SIZE),
|
| 115 |
-
"-hls_flags", "delete_segments+omit_endlist+independent_segments+program_date_time",
|
| 116 |
-
"-hls_segment_type", "mpegts",
|
| 117 |
-
"-hls_segment_filename", seg,
|
| 118 |
-
playlist,
|
| 119 |
-
]
|
| 120 |
-
p = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT)
|
| 121 |
-
# Log output trong background
|
| 122 |
-
def _log_output():
|
| 123 |
-
for line in iter(p.stdout.readline, b''):
|
| 124 |
-
line_s = line.decode("utf-8", errors="ignore").strip()
|
| 125 |
-
if line_s:
|
| 126 |
-
print(f"[DVR:{channel_id}] {line_s}", flush=True)
|
| 127 |
-
threading.Thread(target=_log_output, daemon=True).start()
|
| 128 |
-
return p
|
| 129 |
-
|
| 130 |
-
@router.get("/api/dvr/buffer/start/{channel_id}")
|
| 131 |
-
def buffer_start(channel_id: str):
|
| 132 |
-
channel_id = channel_id.lower().strip()
|
| 133 |
-
with _buf_lock:
|
| 134 |
-
existing = _buffers.get(channel_id)
|
| 135 |
-
if existing and existing.get("popen") and existing["popen"].poll() is None:
|
| 136 |
-
return JSONResponse({"status": "already_running", "channel": channel_id,
|
| 137 |
-
"playlist": f"/api/dvr/buffer/playlist/{channel_id}"})
|
| 138 |
-
urls = _resolve_stream(channel_id)
|
| 139 |
-
if not urls:
|
| 140 |
-
return JSONResponse({"error": "stream not found"}, status_code=404)
|
| 141 |
-
# Chạy buffer trong background thread
|
| 142 |
-
def _run():
|
| 143 |
-
p = _start_buffer(channel_id, urls)
|
| 144 |
-
with _buf_lock:
|
| 145 |
-
if p:
|
| 146 |
-
_buffers[channel_id] = {"popen": p, "urls": urls, "error": None}
|
| 147 |
-
else:
|
| 148 |
-
_buffers[channel_id] = {"popen": None, "urls": urls,
|
| 149 |
-
"error": "all_sources_failed",
|
| 150 |
-
"hint": "FPTPlay/VTVGo CDN block từ server IP. Cần proxy VN."}
|
| 151 |
-
threading.Thread(target=_run, daemon=True).start()
|
| 152 |
-
return JSONResponse({"status": "starting", "channel": channel_id,
|
| 153 |
-
"window_sec": TIMESHIFT_WINDOW_SEC,
|
| 154 |
-
"playlist": f"/api/dvr/buffer/playlist/{channel_id}",
|
| 155 |
-
"note": "Buffer sẵn sau 5-10s. Nhấn DVR 5' lại để xem.",
|
| 156 |
-
"warning": "Nếu buffer fail, FPTPlay/VTVGo CDN có thể block từ server IP."})
|
| 157 |
-
|
| 158 |
-
@router.get("/api/dvr/buffer/stop/{channel_id}")
|
| 159 |
-
def buffer_stop(channel_id: str):
|
| 160 |
-
channel_id = channel_id.lower().strip()
|
| 161 |
-
with _buf_lock:
|
| 162 |
-
entry = _buffers.pop(channel_id, None)
|
| 163 |
-
if entry:
|
| 164 |
-
p = entry.get("popen")
|
| 165 |
-
if p and p.poll() is None:
|
| 166 |
-
p.terminate()
|
| 167 |
-
try:
|
| 168 |
-
p.wait(timeout=5)
|
| 169 |
-
except:
|
| 170 |
-
p.kill()
|
| 171 |
-
return JSONResponse({"status": "stopped", "channel": channel_id})
|
| 172 |
-
return JSONResponse({"status": "not_running", "channel": channel_id})
|
| 173 |
-
|
| 174 |
-
@router.get("/api/dvr/buffer/status/{channel_id}")
|
| 175 |
-
def buffer_status(channel_id: str):
|
| 176 |
-
channel_id = channel_id.lower().strip()
|
| 177 |
-
with _buf_lock:
|
| 178 |
-
entry = _buffers.get(channel_id)
|
| 179 |
-
if not entry:
|
| 180 |
-
# Check if playlist exists from previous session
|
| 181 |
-
pl = os.path.join(BUF_DIR, channel_id, "index.m3u8")
|
| 182 |
-
if os.path.exists(pl):
|
| 183 |
-
return JSONResponse({"active": False, "channel": channel_id, "has_playlist": True,
|
| 184 |
-
"playlist": f"/api/dvr/buffer/playlist/{channel_id}"})
|
| 185 |
-
return JSONResponse({"active": False, "channel": channel_id, "has_playlist": False})
|
| 186 |
-
p = entry.get("popen")
|
| 187 |
-
err = entry.get("error")
|
| 188 |
-
pl = os.path.join(BUF_DIR, channel_id, "index.m3u8")
|
| 189 |
-
segs = []
|
| 190 |
-
if os.path.exists(pl):
|
| 191 |
-
segs = sorted([f for f in os.listdir(os.path.join(BUF_DIR, channel_id)) if f.endswith(".ts")])
|
| 192 |
-
return JSONResponse({
|
| 193 |
-
"active": p is not None and p.poll() is None,
|
| 194 |
-
"channel": channel_id,
|
| 195 |
-
"error": err,
|
| 196 |
-
"hint": entry.get("hint"),
|
| 197 |
-
"segments": len(segs),
|
| 198 |
-
"has_playlist": os.path.exists(pl),
|
| 199 |
-
"playlist": f"/api/dvr/buffer/playlist/{channel_id}" if os.path.exists(pl) else None,
|
| 200 |
-
})
|
| 201 |
-
|
| 202 |
-
@router.get("/api/dvr/buffer/playlist/{channel_id}")
|
| 203 |
-
def buffer_playlist(channel_id: str):
|
| 204 |
-
channel_id = channel_id.lower().strip()
|
| 205 |
-
pl = os.path.join(BUF_DIR, channel_id, "index.m3u8")
|
| 206 |
-
if not os.path.exists(pl):
|
| 207 |
-
return JSONResponse({"error": "buffer not ready"}, status_code=404)
|
| 208 |
-
txt = open(pl, "r", encoding="utf-8", errors="ignore").read()
|
| 209 |
-
# rewrite segment names -> absolute URLs for HLS.js
|
| 210 |
-
out = []
|
| 211 |
-
for line in txt.split("\n"):
|
| 212 |
-
s = line.strip()
|
| 213 |
-
if s and not s.startswith("#") and s.endswith(".ts"):
|
| 214 |
-
# Dùng absolute URL để HLS.js resolve đúng
|
| 215 |
-
out.append(f"/api/dvr/buffer/seg/{channel_id}/{s}")
|
| 216 |
-
else:
|
| 217 |
-
out.append(line)
|
| 218 |
-
# Thêm cache-bust để tránh browser cache cũ
|
| 219 |
-
return Response("\n".join(out), media_type="application/vnd.apple.mpegurl",
|
| 220 |
-
headers={
|
| 221 |
-
"Access-Control-Allow-Origin": "*",
|
| 222 |
-
"Cache-Control": "no-cache, no-store, must-revalidate",
|
| 223 |
-
"Pragma": "no-cache",
|
| 224 |
-
"Expires": "0",
|
| 225 |
-
})
|
| 226 |
-
|
| 227 |
-
@router.get("/api/dvr/buffer/seg/{channel_id}/{seg}")
|
| 228 |
-
def buffer_seg(channel_id: str, seg: str):
|
| 229 |
-
if not re.match(r"^seg_\d+\.ts$", seg):
|
| 230 |
-
return Response(status_code=400)
|
| 231 |
-
path = os.path.join(BUF_DIR, channel_id.lower().strip(), seg)
|
| 232 |
-
if not os.path.exists(path):
|
| 233 |
-
return Response(status_code=404)
|
| 234 |
-
return FileResponse(path, media_type="video/mp2t",
|
| 235 |
-
headers={"Access-Control-Allow-Origin": "*"})
|
| 236 |
-
|
| 237 |
-
# =====================================================================
|
| 238 |
-
# 2) HẸN GIỜ RECORD
|
| 239 |
-
# =====================================================================
|
| 240 |
-
_sched_lock = threading.Lock()
|
| 241 |
-
_active_recs = {} # job_id -> Popen
|
| 242 |
-
|
| 243 |
-
def _load_sched():
|
| 244 |
-
if os.path.exists(SCHED_FILE):
|
| 245 |
-
try:
|
| 246 |
-
return json.load(open(SCHED_FILE, encoding="utf-8"))
|
| 247 |
-
except Exception:
|
| 248 |
-
return []
|
| 249 |
-
return []
|
| 250 |
-
|
| 251 |
-
def _save_sched(items):
|
| 252 |
-
json.dump(items, open(SCHED_FILE, "w", encoding="utf-8"), ensure_ascii=False, indent=2)
|
| 253 |
-
|
| 254 |
-
def _do_record(job):
|
| 255 |
-
"""Ghi 1 job: ưu tiên ffmpeg, fallback sang yt-dlp nếu ffmpeg fail."""
|
| 256 |
-
ch = job["channel"]
|
| 257 |
-
dur = int(job["duration_sec"])
|
| 258 |
-
fname = f"{ch}_{job['id']}.mp4"
|
| 259 |
-
out = os.path.join(REC_DIR, fname)
|
| 260 |
-
# Ưu tiên buffer timeshift nếu đang chạy => bắt được cả 5' "trước đó LIVE"
|
| 261 |
-
buf_pl = os.path.join(BUF_DIR, ch, "index.m3u8")
|
| 262 |
-
if os.path.exists(buf_pl):
|
| 263 |
-
src = buf_pl
|
| 264 |
-
else:
|
| 265 |
-
urls = _resolve_stream(ch)
|
| 266 |
-
src = urls[0] if urls else None
|
| 267 |
-
if not src:
|
| 268 |
-
job["status"] = "error"; job["error"] = "no source"; _save_all(); return
|
| 269 |
-
|
| 270 |
-
# Thử ffmpeg trước
|
| 271 |
-
cmd = [FFMPEG, "-hide_banner", "-loglevel", "warning", "-y"]
|
| 272 |
-
if os.path.exists(buf_pl):
|
| 273 |
-
cmd += ["-live_start_index", "0"]
|
| 274 |
-
cmd += ["-i", src, "-t", str(dur), "-c", "copy", out]
|
| 275 |
-
print(f"[REC] Recording {ch} for {dur}s from {src[:60]}...", flush=True)
|
| 276 |
-
job["status"] = "recording"; _save_all()
|
| 277 |
-
|
| 278 |
-
try:
|
| 279 |
-
p = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.STDOUT)
|
| 280 |
-
_active_recs[job["id"]] = p
|
| 281 |
-
for line in iter(p.stdout.readline, b''):
|
| 282 |
-
line_s = line.decode("utf-8", errors="ignore").strip()
|
| 283 |
-
if line_s:
|
| 284 |
-
print(f"[REC:{ch}] {line_s}", flush=True)
|
| 285 |
-
p.wait()
|
| 286 |
-
_active_recs.pop(job["id"], None)
|
| 287 |
-
if os.path.exists(out) and os.path.getsize(out) > 0:
|
| 288 |
-
job["status"] = "done"
|
| 289 |
-
job["file"] = fname
|
| 290 |
-
print(f"[REC] Done: {out} ({os.path.getsize(out)} bytes)", flush=True)
|
| 291 |
-
else:
|
| 292 |
-
raise Exception("output empty")
|
| 293 |
-
except Exception as e:
|
| 294 |
-
# Fallback: thử yt-dlp nếu có
|
| 295 |
-
if YTDP_PATH:
|
| 296 |
-
print(f"[REC] ffmpeg failed ({e}), trying yt-dlp...", flush=True)
|
| 297 |
-
out_ytdlp = os.path.join(REC_DIR, f"{ch}_{job['id']}_ytdlp.mp4")
|
| 298 |
-
try:
|
| 299 |
-
cmd2 = [YTDP_PATH, "-o", out_ytdlp, "--no-part"]
|
| 300 |
-
# Nếu là HLS live stream, dùng --download-sections
|
| 301 |
-
if src.endswith(".m3u8") or ".m3u8" in src:
|
| 302 |
-
# yt-dlp cần ffmpeg để merge, nhưng vẫn thử
|
| 303 |
-
cmd2 += ["--downloader", "ffmpeg", "--downloader-args", "ffmpeg:-t " + str(dur)]
|
| 304 |
-
else:
|
| 305 |
-
cmd2 += ["--downloader-args", "-t " + str(dur)]
|
| 306 |
-
cmd2.append(src)
|
| 307 |
-
p2 = subprocess.Popen(cmd2, stdout=subprocess.PIPE, stderr=subprocess.STDOUT)
|
| 308 |
-
_active_recs[job["id"]] = p2
|
| 309 |
-
for line in iter(p2.stdout.readline, b''):
|
| 310 |
-
line_s = line.decode("utf-8", errors="ignore").strip()
|
| 311 |
-
if line_s:
|
| 312 |
-
print(f"[REC:{ch} yt-dlp] {line_s}", flush=True)
|
| 313 |
-
p2.wait()
|
| 314 |
-
_active_recs.pop(job["id"], None)
|
| 315 |
-
if os.path.exists(out_ytdlp) and os.path.getsize(out_ytdlp) > 0:
|
| 316 |
-
os.rename(out_ytdlp, out)
|
| 317 |
-
job["status"] = "done"
|
| 318 |
-
job["file"] = fname
|
| 319 |
-
print(f"[REC] yt-dlp Done: {out} ({os.path.getsize(out)} bytes)", flush=True)
|
| 320 |
-
else:
|
| 321 |
-
raise Exception("yt-dlp output empty")
|
| 322 |
-
except Exception as e2:
|
| 323 |
-
_active_recs.pop(job["id"], None)
|
| 324 |
-
job["status"] = "error"
|
| 325 |
-
job["error"] = f"ffmpeg: {e}, yt-dlp: {e2}"
|
| 326 |
-
else:
|
| 327 |
-
_active_recs.pop(job["id"], None)
|
| 328 |
-
job["status"] = "error"
|
| 329 |
-
job["error"] = str(e)
|
| 330 |
-
|
| 331 |
-
# (tuỳ chọn) upload lên HF Dataset để giữ vĩnh viễn
|
| 332 |
-
if job["status"] == "done" and job.get("upload_repo"):
|
| 333 |
-
try:
|
| 334 |
-
from huggingface_hub import HfApi
|
| 335 |
-
HfApi(token=os.environ.get("HF_TOKEN")).upload_file(
|
| 336 |
-
path_or_fileobj=out, path_in_repo=f"recordings/{fname}",
|
| 337 |
-
repo_id=job["upload_repo"], repo_type="dataset")
|
| 338 |
-
job["uploaded"] = True
|
| 339 |
-
except Exception as e:
|
| 340 |
-
job["upload_error"] = str(e)
|
| 341 |
-
_save_all()
|
| 342 |
-
|
| 343 |
-
def _schedule_job(job):
|
| 344 |
-
now = datetime.now(VN_TZ)
|
| 345 |
-
start = datetime.fromisoformat(job["start_iso"])
|
| 346 |
-
if start.tzinfo is None:
|
| 347 |
-
start = start.replace(tzinfo=VN_TZ)
|
| 348 |
-
delay = (start - now).total_seconds()
|
| 349 |
-
if delay < 0:
|
| 350 |
-
delay = 0
|
| 351 |
-
def _runner():
|
| 352 |
-
# auto bật buffer trước giờ ghi để có sẵn 5' lùi (nếu bật pre-buffer)
|
| 353 |
-
if job.get("prebuffer"):
|
| 354 |
-
try: buffer_start(job["channel"])
|
| 355 |
-
except Exception: pass
|
| 356 |
-
_do_record(job)
|
| 357 |
-
t = threading.Timer(delay, _runner)
|
| 358 |
-
t.daemon = True
|
| 359 |
-
t.start()
|
| 360 |
-
job["status"] = "scheduled"
|
| 361 |
-
|
| 362 |
-
def _save_all():
|
| 363 |
-
with _sched_lock:
|
| 364 |
-
_save_sched(_SCHEDULES)
|
| 365 |
-
|
| 366 |
-
_SCHEDULES = _load_sched()
|
| 367 |
-
|
| 368 |
-
@router.post("/api/dvr/record/schedule")
|
| 369 |
-
def record_schedule(payload: dict = Body(...)):
|
| 370 |
-
"""body: {channel, start_iso (giờ VN), end_iso (giờ VN), duration_sec, prebuffer?, upload_repo?}"""
|
| 371 |
-
ch = str(payload.get("channel", "")).lower().strip()
|
| 372 |
-
if not ch:
|
| 373 |
-
return JSONResponse({"error": "missing channel"}, status_code=400)
|
| 374 |
-
start_iso = payload.get("start_iso")
|
| 375 |
-
end_iso = payload.get("end_iso")
|
| 376 |
-
if not start_iso:
|
| 377 |
-
return JSONResponse({"error": "missing start_iso"}, status_code=400)
|
| 378 |
-
# Parse và validate timezone
|
| 379 |
-
try:
|
| 380 |
-
start_dt = datetime.fromisoformat(start_iso)
|
| 381 |
-
if start_dt.tzinfo is None:
|
| 382 |
-
start_dt = start_dt.replace(tzinfo=VN_TZ)
|
| 383 |
-
except Exception:
|
| 384 |
-
return JSONResponse({"error": "invalid start_iso"}, status_code=400)
|
| 385 |
-
# Tính duration từ end_iso nếu có, nếu không dùng duration_sec
|
| 386 |
-
if end_iso:
|
| 387 |
-
try:
|
| 388 |
-
end_dt = datetime.fromisoformat(end_iso)
|
| 389 |
-
if end_dt.tzinfo is None:
|
| 390 |
-
end_dt = end_dt.replace(tzinfo=VN_TZ)
|
| 391 |
-
duration_from_end = int((end_dt - start_dt).total_seconds())
|
| 392 |
-
duration = max(5, min(3600, duration_from_end))
|
| 393 |
-
except Exception:
|
| 394 |
-
duration = int(payload.get("duration_sec", 600))
|
| 395 |
-
else:
|
| 396 |
-
duration = int(payload.get("duration_sec", 600))
|
| 397 |
-
duration = max(5, min(3600, duration))
|
| 398 |
-
job = {
|
| 399 |
-
"id": str(int(time.time() * 1000)),
|
| 400 |
-
"channel": ch,
|
| 401 |
-
"start_iso": start_dt.isoformat(),
|
| 402 |
-
"end_iso": end_iso or (start_dt + timedelta(seconds=duration)).isoformat(),
|
| 403 |
-
"duration_sec": duration,
|
| 404 |
-
"prebuffer": bool(payload.get("prebuffer", True)),
|
| 405 |
-
"upload_repo": payload.get("upload_repo"),
|
| 406 |
-
"status": "pending", "file": None,
|
| 407 |
-
}
|
| 408 |
-
with _sched_lock:
|
| 409 |
-
_SCHEDULES.append(job); _save_sched(_SCHEDULES)
|
| 410 |
-
_schedule_job(job); _save_all()
|
| 411 |
-
return JSONResponse({"status": "ok", "job": job})
|
| 412 |
-
|
| 413 |
-
@router.get("/api/dvr/record/list")
|
| 414 |
-
def record_list():
|
| 415 |
-
return JSONResponse({"jobs": _SCHEDULES})
|
| 416 |
-
|
| 417 |
-
@router.get("/api/dvr/record/file/{job_id}")
|
| 418 |
-
def record_file(job_id: str):
|
| 419 |
-
for j in _SCHEDULES:
|
| 420 |
-
if j["id"] == job_id and j.get("file"):
|
| 421 |
-
path = os.path.join(REC_DIR, j["file"])
|
| 422 |
-
if os.path.exists(path):
|
| 423 |
-
return FileResponse(path, media_type="video/mp4", filename=j["file"])
|
| 424 |
-
return JSONResponse({"error": "not found"}, status_code=404)
|
| 425 |
-
|
| 426 |
-
@router.get("/api/dvr/record/cancel/{job_id}")
|
| 427 |
-
def record_cancel(job_id: str):
|
| 428 |
-
p = _active_recs.get(job_id)
|
| 429 |
-
if p and p.poll() is None:
|
| 430 |
-
p.terminate()
|
| 431 |
-
for j in _SCHEDULES:
|
| 432 |
-
if j["id"] == job_id:
|
| 433 |
-
j["status"] = "cancelled"
|
| 434 |
-
_save_all()
|
| 435 |
-
return JSONResponse({"status": "cancelled", "job_id": job_id})
|
| 436 |
-
|
| 437 |
-
def resume_schedules():
|
| 438 |
-
"""Gọi lúc app khởi động: đặt lại lịch cho job còn ở tương lai."""
|
| 439 |
-
now = datetime.now(VN_TZ)
|
| 440 |
-
for j in _SCHEDULES:
|
| 441 |
-
if j.get("status") in ("pending", "scheduled"):
|
| 442 |
-
try:
|
| 443 |
-
start = datetime.fromisoformat(j["start_iso"])
|
| 444 |
-
if start.tzinfo is None:
|
| 445 |
-
start = start.replace(tzinfo=VN_TZ)
|
| 446 |
-
if start + timedelta(seconds=j["duration_sec"]) > now:
|
| 447 |
-
_schedule_job(j)
|
| 448 |
-
else:
|
| 449 |
-
j["status"] = "missed"
|
| 450 |
-
except Exception:
|
| 451 |
-
pass
|
| 452 |
-
_save_all()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
gitattributes
DELETED
|
@@ -1,35 +0,0 @@
|
|
| 1 |
-
*.7z filter=lfs diff=lfs merge=lfs -text
|
| 2 |
-
*.arrow filter=lfs diff=lfs merge=lfs -text
|
| 3 |
-
*.bin filter=lfs diff=lfs merge=lfs -text
|
| 4 |
-
*.bz2 filter=lfs diff=lfs merge=lfs -text
|
| 5 |
-
*.ckpt filter=lfs diff=lfs merge=lfs -text
|
| 6 |
-
*.ftz filter=lfs diff=lfs merge=lfs -text
|
| 7 |
-
*.gz filter=lfs diff=lfs merge=lfs -text
|
| 8 |
-
*.h5 filter=lfs diff=lfs merge=lfs -text
|
| 9 |
-
*.joblib filter=lfs diff=lfs merge=lfs -text
|
| 10 |
-
*.lfs.* filter=lfs diff=lfs merge=lfs -text
|
| 11 |
-
*.mlmodel filter=lfs diff=lfs merge=lfs -text
|
| 12 |
-
*.model filter=lfs diff=lfs merge=lfs -text
|
| 13 |
-
*.msgpack filter=lfs diff=lfs merge=lfs -text
|
| 14 |
-
*.npy filter=lfs diff=lfs merge=lfs -text
|
| 15 |
-
*.npz filter=lfs diff=lfs merge=lfs -text
|
| 16 |
-
*.onnx filter=lfs diff=lfs merge=lfs -text
|
| 17 |
-
*.ot filter=lfs diff=lfs merge=lfs -text
|
| 18 |
-
*.parquet filter=lfs diff=lfs merge=lfs -text
|
| 19 |
-
*.pb filter=lfs diff=lfs merge=lfs -text
|
| 20 |
-
*.pickle filter=lfs diff=lfs merge=lfs -text
|
| 21 |
-
*.pkl filter=lfs diff=lfs merge=lfs -text
|
| 22 |
-
*.pt filter=lfs diff=lfs merge=lfs -text
|
| 23 |
-
*.pth filter=lfs diff=lfs merge=lfs -text
|
| 24 |
-
*.rar filter=lfs diff=lfs merge=lfs -text
|
| 25 |
-
*.safetensors filter=lfs diff=lfs merge=lfs -text
|
| 26 |
-
saved_model/**/* filter=lfs diff=lfs merge=lfs -text
|
| 27 |
-
*.tar.* filter=lfs diff=lfs merge=lfs -text
|
| 28 |
-
*.tar filter=lfs diff=lfs merge=lfs -text
|
| 29 |
-
*.tflite filter=lfs diff=lfs merge=lfs -text
|
| 30 |
-
*.tgz filter=lfs diff=lfs merge=lfs -text
|
| 31 |
-
*.wasm filter=lfs diff=lfs merge=lfs -text
|
| 32 |
-
*.xz filter=lfs diff=lfs merge=lfs -text
|
| 33 |
-
*.zip filter=lfs diff=lfs merge=lfs -text
|
| 34 |
-
*.zst filter=lfs diff=lfs merge=lfs -text
|
| 35 |
-
*tfevents* filter=lfs diff=lfs merge=lfs -text
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
rebuild
DELETED
|
@@ -1 +0,0 @@
|
|
| 1 |
-
rebuilt at 2026-06-18T09:24:34.973630
|
|
|
|
|
|
restart_trigger
DELETED
|
@@ -1 +0,0 @@
|
|
| 1 |
-
# restart trigger
|
|
|
|
|
|