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