diff --git a/backend/app/controllers/admin/stream_nodes_controller.rb b/backend/app/controllers/admin/stream_nodes_controller.rb index 3a492ba..46e0283 100644 --- a/backend/app/controllers/admin/stream_nodes_controller.rb +++ b/backend/app/controllers/admin/stream_nodes_controller.rb @@ -9,6 +9,7 @@ module Admin @cloud_provider = ENV.fetch("STREAM_CLOUD_PROVIDER", "local_lab") @lab_hosts = lab_dns_snippet @hetzner_configured = ENV["HCLOUD_TOKEN"].present? + @autoscale_metrics = Streams::Autoscaler.metrics end def create diff --git a/backend/app/jobs/streams/autoscaler_job.rb b/backend/app/jobs/streams/autoscaler_job.rb new file mode 100644 index 0000000..a0e44c4 --- /dev/null +++ b/backend/app/jobs/streams/autoscaler_job.rb @@ -0,0 +1,33 @@ +# frozen_string_literal: true + +module Streams + class AutoscalerJob + include Sidekiq::Job + + sidekiq_options retry: 1, queue: "default" + + INTERVAL_SECS = ENV.fetch("STREAM_AUTOSCALE_INTERVAL_SECS", "60").to_i + REDIS_CHAIN_KEY = "streams:autoscaler:chain" + + def self.ensure_chain + return unless redis + return if redis.get(REDIS_CHAIN_KEY) + + redis.set(REDIS_CHAIN_KEY, "1", ex: INTERVAL_SECS * 2) + perform_in(INTERVAL_SECS) + end + + def self.redis + @redis ||= Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0")) + rescue Redis::CannotConnectError + nil + end + + def perform + Streams::Autoscaler.reconcile! + ensure + self.class.redis&.set(REDIS_CHAIN_KEY, "1", ex: INTERVAL_SECS * 2) + self.class.perform_in(INTERVAL_SECS) + end + end +end diff --git a/backend/app/services/streams/autoscaler.rb b/backend/app/services/streams/autoscaler.rb new file mode 100644 index 0000000..07f6d62 --- /dev/null +++ b/backend/app/services/streams/autoscaler.rb @@ -0,0 +1,158 @@ +# frozen_string_literal: true + +module Streams + # Scale-out / warm spare / scale-in dei nodi overflow (lab o Hetzner). + # Kill-switch: STREAM_AUTOSCALE_ENABLED!=1 → no-op. + class Autoscaler + Result = Struct.new(:actions, :metrics, :skipped, :error, keyword_init: true) + + LOCK_KEY = "streams:autoscaler:lock" + + class << self + def enabled? + ENV["STREAM_AUTOSCALE_ENABLED"] == "1" + end + + def soft_free_slots + ENV.fetch("STREAM_AUTOSCALE_SOFT_FREE_SLOTS", "2").to_i + end + + def warm_spare_min + ENV.fetch("STREAM_AUTOSCALE_WARM_SPARE", "1").to_i + end + + def idle_minutes + ENV.fetch("STREAM_AUTOSCALE_IDLE_MINUTES", "30").to_i + end + + def max_overflow_nodes + ENV.fetch("STREAM_AUTOSCALE_MAX_NODES", "5").to_i + end + + def kind + ENV.fetch("STREAM_AUTOSCALE_KIND", "lab") # lab|cloud + end + + def reconcile!(provisioner: nil) + return Result.new(skipped: true, actions: [], metrics: metrics) unless enabled? + + redis = Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0")) + unless redis.set(LOCK_KEY, worker_id, nx: true, ex: 55) + return Result.new(skipped: true, actions: [], metrics: metrics, error: "locked") + end + + begin + new(provisioner: provisioner).reconcile! + ensure + redis.del(LOCK_KEY) + end + end + + def metrics + NodeRegistry.ensure_home_from_env! + nodes = StreamNode.ready.to_a + overflow = StreamNode.where.not(slug: NodeRegistry::HOME_SLUG) + .where(status: %w[ready draining provisioning]).to_a + { + free_slots: nodes.sum(&:free_slots), + ready_nodes: nodes.size, + spare_ready: nodes.count { |n| n.slug != NodeRegistry::HOME_SLUG && n.active_publishers.zero? }, + overflow_nodes: overflow.size, + soft_free_slots: soft_free_slots, + warm_spare_min: warm_spare_min, + max_overflow_nodes: max_overflow_nodes, + enabled: enabled?, + kind: kind + } + end + + def worker_id + ENV.fetch("HOSTNAME", "autoscaler") + end + end + + def initialize(provisioner: nil) + @provisioner = provisioner || NodeProvisioner.new + end + + def reconcile! + actions = [] + m = self.class.metrics + + if need_capacity?(m) && can_provision?(m) + provision_overflow! + actions << :scale_out + m = self.class.metrics + end + + if warm_spare_desired?(m) && m[:spare_ready] < self.class.warm_spare_min && can_provision?(m) + provision_overflow! + actions << :warm_spare + m = self.class.metrics + end + + scale_in_candidates.each do |node| + next if keep_as_warm_spare?(node) + + @provisioner.decommission!(node) + actions << :"scale_in_#{node.slug}" + m = self.class.metrics + rescue NodeProvisioner::BusyError, NodeProvisioner::Error => e + Rails.logger.warn("[Streams::Autoscaler] scale-in #{node.slug}: #{e.message}") + end + + Rails.logger.info("[Streams::Autoscaler] actions=#{actions.inspect} metrics=#{m.inspect}") + Result.new(actions: actions, metrics: m, skipped: false) + end + + private + + def need_capacity?(m) + m[:free_slots] <= self.class.soft_free_slots + end + + def warm_spare_desired?(m) + return false if self.class.warm_spare_min <= 0 + + need_capacity?(m) || overflow_in_use? + end + + def overflow_in_use? + StreamNode.ready.where.not(slug: NodeRegistry::HOME_SLUG).any? { |n| n.active_publishers.positive? } + end + + def can_provision?(m) + m[:overflow_nodes] < self.class.max_overflow_nodes + end + + def provision_overflow! + case self.class.kind + when "cloud" then @provisioner.provision_cloud! + else @provisioner.provision_lab! + end + end + + def scale_in_candidates + StreamNode.where.not(slug: NodeRegistry::HOME_SLUG) + .where(status: %w[ready draining]) + .order(:created_at) + .select { |n| n.active_publishers.zero? && idle_long_enough?(n) } + end + + def idle_long_enough?(node) + idle_since(node) <= self.class.idle_minutes.minutes.ago + end + + def idle_since(node) + last_end = node.stream_sessions.where(status: %w[ended error]).maximum(:ended_at) + last_end || node.created_at + end + + def keep_as_warm_spare?(node) + return false unless warm_spare_desired?(self.class.metrics) + + spares = StreamNode.ready.where.not(slug: NodeRegistry::HOME_SLUG).select { |n| n.active_publishers.zero? } + spares.size <= self.class.warm_spare_min && spares.map(&:id).include?(node.id) + end + end +end diff --git a/backend/app/views/admin/stream_nodes/index.html.erb b/backend/app/views/admin/stream_nodes/index.html.erb index dd8f2a2..d45b759 100644 --- a/backend/app/views/admin/stream_nodes/index.html.erb +++ b/backend/app/views/admin/stream_nodes/index.html.erb @@ -4,6 +4,20 @@ <%= t("admin.stream_nodes.providers", cloud: @cloud_provider, dns: @dns_provider) %>

+

+ <%= t( + "admin.stream_nodes.autoscale", + enabled: (@autoscale_metrics[:enabled] ? "ON" : "OFF"), + free: @autoscale_metrics[:free_slots], + soft: @autoscale_metrics[:soft_free_slots], + spare: @autoscale_metrics[:spare_ready], + warm: @autoscale_metrics[:warm_spare_min], + overflow: @autoscale_metrics[:overflow_nodes], + max: @autoscale_metrics[:max_overflow_nodes], + kind: @autoscale_metrics[:kind] + ) %> +

+

<%= button_to t("admin.stream_nodes.provision_lab"), admin_stream_nodes_path, method: :post, params: { kind: "lab" }, class: "admin-btn" %> <% if @hetzner_configured %> diff --git a/backend/config/initializers/sidekiq.rb b/backend/config/initializers/sidekiq.rb index 75e9ccf..307d522 100644 --- a/backend/config/initializers/sidekiq.rb +++ b/backend/config/initializers/sidekiq.rb @@ -6,6 +6,7 @@ Sidekiq.configure_server do |config| config.on(:startup) do StreamPublisherSyncJob.ensure_chain Ops::HealthMonitorJob.ensure_chain + Streams::AutoscalerJob.ensure_chain end end diff --git a/backend/config/locales/admin.en.yml b/backend/config/locales/admin.en.yml index 4c1da92..ded5346 100644 --- a/backend/config/locales/admin.en.yml +++ b/backend/config/locales/admin.en.yml @@ -118,6 +118,7 @@ en: stream_nodes: title: Streaming nodes providers: "Cloud provider: %{cloud} · DNS provider: %{dns}" + autoscale: "Autoscaler %{enabled} · free_slots=%{free} (soft≤%{soft}) · spare=%{spare}/%{warm} · overflow=%{overflow}/%{max} · kind=%{kind}" provision_lab: Provision lab node provision_cloud: Provision Hetzner node provision_cloud_confirm: "Create a Hetzner Cloud server + DNS record on mltv-stream.net? Billing applies until destroyed." diff --git a/backend/config/locales/admin.it.yml b/backend/config/locales/admin.it.yml index 37fc6ed..1b03118 100644 --- a/backend/config/locales/admin.it.yml +++ b/backend/config/locales/admin.it.yml @@ -118,6 +118,7 @@ it: stream_nodes: title: Nodi streaming providers: "Cloud provider: %{cloud} · DNS provider: %{dns}" + autoscale: "Autoscaler %{enabled} · free_slots=%{free} (soft≤%{soft}) · spare=%{spare}/%{warm} · overflow=%{overflow}/%{max} · kind=%{kind}" provision_lab: Provisiona nodo lab provision_cloud: Provisiona nodo Hetzner provision_cloud_confirm: "Creare un server Hetzner Cloud + record DNS su mltv-stream.net? Verrà addebitato fino allo spegnimento." diff --git a/backend/spec/services/streams/autoscaler_spec.rb b/backend/spec/services/streams/autoscaler_spec.rb new file mode 100644 index 0000000..c795652 --- /dev/null +++ b/backend/spec/services/streams/autoscaler_spec.rb @@ -0,0 +1,125 @@ +# frozen_string_literal: true + +require "rails_helper" + +RSpec.describe Streams::Autoscaler do + def with_env(vars) + previous = vars.keys.index_with { |k| ENV[k] } + vars.each { |k, v| ENV[k] = v } + yield + ensure + previous.each { |k, v| v.nil? ? ENV.delete(k) : ENV[k] = v } + end + + let(:redis) { Redis.new(url: ENV.fetch("REDIS_URL", "redis://redis:6379/0")) } + + before do + redis.del(Streams::Autoscaler::LOCK_KEY) + redis.del(Streams::DnsProviders::Lab::REDIS_KEY) + end + + it "is a no-op when disabled" do + with_env("STREAM_AUTOSCALE_ENABLED" => "0") do + result = described_class.reconcile! + expect(result.skipped).to eq(true) + expect(result.actions).to eq([]) + end + end + + it "scales out and creates a warm spare when free slots are low" do + with_env( + "STREAM_AUTOSCALE_ENABLED" => "1", + "STREAM_AUTOSCALE_SOFT_FREE_SLOTS" => "2", + "STREAM_AUTOSCALE_WARM_SPARE" => "1", + "STREAM_AUTOSCALE_KIND" => "lab", + "STREAM_AUTOSCALE_MAX_NODES" => "5", + "STREAM_CLOUD_PROVIDER" => "local_lab", + "STREAM_DNS_PROVIDER" => "lab", + "STREAM_NODE_HOME_MAX_PUBLISHERS" => "2", + "MEDIAMTX_API_URL" => "http://mtx-home:9997", + "MEDIAMTX_RTMP_URL" => "rtmp://home.example:1935", + "HLS_PUBLIC_URL" => "https://home.example/hls" + ) do + home = Streams::NodeRegistry.ensure_home_from_env! + user = User.create!(email: "as@example.com", name: "A", password: "Password123", role: "coach") + club = Club.create!(name: "AS", sport: "volleyball") + team = club.teams.create!(name: "T", sport: "volleyball", slug: "as-t") + match = team.matches.create!(opponent_name: "X", scheduled_at: 1.hour.from_now) + # Fill home (2/2) → free_slots=0 + 2.times do + StreamSession.create!(match: match, user: user, platform: "matchlivetv", status: "live", stream_node: home) + end + + provisioner = instance_double(Streams::NodeProvisioner) + created = [] + allow(provisioner).to receive(:provision_lab!) do + node = StreamNode.create!( + slug: "ingest-lab-#{created.size + 1}", + hostname: "h#{created.size}.lab", + role: "lab", + status: "ready", + provider: "local", + rtmp_base_url: "rtmp://h:1935", + hls_base_url: "https://h/hls", + api_base_url: "http://h:9997", + max_publishers: 2, + max_relays: 2 + ) + created << node + node + end + + result = described_class.reconcile!(provisioner: provisioner) + expect(result.actions).to include(:scale_out) + expect(created.size).to eq(1) + expect(result.metrics[:free_slots]).to be >= 2 + + # Riempi il nodo overflow → spare_ready=0 ma carico presente → warm spare + StreamSession.create!( + match: match, user: user, platform: "matchlivetv", status: "live", stream_node: created.first + ) + StreamSession.create!( + match: match, user: user, platform: "matchlivetv", status: "live", stream_node: created.first + ) + + result2 = described_class.reconcile!(provisioner: provisioner) + expect(result2.actions).to include(:warm_spare).or include(:scale_out) + expect(created.size).to be >= 2 + end + end + + it "scales in an idle overflow node when warm spare is not required" do + with_env( + "STREAM_AUTOSCALE_ENABLED" => "1", + "STREAM_AUTOSCALE_SOFT_FREE_SLOTS" => "0", + "STREAM_AUTOSCALE_WARM_SPARE" => "0", + "STREAM_AUTOSCALE_IDLE_MINUTES" => "0", + "STREAM_NODE_HOME_MAX_PUBLISHERS" => "6", + "MEDIAMTX_API_URL" => "http://mtx-home:9997", + "MEDIAMTX_RTMP_URL" => "rtmp://home.example:1935", + "HLS_PUBLIC_URL" => "https://home.example/hls" + ) do + Streams::NodeRegistry.ensure_home_from_env! + idle = StreamNode.create!( + slug: "ingest-lab-99", + hostname: "idle.lab", + role: "lab", + status: "ready", + provider: "local", + provider_instance_id: "sim-x", + rtmp_base_url: "rtmp://h:1935", + hls_base_url: "https://h/hls", + api_base_url: "http://h:9997", + max_publishers: 2, + max_relays: 2, + created_at: 2.hours.ago + ) + + provisioner = instance_double(Streams::NodeProvisioner) + expect(provisioner).to receive(:decommission!).with(idle) + + result = described_class.reconcile!(provisioner: provisioner) + expect(result.actions.map(&:to_s)).to include("scale_in_ingest-lab-99") + end + end +end diff --git a/docs/infrastructure/STREAMING_AUTOSCALE.md b/docs/infrastructure/STREAMING_AUTOSCALE.md index 37badc9..fe14507 100644 --- a/docs/infrastructure/STREAMING_AUTOSCALE.md +++ b/docs/infrastructure/STREAMING_AUTOSCALE.md @@ -210,7 +210,7 @@ Cosa il lab **non** replica al 100%: tempi boot Hetzner, vSwitch, RTMP 4G multi- | **0 — Multi-nodo ready** | Registry + URL RTMP/HLS in API (`StreamNode`, `Streams::NodeRegistry`) | **implementata sul branch** | | **1 — Hetzner Cloud + DNS** | `HetznerCloudProvider` + `HetznerDnsProvider`, `mltv-stream.net`, cloud-init, WireGuard (§16) | **implementata sul branch** (WG ops manuale) | | **2 — Routing relay** | Coda `youtube_relay`, sticky owner, cap `RELAY_MAX_CONCURRENT` | **implementata sul branch** | -| **3 — Autoscaler** | Soglie + warm spare (+1) | Automatico | +| **3 — Autoscaler** | Soglie + warm spare (+1), kill-switch `STREAM_AUTOSCALE_ENABLED` | **implementata sul branch** | | **4 — Hardening** | Drain, budget, runbook, kill-switch | Picchi weekend | | **A — Auction (futuro)** | Migrazione control plane + vSwitch al posto di WireGuard | Hardware dedicato Hetzner | | **B2 (opz.)** | Switch object storage | Quando misurato | @@ -250,6 +250,16 @@ Cosa il lab **non** replica al 100%: tempi boot Hetzner, vSwitch, RTMP 4G multi- - Cap per worker: `RELAY_MAX_CONCURRENT` (default 4) + set Redis `youtube_relay:owned:HOSTNAME` - Ensure a capacità piena → requeue 5s invece di avviare un secondo ffmpeg locale +### Fase 3 — dettagli implementati + +- `Streams::Autoscaler` + job chain `Streams::AutoscalerJob` (come health monitor) +- Kill-switch: `STREAM_AUTOSCALE_ENABLED=1` per attivare +- Scale-out se `free_slots ≤ STREAM_AUTOSCALE_SOFT_FREE_SLOTS` +- Warm spare (+N) se carico/soglia e `spare_ready < WARM_SPARE` +- Scale-in overflow idle da `IDLE_MINUTES` (non tocca home; rispetta warm spare) +- Cap `STREAM_AUTOSCALE_MAX_NODES`; kind `lab|cloud` +- Metriche in admin Nodi stream + Ordine di lavoro consigliato sul branch: **0 → L → 1 → 2 → 3 → 4**; **A** quando si decide di lasciare il Proxmox casa. --- diff --git a/infra/.env.production.example b/infra/.env.production.example index 1198f0e..bc61072 100644 --- a/infra/.env.production.example +++ b/infra/.env.production.example @@ -123,3 +123,12 @@ STREAM_NODE_ENV=prod STREAM_CLOUD_PROVIDER=local_lab RELAY_MAX_CONCURRENT=4 YOUTUBE_RELAY_WORKER=1 + +# Autoscaler (kill-switch: 0 finché WireGuard/smoke Cloud non sono OK) +STREAM_AUTOSCALE_ENABLED=0 +STREAM_AUTOSCALE_KIND=lab +STREAM_AUTOSCALE_SOFT_FREE_SLOTS=2 +STREAM_AUTOSCALE_WARM_SPARE=1 +STREAM_AUTOSCALE_IDLE_MINUTES=30 +STREAM_AUTOSCALE_MAX_NODES=5 +STREAM_AUTOSCALE_INTERVAL_SECS=60