2 Commits
Author SHA1 Message Date
afk 92781e4291 Radio station track uploads. 2026-07-16 18:43:01 +02:00
afk f98f83d5df revert 512275a8b8
revert Switch to synchronous Python.

Signed-off-by: Abdulkadir Furkan Şanlı <me@abdulocra.cy>
2026-04-20 11:49:08 +02:00
3 changed files with 134 additions and 35 deletions
+4
View File
@@ -6,6 +6,10 @@ ENV PYTHONUNBUFFERED=1
WORKDIR /usr/src/app
RUN apt-get update && \
apt-get install -y --no-install-recommends curl ffmpeg yt-dlp && \
rm -rf /var/lib/apt/lists/*
COPY requirements.txt ./
RUN uv pip install --system --no-cache -r requirements.txt
+6
View File
@@ -12,3 +12,9 @@ MATRIX_PASSWORD = ""
YOUTUBE_PLAYLIST_TITLE = ""
# YouTube API client secret json path.
YOUTUBE_CLIENT_SECRETS_FILE = ""
# AzuraCast API key.
AZURACAST_API_KEY = ""
# AzuraCast API URL.
AZURACAST_URL = ""
# AzuraCast station ID.
AZURACAST_STATION_ID = ""
+124 -35
View File
@@ -2,8 +2,11 @@
"""ParkerBot"""
import argparse
import asyncio
import datetime
import glob
import html
import json
import os
import pickle
import re
@@ -14,7 +17,7 @@ from google.auth.transport.requests import Request
from google_auth_oauthlib.flow import InstalledAppFlow
from googleapiclient.discovery import build
from googleapiclient import errors
from nio import HttpClient, RoomMessageText, SyncResponse, UploadResponse
from nio import AsyncClient, RoomMessageText, SyncResponse, UploadResponse
DATA_DIR = os.getenv("DATA_DIR", "./")
DB_PATH = os.path.join(DATA_DIR, "parkerbot.sqlite3")
@@ -29,6 +32,10 @@ MATRIX_PASSWORD = os.getenv("MATRIX_PASSWORD")
YOUTUBE_CLIENT_SECRETS_FILE = os.getenv("YOUTUBE_CLIENT_SECRETS_FILE")
YOUTUBE_PLAYLIST_TITLE = os.getenv("YOUTUBE_PLAYLIST_TITLE")
AZURACAST_API_KEY = os.getenv("AZURACAST_API_KEY")
AZURACAST_URL = os.getenv("AZURACAST_URL")
AZURACAST_STATION_ID = os.getenv("AZURACAST_STATION_ID")
def connect_db():
"""Connect to DB and return connection and cursor."""
@@ -204,14 +211,14 @@ def get_video_info(youtube, video_id):
return False, "[Error fetching title]", ""
def send_intro_message(client, sender, room_id):
async def send_intro_message(client, sender, room_id):
"""Sends introduction message in reply to sender, in room with room_id."""
intro_message = (
f"Hi {sender}, I'm ParkerBot! I generate YouTube playlists from links "
"sent to this channel. You can find my source code here: "
"https://git.abdulocra.cy/abdulocracy/parkerbot"
)
client.room_send(
await client.room_send(
room_id=room_id,
message_type="m.room.message",
content={"msgtype": "m.text", "body": intro_message},
@@ -219,11 +226,11 @@ def send_intro_message(client, sender, room_id):
# TODO: Figure out how to properly send GIF, this is broken as shit.
with open("./parker.gif", "rb") as gif_file:
response = client.upload(gif_file, content_type="image/gif")
response = await client.upload(gif_file, content_type="image/gif")
if isinstance(response, UploadResponse):
print("Image was uploaded successfully to server. ")
gif_uri = response.content_uri
client.room_send(
await client.room_send(
room_id=room_id,
message_type="m.room.message",
content={
@@ -237,29 +244,29 @@ def send_intro_message(client, sender, room_id):
print(f"Failed to upload image. Failure response: {response}")
def send_playlist_of_week(client, sender, room_id, playlist_id):
async def send_playlist_of_week(client, sender, room_id, playlist_id):
"""Sends playlist of the week in reply to sender, in room with room_id."""
playlist_link = f"https://www.youtube.com/playlist?list={playlist_id}"
reply_msg = f"{sender}, here's the playlist of the week: {playlist_link}"
client.room_send(
await client.room_send(
room_id=room_id,
message_type="m.room.message",
content={"msgtype": "m.text", "body": reply_msg},
)
def send_playlist_of_all(client, sender, room_id, playlist_id):
async def send_playlist_of_all(client, sender, room_id, playlist_id):
"""Sends playlist of all time in reply to sender, in room with room_id."""
playlist_link = f"https://www.youtube.com/playlist?list={playlist_id}"
reply_msg = f"{sender}, here's the playlist of all time: {playlist_link}"
client.room_send(
await client.room_send(
room_id=room_id,
message_type="m.room.message",
content={"msgtype": "m.text", "body": reply_msg},
)
def message_callback(conn, cursor, youtube, client, room, event):
async def message_callback(conn, cursor, youtube, client, room, event):
"""Event handler for received messages."""
sender = event.sender
if sender != MATRIX_USER:
@@ -282,15 +289,15 @@ def message_callback(conn, cursor, youtube, client, room, event):
recent = abs(current_time - timestamp_sec) < datetime.timedelta(minutes=5)
if body == "!parkerbot" and recent:
send_intro_message(client, sender, room.room_id)
await send_intro_message(client, sender, room.room_id)
return
if body == "!week" and recent:
send_playlist_of_week(client, sender, room.room_id, playlist_id)
await send_playlist_of_week(client, sender, room.room_id, playlist_id)
return
if body == "!all" and recent:
send_playlist_of_all(client, sender, room.room_id, all_playlist_id)
await send_playlist_of_all(client, sender, room.room_id, all_playlist_id)
return
youtube_link_pattern = (
@@ -318,7 +325,7 @@ def message_callback(conn, cursor, youtube, client, room, event):
f"<em><span data-mx-color='#808080'>{escaped_text}</span></em>"
)
client.room_send(
await client.room_send(
room_id=room.room_id,
message_type="m.room.message",
content={
@@ -347,6 +354,8 @@ def message_callback(conn, cursor, youtube, client, room, event):
(playlist_id, message_id, video_id),
)
print(f"Added track to this week's playlist: {link}")
if recent:
asyncio.create_task(process_radio_track(link, video_id, title))
def in_playlist(cursor, video_id, playlist_id):
@@ -382,7 +391,7 @@ def record_message(conn, cursor, sender, link, timestamp):
return cursor.fetchone()[0]
def sync_callback(response):
async def sync_callback(response):
"""Saves Matrix sync token."""
with open(TOKEN_PATH, "w", encoding="utf-8") as f:
f.write(response.next_batch)
@@ -397,9 +406,9 @@ def load_sync_token():
return None
def get_client(conn, cursor, youtube):
async def get_client(conn, cursor, youtube):
"""Returns configured and logged in Matrix client."""
client = HttpClient(MATRIX_SERVER, MATRIX_USER)
client = AsyncClient(MATRIX_SERVER, MATRIX_USER)
client.add_event_callback(
lambda room, event: message_callback(
conn, cursor, youtube, client, room, event
@@ -407,23 +416,23 @@ def get_client(conn, cursor, youtube):
RoomMessageText,
)
client.add_response_callback(sync_callback, SyncResponse)
print(client.login(MATRIX_PASSWORD))
print(await client.login(MATRIX_PASSWORD))
return client
def backwards_sync(conn, cursor, youtube, client, room, start_token):
async def backwards_sync(conn, cursor, youtube, client, room, start_token):
"""Fetch and process historical messages from a given room."""
print("Starting to process channel log...")
from_token = start_token
room_id = room.room_id
while True:
# Fetch room messages
response = client.room_messages(room_id, from_token, direction="b")
response = await client.room_messages(room_id, from_token, direction="b")
# Process each message
for event in response.chunk:
if isinstance(event, RoomMessageText):
message_callback(conn, cursor, youtube, client, room, event)
await message_callback(conn, cursor, youtube, client, room, event)
# Break if there are no more messages to fetch
if not response.end or response.end == from_token:
@@ -433,30 +442,110 @@ def backwards_sync(conn, cursor, youtube, client, room, start_token):
from_token = response.end
def main():
async def process_radio_track(video_link, video_id, title):
"""Downloads native audio and pushes it to AzuraCast API in the background."""
if not AZURACAST_API_KEY:
return
# Sanitize title for the local filesystem
safe_title = "".join(c for c in title if c.isalnum() or c in " -_").strip()
base_filename = f"{safe_title}_{video_id}"
base_filepath = os.path.join(DATA_DIR, base_filename)
print(f"📻 Downloading audio for radio: {title}")
# 1. Download with yt-dlp (Keeping native AAC/Opus format)
dl_cmd = [
"yt-dlp",
"--extract-audio",
"--output",
f"{base_filepath}.%(ext)s",
video_link,
]
process = await asyncio.create_subprocess_exec(
*dl_cmd, stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.DEVNULL
)
await process.communicate()
# Check for the file using glob since the extension (.m4a, .webm, .opus) is dynamic
downloaded_files = glob.glob(f"{base_filepath}.*")
if not downloaded_files:
print(f"❌ Failed to download audio for {title}")
return
filepath = downloaded_files[0]
print(f"📻 Uploading to AzuraCast: {title}")
# 2. Upload to AzuraCast using curl
upload_cmd = [
"curl",
"-s",
"-X",
"POST",
f"{AZURACAST_URL}/api/station/{AZURACAST_STATION_ID}/files",
"-H",
f"X-API-Key: {AZURACAST_API_KEY}",
"-F",
f"file=@{filepath}",
]
curl_proc = await asyncio.create_subprocess_exec(
*upload_cmd, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE
)
stdout, _ = await curl_proc.communicate()
# 3. Extract unique_id and push to the live queue
try:
response = json.loads(stdout.decode())
unique_id = response.get("unique_id")
if unique_id:
req_cmd = [
"curl",
"-s",
"-X",
"POST",
f"{AZURACAST_URL}/api/station/{AZURACAST_STATION_ID}/request/{unique_id}",
"-H",
f"X-API-Key: {AZURACAST_API_KEY}",
]
req_proc = await asyncio.create_subprocess_exec(
*req_cmd,
stdout=asyncio.subprocess.DEVNULL,
stderr=asyncio.subprocess.DEVNULL,
)
await req_proc.communicate()
print(f"✅ Queued on radio: {title}")
else:
print(f"⚠️ Uploaded, but no unique_id returned. Response: {stdout.decode()}")
except json.JSONDecodeError:
print(f"❌ Failed to parse AzuraCast response: {stdout.decode()}")
# Cleanup the local file so your container doesn't bloat
if os.path.exists(filepath):
os.remove(filepath)
async def main():
"""Get DB and Matrix client ready, and start syncing."""
args = parse_arguments()
conn, cursor = connect_db()
define_tables(conn, cursor)
youtube = get_authenticated_service()
client = get_client(conn, cursor, youtube)
client = await get_client(conn, cursor, youtube)
sync_token = load_sync_token()
# This is incredibly dumb and most probably will exceed your YouTube API quota.
if args.backwards_sync:
init_sync = client.sync(30000)
room = client.room_resolve_alias(MATRIX_ROOM)
backwards_sync(conn, cursor, youtube, client, room, init_sync.next_batch)
init_sync = await client.sync(30000)
room = await client.room_resolve_alias(MATRIX_ROOM)
await backwards_sync(conn, cursor, youtube, client, room, init_sync.next_batch)
print("Started syncing...")
while True:
try:
client.sync(timeout=30000, full_state=True, since=sync_token)
sync_token = load_sync_token()
except Exception as e:
print(f"Sync error: {e}")
time.sleep(5)
await client.sync_forever(30000, full_state=True, since=sync_token)
if __name__ == "__main__":
main()
asyncio.run(main())