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
14 changes: 14 additions & 0 deletions pmoves/Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -255,3 +255,17 @@ web-geometry:
xdg-open http://localhost:8087/geometry/ 2>/dev/null || open http://localhost:8087/geometry/ ; \
fi
@$(MAKE) smoke-geometry

.PHONY: discord-smoke
discord-smoke:
@which jq >/dev/null 2>&1 || (echo "jq is required" && exit 1)
@echo "[Discord] Health" && curl -sf http://localhost:8092/healthz | jq -e '.ok==true' >/dev/null && echo OK || (echo FAIL && exit 1)
@echo "[Discord] Publish test" && curl -sS -X POST http://localhost:8092/publish -H 'content-type: application/json' -d '{"content":"PMOVES test ping"}' | jq -e '.ok==true' >/dev/null && echo OK || (echo FAIL && exit 1)

.PHONY: jellyfin-smoke
jellyfin-smoke:
@which jq >/dev/null 2>&1 || (echo "jq is required" && exit 1)
@echo "[Jellyfin] Health" && curl -sf http://localhost:8093/healthz | jq -e '.ok==true' >/dev/null && echo OK || (echo FAIL && exit 1)
@echo "[Jellyfin] Playback URL for latest video" && VID=$$(curl -s 'http://localhost:3000/videos?order=id.desc&select=video_id&limit=1' | jq -r '.[0].video_id // empty'); \
if [ -z "$$VID" ]; then echo "No videos found in Supabase; run yt-emit-smoke first"; exit 0; fi; \
curl -sS -X POST http://localhost:8093/jellyfin/playback-url -H 'content-type: application/json' -d "{\"video_id\":\"$$VID\",\"t\":5}" | jq -e '.ok==true and (.url|length)>0' >/dev/null && echo OK || (echo FAIL && exit 1)
49 changes: 49 additions & 0 deletions pmoves/services/jellyfin-bridge/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -75,3 +75,52 @@ def jellyfin_playback_url(body: Dict[str,Any] = Body(...)):
url = f"{JELLYFIN_URL}/web/index.html#!/details?id={item}&serverId=local&startTime={int(t)}"
return {"ok": True, "url": url}

@app.get("/jellyfin/search")
def jellyfin_search(query: str):
if not (JELLYFIN_URL and JELLYFIN_API_KEY and JELLYFIN_USER_ID):
raise HTTPException(412, 'JELLYFIN_URL, JELLYFIN_API_KEY, and JELLYFIN_USER_ID required')
try:
r = httpx.get(
f"{JELLYFIN_URL}/Users/{JELLYFIN_USER_ID}/Items",
params={"searchTerm": query, "IncludeItemTypes": "Movie,Video"},
headers={"X-Emby-Token": JELLYFIN_API_KEY}, timeout=8
)
r.raise_for_status()
j = r.json()
items = j.get('Items') or []
out = [{"Id": it.get('Id'), "Name": it.get('Name'), "ProductionYear": it.get('ProductionYear')} for it in items]
return {"ok": True, "items": out}
except Exception as e:
raise HTTPException(502, f"jellyfin search error: {e}")

@app.post("/jellyfin/map-by-title")
def jellyfin_map_by_title(body: Dict[str,Any] = Body(...)):
vid = body.get('video_id'); title = body.get('title')
if not vid: raise HTTPException(400, 'video_id required')
if not title:
rows = _supa_get('videos', {'video_id': vid})
if not rows:
raise HTTPException(404, 'video not found')
title = rows[0].get('title')
if not (JELLYFIN_URL and JELLYFIN_API_KEY and JELLYFIN_USER_ID):
raise HTTPException(412, 'JELLYFIN_URL, JELLYFIN_API_KEY, and JELLYFIN_USER_ID required')
# search and pick best match by simple case-insensitive inclusion
r = httpx.get(
f"{JELLYFIN_URL}/Users/{JELLYFIN_USER_ID}/Items",
params={"searchTerm": title, "IncludeItemTypes": "Movie,Video"},
headers={"X-Emby-Token": JELLYFIN_API_KEY}, timeout=8
)
r.raise_for_status()
items = (r.json().get('Items') or [])
best = None
tnorm = (title or '').lower()
for it in items:
name = (it.get('Name') or '').lower()
if tnorm in name or name in tnorm:
best = it; break
if not best and items:
best = items[0]
if not best:
raise HTTPException(404, 'no jellyfin items matched')
_supa_patch('videos', {'video_id': vid}, {"meta": {"jellyfin_item_id": best.get('Id')}})
return {"ok": True, "mapped": {"video_id": vid, "jellyfin_item_id": best.get('Id'), "name": best.get('Name')}}
9 changes: 9 additions & 0 deletions pmoves/services/pmoves-yt/yt.py
Original file line number Diff line number Diff line change
Expand Up @@ -394,6 +394,11 @@ def yt_summarize(body: Dict[str,Any] = Body(...)):
summary = _summarize_ollama(text, style)
# persist into videos + studio_board meta
supa_update('videos', {'video_id': vid}, {'meta': {'gemma': {'style': style, 'provider': provider, 'summary': summary}}})
# emit event for downstream (Discord/NATS)
try:
_publish_event('ingest.summary.ready.v1', {'video_id': vid, 'style': style, 'provider': provider, 'summary': summary[:500]})
except Exception:
pass
return {'ok': True, 'video_id': vid, 'provider': provider, 'style': style, 'summary': summary}

@app.post('/yt/chapters')
Expand All @@ -418,6 +423,10 @@ def yt_chapters(body: Dict[str,Any] = Body(...)):
# fallback: split lines
chapters = [{ 'title': line.strip('- ').strip(), 'blurb': '' } for line in raw.splitlines() if line.strip()][:10]
supa_update('videos', {'video_id': vid}, {'meta': {'gemma': {'chapters': chapters}}})
try:
_publish_event('ingest.chapters.ready.v1', {'video_id': vid, 'n': len(chapters), 'chapters': chapters[:6]})
except Exception:
pass
return {'ok': True, 'video_id': vid, 'chapters': chapters}

# -------------------- Segmentation → JSONL + CGP emit --------------------
Expand Down
72 changes: 60 additions & 12 deletions pmoves/services/publisher-discord/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@
app = FastAPI(title="Publisher-Discord", version="0.1.0")

DISCORD_WEBHOOK_URL = os.environ.get("DISCORD_WEBHOOK_URL", "")
DISCORD_USERNAME = os.environ.get("DISCORD_USERNAME", "PMOVES")
DISCORD_AVATAR_URL = os.environ.get("DISCORD_AVATAR_URL", "")
NATS_URL = os.environ.get("NATS_URL", "nats://nats:4222")
SUBJECTS = os.environ.get("DISCORD_SUBJECTS", "ingest.file.added.v1,ingest.transcript.ready.v1,ingest.summary.ready.v1,ingest.chapters.ready.v1").split(",")

Expand All @@ -16,23 +18,70 @@
async def healthz():
return {"ok": True, "webhook": bool(DISCORD_WEBHOOK_URL)}

async def _post_discord(content: str, embeds: Optional[list]=None):
async def _post_discord(content: Optional[str], embeds: Optional[list]=None, retries: int = 3):
if not DISCORD_WEBHOOK_URL:
return False
payload = {"content": content}
payload = {"username": DISCORD_USERNAME}
if DISCORD_AVATAR_URL:
payload["avatar_url"] = DISCORD_AVATAR_URL
if content:
payload["content"] = content
if embeds:
payload["embeds"] = embeds
async with httpx.AsyncClient(timeout=10) as client:
r = await client.post(DISCORD_WEBHOOK_URL, json=payload)
return r.status_code in (200, 204)
backoff = 1.0
async with httpx.AsyncClient(timeout=15) as client:
for _ in range(max(1, retries)):
r = await client.post(DISCORD_WEBHOOK_URL, json=payload)
if r.status_code in (200, 204):
return True
if r.status_code == 429:
try:
ra = float(r.headers.get("Retry-After", backoff))
except Exception:
ra = backoff
await asyncio.sleep(ra)
backoff = min(backoff * 2.0, 8.0)
continue
if 500 <= r.status_code < 600:
await asyncio.sleep(backoff)
backoff = min(backoff * 2.0, 8.0)
continue
return False
return False

def _format_event(name: str, payload: Dict[str, Any]) -> Dict[str, Any]:
title = f"{name}"
desc = json.dumps(payload)[:1800]
return {
"content": None,
"embeds": [{"title": title, "description": f"```json\n{desc}\n```"}]
}
name = name.strip()
emb = {"title": name, "fields": []}
thumb = None
if name == "ingest.file.added.v1":
title = payload.get("title") or payload.get("key")
emb["title"] = f"Ingest: {title}"
emb["fields"].append({"name":"Bucket", "value": str(payload.get("bucket")), "inline": True})
emb["fields"].append({"name":"Namespace", "value": str(payload.get("namespace")), "inline": True})
if payload.get("video_id"):
emb["fields"].append({"name":"Video ID", "value": str(payload.get("video_id")), "inline": True})
thumb = (payload.get("thumb") if isinstance(payload.get("thumb"), str) else None)
elif name == "ingest.transcript.ready.v1":
emb["title"] = f"Transcript ready: {payload.get('video_id')}"
emb["fields"].append({"name":"Language", "value": str(payload.get("language") or "auto"), "inline": True})
if payload.get("s3_uri"):
emb["fields"].append({"name":"Audio", "value": payload.get("s3_uri"), "inline": False})
elif name == "ingest.summary.ready.v1":
summ = payload.get("summary") or ""
emb["title"] = f"Summary: {payload.get('video_id')}"
emb["description"] = (summ[:1800] + ("…" if len(summ) > 1800 else ""))
elif name == "ingest.chapters.ready.v1":
ch = payload.get("chapters") or []
emb["title"] = f"Chapters: {payload.get('video_id')} ({len(ch)} items)"
if ch:
sample = "\n".join(f"• {c.get('title')}" for c in ch[:6])
emb["description"] = sample
else:
desc = json.dumps(payload)[:1800]
emb["description"] = f"```json\n{desc}\n```"
if thumb:
emb["thumbnail"] = {"url": thumb}
return {"content": None, "embeds": [emb]}

@app.on_event("startup")
async def startup():
Expand Down Expand Up @@ -70,4 +119,3 @@ async def publish_test(body: Dict[str, Any] = Body(...)):
if not ok:
raise HTTPException(502, "discord webhook failed")
return {"ok": True}