UPDATE: Change to how streams are accessed, detected and published to the application
This commit is contained in:
@@ -1,4 +1,5 @@
|
|||||||
worker_processes 1;
|
worker_processes 1;
|
||||||
|
error_log /var/log/nginx/error.log warn;
|
||||||
|
|
||||||
events {
|
events {
|
||||||
worker_connections 1024;
|
worker_connections 1024;
|
||||||
@@ -15,8 +16,8 @@ rtmp {
|
|||||||
deny play all;
|
deny play all;
|
||||||
push rtmp://127.0.0.1:1935/hls-live;
|
push rtmp://127.0.0.1:1935/hls-live;
|
||||||
|
|
||||||
on_publish http://web_server:5000/publish_stream;
|
on_publish http://web_server:5000/init_stream; # if the stream is detected from OBS
|
||||||
on_publish_done http://web_server:5000/end_stream;
|
on_publish_done http://web_server:5000/end_stream; # if the stream is ended on OBS
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -24,7 +25,7 @@ rtmp {
|
|||||||
live on;
|
live on;
|
||||||
|
|
||||||
hls on;
|
hls on;
|
||||||
hls_path /stream_data/;
|
hls_path /stream_data/stream/;
|
||||||
|
|
||||||
allow publish 127.0.0.1;
|
allow publish 127.0.0.1;
|
||||||
deny publish all;
|
deny publish all;
|
||||||
@@ -78,7 +79,7 @@ http {
|
|||||||
|
|
||||||
# The MPEG-TS video chunks are stored in /tmp/hls
|
# The MPEG-TS video chunks are stored in /tmp/hls
|
||||||
location ~ ^/stream/(.+)/(.+\.ts)$ {
|
location ~ ^/stream/(.+)/(.+\.ts)$ {
|
||||||
alias /stream_data/$1/stream/$2;
|
alias /stream_data/stream/$1/$2;
|
||||||
|
|
||||||
# Let the MPEG-TS video chunks be cacheable
|
# Let the MPEG-TS video chunks be cacheable
|
||||||
expires max;
|
expires max;
|
||||||
@@ -86,35 +87,36 @@ http {
|
|||||||
|
|
||||||
# The M3U8 playlists location
|
# The M3U8 playlists location
|
||||||
location ~ ^/stream/(.+)/(.+\.m3u8)$ {
|
location ~ ^/stream/(.+)/(.+\.m3u8)$ {
|
||||||
alias /stream_data/$1/stream/$2;
|
alias /stream_data/stream/$1/$2;
|
||||||
|
|
||||||
# The M3U8 playlists should not be cacheable
|
# The M3U8 playlists should not be cacheable
|
||||||
expires -1d;
|
expires -1d;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#! Unused right now so the following are inaccurate locations
|
||||||
# The thumbnails location
|
# The thumbnails location
|
||||||
location ~ ^/stream/(.+)/thumbnails/(.+\.jpg)$ {
|
# location ~ ^/stream/(.+)/thumbnails/(.+\.jpg)$ {
|
||||||
alias /stream_data/$1/thumbnails/$2;
|
# alias /stream_data/$1/thumbnails/$2;
|
||||||
|
|
||||||
# The thumbnails should not be cacheable
|
# # The thumbnails should not be cacheable
|
||||||
expires -1d;
|
# expires -1d;
|
||||||
}
|
# }
|
||||||
|
|
||||||
# The vods location
|
# # The vods location
|
||||||
location ~ ^/stream/(.+)/vods/(.+\.mp4)$ {
|
# location ~ ^/stream/(.+)/vods/(.+\.mp4)$ {
|
||||||
alias /stream_data/$1/vods/$2;
|
# alias /stream_data/$1/vods/$2;
|
||||||
|
|
||||||
# The vods should not be cacheable
|
# # The vods should not be cacheable
|
||||||
expires -1d;
|
# expires -1d;
|
||||||
}
|
# }
|
||||||
|
|
||||||
location ~ ^/\?token=.*$ {
|
# location ~ ^/\?token=.*$ {
|
||||||
proxy_pass http://frontend:5173;
|
# proxy_pass http://frontend:5173;
|
||||||
proxy_http_version 1.1;
|
# proxy_http_version 1.1;
|
||||||
proxy_set_header Upgrade $http_upgrade;
|
# proxy_set_header Upgrade $http_upgrade;
|
||||||
proxy_set_header Connection "upgrade";
|
# proxy_set_header Connection "upgrade";
|
||||||
proxy_set_header Host $host;
|
# proxy_set_header Host $host;
|
||||||
}
|
# }
|
||||||
|
|
||||||
location / {
|
location / {
|
||||||
proxy_pass http://frontend:5173;
|
proxy_pass http://frontend:5173;
|
||||||
|
|||||||
@@ -72,7 +72,7 @@ def stream_data(streamer_id) -> dict:
|
|||||||
Returns a streamer's current stream data
|
Returns a streamer's current stream data
|
||||||
"""
|
"""
|
||||||
data = get_current_stream_data(streamer_id)
|
data = get_current_stream_data(streamer_id)
|
||||||
|
|
||||||
if session.get('user_id') == streamer_id:
|
if session.get('user_id') == streamer_id:
|
||||||
with Database() as db:
|
with Database() as db:
|
||||||
stream_key = db.fetchone(
|
stream_key = db.fetchone(
|
||||||
@@ -163,10 +163,39 @@ def vods(username):
|
|||||||
|
|
||||||
|
|
||||||
# RTMP Server Routes
|
# RTMP Server Routes
|
||||||
|
|
||||||
|
@stream_bp.route("/init_stream", methods=["POST"])
|
||||||
|
def init_stream():
|
||||||
|
"""
|
||||||
|
Called by NGINX when OBS starts streaming.
|
||||||
|
Creates necessary directories and validates stream key.
|
||||||
|
"""
|
||||||
|
stream_key = request.form.get("name")
|
||||||
|
|
||||||
|
print(f"Stream initialization requested in nginx with key: {stream_key}")
|
||||||
|
|
||||||
|
with Database() as db:
|
||||||
|
# Check if valid stream key and user is allowed to stream
|
||||||
|
user_info = db.fetchone("""SELECT user_id, username, is_live
|
||||||
|
FROM users
|
||||||
|
WHERE stream_key = ?""", (stream_key,))
|
||||||
|
|
||||||
|
if not user_info:
|
||||||
|
print("Unauthorized - Invalid stream key", flush=True)
|
||||||
|
return "Unauthorized - Invalid stream key", 403
|
||||||
|
|
||||||
|
# Create necessary directories
|
||||||
|
username = user_info["username"]
|
||||||
|
create_local_directories(username)
|
||||||
|
|
||||||
|
return "OK", 200
|
||||||
|
|
||||||
|
|
||||||
@stream_bp.route("/publish_stream", methods=["POST"])
|
@stream_bp.route("/publish_stream", methods=["POST"])
|
||||||
def publish_stream():
|
def publish_stream():
|
||||||
"""
|
"""
|
||||||
Authenticates stream from streamer and publishes it to the site
|
Called when user clicks Start Stream in dashboard.
|
||||||
|
Sets up stream in database and starts thumbnail generation.
|
||||||
|
|
||||||
step-by-step:
|
step-by-step:
|
||||||
fetch user info from stream key
|
fetch user info from stream key
|
||||||
@@ -174,27 +203,26 @@ def publish_stream():
|
|||||||
set user as streaming
|
set user as streaming
|
||||||
periodically update thumbnail
|
periodically update thumbnail
|
||||||
"""
|
"""
|
||||||
stream_key = request.form.get("key")
|
data = request.form.get("data")
|
||||||
print("Stream request received")
|
|
||||||
|
|
||||||
# Open database connection
|
|
||||||
with Database() as db:
|
with Database() as db:
|
||||||
# Get user info from stream key
|
user_info = db.fetchone("""SELECT user_id, username, current_stream_title,
|
||||||
user_info = db.fetchone("""SELECT user_id, username, current_stream_title, current_selected_category_id, is_live
|
current_selected_category_id, is_live
|
||||||
FROM users
|
FROM users
|
||||||
WHERE stream_key = ?""", (stream_key,))
|
WHERE stream_key = ?""", (data['stream_key'],))
|
||||||
|
|
||||||
# If stream key is invalid, return unauthorized
|
|
||||||
if not user_info or user_info["is_live"]:
|
if not user_info or user_info["is_live"]:
|
||||||
|
print(
|
||||||
|
"Unauthorized. No user found from Stream key or user is already streaming.", flush=True)
|
||||||
return "Unauthorized", 403
|
return "Unauthorized", 403
|
||||||
|
|
||||||
# Insert stream into database
|
# Insert stream into database
|
||||||
db.execute("""INSERT INTO streams (user_id, title, start_time, num_viewers, category_id)
|
db.execute("""INSERT INTO streams (user_id, title, start_time, num_viewers, category_id)
|
||||||
VALUES (?, ?, ?, ?, ?)""", (user_info["user_id"],
|
VALUES (?, ?, ?, ?, ?)""", (user_info["user_id"],
|
||||||
user_info["current_stream_title"],
|
data["title"],
|
||||||
datetime.now(),
|
datetime.now(),
|
||||||
0,
|
0,
|
||||||
1))
|
get_category_id(data['category_name'])))
|
||||||
|
|
||||||
# Set user as streaming
|
# Set user as streaming
|
||||||
db.execute("""UPDATE users SET is_live = 1 WHERE user_id = ?""",
|
db.execute("""UPDATE users SET is_live = 1 WHERE user_id = ?""",
|
||||||
@@ -203,16 +231,41 @@ def publish_stream():
|
|||||||
username = user_info["username"]
|
username = user_info["username"]
|
||||||
user_id = user_info["user_id"]
|
user_id = user_info["user_id"]
|
||||||
|
|
||||||
# Local file creation
|
|
||||||
create_local_directories(username)
|
|
||||||
|
|
||||||
# Update thumbnail periodically
|
# Update thumbnail periodically
|
||||||
update_thumbnail.delay(user_id,
|
update_thumbnail.delay(user_id,
|
||||||
path_manager.get_stream_file_path(username),
|
path_manager.get_stream_file_path(username),
|
||||||
path_manager.get_thumbnail_file_path(username),
|
path_manager.get_thumbnail_file_path(username),
|
||||||
THUMBNAIL_GENERATION_INTERVAL)
|
THUMBNAIL_GENERATION_INTERVAL)
|
||||||
|
|
||||||
return redirect(f"/{user_info['username']}/stream/")
|
return "OK", 200
|
||||||
|
|
||||||
|
|
||||||
|
@stream_bp.route("/update_stream", methods=["POST"])
|
||||||
|
def update_stream():
|
||||||
|
"""
|
||||||
|
Called by StreamDashboard to update stream info
|
||||||
|
"""
|
||||||
|
# TODO: Add thumbnails (paths) to table, allow user to update thumbnail
|
||||||
|
|
||||||
|
stream_key = request.form.get("key")
|
||||||
|
title = request.form.get("title")
|
||||||
|
category_name = request.form.get("category_name")
|
||||||
|
|
||||||
|
with Database() as db:
|
||||||
|
user_info = db.fetchone("""SELECT user_id, username, is_live
|
||||||
|
FROM users
|
||||||
|
WHERE stream_key = ?""", (stream_key,))
|
||||||
|
|
||||||
|
if not user_info or not user_info["is_live"]:
|
||||||
|
print(
|
||||||
|
"Unauthorized - No user found from stream key or user is not streaming", flush=True)
|
||||||
|
return "Unauthorized", 403
|
||||||
|
|
||||||
|
db.execute("""UPDATE streams
|
||||||
|
SET title = ?, category_id = ?
|
||||||
|
WHERE user_id = ?""", (title, get_category_id(category_name), user_info["user_id"]))
|
||||||
|
|
||||||
|
return "Stream updated", 200
|
||||||
|
|
||||||
|
|
||||||
@stream_bp.route("/end_stream", methods=["POST"])
|
@stream_bp.route("/end_stream", methods=["POST"])
|
||||||
@@ -230,6 +283,12 @@ def end_stream():
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
stream_key = request.form.get("key")
|
stream_key = request.form.get("key")
|
||||||
|
if stream_key is None:
|
||||||
|
stream_key = request.form.get("name")
|
||||||
|
|
||||||
|
if stream_key is None:
|
||||||
|
print("Unauthorized - No stream key provided", flush=True)
|
||||||
|
return "Unauthorized", 403
|
||||||
|
|
||||||
# Open database connection
|
# Open database connection
|
||||||
with Database() as db:
|
with Database() as db:
|
||||||
@@ -244,6 +303,7 @@ def end_stream():
|
|||||||
|
|
||||||
# If stream key is invalid, return unauthorized
|
# If stream key is invalid, return unauthorized
|
||||||
if not user_info:
|
if not user_info:
|
||||||
|
print("Unauthorized - No user found from stream key", flush=True)
|
||||||
return "Unauthorized", 403
|
return "Unauthorized", 403
|
||||||
|
|
||||||
# Remove stream from database
|
# Remove stream from database
|
||||||
@@ -256,9 +316,9 @@ def end_stream():
|
|||||||
|
|
||||||
db.execute("""INSERT INTO vods (user_id, title, datetime, category_id, length, views)
|
db.execute("""INSERT INTO vods (user_id, title, datetime, category_id, length, views)
|
||||||
VALUES (?, ?, ?, ?, ?, ?)""", (user_info["user_id"],
|
VALUES (?, ?, ?, ?, ?, ?)""", (user_info["user_id"],
|
||||||
user_info["current_stream_title"],
|
stream_info["title"],
|
||||||
stream_info["start_time"],
|
stream_info["start_time"],
|
||||||
user_info["current_selected_category_id"],
|
stream_info["category_id"],
|
||||||
stream_length,
|
stream_length,
|
||||||
0))
|
0))
|
||||||
|
|
||||||
|
|||||||
@@ -142,41 +142,14 @@ def get_vod_tags(vod_id: int):
|
|||||||
""", (vod_id,))
|
""", (vod_id,))
|
||||||
return tags
|
return tags
|
||||||
|
|
||||||
def transfer_stream_to_vod(user_id: int):
|
|
||||||
"""
|
|
||||||
Deletes stream from stream table and moves it to VoD table
|
|
||||||
TODO: Add functionaliy to save stream permanently
|
|
||||||
"""
|
|
||||||
|
|
||||||
with Database() as db:
|
|
||||||
stream = db.fetchone("""
|
|
||||||
SELECT * FROM streams WHERE user_id = ?;
|
|
||||||
""", (user_id,))
|
|
||||||
|
|
||||||
if not stream:
|
|
||||||
return False
|
|
||||||
|
|
||||||
## TODO: calculate length in seconds, currently using temp value
|
|
||||||
|
|
||||||
db.execute("""
|
|
||||||
INSERT INTO vods (user_id, title, datetime, category_id, length, views)
|
|
||||||
VALUES (?, ?, ?, ?, ?, ?);
|
|
||||||
""", (stream["user_id"], stream["title"], stream["datetime"], stream["category_id"], 10, stream["num_viewers"]))
|
|
||||||
|
|
||||||
db.execute("""
|
|
||||||
DELETE FROM streams WHERE user_id = ?;
|
|
||||||
""", (user_id,))
|
|
||||||
|
|
||||||
return True
|
|
||||||
|
|
||||||
def create_local_directories(username: str):
|
def create_local_directories(username: str):
|
||||||
"""
|
"""
|
||||||
Create directories for user stream data if they do not exist
|
Create directories for user stream data if they do not exist
|
||||||
"""
|
"""
|
||||||
|
|
||||||
vods_path = f"stream_data/{username}/vods"
|
vods_path = f"stream_data/vods/{username}"
|
||||||
stream_path = f"stream_data/{username}/stream"
|
stream_path = f"stream_data/stream"
|
||||||
thumbnail_path = f"stream_data/{username}/thumbnails"
|
thumbnail_path = f"stream_data/thumbnails/{username}"
|
||||||
|
|
||||||
if not os.path.exists(vods_path):
|
if not os.path.exists(vods_path):
|
||||||
os.makedirs(vods_path)
|
os.makedirs(vods_path)
|
||||||
@@ -188,7 +161,6 @@ def create_local_directories(username: str):
|
|||||||
os.makedirs(thumbnail_path)
|
os.makedirs(thumbnail_path)
|
||||||
|
|
||||||
# Fix permissions
|
# Fix permissions
|
||||||
os.chmod(f"stream_data/{username}", 0o777)
|
|
||||||
os.chmod(vods_path, 0o777)
|
os.chmod(vods_path, 0o777)
|
||||||
os.chmod(stream_path, 0o777)
|
os.chmod(stream_path, 0o777)
|
||||||
os.chmod(thumbnail_path, 0o777)
|
os.chmod(thumbnail_path, 0o777)
|
||||||
|
|||||||
Reference in New Issue
Block a user