Sistemati path API /v1, AppArmor/auth MediaMTX sul nodo e defaults nbg1 così lo smoke Cloud è ripetibile. Co-authored-by: Cursor <cursoragent@cursor.com>
185 lines
5.3 KiB
Ruby
185 lines
5.3 KiB
Ruby
# 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
|
|
# Trailing slash obbligatorio: path assoluti tipo "/servers" altrimenti droppano /v1.
|
|
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", "cpx12"),
|
|
image: ENV.fetch("HCLOUD_IMAGE", "debian-12"),
|
|
location: ENV.fetch("HCLOUD_LOCATION", "nbg1"),
|
|
start_after_create: true,
|
|
labels: default_labels.merge(stringify_labels(labels)),
|
|
ssh_keys: [ENV.fetch("HCLOUD_SSH_KEY", "matchlivetv-stream-hetzner")],
|
|
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
|