require "json" require "net/http" require "uri" module Streams # Relay verso YouTube: legge HLS da MediaMTX e inoltra su RTMPS (-c copy). # Sul nodo cloud ffmpeg è avviato dall'agent locale (immagine CPX) all'avvio diretta. 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 QUEUE_PREFIX = "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 local_node_slug ENV["STREAM_NODE_SLUG"].presence || NodeRegistry::HOME_SLUG end def assigned_node_slug(session) session.stream_node&.slug.presence || NodeRegistry::HOME_SLUG end def node_matches?(session) assigned_node_slug(session) == local_node_slug end def queue_for(session) "#{QUEUE_PREFIX}_#{assigned_node_slug(session)}" end def enqueue_ensure!(session, wait: nil) opts = { queue: queue_for(session) } opts[:wait] = wait if wait YoutubeRelayEnsureJob.set(**opts).perform_later(session.id) end def enqueue_stop!(session) YoutubeRelayStopJob.set(queue: queue_for(session)).perform_later(session.id) end def start(session) return enqueue_ensure!(session) unless worker? && node_matches?(session) 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) && node_matches?(session) stop_on_worker!(session) else enqueue_stop!(session) 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, session_id: session_id) 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? && node_matches?(session) ensure_on_worker!(session) else enqueue_ensure!(session) end end def ensure_on_worker!(session) return :not_worker unless worker? unless node_matches?(session) enqueue_ensure!(session, wait: 2.seconds) Rails.logger.info( "[YoutubeRelay] wrong_node local=#{local_node_slug} assigned=#{assigned_node_slug(session)} session=#{session.id}" ) return :wrong_host end 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, session_id: session.id) 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? enqueue_ensure!(session, wait: 5.seconds) Rails.logger.info("[YoutubeRelay] at capacity worker=#{worker_id} node=#{local_node_slug} 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, session_id: session.id) 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) # Due EnsureJob concorrenti (concurrency Sidekiq > 1) non devono spawnare due ffmpeg. unless redis.set(format("youtube_relay:startlock:%s", session.id), worker_id, nx: true, ex: 20) return pid_for(session.id).to_i end 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 = if relay_agent_enabled? start_via_agent!(session) else spawn_ffmpeg!(youtube_ffmpeg_args(intake_source, output), log_path) end 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} agent=#{relay_agent_url || '-'}" ) schedule_youtube_activate(session) pid rescue Errno::ENOENT => e raise Error, "ffmpeg non disponibile: #{e.message}" end def relay_agent_url ENV["STREAM_NODE_RELAY_AGENT_URL"].presence end def relay_agent_enabled? relay_agent_url.present? end def relay_agent_secret ENV["STREAM_NODE_AGENT_SECRET"].presence || ENV["MEDIAMTX_WEBHOOK_SECRET"].presence || "" end def spawn_ffmpeg!(args, log_path) pid = Process.spawn(*args, %i[out err] => log_path, pgroup: true) Process.detach(pid) pid end def start_via_agent!(session) payload = { "session_id" => session.id, "path" => session.mediamtx_path_name, "rtmps" => "rtmps://a.rtmps.youtube.com/live2/#{session.stream_key}" } res = agent_request(Net::HTTP::Post, "/relays", payload) pid = res["pid"].to_i raise Error, "relay agent spawn failed: #{res.inspect}" if pid <= 0 File.write( log_file(session), "agent=#{relay_agent_url} pid=#{pid} path=#{session.mediamtx_path_name}\n" ) pid end def agent_request(http_class, path, payload = nil) uri = URI.parse("#{relay_agent_url.chomp('/')}#{path}") http = Net::HTTP.new(uri.host, uri.port) http.open_timeout = 5 http.read_timeout = 10 req = http_class.new(uri) req["Authorization"] = "Bearer #{relay_agent_secret}" if relay_agent_secret.present? req["Content-Type"] = "application/json" req.body = JSON.generate(payload) if payload res = http.request(req) body = res.body.present? ? JSON.parse(res.body) : {} unless res.is_a?(Net::HTTPSuccess) raise Error, "relay agent #{path} HTTP #{res.code} #{body.inspect}" end body rescue JSON::ParserError => e raise Error, "relay agent JSON: #{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 + [ "-reconnect", "1", "-reconnect_streamed", "1", "-reconnect_delay_max", "2", "-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) rtmp_base, hls_base = intake_bases(session) if Mediamtx::PublisherOnline.active?(session) return [:rtmp, "#{rtmp_base.chomp('/')}/#{session.mediamtx_path_name}"] end path = session.mediamtx_path_name [:hls, "#{hls_base.chomp('/')}/#{path}/index.m3u8"] end # Sul nodo assegnato ffmpeg legge MediaMTX in loopback (o hostname Docker locale). def intake_bases(session) if relay_agent_enabled? && node_matches?(session) return ["rtmp://127.0.0.1:1935", "http://127.0.0.1:8888"] end if node_matches?(session) [ ENV["STREAM_NODE_LOCAL_RTMP_URL"].presence || session.mediamtx_internal_rtmp_url, ENV["STREAM_NODE_LOCAL_HLS_URL"].presence || session.mediamtx_internal_hls_url ] else [session.mediamtx_internal_rtmp_url, session.mediamtx_internal_hls_url] end 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, session_id: nil) return agent_session_running?(session_id) if relay_agent_enabled? && session_id.present? 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 agent_session_running?(session_id) res = agent_request(Net::HTTP::Get, "/relays/#{session_id}") res["running"] == true rescue Error false end def terminate_pid(pid, session_id: nil) if relay_agent_enabled? && session_id.present? agent_request(Net::HTTP::Delete, "/relays/#{session_id}") return end Process.kill("TERM", -pid.to_i) sleep 0.5 rescue Errno::ESRCH nil else return if relay_agent_enabled? begin Process.kill("KILL", -pid.to_i) rescue Errno::ESRCH nil end end end end end