Files
MatchLiveTv/backend/app/services/streams/youtube_relay.rb
T
eminuxandCursor 0b55d5a8fa Rende sticky e capacizzati i relay YouTube su coda dedicata.
Evita doppi ffmpeg tra host: stop solo sull'owner, ensure con requeue a capacità piena e coda Sidekiq youtube_relay isolata.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-09 19:54:04 +02:00

285 lines
8.8 KiB
Ruby

module Streams
# Relay verso YouTube: legge RTMP/HLS da MediaMTX e inoltra su RTMPS (-c copy). Nessun overlay.
# ffmpeg gira solo sui worker Sidekiq con YOUTUBE_RELAY_WORKER=1 (coda youtube_relay).
class YoutubeRelay
class Error < StandardError; end
REDIS_KEY = "youtube_relay:pid:%s"
OWNER_KEY = "youtube_relay:owner:%s"
OWNED_SET = "youtube_relay:owned:%s"
QUEUE = :youtube_relay
class << self
def worker?
ENV["YOUTUBE_RELAY_WORKER"] == "1"
end
def max_concurrent
ENV.fetch("RELAY_MAX_CONCURRENT", "4").to_i
end
def start(session)
return unless worker?
start_on_worker!(session)
end
# Non cancella owner/pid qui: solo il worker owner deve killare ffmpeg.
def stop(session)
if worker? && owner_is_local?(session.id)
stop_on_worker!(session)
else
YoutubeRelayStopJob.set(queue: QUEUE).perform_later(session.id)
end
true
end
def running?(session_id)
owner = redis.get(format(OWNER_KEY, session_id))
pid = pid_for(session_id)
return false if pid.blank? && owner.blank?
return false if pid.blank? && owner.present? && owner != worker_id && redis.ttl(format(OWNER_KEY, session_id)) <= 30
return false if pid.blank?
return process_alive?(pid) if owner.blank? || owner == worker_id
# Relay su altro host: attivo se lock owner ancora fresco.
redis.ttl(format(OWNER_KEY, session_id)) > 30
end
def ensure_publishing!(session)
return unless session.platform == "youtube"
return if session.terminal?
return if session.stream_key.blank?
if worker?
ensure_on_worker!(session)
else
YoutubeRelayEnsureJob.set(queue: QUEUE).perform_later(session.id)
end
end
def ensure_on_worker!(session)
return :not_worker unless worker?
return unless session.platform == "youtube"
return if session.terminal?
return if session.stream_key.blank?
return unless session.status.in?(%w[live connecting reconnecting paused])
return unless intake_available?(session)
pid = pid_for(session.id)
if pid.present? && !process_alive?(pid.to_i)
clear_local_ownership(session.id)
end
if running?(session.id)
touch_owner!(session.id) if owner_is_local?(session.id)
return :already_running
end
if at_capacity?
YoutubeRelayEnsureJob.set(wait: 5.seconds, queue: QUEUE).perform_later(session.id)
Rails.logger.info("[YoutubeRelay] at capacity worker=#{worker_id} session=#{session.id} requeue")
return :at_capacity
end
last_restart = redis.get(restart_debounce_key(session.id)).to_i
return :debounced if last_restart.positive? && (Time.now.to_i - last_restart) < 5
start_on_worker!(session)
redis.set(restart_debounce_key(session.id), Time.now.to_i, ex: 300)
:started
rescue Error => e
Rails.logger.warn("[YoutubeRelay] ensure_on_worker session=#{session.id}: #{e.message}")
:error
end
# @return [Symbol] :stopped, :wrong_host, :noop
def stop_on_worker!(session)
return :not_worker unless worker?
owner = redis.get(format(OWNER_KEY, session.id))
if owner.present? && owner != worker_id
return :wrong_host
end
pid = pid_for(session.id)
if pid.blank?
clear_local_ownership(session.id)
return :noop
end
terminate_pid(pid)
clear_local_ownership(session.id)
Rails.logger.info("[YoutubeRelay] stopped pid=#{pid} session=#{session.id} worker=#{worker_id}")
:stopped
end
def local_owned_count
redis.scard(format(OWNED_SET, worker_id)).to_i
end
def at_capacity?
local_owned_count >= max_concurrent
end
def owner_is_local?(session_id)
owner = redis.get(format(OWNER_KEY, session_id))
owner.blank? || owner == worker_id
end
private
def start_on_worker!(session)
return unless session.platform == "youtube"
return if session.stream_key.blank?
return if session.terminal?
return unless intake_available?(session)
return pid_for(session.id).to_i if running?(session.id) && owner_is_local?(session.id)
stop_on_worker!(session) if pid_for(session.id).present? && owner_is_local?(session.id)
log_path = log_file(session)
FileUtils.mkdir_p(File.dirname(log_path))
intake_source = mediamtx_intake_source(session)
output = "rtmps://a.rtmps.youtube.com/live2/#{session.stream_key}"
pid = Process.spawn(
*youtube_ffmpeg_args(intake_source, output),
%i[out err] => log_path,
pgroup: true
)
Process.detach(pid)
store_pid(session.id, pid)
claim_ownership!(session.id)
Rails.logger.info("[YoutubeRelay] started pid=#{pid} session=#{session.id} intake=#{intake_source.join(":")} worker=#{worker_id}")
schedule_youtube_activate(session)
pid
rescue Errno::ENOENT => e
raise Error, "ffmpeg non disponibile: #{e.message}"
end
# Remux verso YouTube: video+audio copy (AAC già da app/slate a 48k mono).
# HLS (ADTS) → FLV richiede -bsf:a aac_adtstoasc; RTMP ha già ASC in FLV tags.
def youtube_ffmpeg_args(intake_source, output)
mode, url = intake_source
common_head = [
"ffmpeg", "-nostdin", "-hide_banner", "-loglevel", "warning",
"-fflags", "+genpts+discardcorrupt"
]
copy_tail = ["-c:v", "copy", "-c:a", "copy", "-f", "flv", output]
if mode == :rtmp
return common_head + [
"-analyzeduration", "10000000", "-probesize", "10000000",
"-rw_timeout", "15000000",
"-i", url
] + copy_tail
end
common_head + [
"-rw_timeout", "15000000",
"-live_start_index", "-1",
"-i", url,
"-c:v", "copy",
"-c:a", "copy",
"-bsf:a", "aac_adtstoasc",
"-f", "flv", output
]
end
def mediamtx_intake_source(session)
base = session.mediamtx_internal_rtmp_url
if Mediamtx::PublisherOnline.active?(session)
return [:rtmp, "#{base.chomp('/')}/#{session.mediamtx_path_name}"]
end
path = session.mediamtx_path_name
hls = session.mediamtx_internal_hls_url.chomp("/")
[:hls, "#{hls}/#{path}/index.m3u8"]
end
def intake_available?(session)
return true if Mediamtx::PublisherOnline.active?(session)
info = Mediamtx::Client.for_session(session).list_paths.find { |i| i["name"] == session.mediamtx_path_name }
info && (info["ready"] || info["online"] || info["available"])
rescue StandardError
false
end
def schedule_youtube_activate(session)
Youtube::LivePipeline.schedule_activate!(session, force: true)
end
def worker_id
ENV.fetch("HOSTNAME", "worker")
end
def log_file(session)
Rails.root.join("log", "youtube_relay_#{session.id}.log").to_s
end
def redis
@redis ||= Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0"))
end
def restart_debounce_key(session_id)
format("youtube_relay:debounce:%s", session_id)
end
def claim_ownership!(session_id)
redis.set(format(OWNER_KEY, session_id), worker_id, ex: 48.hours.to_i)
redis.sadd(format(OWNED_SET, worker_id), session_id)
end
def touch_owner!(session_id)
redis.expire(format(OWNER_KEY, session_id), 48.hours.to_i)
redis.expire(format(REDIS_KEY, session_id), 48.hours.to_i)
end
def clear_local_ownership(session_id)
redis.del(format(REDIS_KEY, session_id))
redis.del(format(OWNER_KEY, session_id))
redis.srem(format(OWNED_SET, worker_id), session_id)
end
def store_pid(session_id, pid)
redis.set(format(REDIS_KEY, session_id), pid, ex: 48.hours.to_i)
end
def clear_pid(session_id)
redis.del(format(REDIS_KEY, session_id))
end
def pid_for(session_id)
redis.get(format(REDIS_KEY, session_id))
end
def process_alive?(pid)
stat = File.read("/proc/#{pid.to_i}/stat")
return false if stat.split[2] == "Z"
Process.kill(0, pid.to_i)
true
rescue Errno::ESRCH, Errno::EPERM, Errno::ENOENT
false
end
def terminate_pid(pid)
Process.kill("TERM", -pid.to_i)
sleep 0.5
rescue Errno::ESRCH
nil
else
begin
Process.kill("KILL", -pid.to_i)
rescue Errno::ESRCH
nil
end
end
end
end
end