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