From 318a319608b450ecdd1680b8b3388481894dc17b Mon Sep 17 00:00:00 2001
From: Emiliano Frascaro
+ <%= 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