File size: 6,256 Bytes
5d782cc | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 | from __future__ import annotations
import hashlib
import json
import tempfile
import time
import unittest
from http.server import BaseHTTPRequestHandler, HTTPServer
from pathlib import Path
from threading import Thread
from renderer.automation.templates import resolve_render_template
from renderer.core.config import Settings
from renderer.core.models import TaskResult
from renderer.jobs.manager import JobManager
from renderer.platform.processor import _analysis_command, _parse_analysis
class N8NAutomationTests(unittest.TestCase):
def test_template_variables_preserve_json_types(self) -> None:
result = resolve_render_template(
{"scenes": "{{scenes}}", "title": "Campaign {{campaign.name}}", "enabled": True},
{"scenes": [{"start": 0, "duration": 4}], "campaign": {"name": "July"}},
{"enabled": False},
)
self.assertEqual(result["scenes"], [{"start": 0, "duration": 4}])
self.assertEqual(result["title"], "Campaign July")
self.assertFalse(result["enabled"])
def test_analysis_parsers_return_structured_events(self) -> None:
silence = _parse_analysis(
"silence_detect",
"silence_start: 1.25\nsilence_end: 2.75 | silence_duration: 1.50",
{},
)
black = _parse_analysis("black_detect", "black_start:0 black_end:1.2 black_duration:1.2", {})
loudness = _parse_analysis(
"loudness_analyze",
"I: -20.0 LUFS\nLRA: 4.0 LU\nPeak: -1.2 dBFS",
{},
)
self.assertEqual(silence["events"][0]["duration"], 1.5)
self.assertEqual(black["events"][0]["end"], 1.2)
self.assertEqual(loudness["integrated_lufs"], -20.0)
def test_hls_command_keeps_playlist_segment_prefix(self) -> None:
from renderer.platform.processor import PlatformProcessor
with tempfile.TemporaryDirectory() as directory:
settings = Settings(temp_dir=Path(directory) / "temp", exports_dir=Path(directory) / "exports")
processor = PlatformProcessor(settings)
command = processor._toolkit_command(
"hls", Path("input.mp4"), Path(directory) / "stream.zip", {}, Path(directory)
)
self.assertTrue(any(value.endswith("hls_segment_%05d.ts") for value in command))
self.assertTrue(any(value.endswith("hls_playlist.m3u8") for value in command))
def test_job_records_artifact_and_group_metadata(self) -> None:
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
settings = Settings(
temp_dir=root / "temp",
exports_dir=root / "exports",
jobs_dir=root / "jobs",
storage_dir=root / "storage",
metadata_cache=root / "temp" / "metadata.json",
whisper_model_dir=root / "models",
max_workers=1,
max_retries=1,
)
manager = JobManager(settings)
def handler(job_id: str, log) -> TaskResult:
output = settings.exports_dir / f"{job_id}.json"
output.write_text("automation", encoding="utf-8")
return TaskResult(output_path=output, metrics={"task": "test"})
job_id = manager.submit_task(handler, metadata={"group_id": "group_test", "variant_name": "shorts"})
deadline = time.time() + 5
while time.time() < deadline and manager.get(job_id).state in {"PENDING", "RUNNING"}:
time.sleep(0.02)
record = manager.get(job_id)
manager.executor.shutdown(wait=True)
self.assertEqual(record.state, "COMPLETED")
self.assertEqual(record.metadata["group_id"], "group_test")
artifact = record.metrics["artifact"]
self.assertEqual(artifact["sha256"], hashlib.sha256(b"automation").hexdigest())
def test_webhook_contains_event_and_signature(self) -> None:
received: list[dict] = []
class Handler(BaseHTTPRequestHandler):
def do_POST(self) -> None:
body = self.rfile.read(int(self.headers["Content-Length"]))
received.append({"body": json.loads(body), "signature": self.headers.get("X-Ava2lon-Signature")})
self.send_response(204)
self.end_headers()
def log_message(self, format: str, *args) -> None:
return
server = HTTPServer(("127.0.0.1", 0), Handler)
thread = Thread(target=server.serve_forever, daemon=True)
thread.start()
try:
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
settings = Settings(
temp_dir=root / "temp",
exports_dir=root / "exports",
jobs_dir=root / "jobs",
storage_dir=root / "storage",
metadata_cache=root / "temp" / "metadata.json",
whisper_model_dir=root / "models",
max_workers=1,
max_retries=1,
callback_max_retries=1,
)
manager = JobManager(settings)
def handler(job_id: str, log) -> TaskResult:
output = settings.exports_dir / f"{job_id}.txt"
output.write_text("done", encoding="utf-8")
return TaskResult(output_path=output)
callback = f"http://127.0.0.1:{server.server_port}/callback"
job_id = manager.submit_task(handler, callback_url=callback)
deadline = time.time() + 5
while time.time() < deadline and manager.get(job_id).callback_delivered_at is None:
time.sleep(0.02)
record = manager.get(job_id)
manager.executor.shutdown(wait=True)
finally:
server.shutdown()
server.server_close()
self.assertIsNotNone(record.callback_delivered_at)
self.assertEqual(received[0]["body"]["event"], "task.completed")
self.assertTrue(received[0]["signature"].startswith("sha256="))
if __name__ == "__main__":
unittest.main()
|