Così un CPX nuovo non richiede worker Docker a mano, l'agent è nel cloud-init e gli E2E restano non in elenco. Co-authored-by: Cursor <cursoragent@cursor.com>
454 lines
15 KiB
Ruby
454 lines
15 KiB
Ruby
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
|
|
|
|
# Home può avviare ffmpeg sull'agent del CPX (niente worker Docker per nodo).
|
|
def can_process?(session)
|
|
return false unless worker?
|
|
return true if node_matches?(session)
|
|
|
|
cloud_home_dispatch?(session)
|
|
end
|
|
|
|
def cloud_home_dispatch?(session)
|
|
local_node_slug == NodeRegistry::HOME_SLUG &&
|
|
session.stream_node&.role == "cloud" &&
|
|
session.stream_node.relay_agent_url.present?
|
|
end
|
|
|
|
def queue_for(session)
|
|
node = session.stream_node
|
|
if node&.role == "cloud" && node.relay_agent_url.present?
|
|
return "#{QUEUE_PREFIX}_cloud"
|
|
end
|
|
|
|
"#{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 can_process?(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) && can_process?(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 can_process?(session)
|
|
ensure_on_worker!(session)
|
|
else
|
|
enqueue_ensure!(session)
|
|
end
|
|
end
|
|
|
|
def ensure_on_worker!(session)
|
|
return :not_worker unless worker?
|
|
unless can_process?(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_for?(session)
|
|
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}")
|
|
enqueue_ensure!(session, wait: 5.seconds) unless session.terminal?
|
|
: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
|
|
|
|
# I relay sull'agent CPX non consumano ffmpeg locale: non applicare il cap del worker home.
|
|
def at_capacity_for?(session)
|
|
return false if cloud_home_dispatch?(session)
|
|
|
|
at_capacity?
|
|
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 agent_dispatch?(session)
|
|
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=#{agent_base_url(session) || '-'}"
|
|
)
|
|
schedule_youtube_activate(session)
|
|
pid
|
|
rescue Errno::ENOENT => e
|
|
raise Error, "ffmpeg non disponibile: #{e.message}"
|
|
end
|
|
|
|
def agent_base_url(session)
|
|
return if session.blank?
|
|
|
|
ENV["STREAM_NODE_RELAY_AGENT_URL"].presence || session.stream_node&.relay_agent_url
|
|
end
|
|
|
|
def agent_dispatch?(session)
|
|
agent_base_url(session).present?
|
|
end
|
|
|
|
def relay_agent_secret
|
|
ENV["STREAM_NODE_AGENT_SECRET"].presence || "mediamtx_webhook_dev_secret"
|
|
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, session: session)
|
|
pid = res["pid"].to_i
|
|
raise Error, "relay agent spawn failed: #{res.inspect}" if pid <= 0
|
|
|
|
File.write(
|
|
log_file(session),
|
|
"agent=#{agent_base_url(session)} pid=#{pid} path=#{session.mediamtx_path_name}\n"
|
|
)
|
|
pid
|
|
end
|
|
|
|
def agent_request(http_class, path, payload = nil, session: nil, session_id: nil)
|
|
session ||= StreamSession.find_by(id: session_id) if session_id.present?
|
|
base = agent_base_url(session)
|
|
raise Error, "relay agent URL mancante" if base.blank?
|
|
|
|
uri = URI.parse("#{base.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 agent_dispatch?(session) && (node_matches?(session) || cloud_home_dispatch?(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)
|
|
session = session_id.present? ? StreamSession.find_by(id: session_id) : nil
|
|
return agent_session_running?(session) if session && agent_dispatch?(session)
|
|
|
|
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)
|
|
res = agent_request(Net::HTTP::Get, "/relays/#{session.id}", session: session)
|
|
res["running"] == true
|
|
rescue Error
|
|
false
|
|
end
|
|
|
|
def terminate_pid(pid, session_id: nil)
|
|
session = session_id.present? ? StreamSession.find_by(id: session_id) : nil
|
|
if session && agent_dispatch?(session)
|
|
agent_request(Net::HTTP::Delete, "/relays/#{session_id}", session: session)
|
|
return
|
|
end
|
|
|
|
Process.kill("TERM", -pid.to_i)
|
|
sleep 0.5
|
|
rescue Errno::ESRCH
|
|
nil
|
|
else
|
|
return if session && agent_dispatch?(session)
|
|
|
|
begin
|
|
Process.kill("KILL", -pid.to_i)
|
|
rescue Errno::ESRCH
|
|
nil
|
|
end
|
|
end
|
|
end
|
|
end
|
|
end
|