Add /force: host oversized videos as temporary 24h links

When a posted video is too long (>1h, probed up front) or a normal download
fails as too-big/timeout, the bot offers /force. Replying /force downloads the
full video off-loop (separate slot, no size cap, --max-filesize 10G), uploads it
via rsync/SSH to a web docroot, and posts a tokenized link auto-deleted after 24h
by an hourly sweep (also evicts oldest past 10 GB; remote-clock TTL; only ever
removes its own token dirs).

Hosting target is env-configured (FORCE_REMOTE/FORCE_REMOTE_DIR/FORCE_BASE_URL/
FORCE_SSH_PORT/FORCE_SSH_KEY) so published code carries no infra details; the
feature self-disables when unset (probe is skipped too, so normal posting is
unchanged). Staged on disk (not tmpfs), self-heals orphaned staging dirs, and
re-arms the offer on transient failure. See signal-bot.service.example.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
2026-06-15 20:42:37 -04:00
parent fd0da1a59e
commit 12b1f2f03d
2 changed files with 343 additions and 12 deletions
+327 -5
View File
@@ -5,9 +5,14 @@ import json
import logging import logging
import os import os
import re import re
import secrets
import shlex
import shutil
import subprocess import subprocess
import tempfile import tempfile
import time import time
from datetime import datetime
from urllib.parse import quote
from signalbot import Command, Context, SignalBot from signalbot import Command, Context, SignalBot
from signalbot.command import regex_triggered, triggered from signalbot.command import regex_triggered, triggered
@@ -33,6 +38,29 @@ COOKIES = os.path.join(os.path.dirname(os.path.abspath(__file__)), "cookies.txt"
STATE_DB = os.path.join(os.path.dirname(os.path.abspath(__file__)), "bot-state.db") STATE_DB = os.path.join(os.path.dirname(os.path.abspath(__file__)), "bot-state.db")
ADMIN_NUMBERS = {n.strip() for n in os.environ.get("BOT_ADMINS", "").split(",") if n.strip()} ADMIN_NUMBERS = {n.strip() for n in os.environ.get("BOT_ADMINS", "").split(",") if n.strip()}
# --- /force: host an oversized video as a temporary download link ---
FORCE_DURATION_THRESHOLD = 3600 # offer /force for videos longer than this (1 hr)
FORCE_OFFER_TTL = 600 # a /force offer stays valid for 10 min
FORCE_TIMEOUT = 1800 # max seconds for a forced full download (30 min)
FORCE_MAX_FILESIZE = "10G" # yt-dlp aborts a single download larger than this
RSYNC_TIMEOUT = 1800 # max seconds to upload to the web box
FORCE_STAGE_ROOT = "/var/tmp/signal-bot-force" # disk-backed staging (/tmp is tmpfs/RAM)
MEDIA_TTL = 24 * 3600 # delete hosted files after 24 h
MEDIA_MAX_BYTES = 10 * 1024 * 1024 * 1024 # evict oldest once the hosted folder exceeds 10 GB
TOKEN_RE = re.compile(r"^[0-9a-f]{24}$") # cleanup only ever touches dirs we created
# /force hosting target — configured via env so the published code carries no
# infra details. The systemd unit supplies the real values; unset = feature off.
FORCE_REMOTE = os.environ.get("FORCE_REMOTE", "") # ssh target, e.g. "user@host"
FORCE_REMOTE_DIR = os.environ.get("FORCE_REMOTE_DIR", "") # web docroot on that host
FORCE_BASE_URL = os.environ.get("FORCE_BASE_URL", "").rstrip("/") # public URL for that dir
FORCE_SSH_PORT = os.environ.get("FORCE_SSH_PORT", "22")
FORCE_SSH_KEY = os.environ.get("FORCE_SSH_KEY", "") # ssh identity file (optional)
FORCE_ENABLED = bool(FORCE_REMOTE and FORCE_REMOTE_DIR and FORCE_BASE_URL)
_SSH_OPTS = ["-p", FORCE_SSH_PORT, "-o", "BatchMode=yes", "-o", "StrictHostKeyChecking=accept-new",
*(["-i", FORCE_SSH_KEY] if FORCE_SSH_KEY else [])]
_RSYNC_RSH = "ssh " + " ".join(shlex.quote(o) for o in _SSH_OPTS)
logging.basicConfig( logging.basicConfig(
level=logging.INFO, level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(message)s", format="%(asctime)s [%(levelname)s] %(message)s",
@@ -67,6 +95,19 @@ _storage = None
_RECENT_URLS_KEY = "recent_urls" _RECENT_URLS_KEY = "recent_urls"
_RECENT_SEP = "\x1f" # joins (group, url) into one JSON-safe storage key _RECENT_SEP = "\x1f" # joins (group, url) into one JSON-safe storage key
# (group_id, sender) -> {"url","title","duration","time"}: a pending /force offer.
pending_force = {}
# Separate 1-at-a-time slot for forced downloads so a long one can't starve clips.
_force_semaphore = None
def _get_force_semaphore() -> asyncio.Semaphore:
global _force_semaphore
if _force_semaphore is None:
_force_semaphore = asyncio.Semaphore(1)
return _force_semaphore
def _get_job_semaphore() -> asyncio.Semaphore: def _get_job_semaphore() -> asyncio.Semaphore:
global _job_semaphore global _job_semaphore
@@ -154,6 +195,24 @@ def _mark_url_handled(group_id, url) -> None:
_persist_recent_urls() _persist_recent_urls()
def _sweep_pending_force() -> None:
now = time.monotonic()
for k in [k for k, v in pending_force.items() if now - v["time"] > FORCE_OFFER_TTL]:
del pending_force[k]
def _set_pending_force(group_id, sender, url, title, duration) -> None:
_sweep_pending_force()
pending_force[(group_id, sender)] = {
"url": url, "title": title, "duration": duration, "time": time.monotonic(),
}
def _pop_pending_force(group_id, sender):
_sweep_pending_force()
return pending_force.pop((group_id, sender), None)
async def _safe_reply(c: Context, text: str) -> None: async def _safe_reply(c: Context, text: str) -> None:
try: try:
await c.reply(text) await c.reply(text)
@@ -346,17 +405,53 @@ class VideoCommand(Command):
key = (group, url) key = (group, url)
_inflight.add(key) _inflight.add(key)
try: try:
# Clip path: a bounded window, just grab it (no oversize concern).
if clip is not None: if clip is not None:
log.info("Clipping %s to window %d-%ds", url, clip[0], clip[1]) log.info("Clipping %s to window %d-%ds", url, clip[0], clip[1])
async with _get_job_semaphore(): async with _get_job_semaphore():
b64, err, silent = await asyncio.to_thread(_produce_video, url, clip) b64, err, silent = await asyncio.to_thread(_produce_video, url, clip)
await self._deliver(c, group, sender, url, b64, err, silent,
title=None, duration=None, allow_force=False)
return
# Full-video path. Only when /force hosting is enabled do we probe
# duration first (to fast-path an over-long video to a /force offer);
# otherwise behave like before and just download.
async with _get_job_semaphore():
if FORCE_ENABLED:
duration, title = await asyncio.to_thread(_probe_metadata, url)
too_long = duration is not None and duration > FORCE_DURATION_THRESHOLD
else:
duration, title, too_long = None, None, False
if too_long:
b64, err, silent = None, "", False
else:
b64, err, silent = await asyncio.to_thread(_produce_video, url, clip)
if too_long:
_mark_url_handled(group, url)
await _too_big_response(c, group, sender, url, title, duration)
return
await self._deliver(c, group, sender, url, b64, err, silent,
title=title, duration=duration, allow_force=True)
except Exception as e: # noqa: BLE001
log.exception("Unexpected error handling %s: %s", url, e)
await _safe_reply(c, "Something went wrong handling that video.")
finally:
_inflight.discard(key)
async def _deliver(self, c: Context, group, sender, url, b64, err, silent,
*, title, duration, allow_force) -> None:
if silent: if silent:
# A genuine no-media determination: record it so an edit replay # A genuine no-media determination: record it so an edit replay
# doesn't pointlessly re-run yt-dlp. # doesn't pointlessly re-run yt-dlp.
_mark_url_handled(group, url) _mark_url_handled(group, url)
return return
if err: if err:
if allow_force and _is_too_big_error(err):
_mark_url_handled(group, url)
await _too_big_response(c, group, sender, url, title, duration)
else:
# Hard error: do NOT mark handled, so a corrective edit can retry. # Hard error: do NOT mark handled, so a corrective edit can retry.
await _safe_reply(c, f"Couldn't grab that video: {err}") await _safe_reply(c, f"Couldn't grab that video: {err}")
return return
@@ -364,11 +459,6 @@ class VideoCommand(Command):
_mark_url_handled(group, url) _mark_url_handled(group, url)
_set_video(group, sender, b64) _set_video(group, sender, b64)
await _safe_send_video(c, b64) await _safe_send_video(c, b64)
except Exception as e: # noqa: BLE001
log.exception("Unexpected error handling %s: %s", url, e)
await _safe_reply(c, "Something went wrong handling that video.")
finally:
_inflight.discard(key)
def _run_ytdlp(url: str, outpath: str, tmpdir: str, def _run_ytdlp(url: str, outpath: str, tmpdir: str,
@@ -481,6 +571,189 @@ def _produce_video(url: str, clip: tuple[int, int] | None) -> tuple[str | None,
return base64.b64encode(f.read()).decode("utf-8"), None, False return base64.b64encode(f.read()).decode("utf-8"), None, False
def _probe_metadata(url: str) -> tuple[int | None, str | None]:
"""Fetch (duration_seconds, title) without downloading. Blocking — use to_thread."""
cmd = [
YTDLP, "--no-playlist", "--quiet", "--no-warnings",
"--js-runtimes", "node",
"--print", "%(duration)s\n%(title)s",
*(["--cookies", COOKIES] if os.path.exists(COOKIES) else []),
"--", url,
]
try:
r = subprocess.run(cmd, capture_output=True, text=True, timeout=90)
except subprocess.TimeoutExpired:
return None, None
if r.returncode != 0:
return None, None
lines = r.stdout.splitlines()
duration = None
if lines:
try:
duration = int(float(lines[0]))
except ValueError:
duration = None
title = lines[1].strip() if len(lines) > 1 and lines[1].strip() not in ("", "NA") else None
return duration, title
def _fmt_duration(seconds: int | None) -> str:
if not seconds:
return ""
h, rem = divmod(int(seconds), 3600)
m, s = divmod(rem, 60)
return f"{h}h{m:02d}m" if h else f"{m}m{s:02d}s"
def _too_big_offer(title: str | None, duration: int | None) -> str:
name = f' "{title}"' if title else ""
dur = f" ({_fmt_duration(duration)})" if duration else ""
return (
f"That video{name}{dur} is too big to post here. "
f"Reply `/force` within 10 min and I'll download it and give you a temporary "
f"link (auto-deleted in 24h)."
)
# Bot-generated oversize/slow sentinels only — NOT yt-dlp's transient network
# "read operation timed out", which should stay a retryable hard error.
_TOO_BIG_PATTERNS = ("too large", "after re-encode", "yt-dlp timed out after", "ffmpeg timed out")
def _is_too_big_error(err: str) -> bool:
e = (err or "").lower()
return any(p in e for p in _TOO_BIG_PATTERNS)
async def _too_big_response(c: Context, group, sender, url, title, duration) -> None:
"""Reply to an oversized video: offer /force if hosting is set up, else just note it."""
if FORCE_ENABLED:
_set_pending_force(group, sender, url, title, duration)
await _safe_reply(c, _too_big_offer(title, duration))
else:
name = f' "{title}"' if title else ""
dur = f" ({_fmt_duration(duration)})" if duration else ""
await _safe_reply(c, f"That video{name}{dur} is too big to post here.")
def _force_download_and_publish(url: str, token: str) -> tuple[str | None, int | None, str | None]:
"""Blocking: download the full video, upload it to the web box, return
(public_url, size_mb, error). Stages on disk (/tmp is RAM) and cleans up."""
os.makedirs(FORCE_STAGE_ROOT, exist_ok=True)
# Self-heal: drop staging dirs orphaned by a hard crash (only one force at a time).
for old in os.listdir(FORCE_STAGE_ROOT):
if old.startswith("force-"):
shutil.rmtree(os.path.join(FORCE_STAGE_ROOT, old), ignore_errors=True)
stage = tempfile.mkdtemp(dir=FORCE_STAGE_ROOT, prefix="force-")
try:
tokdir = os.path.join(stage, token)
os.makedirs(tokdir)
cmd = [
YTDLP, "--no-playlist", "--no-warnings", "--js-runtimes", "node",
"-f", "best[ext=mp4]/best", "--merge-output-format", "mp4",
"--max-filesize", FORCE_MAX_FILESIZE,
*(["--cookies", COOKIES] if os.path.exists(COOKIES) else []),
"-o", os.path.join(tokdir, "video.%(ext)s"),
"--", url,
]
try:
r = subprocess.run(cmd, capture_output=True, text=True, timeout=FORCE_TIMEOUT)
except subprocess.TimeoutExpired:
return None, None, f"timed out after {FORCE_TIMEOUT // 60} min"
files = [f for f in os.listdir(tokdir) if os.path.isfile(os.path.join(tokdir, f))]
if r.returncode != 0 or not files:
# yt-dlp prints the max-filesize abort to stdout (and exits 0), so check both.
blob = (r.stdout + r.stderr).lower()
if "larger than" in blob or "max-filesize" in blob:
return None, None, f"video is larger than the {FORCE_MAX_FILESIZE}B hosting limit"
return None, None, _summarize_ytdlp_error(r.stderr)
# Pick the actual video (largest file), not an arbitrary listdir entry.
fname = max(files, key=lambda f: os.path.getsize(os.path.join(tokdir, f)))
size_mb = os.path.getsize(os.path.join(tokdir, fname)) // (1024 * 1024)
up = subprocess.run(
["rsync", "-a", "--chmod=F644,D755", "-e", _RSYNC_RSH,
tokdir, f"{FORCE_REMOTE}:{FORCE_REMOTE_DIR}/"],
capture_output=True, text=True, timeout=RSYNC_TIMEOUT,
)
if up.returncode != 0:
return None, None, f"upload failed ({up.stderr.strip()[-160:]})"
return f"{FORCE_BASE_URL}/{token}/{quote(fname)}", size_mb, None
except Exception as e: # noqa: BLE001
log.exception("forced download failed for %s", url)
return None, None, f"unexpected error ({e})"
finally:
shutil.rmtree(stage, ignore_errors=True)
def _cleanup_media() -> None:
"""Delete hosted videos older than MEDIA_TTL and evict the oldest past
MEDIA_MAX_BYTES. Only ever touches token dirs we created. Blocking."""
ssh_base = ["ssh", *_SSH_OPTS, FORCE_REMOTE]
# Emit "age_seconds<TAB>bytes<TAB>name" using the REMOTE clock, so the TTL
# decision doesn't depend on this host's clock matching the web box's.
list_script = (
f"cd {shlex.quote(FORCE_REMOTE_DIR)} 2>/dev/null || exit 0; now=$(date +%s); "
f'for d in */; do [ -d "$d" ] || continue; n="${{d%/}}"; '
f'printf "%s\\t%s\\t%s\\n" "$((now - $(stat -c %Y "$d")))" '
f'"$(du -sb "$d" 2>/dev/null | cut -f1)" "$n"; done'
)
try:
r = subprocess.run([*ssh_base, list_script], capture_output=True, text=True, timeout=120)
except subprocess.TimeoutExpired:
log.warning("media cleanup: listing timed out")
return
if r.returncode != 0:
log.warning("media cleanup: listing failed (%s)", r.stderr.strip()[-160:])
return
entries = [] # (age_seconds, bytes, name)
for line in r.stdout.splitlines():
parts = line.split("\t")
if len(parts) != 3 or not TOKEN_RE.match(parts[2]):
continue
try:
age = int(parts[0])
except ValueError:
continue # no usable age -> can't make a TTL decision, skip this round
# A du failure leaves the size field empty; treat as 0 so TTL can still fire.
nbytes = int(parts[1]) if parts[1].isdigit() else 0
entries.append((age, nbytes, parts[2]))
doomed = {name for age, _b, name in entries if age > MEDIA_TTL}
fresh = sorted(e for e in entries if e[2] not in doomed) # ascending age = newest first
total = 0
for _age, nbytes, name in fresh:
total += nbytes
if total > MEDIA_MAX_BYTES:
doomed.add(name)
doomed = [n for n in doomed if TOKEN_RE.match(n)]
if not doomed:
return
paths = " ".join(shlex.quote(f"{FORCE_REMOTE_DIR}/{n}") for n in doomed)
try:
rm = subprocess.run([*ssh_base, f"rm -rf -- {paths}"],
capture_output=True, text=True, timeout=120)
except subprocess.TimeoutExpired:
log.warning("media cleanup: rm timed out")
return
if rm.returncode == 0:
log.info("media cleanup: removed %d hosted video(s)", len(doomed))
else:
log.warning("media cleanup: rm failed (%s)", rm.stderr.strip()[-160:])
async def _cleanup_media_job() -> None:
try:
await asyncio.to_thread(_cleanup_media)
except Exception as e: # noqa: BLE001
log.warning("media cleanup job error: %s", e)
def _reencode(input_file: str, tmpdir: str) -> tuple[str | None, str]: def _reencode(input_file: str, tmpdir: str) -> tuple[str | None, str]:
"""Re-encode video with ffmpeg to fit under MAX_FILE_SIZE. """Re-encode video with ffmpeg to fit under MAX_FILE_SIZE.
@@ -743,6 +1016,43 @@ class ReverseCommand(Command):
await _safe_send_video(c, out) await _safe_send_video(c, out)
class ForceCommand(Command):
"""`/force` downloads a video the bot refused as too big and hosts it as a
temporary link (auto-deleted in 24h)."""
async def handle(self, c: Context) -> None:
if not c.message.is_group():
return
if (c.message.text or "").strip().lower() != "/force":
return
if not FORCE_ENABLED:
await _safe_reply(c, "Big-video hosting isn't set up, sorry.")
return
group = c.message.group
sender = _sender_number(c.message)
pend = _pop_pending_force(group, sender)
if not pend:
await _safe_reply(c, "Nothing to force (offers expire after 10 min). Share the video link again and I'll re-offer `/force` if it's too big.")
return
url = pend["url"]
title = pend.get("title")
duration = pend.get("duration")
nm = f' "{title}"' if title else ""
dur = f" ({_fmt_duration(duration)})" if duration else ""
await _safe_reply(c, f"Downloading{nm}{dur}… this may take a while.")
token = secrets.token_hex(12)
async with _get_force_semaphore():
public_url, size_mb, err = await asyncio.to_thread(_force_download_and_publish, url, token)
if err:
# Re-arm the offer so a transient failure (network/rsync) can be retried.
_set_pending_force(group, sender, url, title, duration)
await _safe_reply(c, f"Couldn't download that video: {err}. Reply `/force` to try again.")
return
await _safe_reply(c, f"{public_url}\n({size_mb} MB — will be deleted in 24 hours)")
HELP_TEXT = f"""🎬 Video bot — what I can do HELP_TEXT = f"""🎬 Video bot — what I can do
Share a video link (X/Twitter, Instagram, YouTube, TikTok) and I'll post the video back to the group. Share a video link (X/Twitter, Instagram, YouTube, TikTok) and I'll post the video back to the group.
@@ -760,6 +1070,8 @@ A YouTube link with a timestamp → I post a {CLIP_DURATION}s clip starting at t
/rev — reverse the last video. /rev — reverse the last video.
/force — if a video's too big to post (e.g. over an hour), I'll offer this; reply `/force` and I'll download it and post a temporary link (auto-deleted in 24h).
/help — show this message. /help — show this message.
(In a DM, admins can run /cookies to refresh Instagram login cookies.)""" (In a DM, admins can run /cookies to refresh Instagram login cookies.)"""
@@ -865,8 +1177,18 @@ def main():
bot.register(VideoCommand(), contacts=False, groups=True) bot.register(VideoCommand(), contacts=False, groups=True)
bot.register(ReverseCommand(), contacts=False, groups=True) bot.register(ReverseCommand(), contacts=False, groups=True)
bot.register(SpeedCommand(), contacts=False, groups=True) bot.register(SpeedCommand(), contacts=False, groups=True)
bot.register(ForceCommand(), contacts=False, groups=True)
bot.register(HelpCommand(), contacts=False, groups=True) bot.register(HelpCommand(), contacts=False, groups=True)
bot.register(CookiesCommand(), contacts=True, groups=False) bot.register(CookiesCommand(), contacts=True, groups=False)
# Sweep hosted /force videos hourly (delete >24h old, evict oldest past 10 GB),
# starting shortly after launch. Only when /force hosting is configured.
if FORCE_ENABLED:
bot.scheduler.add_job(_cleanup_media_job, "interval", hours=1, next_run_time=datetime.now())
log.info("/force hosting enabled -> %s", FORCE_BASE_URL)
else:
log.info("/force hosting disabled (FORCE_REMOTE/FORCE_REMOTE_DIR/FORCE_BASE_URL unset)")
log.info("Starting Signal video bot...") log.info("Starting Signal video bot...")
bot.start() bot.start()
+9
View File
@@ -12,6 +12,15 @@ Environment=SIGNAL_PHONE_NUMBER=+15551234567
Environment=SIGNAL_SERVICE=127.0.0.1:8080 Environment=SIGNAL_SERVICE=127.0.0.1:8080
# Optional: comma-separated numbers allowed to run /cookies (admin command). # Optional: comma-separated numbers allowed to run /cookies (admin command).
# Environment=BOT_ADMINS=+15551234567 # Environment=BOT_ADMINS=+15551234567
# Optional: enable /force, which hosts oversized videos as temporary links.
# The bot rsyncs the file (over SSH) into a web-served directory and posts the URL;
# an hourly job deletes anything older than 24h and evicts oldest past 10 GB.
# Leave unset to disable /force. All three of REMOTE/REMOTE_DIR/BASE_URL are required.
# Environment=FORCE_REMOTE=user@webhost # ssh target that can write the docroot
# Environment=FORCE_REMOTE_DIR=/var/www/site/media # docroot on that host
# Environment=FORCE_BASE_URL=https://site/media # public URL serving that docroot
# Environment=FORCE_SSH_PORT=22 # ssh port (default 22)
# Environment=FORCE_SSH_KEY=/home/YOUR_USER/.ssh/id_ed25519 # ssh identity (optional)
ExecStart=/home/YOUR_USER/signal-bot/venv/bin/python bot.py ExecStart=/home/YOUR_USER/signal-bot/venv/bin/python bot.py
# always (not on-failure): the bot's library swallows exceptions in its own run # always (not on-failure): the bot's library swallows exceptions in its own run
# loop and never exits non-zero even when wedged, so on-failure would never bounce it. # loop and never exits non-zero even when wedged, so on-failure would never bounce it.