Compare commits
43
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f02782d0f0 | ||
|
|
fb82b74e92 | ||
|
|
a07f231417 | ||
|
|
43d5003de8 | ||
|
|
eadfc29da3 | ||
|
|
50d3e2cf16 | ||
|
|
55e71ca838 | ||
|
|
5b723bf844 | ||
|
|
e502844b88 | ||
|
|
0702cb699c | ||
|
|
e09b3ffcf8 | ||
|
|
41d4235258 | ||
|
|
a51d5e8da5 | ||
|
|
5bcb4170c7 | ||
|
|
e436f5ece4 | ||
|
|
1e4e0560f1 | ||
|
|
841a6a15be | ||
|
|
dd94bf7e66 | ||
|
|
baa6283a15 | ||
|
|
bfb3115a64 | ||
|
|
db8b14ba4e | ||
|
|
46244fe066 | ||
|
|
3ed70b6999 | ||
|
|
b05b1acd7a | ||
|
|
c777853c97 | ||
|
|
9225318e1a | ||
|
|
061b7fad8b | ||
|
|
be7957c078 | ||
|
|
91667e1bfd | ||
|
|
1fecccaafd | ||
|
|
d100c8030f | ||
|
|
f1906517a9 | ||
|
|
00fde25c50 | ||
|
|
4c5af99efa | ||
|
|
8bf12e7721 | ||
|
|
5a0ca8ef1c | ||
|
|
714602f3da | ||
|
|
baac3a512b | ||
|
|
559284f0b2 | ||
|
|
e5fc925bae | ||
|
|
318a319608 | ||
|
|
0b55d5a8fa | ||
|
|
c3878cdc6d |
@@ -55,6 +55,10 @@ class SessionChannel < ApplicationCable::Channel
|
||||
Sessions::Pause.new(session).call
|
||||
when "resume_stream"
|
||||
Sessions::Resume.new(session).call
|
||||
when "mute_audio"
|
||||
Sessions::SetAudioMute.new(session, muted: true).call
|
||||
when "unmute_audio"
|
||||
Sessions::SetAudioMute.new(session, muted: false).call
|
||||
when "set_min_quality"
|
||||
Sessions::SetMinQuality.new(session, preset: data["min_quality_preset"]).call
|
||||
end
|
||||
|
||||
@@ -0,0 +1,76 @@
|
||||
module Admin
|
||||
class AnnouncementsController < Admin::BaseController
|
||||
before_action :set_announcement, only: %i[edit update destroy]
|
||||
|
||||
def index
|
||||
@announcements = AppAnnouncement.newest
|
||||
end
|
||||
|
||||
def new
|
||||
@announcement = AppAnnouncement.new(
|
||||
kind: "info",
|
||||
severity: "info",
|
||||
status: "draft",
|
||||
show_on_android: true,
|
||||
show_on_ios: true,
|
||||
show_on_web_public: true,
|
||||
show_on_web_private: true,
|
||||
dismissible: true
|
||||
)
|
||||
end
|
||||
|
||||
def create
|
||||
@announcement = AppAnnouncement.new(announcement_params)
|
||||
@announcement.created_by_admin = current_admin_account
|
||||
if @announcement.save
|
||||
redirect_to admin_announcements_path, notice: t("admin.flash.announcement_created")
|
||||
else
|
||||
flash.now[:alert] = @announcement.errors.full_messages.join(", ")
|
||||
render :new, status: :unprocessable_entity
|
||||
end
|
||||
end
|
||||
|
||||
def edit; end
|
||||
|
||||
def update
|
||||
if @announcement.update(announcement_params)
|
||||
redirect_to admin_announcements_path, notice: t("admin.flash.announcement_updated")
|
||||
else
|
||||
flash.now[:alert] = @announcement.errors.full_messages.join(", ")
|
||||
render :edit, status: :unprocessable_entity
|
||||
end
|
||||
end
|
||||
|
||||
def destroy
|
||||
@announcement.destroy!
|
||||
redirect_to admin_announcements_path, notice: t("admin.flash.announcement_destroyed")
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def set_announcement
|
||||
@announcement = AppAnnouncement.find(params[:id])
|
||||
end
|
||||
|
||||
def announcement_params
|
||||
params.require(:app_announcement).permit(
|
||||
:kind, :severity, :status, :title, :body,
|
||||
:action_url, :action_label, :dismissible, :starts_at, :ends_at,
|
||||
:show_on_android, :show_on_ios, :show_on_web_public, :show_on_web_private,
|
||||
translations: {
|
||||
en: AppAnnouncement::COPY_KEYS,
|
||||
fr: AppAnnouncement::COPY_KEYS,
|
||||
de: AppAnnouncement::COPY_KEYS,
|
||||
es: AppAnnouncement::COPY_KEYS
|
||||
}
|
||||
).tap do |permitted|
|
||||
%i[dismissible show_on_android show_on_ios show_on_web_public show_on_web_private].each do |key|
|
||||
permitted[key] = ActiveModel::Type::Boolean.new.cast(permitted[key])
|
||||
end
|
||||
%i[starts_at ends_at action_url action_label].each do |key|
|
||||
permitted[key] = nil if permitted[key].blank?
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -2,6 +2,7 @@ module Admin
|
||||
class BillingController < BaseController
|
||||
def index
|
||||
@pending_payments = pending_scope.recent
|
||||
@pending_transfers = transfer_scope
|
||||
@completed_payments = Billing::Payment.with_invoice_pdf.includes(:club, :invoice).recent.limit(40)
|
||||
@filter_club = Club.find_by(id: params[:club_id]) if params[:club_id].present?
|
||||
@clubs = Club.order(:name)
|
||||
@@ -18,6 +19,25 @@ module Admin
|
||||
alert: e.message
|
||||
end
|
||||
|
||||
def confirm_transfer
|
||||
order = Billing::TransferOrder.find(params[:id])
|
||||
Billing::ConfirmBankTransfer.call(
|
||||
order: order,
|
||||
admin: current_admin_account,
|
||||
pdf: params[:pdf]
|
||||
)
|
||||
redirect_to admin_billing_path(club_id: order.club_id),
|
||||
notice: t("admin.flash.transfer_confirmed", plan: order.plan.name, club: order.club.name)
|
||||
rescue Billing::ConfirmBankTransfer::Error, Billing::AttachPaymentInvoice::Error, Billing::IssueInvoice::Error => e
|
||||
redirect_to admin_billing_path, alert: e.message
|
||||
end
|
||||
|
||||
def cancel_transfer
|
||||
order = Billing::TransferOrder.find(params[:id])
|
||||
order.cancel!
|
||||
redirect_to admin_billing_path(club_id: order.club_id), notice: t("admin.flash.transfer_cancelled")
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def pending_scope
|
||||
@@ -26,6 +46,12 @@ module Admin
|
||||
scope
|
||||
end
|
||||
|
||||
def transfer_scope
|
||||
scope = Billing::TransferOrder.awaiting_payment.includes(:club, :requested_by_user)
|
||||
scope = scope.where(club_id: params[:club_id]) if params[:club_id].present?
|
||||
scope
|
||||
end
|
||||
|
||||
def billing_redirect_params(payment)
|
||||
{ club_id: payment.club_id, anchor: "payment-#{payment.id}" }.compact
|
||||
end
|
||||
|
||||
@@ -1,9 +1,9 @@
|
||||
module Admin
|
||||
class ClubsController < BaseController
|
||||
before_action :set_club, only: %i[show grant_comped revoke_comped]
|
||||
before_action :set_club, only: %i[show grant_comped revoke_comped set_quote revoke_quote]
|
||||
|
||||
def index
|
||||
@clubs = Club.includes(:teams, subscription: %i[plan admin_comped_by])
|
||||
@clubs = Club.includes(:teams, :billing_quote, subscription: %i[plan admin_comped_by])
|
||||
.order(:name)
|
||||
end
|
||||
|
||||
@@ -11,6 +11,7 @@ module Admin
|
||||
@subscription = @club.subscription || @club.build_subscription(plan: Plan["free"], status: "active")
|
||||
@plans = Plan.ordered.reject { |p| p.slug == "free" }
|
||||
@teams = @club.teams.order(:name)
|
||||
@quote = @club.active_billing_quote
|
||||
end
|
||||
|
||||
def grant_comped
|
||||
@@ -32,6 +33,32 @@ module Admin
|
||||
redirect_back_or_club alert: e.message
|
||||
end
|
||||
|
||||
def set_quote
|
||||
quote = Billing::SetClubQuote.upsert(
|
||||
club: @club,
|
||||
plan_slug: params.require(:plan_slug),
|
||||
interval: params[:interval],
|
||||
amount_euros: params[:amount_euros],
|
||||
note: params[:note],
|
||||
admin: current_admin_account
|
||||
)
|
||||
redirect_back_or_club notice: t(
|
||||
"admin.flash.quote_saved",
|
||||
club: @club.name,
|
||||
plan: quote.plan.name,
|
||||
amount: quote.formatted_amount
|
||||
)
|
||||
rescue Billing::SetClubQuote::Error, ActionController::ParameterMissing => e
|
||||
redirect_back_or_club alert: e.message
|
||||
end
|
||||
|
||||
def revoke_quote
|
||||
Billing::SetClubQuote.revoke(club: @club, admin: current_admin_account)
|
||||
redirect_back_or_club notice: t("admin.flash.quote_revoked", club: @club.name)
|
||||
rescue Billing::SetClubQuote::Error => e
|
||||
redirect_back_or_club alert: e.message
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def set_club
|
||||
|
||||
@@ -5,7 +5,7 @@ module Admin
|
||||
@host = HostMetrics.new.sample!
|
||||
@active_sessions = StreamSession
|
||||
.where(status: DashboardStats::ACTIVE_STATUSES)
|
||||
.includes(match: :team)
|
||||
.includes(:stream_node, match: :team)
|
||||
.order(started_at: :desc)
|
||||
@teams = Team.includes(:matches).order(:name).limit(8)
|
||||
end
|
||||
|
||||
@@ -1,11 +1,11 @@
|
||||
module Admin
|
||||
class SessionsController < Admin::BaseController
|
||||
def index
|
||||
@sessions = StreamSession.includes(:match, :user).order(created_at: :desc).limit(50)
|
||||
@sessions = StreamSession.includes(:stream_node, :user, match: :team).order(created_at: :desc).limit(50)
|
||||
end
|
||||
|
||||
def show
|
||||
@session = StreamSession.find(params[:id])
|
||||
@session = StreamSession.includes(:stream_node, match: :team).find(params[:id])
|
||||
@events = @session.stream_events.recent.limit(50)
|
||||
end
|
||||
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
module Admin
|
||||
class StreamNodesController < Admin::BaseController
|
||||
def index
|
||||
Streams::NodeRegistry.ensure_home_from_env!
|
||||
@nodes = StreamNode.order(:role, :slug)
|
||||
@dns_provider = ENV.fetch("STREAM_DNS_PROVIDER", "lab")
|
||||
@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
|
||||
kind = params[:kind].to_s
|
||||
node =
|
||||
if kind == "cloud"
|
||||
Streams::NodeProvisioner.new.provision_cloud!
|
||||
else
|
||||
Streams::NodeProvisioner.new.provision_lab!
|
||||
end
|
||||
redirect_to admin_stream_nodes_path,
|
||||
notice: t("admin.flash.stream_node_created", slug: node.slug)
|
||||
rescue Streams::NodeProvisioner::Error, Streams::CloudProviders::Error, Streams::DnsProviders::Error,
|
||||
KeyError => e
|
||||
redirect_to admin_stream_nodes_path, alert: e.message
|
||||
end
|
||||
|
||||
def drain
|
||||
node = StreamNode.find(params[:id])
|
||||
Streams::NodeProvisioner.new.drain!(node)
|
||||
redirect_to admin_stream_nodes_path, notice: t("admin.flash.stream_node_draining", slug: node.slug)
|
||||
rescue Streams::NodeProvisioner::Error => e
|
||||
redirect_to admin_stream_nodes_path, alert: e.message
|
||||
end
|
||||
|
||||
def destroy
|
||||
node = StreamNode.find(params[:id])
|
||||
Streams::NodeProvisioner.new.decommission!(node)
|
||||
redirect_to admin_stream_nodes_path, notice: t("admin.flash.stream_node_destroyed", slug: node.slug)
|
||||
rescue Streams::NodeProvisioner::BusyError, Streams::NodeProvisioner::Error,
|
||||
Streams::CloudProviders::Error, Streams::DnsProviders::Error => e
|
||||
redirect_to admin_stream_nodes_path, alert: e.message
|
||||
end
|
||||
|
||||
def kill_switch
|
||||
Streams::Autoscaler.engage_kill_switch!
|
||||
redirect_to admin_stream_nodes_path, notice: t("admin.flash.autoscale_kill_on")
|
||||
end
|
||||
|
||||
def clear_kill_switch
|
||||
Streams::Autoscaler.clear_kill_switch!
|
||||
redirect_to admin_stream_nodes_path, notice: t("admin.flash.autoscale_kill_off")
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def lab_dns_snippet
|
||||
return "" unless @dns_provider == "lab"
|
||||
|
||||
Streams::DnsProviders::Lab.new.hosts_file_snippet
|
||||
rescue Redis::BaseError
|
||||
""
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,13 @@
|
||||
module Api
|
||||
module V1
|
||||
class AnnouncementsController < BaseController
|
||||
skip_before_action :authenticate_request!
|
||||
|
||||
def index
|
||||
platform = params[:platform].to_s.presence
|
||||
announcements = AppAnnouncement.for_platform(platform).newest
|
||||
render json: announcements.map { |item| item.as_api_json(platform: platform) }
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -2,6 +2,7 @@ module Api
|
||||
module V1
|
||||
class BaseController < ApplicationController
|
||||
rescue_from Teams::EntitlementError, with: :render_entitlement_error
|
||||
rescue_from BrandingAttachments::CoverUploadError, with: :render_cover_upload_error
|
||||
rescue_from Youtube::BroadcastService::Error, with: :render_youtube_error
|
||||
|
||||
private
|
||||
@@ -20,6 +21,10 @@ module Api
|
||||
billing_url: error.billing_url
|
||||
}, status: :forbidden
|
||||
end
|
||||
|
||||
def render_cover_upload_error(error)
|
||||
render json: { error: error.message, error_code: "cover_upload_invalid" }, status: :unprocessable_entity
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -2,6 +2,7 @@ module Api
|
||||
module V1
|
||||
class MatchesController < BaseController
|
||||
include Api::UrlHelper
|
||||
include Api::CoverJson
|
||||
include BrandingAttachments
|
||||
|
||||
before_action :set_team, only: %i[index create]
|
||||
@@ -21,6 +22,7 @@ module Api
|
||||
normalize_scoring_rules!(attrs)
|
||||
match = @team.matches.create!(attrs)
|
||||
attach_opponent_logo(match)
|
||||
attach_match_cover(match)
|
||||
render json: match_json(match), status: :created
|
||||
end
|
||||
|
||||
@@ -34,6 +36,7 @@ module Api
|
||||
normalize_opponent_color!(attrs)
|
||||
@match.update!(attrs)
|
||||
attach_opponent_logo(@match)
|
||||
attach_match_cover(@match)
|
||||
render json: match_json(@match)
|
||||
end
|
||||
|
||||
@@ -98,6 +101,11 @@ module Api
|
||||
match.opponent_logo_file.attach(file) if file.present?
|
||||
end
|
||||
|
||||
def attach_match_cover(match)
|
||||
attach_cover_image(match, entitlements: match.team.entitlements)
|
||||
enqueue_cover_slate_generation(match)
|
||||
end
|
||||
|
||||
def match_json(match, detail: false)
|
||||
active = active_session_for(match)
|
||||
team = match.team
|
||||
@@ -123,6 +131,7 @@ module Api
|
||||
home_logo_url: api_absolute_url(team.effective_logo_url),
|
||||
opponent_primary_color: match.effective_opponent_primary_color,
|
||||
opponent_logo_url: api_absolute_url(match.opponent_logo_url),
|
||||
**match_cover_json(match),
|
||||
active_session_id: active&.id,
|
||||
active_session_status: active&.status,
|
||||
stream_completed: match.stream_completed?,
|
||||
|
||||
@@ -34,6 +34,12 @@ module Api
|
||||
render json: session_json(@session)
|
||||
end
|
||||
|
||||
def audio_mute
|
||||
muted = params.key?(:muted) ? params[:muted] : true
|
||||
Sessions::SetAudioMute.new(@session, muted: muted).call
|
||||
render json: session_json(@session)
|
||||
end
|
||||
|
||||
def score
|
||||
Scoring::SyncState.new(session: @session, payload: score_sync_params).call
|
||||
render json: session_json(@session, detail: true)
|
||||
@@ -193,6 +199,7 @@ module Api
|
||||
platform: session.platform,
|
||||
rtmp_ingest_url: session.rtmp_ingest_url,
|
||||
hls_playback_url: session.hls_playback_url,
|
||||
stream_node: session.stream_node&.slug,
|
||||
watch_page_url: session.matchlivetv_platform? ? session.watch_page_url : nil,
|
||||
share_url: session.share_url,
|
||||
youtube_watch_url: session.youtube_watch_url,
|
||||
@@ -203,7 +210,8 @@ module Api
|
||||
min_quality_preset: session.min_quality_preset,
|
||||
target_bitrate: session.target_bitrate,
|
||||
started_at: session.started_at,
|
||||
disconnection_count: session.disconnection_count
|
||||
disconnection_count: session.disconnection_count,
|
||||
audio_muted: session.audio_muted
|
||||
}
|
||||
if detail
|
||||
data[:score] = session.score_state&.as_cable_payload
|
||||
|
||||
@@ -137,7 +137,8 @@ module Api
|
||||
youtube_mode: ent.plan.youtube_mode,
|
||||
staff_role: current_user.staff_role_for(team),
|
||||
can_stream: current_user.can_stream_for?(team),
|
||||
can_connect_youtube: team.club.owned_by?(current_user) && ent.premium_full? && ent.youtube_enabled?
|
||||
can_connect_youtube: team.club.owned_by?(current_user) && ent.premium_full? && ent.youtube_enabled?,
|
||||
custom_cover_enabled: ent.can_use_custom_cover?
|
||||
}
|
||||
if detail
|
||||
data[:members] = team.user_teams.includes(:user).where.not(role: "owner").map do |ut|
|
||||
|
||||
@@ -0,0 +1,16 @@
|
||||
module Api
|
||||
module CoverJson
|
||||
extend ActiveSupport::Concern
|
||||
|
||||
private
|
||||
|
||||
def match_cover_json(match)
|
||||
resolver = Streams::CoverResolver.new(match)
|
||||
{
|
||||
effective_cover_url: api_absolute_url(resolver.effective_cover_url),
|
||||
cover_source: resolver.cover_source,
|
||||
custom_cover_enabled: match.team.entitlements.can_use_custom_cover?
|
||||
}
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -8,6 +8,48 @@ module BrandingAttachments
|
||||
record.logo_file.attach(file) if file.present?
|
||||
end
|
||||
|
||||
def attach_cover_image(record, entitlements:)
|
||||
remove_flag = params[:remove_cover].to_s == "1" || params.dig(:match, :remove_cover).to_s == "1"
|
||||
if remove_flag
|
||||
record.cover_image.purge if record.cover_image.attached?
|
||||
record.cover_slate.purge if record.cover_slate.attached?
|
||||
return
|
||||
end
|
||||
|
||||
file = params.dig(:branding, :cover_image) || params[:cover_image]
|
||||
return if file.blank?
|
||||
|
||||
entitlements.assert_can_upload_cover!
|
||||
validate_cover_file!(file)
|
||||
record.cover_slate.purge if record.cover_slate.attached?
|
||||
record.cover_image.attach(file)
|
||||
end
|
||||
|
||||
def validate_cover_file!(file)
|
||||
allowed_types = Coverable::COVER_IMAGE_TYPES
|
||||
|
||||
# Rileva MIME dal contenuto (magic bytes) ignorando il declared_type del client.
|
||||
io = file.respond_to?(:tempfile) ? file.tempfile : file
|
||||
io.rewind if io.respond_to?(:rewind)
|
||||
detected = Marcel::MimeType.for(io, name: file.respond_to?(:original_filename) ? file.original_filename : nil)
|
||||
io.rewind if io.respond_to?(:rewind)
|
||||
|
||||
unless detected.in?(allowed_types)
|
||||
raise CoverUploadError, I18n.t("coverable.errors.invalid_type")
|
||||
end
|
||||
|
||||
byte_size = file.respond_to?(:size) ? file.size : File.size(file.path)
|
||||
if byte_size > Coverable::MAX_COVER_BYTES
|
||||
raise CoverUploadError, I18n.t("coverable.errors.too_large")
|
||||
end
|
||||
end
|
||||
|
||||
class CoverUploadError < StandardError; end
|
||||
|
||||
def enqueue_cover_slate_generation(record)
|
||||
GenerateCoverSlateJob.enqueue_for(record) if record.cover_image.attached?
|
||||
end
|
||||
|
||||
def normalize_hex_color(value, fallback)
|
||||
v = value.to_s.strip
|
||||
v = "##{v}" if v.match?(/\A[0-9A-Fa-f]{6}\z/)
|
||||
|
||||
@@ -4,11 +4,20 @@ module MediamtxPlayback
|
||||
private
|
||||
|
||||
def mediamtx_paths_index
|
||||
@mediamtx_paths_index ||= Mediamtx::Client.new.list_paths.index_by { |i| i["name"] }
|
||||
mediamtx_paths_index_for_url(MatchLiveTv.mediamtx_api_url)
|
||||
end
|
||||
|
||||
def mediamtx_paths_index_for(session)
|
||||
mediamtx_paths_index_for_url(session.mediamtx_api_base_url)
|
||||
end
|
||||
|
||||
def mediamtx_paths_index_for_url(api_url)
|
||||
@mediamtx_paths_by_origin ||= {}
|
||||
@mediamtx_paths_by_origin[api_url] ||= Mediamtx::Client.new(base_url: api_url).list_paths.index_by { |i| i["name"] }
|
||||
end
|
||||
|
||||
def mediamtx_path_info(session, path_name: mediamtx_playback_path_name(session))
|
||||
mediamtx_paths_index[path_name]
|
||||
mediamtx_paths_index_for(session)[path_name]
|
||||
end
|
||||
|
||||
def mediamtx_playback_path_name(session)
|
||||
@@ -28,7 +37,7 @@ module MediamtxPlayback
|
||||
end
|
||||
|
||||
def mediamtx_publisher_online?(session)
|
||||
info = mediamtx_paths_index[session.mediamtx_path_name]
|
||||
info = mediamtx_paths_index_for(session)[session.mediamtx_path_name]
|
||||
info && info["online"]
|
||||
end
|
||||
end
|
||||
|
||||
@@ -20,7 +20,8 @@ class HlsProxyController < ActionController::Base
|
||||
def proxy_mediamtx_path(upstream_path)
|
||||
return head :not_found if upstream_path.blank?
|
||||
|
||||
upstream = "#{MatchLiveTv.mediamtx_hls_url}/#{upstream_path}"
|
||||
origin = hls_origin_for(upstream_path)
|
||||
upstream = "#{origin}/#{upstream_path}"
|
||||
upstream = "#{upstream}?#{request.query_string}" if request.query_string.present?
|
||||
cookie = request.headers["Cookie"].presence || "cookieCheck=1"
|
||||
|
||||
@@ -30,7 +31,7 @@ class HlsProxyController < ActionController::Base
|
||||
if response.status.in?([301, 302, 307, 308])
|
||||
location = response.headers["location"].to_s
|
||||
cookie = cookie_from_set_header(response.headers["set-cookie"]).presence || cookie
|
||||
upstream = resolve_upstream_url(location)
|
||||
upstream = resolve_upstream_url(location, origin)
|
||||
next
|
||||
end
|
||||
|
||||
@@ -43,10 +44,21 @@ class HlsProxyController < ActionController::Base
|
||||
head :bad_gateway
|
||||
end
|
||||
|
||||
def resolve_upstream_url(location)
|
||||
def hls_origin_for(upstream_path)
|
||||
session_id = upstream_path[/\Alive\/match_([0-9a-f-]{36})/i, 1]
|
||||
if session_id
|
||||
session = StreamSession.find_by(id: session_id)
|
||||
node_hls = session&.stream_node&.internal_hls_url.presence
|
||||
return node_hls.sub(%r{/$}, "") if node_hls.present?
|
||||
end
|
||||
|
||||
MatchLiveTv.mediamtx_hls_url.sub(%r{/$}, "")
|
||||
end
|
||||
|
||||
def resolve_upstream_url(location, origin)
|
||||
return location if location.start_with?("http://", "https://")
|
||||
|
||||
base = MatchLiveTv.mediamtx_hls_url.sub(%r{/$}, "")
|
||||
base = origin.to_s.sub(%r{/$}, "")
|
||||
location.start_with?("/") ? "#{base}#{location}" : "#{base}/#{location}"
|
||||
end
|
||||
|
||||
|
||||
@@ -11,11 +11,15 @@ module Public
|
||||
@club.assign_attributes(billing_profile_params)
|
||||
if @club.save(context: :billing_profile)
|
||||
if premium_checkout_return_params.present? && @club.billing_profile_complete?
|
||||
redirect_to public_club_checkout_path(
|
||||
@club,
|
||||
plan: premium_checkout_return_params[:plan],
|
||||
interval: premium_checkout_return_params[:interval]
|
||||
), notice: t("flash.club_billing.profile_saved_proceed_payment")
|
||||
if @club.active_billing_quote.present?
|
||||
redirect_to public_club_billing_path(@club), notice: t("flash.club_billing.profile_updated")
|
||||
else
|
||||
redirect_to public_club_checkout_path(
|
||||
@club,
|
||||
plan: premium_checkout_return_params[:plan],
|
||||
interval: premium_checkout_return_params[:interval]
|
||||
), notice: t("flash.club_billing.profile_saved_proceed_payment")
|
||||
end
|
||||
else
|
||||
redirect_to public_club_billing_path(@club), notice: t("flash.club_billing.profile_updated")
|
||||
end
|
||||
@@ -42,6 +46,29 @@ module Public
|
||||
redirect_to public_club_billing_path(@club), alert: t("flash.club_billing.stripe_error", message: e.message)
|
||||
end
|
||||
|
||||
def request_bank_transfer
|
||||
if @club.subscription&.admin_comped?
|
||||
redirect_to public_club_billing_path(@club), alert: t("flash.clubs.comped_change_denied")
|
||||
return
|
||||
end
|
||||
|
||||
unless @club.billing_profile_complete?
|
||||
redirect_to public_club_billing_profile_path(@club, plan: params[:plan], interval: params[:interval]),
|
||||
alert: t("flash.clubs.complete_billing_first")
|
||||
return
|
||||
end
|
||||
|
||||
Billing::RequestBankTransfer.call(
|
||||
club: @club,
|
||||
user: current_user,
|
||||
plan_slug: params[:plan],
|
||||
interval: params[:interval]
|
||||
)
|
||||
redirect_to public_club_billing_path(@club), notice: t("flash.club_billing.bank_transfer_requested")
|
||||
rescue Billing::RequestBankTransfer::Error, ArgumentError => e
|
||||
redirect_to public_club_billing_path(@club), alert: e.message
|
||||
end
|
||||
|
||||
def download_invoice
|
||||
invoice = @club.billing_invoices.find(params[:invoice_id])
|
||||
unless invoice.pdf.attached?
|
||||
|
||||
@@ -23,23 +23,28 @@ module Public
|
||||
ClubMembership.create!(user: current_user, club: club, role: "owner")
|
||||
|
||||
first_team_name = params.dig(:first_team, :name).presence || t("club.new.default_first_team_name")
|
||||
club.teams.create!(
|
||||
first_team = club.teams.create!(
|
||||
name: first_team_name,
|
||||
sport: club.sport
|
||||
)
|
||||
|
||||
plan = params[:plan].presence_in(%w[free premium_light premium_full]) || "free"
|
||||
Billing::AssignPlan.call(club: club, plan_slug: plan)
|
||||
desired_plan = params[:plan].presence_in(%w[free premium_light premium_full]) || "free"
|
||||
Billing::AssignPlan.call(club: club, plan_slug: "free")
|
||||
Teams::StaffAssignment.designate_club_owner!(team: first_team, user: current_user)
|
||||
|
||||
if plan.in?(%w[premium_light premium_full]) && MatchLiveTv.stripe_enabled?
|
||||
unless club.billing_profile_complete?
|
||||
redirect_to public_club_billing_profile_path(club, plan: plan, interval: checkout_interval_param),
|
||||
alert: t("flash.clubs.complete_billing_first")
|
||||
return
|
||||
if desired_plan.in?(%w[premium_light premium_full])
|
||||
interval = Billing::Stripe::PriceCatalog::DEFAULT_INTERVAL
|
||||
if MatchLiveTv.stripe_enabled?
|
||||
unless club.billing_profile_complete?
|
||||
redirect_to public_club_billing_profile_path(club, plan: desired_plan, interval: interval),
|
||||
alert: t("flash.clubs.complete_billing_first")
|
||||
return
|
||||
end
|
||||
|
||||
redirect_to public_club_checkout_path(club, plan: desired_plan, interval: interval)
|
||||
else
|
||||
redirect_to public_club_billing_path(club), notice: t("flash.clubs.created_complete_subscription")
|
||||
end
|
||||
|
||||
interval = checkout_interval_param
|
||||
redirect_to public_club_checkout_path(club, plan: plan, interval: interval)
|
||||
else
|
||||
redirect_to public_club_path(club), notice: t("flash.clubs.created")
|
||||
end
|
||||
@@ -59,17 +64,26 @@ module Public
|
||||
|
||||
def edit
|
||||
require_club_owner!(@club)
|
||||
@entitlements = @club.teams.first&.entitlements
|
||||
end
|
||||
|
||||
def update
|
||||
require_club_owner!(@club)
|
||||
@entitlements = @club.teams.first!.entitlements
|
||||
@club.assign_attributes(club_params)
|
||||
attach_branding_logo(@club)
|
||||
attach_cover_image(@club, entitlements: @entitlements)
|
||||
@club.save!
|
||||
enqueue_cover_slate_generation(@club)
|
||||
redirect_to public_club_path(@club), notice: t("flash.clubs.updated")
|
||||
rescue ActiveRecord::RecordInvalid => e
|
||||
@entitlements = @club.teams.first&.entitlements
|
||||
flash.now[:alert] = e.record.errors.full_messages.join(", ")
|
||||
render :edit, status: :unprocessable_entity
|
||||
rescue Teams::EntitlementError, BrandingAttachments::CoverUploadError => e
|
||||
@entitlements = @club.teams.first&.entitlements
|
||||
flash.now[:alert] = e.message
|
||||
render :edit, status: :unprocessable_entity
|
||||
end
|
||||
|
||||
def billing
|
||||
@@ -82,6 +96,8 @@ module Public
|
||||
@entitlements = @team.entitlements
|
||||
@plans = Plan.ordered
|
||||
@payments = @club.billing_payments.recent.includes(:invoice).limit(50)
|
||||
@quote = @club.active_billing_quote
|
||||
@pending_transfer = @club.billing_transfer_orders.awaiting_payment.first
|
||||
end
|
||||
|
||||
def checkout
|
||||
@@ -92,6 +108,12 @@ module Public
|
||||
return
|
||||
end
|
||||
|
||||
if @club.active_billing_quote.present?
|
||||
redirect_to public_club_billing_path(@club),
|
||||
alert: t("flash.club_billing.quote_checkout_denied")
|
||||
return
|
||||
end
|
||||
|
||||
unless MatchLiveTv.stripe_enabled?
|
||||
redirect_to public_club_billing_path(@club), alert: t("flash.clubs.stripe_not_configured")
|
||||
return
|
||||
|
||||
@@ -0,0 +1,46 @@
|
||||
module Public
|
||||
class ContactsController < WebBaseController
|
||||
def show
|
||||
@inquiry = empty_inquiry
|
||||
end
|
||||
|
||||
def create
|
||||
result = Contacts::SubmitInquiry.call(params: inquiry_params, ip: request.remote_ip)
|
||||
|
||||
if result.ok
|
||||
redirect_to public_contatti_path, notice: t("flash.contacts.sent")
|
||||
return
|
||||
end
|
||||
|
||||
@inquiry = inquiry_params
|
||||
flash.now[:alert] = t("flash.contacts.#{result.error}")
|
||||
render :show, status: result.error == :throttled ? :too_many_requests : :unprocessable_entity
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def inquiry_params
|
||||
{
|
||||
name: params[:name].to_s,
|
||||
email: params[:email].to_s,
|
||||
club_name: params[:club_name].to_s,
|
||||
topic: params[:topic].to_s,
|
||||
message: params[:message].to_s,
|
||||
website: params[:website].to_s,
|
||||
accept_privacy: params[:accept_privacy].to_s
|
||||
}
|
||||
end
|
||||
|
||||
def empty_inquiry
|
||||
{
|
||||
name: "",
|
||||
email: "",
|
||||
club_name: "",
|
||||
topic: "commercial",
|
||||
message: "",
|
||||
website: "",
|
||||
accept_privacy: ""
|
||||
}
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -51,6 +51,7 @@ module Public
|
||||
@session.status.in?(%w[connecting live reconnecting paused])
|
||||
)
|
||||
@publisher_online = !@stream_closed && mediamtx_publisher_online?(@session)
|
||||
@cover_poster_url = Streams::CoverResolver.new(@match).effective_cover_url
|
||||
end
|
||||
|
||||
def status
|
||||
|
||||
@@ -7,7 +7,7 @@ module Public
|
||||
layout "regia"
|
||||
|
||||
protect_from_forgery with: :null_session
|
||||
skip_before_action :verify_authenticity_token, only: %i[score pause resume stop min_quality]
|
||||
skip_before_action :verify_authenticity_token, only: %i[score pause resume stop audio_mute min_quality]
|
||||
|
||||
before_action :set_session_from_token
|
||||
|
||||
@@ -61,6 +61,12 @@ module Public
|
||||
render json: status_payload
|
||||
end
|
||||
|
||||
def audio_mute
|
||||
muted = params.key?(:muted) ? params[:muted] : !@session.audio_muted
|
||||
Sessions::SetAudioMute.new(@session, muted: muted).call
|
||||
render json: status_payload
|
||||
end
|
||||
|
||||
def min_quality
|
||||
Sessions::SetMinQuality.new(@session, preset: params[:min_quality_preset]).call
|
||||
render json: status_payload
|
||||
@@ -90,6 +96,7 @@ module Public
|
||||
{
|
||||
status: @session.status,
|
||||
paused: @session.paused?,
|
||||
audio_muted: @session.audio_muted,
|
||||
stream_closed: closed,
|
||||
live: !closed && (@session.live? || @session.paused? || publisher_online),
|
||||
on_air: playable || (!closed && @session.status.in?(%w[connecting live reconnecting paused])),
|
||||
|
||||
@@ -9,6 +9,7 @@ module Public
|
||||
|
||||
protect_from_forgery with: :exception
|
||||
before_action :set_site_locale
|
||||
before_action :load_site_announcements
|
||||
|
||||
private
|
||||
|
||||
@@ -40,5 +41,27 @@ module Public
|
||||
def logged_in?
|
||||
current_user.present?
|
||||
end
|
||||
|
||||
PRIVATE_WEB_CONTROLLERS = %w[
|
||||
accounts clubs teams club_recordings club_billing
|
||||
club_matches matches team_roster_members
|
||||
].freeze
|
||||
|
||||
def load_site_announcements
|
||||
@site_announcements = if skip_site_announcements?
|
||||
AppAnnouncement.none
|
||||
else
|
||||
channel = private_web_area? ? :web_private : :web_public
|
||||
AppAnnouncement.for_channel(channel).newest
|
||||
end
|
||||
end
|
||||
|
||||
def skip_site_announcements?
|
||||
controller_name == "regia"
|
||||
end
|
||||
|
||||
def private_web_area?
|
||||
PRIVATE_WEB_CONTROLLERS.include?(controller_name)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -8,6 +8,7 @@ module Public
|
||||
{ loc: "#{base}/prezzi", changefreq: "monthly", priority: "0.9" },
|
||||
{ loc: "#{base}/pallavolo-giovanile", changefreq: "monthly", priority: "0.85" },
|
||||
{ loc: "#{base}/faq", changefreq: "monthly", priority: "0.8" },
|
||||
{ loc: "#{base}/contatti", changefreq: "monthly", priority: "0.6" },
|
||||
{ loc: "#{base}/live", changefreq: "hourly", priority: "0.85" },
|
||||
{ loc: "#{base}/squadre", changefreq: "daily", priority: "0.85" },
|
||||
{ loc: "#{base}/privacy", changefreq: "yearly", priority: "0.3" },
|
||||
|
||||
@@ -19,6 +19,7 @@ module Public
|
||||
require_club_owner!(@club)
|
||||
team = @club.teams.create!(team_params)
|
||||
attach_branding_logo(team)
|
||||
Teams::StaffAssignment.designate_club_owner!(team: team, user: current_user)
|
||||
redirect_to public_club_path(@club), notice: t("flash.teams.added", name: team.name)
|
||||
rescue ActiveRecord::RecordInvalid => e
|
||||
flash.now[:alert] = e.record.errors.full_messages.join(", ")
|
||||
@@ -41,24 +42,33 @@ module Public
|
||||
def edit
|
||||
require_club_owner_for_team!(@team)
|
||||
@club = @team.club
|
||||
@entitlements = @team.entitlements
|
||||
end
|
||||
|
||||
def update
|
||||
require_club_owner_for_team!(@team)
|
||||
@team.assign_attributes(team_params)
|
||||
attach_branding_logo(@team)
|
||||
attach_cover_image(@team, entitlements: @team.entitlements)
|
||||
attach_team_photo(@team)
|
||||
@team.save!
|
||||
enqueue_cover_slate_generation(@team)
|
||||
redirect_to public_team_details_path(@team), notice: t("flash.teams.updated")
|
||||
rescue ActiveRecord::RecordInvalid => e
|
||||
@club = @team.club
|
||||
@entitlements = @team.entitlements
|
||||
flash.now[:alert] = e.record.errors.full_messages.join(", ")
|
||||
render :edit, status: :unprocessable_entity
|
||||
rescue Teams::EntitlementError, BrandingAttachments::CoverUploadError => e
|
||||
@club = @team.club
|
||||
@entitlements = @team.entitlements
|
||||
flash.now[:alert] = e.message
|
||||
render :edit, status: :unprocessable_entity
|
||||
end
|
||||
|
||||
def invite
|
||||
require_club_owner_for_team!(@team)
|
||||
@entitlements = @team.entitlements
|
||||
redirect_to public_team_details_path(@team, anchor: "invita-trasmissione")
|
||||
end
|
||||
|
||||
def assign_self_staff
|
||||
@@ -87,18 +97,32 @@ module Public
|
||||
email = params[:email]&.downcase&.strip
|
||||
Teams::StaffEmailValidator.assert_available!(team: @team, email: email, staff_kind: staff_kind)
|
||||
token = TeamInvitation.generate_token
|
||||
@team.team_invitations.create!(
|
||||
invitation = @team.team_invitations.create!(
|
||||
email: email,
|
||||
token_digest: Digest::SHA256.hexdigest(token),
|
||||
role: "member",
|
||||
staff_kind: staff_kind,
|
||||
expires_at: 7.days.from_now
|
||||
)
|
||||
@invite_url = join_public_invitation_url(token: token)
|
||||
flash.now[:notice] = t("flash.teams.invite_link_generated")
|
||||
render :invite
|
||||
invite_url = public_invitation_url(token: token)
|
||||
flash[:invite_url] = invite_url
|
||||
|
||||
begin
|
||||
Teams::InvitationMailer.transmission_invite(
|
||||
team: @team,
|
||||
invitation: invitation,
|
||||
invite_url: invite_url,
|
||||
invited_by: current_user
|
||||
).deliver_now
|
||||
flash[:notice] = t("flash.teams.invite_email_sent", email: email)
|
||||
rescue StandardError => e
|
||||
Rails.logger.error("[invite_email] #{e.class}: #{e.message}")
|
||||
flash[:alert] = t("flash.teams.invite_email_failed", email: email)
|
||||
end
|
||||
|
||||
redirect_to public_team_details_path(@team, anchor: "invito-generato")
|
||||
rescue Teams::EntitlementError, Teams::StaffAssignmentError => e
|
||||
redirect_to public_team_invite_path(@team), alert: e.message
|
||||
redirect_to public_team_details_path(@team, anchor: "invita-trasmissione"), alert: e.message
|
||||
end
|
||||
|
||||
def remove_member
|
||||
|
||||
@@ -30,6 +30,19 @@ module AdminHelper
|
||||
links
|
||||
end
|
||||
|
||||
def admin_session_ingest_badge_class(node)
|
||||
case node.role
|
||||
when "home" then "badge--ingest-home"
|
||||
when "lab" then "badge--ingest-lab"
|
||||
when "cloud" then "badge--ingest-cloud"
|
||||
else "badge--paused"
|
||||
end
|
||||
end
|
||||
|
||||
def admin_session_ingest_role_label(node)
|
||||
I18n.t("admin.sessions.ingest.role.#{node.role}", default: node.role.to_s.humanize)
|
||||
end
|
||||
|
||||
def admin_regia_expires_label(iso_time)
|
||||
return if iso_time.blank?
|
||||
|
||||
@@ -37,4 +50,31 @@ module AdminHelper
|
||||
rescue ArgumentError, TypeError
|
||||
nil
|
||||
end
|
||||
|
||||
def announcement_status_badge(item)
|
||||
case item.status
|
||||
when "published" then item.severity == "critical" ? "live" : "ready"
|
||||
when "archived" then "paused"
|
||||
else "connecting"
|
||||
end
|
||||
end
|
||||
|
||||
def announcement_window_label(item)
|
||||
start_label = item.starts_at&.in_time_zone&.strftime("%d/%m %H:%M")
|
||||
end_label = item.ends_at&.in_time_zone&.strftime("%d/%m %H:%M")
|
||||
if start_label && end_label
|
||||
"#{start_label} → #{end_label}"
|
||||
elsif start_label
|
||||
I18n.t("admin.announcements.window.from", time: start_label)
|
||||
elsif end_label
|
||||
I18n.t("admin.announcements.window.until", time: end_label)
|
||||
else
|
||||
I18n.t("admin.announcements.window.always")
|
||||
end
|
||||
end
|
||||
|
||||
def announcement_channels_label(item)
|
||||
labels = item.selected_channels.map { |key| I18n.t("admin.announcements.channels.#{key}") }
|
||||
labels.presence&.join(" · ") || I18n.t("admin.common.dash")
|
||||
end
|
||||
end
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
module LegalHelper
|
||||
def legal_last_updated
|
||||
"3 giugno 2026"
|
||||
"20 agosto 2026"
|
||||
end
|
||||
|
||||
def cookie_policy_last_updated
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
module Public
|
||||
module BillingHelper
|
||||
def plan_billing_action(current_slug:, target_plan:, stripe_subscription_active:, current_interval: nil, subscription: nil)
|
||||
def plan_billing_action(current_slug:, target_plan:, stripe_subscription_active:, current_interval: nil, subscription: nil, club: nil)
|
||||
if target_plan.slug == "free"
|
||||
return { kind: :current, label: I18n.t("billing.actions.current_plan") } if current_slug == "free"
|
||||
return { kind: :none } if stripe_subscription_active
|
||||
@@ -8,17 +8,50 @@ module Public
|
||||
return { kind: :contact, label: I18n.t("billing.actions.contact_for_free") }
|
||||
end
|
||||
|
||||
quote = club&.active_billing_quote
|
||||
if quote
|
||||
if quote.plan_slug == target_plan.slug
|
||||
if current_paid_plan?(subscription, target_plan)
|
||||
return {
|
||||
kind: :current,
|
||||
label: "#{I18n.t('billing.actions.current_plan')} — #{quote.price_label}"
|
||||
}
|
||||
end
|
||||
|
||||
return { kind: :quoted, plan: target_plan, quote: quote, intervals: [quote.billing_interval] }
|
||||
end
|
||||
|
||||
return { kind: :quoted_other, label: I18n.t("billing.bank_transfer.quoted_other") }
|
||||
end
|
||||
|
||||
intervals = bank_transfer_intervals_for(target_plan)
|
||||
unless MatchLiveTv.stripe_enabled?
|
||||
if current_paid_plan?(subscription, target_plan)
|
||||
return { kind: :current, label: current_plan_label(target_plan, subscription) }
|
||||
end
|
||||
if MatchLiveTv.bank_transfer_configured?
|
||||
return { kind: :bank_only, plan: target_plan, intervals: intervals }
|
||||
end
|
||||
|
||||
return { kind: :disabled, label: I18n.t("billing.actions.stripe_not_configured") }
|
||||
end
|
||||
|
||||
intervals = Billing::Stripe::PriceCatalog.available_intervals(plan_slug: target_plan.slug)
|
||||
return { kind: :disabled, label: I18n.t("billing.actions.stripe_prices_not_configured") } if intervals.empty?
|
||||
stripe_intervals = Billing::Stripe::PriceCatalog.available_intervals(plan_slug: target_plan.slug)
|
||||
if stripe_intervals.empty?
|
||||
if current_paid_plan?(subscription, target_plan)
|
||||
return { kind: :current, label: current_plan_label(target_plan, subscription) }
|
||||
end
|
||||
if MatchLiveTv.bank_transfer_configured?
|
||||
return { kind: :bank_only, plan: target_plan, intervals: intervals }
|
||||
end
|
||||
|
||||
return { kind: :disabled, label: I18n.t("billing.actions.stripe_prices_not_configured") }
|
||||
end
|
||||
|
||||
active_interval = current_interval.presence || Billing::Stripe::PriceCatalog::DEFAULT_INTERVAL
|
||||
|
||||
if current_slug == target_plan.slug && stripe_subscription_active && !subscription&.plan_change_pending?
|
||||
other_intervals = intervals - [active_interval]
|
||||
other_intervals = stripe_intervals - [active_interval]
|
||||
if other_intervals.empty?
|
||||
label = "#{I18n.t('billing.actions.current_plan')} — #{Billing::Stripe::PriceCatalog.label(plan_slug: target_plan.slug, interval: active_interval)}"
|
||||
return { kind: :current, label: label }
|
||||
@@ -34,13 +67,17 @@ module Public
|
||||
}
|
||||
end
|
||||
|
||||
if current_paid_plan?(subscription, target_plan)
|
||||
return { kind: :current, label: current_plan_label(target_plan, subscription) }
|
||||
end
|
||||
|
||||
if current_slug == "free" || !stripe_subscription_active
|
||||
{ kind: :checkout_options, plan: target_plan, intervals: intervals, subscription: subscription }
|
||||
{ kind: :checkout_options, plan: target_plan, intervals: stripe_intervals.presence || intervals, subscription: subscription }
|
||||
else
|
||||
{
|
||||
kind: :change_options,
|
||||
plan: target_plan,
|
||||
intervals: intervals,
|
||||
intervals: stripe_intervals,
|
||||
current_slug: current_slug,
|
||||
current_interval: active_interval,
|
||||
subscription: subscription
|
||||
@@ -57,6 +94,10 @@ module Public
|
||||
Billing::Stripe::PlanChangeMessages.button_label(plan: plan, interval: interval, kind: btn_kind)
|
||||
end
|
||||
|
||||
def plan_cta_name(plan)
|
||||
I18n.t("pages.plans.cta_name.#{plan.slug}", default: plan.name)
|
||||
end
|
||||
|
||||
def billing_plan_change_info_lines(subscription: nil)
|
||||
Billing::Stripe::PlanChangeMessages.billing_info_lines(subscription: subscription)
|
||||
end
|
||||
@@ -92,5 +133,41 @@ module Public
|
||||
I18n.t("billing.actions.profile_missing", fields: missing.join(", "))
|
||||
end
|
||||
|
||||
def bank_transfer_intervals_for(plan)
|
||||
Billing::Stripe::PriceCatalog.catalog_intervals(plan_slug: plan.slug)
|
||||
end
|
||||
|
||||
def bank_transfer_price_label(plan, interval, quote: nil)
|
||||
if quote&.matches?(plan.slug, interval)
|
||||
quote.price_label
|
||||
else
|
||||
Billing::Stripe::PriceCatalog.format_interval_price(
|
||||
Billing::Stripe::PriceCatalog.amount_cents(plan_slug: plan.slug, interval: interval),
|
||||
interval
|
||||
)
|
||||
end
|
||||
end
|
||||
|
||||
def show_bank_transfer_for?(club:, plan:, interval:, quote: nil, pending_transfer: nil)
|
||||
return false unless club.present? && plan.slug != "free"
|
||||
return false unless MatchLiveTv.bank_transfer_configured?
|
||||
return false if quote && !quote.matches?(plan.slug, interval)
|
||||
return false if pending_transfer&.awaiting_payment?
|
||||
return false if current_paid_plan?(club.subscription, plan)
|
||||
|
||||
true
|
||||
end
|
||||
|
||||
def current_paid_plan?(subscription, target_plan)
|
||||
return false unless subscription&.active? && subscription.premium?
|
||||
return false if subscription.plan_change_pending?
|
||||
|
||||
subscription.plan.slug == target_plan.slug
|
||||
end
|
||||
|
||||
def current_plan_label(target_plan, subscription)
|
||||
interval = subscription&.billing_interval.presence || Billing::Stripe::PriceCatalog::DEFAULT_INTERVAL
|
||||
"#{I18n.t('billing.actions.current_plan')} — #{Billing::Stripe::PriceCatalog.label(plan_slug: target_plan.slug, interval: interval)}"
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -44,11 +44,26 @@ module Public
|
||||
partialsPrefix: t("regia.board.partials_prefix"),
|
||||
resumed: t("regia.js.resumed"),
|
||||
pausedCover: t("regia.js.paused_cover"),
|
||||
muted: t("regia.js.muted"),
|
||||
unmuted: t("regia.js.unmuted"),
|
||||
closed: t("regia.js.closed"),
|
||||
resumeError: t("regia.js.resume_error"),
|
||||
pauseError: t("regia.js.pause_error"),
|
||||
muteError: t("regia.js.mute_error"),
|
||||
closeError: t("regia.js.close_error"),
|
||||
closeConfirm: t("regia.js.close_confirm"),
|
||||
pauseConfirmTitle: t("regia.js.pause_confirm_title"),
|
||||
pauseConfirmBody: t("regia.js.pause_confirm_body"),
|
||||
resumeConfirmTitle: t("regia.js.resume_confirm_title"),
|
||||
resumeConfirmBody: t("regia.js.resume_confirm_body"),
|
||||
muteConfirmTitle: t("regia.js.mute_confirm_title"),
|
||||
muteConfirmBody: t("regia.js.mute_confirm_body"),
|
||||
unmuteConfirmTitle: t("regia.js.unmute_confirm_title"),
|
||||
unmuteConfirmBody: t("regia.js.unmute_confirm_body"),
|
||||
advancePeriodTitle: t("regia.js.advance_period_title"),
|
||||
advancePeriodBody: t("regia.js.advance_period_body"),
|
||||
confirmAction: t("regia.modal.confirm_action"),
|
||||
stopConfirmTitle: t("regia.js.stop_confirm_title"),
|
||||
linkUnavailable: t("regia.js.link_unavailable"),
|
||||
linkCopied: t("regia.js.link_copied"),
|
||||
copyPrompt: t("regia.js.copy_prompt"),
|
||||
@@ -64,6 +79,8 @@ module Public
|
||||
previewOffHint: t("regia.preview_off_hint"),
|
||||
resumeLabel: t("regia.resume"),
|
||||
pauseLabel: t("regia.pause"),
|
||||
muteLabel: t("regia.mute"),
|
||||
unmuteLabel: t("regia.unmute"),
|
||||
endedBadge: t("regia.status.ended"),
|
||||
pausedBadge: t("regia.status.paused"),
|
||||
liveBadge: t("regia.status.live"),
|
||||
|
||||
@@ -6,7 +6,7 @@ class CleanupExpiredSessionsJob
|
||||
.where("updated_at < ?", 6.hours.ago)
|
||||
.find_each do |session|
|
||||
session.fail! if session.may_fail?
|
||||
Mediamtx::Client.new.delete_path(session)
|
||||
Mediamtx::Client.for_session(session).delete_path(session)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -6,12 +6,37 @@ class ExpireEndedSubscriptionsJob
|
||||
free_plan = Plan["free"]
|
||||
|
||||
Subscription.where(cancel_at_period_end: true)
|
||||
.where("current_period_end <= ?", Time.current)
|
||||
.where.not(plan_id: free_plan.id)
|
||||
.find_each do |sub|
|
||||
Billing::Stripe::FinalizeSubscription.call(club: sub.club)
|
||||
.where("current_period_end <= ?", Time.current)
|
||||
.where.not(plan_id: free_plan.id)
|
||||
.where(admin_comped: false)
|
||||
.find_each do |sub|
|
||||
expire!(sub)
|
||||
rescue StandardError => e
|
||||
Rails.logger.warn("[ExpireEndedSubscriptions] club=#{sub.club_id} #{e.message}")
|
||||
end
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def expire!(sub)
|
||||
if sub.stripe_subscription_id.present?
|
||||
Billing::Stripe::FinalizeSubscription.call(club: sub.club)
|
||||
else
|
||||
Billing::AssignPlan.call(
|
||||
club: sub.club,
|
||||
plan_slug: "free",
|
||||
status: "active",
|
||||
stripe_attrs: {
|
||||
stripe_subscription_id: nil,
|
||||
stripe_schedule_id: nil,
|
||||
current_period_start: nil,
|
||||
current_period_end: nil,
|
||||
cancel_at_period_end: false,
|
||||
billing_interval: nil,
|
||||
pending_plan_id: nil,
|
||||
pending_billing_interval: nil
|
||||
}
|
||||
)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -0,0 +1,19 @@
|
||||
class GenerateCoverSlateJob < ApplicationJob
|
||||
queue_as :default
|
||||
|
||||
def self.enqueue_for(record)
|
||||
return unless record&.cover_image&.attached?
|
||||
|
||||
queue = record.is_a?(Match) ? :critical : :default
|
||||
set(queue: queue).perform_later(record.class.name, record.id)
|
||||
end
|
||||
|
||||
def perform(class_name, record_id)
|
||||
record = class_name.constantize.find_by(id: record_id)
|
||||
return unless record&.cover_image&.attached?
|
||||
|
||||
Streams::GenerateCoverSlate.call(record)
|
||||
rescue Streams::GenerateCoverSlate::Error => e
|
||||
Rails.logger.warn("[GenerateCoverSlateJob] #{class_name}##{record_id} #{e.message}")
|
||||
end
|
||||
end
|
||||
@@ -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
|
||||
@@ -1,6 +1,6 @@
|
||||
# Avvia/riavvia il relay YouTube solo nel container sidekiq (YOUTUBE_RELAY_WORKER=1).
|
||||
# Avvia/riavvia il relay YouTube sul worker del nodo (coda youtube_relay_<slug>).
|
||||
class YoutubeRelayEnsureJob < ApplicationJob
|
||||
queue_as :default
|
||||
queue_as Streams::YoutubeRelay::QUEUE
|
||||
|
||||
def perform(session_id)
|
||||
session = StreamSession.find_by(id: session_id)
|
||||
|
||||
@@ -1,10 +1,23 @@
|
||||
# Ferma il relay solo sull'owner. Se il job gira su un altro host, requeue breve.
|
||||
class YoutubeRelayStopJob < ApplicationJob
|
||||
queue_as :default
|
||||
queue_as Streams::YoutubeRelay::QUEUE
|
||||
|
||||
def perform(session_id)
|
||||
discard_on ActiveJob::DeserializationError
|
||||
|
||||
def perform(session_id, attempts = 0)
|
||||
session = StreamSession.find_by(id: session_id)
|
||||
return unless session
|
||||
|
||||
Streams::YoutubeRelay.stop_on_worker!(session) if Streams::YoutubeRelay.worker?
|
||||
unless Streams::YoutubeRelay.worker?
|
||||
# Solo i worker relay processano questa coda in modo utile.
|
||||
return
|
||||
end
|
||||
|
||||
result = Streams::YoutubeRelay.stop_on_worker!(session)
|
||||
return unless result == :wrong_host
|
||||
return if attempts >= 30
|
||||
|
||||
self.class.set(wait: 2.seconds, queue: Streams::YoutubeRelay.queue_for(session))
|
||||
.perform_later(session_id, attempts + 1)
|
||||
end
|
||||
end
|
||||
|
||||
@@ -0,0 +1,35 @@
|
||||
module Billing
|
||||
class BankTransferMailer < ApplicationMailer
|
||||
def instructions
|
||||
@order = params[:order]
|
||||
@club = @order.club
|
||||
@iban = MatchLiveTv.bank_transfer_iban
|
||||
@holder = MatchLiveTv.bank_transfer_account_holder
|
||||
@bank_name = MatchLiveTv.bank_transfer_bank_name
|
||||
@bic = MatchLiveTv.bank_transfer_bic
|
||||
@proof_email = MatchLiveTv.bank_transfer_proof_email
|
||||
|
||||
mail(
|
||||
to: recipient_email,
|
||||
subject: t("mailers.bank_transfer.instructions.subject", plan: @order.plan.name)
|
||||
)
|
||||
end
|
||||
|
||||
def plan_activated
|
||||
@order = params[:order]
|
||||
@club = @order.club
|
||||
@subscription = @club.subscription
|
||||
|
||||
mail(
|
||||
to: recipient_email,
|
||||
subject: t("mailers.bank_transfer.plan_activated.subject", plan: @order.plan.name)
|
||||
)
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def recipient_email
|
||||
@club.billing_email.presence || @club.owner&.email
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -5,17 +5,37 @@ module Billing
|
||||
@club = @invoice.club
|
||||
|
||||
I18n.with_locale(I18n.locale) do
|
||||
prefix = t("mailers.invoice.attachment_prefix")
|
||||
attachments["#{prefix}-#{@invoice.number.parameterize}.pdf"] = {
|
||||
mime_type: "application/pdf",
|
||||
content: @invoice.pdf.download
|
||||
}
|
||||
|
||||
attach_invoice_pdf!
|
||||
mail(
|
||||
to: @club.billing_email,
|
||||
subject: t("mailers.invoice.subject", number: @invoice.display_number)
|
||||
)
|
||||
end
|
||||
end
|
||||
|
||||
def plan_activated_with_invoice
|
||||
@invoice = params[:invoice]
|
||||
@club = @invoice.club
|
||||
|
||||
I18n.with_locale(I18n.locale) do
|
||||
attach_invoice_pdf!
|
||||
mail(
|
||||
to: @club.billing_email,
|
||||
subject: t("mailers.invoice_activated.subject",
|
||||
plan: @club.subscription&.plan&.name || @invoice.display_number,
|
||||
number: @invoice.display_number)
|
||||
)
|
||||
end
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def attach_invoice_pdf!
|
||||
prefix = t("mailers.invoice.attachment_prefix")
|
||||
attachments["#{prefix}-#{@invoice.number.parameterize}.pdf"] = {
|
||||
mime_type: "application/pdf",
|
||||
content: @invoice.pdf.download
|
||||
}
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -0,0 +1,25 @@
|
||||
class ContactMailer < ApplicationMailer
|
||||
def inquiry(name:, email:, club_name:, topic:, message:)
|
||||
@name = name
|
||||
@email = email
|
||||
@club_name = club_name.presence
|
||||
@topic = topic
|
||||
@message = message
|
||||
|
||||
mail(
|
||||
to: recipient_for(topic),
|
||||
reply_to: email,
|
||||
subject: t("mailers.contact.subject.#{topic}", name: name)
|
||||
)
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def recipient_for(topic)
|
||||
if topic == "privacy"
|
||||
MatchLiveTv.privacy_controller_email
|
||||
else
|
||||
MatchLiveTv.support_email
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,25 @@
|
||||
module Teams
|
||||
class InvitationMailer < ApplicationMailer
|
||||
default from: -> { MatchLiveTv.mail_from }
|
||||
|
||||
def transmission_invite(team:, invitation:, invite_url:, invited_by:)
|
||||
@team = team
|
||||
@club = team.club
|
||||
@invitation = invitation
|
||||
@invite_url = invite_url
|
||||
@invited_by = invited_by
|
||||
@expires_on = I18n.l(invitation.expires_at.to_date, format: :long)
|
||||
|
||||
I18n.with_locale(I18n.locale) do
|
||||
mail(
|
||||
to: invitation.email,
|
||||
subject: t(
|
||||
"mailers.transmission_invite.subject",
|
||||
team: @team.name,
|
||||
club: @club.name
|
||||
)
|
||||
)
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -2,6 +2,7 @@ class AdminAccount < ApplicationRecord
|
||||
include PasswordComplexity
|
||||
|
||||
has_secure_password
|
||||
has_many :app_announcements, foreign_key: :created_by_admin_id, dependent: :nullify, inverse_of: :created_by_admin
|
||||
|
||||
validates :username, presence: true, uniqueness: true
|
||||
end
|
||||
|
||||
@@ -0,0 +1,240 @@
|
||||
class AppAnnouncement < ApplicationRecord
|
||||
KINDS = %w[info maintenance update].freeze
|
||||
SEVERITIES = %w[info warning critical].freeze
|
||||
STATUSES = %w[draft published archived].freeze
|
||||
CHANNELS = {
|
||||
android: :show_on_android,
|
||||
ios: :show_on_ios,
|
||||
web_public: :show_on_web_public,
|
||||
web_private: :show_on_web_private
|
||||
}.freeze
|
||||
COPY_KEYS = %w[title body action_url action_label].freeze
|
||||
URL_KEYS = %w[action_url].freeze
|
||||
TITLE_KEYS = %w[title].freeze
|
||||
BODY_KEYS = %w[body].freeze
|
||||
LABEL_KEYS = %w[action_label].freeze
|
||||
LOCALE_FALLBACKS = {
|
||||
"it" => %w[it],
|
||||
"en" => %w[en it],
|
||||
"fr" => %w[fr en it],
|
||||
"de" => %w[de en it],
|
||||
"es" => %w[es en it]
|
||||
}.freeze
|
||||
|
||||
belongs_to :created_by_admin, class_name: "AdminAccount", optional: true
|
||||
|
||||
validates :kind, :severity, :status, :title, :body, presence: true
|
||||
validates :kind, inclusion: { in: KINDS }
|
||||
validates :severity, inclusion: { in: SEVERITIES }
|
||||
validates :status, inclusion: { in: STATUSES }
|
||||
validates :title, length: { maximum: 120 }, allow_blank: true
|
||||
validates :body, length: { maximum: 2000 }, allow_blank: true
|
||||
validates :action_label, length: { maximum: 40 }, allow_blank: true
|
||||
URL_KEYS.each do |field|
|
||||
validates field, format: { with: /\Ahttps?:\/\/.+/i, allow_blank: true }
|
||||
end
|
||||
validate :ends_at_after_starts_at
|
||||
validate :at_least_one_channel
|
||||
validate :translations_are_valid
|
||||
before_validation :normalize_translations
|
||||
|
||||
scope :newest, -> { order(created_at: :desc) }
|
||||
scope :published, -> { where(status: "published") }
|
||||
scope :active_now, lambda {
|
||||
now = Time.current
|
||||
published
|
||||
.where("starts_at IS NULL OR starts_at <= ?", now)
|
||||
.where("ends_at IS NULL OR ends_at > ?", now)
|
||||
}
|
||||
|
||||
def self.for_channel(channel)
|
||||
column = CHANNELS[channel.to_sym]
|
||||
raise ArgumentError, "unknown announcement channel: #{channel}" if column.blank?
|
||||
|
||||
active_now.where(column => true)
|
||||
end
|
||||
|
||||
def self.for_app
|
||||
active_now.where("show_on_android = TRUE OR show_on_ios = TRUE")
|
||||
end
|
||||
|
||||
def self.for_platform(platform)
|
||||
case platform.to_s
|
||||
when "android" then for_channel(:android)
|
||||
when "ios" then for_channel(:ios)
|
||||
else for_app
|
||||
end
|
||||
end
|
||||
|
||||
def published?
|
||||
status == "published"
|
||||
end
|
||||
|
||||
def selected_channels
|
||||
CHANNELS.each_with_object([]) do |(key, column), keys|
|
||||
keys << key if public_send(column)
|
||||
end
|
||||
end
|
||||
|
||||
def copy_values_for(locale)
|
||||
loc = locale.to_s
|
||||
if default_locale?(loc)
|
||||
COPY_KEYS.index_with { |key| public_send(key).to_s.presence }
|
||||
else
|
||||
translation_hash(loc)
|
||||
end
|
||||
end
|
||||
|
||||
def filled_locale_codes
|
||||
LocaleResolver.available.select { |loc| locale_has_copy?(loc) }.map(&:to_s)
|
||||
end
|
||||
|
||||
def as_api_json(platform: nil, locale: I18n.locale)
|
||||
{
|
||||
id: id,
|
||||
kind: kind,
|
||||
severity: severity,
|
||||
title: copy_for("title", locale: locale),
|
||||
body: copy_for("body", locale: locale),
|
||||
dismissible: dismissible,
|
||||
action_url: resolved_action_url(platform, locale: locale),
|
||||
action_label: copy_for("action_label", locale: locale).presence,
|
||||
starts_at: starts_at&.iso8601,
|
||||
ends_at: ends_at&.iso8601
|
||||
}
|
||||
end
|
||||
|
||||
def localized_title(locale: I18n.locale)
|
||||
copy_for("title", locale: locale)
|
||||
end
|
||||
|
||||
def localized_body(locale: I18n.locale)
|
||||
copy_for("body", locale: locale)
|
||||
end
|
||||
|
||||
def localized_action_label(locale: I18n.locale)
|
||||
copy_for("action_label", locale: locale)
|
||||
end
|
||||
|
||||
def resolved_action_url(platform, locale: I18n.locale)
|
||||
url = copy_for("action_url", locale: locale)
|
||||
return url if url.present?
|
||||
return nil unless kind == "update"
|
||||
|
||||
case platform.to_s
|
||||
when "android" then MatchLiveTv.play_store_url
|
||||
when "ios" then MatchLiveTv.app_store_url
|
||||
end
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def copy_for(field, locale: I18n.locale)
|
||||
field = field.to_s
|
||||
locale_chain(locale).each do |loc|
|
||||
value = value_for(loc, field)
|
||||
return value if value.present?
|
||||
end
|
||||
nil
|
||||
end
|
||||
|
||||
def value_for(locale, key)
|
||||
copy_values_for(locale)[key].presence
|
||||
end
|
||||
|
||||
def translation_hash(locale)
|
||||
raw = translations.is_a?(Hash) ? translations : {}
|
||||
values = raw.stringify_keys[locale.to_s]
|
||||
return {} unless values.is_a?(Hash)
|
||||
|
||||
values.stringify_keys
|
||||
end
|
||||
|
||||
def locale_chain(locale)
|
||||
loc = (LocaleResolver.normalize(locale) || I18n.default_locale).to_s
|
||||
LOCALE_FALLBACKS[loc] || [loc, I18n.default_locale.to_s].uniq
|
||||
end
|
||||
|
||||
def locale_has_copy?(locale)
|
||||
values = copy_values_for(locale)
|
||||
values["title"].present? || values["body"].present?
|
||||
end
|
||||
|
||||
def default_locale?(locale)
|
||||
locale.to_s == I18n.default_locale.to_s
|
||||
end
|
||||
|
||||
def normalize_translations
|
||||
self.translations = self.class.sanitize_translations(translations)
|
||||
end
|
||||
|
||||
def self.sanitize_translations(value)
|
||||
raw = coerce_translation_hash(value)
|
||||
LocaleResolver.available.each_with_object({}) do |loc, acc|
|
||||
next if loc.to_s == I18n.default_locale.to_s
|
||||
|
||||
payload = raw[loc.to_s]
|
||||
next unless payload.is_a?(Hash)
|
||||
|
||||
cleaned = COPY_KEYS.each_with_object({}) do |key, fields|
|
||||
text = payload[key].to_s.strip
|
||||
fields[key] = text if text.present?
|
||||
end
|
||||
acc[loc.to_s] = cleaned if cleaned.any?
|
||||
end
|
||||
end
|
||||
|
||||
def self.coerce_translation_hash(value)
|
||||
hash = if value.is_a?(ActionController::Parameters)
|
||||
value.to_unsafe_h
|
||||
elsif value.is_a?(Hash)
|
||||
value
|
||||
else
|
||||
{}
|
||||
end
|
||||
hash.deep_stringify_keys
|
||||
end
|
||||
private_class_method :coerce_translation_hash
|
||||
|
||||
def at_least_one_channel
|
||||
return if show_on_android? || show_on_ios? || show_on_web_public? || show_on_web_private?
|
||||
|
||||
errors.add(:base, :no_channel)
|
||||
end
|
||||
|
||||
def ends_at_after_starts_at
|
||||
return if starts_at.blank? || ends_at.blank?
|
||||
return if ends_at > starts_at
|
||||
|
||||
errors.add(:ends_at, :invalid)
|
||||
end
|
||||
|
||||
def translations_are_valid
|
||||
translation_entries.each do |locale, payload|
|
||||
TITLE_KEYS.each do |key|
|
||||
errors.add(:base, :translation_too_long, locale: locale, field: key) if payload[key].to_s.length > 120
|
||||
end
|
||||
BODY_KEYS.each do |key|
|
||||
errors.add(:base, :translation_too_long, locale: locale, field: key) if payload[key].to_s.length > 2000
|
||||
end
|
||||
LABEL_KEYS.each do |key|
|
||||
errors.add(:base, :translation_too_long, locale: locale, field: key) if payload[key].to_s.length > 40
|
||||
end
|
||||
URL_KEYS.each do |key|
|
||||
url = payload[key].to_s
|
||||
next if url.blank? || url.match?(/\Ahttps?:\/\/.+/i)
|
||||
|
||||
errors.add(:base, :translation_invalid_url, locale: locale)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
def translation_entries
|
||||
raw = translations.is_a?(Hash) ? translations : {}
|
||||
raw.stringify_keys.filter_map do |locale, payload|
|
||||
next unless payload.is_a?(Hash)
|
||||
|
||||
[locale, payload.stringify_keys]
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,34 @@
|
||||
module Billing
|
||||
class ClubQuote < ApplicationRecord
|
||||
self.table_name = "billing_club_quotes"
|
||||
|
||||
PLAN_SLUGS = %w[premium_light premium_full].freeze
|
||||
INTERVALS = Billing::Stripe::PriceCatalog::INTERVALS
|
||||
|
||||
belongs_to :club
|
||||
belongs_to :created_by_admin, class_name: "AdminAccount", optional: true
|
||||
|
||||
validates :plan_slug, inclusion: { in: PLAN_SLUGS }
|
||||
validates :billing_interval, inclusion: { in: INTERVALS }
|
||||
validates :amount_cents, numericality: { greater_than: 0 }
|
||||
validates :currency, presence: true
|
||||
|
||||
scope :active, -> { where(active: true) }
|
||||
|
||||
def matches?(plan_slug, interval)
|
||||
self.plan_slug == plan_slug.to_s && billing_interval == interval.to_s
|
||||
end
|
||||
|
||||
def plan
|
||||
Plan[plan_slug]
|
||||
end
|
||||
|
||||
def formatted_amount
|
||||
Billing::Stripe::PriceCatalog.format_eur(amount_cents)
|
||||
end
|
||||
|
||||
def price_label
|
||||
Billing::Stripe::PriceCatalog.format_interval_price(amount_cents, billing_interval)
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -2,14 +2,17 @@ module Billing
|
||||
class Payment < ApplicationRecord
|
||||
self.table_name = "billing_payments"
|
||||
|
||||
STATUSES = %w[paid failed refunded].freeze
|
||||
STATUSES = %w[pending paid failed refunded].freeze
|
||||
PROVIDERS = %w[stripe bank_transfer].freeze
|
||||
|
||||
belongs_to :club
|
||||
has_one :invoice, class_name: "Billing::Invoice", foreign_key: :billing_payment_id, dependent: :nullify
|
||||
has_one :transfer_order, class_name: "Billing::TransferOrder", foreign_key: :billing_payment_id, dependent: :nullify
|
||||
|
||||
validates :amount_cents, numericality: { greater_than: 0 }
|
||||
validates :currency, presence: true
|
||||
validates :status, inclusion: { in: STATUSES }
|
||||
validates :provider, inclusion: { in: PROVIDERS }
|
||||
validates :stripe_invoice_id, uniqueness: true, allow_nil: true
|
||||
|
||||
scope :recent, -> { order(paid_at: :desc, created_at: :desc) }
|
||||
@@ -62,7 +65,12 @@ module Billing
|
||||
end
|
||||
|
||||
def display_status
|
||||
{ "paid" => "Pagato", "failed" => "Non riuscito", "refunded" => "Rimborsato" }[status] || status
|
||||
{
|
||||
"pending" => "In attesa di bonifico",
|
||||
"paid" => "Pagato",
|
||||
"failed" => "Non riuscito",
|
||||
"refunded" => "Rimborsato"
|
||||
}[status] || status
|
||||
end
|
||||
|
||||
def invoice_for_display
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
module Billing
|
||||
class TransferOrder < ApplicationRecord
|
||||
self.table_name = "billing_transfer_orders"
|
||||
|
||||
KINDS = %w[list_price commercial_quote].freeze
|
||||
STATUSES = %w[awaiting_payment paid cancelled].freeze
|
||||
PLAN_SLUGS = ClubQuote::PLAN_SLUGS
|
||||
INTERVALS = ClubQuote::INTERVALS
|
||||
|
||||
belongs_to :club
|
||||
belongs_to :billing_club_quote, class_name: "Billing::ClubQuote", optional: true
|
||||
belongs_to :billing_payment, class_name: "Billing::Payment", optional: true
|
||||
belongs_to :requested_by_user, class_name: "User", optional: true
|
||||
belongs_to :confirmed_by_admin, class_name: "AdminAccount", optional: true
|
||||
|
||||
validates :plan_slug, inclusion: { in: PLAN_SLUGS }
|
||||
validates :billing_interval, inclusion: { in: INTERVALS }
|
||||
validates :amount_cents, numericality: { greater_than: 0 }
|
||||
validates :kind, inclusion: { in: KINDS }
|
||||
validates :status, inclusion: { in: STATUSES }
|
||||
validates :reference_code, presence: true, uniqueness: true
|
||||
validates :currency, presence: true
|
||||
|
||||
scope :awaiting_payment, -> { where(status: "awaiting_payment").order(created_at: :desc) }
|
||||
scope :recent, -> { order(created_at: :desc) }
|
||||
|
||||
def awaiting_payment?
|
||||
status == "awaiting_payment"
|
||||
end
|
||||
|
||||
def paid?
|
||||
status == "paid"
|
||||
end
|
||||
|
||||
def quoted?
|
||||
kind == "commercial_quote"
|
||||
end
|
||||
|
||||
def plan
|
||||
Plan[plan_slug]
|
||||
end
|
||||
|
||||
def formatted_amount
|
||||
format("%.2f €", amount_cents / 100.0)
|
||||
end
|
||||
|
||||
def price_label
|
||||
Billing::Stripe::PriceCatalog.format_interval_price(amount_cents, billing_interval)
|
||||
end
|
||||
|
||||
def payment_causal
|
||||
"MLTV #{reference_code} #{club.name}".truncate(140, omission: "")
|
||||
end
|
||||
|
||||
def cancel!
|
||||
return self unless awaiting_payment?
|
||||
|
||||
transaction do
|
||||
update!(status: "cancelled", cancelled_at: Time.current)
|
||||
if billing_payment&.status == "pending"
|
||||
billing_payment.update!(status: "failed")
|
||||
end
|
||||
end
|
||||
self
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -1,5 +1,6 @@
|
||||
class Club < ApplicationRecord
|
||||
include Brandable
|
||||
include Coverable
|
||||
include ClubBillingProfile
|
||||
|
||||
has_many :club_memberships, dependent: :destroy
|
||||
@@ -7,8 +8,11 @@ class Club < ApplicationRecord
|
||||
has_many :teams, dependent: :destroy
|
||||
has_one :youtube_credential, dependent: :destroy
|
||||
has_one :subscription, dependent: :destroy
|
||||
has_one :billing_quote, -> { where(active: true) }, class_name: "Billing::ClubQuote", inverse_of: :club
|
||||
has_many :billing_quotes, class_name: "Billing::ClubQuote", dependent: :destroy, inverse_of: :club
|
||||
has_many :billing_payments, class_name: "Billing::Payment", dependent: :destroy
|
||||
has_many :billing_invoices, class_name: "Billing::Invoice", dependent: :destroy
|
||||
has_many :billing_transfer_orders, class_name: "Billing::TransferOrder", dependent: :destroy
|
||||
|
||||
validates :name, presence: true
|
||||
validates :sport, presence: true
|
||||
@@ -23,6 +27,14 @@ class Club < ApplicationRecord
|
||||
club_memberships.exists?(user: user, role: "owner")
|
||||
end
|
||||
|
||||
def active_billing_quote
|
||||
billing_quote
|
||||
end
|
||||
|
||||
def pending_transfer_order
|
||||
billing_transfer_orders.awaiting_payment.first
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def branding_parent
|
||||
|
||||
@@ -0,0 +1,50 @@
|
||||
module Coverable
|
||||
extend ActiveSupport::Concern
|
||||
|
||||
COVER_IMAGE_TYPES = %w[image/png image/jpeg image/webp].freeze
|
||||
MAX_COVER_BYTES = 5.megabytes
|
||||
|
||||
included do
|
||||
has_one_attached :cover_image
|
||||
has_one_attached :cover_slate
|
||||
|
||||
validate :cover_image_file_type, if: -> { cover_image.attached? }
|
||||
validate :cover_image_file_size, if: -> { cover_image.attached? }
|
||||
end
|
||||
|
||||
def cover_image_attached?
|
||||
cover_image.attached?
|
||||
end
|
||||
|
||||
def cover_url
|
||||
return unless cover_image.attached?
|
||||
|
||||
Rails.application.routes.url_helpers.rails_blob_path(cover_image, only_path: true)
|
||||
end
|
||||
|
||||
def effective_cover_url
|
||||
if cover_image.attached?
|
||||
cover_url
|
||||
else
|
||||
cover_parent&.effective_cover_url
|
||||
end
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def cover_parent
|
||||
nil
|
||||
end
|
||||
|
||||
def cover_image_file_type
|
||||
return if cover_image.content_type.in?(COVER_IMAGE_TYPES)
|
||||
|
||||
errors.add(:cover_image, I18n.t("coverable.errors.invalid_type"))
|
||||
end
|
||||
|
||||
def cover_image_file_size
|
||||
return if cover_image.byte_size <= MAX_COVER_BYTES
|
||||
|
||||
errors.add(:cover_image, I18n.t("coverable.errors.too_large"))
|
||||
end
|
||||
end
|
||||
@@ -1,4 +1,6 @@
|
||||
class Match < ApplicationRecord
|
||||
include Coverable
|
||||
|
||||
belongs_to :team
|
||||
has_many :stream_sessions, dependent: :destroy
|
||||
|
||||
@@ -174,4 +176,8 @@ class Match < ApplicationRecord
|
||||
rescue ArgumentError => e
|
||||
errors.add(:scoring_rules, e.message)
|
||||
end
|
||||
|
||||
def cover_parent
|
||||
nil
|
||||
end
|
||||
end
|
||||
|
||||
@@ -4,7 +4,7 @@ module Ops
|
||||
|
||||
KINDS = %w[
|
||||
disk_space recordings_size service_down http_public http_rails rails_latency
|
||||
sidekiq_stale sidekiq_dead log_pattern garage_storage
|
||||
sidekiq_stale sidekiq_dead log_pattern garage_storage stream_overflow
|
||||
].freeze
|
||||
SEVERITIES = %w[critical warning info].freeze
|
||||
STATUSES = %w[open acknowledged resolved].freeze
|
||||
|
||||
@@ -84,6 +84,10 @@ class Plan < ApplicationRecord
|
||||
feature("youtube_enabled") == true
|
||||
end
|
||||
|
||||
def custom_cover_enabled?
|
||||
feature("custom_cover") == true
|
||||
end
|
||||
|
||||
def premium?
|
||||
slug.in?(PREMIUM_SLUGS)
|
||||
end
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
class StreamEvent < ApplicationRecord
|
||||
EVENT_TYPES = %w[
|
||||
connected disconnected reconnected quality_changed error paused resumed
|
||||
started ended network_test pairing youtube_ready
|
||||
started ended network_test pairing youtube_ready audio_muted audio_unmuted
|
||||
].freeze
|
||||
|
||||
belongs_to :stream_session
|
||||
|
||||
@@ -0,0 +1,54 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
# Nodo streaming (MediaMTX [+ relay ffmpeg]). Fase 0: registry + assignment URL.
|
||||
class StreamNode < ApplicationRecord
|
||||
ROLES = %w[home cloud lab].freeze
|
||||
STATUSES = %w[provisioning ready draining offline error].freeze
|
||||
PROVIDERS = %w[local proxmox_lab hetzner].freeze
|
||||
|
||||
has_many :stream_sessions, dependent: :nullify
|
||||
|
||||
validates :slug, presence: true, uniqueness: true
|
||||
validates :hostname, presence: true
|
||||
validates :role, inclusion: { in: ROLES }
|
||||
validates :status, inclusion: { in: STATUSES }
|
||||
validates :provider, inclusion: { in: PROVIDERS }
|
||||
validates :rtmp_base_url, :hls_base_url, :api_base_url, presence: true
|
||||
validates :max_publishers, :max_relays, numericality: { greater_than: 0 }
|
||||
|
||||
scope :ready, -> { where(status: "ready") }
|
||||
scope :allocatable, -> { ready }
|
||||
|
||||
# Sessioni che occupano uno slot path MediaMTX su questo nodo.
|
||||
def occupying_sessions
|
||||
stream_sessions.where(status: %w[idle connecting live reconnecting paused])
|
||||
end
|
||||
|
||||
def active_publishers
|
||||
occupying_sessions.count
|
||||
end
|
||||
|
||||
def free_slots
|
||||
[max_publishers - active_publishers, 0].max
|
||||
end
|
||||
|
||||
def allocatable?
|
||||
status == "ready" && free_slots.positive?
|
||||
end
|
||||
|
||||
# Agent ffmpeg sul CPX (cloud-init :9100). Home Sidekiq lo chiama senza worker Docker per nodo.
|
||||
def relay_agent_url
|
||||
return unless role == "cloud"
|
||||
|
||||
ip = (metadata || {})["public_ip"].presence
|
||||
if ip.blank?
|
||||
host = URI.parse(api_base_url.to_s).host
|
||||
ip = host if host.present? && host != "localhost"
|
||||
end
|
||||
return if ip.blank? || ip == "127.0.0.1"
|
||||
|
||||
"http://#{ip}:#{ENV.fetch("STREAM_NODE_AGENT_PORT", "9100")}"
|
||||
rescue URI::InvalidURIError
|
||||
nil
|
||||
end
|
||||
end
|
||||
@@ -7,6 +7,7 @@ class StreamSession < ApplicationRecord
|
||||
|
||||
belongs_to :match
|
||||
belongs_to :user
|
||||
belongs_to :stream_node, optional: true
|
||||
has_many :stream_events, dependent: :destroy
|
||||
has_one :score_state, dependent: :destroy
|
||||
has_many :device_states, dependent: :destroy
|
||||
@@ -75,7 +76,7 @@ class StreamSession < ApplicationRecord
|
||||
def rtmp_ingest_url
|
||||
# RootEncoder richiede rtmp://host:port/app/stream (due segmenti).
|
||||
# MediaMTX path = live/match_{uuid} (no ?token= nel path).
|
||||
"#{MatchLiveTv.mediamtx_rtmp_url}/#{mediamtx_path_name}"
|
||||
"#{rtmp_base_url.chomp('/')}/#{mediamtx_path_name}"
|
||||
end
|
||||
|
||||
def mediamtx_path_name
|
||||
@@ -91,8 +92,29 @@ class StreamSession < ApplicationRecord
|
||||
end
|
||||
|
||||
def hls_playback_url
|
||||
base = MatchLiveTv.hls_public_url.chomp("/")
|
||||
"#{base}/#{effective_hls_path_name}/index.m3u8"
|
||||
"#{hls_base_url.chomp('/')}/#{effective_hls_path_name}/index.m3u8"
|
||||
end
|
||||
|
||||
def rtmp_base_url
|
||||
stream_node&.rtmp_base_url.presence || MatchLiveTv.mediamtx_rtmp_url
|
||||
end
|
||||
|
||||
def hls_base_url
|
||||
stream_node&.hls_base_url.presence || MatchLiveTv.hls_public_url
|
||||
end
|
||||
|
||||
def mediamtx_api_base_url
|
||||
stream_node&.api_base_url.presence || MatchLiveTv.mediamtx_api_url
|
||||
end
|
||||
|
||||
def mediamtx_internal_rtmp_url
|
||||
stream_node&.internal_rtmp_url.presence ||
|
||||
ENV.fetch("MEDIAMTX_INTERNAL_RTMP_URL", "rtmp://mediamtx:1935")
|
||||
end
|
||||
|
||||
def mediamtx_internal_hls_url
|
||||
stream_node&.internal_hls_url.presence ||
|
||||
ENV.fetch("MEDIAMTX_HLS_URL", "http://mediamtx:8888")
|
||||
end
|
||||
|
||||
def effective_hls_path_name
|
||||
|
||||
@@ -35,4 +35,8 @@ class Subscription < ApplicationRecord
|
||||
def admin_comped?
|
||||
admin_comped
|
||||
end
|
||||
|
||||
def bank_transfer?
|
||||
premium? && stripe_subscription_id.blank? && !admin_comped? && current_period_end.present?
|
||||
end
|
||||
end
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
class Team < ApplicationRecord
|
||||
include Brandable
|
||||
include Coverable
|
||||
|
||||
belongs_to :club
|
||||
has_many :user_teams, dependent: :destroy
|
||||
@@ -99,6 +100,10 @@ class Team < ApplicationRecord
|
||||
club
|
||||
end
|
||||
|
||||
def cover_parent
|
||||
club
|
||||
end
|
||||
|
||||
def assign_slug
|
||||
self.slug = Teams::GenerateSlug.call(self) if slug.blank?
|
||||
end
|
||||
|
||||
@@ -4,19 +4,33 @@ class YoutubeCredential < ApplicationRecord
|
||||
attr_encrypted :access_token,
|
||||
key: :encryption_key,
|
||||
attribute: "access_token_encrypted",
|
||||
mode: :single_iv_salt
|
||||
mode: :single_iv_and_salt,
|
||||
algorithm: "aes-256-cbc",
|
||||
iv: :encryption_iv
|
||||
attr_encrypted :refresh_token,
|
||||
key: :encryption_key,
|
||||
attribute: "refresh_token_encrypted",
|
||||
mode: :single_iv_salt
|
||||
mode: :single_iv_and_salt,
|
||||
algorithm: "aes-256-cbc",
|
||||
iv: :encryption_iv
|
||||
|
||||
def expired?
|
||||
expires_at.present? && expires_at < Time.current
|
||||
end
|
||||
|
||||
def usable?
|
||||
refresh_token.present?
|
||||
rescue ArgumentError, OpenSSL::Cipher::CipherError, NoMethodError
|
||||
false
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def encryption_key
|
||||
Rails.application.secret_key_base[0, 32]
|
||||
end
|
||||
|
||||
def encryption_iv
|
||||
encryption_key[0, 16]
|
||||
end
|
||||
end
|
||||
|
||||
@@ -23,6 +23,7 @@ module Billing
|
||||
raise Error, "Piano non valido" unless @plan_slug.in?(VALID_PLANS)
|
||||
|
||||
release_stripe_schedule!
|
||||
cancel_awaiting_transfers!
|
||||
|
||||
AssignPlan.call(
|
||||
club: @club,
|
||||
@@ -97,5 +98,9 @@ module Billing
|
||||
ensure
|
||||
sub&.update!(stripe_schedule_id: nil, pending_plan_id: nil, pending_billing_interval: nil)
|
||||
end
|
||||
|
||||
def cancel_awaiting_transfers!
|
||||
@club.billing_transfer_orders.awaiting_payment.find_each(&:cancel!)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -2,13 +2,14 @@ module Billing
|
||||
class AttachPaymentInvoice
|
||||
class Error < StandardError; end
|
||||
|
||||
def self.call(payment:, pdf:)
|
||||
new(payment: payment, pdf: pdf).call
|
||||
def self.call(payment:, pdf:, mailer_action: :invoice_pdf)
|
||||
new(payment: payment, pdf: pdf, mailer_action: mailer_action).call
|
||||
end
|
||||
|
||||
def initialize(payment:, pdf:)
|
||||
def initialize(payment:, pdf:, mailer_action: :invoice_pdf)
|
||||
@payment = payment
|
||||
@pdf = pdf
|
||||
@mailer_action = mailer_action
|
||||
end
|
||||
|
||||
def call
|
||||
@@ -17,7 +18,7 @@ module Billing
|
||||
|
||||
club = @payment.club
|
||||
invoice = @payment.invoice || build_invoice!(club)
|
||||
IssueInvoice.call(invoice: invoice, pdf: @pdf)
|
||||
IssueInvoice.call(invoice: invoice, pdf: @pdf, mailer_action: @mailer_action)
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
@@ -0,0 +1,102 @@
|
||||
module Billing
|
||||
class ConfirmBankTransfer
|
||||
class Error < StandardError; end
|
||||
|
||||
def self.call(order:, admin:, pdf: nil)
|
||||
new(order: order, admin: admin, pdf: pdf).call
|
||||
end
|
||||
|
||||
def initialize(order:, admin:, pdf: nil)
|
||||
@order = order
|
||||
@admin = admin
|
||||
@pdf = pdf
|
||||
end
|
||||
|
||||
def call
|
||||
raise Error, "Bonifico già gestito" unless @order.awaiting_payment?
|
||||
|
||||
club = @order.club
|
||||
payment = @order.billing_payment
|
||||
raise Error, "Pagamento collegato mancante" if payment.blank?
|
||||
|
||||
ApplicationRecord.transaction do
|
||||
cancel_existing_stripe!(club.subscription)
|
||||
|
||||
period_start, period_end = period_bounds(club.subscription)
|
||||
AssignPlan.call(
|
||||
club: club,
|
||||
plan_slug: @order.plan_slug,
|
||||
status: "active",
|
||||
stripe_attrs: {
|
||||
stripe_subscription_id: nil,
|
||||
stripe_schedule_id: nil,
|
||||
pending_plan_id: nil,
|
||||
pending_billing_interval: nil,
|
||||
billing_interval: @order.billing_interval,
|
||||
current_period_start: period_start,
|
||||
current_period_end: period_end,
|
||||
cancel_at_period_end: true,
|
||||
admin_comped: false,
|
||||
admin_comped_reason: nil,
|
||||
admin_comped_at: nil,
|
||||
admin_comped_by_id: nil
|
||||
}
|
||||
)
|
||||
|
||||
payment.update!(status: "paid", paid_at: Time.current, provider: "bank_transfer")
|
||||
@order.update!(
|
||||
status: "paid",
|
||||
confirmed_by_admin: @admin,
|
||||
confirmed_at: Time.current
|
||||
)
|
||||
end
|
||||
|
||||
deliver_activation!(payment)
|
||||
@order.reload
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def period_bounds(subscription)
|
||||
start_at = Time.current
|
||||
if subscription&.premium? &&
|
||||
!subscription.admin_comped? &&
|
||||
subscription.plan.slug == @order.plan_slug &&
|
||||
subscription.current_period_end.present? &&
|
||||
subscription.current_period_end > Time.current
|
||||
start_at = subscription.current_period_end
|
||||
end
|
||||
|
||||
end_at = @order.billing_interval == "yearly" ? start_at.advance(years: 1) : start_at.advance(months: 1)
|
||||
[start_at, end_at]
|
||||
end
|
||||
|
||||
def cancel_existing_stripe!(subscription)
|
||||
return if subscription.blank? || subscription.stripe_subscription_id.blank?
|
||||
return unless MatchLiveTv.stripe_enabled?
|
||||
|
||||
::Stripe::Subscription.cancel(subscription.stripe_subscription_id)
|
||||
rescue ::Stripe::InvalidRequestError => e
|
||||
Rails.logger.warn("[BankTransfer] stripe cancel club=#{subscription.club_id} #{e.message}")
|
||||
end
|
||||
|
||||
def deliver_activation!(payment)
|
||||
if pdf_present?
|
||||
AttachPaymentInvoice.call(
|
||||
payment: payment.reload,
|
||||
pdf: @pdf,
|
||||
mailer_action: :plan_activated_with_invoice
|
||||
)
|
||||
else
|
||||
MatchLiveTv.deliver_mail(BankTransferMailer.with(order: @order.reload).plan_activated)
|
||||
end
|
||||
end
|
||||
|
||||
def pdf_present?
|
||||
return false if @pdf.blank?
|
||||
return @pdf.present? unless @pdf.respond_to?(:tempfile)
|
||||
|
||||
@pdf.original_filename.present?
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,26 @@
|
||||
module Billing
|
||||
class EuroAmount
|
||||
class Error < StandardError; end
|
||||
|
||||
def self.to_cents(value)
|
||||
raw = value.to_s.strip
|
||||
raise Error, "Indica l'importo in euro" if raw.blank?
|
||||
|
||||
normalized = raw.gsub(/\s+/, "")
|
||||
if normalized.match?(/\A\d{1,3}(\.\d{3})*,\d{1,2}\z/)
|
||||
normalized = normalized.gsub(".", "").tr(",", ".")
|
||||
elsif normalized.match?(/\A\d+,\d{1,2}\z/)
|
||||
normalized = normalized.tr(",", ".")
|
||||
elsif normalized.match?(/\A\d{1,3}(,\d{3})*\.\d{1,2}\z/)
|
||||
normalized = normalized.gsub(",", "")
|
||||
end
|
||||
|
||||
raise Error, "Importo non valido" unless normalized.match?(/\A\d+(\.\d{1,2})?\z/)
|
||||
|
||||
cents = (BigDecimal(normalized) * 100).round
|
||||
raise Error, "L'importo deve essere maggiore di zero" unless cents.positive?
|
||||
|
||||
cents.to_i
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -2,13 +2,17 @@ module Billing
|
||||
class IssueInvoice
|
||||
class Error < StandardError; end
|
||||
|
||||
def self.call(invoice:, pdf: nil)
|
||||
new(invoice: invoice, pdf: pdf).call
|
||||
MAILER_ACTIONS = %i[invoice_pdf plan_activated_with_invoice].freeze
|
||||
|
||||
def self.call(invoice:, pdf: nil, mailer_action: :invoice_pdf)
|
||||
new(invoice: invoice, pdf: pdf, mailer_action: mailer_action).call
|
||||
end
|
||||
|
||||
def initialize(invoice:, pdf: nil)
|
||||
def initialize(invoice:, pdf: nil, mailer_action: :invoice_pdf)
|
||||
@invoice = invoice
|
||||
@pdf = pdf
|
||||
@mailer_action = mailer_action.to_sym
|
||||
raise Error, "Azione email non valida" unless MAILER_ACTIONS.include?(@mailer_action)
|
||||
end
|
||||
|
||||
def call
|
||||
@@ -22,8 +26,10 @@ module Billing
|
||||
|
||||
@invoice.update!(status: "issued")
|
||||
|
||||
Billing::InvoiceMailer.with(invoice: @invoice).invoice_pdf.deliver_now
|
||||
@invoice.update!(status: "sent", emailed_at: Time.current)
|
||||
mail = Billing::InvoiceMailer.with(invoice: @invoice).public_send(@mailer_action)
|
||||
if MatchLiveTv.deliver_mail(mail)
|
||||
@invoice.update!(status: "sent", emailed_at: Time.current)
|
||||
end
|
||||
|
||||
@invoice
|
||||
end
|
||||
|
||||
@@ -5,9 +5,13 @@ module Billing
|
||||
|
||||
AMOUNT_HINTS = {
|
||||
500 => %w[premium_light monthly],
|
||||
790 => %w[premium_light monthly],
|
||||
4000 => %w[premium_light yearly],
|
||||
5900 => %w[premium_light yearly],
|
||||
2000 => %w[premium_full monthly],
|
||||
20_000 => %w[premium_full yearly]
|
||||
2490 => %w[premium_full monthly],
|
||||
20_000 => %w[premium_full yearly],
|
||||
19_900 => %w[premium_full yearly]
|
||||
}.freeze
|
||||
|
||||
class << self
|
||||
|
||||
@@ -0,0 +1,95 @@
|
||||
module Billing
|
||||
class RequestBankTransfer
|
||||
class Error < StandardError; end
|
||||
|
||||
PLAN_SLUGS = %w[premium_light premium_full].freeze
|
||||
|
||||
def self.call(club:, user:, plan_slug:, interval:)
|
||||
new(club: club, user: user, plan_slug: plan_slug, interval: interval).call
|
||||
end
|
||||
|
||||
def initialize(club:, user:, plan_slug:, interval:)
|
||||
@club = club
|
||||
@user = user
|
||||
@plan_slug = plan_slug.to_s
|
||||
@interval = interval
|
||||
end
|
||||
|
||||
def call
|
||||
raise Error, "Bonifico non configurato sul server" unless MatchLiveTv.bank_transfer_configured?
|
||||
raise Error, "Piano non valido" unless @plan_slug.in?(PLAN_SLUGS)
|
||||
raise Error, "Completa i dati di fatturazione prima di richiedere il bonifico." unless @club.billing_profile_complete?
|
||||
raise Error, "Il piano è un abbonamento omaggio. Contatta il supporto per passarlo a pagamento." if @club.subscription&.admin_comped?
|
||||
|
||||
@interval = Billing::Stripe::PriceCatalog.normalize_interval(@interval)
|
||||
quote = @club.active_billing_quote
|
||||
amount_cents, kind = resolve_amount(quote)
|
||||
|
||||
existing = @club.billing_transfer_orders.awaiting_payment.first
|
||||
if existing
|
||||
if existing.plan_slug == @plan_slug && existing.billing_interval == @interval && existing.amount_cents == amount_cents
|
||||
return existing
|
||||
end
|
||||
|
||||
existing.cancel!
|
||||
end
|
||||
|
||||
order = nil
|
||||
ApplicationRecord.transaction do
|
||||
payment = @club.billing_payments.create!(
|
||||
provider: "bank_transfer",
|
||||
amount_cents: amount_cents,
|
||||
currency: "eur",
|
||||
status: "pending",
|
||||
plan_slug: @plan_slug,
|
||||
description: payment_description(kind, amount_cents)
|
||||
)
|
||||
order = @club.billing_transfer_orders.create!(
|
||||
billing_club_quote: kind == "commercial_quote" ? quote : nil,
|
||||
billing_payment: payment,
|
||||
plan_slug: @plan_slug,
|
||||
billing_interval: @interval,
|
||||
amount_cents: amount_cents,
|
||||
currency: "eur",
|
||||
kind: kind,
|
||||
status: "awaiting_payment",
|
||||
reference_code: generate_reference_code,
|
||||
requested_by_user: @user
|
||||
)
|
||||
end
|
||||
|
||||
# Le istruzioni (IBAN/causale) le manda a mano l'admin da /admin/billing.
|
||||
Rails.logger.info("[BankTransfer] ordine #{order.reference_code} in attesa, nessuna mail istruzioni")
|
||||
order
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def resolve_amount(quote)
|
||||
if quote
|
||||
unless quote.matches?(@plan_slug, @interval)
|
||||
raise Error, "Per questa società è attivo un prezzo concordato su #{quote.plan.name} (#{quote.price_label})."
|
||||
end
|
||||
|
||||
return [quote.amount_cents, "commercial_quote"]
|
||||
end
|
||||
|
||||
[Billing::Stripe::PriceCatalog.amount_cents(plan_slug: @plan_slug, interval: @interval), "list_price"]
|
||||
end
|
||||
|
||||
def payment_description(kind, amount_cents)
|
||||
label = Billing::Stripe::PriceCatalog.format_interval_price(amount_cents, @interval)
|
||||
suffix = kind == "commercial_quote" ? "prezzo concordato, bonifico" : "bonifico"
|
||||
"#{Plan[@plan_slug].name} — #{label} (#{suffix})"
|
||||
end
|
||||
|
||||
def generate_reference_code
|
||||
8.times do
|
||||
code = "MLTV-#{SecureRandom.alphanumeric(6).upcase}"
|
||||
return code unless TransferOrder.exists?(reference_code: code)
|
||||
end
|
||||
|
||||
raise Error, "Impossibile generare il riferimento del bonifico"
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,74 @@
|
||||
module Billing
|
||||
class SetClubQuote
|
||||
class Error < StandardError; end
|
||||
|
||||
PLAN_SLUGS = %w[premium_light premium_full].freeze
|
||||
|
||||
def self.upsert(club:, plan_slug:, interval:, amount_euros:, note:, admin:)
|
||||
new(club: club, plan_slug: plan_slug, interval: interval, amount_euros: amount_euros, note: note, admin: admin).upsert
|
||||
end
|
||||
|
||||
def self.revoke(club:, admin:)
|
||||
new(club: club, admin: admin).revoke
|
||||
end
|
||||
|
||||
def initialize(club:, plan_slug: nil, interval: nil, amount_euros: nil, note: nil, admin: nil)
|
||||
@club = club
|
||||
@plan_slug = plan_slug.to_s.presence
|
||||
@interval = interval
|
||||
@amount_euros = amount_euros
|
||||
@note = note.to_s.strip.presence
|
||||
@admin = admin
|
||||
end
|
||||
|
||||
def upsert
|
||||
raise Error, "Piano non valido" unless @plan_slug.in?(PLAN_SLUGS)
|
||||
|
||||
interval = Billing::Stripe::PriceCatalog.normalize_interval(@interval)
|
||||
amount_cents = EuroAmount.to_cents(@amount_euros)
|
||||
quote = @club.billing_quotes.active.first || @club.billing_quotes.build
|
||||
|
||||
ApplicationRecord.transaction do
|
||||
quote.assign_attributes(
|
||||
plan_slug: @plan_slug,
|
||||
billing_interval: interval,
|
||||
amount_cents: amount_cents,
|
||||
currency: "eur",
|
||||
note: @note,
|
||||
active: true,
|
||||
created_by_admin: @admin || quote.created_by_admin
|
||||
)
|
||||
quote.save!
|
||||
|
||||
cancel_incompatible_orders!(quote)
|
||||
end
|
||||
|
||||
quote
|
||||
rescue EuroAmount::Error, ArgumentError => e
|
||||
raise Error, e.message
|
||||
end
|
||||
|
||||
def revoke
|
||||
quote = @club.active_billing_quote
|
||||
raise Error, "Nessun prezzo concordato attivo" if quote.blank?
|
||||
|
||||
ApplicationRecord.transaction do
|
||||
quote.update!(active: false)
|
||||
@club.billing_transfer_orders.awaiting_payment.where(kind: "commercial_quote").find_each(&:cancel!)
|
||||
end
|
||||
quote
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def cancel_incompatible_orders!(quote)
|
||||
@club.billing_transfer_orders.awaiting_payment.find_each do |order|
|
||||
next if order.plan_slug == quote.plan_slug &&
|
||||
order.billing_interval == quote.billing_interval &&
|
||||
order.amount_cents == quote.amount_cents
|
||||
|
||||
order.cancel!
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -23,10 +23,11 @@ module Billing
|
||||
|
||||
def button_label(plan:, interval:, kind:)
|
||||
price = PriceCatalog.label(plan_slug: plan.slug, interval: interval)
|
||||
name = I18n.t("pages.plans.cta_name.#{plan.slug}", default: plan.name)
|
||||
if kind == :checkout
|
||||
I18n.t("billing.messages.activate_button", plan: plan.name, price: price)
|
||||
I18n.t("billing.messages.activate_button", plan: name, price: price)
|
||||
else
|
||||
I18n.t("billing.messages.switch_button", plan: plan.name, price: price)
|
||||
I18n.t("billing.messages.switch_button", plan: name, price: price)
|
||||
end
|
||||
end
|
||||
|
||||
|
||||
@@ -4,9 +4,16 @@ module Billing
|
||||
INTERVALS = %w[monthly yearly].freeze
|
||||
DEFAULT_INTERVAL = "yearly"
|
||||
|
||||
PRICING = {
|
||||
"premium_light" => { "yearly" => "€40/anno", "monthly" => "€5/mese" },
|
||||
"premium_full" => { "yearly" => "€200/anno", "monthly" => "€20/mese" }
|
||||
# Importi in centesimi (EUR). Il yearly è il prezzo lancio addebitato.
|
||||
AMOUNTS = {
|
||||
"premium_light" => {
|
||||
"yearly" => { charge: 5900, list: 7900, equivalent_monthly: 492 },
|
||||
"monthly" => { charge: 790 }
|
||||
},
|
||||
"premium_full" => {
|
||||
"yearly" => { charge: 19_900, list: 24_900, equivalent_monthly: 1658 },
|
||||
"monthly" => { charge: 2490 }
|
||||
}
|
||||
}.freeze
|
||||
|
||||
class << self
|
||||
@@ -30,7 +37,23 @@ module Billing
|
||||
end
|
||||
|
||||
def label(plan_slug:, interval:)
|
||||
PRICING.dig(plan_slug.to_s, normalize_interval(interval)) || normalize_interval(interval)
|
||||
interval = normalize_interval(interval)
|
||||
format_interval_price(charge_cents(plan_slug, interval), interval)
|
||||
end
|
||||
|
||||
def list_label(plan_slug:, interval: "yearly")
|
||||
interval = normalize_interval(interval)
|
||||
cents = list_cents(plan_slug, interval)
|
||||
return nil if cents.blank?
|
||||
|
||||
format_interval_price(cents, interval)
|
||||
end
|
||||
|
||||
def equivalent_monthly_label(plan_slug:)
|
||||
cents = AMOUNTS.dig(plan_slug.to_s, "yearly", :equivalent_monthly)
|
||||
return nil if cents.blank?
|
||||
|
||||
format_interval_price(cents, "monthly")
|
||||
end
|
||||
|
||||
def available_intervals(plan_slug:)
|
||||
@@ -52,8 +75,45 @@ module Billing
|
||||
nil
|
||||
end
|
||||
|
||||
def format_eur(cents)
|
||||
cents = cents.to_i
|
||||
whole, frac = cents.divmod(100)
|
||||
if frac.zero?
|
||||
"€#{whole}"
|
||||
else
|
||||
sep = I18n.locale.to_s == "en" ? "." : ","
|
||||
"€#{whole}#{sep}#{frac.to_s.rjust(2, '0')}"
|
||||
end
|
||||
end
|
||||
|
||||
def catalog_intervals(plan_slug:)
|
||||
INTERVALS.select { |interval| AMOUNTS.dig(plan_slug.to_s, interval, :charge).present? }
|
||||
end
|
||||
|
||||
def amount_cents(plan_slug:, interval:)
|
||||
interval = normalize_interval(interval)
|
||||
cents = charge_cents(plan_slug, interval)
|
||||
raise ArgumentError, "Prezzo listino non disponibile per #{plan_slug} (#{interval})" if cents.blank?
|
||||
|
||||
cents
|
||||
end
|
||||
|
||||
def format_interval_price(cents, interval)
|
||||
return nil if cents.blank?
|
||||
|
||||
I18n.t("billing.prices.#{interval}", amount: format_eur(cents))
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def charge_cents(plan_slug, interval)
|
||||
AMOUNTS.dig(plan_slug.to_s, interval, :charge)
|
||||
end
|
||||
|
||||
def list_cents(plan_slug, interval)
|
||||
AMOUNTS.dig(plan_slug.to_s, interval, :list)
|
||||
end
|
||||
|
||||
def price_id_for(plan_slug, interval)
|
||||
case [plan_slug, interval]
|
||||
when %w[premium_light monthly]
|
||||
|
||||
@@ -0,0 +1,102 @@
|
||||
module Contacts
|
||||
class SubmitInquiry
|
||||
TOPICS = %w[commercial support privacy].freeze
|
||||
LIMIT = 5
|
||||
WINDOW = 1.hour
|
||||
NAME_MAX = 120
|
||||
CLUB_MAX = 160
|
||||
MESSAGE_MIN = 10
|
||||
MESSAGE_MAX = 4_000
|
||||
|
||||
Result = Struct.new(:ok, :error, keyword_init: true)
|
||||
|
||||
def self.call(params:, ip:)
|
||||
new(params: params, ip: ip).call
|
||||
end
|
||||
|
||||
def initialize(params:, ip:)
|
||||
@params = params
|
||||
@ip = ip.to_s.presence || "unknown"
|
||||
end
|
||||
|
||||
def call
|
||||
if honeypot_filled?
|
||||
record_attempt
|
||||
return Result.new(ok: true, error: nil)
|
||||
end
|
||||
|
||||
return Result.new(ok: false, error: :throttled) if throttled?
|
||||
|
||||
error = validate
|
||||
return Result.new(ok: false, error: error) if error
|
||||
|
||||
record_attempt
|
||||
deliver
|
||||
Result.new(ok: true, error: nil)
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def honeypot_filled?
|
||||
@params[:website].to_s.strip.present?
|
||||
end
|
||||
|
||||
def throttled?
|
||||
current_count >= LIMIT
|
||||
end
|
||||
|
||||
def validate
|
||||
return :privacy_required unless @params[:accept_privacy].to_s == "1"
|
||||
return :invalid if name.blank? || name.length > NAME_MAX
|
||||
return :invalid_email if email.blank? || email !~ URI::MailTo::EMAIL_REGEXP
|
||||
return :invalid if club_name.length > CLUB_MAX
|
||||
return :invalid unless TOPICS.include?(topic)
|
||||
return :invalid_message if message.length < MESSAGE_MIN || message.length > MESSAGE_MAX
|
||||
|
||||
nil
|
||||
end
|
||||
|
||||
def deliver
|
||||
mail = ContactMailer.inquiry(
|
||||
name: name,
|
||||
email: email,
|
||||
club_name: club_name,
|
||||
topic: topic,
|
||||
message: message
|
||||
)
|
||||
MatchLiveTv.deliver_mail(mail)
|
||||
end
|
||||
|
||||
def record_attempt
|
||||
Rails.cache.write(cache_key, current_count + 1, expires_in: WINDOW, raw: true)
|
||||
end
|
||||
|
||||
def current_count
|
||||
Rails.cache.read(cache_key, raw: true).to_i
|
||||
end
|
||||
|
||||
def cache_key
|
||||
"contacts:inquiry:#{@ip}"
|
||||
end
|
||||
|
||||
def name
|
||||
@params[:name].to_s.strip
|
||||
end
|
||||
|
||||
def email
|
||||
@params[:email].to_s.strip.downcase
|
||||
end
|
||||
|
||||
def club_name
|
||||
@params[:club_name].to_s.strip
|
||||
end
|
||||
|
||||
def topic
|
||||
@params[:topic].to_s.strip
|
||||
end
|
||||
|
||||
def message
|
||||
@params[:message].to_s.strip
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -4,7 +4,12 @@ module Mediamtx
|
||||
class Client
|
||||
class Error < StandardError; end
|
||||
|
||||
def self.for_session(session)
|
||||
new(base_url: session.mediamtx_api_base_url)
|
||||
end
|
||||
|
||||
def initialize(base_url: MatchLiveTv.mediamtx_api_url)
|
||||
@base_url = base_url
|
||||
@conn = Faraday.new(url: base_url) do |f|
|
||||
f.request :json
|
||||
f.response :json
|
||||
@@ -12,17 +17,18 @@ module Mediamtx
|
||||
end
|
||||
end
|
||||
|
||||
attr_reader :base_url
|
||||
|
||||
def create_path(session)
|
||||
path = session.mediamtx_path_name
|
||||
# record: false finché non c'è publisher — con alwaysAvailable MediaMTX registrerebbe
|
||||
# solo la slate (schermo nero) in pausa/attesa.
|
||||
# YouTube: niente slate sul path camera (maschera il video al relay ffmpeg).
|
||||
# solo la slate in pausa/attesa. Slate resta accesa anche su YouTube (copertina se l'app cade).
|
||||
body = recording_body(session, enabled: false).merge(
|
||||
source: "publisher",
|
||||
overridePublisher: true
|
||||
)
|
||||
body[:alwaysAvailable] = true
|
||||
body[:alwaysAvailableFile] = slate_file_path
|
||||
body[:alwaysAvailableFile] = slate_file_path(session)
|
||||
# YouTube: telefono → MediaMTX; relay copy verso RTMPS in sidekiq.
|
||||
response = @conn.post("/v3/config/paths/add/#{CGI.escape(path)}", body)
|
||||
unless response.success?
|
||||
@@ -68,14 +74,55 @@ module Mediamtx
|
||||
set_always_available(session, enabled: true)
|
||||
end
|
||||
|
||||
# Slate alwaysAvailable: copertina sullo stesso path quando il telefono è offline.
|
||||
# Disattivare quando il publisher è in onda; riattivare in pausa/disconnessione.
|
||||
# Garantisce path config + slate: se MediaMTX ha perso la config (restart → all_others),
|
||||
# in pausa non resterebbe nessuno stream verso HLS/YouTube.
|
||||
def ensure_path_with_slate!(session)
|
||||
path = session.mediamtx_path_name
|
||||
response = @conn.get("/v3/config/paths/get/#{CGI.escape(path)}")
|
||||
missing = response.status == 404 ||
|
||||
(response.body.is_a?(Hash) && response.body["error"].present?)
|
||||
|
||||
if missing
|
||||
forget_always_available(path)
|
||||
begin
|
||||
create_path(session)
|
||||
rescue Error
|
||||
# Path creato in parallelo o già presente: forza solo la slate.
|
||||
forget_always_available(path)
|
||||
set_always_available(session, enabled: true)
|
||||
end
|
||||
return true
|
||||
end
|
||||
|
||||
conf = response.body.is_a?(Hash) ? response.body : {}
|
||||
desired = slate_file_path(session)
|
||||
if conf["alwaysAvailable"] == true && conf["alwaysAvailableFile"].to_s == desired.to_s
|
||||
remember_always_available(path, enabled: true)
|
||||
return true
|
||||
end
|
||||
|
||||
forget_always_available(path)
|
||||
set_always_available(session, enabled: true)
|
||||
end
|
||||
|
||||
# Spegnere record solo se attivo: ogni PATCH path ricarica MediaMTX e distrugge i muxer HLS.
|
||||
def disable_recording_if_active!(session)
|
||||
path = session.mediamtx_path_name
|
||||
response = @conn.get("/v3/config/paths/get/#{CGI.escape(path)}")
|
||||
return true unless response.success?
|
||||
return true unless response.body.is_a?(Hash) && response.body["record"] == true
|
||||
|
||||
set_path_recording(session, enabled: false)
|
||||
end
|
||||
|
||||
# Slate alwaysAvailable: copertina sullo stesso path quando il telefono è offline.
|
||||
# Resta accesa anche con publisher in onda (MediaMTX usa il publisher se presente).
|
||||
def set_always_available(session, enabled:)
|
||||
path = session.mediamtx_path_name
|
||||
return true if always_available_remembered?(path, enabled: enabled)
|
||||
|
||||
body = if enabled
|
||||
{ alwaysAvailable: true, alwaysAvailableFile: slate_file_path }
|
||||
{ alwaysAvailable: true, alwaysAvailableFile: slate_file_path(session) }
|
||||
else
|
||||
{ alwaysAvailable: false }
|
||||
end
|
||||
@@ -98,6 +145,16 @@ module Mediamtx
|
||||
[]
|
||||
end
|
||||
|
||||
def list_rtmp_conns
|
||||
response = @conn.get("/v3/rtmpconns/list")
|
||||
return [] unless response.success?
|
||||
|
||||
body = response.body
|
||||
body.is_a?(Hash) ? (body["items"] || []) : []
|
||||
rescue Error, Faraday::Error
|
||||
[]
|
||||
end
|
||||
|
||||
def online_path_names
|
||||
Set.new(list_paths.filter_map { |item| item["name"] if item["online"] })
|
||||
end
|
||||
@@ -157,7 +214,13 @@ module Mediamtx
|
||||
SCRIPT
|
||||
end
|
||||
|
||||
def slate_file_path
|
||||
def slate_file_path(session)
|
||||
return default_slate_file_path if session.blank?
|
||||
|
||||
Streams::SlateDistributor.slate_path_for_session(session)
|
||||
end
|
||||
|
||||
def default_slate_file_path
|
||||
ENV.fetch("MEDIAMTX_SLATE_FILE", "/slates/offline.mp4")
|
||||
end
|
||||
|
||||
|
||||
@@ -4,17 +4,37 @@ module Mediamtx
|
||||
module_function
|
||||
|
||||
def active?(session)
|
||||
return true if rtmp_publisher?(session)
|
||||
|
||||
active_path?(path_info(session))
|
||||
end
|
||||
|
||||
def path_info(session)
|
||||
Client.new.list_paths.find { |i| i["name"] == session.mediamtx_path_name }
|
||||
Client.for_session(session).list_paths.find { |i| i["name"] == session.mediamtx_path_name }
|
||||
end
|
||||
|
||||
# MediaMTX <=1.19: online + source.type=rtmpConn.
|
||||
# MediaMTX 1.20+: online/source spesso null anche con publisher; usare rtmpconns.
|
||||
def active_path?(info)
|
||||
return false unless info
|
||||
return true if info["online"] == true && rtmp_source?(info.dig("source", "type"))
|
||||
|
||||
info["online"] == true && info.dig("source", "type") == "rtmpConn"
|
||||
false
|
||||
end
|
||||
|
||||
def rtmp_publisher?(session)
|
||||
path = session.mediamtx_path_name.to_s
|
||||
return false if path.blank?
|
||||
|
||||
Client.for_session(session).list_rtmp_conns.any? do |conn|
|
||||
conn_path = conn["path"].to_s.sub(%r{\A/}, "")
|
||||
next false unless conn_path == path
|
||||
|
||||
state = conn["state"].to_s
|
||||
state.empty? || state == "publish" || state == "idle"
|
||||
end
|
||||
rescue StandardError
|
||||
false
|
||||
end
|
||||
|
||||
def h264_video?(info)
|
||||
@@ -26,12 +46,14 @@ module Mediamtx
|
||||
end
|
||||
|
||||
def video_publishing?(session)
|
||||
info = path_info(session)
|
||||
return false unless active_path?(info)
|
||||
# Slate alwaysAvailable ha H264 ma non è il telefono.
|
||||
return false unless info.dig("source", "type") == "rtmpConn"
|
||||
return false unless active?(session)
|
||||
|
||||
info = path_info(session)
|
||||
h264_video?(info)
|
||||
end
|
||||
|
||||
def rtmp_source?(type)
|
||||
type.to_s.match?(/\Artmps?Conn\z/)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -12,14 +12,13 @@ module Mediamtx
|
||||
return @session if @session.terminal?
|
||||
|
||||
path_info = Mediamtx::PublisherOnline.path_info(@session)
|
||||
publisher_online = Mediamtx::PublisherOnline.active_path?(path_info)
|
||||
publisher_online = Mediamtx::PublisherOnline.active?(@session)
|
||||
|
||||
if publisher_online
|
||||
clear_publisher_misses!(@session.id)
|
||||
if @session.paused?
|
||||
# RTMP ancora connesso in pausa: non forzare live/reconnect.
|
||||
else
|
||||
enable_live_path_once!(@session)
|
||||
if @session.may_go_live?
|
||||
@session.go_live!
|
||||
@session.reload
|
||||
@@ -75,27 +74,13 @@ module Mediamtx
|
||||
@session.update!(timeout_job_id: job)
|
||||
end
|
||||
|
||||
# Compatibilità con webhook / controller.
|
||||
def self.schedule_youtube_pipeline!(session, force: false)
|
||||
Youtube::LivePipeline.schedule!(session, force: force)
|
||||
end
|
||||
|
||||
def enable_live_path_once!(session)
|
||||
return unless session.platform == "youtube"
|
||||
|
||||
key = format("youtube:slate_disabled:%s", session.id)
|
||||
return unless redis.set(key, "1", nx: true, ex: 48.hours.to_i)
|
||||
|
||||
Client.new.set_always_available(session, enabled: false)
|
||||
rescue Client::Error => e
|
||||
redis.del(format("youtube:slate_disabled:%s", session.id))
|
||||
Rails.logger.warn("[PublisherSync] disable slate session=#{session.id}: #{e.message}")
|
||||
end
|
||||
|
||||
def restore_slate_path!(session)
|
||||
return if session.platform == "matchlivetv"
|
||||
|
||||
Client.new.set_always_available(session, enabled: true)
|
||||
# Anche su matchlivetv: dopo restart MediaMTX la path può cadere su all_others senza slate.
|
||||
Client.for_session(session).ensure_path_with_slate!(session)
|
||||
rescue Client::Error => e
|
||||
Rails.logger.warn("[PublisherSync] enable slate session=#{session.id}: #{e.message}")
|
||||
end
|
||||
@@ -128,7 +113,7 @@ module Mediamtx
|
||||
return
|
||||
end
|
||||
|
||||
Client.new.set_path_recording(session, enabled: enabled)
|
||||
Client.for_session(session).set_path_recording(session, enabled: enabled)
|
||||
redis.set(key, desired, ex: 48.hours.to_i)
|
||||
mark_recording_patch!(session.id)
|
||||
rescue Client::Error => e
|
||||
@@ -160,6 +145,14 @@ module Mediamtx
|
||||
nil
|
||||
end
|
||||
|
||||
def self.remember_recording_off!(session_id)
|
||||
redis = Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0"))
|
||||
redis.set("mediamtx:recording:#{session_id}", "0", ex: 48.hours.to_i)
|
||||
redis.set(format("mediamtx:recording:patched_at:%s", session_id), Time.current.to_i, ex: 48.hours.to_i)
|
||||
rescue Redis::BaseError
|
||||
nil
|
||||
end
|
||||
|
||||
def recording_state_key(session_id)
|
||||
"mediamtx:recording:#{session_id}"
|
||||
end
|
||||
@@ -183,7 +176,7 @@ module Mediamtx
|
||||
key = format("youtube_relay:sched:%s", session.id)
|
||||
return unless redis.set(key, "1", nx: true, ex: 10)
|
||||
|
||||
YoutubeRelayEnsureJob.perform_later(session.id)
|
||||
Streams::YoutubeRelay.enqueue_ensure!(session)
|
||||
end
|
||||
|
||||
def redis
|
||||
|
||||
@@ -44,7 +44,8 @@ module Ops
|
||||
check_sidekiq_heartbeat,
|
||||
check_sidekiq_dead,
|
||||
check_http_rails,
|
||||
check_rails_latency
|
||||
check_rails_latency,
|
||||
check_stream_overflow
|
||||
]
|
||||
findings << check_http_public if public_check_due?
|
||||
findings
|
||||
@@ -161,6 +162,58 @@ module Ops
|
||||
fail_finding("garage_storage", "warning", "garage_storage:head", "Garage storage non raggiungibile", e.message)
|
||||
end
|
||||
|
||||
def check_stream_overflow
|
||||
return ok_finding("stream_overflow:skip", "Nodi stream non migrati") unless ActiveRecord::Base.connection.data_source_exists?("stream_nodes")
|
||||
|
||||
orphan_hours = ENV.fetch("STREAM_OVERFLOW_ORPHAN_HOURS", "3").to_i
|
||||
orphans = StreamNode.where.not(slug: "home").where(status: %w[ready draining]).select do |n|
|
||||
n.active_publishers.zero? && n.created_at < orphan_hours.hours.ago
|
||||
end
|
||||
metrics = Streams::Autoscaler.metrics
|
||||
over_budget = !metrics[:within_budget]
|
||||
at_max = metrics[:overflow_nodes] >= metrics[:max_overflow_nodes] && metrics[:free_slots] <= metrics[:soft_free_slots]
|
||||
|
||||
if orphans.any?
|
||||
return Finding.new(
|
||||
kind: "stream_overflow",
|
||||
severity: "warning",
|
||||
healthy: false,
|
||||
title: "Nodi stream overflow idle",
|
||||
message: "#{orphans.size} nodo/i idle da >#{orphan_hours}h: #{orphans.map(&:slug).join(', ')}",
|
||||
metadata: { "slugs" => orphans.map(&:slug) },
|
||||
fingerprint: "stream_overflow:orphan_idle"
|
||||
)
|
||||
end
|
||||
|
||||
if over_budget
|
||||
return Finding.new(
|
||||
kind: "stream_overflow",
|
||||
severity: "warning",
|
||||
healthy: false,
|
||||
title: "Budget overflow streaming",
|
||||
message: "Stima €#{metrics[:estimated_monthly_eur]}/mese > budget €#{metrics[:monthly_budget_eur]}",
|
||||
metadata: metrics.transform_keys(&:to_s),
|
||||
fingerprint: "stream_overflow:budget"
|
||||
)
|
||||
end
|
||||
|
||||
if at_max
|
||||
return Finding.new(
|
||||
kind: "stream_overflow",
|
||||
severity: "warning",
|
||||
healthy: false,
|
||||
title: "Capacità stream al massimo",
|
||||
message: "overflow=#{metrics[:overflow_nodes]}/#{metrics[:max_overflow_nodes]} free_slots=#{metrics[:free_slots]}",
|
||||
metadata: metrics.transform_keys(&:to_s),
|
||||
fingerprint: "stream_overflow:at_max"
|
||||
)
|
||||
end
|
||||
|
||||
ok_finding("stream_overflow:ok", "Overflow streaming OK")
|
||||
rescue StandardError => e
|
||||
fail_finding("stream_overflow", "warning", "stream_overflow:error", "Check overflow fallito", e.message)
|
||||
end
|
||||
|
||||
def check_sidekiq_heartbeat
|
||||
redis = Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0"))
|
||||
last = redis.get(Ops::HealthMonitorJob::HEARTBEAT_KEY).to_i
|
||||
|
||||
@@ -95,7 +95,7 @@ module Recordings
|
||||
def cleanup_mediamtx_path(session)
|
||||
return unless session.status.in?(%w[ended error])
|
||||
|
||||
Mediamtx::Client.new.delete_path(session)
|
||||
Mediamtx::Client.for_session(session).delete_path(session)
|
||||
rescue Mediamtx::Client::Error => e
|
||||
@logger.warn("[Recordings::CleanupLocal] delete_path #{session.id}: #{e.message}")
|
||||
end
|
||||
|
||||
@@ -95,7 +95,7 @@ module Recordings
|
||||
|
||||
def cleanup_mediamtx_path
|
||||
Mediamtx::PublisherSync.forget_recording_state!(@session.id)
|
||||
Mediamtx::Client.new.delete_path(@session)
|
||||
Mediamtx::Client.for_session(@session).delete_path(@session)
|
||||
rescue Mediamtx::Client::Error => e
|
||||
Rails.logger.warn("[Recordings::UploadFromSession] delete_path: #{e.message}")
|
||||
end
|
||||
|
||||
@@ -34,11 +34,21 @@ module Sessions
|
||||
end
|
||||
|
||||
StreamSession.transaction do
|
||||
session.stream_node = Streams::NodeRegistry.allocate!
|
||||
session.save!
|
||||
Scoring::Engine.ensure_score_for(session)
|
||||
mtx = Mediamtx::Client.new
|
||||
mtx.create_path(session)
|
||||
log_event(session, "pairing", { created: true, platform: session.platform })
|
||||
Streams::CoverSlateEnsurer.ensure_for!(@match)
|
||||
Streams::SlateDistributor.ensure_for!(session)
|
||||
Mediamtx::Client.for_session(session).create_path(session)
|
||||
log_event(
|
||||
session,
|
||||
"pairing",
|
||||
{
|
||||
created: true,
|
||||
platform: session.platform,
|
||||
stream_node: session.stream_node&.slug
|
||||
}
|
||||
)
|
||||
end
|
||||
|
||||
if session.platform == "youtube"
|
||||
@@ -46,6 +56,8 @@ module Sessions
|
||||
end
|
||||
|
||||
session
|
||||
rescue Streams::NodeRegistry::NoCapacityError => e
|
||||
raise Teams::EntitlementError.new(e.message, code: "stream_capacity_exhausted")
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
@@ -7,11 +7,13 @@ module Sessions
|
||||
def call
|
||||
cancel_timeout_job
|
||||
@session.pause! if @session.may_pause?
|
||||
Mediamtx::PublisherSync.forget_recording_state!(@session.id)
|
||||
begin
|
||||
mtx = Mediamtx::Client.new
|
||||
mtx.set_path_recording(@session, enabled: false)
|
||||
mtx.set_always_available(@session, enabled: true)
|
||||
mtx = Mediamtx::Client.for_session(@session)
|
||||
# Ricrea/riallinea slate prima di spegnere l'RTMP del telefono.
|
||||
mtx.ensure_path_with_slate!(@session)
|
||||
# PATCH recording solo se era attivo (evita kill muxer HLS a ogni pausa).
|
||||
mtx.disable_recording_if_active!(@session)
|
||||
Mediamtx::PublisherSync.remember_recording_off!(@session.id)
|
||||
rescue Mediamtx::Client::Error => e
|
||||
Rails.logger.warn("[Sessions::Pause] MediaMTX: #{e.message}")
|
||||
end
|
||||
|
||||
@@ -0,0 +1,28 @@
|
||||
module Sessions
|
||||
# Silenzia o riattiva il microfono della camera in onda.
|
||||
# Lo stream continua a inviare AAC silenzioso (YouTube/MediaMTX richiedono la traccia audio).
|
||||
class SetAudioMute
|
||||
def initialize(session, muted:)
|
||||
@session = session
|
||||
@muted = ActiveModel::Type::Boolean.new.cast(muted)
|
||||
end
|
||||
|
||||
def call
|
||||
return @session if @session.audio_muted == @muted
|
||||
|
||||
@session.update!(audio_muted: @muted)
|
||||
action = @muted ? "mute_audio" : "unmute_audio"
|
||||
event_type = @muted ? "audio_muted" : "audio_unmuted"
|
||||
log_event(event_type)
|
||||
SessionChannel.broadcast_message(@session, { type: "command", action: action, muted: @muted })
|
||||
SessionChannel.broadcast_message(@session, { type: "stream_event", event: "audio_muted", muted: @muted })
|
||||
@session
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def log_event(type)
|
||||
@session.stream_events.create!(event_type: type, occurred_at: Time.current, metadata: { muted: @muted })
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -33,7 +33,7 @@ module Sessions
|
||||
end
|
||||
|
||||
def remove_mediamtx_paths!
|
||||
Mediamtx::Client.new.delete_path(@session)
|
||||
Mediamtx::Client.for_session(@session).delete_path(@session)
|
||||
rescue Mediamtx::Client::Error => e
|
||||
Rails.logger.warn("[Sessions::Stop] delete_path #{@session.id}: #{e.message}")
|
||||
end
|
||||
|
||||
@@ -0,0 +1,241 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
module Streams
|
||||
# Scale-out / warm spare / scale-in dei nodi overflow (lab o Hetzner).
|
||||
# Kill-switch: STREAM_AUTOSCALE_ENABLED!=1 OPPURE Redis streams:autoscaler:kill_switch=1.
|
||||
class Autoscaler
|
||||
Result = Struct.new(:actions, :metrics, :skipped, :error, keyword_init: true)
|
||||
|
||||
LOCK_KEY = "streams:autoscaler:lock"
|
||||
KILL_SWITCH_KEY = "streams:autoscaler:kill_switch"
|
||||
|
||||
class << self
|
||||
def enabled?
|
||||
return false if kill_switch_engaged?
|
||||
return false unless ENV["STREAM_AUTOSCALE_ENABLED"] == "1"
|
||||
|
||||
true
|
||||
end
|
||||
|
||||
def kill_switch_engaged?
|
||||
redis_get(KILL_SWITCH_KEY) == "1"
|
||||
end
|
||||
|
||||
def engage_kill_switch!
|
||||
redis_set(KILL_SWITCH_KEY, "1")
|
||||
end
|
||||
|
||||
def clear_kill_switch!
|
||||
redis_del(KILL_SWITCH_KEY)
|
||||
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 allow_cloud?
|
||||
ENV["STREAM_AUTOSCALE_ALLOW_CLOUD"] == "1" && ENV["HCLOUD_TOKEN"].present?
|
||||
end
|
||||
|
||||
def node_eur_per_hour
|
||||
ENV.fetch("STREAM_AUTOSCALE_NODE_EUR_PER_HOUR", "0.015").to_f
|
||||
end
|
||||
|
||||
def monthly_budget_eur
|
||||
ENV.fetch("STREAM_AUTOSCALE_MONTHLY_BUDGET_EUR", "40").to_f
|
||||
end
|
||||
|
||||
def estimated_monthly_eur(overflow_count = nil)
|
||||
count = overflow_count || metrics[:overflow_nodes]
|
||||
# Worst case: nodi sempre accesi 24/7
|
||||
(count * node_eur_per_hour * 24 * 30).round(2)
|
||||
end
|
||||
|
||||
def within_budget?(overflow_count = nil)
|
||||
estimated_monthly_eur(overflow_count) <= monthly_budget_eur
|
||||
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?,
|
||||
env_enabled: ENV["STREAM_AUTOSCALE_ENABLED"] == "1",
|
||||
kill_switch: kill_switch_engaged?,
|
||||
kind: kind,
|
||||
allow_cloud: allow_cloud?,
|
||||
estimated_monthly_eur: estimated_monthly_eur(overflow.size),
|
||||
monthly_budget_eur: monthly_budget_eur,
|
||||
within_budget: within_budget?(overflow.size)
|
||||
}
|
||||
end
|
||||
|
||||
def worker_id
|
||||
ENV.fetch("HOSTNAME", "autoscaler")
|
||||
end
|
||||
|
||||
def redis_get(key)
|
||||
Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0")).get(key)
|
||||
rescue Redis::BaseError
|
||||
nil
|
||||
end
|
||||
|
||||
def redis_set(key, value)
|
||||
Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0")).set(key, value)
|
||||
end
|
||||
|
||||
def redis_del(key)
|
||||
Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0")).del(key)
|
||||
end
|
||||
end
|
||||
|
||||
def initialize(provisioner: nil)
|
||||
@provisioner = provisioner || NodeProvisioner.new
|
||||
@provisioned_this_round = []
|
||||
end
|
||||
|
||||
def reconcile!
|
||||
actions = []
|
||||
m = self.class.metrics
|
||||
|
||||
if need_capacity?(m) && can_provision?(m)
|
||||
provision_overflow!
|
||||
actions << :scale_out
|
||||
m = self.class.metrics
|
||||
elsif need_capacity?(m) && !can_provision?(m)
|
||||
actions << :blocked_capacity
|
||||
Rails.logger.warn("[Streams::Autoscaler] capacity needed but blocked metrics=#{m.inspect}")
|
||||
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)
|
||||
next if @provisioned_this_round.include?(node.id)
|
||||
|
||||
safe_scale_in!(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)
|
||||
return false if m[:overflow_nodes] >= self.class.max_overflow_nodes
|
||||
return false unless self.class.within_budget?(m[:overflow_nodes] + 1)
|
||||
return false if self.class.kind == "cloud" && !self.class.allow_cloud?
|
||||
|
||||
true
|
||||
end
|
||||
|
||||
def provision_overflow!
|
||||
node =
|
||||
case self.class.kind
|
||||
when "cloud"
|
||||
raise NodeProvisioner::Error, "Cloud autoscale disabilitato (STREAM_AUTOSCALE_ALLOW_CLOUD / HCLOUD_TOKEN)" unless self.class.allow_cloud?
|
||||
|
||||
@provisioner.provision_cloud!
|
||||
else
|
||||
@provisioner.provision_lab!
|
||||
end
|
||||
@provisioned_this_round << node.id if node
|
||||
node
|
||||
end
|
||||
|
||||
def safe_scale_in!(node)
|
||||
@provisioner.drain!(node) unless node.status == "draining"
|
||||
node.reload
|
||||
raise NodeProvisioner::BusyError, "sessioni ancora attive" if node.occupying_sessions.exists?
|
||||
raise NodeProvisioner::Error, "idle insufficiente" unless idle_long_enough?(node)
|
||||
|
||||
@provisioner.decommission!(node)
|
||||
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
|
||||
@@ -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,200 @@
|
||||
# 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
|
||||
|
||||
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
|
||||
if path.present?
|
||||
raise Error, "HCLOUD_USER_DATA_FILE non leggibile nel container: #{path}" unless File.file?(path)
|
||||
|
||||
return File.read(path)
|
||||
end
|
||||
|
||||
inline = ENV["HCLOUD_USER_DATA"].presence
|
||||
raise Error, "Manca cloud-init: imposta HCLOUD_USER_DATA_FILE (montato) o HCLOUD_USER_DATA" if inline.blank?
|
||||
|
||||
inline
|
||||
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 = {})
|
||||
tries = 0
|
||||
begin
|
||||
response = conn.get(path) do |req|
|
||||
req.headers.update(auth_headers)
|
||||
req.params.update(params)
|
||||
end
|
||||
unwrap!(response)
|
||||
rescue Faraday::SSLError, Faraday::ConnectionFailed, Faraday::TimeoutError
|
||||
tries += 1
|
||||
raise if tries >= 5
|
||||
|
||||
sleep 2
|
||||
retry
|
||||
end
|
||||
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,77 @@
|
||||
module Streams
|
||||
# Catena effective cover: partita → squadra → società → default Match Live TV.
|
||||
class CoverResolver
|
||||
Result = Struct.new(:record, :source, :cover_url, :slate_local_path, keyword_init: true)
|
||||
|
||||
DEFAULT_COVER_URL = "/images/copertina-canale.png".freeze
|
||||
|
||||
def initialize(match)
|
||||
@match = match
|
||||
@team = match.team
|
||||
@club = @team.club
|
||||
@entitlements = @team.entitlements
|
||||
end
|
||||
|
||||
def effective_cover_url
|
||||
resolve&.cover_url || default_cover_url
|
||||
end
|
||||
|
||||
def cover_source
|
||||
resolve&.source || "default"
|
||||
end
|
||||
|
||||
def resolve
|
||||
return nil unless @entitlements.can_use_custom_cover?
|
||||
|
||||
candidates.each do |record, source|
|
||||
next unless record.cover_image.attached?
|
||||
next unless record.cover_slate.attached?
|
||||
|
||||
path = CoverSlatePaths.local_path_for(record)
|
||||
next if path.blank?
|
||||
|
||||
return Result.new(
|
||||
record: record,
|
||||
source: source,
|
||||
cover_url: blob_path(record.cover_image),
|
||||
slate_local_path: path
|
||||
)
|
||||
end
|
||||
|
||||
nil
|
||||
end
|
||||
|
||||
# Restituisce il record con cover_image nella catena anche se la slate non è ancora su disco.
|
||||
# Usato da CoverSlateEnsurer per sapere cosa generare al go-live.
|
||||
def pending_cover_record
|
||||
return nil unless @entitlements.can_use_custom_cover?
|
||||
|
||||
candidates.each do |record, _source|
|
||||
return record if record.cover_image.attached?
|
||||
end
|
||||
nil
|
||||
end
|
||||
|
||||
def default_cover_url
|
||||
DEFAULT_COVER_URL
|
||||
end
|
||||
|
||||
def default_slate_path
|
||||
ENV.fetch("MEDIAMTX_SLATE_FILE", "/slates/offline.mp4")
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def candidates
|
||||
[
|
||||
[@match, "match"],
|
||||
[@team, "team"],
|
||||
[@club, "club"]
|
||||
]
|
||||
end
|
||||
|
||||
def blob_path(attachment)
|
||||
Rails.application.routes.url_helpers.rails_blob_path(attachment, only_path: true)
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,39 @@
|
||||
module Streams
|
||||
# Garantisce che la slate MP4 esista per il record effettivo della partita (sync al go-live).
|
||||
class CoverSlateEnsurer
|
||||
def self.ensure_for!(match)
|
||||
new(match).ensure!
|
||||
end
|
||||
|
||||
def initialize(match)
|
||||
@match = match
|
||||
@resolver = CoverResolver.new(match)
|
||||
end
|
||||
|
||||
def ensure!
|
||||
record = find_pending_cover_record
|
||||
return unless record
|
||||
return if slate_ready?(record)
|
||||
|
||||
GenerateCoverSlate.call(record)
|
||||
rescue GenerateCoverSlate::Error => e
|
||||
Rails.logger.warn("[CoverSlateEnsurer] match=#{@match.id} #{e.message}")
|
||||
nil
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
# Cerca il record con cover_image nella catena di ereditarietà,
|
||||
# anche se la slate non è ancora pronta (a differenza di resolve che richiede la slate su disco).
|
||||
def find_pending_cover_record
|
||||
@resolver.pending_cover_record
|
||||
end
|
||||
|
||||
def slate_ready?(record)
|
||||
return false unless record.cover_slate.attached?
|
||||
|
||||
path = CoverSlatePaths.local_path_for(record)
|
||||
path.present? && File.exist?(path)
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,35 @@
|
||||
module Streams
|
||||
class CoverSlatePaths
|
||||
def self.slates_custom_root
|
||||
ENV.fetch("MEDIAMTX_SLATES_CUSTOM_DIR", "/slates/custom")
|
||||
end
|
||||
|
||||
# ActiveStorage checksum is base64 e può contenere "/", "+", "=" — non usabile come path.
|
||||
def self.safe_digest(checksum)
|
||||
checksum.to_s.tr("+/", "-_").delete("=")
|
||||
end
|
||||
|
||||
def self.filename_for(record)
|
||||
return nil unless record.cover_slate.attached?
|
||||
|
||||
digest = safe_digest(record.cover_slate.blob.checksum)
|
||||
prefix = record.class.name.underscore
|
||||
"#{prefix}-#{record.id}-#{digest}.mp4"
|
||||
end
|
||||
|
||||
def self.expected_path(record)
|
||||
name = filename_for(record)
|
||||
return nil if name.blank?
|
||||
|
||||
File.join(slates_custom_root, name)
|
||||
end
|
||||
|
||||
def self.local_path_for(record)
|
||||
path = expected_path(record)
|
||||
return nil if path.blank?
|
||||
return path if File.exist?(path)
|
||||
|
||||
nil
|
||||
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,106 @@
|
||||
# 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
|
||||
# Trailing slash obbligatorio: path assoluti altrimenti droppano /v1.
|
||||
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,121 @@
|
||||
require "open3"
|
||||
|
||||
module Streams
|
||||
# Converte cover_image in MP4 slate (H.264 720p + AAC mono), allineato a infra/scripts/generate_slate.sh.
|
||||
class GenerateCoverSlate
|
||||
class Error < StandardError; end
|
||||
|
||||
SLATE_DURATION_SECS = 30
|
||||
VIDEO_FILTER = "scale=1280:720:force_original_aspect_ratio=decrease,pad=1280:720:(ow-iw)/2:(oh-ih)/2:color=0x0a0a0e".freeze
|
||||
|
||||
def self.call(record)
|
||||
new(record).call
|
||||
end
|
||||
|
||||
def initialize(record)
|
||||
@record = record
|
||||
end
|
||||
|
||||
def call
|
||||
return @record unless @record.cover_image.attached?
|
||||
|
||||
@image_temp = download_cover_image
|
||||
@mp4_temp = encode_slate(@image_temp)
|
||||
attach_slate(@mp4_temp)
|
||||
write_to_slates_dir(@mp4_temp)
|
||||
@record
|
||||
rescue Error
|
||||
raise
|
||||
rescue StandardError => e
|
||||
raise Error, e.message
|
||||
ensure
|
||||
cleanup_tempfiles
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def download_cover_image
|
||||
ext = extension_for(@record.cover_image.content_type)
|
||||
path = temp_path("cover-in", ext)
|
||||
@record.cover_image.blob.open(tmpdir: Dir.tmpdir) do |file|
|
||||
FileUtils.cp(file.path, path)
|
||||
end
|
||||
path
|
||||
end
|
||||
|
||||
def encode_slate(image_path)
|
||||
out = temp_path("cover-slate", ".mp4")
|
||||
cmd = ffmpeg_args(image_path, out)
|
||||
stdout, stderr, status = Open3.capture3(*cmd)
|
||||
unless status.success? && File.exist?(out) && File.size(out).positive?
|
||||
detail = stderr.to_s.strip.presence || stdout.to_s.strip
|
||||
raise Error, "ffmpeg slate failed: #{detail}"
|
||||
end
|
||||
|
||||
out
|
||||
end
|
||||
|
||||
def ffmpeg_args(input, output)
|
||||
[
|
||||
"ffmpeg", "-nostdin", "-hide_banner", "-loglevel", "error",
|
||||
"-loop", "1", "-framerate", "30", "-i", input,
|
||||
"-f", "lavfi", "-i", "anullsrc=r=48000:cl=mono",
|
||||
"-filter:v", VIDEO_FILTER,
|
||||
"-map", "0:v", "-map", "1:a",
|
||||
"-t", SLATE_DURATION_SECS.to_s,
|
||||
"-c:v", "libx264", "-pix_fmt", "yuv420p", "-profile:v", "baseline", "-level", "3.1",
|
||||
"-x264-params", "keyint=30:min-keyint=30:scenecut=0:bframes=0",
|
||||
"-g", "30", "-keyint_min", "30", "-force_key_frames", "expr:gte(t,n_forced*1)",
|
||||
"-preset", "fast",
|
||||
"-c:a", "aac", "-b:a", "128k", "-ac", "1", "-shortest", "-y", output
|
||||
]
|
||||
end
|
||||
|
||||
def attach_slate(mp4_path)
|
||||
@record.cover_slate.purge if @record.cover_slate.attached?
|
||||
File.open(mp4_path, "rb") do |io|
|
||||
@record.cover_slate.attach(
|
||||
io: io,
|
||||
filename: slate_filename,
|
||||
content_type: "video/mp4"
|
||||
)
|
||||
end
|
||||
@record.reload
|
||||
end
|
||||
|
||||
def write_to_slates_dir(mp4_path)
|
||||
dest = CoverSlatePaths.expected_path(@record)
|
||||
return unless dest
|
||||
|
||||
FileUtils.mkdir_p(File.dirname(dest))
|
||||
FileUtils.cp(mp4_path, dest)
|
||||
rescue Errno::EACCES, Errno::EROFS => e
|
||||
Rails.logger.info("[GenerateCoverSlate] skip disk copy #{dest}: #{e.class}")
|
||||
end
|
||||
|
||||
def slate_filename
|
||||
"#{@record.class.name.underscore}-#{@record.id}.mp4"
|
||||
end
|
||||
|
||||
def extension_for(content_type)
|
||||
case content_type
|
||||
when "image/png" then ".png"
|
||||
when "image/webp" then ".webp"
|
||||
else ".jpg"
|
||||
end
|
||||
end
|
||||
|
||||
def temp_path(prefix, ext)
|
||||
path = File.join(Dir.tmpdir, "#{prefix}-#{@record.class.name}-#{@record.id}-#{SecureRandom.hex(4)}#{ext}")
|
||||
@tempfiles ||= []
|
||||
@tempfiles << path
|
||||
path
|
||||
end
|
||||
|
||||
def cleanup_tempfiles
|
||||
Array(@tempfiles).each do |path|
|
||||
FileUtils.rm_f(path)
|
||||
end
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,184 @@
|
||||
# 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"])
|
||||
urls = urls_for_node(role: role, home: home, hostname: hostname, simulated: simulated,
|
||||
private_ip: private_ip, public_ip: ip, use_node_hostname: use_node_hostname)
|
||||
|
||||
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: urls.fetch(:rtmp_base_url),
|
||||
hls_base_url: urls.fetch(:hls_base_url),
|
||||
api_base_url: urls.fetch(:api_base_url),
|
||||
internal_rtmp_url: urls.fetch(:internal_rtmp_url),
|
||||
internal_hls_url: urls.fetch(:internal_hls_url),
|
||||
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 urls_for_node(role:, home:, hostname:, simulated:, private_ip:, public_ip: nil, use_node_hostname:)
|
||||
if role == "lab" && ENV["MEDIAMTX_LAB_API_URL"].present?
|
||||
return {
|
||||
api_base_url: ENV.fetch("MEDIAMTX_LAB_API_URL"),
|
||||
internal_rtmp_url: ENV.fetch("MEDIAMTX_LAB_INTERNAL_RTMP_URL", "rtmp://mediamtx_lab:1935"),
|
||||
internal_hls_url: ENV.fetch("MEDIAMTX_LAB_HLS_URL", "http://mediamtx_lab:8888"),
|
||||
rtmp_base_url: ENV.fetch("MEDIAMTX_LAB_RTMP_URL", "rtmp://127.0.0.1:11935"),
|
||||
# HLS pubblico resta sul proxy del sito: il controller instrada al MediaMTX del nodo.
|
||||
hls_base_url: home.hls_base_url
|
||||
}
|
||||
end
|
||||
|
||||
control_ip = control_ip_for(role: role, private_ip: private_ip, public_ip: public_ip)
|
||||
{
|
||||
api_base_url: simulated ? home.api_base_url : "http://#{control_ip}:9997",
|
||||
internal_rtmp_url: simulated ? home.internal_rtmp_url : "rtmp://#{control_ip}:1935",
|
||||
internal_hls_url: simulated ? home.internal_hls_url : "http://#{control_ip}:8888",
|
||||
rtmp_base_url: use_node_hostname ? "rtmp://#{hostname}:1935" : home.rtmp_base_url,
|
||||
# Player del sito: proxy Rails/edge. HTTPS diretto sul nodo arriva dopo TLS/Caddy.
|
||||
hls_base_url: home.hls_base_url
|
||||
}
|
||||
end
|
||||
|
||||
# Senza WireGuard/network privata Rails deve parlare con l'IPv4 pubblico (API+HLS).
|
||||
def control_ip_for(role:, private_ip:, public_ip:)
|
||||
public_ip = public_ip.presence
|
||||
return private_ip if role != "cloud" || public_ip.blank?
|
||||
return public_ip if ENV["STREAM_CLOUD_PUBLIC_CONTROL"] == "1"
|
||||
return public_ip if private_ip.blank? || private_ip == public_ip
|
||||
|
||||
private_ip
|
||||
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,66 @@
|
||||
# frozen_string_literal: true
|
||||
|
||||
module Streams
|
||||
# Assegna un StreamNode a una nuova sessione.
|
||||
# Preferisce home finché ha slot; poi least-loaded tra i nodi overflow 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|
|
||||
attrs = {
|
||||
hostname: home_hostname,
|
||||
role: "home",
|
||||
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
|
||||
}
|
||||
# Non sovrascrivere drain/offline/error (ops) a ogni allocate!
|
||||
attrs[:status] = "ready" if node.new_record? || !%w[draining offline error].include?(node.status)
|
||||
node.assign_attributes(attrs)
|
||||
node.save!
|
||||
end
|
||||
end
|
||||
|
||||
def allocate!
|
||||
ensure_home_from_env!
|
||||
|
||||
# Overflow: riempi prima home; i nodi cloud/lab sono solo quando home è pieno
|
||||
# (altrimenti lo warm spare ruberebbe tutte le sessioni).
|
||||
home = StreamNode.find_by(slug: HOME_SLUG)
|
||||
return home if home&.allocatable?
|
||||
|
||||
node = StreamNode.ready
|
||||
.where.not(slug: HOME_SLUG)
|
||||
.to_a
|
||||
.select(&:allocatable?)
|
||||
.min_by { |n| [n.active_publishers, 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
|
||||
@@ -0,0 +1,105 @@
|
||||
module Streams
|
||||
# Assicura che la slate MP4 custom sia sul nodo assegnato (home/lab locale, CPX via agent).
|
||||
class SlateDistributor
|
||||
class Error < StandardError; end
|
||||
|
||||
def self.slate_path_for_session(session)
|
||||
new(session).mediamtx_slate_path
|
||||
end
|
||||
|
||||
def self.ensure_for!(session)
|
||||
new(session).ensure!
|
||||
end
|
||||
|
||||
def initialize(session)
|
||||
@session = session
|
||||
@match = session.match
|
||||
@node = session.stream_node
|
||||
@resolver = CoverResolver.new(@match)
|
||||
end
|
||||
|
||||
def mediamtx_slate_path
|
||||
ensure! unless defined?(@mediamtx_path)
|
||||
@mediamtx_path
|
||||
end
|
||||
|
||||
def ensure!
|
||||
result = @resolver.resolve
|
||||
unless result
|
||||
@mediamtx_path = default_slate_path
|
||||
return @mediamtx_path
|
||||
end
|
||||
|
||||
host_path = result.slate_local_path
|
||||
filename = CoverSlatePaths.filename_for(result.record)
|
||||
raise Error, "slate filename missing" if filename.blank?
|
||||
|
||||
# Path completo (non File.basename): il checksum AS può contenere "/" e spezzerebbe il path.
|
||||
mediamtx_path = CoverSlatePaths.expected_path(result.record)
|
||||
|
||||
if cloud_node?
|
||||
push_to_agent!(filename, result.record, host_path)
|
||||
else
|
||||
verify_local_slate!(host_path)
|
||||
end
|
||||
|
||||
@mediamtx_path = mediamtx_path
|
||||
rescue StandardError => e
|
||||
# Agent cloud down / timeout / disk missing: non bloccare il go-live, usa slate default.
|
||||
Rails.logger.warn("[SlateDistributor] session=#{@session.id} #{e.class}: #{e.message}")
|
||||
@mediamtx_path = default_slate_path
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def cloud_node?
|
||||
@node&.role == "cloud" && agent_url.present?
|
||||
end
|
||||
|
||||
def agent_url
|
||||
ENV["STREAM_NODE_RELAY_AGENT_URL"].presence || @node&.relay_agent_url
|
||||
end
|
||||
|
||||
def agent_secret
|
||||
ENV["STREAM_NODE_AGENT_SECRET"].presence || "mediamtx_webhook_dev_secret"
|
||||
end
|
||||
|
||||
def push_to_agent!(filename, record, host_path)
|
||||
bytes = read_slate_bytes(record, host_path)
|
||||
raise Error, "slate bytes missing for #{filename}" if bytes.blank?
|
||||
|
||||
uri = URI.parse("#{agent_url.chomp('/')}/slates/#{CGI.escape(filename)}")
|
||||
http = Net::HTTP.new(uri.host, uri.port)
|
||||
http.open_timeout = 5
|
||||
http.read_timeout = 120
|
||||
req = Net::HTTP::Put.new(uri)
|
||||
req["Authorization"] = "Bearer #{agent_secret}" if agent_secret.present?
|
||||
req["Content-Type"] = "video/mp4"
|
||||
req.body = bytes
|
||||
res = http.request(req)
|
||||
return if res.is_a?(Net::HTTPSuccess)
|
||||
|
||||
raise Error, "agent PUT /slates/#{filename} HTTP #{res.code} #{res.body.to_s.truncate(200)}"
|
||||
end
|
||||
|
||||
def read_slate_bytes(record, host_path)
|
||||
return File.binread(host_path) if host_path.present? && File.exist?(host_path)
|
||||
|
||||
return unless record.cover_slate.attached?
|
||||
|
||||
record.cover_slate.blob.open(tmpdir: Dir.tmpdir) do |file|
|
||||
return File.binread(file.path)
|
||||
end
|
||||
end
|
||||
|
||||
def verify_local_slate!(host_path)
|
||||
return if host_path.present? && File.exist?(host_path)
|
||||
|
||||
raise Error, "slate missing on disk: #{host_path}"
|
||||
end
|
||||
|
||||
def default_slate_path
|
||||
ENV.fetch("MEDIAMTX_SLATE_FILE", "/slates/offline.mp4")
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -1,29 +1,85 @@
|
||||
require "json"
|
||||
require "net/http"
|
||||
require "uri"
|
||||
|
||||
module Streams
|
||||
# Relay verso YouTube: legge RTMP/HLS da MediaMTX e inoltra su RTMPS (-c copy). Nessun overlay.
|
||||
# ffmpeg gira solo nel container sidekiq (YOUTUBE_RELAY_WORKER=1).
|
||||
# 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 unless worker?
|
||||
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)
|
||||
clear_pid(session.id)
|
||||
redis.del(format(OWNER_KEY, session.id))
|
||||
if worker?
|
||||
if worker? && owner_is_local?(session.id) && can_process?(session)
|
||||
stop_on_worker!(session)
|
||||
else
|
||||
YoutubeRelayStopJob.perform_later(session.id)
|
||||
enqueue_stop!(session)
|
||||
end
|
||||
true
|
||||
end
|
||||
@@ -31,11 +87,13 @@ module Streams
|
||||
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) if owner.blank? || owner == worker_id
|
||||
return process_alive?(pid, session_id: session_id) if owner.blank? || owner == worker_id
|
||||
|
||||
# Relay avviato in un altro container: consideralo attivo se il lock è recente.
|
||||
# Relay su altro host: attivo se lock owner ancora fresco.
|
||||
redis.ttl(format(OWNER_KEY, session_id)) > 30
|
||||
end
|
||||
|
||||
@@ -44,15 +102,22 @@ module Streams
|
||||
return if session.terminal?
|
||||
return if session.stream_key.blank?
|
||||
|
||||
if worker?
|
||||
if can_process?(session)
|
||||
ensure_on_worker!(session)
|
||||
else
|
||||
YoutubeRelayEnsureJob.perform_later(session.id)
|
||||
enqueue_ensure!(session)
|
||||
end
|
||||
end
|
||||
|
||||
def ensure_on_worker!(session)
|
||||
return unless worker?
|
||||
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?
|
||||
@@ -60,28 +125,72 @@ module Streams
|
||||
return unless intake_available?(session)
|
||||
|
||||
pid = pid_for(session.id)
|
||||
clear_pid(session.id) if pid.present? && !process_alive?(pid.to_i)
|
||||
if pid.present? && !process_alive?(pid.to_i, session_id: session.id)
|
||||
clear_local_ownership(session.id)
|
||||
end
|
||||
|
||||
return if running?(session.id)
|
||||
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 if last_restart.positive? && (Time.now.to_i - last_restart) < 5
|
||||
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)
|
||||
pid = pid_for(session.id)
|
||||
return false if pid.blank?
|
||||
return :not_worker unless worker?
|
||||
|
||||
terminate_pid(pid)
|
||||
clear_pid(session.id)
|
||||
redis.del(format(OWNER_KEY, session.id))
|
||||
Rails.logger.info("[YoutubeRelay] stopped pid=#{pid} session=#{session.id}")
|
||||
true
|
||||
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
|
||||
@@ -92,30 +201,98 @@ module Streams
|
||||
return if session.terminal?
|
||||
return unless intake_available?(session)
|
||||
|
||||
return pid_for(session.id).to_i if running?(session.id)
|
||||
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?
|
||||
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 = Process.spawn(
|
||||
*youtube_ffmpeg_args(intake_source, output),
|
||||
%i[out err] => log_path,
|
||||
pgroup: true
|
||||
)
|
||||
Process.detach(pid)
|
||||
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)
|
||||
redis.set(format(OWNER_KEY, session.id), worker_id, ex: 48.hours.to_i)
|
||||
Rails.logger.info("[YoutubeRelay] started pid=#{pid} session=#{session.id} intake=#{intake_source.join(":")}")
|
||||
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)
|
||||
@@ -135,6 +312,9 @@ module Streams
|
||||
end
|
||||
|
||||
common_head + [
|
||||
"-reconnect", "1",
|
||||
"-reconnect_streamed", "1",
|
||||
"-reconnect_delay_max", "2",
|
||||
"-rw_timeout", "15000000",
|
||||
"-live_start_index", "-1",
|
||||
"-i", url,
|
||||
@@ -146,20 +326,35 @@ module Streams
|
||||
end
|
||||
|
||||
def mediamtx_intake_source(session)
|
||||
base = ENV.fetch("MEDIAMTX_INTERNAL_RTMP_URL", "rtmp://mediamtx:1935")
|
||||
rtmp_base, hls_base = intake_bases(session)
|
||||
if Mediamtx::PublisherOnline.active?(session)
|
||||
return [:rtmp, "#{base.chomp('/')}/#{session.mediamtx_path_name}"]
|
||||
return [:rtmp, "#{rtmp_base.chomp('/')}/#{session.mediamtx_path_name}"]
|
||||
end
|
||||
|
||||
path = session.mediamtx_path_name
|
||||
hls = ENV.fetch("MEDIAMTX_HLS_URL", "http://mediamtx:8888").chomp("/")
|
||||
[:hls, "#{hls}/#{path}/index.m3u8"]
|
||||
[: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.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
|
||||
@@ -185,6 +380,22 @@ module Streams
|
||||
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
|
||||
@@ -197,7 +408,10 @@ module Streams
|
||||
redis.get(format(REDIS_KEY, session_id))
|
||||
end
|
||||
|
||||
def process_alive?(pid)
|
||||
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"
|
||||
|
||||
@@ -207,12 +421,27 @@ module Streams
|
||||
false
|
||||
end
|
||||
|
||||
def terminate_pid(pid)
|
||||
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
|
||||
|
||||
@@ -35,6 +35,20 @@ module Teams
|
||||
active? && plan.phone_download_enabled?
|
||||
end
|
||||
|
||||
def can_use_custom_cover?
|
||||
premium_full? && plan.custom_cover_enabled?
|
||||
end
|
||||
|
||||
def assert_can_upload_cover!
|
||||
unless can_use_custom_cover?
|
||||
raise EntitlementError.new(
|
||||
"La copertina sponsor personalizzata richiede Premium Full",
|
||||
code: "premium_full_required",
|
||||
billing_url: billing_url
|
||||
)
|
||||
end
|
||||
end
|
||||
|
||||
def max_staff_transmission
|
||||
plan.max_staff_transmission
|
||||
end
|
||||
@@ -225,7 +239,8 @@ module Teams
|
||||
recording_retention_days: recording_retention_days,
|
||||
phone_download_enabled: phone_download_enabled?,
|
||||
youtube_enabled: youtube_enabled?,
|
||||
youtube_mode: plan.youtube_mode
|
||||
youtube_mode: plan.youtube_mode,
|
||||
custom_cover_enabled: can_use_custom_cover?
|
||||
}
|
||||
end
|
||||
|
||||
|
||||
@@ -1,13 +1,4 @@
|
||||
module Teams
|
||||
class StaffAssignmentError < StandardError
|
||||
attr_reader :message
|
||||
|
||||
def initialize(message)
|
||||
@message = message
|
||||
super(message)
|
||||
end
|
||||
end
|
||||
|
||||
class StaffAssignment
|
||||
def self.call(team:, user:, staff_kind: "transmission", membership: nil)
|
||||
new(team: team, user: user, staff_kind: staff_kind, membership: membership).call
|
||||
@@ -32,32 +23,12 @@ module Teams
|
||||
ut.update!(staff_kind: "transmission")
|
||||
ut
|
||||
end
|
||||
end
|
||||
|
||||
class StaffEmailValidator
|
||||
def self.assert_available!(team:, email:, except_user: nil, staff_kind: nil)
|
||||
new(team: team, email: email, except_user: except_user).assert_available!
|
||||
end
|
||||
|
||||
def initialize(team:, email:, except_user: nil)
|
||||
@team = team
|
||||
@email = email.to_s.downcase.strip
|
||||
@except_user = except_user
|
||||
end
|
||||
|
||||
def assert_available!
|
||||
user_scope = @team.user_teams.joins(:user).where(users: { email: @email }).where.not(staff_kind: nil)
|
||||
user_scope = user_scope.where.not(user_id: @except_user.id) if @except_user
|
||||
if user_scope.exists?
|
||||
raise StaffAssignmentError,
|
||||
"L'email #{@email} è già assegnata come responsabile trasmissione per questa squadra."
|
||||
end
|
||||
|
||||
inv = @team.team_invitations.pending.where("LOWER(email) = ?", @email)
|
||||
return unless inv.exists?
|
||||
|
||||
raise StaffAssignmentError,
|
||||
"Esiste già un invito in sospeso per #{@email}."
|
||||
def self.designate_club_owner!(team:, user:)
|
||||
membership = user.user_teams.find_or_initialize_by(team: team)
|
||||
membership.role = "member" if membership.new_record?
|
||||
membership.save!
|
||||
call(team: team, user: user, membership: membership)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -0,0 +1,10 @@
|
||||
module Teams
|
||||
class StaffAssignmentError < StandardError
|
||||
attr_reader :message
|
||||
|
||||
def initialize(message)
|
||||
@message = message
|
||||
super(message)
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,28 @@
|
||||
module Teams
|
||||
class StaffEmailValidator
|
||||
def self.assert_available!(team:, email:, except_user: nil, staff_kind: nil)
|
||||
new(team: team, email: email, except_user: except_user).assert_available!
|
||||
end
|
||||
|
||||
def initialize(team:, email:, except_user: nil)
|
||||
@team = team
|
||||
@email = email.to_s.downcase.strip
|
||||
@except_user = except_user
|
||||
end
|
||||
|
||||
def assert_available!
|
||||
user_scope = @team.user_teams.joins(:user).where(users: { email: @email }).where.not(staff_kind: nil)
|
||||
user_scope = user_scope.where.not(user_id: @except_user.id) if @except_user
|
||||
if user_scope.exists?
|
||||
raise StaffAssignmentError,
|
||||
"L'email #{@email} è già assegnata come responsabile trasmissione per questa squadra."
|
||||
end
|
||||
|
||||
inv = @team.team_invitations.pending.where("LOWER(email) = ?", @email)
|
||||
return unless inv.exists?
|
||||
|
||||
raise StaffAssignmentError,
|
||||
"Esiste già un invito in sospeso per #{@email}."
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -68,13 +68,13 @@ module Webhooks
|
||||
def enable_recording(session)
|
||||
return unless session.match.team.entitlements.recording_enabled_for_mediamtx?
|
||||
|
||||
Mediamtx::Client.new.set_path_recording(session, enabled: true)
|
||||
Mediamtx::Client.for_session(session).set_path_recording(session, enabled: true)
|
||||
rescue Mediamtx::Client::Error => e
|
||||
Rails.logger.warn("[MediamtxHandler] enable recording: #{e.message}")
|
||||
end
|
||||
|
||||
def disable_recording(session)
|
||||
Mediamtx::Client.new.set_path_recording(session, enabled: false)
|
||||
Mediamtx::Client.for_session(session).set_path_recording(session, enabled: false)
|
||||
rescue Mediamtx::Client::Error => e
|
||||
Rails.logger.warn("[MediamtxHandler] disable recording: #{e.message}")
|
||||
end
|
||||
|
||||
@@ -32,7 +32,8 @@ module Youtube
|
||||
end
|
||||
|
||||
def team_channel_available?
|
||||
@ent.youtube_enabled? && @ent.premium_full? && @mode == "team" && @team.club.youtube_credential.present?
|
||||
@ent.youtube_enabled? && @ent.premium_full? && @mode == "team" &&
|
||||
@team.club.youtube_credential&.usable?
|
||||
end
|
||||
|
||||
def effective_channel
|
||||
|
||||
@@ -63,7 +63,6 @@ module Youtube
|
||||
return
|
||||
end
|
||||
|
||||
Mediamtx::Client.new.set_always_available(session, enabled: false)
|
||||
session.go_live! if session.may_go_live?
|
||||
session.reconnect! if session.reconnecting? && session.may_reconnect?
|
||||
|
||||
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user