| 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() |
|
|