Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
61 changes: 35 additions & 26 deletions app.py
Original file line number Diff line number Diff line change
Expand Up @@ -146,10 +146,22 @@
}

.badge-kami {
background: linear-gradient(135deg, rgba(245, 158, 11, 0.15), rgba(168, 85, 247, 0.15));
border-color: rgba(245, 158, 11, 0.45);
color: #FDE047;
box-shadow: 0 0 12px rgba(245, 158, 11, 0.25);
background: linear-gradient(135deg, rgba(245, 158, 11, 0.20), rgba(217, 119, 6, 0.10));
border-color: rgba(251, 191, 36, 0.58);
color: #FCD34D;
box-shadow: 0 0 14px rgba(245, 158, 11, 0.28), inset 0 1px 0 rgba(255, 255, 255, 0.08);
white-space: nowrap;
}

@media (max-width: 640px) {
#studio-navbar {
padding: 8px 10px !important;
}

.badge-kami {
padding: 4px 9px;
font-size: 0.75rem;
}
}

.badge-pro {
Expand Down Expand Up @@ -638,6 +650,22 @@
"story_script": 50000,
}

DEVELOPER_ATTRIBUTION = "Developer: Kamran Ashraf"

def studio_navbar_html(gemini_chip: str) -> str:
return f"""
<div style="display:flex; align-items:center; justify-content:space-between; width:100%; flex-wrap:wrap; gap:10px;">
<div style="display:flex; align-items:center; gap:12px;">
<span class="studio-logo">🎬 VideoStudio Pro</span>
<span class="badge-chip badge-4k">4K Ultra HD</span>
</div>
<div style="display:flex; align-items:center; gap:10px; flex-wrap:wrap;">
<span class="badge-chip badge-kami">{DEVELOPER_ATTRIBUTION}</span>
{gemini_chip}
</div>
</div>
"""

def restore_browser_draft(payload: str) -> tuple:
try:
data = json.loads(payload) if isinstance(payload, str) else {}
Expand Down Expand Up @@ -753,12 +781,6 @@ def duration_hint(sec) -> str:
return (f"<span style='color:#38BDF8;font-size:0.82rem;'>🎞️ {int(float(sec))}s 4K · "
f"{n} AI shots chained & crossfaded</span>")

def job_progress_percent(progress: float) -> int:
try:
return max(0, min(100, int(float(progress) * 100)))
except (TypeError, ValueError, OverflowError):
return 0

def queue_wait_label(job: JobStatus, jobs: List[JobStatus], now: Optional[float] = None) -> str:
if job.stage != Stage.QUEUED:
return ""
Expand Down Expand Up @@ -1044,7 +1066,6 @@ def generate_ai_video_live(
time.sleep(1.0)
continue

pct = job_progress_percent(status.progress)
elapsed = int(time.time() - start_t)
card_html = render_job_card(
job_id,
Expand All @@ -1069,7 +1090,7 @@ def generate_ai_video_live(
yield "❌ Job failed. See the shared status panel for technical details.", card_html, None, None, job_id
break

yield f"⏳ Rendering 4K ({pct}%)...", card_html, None, None, job_id
yield "⏳ Rendering 4K in progress...", card_html, None, None, job_id
time.sleep(1.0)

def generate_storyboard_video_live(
Expand Down Expand Up @@ -1105,7 +1126,6 @@ def generate_storyboard_video_live(
time.sleep(1.0)
continue

pct = job_progress_percent(status.progress)
elapsed = int(time.time() - start_t)
card_html = render_job_card(
job_id,
Expand All @@ -1129,7 +1149,7 @@ def generate_storyboard_video_live(
yield "❌ Storyboard job failed. See the shared status panel for technical details.", card_html, None, job_id
break

yield f"⏳ Building Storyboard ({pct}%)...", card_html, None, job_id
yield "⏳ Building Storyboard in progress...", card_html, None, job_id
time.sleep(1.0)

def cancel_active_job_btn(job_id: str) -> str:
Expand Down Expand Up @@ -1224,18 +1244,7 @@ def build_app() -> gr.Blocks:

# Slim High-End Studio Top Bar
with gr.Group(elem_id="studio-navbar"):
gr.HTML(f"""
<div style="display:flex; align-items:center; justify-content:space-between; width:100%; flex-wrap:wrap; gap:10px;">
<div style="display:flex; align-items:center; gap:12px;">
<span class="studio-logo">🎬 VideoStudio Pro</span>
<span class="badge-chip badge-4k">4K Ultra HD</span>
</div>
<div style="display:flex; align-items:center; gap:10px; flex-wrap:wrap;">
<span class="badge-chip badge-kami">👑 Kami</span>
{gemini_chip}
</div>
</div>
""")
gr.HTML(studio_navbar_html(gemini_chip))

job_progress_panel = gr.HTML(render_current_job_panel())
with gr.Row():
Expand Down
115 changes: 109 additions & 6 deletions tests/test_core.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,6 @@
import json
import os
import threading
import tempfile
import unittest
from pathlib import Path
Expand All @@ -9,11 +11,11 @@
DRAFT_DEFAULTS,
clear_browser_draft,
draft_field_save_js,
job_progress_percent,
readable_job_message,
render_job_card,
restore_browser_draft,
restore_selected_tab_js,
studio_navbar_html,
sync_restored_job_outputs,
)
from video_studio import (
Expand Down Expand Up @@ -162,6 +164,13 @@ def test_active_job_shows_truthful_indeterminate_progress_and_step(self):
self.assertNotIn("42%", card)
self.assertIn("is-indeterminate", card)

def test_shared_navbar_displays_gold_developer_attribution_once(self):
navbar = studio_navbar_html("<span class='badge-chip badge-pro'>Connected</span>")
self.assertEqual(navbar.count("Developer: Kamran Ashraf"), 1)
self.assertIn("badge-kami", navbar)
self.assertIn("#FCD34D", app.CSS)
self.assertIn("@media (max-width: 640px)", app.CSS)

def test_browser_draft_restore_validates_values_and_never_includes_credentials(self):
values = restore_browser_draft(
json.dumps(
Expand Down Expand Up @@ -269,11 +278,30 @@ def test_restored_job_sync_reconnects_result_without_reloading_same_video(self):
self.assertEqual(repeated[3], app.gr.skip())
self.assertEqual(repeated[8], str(video))

def test_invalid_progress_is_safe(self):
self.assertEqual(job_progress_percent(float("nan")), 0)
self.assertEqual(job_progress_percent("invalid"), 0)
self.assertEqual(job_progress_percent(-0.5), 0)
self.assertEqual(job_progress_percent(2), 100)
def test_generation_status_does_not_claim_unmeasured_percentages(self):
queue = app.STUDIO.queue
status = JobStatus(
job_id="active",
stage=Stage.GENERATING,
progress=0.42,
message="Generating AI clip 2/4...",
settings=JobSettings(job_id="active", prompt="A rendering video"),
)
with (
patch.object(queue, "submit", return_value="active"),
patch.object(queue, "get_status", return_value=status),
patch.object(queue, "all_jobs", return_value=[status]),
patch.object(app.time, "sleep"),
):
update = next(
app.generate_ai_video_live(
"A rendering video", "", "Cinematic", "4K", "16:9", 6,
"English", "Female", "", "Ambient", True, "Classic White",
"None", False, True, False, -1,
)
)
self.assertEqual(update[0], "⏳ Rendering 4K in progress...")
self.assertNotIn("%", update[0])
self.assertEqual(Stage.from_str("FAILED"), Stage.FAILED)

def test_view_only_queue_rejects_mutations(self):
Expand All @@ -297,6 +325,81 @@ def test_invalid_worker_lock_fails_closed(self):
self.assertFalse(queue._acquire_worker_lock())
self.assertEqual(queue._lock_path.read_text(encoding="utf-8"), "not-a-process-id")

def test_concurrent_queue_saves_serialize_shared_temp_file_writes(self):
with tempfile.TemporaryDirectory() as tmp:
queue = object.__new__(JobQueue)
queue.queue_file = Path(tmp) / "jobs.json"
queue._lock = threading.RLock()
queue._last_persist = 0.0
status = JobStatus(
job_id="saved",
message="first snapshot",
settings=JobSettings(job_id="saved"),
)
queue._jobs = {"saved": status}

temp_path = queue.queue_file.with_suffix(f".{os.getpid()}.tmp")
first_write_started = threading.Event()
second_call_started = threading.Event()
second_write_started = threading.Event()
release_first_write = threading.Event()
writer_lock = threading.Lock()
write_count = 0
write_errors = []
original_write_text = Path.write_text

def controlled_write_text(path, data, *args, **kwargs):
nonlocal write_count
if path == temp_path:
with writer_lock:
write_count += 1
current_write = write_count
if current_write == 1:
first_write_started.set()
if not release_first_write.wait(timeout=5):
raise TimeoutError("Timed out waiting to release the first queue write.")
else:
second_write_started.set()
return original_write_text(path, data, *args, **kwargs)

def save():
try:
queue._save(force=True)
except Exception as exc:
write_errors.append(exc)

def save_newer():
second_call_started.set()
try:
with queue._lock:
status.message = "second snapshot"
queue._save(force=True)
except Exception as exc:
write_errors.append(exc)

with patch.object(Path, "write_text", controlled_write_text):
first = threading.Thread(target=save)
first.start()
self.assertTrue(first_write_started.wait(timeout=2))

second = threading.Thread(target=save_newer)
second.start()
self.assertTrue(second_call_started.wait(timeout=2))
self.assertFalse(second_write_started.wait(timeout=0.2))

release_first_write.set()
first.join(timeout=2)
second.join(timeout=2)

self.assertFalse(first.is_alive())
self.assertFalse(second.is_alive())
self.assertEqual(write_errors, [])
self.assertTrue(second_write_started.is_set())
self.assertEqual(
json.loads(queue.queue_file.read_text(encoding="utf-8"))["saved"]["message"],
"second snapshot",
)


if __name__ == "__main__":
unittest.main()
23 changes: 12 additions & 11 deletions video_studio.py
Original file line number Diff line number Diff line change
Expand Up @@ -1693,7 +1693,8 @@ def __init__(self, studio: "VideoStudio"):
self.queue_file = DEFAULT_OUTPUT_DIR / "queue" / "jobs.json"
self.queue_file.parent.mkdir(parents=True, exist_ok=True)
self._lock = threading.RLock() # re-entrant: _save() is called
# from inside methods that already hold the lock
# from inside methods that already hold the lock. Queue writes stay
# under this lock too, so concurrent snapshots cannot share a temp file.
self._jobs: Dict[str, JobStatus] = {}
self._cancelled_jobs: set[str] = set()
self._active_thread: Optional[threading.Thread] = None
Expand Down Expand Up @@ -1848,16 +1849,16 @@ def _save(self, force: bool = False):
"not_before": js.not_before,
}
payload = json.dumps(data, indent=2)
tmp = self.queue_file.with_suffix(f".{os.getpid()}.tmp")
for attempt in range(5):
try:
tmp.write_text(payload, encoding="utf-8")
tmp.replace(self.queue_file) # atomic
return
except (PermissionError, OSError) as exc:
time.sleep(0.2 * (attempt + 1))
last = exc
log.warning("Queue persist failed after retries: %s", last)
tmp = self.queue_file.with_suffix(f".{os.getpid()}.tmp")
for attempt in range(5):
try:
tmp.write_text(payload, encoding="utf-8")
tmp.replace(self.queue_file) # atomic
return
except (PermissionError, OSError) as exc:
time.sleep(0.2 * (attempt + 1))
last = exc
log.warning("Queue persist failed after retries: %s", last)

def submit(self, settings: JobSettings) -> str:
if not self.is_worker:
Expand Down
Loading