Aggiunge registry nodi stream e provisioning Hetzner/lab per lo scale-out.
Prepara l'architettura multi-nodo (MediaMTX+ffmpeg) con assignment URL per sessione, admin di provision/drain e provider Cloud/DNS astratti verso mltv-stream.net. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -0,0 +1,15 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
module Streams
|
||||
module CloudProviders
|
||||
def self.build(name = ENV.fetch("STREAM_CLOUD_PROVIDER", "local_lab"))
|
||||
case name.to_s
|
||||
when "local_lab" then LocalLab.new
|
||||
when "proxmox_lab" then ProxmoxLab.new
|
||||
when "hetzner" then Hetzner.new
|
||||
else
|
||||
raise Error, "STREAM_CLOUD_PROVIDER sconosciuto: #{name}"
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,35 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
module Streams
|
||||
module CloudProviders
|
||||
class Error < StandardError; end
|
||||
|
||||
# Descrittore restituito da create_node / list.
|
||||
Instance = Struct.new(
|
||||
:id, :name, :public_ip, :private_ip, :status, :raw,
|
||||
keyword_init: true
|
||||
)
|
||||
|
||||
class Base
|
||||
def create_node(name:, labels: {})
|
||||
raise NotImplementedError
|
||||
end
|
||||
|
||||
def destroy_node(instance_id)
|
||||
raise NotImplementedError
|
||||
end
|
||||
|
||||
def list_nodes(labels: {})
|
||||
raise NotImplementedError
|
||||
end
|
||||
|
||||
def wait_until_running(instance_id, timeout: 120)
|
||||
raise NotImplementedError
|
||||
end
|
||||
|
||||
def public_ip(instance_id)
|
||||
raise NotImplementedError
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,183 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
require "faraday"
|
||||
|
||||
module Streams
|
||||
module CloudProviders
|
||||
# Hetzner Cloud — create/destroy server per nodi stream.
|
||||
#
|
||||
# ENV:
|
||||
# HCLOUD_TOKEN (obbligatorio)
|
||||
# HCLOUD_LOCATION (default fsn1)
|
||||
# HCLOUD_SERVER_TYPE (default cpx21)
|
||||
# HCLOUD_IMAGE (default debian-12)
|
||||
# HCLOUD_SSH_KEY (nome chiave, default matchlivetv-stream)
|
||||
# HCLOUD_NETWORK_ID (opzionale, private network / WireGuard prep)
|
||||
# HCLOUD_USER_DATA_FILE (opzionale, cloud-init path)
|
||||
class Hetzner < Base
|
||||
API = "https://api.hetzner.cloud/v1"
|
||||
|
||||
def initialize(token: ENV.fetch("HCLOUD_TOKEN"), conn: nil)
|
||||
@token = token
|
||||
@conn = conn
|
||||
end
|
||||
|
||||
def create_node(name:, labels: {})
|
||||
body = {
|
||||
name: name,
|
||||
server_type: ENV.fetch("HCLOUD_SERVER_TYPE", "cpx21"),
|
||||
image: ENV.fetch("HCLOUD_IMAGE", "debian-12"),
|
||||
location: ENV.fetch("HCLOUD_LOCATION", "fsn1"),
|
||||
start_after_create: true,
|
||||
labels: default_labels.merge(stringify_labels(labels)),
|
||||
ssh_keys: [ENV.fetch("HCLOUD_SSH_KEY", "matchlivetv-stream")],
|
||||
public_net: {
|
||||
enable_ipv4: true,
|
||||
enable_ipv6: false
|
||||
}
|
||||
}
|
||||
network_id = ENV["HCLOUD_NETWORK_ID"].presence
|
||||
body[:networks] = [network_id.to_i] if network_id
|
||||
user_data = cloud_init_user_data
|
||||
body[:user_data] = user_data if user_data.present?
|
||||
|
||||
data = post("/servers", body)
|
||||
server = data["server"] || {}
|
||||
action = data["action"]
|
||||
wait_action!(action) if action
|
||||
instance = wait_until_running(server["id"].to_s)
|
||||
instance.name = name
|
||||
instance
|
||||
end
|
||||
|
||||
def destroy_node(instance_id)
|
||||
delete("/servers/#{instance_id}")
|
||||
true
|
||||
end
|
||||
|
||||
def list_nodes(labels: {})
|
||||
params = {}
|
||||
label_selector = labels.map { |k, v| "#{k}=#{v}" }.join(",")
|
||||
params[:label_selector] = label_selector if label_selector.present?
|
||||
params[:label_selector] ||= "matchlivetv=true,role=stream-node"
|
||||
|
||||
data = get("/servers", params)
|
||||
Array(data["servers"]).map { |s| instance_from_server(s) }
|
||||
end
|
||||
|
||||
def wait_until_running(instance_id, timeout: 180)
|
||||
deadline = Time.now + timeout
|
||||
loop do
|
||||
data = get("/servers/#{instance_id}")
|
||||
server = data["server"]
|
||||
status = server["status"]
|
||||
if status == "running"
|
||||
return instance_from_server(server)
|
||||
end
|
||||
raise Error, "Timeout attesa server Hetzner #{instance_id} (status=#{status})" if Time.now >= deadline
|
||||
|
||||
sleep 3
|
||||
end
|
||||
end
|
||||
|
||||
def public_ip(instance_id)
|
||||
wait_until_running(instance_id).public_ip
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def default_labels
|
||||
{
|
||||
"matchlivetv" => "true",
|
||||
"role" => "stream-node",
|
||||
"env" => ENV.fetch("STREAM_NODE_ENV", "prod")
|
||||
}
|
||||
end
|
||||
|
||||
def stringify_labels(labels)
|
||||
labels.to_h.transform_keys(&:to_s).transform_values(&:to_s)
|
||||
end
|
||||
|
||||
def instance_from_server(server)
|
||||
public_ip = server.dig("public_net", "ipv4", "ip")
|
||||
private_ip = Array(server["private_net"]).first&.dig("ip")
|
||||
Instance.new(
|
||||
id: server["id"].to_s,
|
||||
name: server["name"],
|
||||
public_ip: public_ip,
|
||||
private_ip: private_ip.presence || public_ip,
|
||||
status: server["status"],
|
||||
raw: server
|
||||
)
|
||||
end
|
||||
|
||||
def cloud_init_user_data
|
||||
path = ENV["HCLOUD_USER_DATA_FILE"].presence
|
||||
return File.read(path) if path && File.file?(path)
|
||||
|
||||
ENV["HCLOUD_USER_DATA"].presence
|
||||
end
|
||||
|
||||
def conn
|
||||
@conn ||= Faraday.new(url: API) do |f|
|
||||
f.request :json
|
||||
f.response :json, content_type: /\bjson$/
|
||||
f.adapter Faraday.default_adapter
|
||||
end
|
||||
end
|
||||
|
||||
def auth_headers
|
||||
{ "Authorization" => "Bearer #{@token}" }
|
||||
end
|
||||
|
||||
def get(path, params = {})
|
||||
response = conn.get(path) do |req|
|
||||
req.headers.update(auth_headers)
|
||||
req.params.update(params)
|
||||
end
|
||||
unwrap!(response)
|
||||
end
|
||||
|
||||
def post(path, body)
|
||||
response = conn.post(path) do |req|
|
||||
req.headers.update(auth_headers)
|
||||
req.body = body
|
||||
end
|
||||
unwrap!(response)
|
||||
end
|
||||
|
||||
def delete(path)
|
||||
response = conn.delete(path) do |req|
|
||||
req.headers.update(auth_headers)
|
||||
end
|
||||
return {} if response.status == 204
|
||||
|
||||
unwrap!(response)
|
||||
end
|
||||
|
||||
def unwrap!(response)
|
||||
unless response.success?
|
||||
raise Error, "Hetzner Cloud API #{response.status}: #{response.body.inspect}"
|
||||
end
|
||||
|
||||
response.body.is_a?(Hash) ? response.body : {}
|
||||
end
|
||||
|
||||
def wait_action!(action, timeout: 120)
|
||||
return unless action.is_a?(Hash) && action["id"]
|
||||
|
||||
deadline = Time.now + timeout
|
||||
id = action["id"]
|
||||
loop do
|
||||
data = get("/actions/#{id}")
|
||||
status = data.dig("action", "status")
|
||||
return if status == "success"
|
||||
raise Error, "Hetzner action #{id} failed: #{data.inspect}" if status == "error"
|
||||
raise Error, "Timeout action Hetzner #{id}" if Time.now >= deadline
|
||||
|
||||
sleep 2
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,45 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
module Streams
|
||||
module CloudProviders
|
||||
# Lab senza API Proxmox: simula create/destroy e riusa MediaMTX home per i path.
|
||||
# Utile per testare registry, assignment e admin senza secondi host.
|
||||
class LocalLab < Base
|
||||
def create_node(name:, labels: {})
|
||||
Instance.new(
|
||||
id: "sim-#{name}",
|
||||
name: name,
|
||||
public_ip: labels[:public_ip].presence || "127.0.0.1",
|
||||
private_ip: labels[:private_ip].presence || "127.0.0.1",
|
||||
status: "running",
|
||||
raw: { simulated: true, labels: labels }
|
||||
)
|
||||
end
|
||||
|
||||
def destroy_node(instance_id)
|
||||
true
|
||||
end
|
||||
|
||||
def list_nodes(labels: {})
|
||||
StreamNode.where(provider: "local", role: "lab").map do |node|
|
||||
Instance.new(
|
||||
id: node.provider_instance_id,
|
||||
name: node.slug,
|
||||
public_ip: node.metadata["public_ip"],
|
||||
private_ip: node.metadata["private_ip"],
|
||||
status: node.status == "ready" ? "running" : node.status,
|
||||
raw: node.metadata
|
||||
)
|
||||
end
|
||||
end
|
||||
|
||||
def wait_until_running(instance_id, timeout: 120)
|
||||
Instance.new(id: instance_id, name: instance_id, status: "running")
|
||||
end
|
||||
|
||||
def public_ip(instance_id)
|
||||
"127.0.0.1"
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,162 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
require "faraday"
|
||||
|
||||
module Streams
|
||||
module CloudProviders
|
||||
# Clone/start/stop di VM template su Proxmox VE (API token).
|
||||
#
|
||||
# ENV richiesti:
|
||||
# PROXMOX_API_URL, PROXMOX_TOKEN_ID, PROXMOX_TOKEN_SECRET,
|
||||
# PROXMOX_NODE, PROXMOX_TEMPLATE_VMID
|
||||
class ProxmoxLab < Base
|
||||
def initialize(
|
||||
api_url: ENV.fetch("PROXMOX_API_URL"),
|
||||
token_id: ENV.fetch("PROXMOX_TOKEN_ID"),
|
||||
token_secret: ENV.fetch("PROXMOX_TOKEN_SECRET"),
|
||||
node: ENV.fetch("PROXMOX_NODE"),
|
||||
template_vmid: ENV.fetch("PROXMOX_TEMPLATE_VMID"),
|
||||
verify_ssl: ENV.fetch("PROXMOX_VERIFY_SSL", "false") == "true"
|
||||
)
|
||||
@api_url = api_url.to_s.chomp("/")
|
||||
@token_id = token_id
|
||||
@token_secret = token_secret
|
||||
@node = node
|
||||
@template_vmid = template_vmid.to_i
|
||||
@verify_ssl = verify_ssl
|
||||
end
|
||||
|
||||
def create_node(name:, labels: {})
|
||||
newid = next_vmid
|
||||
post("/nodes/#{@node}/qemu/#{@template_vmid}/clone", {
|
||||
newid: newid,
|
||||
name: name,
|
||||
full: 1,
|
||||
target: @node
|
||||
})
|
||||
post("/nodes/#{@node}/qemu/#{newid}/status/start", {})
|
||||
wait_until_running(newid.to_s)
|
||||
ip = public_ip(newid.to_s)
|
||||
Instance.new(
|
||||
id: newid.to_s,
|
||||
name: name,
|
||||
public_ip: ip,
|
||||
private_ip: ip,
|
||||
status: "running",
|
||||
raw: { node: @node, vmid: newid, labels: labels }
|
||||
)
|
||||
end
|
||||
|
||||
def destroy_node(instance_id)
|
||||
vmid = instance_id.to_i
|
||||
begin
|
||||
post("/nodes/#{@node}/qemu/#{vmid}/status/stop", { timeout: 30 })
|
||||
rescue Error
|
||||
# già spenta
|
||||
end
|
||||
sleep 2
|
||||
delete("/nodes/#{@node}/qemu/#{vmid}", { purge: 1 })
|
||||
true
|
||||
end
|
||||
|
||||
def list_nodes(labels: {})
|
||||
items = get("/nodes/#{@node}/qemu")
|
||||
Array(items).filter_map do |row|
|
||||
name = row["name"].to_s
|
||||
next unless name.start_with?("mltv-stream-") || name.start_with?("ingest-")
|
||||
|
||||
Instance.new(
|
||||
id: row["vmid"].to_s,
|
||||
name: name,
|
||||
public_ip: nil,
|
||||
private_ip: nil,
|
||||
status: row["status"],
|
||||
raw: row
|
||||
)
|
||||
end
|
||||
end
|
||||
|
||||
def wait_until_running(instance_id, timeout: 180)
|
||||
deadline = Time.now + timeout
|
||||
loop do
|
||||
status = get("/nodes/#{@node}/qemu/#{instance_id}/status/current")
|
||||
return Instance.new(id: instance_id.to_s, status: "running", raw: status) if status["status"] == "running"
|
||||
raise Error, "Timeout attesa VM #{instance_id}" if Time.now >= deadline
|
||||
|
||||
sleep 3
|
||||
end
|
||||
end
|
||||
|
||||
def public_ip(instance_id)
|
||||
agent = get("/nodes/#{@node}/qemu/#{instance_id}/agent/network-get-interfaces")
|
||||
interfaces = agent.is_a?(Hash) ? agent["result"] : nil
|
||||
Array(interfaces).each do |iface|
|
||||
Array(iface["ip-addresses"]).each do |addr|
|
||||
ip = addr["ip-address"].to_s
|
||||
next if ip.blank? || ip.start_with?("127.") || ip.include?(":")
|
||||
|
||||
return ip
|
||||
end
|
||||
end
|
||||
ENV["STREAM_LAB_FALLBACK_IP"].presence || "127.0.0.1"
|
||||
rescue Error
|
||||
ENV["STREAM_LAB_FALLBACK_IP"].presence || "127.0.0.1"
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def next_vmid
|
||||
used = Array(get("/cluster/resources", type: "vm")).map { |r| r["vmid"].to_i }
|
||||
candidate = ENV.fetch("PROXMOX_VMID_START", "9100").to_i
|
||||
candidate += 1 while used.include?(candidate)
|
||||
candidate
|
||||
end
|
||||
|
||||
def conn
|
||||
@conn ||= Faraday.new(url: "#{@api_url}/api2/json") do |f|
|
||||
f.request :url_encoded
|
||||
f.response :json, content_type: /\bjson$/
|
||||
f.adapter Faraday.default_adapter
|
||||
f.ssl[:verify] = @verify_ssl
|
||||
end
|
||||
end
|
||||
|
||||
def auth_headers
|
||||
{ "Authorization" => "PVEAPIToken=#{@token_id}=#{@token_secret}" }
|
||||
end
|
||||
|
||||
def get(path, params = {})
|
||||
response = conn.get(path) do |req|
|
||||
req.headers.update(auth_headers)
|
||||
req.params.update(params)
|
||||
end
|
||||
unwrap!(response)
|
||||
end
|
||||
|
||||
def post(path, body = {})
|
||||
response = conn.post(path) do |req|
|
||||
req.headers.update(auth_headers)
|
||||
req.body = body
|
||||
end
|
||||
unwrap!(response)
|
||||
end
|
||||
|
||||
def delete(path, params = {})
|
||||
response = conn.delete(path) do |req|
|
||||
req.headers.update(auth_headers)
|
||||
req.params.update(params)
|
||||
end
|
||||
unwrap!(response)
|
||||
end
|
||||
|
||||
def unwrap!(response)
|
||||
unless response.success?
|
||||
raise Error, "Proxmox API #{response.status}: #{response.body.inspect}"
|
||||
end
|
||||
|
||||
body = response.body
|
||||
body.is_a?(Hash) && body.key?("data") ? body["data"] : body
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,14 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
module Streams
|
||||
module DnsProviders
|
||||
def self.build(name = ENV.fetch("STREAM_DNS_PROVIDER", "lab"))
|
||||
case name.to_s
|
||||
when "lab" then Lab.new
|
||||
when "hetzner" then Hetzner.new
|
||||
else
|
||||
raise Error, "STREAM_DNS_PROVIDER sconosciuto: #{name}"
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,21 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
module Streams
|
||||
module DnsProviders
|
||||
class Error < StandardError; end
|
||||
|
||||
class Base
|
||||
def upsert_a(name, ip)
|
||||
raise NotImplementedError
|
||||
end
|
||||
|
||||
def delete_a(name)
|
||||
raise NotImplementedError
|
||||
end
|
||||
|
||||
def resolve(name)
|
||||
raise NotImplementedError
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,105 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
require "faraday"
|
||||
require "cgi"
|
||||
|
||||
module Streams
|
||||
module DnsProviders
|
||||
# Hetzner Cloud DNS (Console) — zone mltv-stream.net.
|
||||
#
|
||||
# ENV:
|
||||
# HCLOUD_TOKEN (stesso del Cloud)
|
||||
# STREAM_DNS_ZONE (default mltv-stream.net)
|
||||
# STREAM_DNS_TTL (default 60)
|
||||
class Hetzner < Base
|
||||
API = "https://api.hetzner.cloud/v1"
|
||||
|
||||
def initialize(token: ENV.fetch("HCLOUD_TOKEN"), zone: nil, conn: nil)
|
||||
@token = token
|
||||
@zone = zone || ENV.fetch("STREAM_DNS_ZONE", "mltv-stream.net")
|
||||
@ttl = ENV.fetch("STREAM_DNS_TTL", "60").to_i
|
||||
@conn = conn
|
||||
end
|
||||
|
||||
def upsert_a(name, ip)
|
||||
rr_name = relative_name(name)
|
||||
delete_a(name)
|
||||
post("/zones/#{CGI.escape(@zone)}/rrsets", {
|
||||
name: rr_name,
|
||||
type: "A",
|
||||
ttl: @ttl,
|
||||
records: [{ value: ip.to_s, comment: "matchlivetv stream-node" }],
|
||||
labels: { "matchlivetv" => "true", "role" => "stream-node" }
|
||||
})
|
||||
true
|
||||
end
|
||||
|
||||
def delete_a(name)
|
||||
rr_name = relative_name(name)
|
||||
encoded = CGI.escape(rr_name)
|
||||
response = conn.delete("/zones/#{CGI.escape(@zone)}/rrsets/#{encoded}/A") do |req|
|
||||
req.headers.update(auth_headers)
|
||||
end
|
||||
return true if response.status == 404 || response.status == 204 || response.success?
|
||||
|
||||
raise Error, "Hetzner DNS API #{response.status}: #{response.body.inspect}"
|
||||
end
|
||||
|
||||
def resolve(name)
|
||||
rr_name = relative_name(name)
|
||||
data = get("/zones/#{CGI.escape(@zone)}/rrsets", name: rr_name, type: "A")
|
||||
rrset = Array(data["rrsets"]).first
|
||||
Array(rrset&.dig("records")).first&.dig("value")
|
||||
rescue Error
|
||||
nil
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def relative_name(name)
|
||||
host = name.to_s.strip.downcase.delete_suffix(".")
|
||||
suffix = ".#{@zone}"
|
||||
return "@" if host == @zone
|
||||
return host.delete_suffix(suffix) if host.end_with?(suffix)
|
||||
|
||||
host
|
||||
end
|
||||
|
||||
def conn
|
||||
@conn ||= Faraday.new(url: API) do |f|
|
||||
f.request :json
|
||||
f.response :json, content_type: /\bjson$/
|
||||
f.adapter Faraday.default_adapter
|
||||
end
|
||||
end
|
||||
|
||||
def auth_headers
|
||||
{ "Authorization" => "Bearer #{@token}" }
|
||||
end
|
||||
|
||||
def get(path, params = {})
|
||||
response = conn.get(path) do |req|
|
||||
req.headers.update(auth_headers)
|
||||
req.params.update(params)
|
||||
end
|
||||
unwrap!(response)
|
||||
end
|
||||
|
||||
def post(path, body)
|
||||
response = conn.post(path) do |req|
|
||||
req.headers.update(auth_headers)
|
||||
req.body = body
|
||||
end
|
||||
unwrap!(response)
|
||||
end
|
||||
|
||||
def unwrap!(response)
|
||||
unless response.success?
|
||||
raise Error, "Hetzner DNS API #{response.status}: #{response.body.inspect}"
|
||||
end
|
||||
|
||||
response.body.is_a?(Hash) ? response.body : {}
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,47 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
module Streams
|
||||
module DnsProviders
|
||||
# DNS lab in Redis (e dump hosts). Nessuna chiamata al registrar.
|
||||
class Lab < Base
|
||||
REDIS_KEY = "stream_dns:a_records"
|
||||
|
||||
def initialize(redis: nil)
|
||||
@redis = redis
|
||||
end
|
||||
|
||||
def upsert_a(name, ip)
|
||||
host = normalize(name)
|
||||
redis.hset(REDIS_KEY, host, ip.to_s)
|
||||
true
|
||||
end
|
||||
|
||||
def delete_a(name)
|
||||
redis.hdel(REDIS_KEY, normalize(name))
|
||||
true
|
||||
end
|
||||
|
||||
def resolve(name)
|
||||
redis.hget(REDIS_KEY, normalize(name))
|
||||
end
|
||||
|
||||
def all_records
|
||||
redis.hgetall(REDIS_KEY)
|
||||
end
|
||||
|
||||
def hosts_file_snippet
|
||||
all_records.sort.map { |host, ip| "#{ip}\t#{host}" }.join("\n")
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def normalize(name)
|
||||
name.to_s.strip.downcase.delete_suffix(".")
|
||||
end
|
||||
|
||||
def redis
|
||||
@redis ||= Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0"))
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,152 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
module Streams
|
||||
# Provisiona / decommissiona nodi stream (lab o cloud) e aggiorna DNS + registry.
|
||||
class NodeProvisioner
|
||||
class Error < StandardError; end
|
||||
class BusyError < Error; end
|
||||
|
||||
LAB_DNS_SUFFIX = -> { ENV.fetch("STREAM_LAB_DNS_SUFFIX", "lab.mltv-stream.net") }
|
||||
CLOUD_DNS_SUFFIX = -> { ENV.fetch("STREAM_CLOUD_DNS_SUFFIX", ENV.fetch("STREAM_DNS_ZONE", "mltv-stream.net")) }
|
||||
|
||||
def initialize(cloud: nil, dns: nil)
|
||||
@cloud = cloud
|
||||
@dns = dns
|
||||
end
|
||||
|
||||
def provision_lab!(prefix: "ingest-lab")
|
||||
provision!(
|
||||
prefix: prefix,
|
||||
role: "lab",
|
||||
dns_suffix: LAB_DNS_SUFFIX.call,
|
||||
cloud: cloud_provider(ENV.fetch("STREAM_CLOUD_PROVIDER", "local_lab")),
|
||||
dns: dns_provider(ENV.fetch("STREAM_DNS_PROVIDER", "lab")),
|
||||
max: ENV.fetch("STREAM_LAB_MAX_PUBLISHERS", "2").to_i,
|
||||
use_node_hostname: ENV["STREAM_LAB_USE_NODE_HOSTNAME"] == "1"
|
||||
)
|
||||
end
|
||||
|
||||
def provision_cloud!(prefix: "ingest")
|
||||
provision!(
|
||||
prefix: prefix,
|
||||
role: "cloud",
|
||||
dns_suffix: CLOUD_DNS_SUFFIX.call,
|
||||
cloud: cloud_provider("hetzner"),
|
||||
dns: dns_provider("hetzner"),
|
||||
max: ENV.fetch("STREAM_CLOUD_MAX_PUBLISHERS", "4").to_i,
|
||||
use_node_hostname: true
|
||||
)
|
||||
end
|
||||
|
||||
def decommission!(node)
|
||||
raise Error, "Non si può decommissionare il nodo home" if node.slug == Streams::NodeRegistry::HOME_SLUG
|
||||
if node.occupying_sessions.exists?
|
||||
raise BusyError, "Nodo #{node.slug} ha ancora sessioni attive"
|
||||
end
|
||||
|
||||
node.update!(status: "draining")
|
||||
cloud = cloud_for_node(node)
|
||||
dns = dns_for_node(node)
|
||||
cloud.destroy_node(node.provider_instance_id) if node.provider_instance_id.present?
|
||||
dns.delete_a(node.hostname) if node.hostname.present?
|
||||
node.destroy!
|
||||
true
|
||||
end
|
||||
|
||||
def drain!(node)
|
||||
raise Error, "Non si può mettere in drain il nodo home" if node.slug == Streams::NodeRegistry::HOME_SLUG
|
||||
|
||||
node.update!(status: "draining")
|
||||
node
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def provision!(prefix:, role:, dns_suffix:, cloud:, dns:, max:, use_node_hostname:)
|
||||
Streams::NodeRegistry.ensure_home_from_env!
|
||||
home = StreamNode.find_by!(slug: Streams::NodeRegistry::HOME_SLUG)
|
||||
slug = next_slug(prefix)
|
||||
hostname = "#{slug}.#{dns_suffix}"
|
||||
|
||||
instance = cloud.create_node(
|
||||
name: "mltv-stream-#{slug}",
|
||||
labels: { role: "stream-node", env: role == "cloud" ? "prod" : "lab" }
|
||||
)
|
||||
ip = instance.public_ip.presence || "127.0.0.1"
|
||||
private_ip = instance.private_ip.presence || ip
|
||||
dns.upsert_a(hostname, ip)
|
||||
|
||||
simulated = instance.raw.is_a?(Hash) && (instance.raw[:simulated] || instance.raw["simulated"])
|
||||
api_base = simulated ? home.api_base_url : "http://#{private_ip}:9997"
|
||||
internal_rtmp = simulated ? home.internal_rtmp_url : "rtmp://#{private_ip}:1935"
|
||||
internal_hls = simulated ? home.internal_hls_url : "http://#{private_ip}:8888"
|
||||
|
||||
StreamNode.create!(
|
||||
slug: slug,
|
||||
hostname: hostname,
|
||||
role: role,
|
||||
status: "ready",
|
||||
provider: provider_name_for(cloud, role: role),
|
||||
provider_instance_id: instance.id,
|
||||
rtmp_base_url: use_node_hostname ? "rtmp://#{hostname}:1935" : home.rtmp_base_url,
|
||||
hls_base_url: use_node_hostname ? "https://#{hostname}/hls" : home.hls_base_url,
|
||||
api_base_url: api_base,
|
||||
internal_rtmp_url: internal_rtmp,
|
||||
internal_hls_url: internal_hls,
|
||||
max_publishers: max,
|
||||
max_relays: max,
|
||||
last_health_at: Time.current,
|
||||
metadata: {
|
||||
"public_ip" => ip,
|
||||
"private_ip" => private_ip,
|
||||
"simulated" => simulated,
|
||||
"cloud_raw" => instance.raw
|
||||
}
|
||||
)
|
||||
end
|
||||
|
||||
def next_slug(prefix)
|
||||
used = StreamNode.where("slug LIKE ?", "#{prefix}-%").pluck(:slug)
|
||||
n = 1
|
||||
loop do
|
||||
candidate = format("%s-%02d", prefix, n)
|
||||
return candidate unless used.include?(candidate)
|
||||
|
||||
n += 1
|
||||
end
|
||||
end
|
||||
|
||||
def cloud_provider(name)
|
||||
@cloud || Streams::CloudProviders.build(name)
|
||||
end
|
||||
|
||||
def dns_provider(name)
|
||||
@dns || Streams::DnsProviders.build(name)
|
||||
end
|
||||
|
||||
def provider_name_for(cloud, role: nil)
|
||||
return "hetzner" if role.to_s == "cloud"
|
||||
|
||||
case cloud
|
||||
when Streams::CloudProviders::ProxmoxLab then "proxmox_lab"
|
||||
when Streams::CloudProviders::Hetzner then "hetzner"
|
||||
else "local"
|
||||
end
|
||||
end
|
||||
|
||||
def cloud_for_node(node)
|
||||
case node.provider
|
||||
when "hetzner" then Streams::CloudProviders::Hetzner.new
|
||||
when "proxmox_lab" then Streams::CloudProviders::ProxmoxLab.new
|
||||
else Streams::CloudProviders::LocalLab.new
|
||||
end
|
||||
end
|
||||
|
||||
def dns_for_node(node)
|
||||
case node.provider
|
||||
when "hetzner" then Streams::DnsProviders::Hetzner.new
|
||||
else Streams::DnsProviders::Lab.new
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,57 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
module Streams
|
||||
# Assegna un StreamNode a una nuova sessione (least-loaded tra i ready).
|
||||
# Garantisce il nodo "home" derivato dagli ENV MediaMTX attuali.
|
||||
class NodeRegistry
|
||||
class NoCapacityError < StandardError; end
|
||||
|
||||
HOME_SLUG = "home"
|
||||
|
||||
class << self
|
||||
def ensure_home_from_env!
|
||||
StreamNode.find_or_initialize_by(slug: HOME_SLUG).tap do |node|
|
||||
node.assign_attributes(
|
||||
hostname: home_hostname,
|
||||
role: "home",
|
||||
status: "ready",
|
||||
provider: "local",
|
||||
rtmp_base_url: MatchLiveTv.mediamtx_rtmp_url,
|
||||
hls_base_url: MatchLiveTv.hls_public_url,
|
||||
api_base_url: MatchLiveTv.mediamtx_api_url,
|
||||
internal_rtmp_url: ENV.fetch("MEDIAMTX_INTERNAL_RTMP_URL", "rtmp://mediamtx:1935"),
|
||||
internal_hls_url: ENV.fetch("MEDIAMTX_HLS_URL", "http://mediamtx:8888"),
|
||||
max_publishers: ENV.fetch("STREAM_NODE_HOME_MAX_PUBLISHERS", "6").to_i,
|
||||
max_relays: ENV.fetch("STREAM_NODE_HOME_MAX_RELAYS", "6").to_i
|
||||
)
|
||||
node.save!
|
||||
end
|
||||
end
|
||||
|
||||
def allocate!
|
||||
ensure_home_from_env!
|
||||
|
||||
node = StreamNode.ready
|
||||
.to_a
|
||||
.select(&:allocatable?)
|
||||
.min_by { |n| [n.active_publishers, n.role == "home" ? 1 : 0, n.slug] }
|
||||
|
||||
raise NoCapacityError, "Nessun nodo streaming con slot liberi" if node.nil?
|
||||
|
||||
node
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def home_hostname
|
||||
ENV["STREAM_NODE_HOME_HOSTNAME"].presence ||
|
||||
begin
|
||||
uri = URI.parse(MatchLiveTv.mediamtx_rtmp_url.sub(/\Artmps?:\/\//, "http://"))
|
||||
uri.host.presence
|
||||
rescue URI::InvalidURIError
|
||||
nil
|
||||
end || "home"
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -146,20 +146,20 @@ module Streams
|
||||
end
|
||||
|
||||
def mediamtx_intake_source(session)
|
||||
base = ENV.fetch("MEDIAMTX_INTERNAL_RTMP_URL", "rtmp://mediamtx:1935")
|
||||
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 = ENV.fetch("MEDIAMTX_HLS_URL", "http://mediamtx:8888").chomp("/")
|
||||
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.new.list_paths.find { |i| i["name"] == session.mediamtx_path_name }
|
||||
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
|
||||
|
||||
Reference in New Issue
Block a user