mirror of
https://github.com/halfwaystudent/douyin-sparkflow.git
synced 2026-09-05 07:27:53 +08:00
fix: prevent duplicate scheduled sends
This commit is contained in:
@@ -61,11 +61,11 @@ def _douyin_network_mode():
|
||||
|
||||
|
||||
def douyin_network_modes():
|
||||
# Direct is the default; Mihomo is used only when explicitly selected.
|
||||
# Direct is the default; Mihomo is the fallback unless explicitly selected.
|
||||
mode = _douyin_network_mode()
|
||||
if mode == "mihomo":
|
||||
return ("mihomo",)
|
||||
return ("direct",)
|
||||
return ("direct", "mihomo")
|
||||
|
||||
|
||||
def _douyin_browser_proxy(network_mode=None):
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import re
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
@@ -5,6 +6,61 @@ from datetime import datetime
|
||||
from pathlib import Path
|
||||
|
||||
|
||||
LEGACY_DOCKER_COMMAND_MARKERS = (
|
||||
"docker ps --format",
|
||||
"docker exec",
|
||||
)
|
||||
|
||||
|
||||
def migrate_legacy_line(line):
|
||||
"""Replace old Docker-in-Docker task commands with a local task runner."""
|
||||
if not line.strip() or line.lstrip().startswith("#"):
|
||||
return line
|
||||
if not all(marker in line for marker in LEGACY_DOCKER_COMMAND_MARKERS):
|
||||
return line
|
||||
|
||||
parts = line.split(maxsplit=5)
|
||||
if len(parts) < 6:
|
||||
return line
|
||||
|
||||
schedule = " ".join(parts[:5])
|
||||
command = parts[5]
|
||||
redirect_match = re.search(r"(\s+>>\s+\S+(?:\s+2>&1)?)\s*$", command)
|
||||
redirect = redirect_match.group(1) if redirect_match else ""
|
||||
|
||||
if "SPARKFLOW_MANUAL_UNSENT_ONLY=1" in command:
|
||||
env_prefix = (
|
||||
"SPARKFLOW_TRIGGER_LABEL='unsent fallback' "
|
||||
"SPARKFLOW_MANUAL_RUN=1 "
|
||||
"SPARKFLOW_MANUAL_UNSENT_ONLY=1 "
|
||||
"PYTHONUNBUFFERED=1"
|
||||
)
|
||||
else:
|
||||
env_prefix = "SPARKFLOW_TRIGGER_LABEL='scheduled send'"
|
||||
|
||||
return f"{schedule} env {env_prefix} bash /app/scripts/run_scheduled_task.sh{redirect}"
|
||||
|
||||
|
||||
def migrate_legacy_crontab(path):
|
||||
"""Rewrite persisted legacy task lines once, preserving other cron entries."""
|
||||
if not path.exists():
|
||||
return False
|
||||
|
||||
raw = path.read_text(encoding="utf-8-sig", errors="replace")
|
||||
lines = raw.splitlines()
|
||||
migrated_lines = [migrate_legacy_line(line) for line in lines]
|
||||
if migrated_lines == lines:
|
||||
return False
|
||||
|
||||
updated = "\n".join(migrated_lines)
|
||||
if updated:
|
||||
updated += "\n"
|
||||
temporary = path.with_name(f".{path.name}.migration.tmp")
|
||||
temporary.write_text(updated, encoding="utf-8")
|
||||
temporary.replace(path)
|
||||
return True
|
||||
|
||||
|
||||
def expand_field(field, minimum, maximum, current):
|
||||
values = set()
|
||||
for part in str(field).split(","):
|
||||
@@ -68,6 +124,8 @@ def read_crontab(path):
|
||||
|
||||
def run_loop(crontab_path):
|
||||
crontab_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
if migrate_legacy_crontab(crontab_path):
|
||||
print(f"[cron_runner] migrated legacy task commands in {crontab_path}", flush=True)
|
||||
last_minute_key = None
|
||||
print(f"[cron_runner] watching {crontab_path}", flush=True)
|
||||
while True:
|
||||
|
||||
@@ -0,0 +1,21 @@
|
||||
#!/usr/bin/env bash
|
||||
|
||||
set -u
|
||||
|
||||
script_dir="$(CDPATH= cd -- "$(dirname -- "${BASH_SOURCE[0]}")" && pwd)"
|
||||
app_root="$(CDPATH= cd -- "$script_dir/.." && pwd)"
|
||||
trigger_label="${SPARKFLOW_TRIGGER_LABEL:-scheduled task}"
|
||||
|
||||
echo "[AUTO_TRIGGER] $(date -Iseconds) ${trigger_label} start"
|
||||
|
||||
cd "$app_root"
|
||||
cd_rc=$?
|
||||
if [ "$cd_rc" -ne 0 ]; then
|
||||
echo "[AUTO_TRIGGER] $(date -Iseconds) ${trigger_label} failed to enter app directory rc=${cd_rc}"
|
||||
exit "$cd_rc"
|
||||
fi
|
||||
|
||||
python main.py --doTask
|
||||
task_rc=$?
|
||||
echo "[AUTO_TRIGGER] $(date -Iseconds) ${trigger_label} exit rc=${task_rc}"
|
||||
exit "$task_rc"
|
||||
@@ -33,6 +33,18 @@ class DeploymentContractTests(unittest.TestCase):
|
||||
self.assertTrue((SOURCE_ROOT / "config.example.json").is_file())
|
||||
self.assertIn("config.json", (SOURCE_ROOT / ".gitignore").read_text(encoding="utf-8"))
|
||||
|
||||
def test_scheduler_uses_local_task_runner_and_migrates_legacy_commands(self):
|
||||
compose = (REPO_ROOT / "docker-compose.yml").read_text(encoding="utf-8")
|
||||
runner_path = SOURCE_ROOT / "scripts" / "run_scheduled_task.sh"
|
||||
runner = runner_path.read_text(encoding="utf-8")
|
||||
cron_runner = (SOURCE_ROOT / "scripts" / "cron_runner.py").read_text(encoding="utf-8")
|
||||
scheduler = compose.split(" scheduler:", 1)[1].split("\n task:", 1)[0]
|
||||
|
||||
self.assertTrue(runner_path.is_file())
|
||||
self.assertIn("python main.py --doTask", runner)
|
||||
self.assertIn("migrate_legacy_crontab", cron_runner)
|
||||
self.assertNotIn("/var/run/docker.sock:/var/run/docker.sock", scheduler)
|
||||
|
||||
def test_compose_runtime_mounts_follow_least_privilege(self):
|
||||
text = (REPO_ROOT / "docker-compose.yml").read_text(encoding="utf-8")
|
||||
web = text.split(" web:", 1)[1].split("\n login-desktop:", 1)[0]
|
||||
@@ -41,7 +53,7 @@ class DeploymentContractTests(unittest.TestCase):
|
||||
|
||||
self.assertIn("/var/run/docker.sock:/var/run/docker.sock", web)
|
||||
self.assertIn(".:/opt/douyin-sparkflow", web)
|
||||
self.assertIn("/var/run/docker.sock:/var/run/docker.sock", scheduler)
|
||||
self.assertNotIn("/var/run/docker.sock:/var/run/docker.sock", scheduler)
|
||||
self.assertNotIn(".:/opt/douyin-sparkflow", scheduler)
|
||||
self.assertNotIn("/var/run/docker.sock:/var/run/docker.sock", task)
|
||||
self.assertNotIn(".:/opt/douyin-sparkflow", task)
|
||||
@@ -67,6 +79,23 @@ class DeploymentContractTests(unittest.TestCase):
|
||||
self.assertEqual(len(lines), 1)
|
||||
self.assertTrue(lines[0].startswith("*/20 "))
|
||||
|
||||
def test_legacy_cron_command_migrates_without_touching_schedule(self):
|
||||
from scripts.cron_runner import migrate_legacy_line
|
||||
|
||||
legacy = (
|
||||
"20 18 * * * /bin/bash -lc 'docker ps --format \"{{.Names}}\" | "
|
||||
"grep douyin-web; docker exec douyin-web sh -lc \"cd /app && "
|
||||
"env SPARKFLOW_MANUAL_RUN=1 SPARKFLOW_MANUAL_UNSENT_ONLY=1 "
|
||||
"python main.py --doTask\"' >> /var/log/douyin-sparkflow.log 2>&1"
|
||||
)
|
||||
migrated = migrate_legacy_line(legacy)
|
||||
|
||||
self.assertTrue(migrated.startswith("20 18 * * * "))
|
||||
self.assertIn("run_scheduled_task.sh", migrated)
|
||||
self.assertIn("SPARKFLOW_MANUAL_UNSENT_ONLY=1", migrated)
|
||||
self.assertNotIn("docker ps", migrated)
|
||||
self.assertNotIn("docker exec", migrated)
|
||||
|
||||
def test_build_proxy_does_not_leak_into_runtime_and_runtime_proxy_is_explicit(self):
|
||||
dockerfile = (SOURCE_ROOT / "Dockerfile.server").read_text(encoding="utf-8")
|
||||
compose = (REPO_ROOT / "docker-compose.yml").read_text(encoding="utf-8")
|
||||
@@ -101,6 +130,9 @@ class DeploymentContractTests(unittest.TestCase):
|
||||
windows = (REPO_ROOT / "deploy" / "install-local.ps1").read_text(encoding="utf-8")
|
||||
self.assertIn("runtime_config_backup", server)
|
||||
self.assertIn("Restored runtime config.json", server)
|
||||
self.assertIn("remove_legacy_host_cron", server)
|
||||
self.assertIn("host-crontab-", server)
|
||||
self.assertIn("docker ps --format", server)
|
||||
self.assertNotIn("bash ./refresh_proxy.sh", windows)
|
||||
self.assertIn("Initialize-ProxyConfig", windows)
|
||||
|
||||
@@ -124,15 +156,15 @@ class DeploymentContractTests(unittest.TestCase):
|
||||
env_example = (REPO_ROOT / ".env.example").read_text(encoding="utf-8")
|
||||
server = (SOURCE_ROOT / "login_desktop_server.py").read_text(encoding="utf-8")
|
||||
login_block = compose.split(" login-desktop:", 1)[1].split(" scheduler:", 1)[0]
|
||||
self.assertIn("LOGIN_DESKTOP_PROXY_MODE: ${LOGIN_DESKTOP_PROXY_MODE:-direct}", login_block)
|
||||
self.assertIn("LOGIN_DESKTOP_PROXY_MODE: ${LOGIN_DESKTOP_PROXY_MODE:-auto}", login_block)
|
||||
self.assertIn("LOGIN_DESKTOP_PROXY: ${LOGIN_DESKTOP_PROXY:-http://proxy:7890}", login_block)
|
||||
self.assertNotIn("HTTP_PROXY: http://proxy:7890", login_block)
|
||||
self.assertIn("LOGIN_DESKTOP_PROXY_MODE=direct", env_example)
|
||||
self.assertIn("LOGIN_DESKTOP_PROXY_MODE=auto", env_example)
|
||||
self.assertIn('candidates.append(("direct", None))', server)
|
||||
self.assertIn('candidates.append(("proxy", LOGIN_PROXY_SERVER))', server)
|
||||
self.assertIn('"--no-proxy-server"', server)
|
||||
self.assertIn('"/preflight"', server)
|
||||
self.assertIn('LOGIN_DESKTOP_PROXY_MODE: ${LOGIN_DESKTOP_PROXY_MODE:-direct}', login_block)
|
||||
self.assertIn('LOGIN_DESKTOP_PROXY_MODE: ${LOGIN_DESKTOP_PROXY_MODE:-auto}', login_block)
|
||||
dashboard = (SOURCE_ROOT / "webui" / "templates" / "dashboard.html").read_text(encoding="utf-8")
|
||||
self.assertIn('name="douyin_network_mode"', dashboard)
|
||||
self.assertIn('name="douyin_proxy_url"', dashboard)
|
||||
@@ -176,4 +208,4 @@ class DeploymentContractTests(unittest.TestCase):
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
unittest.main()
|
||||
unittest.main()
|
||||
|
||||
@@ -305,7 +305,8 @@ class WebUiSafetyTests(unittest.TestCase):
|
||||
self.assertIn("*/20 10-17 * * *", text)
|
||||
self.assertIn("0 18 * * *", text)
|
||||
self.assertIn("20 18 * * *", text)
|
||||
self.assertIn("docker exec", text)
|
||||
self.assertIn("run_scheduled_task.sh", text)
|
||||
self.assertNotIn("docker exec", text)
|
||||
|
||||
def test_overview_api_requires_authentication_and_disables_cache(self):
|
||||
client = TestClient(app_module.app)
|
||||
|
||||
@@ -23,6 +23,7 @@ TASK_SCHEDULE_MARKERS = (
|
||||
"docker compose run --rm task",
|
||||
"docker compose run --rm douyin",
|
||||
"main.py --doTask",
|
||||
"run_scheduled_task.sh",
|
||||
)
|
||||
HOST_CRONTAB_PATH = Path("/host-spool-cron/root")
|
||||
WINDOWED_SCHEDULE_RE = re.compile(r"^(\d{2}):(\d{2})-(\d{2}):(\d{2})/(\d+)m$", re.IGNORECASE)
|
||||
@@ -199,20 +200,9 @@ def _compose_env_args(extra_env=None):
|
||||
|
||||
def build_scheduled_task_command(extra_env=None, trigger_label="scheduled send"):
|
||||
if running_in_container():
|
||||
task_command = _with_env_prefix("python main.py --doTask", extra_env)
|
||||
return (
|
||||
"/bin/bash -lc 'timestamp=$(date -Iseconds); "
|
||||
f"echo \"[AUTO_TRIGGER] $timestamp {trigger_label} start\"; "
|
||||
"container=$(docker ps --format \"{{.Names}}\" | "
|
||||
"grep -E \"^(douyin-web-hostfix|douyin-web)$\" | head -n 1); "
|
||||
"if [ -z \"$container\" ]; then "
|
||||
"echo \"[AUTO_TRIGGER] $timestamp no matching container found\"; "
|
||||
"exit 1; "
|
||||
"fi; "
|
||||
"echo \"[AUTO_TRIGGER] $timestamp container=$container\"; "
|
||||
"docker exec \"$container\" sh -lc "
|
||||
f"\"cd /app && {task_command}\"'"
|
||||
)
|
||||
task_env = dict(extra_env or {})
|
||||
task_env["SPARKFLOW_TRIGGER_LABEL"] = trigger_label
|
||||
return _with_env_prefix("bash /app/scripts/run_scheduled_task.sh", task_env)
|
||||
if compose_file_path():
|
||||
compose_root_quoted = shlex.quote(str(compose_root()))
|
||||
compose_env_args = _compose_env_args(extra_env)
|
||||
|
||||
Reference in New Issue
Block a user