Compare commits

..
Author SHA1 Message Date
eminuxandCursor 79850dfc2c Non far crashare il reset password se SMTP rifiuta il destinatario.
Un account di test su dominio .test faceva 500; ora l'invio fallito viene loggato e l'API resta ok.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-31 12:00:22 +02:00
eminuxandCursor 645807b853 Evita 500 senza SMTP e chiude gli incidenti overflow stale.
Reset password e mail replay usano deliver_mail; il health check overflow risolve tutte le fingerprint del kind, non solo quella sana.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-31 08:55:36 +02:00
eminuxandCursor b08049cf69 Fix home 500: usa request.cookie_jar nella partial analytics suppress.
L'action cookies di PagesController mascherava il helper cookies nelle viste marketing.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-28 19:31:05 +02:00
eminuxandCursor 29c8afb94d Admin analytics: heatmap per device, anteprima staff e analisi costi.
Separa heatmap mobile/desktop, opt-out analytics per operatori, dashboard costi con KPI e trend mensili.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-28 19:25:38 +02:00
eminuxandCursor e0e07e5316 Non forza enableEmbed in creazione live: YouTube lo rifiuta sul canale piattaforma.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-28 15:15:35 +02:00
eminuxandCursor 42dadb993d Abilita l'embed YouTube sulle dirette e tiene la copia locale se il VOD non è incorporabile.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-28 15:05:12 +02:00
eminuxandCursor 05ef56c56d Aggiunge replay YouTube temporaneo, pull registrazioni dai CPX e snapshot ingest in admin.
Così overflow Hetzner e VOD YouTube restano in archivio dopo lo spegnimento del nodo, e la colonna ingest non si svuota.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-28 13:01:41 +02:00
eminuxandCursor 35dfa923e3 Usa packageName dinamico negli E2E Android per la variante collaudo.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-26 10:22:20 +02:00
eminuxandCursor 9d8b35c06c Salva telemetria client (OS, app, device, operatore) sulle sessioni per il debug admin.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-26 09:43:43 +02:00
eminuxandCursor 873e0ea55c Fix race analytics: retry su unique violation in aggregazione heatmap.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-21 16:23:02 +02:00
eminuxandCursor 4573edfedc Aggiunge i link ai profili social ufficiali nel footer e in Contatti.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-21 16:10:18 +02:00
eminuxandCursor b301868774 Aggiunge regola Cursor per scalare i soft limit ingest col crescere dei clienti.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-21 15:11:25 +02:00
eminuxandCursor b87a0f7bc0 Fix race create sessione su CPX: ready solo dopo MediaMTX.
I nodi cloud restano in provisioning finché :9997 risponde; retry su create_path e 503 retryable se l’ingest è ancora irraggiungibile.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-21 14:06:47 +02:00
eminuxandCursor 70fccc493e Usa screenshot di pagina come sfondo heatmap al posto dell’iframe.
Cattura periodica client-side (con consenso), storage ActiveStorage e overlay allineato all’immagine.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 23:14:47 +02:00
eminuxandCursor 406a90ae5b Bumpe site-analytics.js a v=3 per bypassare la cache NPM su collaudo.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 23:08:10 +02:00
eminuxandCursor 8e7bbdf8b2 Bumpe cache-bust della heatmap JS su collaudo (evita 404 NPM su ?v=1).
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 23:07:07 +02:00
eminuxandCursor b84321346f Upgrade heatmaps stile Hotjar: movimenti mouse e anteprima pagina live.
Traccia i movimenti aggregati, overlay a gradienti sull’iframe della pagina e toggle click/move in admin.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 23:03:12 +02:00
eminuxandCursor 8a07c5d569 Aggiunge heatmaps first-party (click/scroll) con pagina admin Analytics.
Raccoglie eventi aggregati dietro consenso cookie, senza PII né tracking su admin/regia/live/replay; GA4 resta invariato.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 22:50:15 +02:00
eminuxandCursor cc96c0396a Migliora sessioni admin e rende robusto il logout web.
Aggiunge filtri/colonne (società, orari) e dettaglio leggibile; evita 422 CSRF su logout e non cancella le cover slate custom in sync prod.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 20:41:12 +02:00
eminuxandCursor 79c48b226e Copia AAB/APK di release in native/android/dist con nome versionato.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 16:53:53 +02:00
eminuxandCursor c9b75065bc Bump Android a 2.0.12-native (versionCode 33) per Play Store.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 16:50:20 +02:00
eminuxandCursor 4ca9f3af8b Merge branch 'collaudo': avvisi multilanguage e allineamento sessioni admin.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 16:45:46 +02:00
eminuxandCursor f02782d0f0 Aggiunge traduzioni it/en/fr/de/es agli avvisi e toglie gli override Android/iOS.
Un avviso ha un testo per lingua; per differenziare le piattaforme si creano due avvisi. L'app riceve già il copy risolto da Accept-Language.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 16:26:30 +02:00
eminuxandCursor fb82b74e92 Allinea verticalmente le righe della tabella sessioni admin.
Il flex era sulla td e spezzava vertical-align; ora sta nel contenuto e l'ingest è su una riga.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 16:24:18 +02:00
eminuxandCursor 60dfe7a75b Allinea verticalmente le righe della tabella sessioni admin.
Il flex era sulla td e spezzava vertical-align; ora sta nel contenuto e l'ingest è su una riga.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 16:21:03 +02:00
eminuxandCursor a07f231417 Permette titolo, testo e link diversi per Android e iOS negli avvisi.
Separa i canali app e mostra il nodo ingest nelle sessioni admin.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 16:09:04 +02:00
eminuxandCursor 43d5003de8 Completa UI e test del form contatti.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 15:39:11 +02:00
eminuxandCursor eadfc29da3 Corregge YAML form contatti: quota form_lead con due punti.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 15:37:45 +02:00
eminuxandCursor 50d3e2cf16 Completa traduzioni privacy, mailer e flash del form contatti.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 15:31:52 +02:00
eminuxandCursor 55e71ca838 Allinea privacy DE/ES/FR al form contatti.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 15:31:07 +02:00
eminuxandCursor 5b723bf844 Aggiunge annunci in-app/sito, contatti e migliora UI copertina sul portale.
Include admin/API/banner nativi, form contatti e reset copertina standard con layout edit più ampio.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 15:30:49 +02:00
eminuxandCursor e502844b88 fix: non bloccare go-live se agent slate Hetzner è down
Se il push della cover custom all'agent fallisce (timeout/connection refused),
usa lo slate di default invece di abortire la creazione sessione.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 13:26:52 +02:00
eminuxandCursor 0702cb699c Rendi affidabile pausa con slate e conferma i controlli laterali.
Evita buchi HLS/YouTube in pausa (path slate ricreata, reload player, no patch recording inutili) e richiede conferma su tasti laterali app/regia.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 12:44:23 +02:00
eminuxandCursor e09b3ffcf8 Corregge validazione upload copertina e generazione slate al go-live.
Rifiuta file troppo grandi o non-immagine via Marcel prima di attach, e genera la slate in sync anche quando cover_image esiste ma cover_slate non è ancora pronta.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-20 09:51:33 +02:00
eminuxandCursor 41d4235258 Aggiunge copertina sponsor custom solo Premium Full, con override in wizard e slate sui nodi di streaming.
Permette upload da portale e app con ereditarietà partita→squadra→società, generazione MP4 via ffmpeg e distribuzione su home lab e nodi CPX per pause e assenza segnale.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-19 23:47:17 +02:00
eminuxandCursor a51d5e8da5 Non inviare in automatico le istruzioni di bonifico.
IBAN e causale restano da comunicare a mano dall'admin; la richiesta registra solo l'ordine in attesa.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-19 18:40:45 +02:00
eminuxandCursor 5bcb4170c7 Mostra il piano attivo dopo il bonifico e allinea il flusso società.
Dopo la conferma admin la card restava su «Paga con bonifico»; la società parte da Free, l'owner è staff trasmissione e le mail non bloccano se manca SMTP.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-19 18:31:39 +02:00
eminuxandCursor e436f5ece4 Aggiunge pagamento con bonifico e importo concordato per società.
Il piano si attiva solo dopo la conferma admin; Stripe resta a listino se non c'è un prezzo commerciale.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-18 23:20:43 +02:00
eminuxandCursor 1e4e0560f1 Allinea le spec alla policy password e prepara i piani di test.
Senza questo la suite locale fallisce sulla complessità password e sui piani mancanti.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-18 23:20:31 +02:00
eminuxandCursor 841a6a15be Avvia i relay YouTube overflow dal Sidekiq home e passa STREAM/HCLOUD in produzione.
Così un CPX nuovo non richiede worker Docker a mano, l'agent è nel cloud-init e gli E2E restano non in elenco.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-18 08:57:56 +02:00
eminuxandCursor dd94bf7e66 Riprova le GET Hetzner su errori SSL transitori e aggiunge i loghi con nome.
Durante il wait del CPX l'API chiudeva la connessione TLS; i tre PNG servono le varianti dark/white del marchio.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-17 22:56:46 +02:00
eminuxandCursor baa6283a15 Mantiene la copertina YouTube se il telefono cade e avvia ffmpeg sul nodo, non via SSH.
Il relay HLS→RTMPS resta sul worker/agent del nodo assegnato, così tre dirette contemporanee restano in onda sul sito e sul canale della società senza schermo nero.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-16 20:41:46 +02:00
eminux bfb3115a64 Merge branch 'main' into collaudo
Porta in collaudo l'invito trasmissione con email automatica.
2026-08-15 20:02:40 +02:00
eminuxandCursor db8b14ba4e Snellisci l'invito a trasmettere e invia l'email automatica.
Il form invito è sulla pagina squadra, genera il link e manda la mail; in caso di errore SMTP resta la copia manuale.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-15 20:02:25 +02:00
eminux 46244fe066 Merge branch 'main' into collaudo
Porta su collaudo i prezzi di lancio, la pagina Funzionalità e i badge store.
2026-08-14 23:09:20 +02:00
eminuxandCursor 3ed70b6999 Aggiorna marketing: prezzi lancio, Funzionalità narrativa e link store.
Allinea i piani Light/Full ai prezzi di lancio, ripensa /funzionalita per vendere il prodotto e aggiunge i badge App Store/Play nella home.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-14 23:00:17 +02:00
eminuxandCursor b05b1acd7a Usa l'hostname pubblico per l'ingest RTMP di collaudo.
I telefoni in 4G non raggiungono l'IP LAN; la NAT sulla 11935 espone collaudo.matchlivetv.it.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-13 23:26:34 +02:00
eminuxandCursor c777853c97 Rende autonomo lo spec Hetzner fornendo cloud-init inline.
Così il test di create_node non dipende dal file montato in collaudo.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-13 20:24:37 +02:00
eminuxandCursor 9225318e1a Aggiunge il flavor Android collaudo installabile accanto all'app di produzione.
Usa un applicationId distinto, icona con fascia gialla e API su collaudo.matchlivetv.it.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-13 20:19:20 +02:00
eminuxandCursor 061b7fad8b Merge branch 'feature/adaptive-bitrate' into collaudo
Integra bitrate adattivo e telemetria, con il mute audio come icona nella toolbar della regia.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-13 20:16:26 +02:00
eminuxandCursor be7957c078 Merge branch 'feature/stream-audio-mute' into collaudo
Integra il mute del microfono in diretta per il collaudo.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-13 20:10:27 +02:00
eminux 91667e1bfd Merge branch 'feature/streaming-autoscale-hetzner' into collaudo
Integra l'autoscaler streaming Hetzner nell'ambiente di collaudo.
2026-08-13 20:10:11 +02:00
eminuxandCursor 1fecccaafd Aggiunge il mute del microfono in diretta per evitare claim YouTube sulla musica in palestra.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-12 22:38:05 +02:00
eminux d100c8030f Merge branch 'main' into feature/streaming-autoscale-hetzner 2026-08-11 19:45:27 +02:00
eminuxandCursor f1906517a9 Rende leggibile la pagina admin Nodi streaming con KPI e toolbar.
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-11 19:04:12 +02:00
eminuxandCursor 00fde25c50 Evita scale-in del nodo appena creato nello stesso reconcile.
Con IDLE_MINUTES=0 (o race stretta) lo scale-out veniva annullato subito; lo smoke AutoscalerJob Cloud su collaudo ora completa scale-out e scale-in.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-11 18:42:35 +02:00
eminuxandCursor 4c5af99efa Allinea slate Cloud a AAC 48k mono e preferisce home in allocate.
Senza questo il telefono (mono 48k) veniva rifiutato sui nodi Hetzner e lo warm spare rubava le sessioni a home.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-11 18:27:13 +02:00
eminuxandCursor 8bf12e7721 Bump versione Android a 2.0.11-native (versionCode 32).
Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-11 18:10:08 +02:00
eminuxandCursor 5a0ca8ef1c Corregge cloud-init montato, PublisherOnline multi-nodo e E2E locale-aware.
Senza volume stream-node i nodi Hetzner nascevano senza MediaMTX; sync live usa ora Client.for_session e rtmpconns (MediaMTX 1.20).

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-11 08:48:28 +02:00
eminuxandCursor 714602f3da Aggiunge overlay Compose per collaudo con RTMP su porta 11935.
Isola MediaMTX dal :1935 di produzione sullo stesso IP pubblico e documenta env/helper per il container Proxmox di test.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-10 20:29:47 +02:00
eminuxandCursor baac3a512b Preserva drain/offline su home e genera slate offline nel cloud-init.
Evita che ensure_home_from_env! annulli il drain a ogni allocate, e prepara /slates/offline.mp4 sul nodo Cloud per create_path MediaMTX.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-10 12:29:23 +02:00
eminuxandCursor 559284f0b2 Corregge il provisioning Hetzner: Faraday, cloud-init e default cpx12.
Sistemati path API /v1, AppArmor/auth MediaMTX sul nodo e defaults nbg1 così lo smoke Cloud è ripetibile.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-10 12:20:36 +02:00
eminuxandCursor e5fc925bae Completa il hardening autoscaler (Fase 4): budget, kill-switch e runbook.
Drain sicuro in scale-in, alert Ops su overflow e controlli admin senza abilitare il deploy in produzione.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-09 20:00:15 +02:00
eminuxandCursor 318a319608 Aggiunge l'autoscaler streaming con warm spare e kill-switch.
Scala i nodi overflow quando gli slot calano, mantiene una spare a caldo sotto carico e spegne gli idle, disattivabile con STREAM_AUTOSCALE_ENABLED.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-09 19:56:58 +02:00
eminuxandCursor 0b55d5a8fa Rende sticky e capacizzati i relay YouTube su coda dedicata.
Evita doppi ffmpeg tra host: stop solo sull'owner, ensure con requeue a capacità piena e coda Sidekiq youtube_relay isolata.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-09 19:54:04 +02:00
eminuxandCursor c3878cdc6d Aggiunge registry nodi stream e provisioning Hetzner/lab per lo scale-out.
Prepara l'architettura multi-nodo (MediaMTX+ffmpeg) con assignment URL per sessione, admin di provision/drain e provider Cloud/DNS astratti verso mltv-stream.net.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-09 19:52:47 +02:00
473 changed files with 24075 additions and 1118 deletions
@@ -0,0 +1,22 @@
---
description: Scalare soft limit ingest/autoscaler man mano che crescono i clienti
alwaysApply: true
---
# Soft limit stream / autoscaler
Con la crescita del numero di clienti (e delle dirette concorrenti attese), **alzare i soft limit** prima di restare a corto di capacità.
Config rilevante (prod, tipicamente `infra/.env` + `StreamNode`):
| Parametro | Ruolo oggi (post load test 2026-08) |
|-----------|--------------------------------------|
| `STREAM_NODE_HOME_MAX_PUBLISHERS` | Soft home (es. 6) |
| `STREAM_CLOUD_MAX_PUBLISHERS` | Soft per CPX (es. 4; cpx12 ha tenuto 8 in probe) |
| `STREAM_AUTOSCALE_MAX_OVERFLOW` / max overflow nodes | Quanti CPX in parallelo (es. 3 → tetto cluster ≈ home + N×cloud) |
| `STREAM_AUTOSCALE_SOFT_FREE_SLOTS` | Anticipo scale-out |
| Piano `concurrent_streams_limit` | Tetto **per club** (Premium Full = 10): indipendente dal cluster |
Capienza cluster soft ≈ `home_max + max_overflow × cloud_max` (es. 6+3×4 = **18**).
Quando si parla di capacity planning, deploy autoscale, o “troppe dirette”, ricordare di rivedere questi valori (e il limite piano) in base ai clienti reali — non lasciare i soft limit di collaudo/early-prod a lungo.
+4
View File
@@ -55,6 +55,10 @@ class SessionChannel < ApplicationCable::Channel
Sessions::Pause.new(session).call Sessions::Pause.new(session).call
when "resume_stream" when "resume_stream"
Sessions::Resume.new(session).call 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" when "set_min_quality"
Sessions::SetMinQuality.new(session, preset: data["min_quality_preset"]).call Sessions::SetMinQuality.new(session, preset: data["min_quality_preset"]).call
end end
@@ -0,0 +1,164 @@
# frozen_string_literal: true
module Admin
class AnalyticsController < Admin::BaseController
def index
@filters = {
from: parse_date(params[:from]) || 7.days.ago.to_date,
to: parse_date(params[:to]) || Time.zone.today,
device: params[:device].presence
}
@analytics_preview_active = analytics_preview_active?
scope = AnalyticsPageStat.where(day: @filters[:from]..@filters[:to])
scope = scope.where(device: @filters[:device]) if @filters[:device].present? && AnalyticsEvent::DEVICES.include?(@filters[:device])
cell_scope = AnalyticsPageCell.where(day: @filters[:from]..@filters[:to])
cell_scope = cell_scope.where(device: @filters[:device]) if @filters[:device].present? && AnalyticsEvent::DEVICES.include?(@filters[:device])
clicks_by_path = cell_scope.group(:page_path).sum(:click_count)
moves_by_path = cell_scope.group(:page_path).sum(:move_count)
pageviews_by_path = scope.group(:page_path).sum(:pageview_count)
scroll_samples_by_path = scope.group(:page_path).sum(:scroll_samples)
scroll_sum_by_path = scope.group(:page_path).sum(:scroll_sum_pct)
max_scroll_by_path = scope.group(:page_path).maximum(:max_scroll_pct)
paths = (pageviews_by_path.keys + clicks_by_path.keys + moves_by_path.keys).uniq
@pages = paths.map do |path|
samples = scroll_samples_by_path[path].to_i
{
page_path: path,
pageviews: pageviews_by_path[path].to_i,
clicks: clicks_by_path[path].to_i,
moves: moves_by_path[path].to_i,
scroll_samples: samples,
scroll_sum: scroll_sum_by_path[path].to_i,
max_scroll: max_scroll_by_path[path].to_i
}
end.sort_by { |r| [-r[:pageviews], -r[:moves], -r[:clicks], r[:page_path]] }
end
def show
@page_path = params[:page_path].to_s
redirect_to admin_analytics_path, alert: t("admin.analytics.missing_path") and return if @page_path.blank?
from = parse_date(params[:from]) || 7.days.ago.to_date
to = parse_date(params[:to]) || Time.zone.today
@device_tab_stats = device_tab_stats(@page_path, from, to)
device = resolve_heatmap_device(@device_tab_stats, params[:device].presence)
@filters = {
from: from,
to: to,
device: device,
layer: params[:layer].to_s
}
cells = AnalyticsPageCell.where(
page_path: @page_path,
day: @filters[:from]..@filters[:to],
device: @filters[:device]
)
@click_total = cells.sum(:click_count)
@move_total = cells.sum(:move_count)
@filters[:layer] =
if %w[move click].include?(@filters[:layer])
@filters[:layer]
elsif @move_total.positive?
"move"
else
"click"
end
counter = @filters[:layer] == "click" ? :click_count : :move_count
@cells = cells.group(:cell_x, :cell_y).sum(counter)
@max_weight = @cells.values.max.to_i
@total_points = @cells.values.sum
stats = AnalyticsPageStat.where(
page_path: @page_path,
day: @filters[:from]..@filters[:to],
device: @filters[:device]
)
@pageviews = stats.sum(:pageview_count)
@scroll_samples = stats.sum(:scroll_samples)
@scroll_sum = stats.sum(:scroll_sum_pct)
@max_scroll = stats.maximum(:max_scroll_pct).to_i
@avg_scroll = @scroll_samples.positive? ? (@scroll_sum.to_f / @scroll_samples).round : 0
@grid = AnalyticsPageCell::GRID_SIZE
@snapshot = find_snapshot(@page_path, @filters[:device])
end
def preview_enable
Analytics::Suppress.enable!(cookies)
redirect_to preview_return_to(params[:return_to]), notice: t("admin.analytics.preview.enabled")
end
def preview_disable
Analytics::Suppress.disable!(cookies)
redirect_to admin_analytics_path, notice: t("admin.analytics.preview.disabled")
end
private
def analytics_preview_active?
Analytics::Suppress.active?(cookies[Analytics::Suppress::COOKIE_NAME])
end
def preview_return_to(value)
path = value.to_s.strip
return root_path if path.blank?
return path if path.start_with?("/") && !path.start_with?("//")
admin_analytics_path
end
def parse_date(value)
return nil if value.blank?
Date.parse(value.to_s)
rescue ArgumentError, TypeError
nil
end
def device_tab_stats(page_path, from, to)
cell_totals = AnalyticsPageCell.where(page_path: page_path, day: from..to)
.group(:device)
.pluck(
:device,
Arel.sql("SUM(click_count)"),
Arel.sql("SUM(move_count)")
)
cell_by_device = cell_totals.to_h { |device, clicks, moves| [device, { clicks: clicks.to_i, moves: moves.to_i }] }
pageview_totals = AnalyticsPageStat.where(page_path: page_path, day: from..to)
.group(:device)
.sum(:pageview_count)
AnalyticsEvent::DEVICES.index_with do |device|
cells = cell_by_device[device] || { clicks: 0, moves: 0 }
cells.merge(pageviews: pageview_totals[device].to_i)
end
end
def resolve_heatmap_device(tab_stats, requested)
if requested.present? && AnalyticsEvent::DEVICES.include?(requested)
return requested
end
AnalyticsEvent::DEVICES.max_by do |device|
stats = tab_stats[device]
stats[:clicks] + stats[:moves] + stats[:pageviews]
end
end
def find_snapshot(page_path, device)
return nil unless AnalyticsEvent::DEVICES.include?(device)
AnalyticsPageSnapshot.where(page_path: page_path, device: device)
.order(captured_at: :desc)
.detect { |snapshot| snapshot.image.attached? }
end
end
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
@@ -1,6 +1,8 @@
module Admin module Admin
class AuthController < Admin::BaseController class AuthController < Admin::BaseController
skip_before_action :require_admin_login, only: %i[new create] skip_before_action :require_admin_login, only: %i[new create]
# Same as public logout: tolerate stale CSRF after long-lived admin tabs.
skip_before_action :verify_authenticity_token, only: :destroy
def new def new
redirect_to admin_root_path if admin_logged_in? redirect_to admin_root_path if admin_logged_in?
@@ -2,6 +2,7 @@ module Admin
class BillingController < BaseController class BillingController < BaseController
def index def index
@pending_payments = pending_scope.recent @pending_payments = pending_scope.recent
@pending_transfers = transfer_scope
@completed_payments = Billing::Payment.with_invoice_pdf.includes(:club, :invoice).recent.limit(40) @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? @filter_club = Club.find_by(id: params[:club_id]) if params[:club_id].present?
@clubs = Club.order(:name) @clubs = Club.order(:name)
@@ -18,6 +19,25 @@ module Admin
alert: e.message alert: e.message
end 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 private
def pending_scope def pending_scope
@@ -26,6 +46,12 @@ module Admin
scope scope
end 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) def billing_redirect_params(payment)
{ club_id: payment.club_id, anchor: "payment-#{payment.id}" }.compact { club_id: payment.club_id, anchor: "payment-#{payment.id}" }.compact
end end
@@ -1,9 +1,9 @@
module Admin module Admin
class ClubsController < BaseController 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 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) .order(:name)
end end
@@ -11,6 +11,7 @@ module Admin
@subscription = @club.subscription || @club.build_subscription(plan: Plan["free"], status: "active") @subscription = @club.subscription || @club.build_subscription(plan: Plan["free"], status: "active")
@plans = Plan.ordered.reject { |p| p.slug == "free" } @plans = Plan.ordered.reject { |p| p.slug == "free" }
@teams = @club.teams.order(:name) @teams = @club.teams.order(:name)
@quote = @club.active_billing_quote
end end
def grant_comped def grant_comped
@@ -32,6 +33,32 @@ module Admin
redirect_back_or_club alert: e.message redirect_back_or_club alert: e.message
end 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 private
def set_club def set_club
@@ -0,0 +1,68 @@
# frozen_string_literal: true
module Admin
class CostEntriesController < Admin::BaseController
before_action :set_entry, only: %i[edit update destroy]
def new
@entry = PlatformCostEntry.new(
month: parse_month(params[:month]) || Time.zone.today.beginning_of_month,
label: PlatformCostEntry::DEFAULT_LABEL
)
end
def create
@entry = PlatformCostEntry.new(entry_attributes)
if @entry.save
redirect_to admin_costs_path(month: month_param(@entry.month)), notice: t("admin.flash.cost_entry_created")
else
flash.now[:alert] = @entry.errors.full_messages.join(", ")
render :new, status: :unprocessable_entity
end
end
def edit; end
def update
if @entry.update(entry_attributes)
redirect_to admin_costs_path(month: month_param(@entry.month)), notice: t("admin.flash.cost_entry_updated")
else
flash.now[:alert] = @entry.errors.full_messages.join(", ")
render :edit, status: :unprocessable_entity
end
end
def destroy
month = @entry.month
@entry.destroy!
redirect_to admin_costs_path(month: month_param(month)), notice: t("admin.flash.cost_entry_destroyed")
end
private
def set_entry
@entry = PlatformCostEntry.find(params[:id])
end
def entry_attributes
attrs = params.require(:platform_cost_entry).permit(:month, :label, :amount_euros, :notes)
if attrs[:amount_euros].present?
attrs[:amount_cents] = Billing::EuroAmount.to_cents(attrs.delete(:amount_euros))
end
attrs[:notes] = nil if attrs[:notes].blank?
attrs
end
def parse_month(value)
return nil if value.blank?
Date.strptime(value.to_s, "%Y-%m").beginning_of_month
rescue ArgumentError, TypeError
nil
end
def month_param(date)
date.strftime("%Y-%m")
end
end
end
@@ -0,0 +1,30 @@
# frozen_string_literal: true
module Admin
class CostsController < Admin::BaseController
def index
@month = parse_month(params[:month]) || Time.zone.today.beginning_of_month
analytics = Admin::CostAnalytics.new(month: @month)
@summary = analytics.summary
@trend = analytics.trend
@clubs = analytics.club_breakdown
@entries = PlatformCostEntry.for_month(@month).ordered
@month_options = month_options(@month)
end
private
def parse_month(value)
return nil if value.blank?
Date.strptime(value.to_s, "%Y-%m").beginning_of_month
rescue ArgumentError, TypeError
nil
end
def month_options(selected)
start = selected - 23.months
(0..23).map { |i| start + i.months }.reverse
end
end
end
@@ -5,7 +5,7 @@ module Admin
@host = HostMetrics.new.sample! @host = HostMetrics.new.sample!
@active_sessions = StreamSession @active_sessions = StreamSession
.where(status: DashboardStats::ACTIVE_STATUSES) .where(status: DashboardStats::ACTIVE_STATUSES)
.includes(match: :team) .includes(:stream_node, match: :team)
.order(started_at: :desc) .order(started_at: :desc)
@teams = Team.includes(:matches).order(:name).limit(8) @teams = Team.includes(:matches).order(:name).limit(8)
end end
@@ -1,12 +1,23 @@
module Admin module Admin
class SessionsController < Admin::BaseController class SessionsController < Admin::BaseController
def index def index
@sessions = StreamSession.includes(:match, :user).order(created_at: :desc).limit(50) @filters = session_filters
scope = filtered_sessions
@total_count = scope.count
@sessions = scope
.includes(:stream_node, :user, match: { team: :club })
.order(Arel.sql("COALESCE(stream_sessions.started_at, stream_sessions.created_at) DESC"))
.limit(100)
@clubs = Club.order(:name)
@nodes = StreamNode.order(:slug)
end end
def show def show
@session = StreamSession.find(params[:id]) @session = StreamSession
@events = @session.stream_events.recent.limit(50) .includes(:stream_node, :user, :device_states, :recording, match: { team: :club })
.find(params[:id])
@events = @session.stream_events.recent.limit(100)
@club = @session.match.team.club
end end
def stop def stop
@@ -39,5 +50,73 @@ module Admin
regia_expires_at: @session.regia_token_expires_at&.iso8601 regia_expires_at: @session.regia_token_expires_at&.iso8601
} }
end end
private
def session_filters
{
q: params[:q].to_s.strip.presence,
status: params[:status].to_s.strip.presence,
platform: params[:platform].to_s.strip.presence,
club_id: params[:club_id].to_s.strip.presence,
stream_node_id: params[:stream_node_id].to_s.strip.presence,
from: params[:from].to_s.strip.presence,
to: params[:to].to_s.strip.presence
}
end
def filtered_sessions
scope = StreamSession.left_outer_joins(match: { team: :club })
if @filters[:q].present?
term = "%#{ActiveRecord::Base.sanitize_sql_like(@filters[:q])}%"
scope = scope.where(
"teams.name ILIKE :term OR matches.opponent_name ILIKE :term OR clubs.name ILIKE :term OR matches.location ILIKE :term",
term: term
)
end
if @filters[:status].present? && StreamSession::STATUSES.include?(@filters[:status])
scope = scope.where(stream_sessions: { status: @filters[:status] })
end
if @filters[:platform].present? && StreamSession::LIVE_PLATFORMS.include?(@filters[:platform])
scope = scope.where(stream_sessions: { platform: @filters[:platform] })
end
if @filters[:club_id].present?
scope = scope.where(teams: { club_id: @filters[:club_id] })
end
if @filters[:stream_node_id].present?
node = StreamNode.find_by(id: @filters[:stream_node_id])
if node
scope = scope.where(
"stream_sessions.stream_node_id = :id OR stream_sessions.ingest_slug = :slug",
id: node.id,
slug: node.slug
)
end
end
if (from_time = parse_filter_date(@filters[:from], end_of_day: false))
scope = scope.where(
"COALESCE(stream_sessions.started_at, stream_sessions.created_at) >= ?",
from_time
)
end
if (to_time = parse_filter_date(@filters[:to], end_of_day: true))
scope = scope.where(
"COALESCE(stream_sessions.started_at, stream_sessions.created_at) <= ?",
to_time
)
end
scope
end
def parse_filter_date(value, end_of_day:)
return nil if value.blank?
date = Date.parse(value)
end_of_day ? date.end_of_day : date.beginning_of_day
rescue ArgumentError, TypeError
nil
end
end end
end 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,45 @@
# frozen_string_literal: true
module Analytics
class EventsController < ActionController::API
include ActionController::Cookies
def create
if analytics_suppressed?
return render json: { accepted: 0, rejected: 0, suppressed: true }, status: :accepted
end
payload = parse_payload
result = Analytics::Ingest.new(events: payload, remote_ip: request.remote_ip).call
if result.rate_limited
return head :too_many_requests
end
render json: { accepted: result.accepted, rejected: result.rejected }, status: :accepted
end
private
def analytics_suppressed?
Analytics::Suppress.active?(cookies[Analytics::Suppress::COOKIE_NAME])
end
def parse_payload
body = request.request_parameters
return body["events"] if body.is_a?(Hash) && body["events"].is_a?(Array)
return body if body.is_a?(Array)
raw = request.raw_post
return [] if raw.blank?
parsed = JSON.parse(raw)
return parsed["events"] if parsed.is_a?(Hash) && parsed["events"].is_a?(Array)
return parsed if parsed.is_a?(Array)
[]
rescue JSON::ParserError
[]
end
end
end
@@ -0,0 +1,33 @@
# frozen_string_literal: true
module Analytics
class SnapshotsController < ActionController::API
include ActionController::Cookies
def create
if analytics_suppressed?
return render json: { ok: true, skipped: true, suppressed: true }, status: :accepted
end
result = Analytics::SnapshotIngest.new(
path: params[:path] || params[:page_path],
device: params[:device],
width: params[:width],
height: params[:height],
image: params[:image],
remote_ip: request.remote_ip
).call
return head :too_many_requests if result.rate_limited
return render json: { ok: false, error: result.error }, status: :unprocessable_content unless result.ok
render json: { ok: true, skipped: result.skipped }, status: :accepted
end
private
def analytics_suppressed?
Analytics::Suppress.active?(cookies[Analytics::Suppress::COOKIE_NAME])
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,8 @@ module Api
module V1 module V1
class BaseController < ApplicationController class BaseController < ApplicationController
rescue_from Teams::EntitlementError, with: :render_entitlement_error rescue_from Teams::EntitlementError, with: :render_entitlement_error
rescue_from Streams::IngestUnavailableError, with: :render_ingest_unavailable
rescue_from BrandingAttachments::CoverUploadError, with: :render_cover_upload_error
rescue_from Youtube::BroadcastService::Error, with: :render_youtube_error rescue_from Youtube::BroadcastService::Error, with: :render_youtube_error
private private
@@ -20,6 +22,17 @@ module Api
billing_url: error.billing_url billing_url: error.billing_url
}, status: :forbidden }, status: :forbidden
end end
def render_ingest_unavailable(error)
render json: {
error: error.message,
error_code: error.code
}, status: :service_unavailable
end
def render_cover_upload_error(error)
render json: { error: error.message, error_code: "cover_upload_invalid" }, status: :unprocessable_entity
end
end end
end end
end end
@@ -2,6 +2,7 @@ module Api
module V1 module V1
class MatchesController < BaseController class MatchesController < BaseController
include Api::UrlHelper include Api::UrlHelper
include Api::CoverJson
include BrandingAttachments include BrandingAttachments
before_action :set_team, only: %i[index create] before_action :set_team, only: %i[index create]
@@ -21,6 +22,7 @@ module Api
normalize_scoring_rules!(attrs) normalize_scoring_rules!(attrs)
match = @team.matches.create!(attrs) match = @team.matches.create!(attrs)
attach_opponent_logo(match) attach_opponent_logo(match)
attach_match_cover(match)
render json: match_json(match), status: :created render json: match_json(match), status: :created
end end
@@ -34,6 +36,7 @@ module Api
normalize_opponent_color!(attrs) normalize_opponent_color!(attrs)
@match.update!(attrs) @match.update!(attrs)
attach_opponent_logo(@match) attach_opponent_logo(@match)
attach_match_cover(@match)
render json: match_json(@match) render json: match_json(@match)
end end
@@ -98,6 +101,11 @@ module Api
match.opponent_logo_file.attach(file) if file.present? match.opponent_logo_file.attach(file) if file.present?
end 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) def match_json(match, detail: false)
active = active_session_for(match) active = active_session_for(match)
team = match.team team = match.team
@@ -123,6 +131,7 @@ module Api
home_logo_url: api_absolute_url(team.effective_logo_url), home_logo_url: api_absolute_url(team.effective_logo_url),
opponent_primary_color: match.effective_opponent_primary_color, opponent_primary_color: match.effective_opponent_primary_color,
opponent_logo_url: api_absolute_url(match.opponent_logo_url), opponent_logo_url: api_absolute_url(match.opponent_logo_url),
**match_cover_json(match),
active_session_id: active&.id, active_session_id: active&.id,
active_session_status: active&.status, active_session_status: active&.status,
stream_completed: match.stream_completed?, stream_completed: match.stream_completed?,
@@ -97,10 +97,11 @@ module Api
replay_url: recording.replay_url, replay_url: recording.replay_url,
playback_url: recording.playback_stream_url, playback_url: recording.playback_stream_url,
thumbnail_url: recording.thumbnail_url, thumbnail_url: recording.thumbnail_url,
download_enabled: ent.phone_download_enabled?, download_enabled: ent.phone_download_enabled? && recording.storage_key.present?,
youtube_video_id: recording.youtube_video_id, youtube_video_id: recording.youtube_video_id,
youtube_watch_url: recording.youtube_watch_url, youtube_watch_url: recording.youtube_watch_url,
youtube_publish_enabled: ent.premium_full? && ent.youtube_enabled?, youtube_publish_enabled: ent.premium_full? && ent.youtube_enabled? && recording.storage_key.present? && recording.youtube_video_id.blank?,
replay_source: recording.replay_source,
source_platform: recording.source_platform, source_platform: recording.source_platform,
source_platform_label: recording.source_platform_label, source_platform_label: recording.source_platform_label,
expires_at: recording.expires_at, expires_at: recording.expires_at,
@@ -34,6 +34,12 @@ module Api
render json: session_json(@session) render json: session_json(@session)
end 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 def score
Scoring::SyncState.new(session: @session, payload: score_sync_params).call Scoring::SyncState.new(session: @session, payload: score_sync_params).call
render json: session_json(@session, detail: true) render json: session_json(@session, detail: true)
@@ -71,6 +77,7 @@ module Api
thermal_state: sanitized_thermal_state, thermal_state: sanitized_thermal_state,
last_seen_at: Time.current last_seen_at: Time.current
) )
Sessions::ApplyClientInfo.call(@session, params[:client]) if params[:client].present?
sync_publisher_when_streaming!(params[:fps].to_f) sync_publisher_when_streaming!(params[:fps].to_f)
SessionChannel.broadcast_message(@session, state.as_cable_payload) SessionChannel.broadcast_message(@session, state.as_cable_payload)
head :no_content head :no_content
@@ -174,7 +181,12 @@ module Api
end end
def session_params def session_params
params.permit(:platform, :privacy_status, :quality_preset, :target_bitrate, :target_fps, :youtube_channel) params.permit(
:platform, :privacy_status, :quality_preset, :target_bitrate, :target_fps, :youtube_channel,
client: %i[os client_os app_version app_build version build build_number
device_manufacturer manufacturer device_model model
os_version system_version carrier network_operator operator]
)
end end
def score_sync_params def score_sync_params
@@ -193,6 +205,7 @@ module Api
platform: session.platform, platform: session.platform,
rtmp_ingest_url: session.rtmp_ingest_url, rtmp_ingest_url: session.rtmp_ingest_url,
hls_playback_url: session.hls_playback_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, watch_page_url: session.matchlivetv_platform? ? session.watch_page_url : nil,
share_url: session.share_url, share_url: session.share_url,
youtube_watch_url: session.youtube_watch_url, youtube_watch_url: session.youtube_watch_url,
@@ -203,7 +216,8 @@ module Api
min_quality_preset: session.min_quality_preset, min_quality_preset: session.min_quality_preset,
target_bitrate: session.target_bitrate, target_bitrate: session.target_bitrate,
started_at: session.started_at, started_at: session.started_at,
disconnection_count: session.disconnection_count disconnection_count: session.disconnection_count,
audio_muted: session.audio_muted
} }
if detail if detail
data[:score] = session.score_state&.as_cable_payload data[:score] = session.score_state&.as_cable_payload
@@ -137,7 +137,8 @@ module Api
youtube_mode: ent.plan.youtube_mode, youtube_mode: ent.plan.youtube_mode,
staff_role: current_user.staff_role_for(team), staff_role: current_user.staff_role_for(team),
can_stream: current_user.can_stream_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 if detail
data[:members] = team.user_teams.includes(:user).where.not(role: "owner").map do |ut| data[:members] = team.user_teams.includes(:user).where.not(role: "owner").map do |ut|
@@ -170,12 +171,13 @@ module Api
replay_url: recording.replay_url, replay_url: recording.replay_url,
playback_url: recording.playback_stream_url, playback_url: recording.playback_stream_url,
thumbnail_url: recording.thumbnail_url, thumbnail_url: recording.thumbnail_url,
download_enabled: ent.phone_download_enabled?, download_enabled: ent.phone_download_enabled? && recording.storage_key.present?,
view_count: recording.view_count, view_count: recording.view_count,
views_label: recording.views_label, views_label: recording.views_label,
youtube_video_id: recording.youtube_video_id, youtube_video_id: recording.youtube_video_id,
youtube_watch_url: recording.youtube_watch_url, youtube_watch_url: recording.youtube_watch_url,
youtube_publish_enabled: ent.premium_full? && ent.youtube_enabled?, youtube_publish_enabled: ent.premium_full? && ent.youtube_enabled? && recording.storage_key.present? && recording.youtube_video_id.blank?,
replay_source: recording.replay_source,
source_platform: recording.source_platform, source_platform: recording.source_platform,
source_platform_label: recording.source_platform_label, source_platform_label: recording.source_platform_label,
expires_at: recording.expires_at, expires_at: recording.expires_at,
@@ -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? record.logo_file.attach(file) if file.present?
end 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) def normalize_hex_color(value, fallback)
v = value.to_s.strip v = value.to_s.strip
v = "##{v}" if v.match?(/\A[0-9A-Fa-f]{6}\z/) v = "##{v}" if v.match?(/\A[0-9A-Fa-f]{6}\z/)
@@ -4,11 +4,20 @@ module MediamtxPlayback
private private
def mediamtx_paths_index 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 end
def mediamtx_path_info(session, path_name: mediamtx_playback_path_name(session)) 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 end
def mediamtx_playback_path_name(session) def mediamtx_playback_path_name(session)
@@ -28,7 +37,7 @@ module MediamtxPlayback
end end
def mediamtx_publisher_online?(session) 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"] info && info["online"]
end end
end end
@@ -20,7 +20,8 @@ class HlsProxyController < ActionController::Base
def proxy_mediamtx_path(upstream_path) def proxy_mediamtx_path(upstream_path)
return head :not_found if upstream_path.blank? 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? upstream = "#{upstream}?#{request.query_string}" if request.query_string.present?
cookie = request.headers["Cookie"].presence || "cookieCheck=1" cookie = request.headers["Cookie"].presence || "cookieCheck=1"
@@ -30,7 +31,7 @@ class HlsProxyController < ActionController::Base
if response.status.in?([301, 302, 307, 308]) if response.status.in?([301, 302, 307, 308])
location = response.headers["location"].to_s location = response.headers["location"].to_s
cookie = cookie_from_set_header(response.headers["set-cookie"]).presence || cookie cookie = cookie_from_set_header(response.headers["set-cookie"]).presence || cookie
upstream = resolve_upstream_url(location) upstream = resolve_upstream_url(location, origin)
next next
end end
@@ -43,10 +44,21 @@ class HlsProxyController < ActionController::Base
head :bad_gateway head :bad_gateway
end 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://") 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}" location.start_with?("/") ? "#{base}#{location}" : "#{base}/#{location}"
end end
@@ -11,11 +11,15 @@ module Public
@club.assign_attributes(billing_profile_params) @club.assign_attributes(billing_profile_params)
if @club.save(context: :billing_profile) if @club.save(context: :billing_profile)
if premium_checkout_return_params.present? && @club.billing_profile_complete? if premium_checkout_return_params.present? && @club.billing_profile_complete?
redirect_to public_club_checkout_path( if @club.active_billing_quote.present?
@club, redirect_to public_club_billing_path(@club), notice: t("flash.club_billing.profile_updated")
plan: premium_checkout_return_params[:plan], else
interval: premium_checkout_return_params[:interval] redirect_to public_club_checkout_path(
), notice: t("flash.club_billing.profile_saved_proceed_payment") @club,
plan: premium_checkout_return_params[:plan],
interval: premium_checkout_return_params[:interval]
), notice: t("flash.club_billing.profile_saved_proceed_payment")
end
else else
redirect_to public_club_billing_path(@club), notice: t("flash.club_billing.profile_updated") redirect_to public_club_billing_path(@club), notice: t("flash.club_billing.profile_updated")
end end
@@ -42,6 +46,29 @@ module Public
redirect_to public_club_billing_path(@club), alert: t("flash.club_billing.stripe_error", message: e.message) redirect_to public_club_billing_path(@club), alert: t("flash.club_billing.stripe_error", message: e.message)
end 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 def download_invoice
invoice = @club.billing_invoices.find(params[:invoice_id]) invoice = @club.billing_invoices.find(params[:invoice_id])
unless invoice.pdf.attached? unless invoice.pdf.attached?
@@ -23,23 +23,28 @@ module Public
ClubMembership.create!(user: current_user, club: club, role: "owner") 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") 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, name: first_team_name,
sport: club.sport sport: club.sport
) )
plan = params[:plan].presence_in(%w[free premium_light premium_full]) || "free" desired_plan = params[:plan].presence_in(%w[free premium_light premium_full]) || "free"
Billing::AssignPlan.call(club: club, plan_slug: plan) 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? if desired_plan.in?(%w[premium_light premium_full])
unless club.billing_profile_complete? interval = Billing::Stripe::PriceCatalog::DEFAULT_INTERVAL
redirect_to public_club_billing_profile_path(club, plan: plan, interval: checkout_interval_param), if MatchLiveTv.stripe_enabled?
alert: t("flash.clubs.complete_billing_first") unless club.billing_profile_complete?
return 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 end
interval = checkout_interval_param
redirect_to public_club_checkout_path(club, plan: plan, interval: interval)
else else
redirect_to public_club_path(club), notice: t("flash.clubs.created") redirect_to public_club_path(club), notice: t("flash.clubs.created")
end end
@@ -59,17 +64,26 @@ module Public
def edit def edit
require_club_owner!(@club) require_club_owner!(@club)
@entitlements = @club.teams.first&.entitlements
end end
def update def update
require_club_owner!(@club) require_club_owner!(@club)
@entitlements = @club.teams.first!.entitlements
@club.assign_attributes(club_params) @club.assign_attributes(club_params)
attach_branding_logo(@club) attach_branding_logo(@club)
attach_cover_image(@club, entitlements: @entitlements)
@club.save! @club.save!
enqueue_cover_slate_generation(@club)
redirect_to public_club_path(@club), notice: t("flash.clubs.updated") redirect_to public_club_path(@club), notice: t("flash.clubs.updated")
rescue ActiveRecord::RecordInvalid => e rescue ActiveRecord::RecordInvalid => e
@entitlements = @club.teams.first&.entitlements
flash.now[:alert] = e.record.errors.full_messages.join(", ") flash.now[:alert] = e.record.errors.full_messages.join(", ")
render :edit, status: :unprocessable_entity 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 end
def billing def billing
@@ -82,6 +96,8 @@ module Public
@entitlements = @team.entitlements @entitlements = @team.entitlements
@plans = Plan.ordered @plans = Plan.ordered
@payments = @club.billing_payments.recent.includes(:invoice).limit(50) @payments = @club.billing_payments.recent.includes(:invoice).limit(50)
@quote = @club.active_billing_quote
@pending_transfer = @club.billing_transfer_orders.awaiting_payment.first
end end
def checkout def checkout
@@ -92,6 +108,12 @@ module Public
return return
end 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? unless MatchLiveTv.stripe_enabled?
redirect_to public_club_billing_path(@club), alert: t("flash.clubs.stripe_not_configured") redirect_to public_club_billing_path(@club), alert: t("flash.clubs.stripe_not_configured")
return 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]) @session.status.in?(%w[connecting live reconnecting paused])
) )
@publisher_online = !@stream_closed && mediamtx_publisher_online?(@session) @publisher_online = !@stream_closed && mediamtx_publisher_online?(@session)
@cover_poster_url = Streams::CoverResolver.new(@match).effective_cover_url
end end
def status def status
@@ -7,7 +7,7 @@ module Public
layout "regia" layout "regia"
protect_from_forgery with: :null_session 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 before_action :set_session_from_token
@@ -61,6 +61,12 @@ module Public
render json: status_payload render json: status_payload
end 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 def min_quality
Sessions::SetMinQuality.new(@session, preset: params[:min_quality_preset]).call Sessions::SetMinQuality.new(@session, preset: params[:min_quality_preset]).call
render json: status_payload render json: status_payload
@@ -90,6 +96,7 @@ module Public
{ {
status: @session.status, status: @session.status,
paused: @session.paused?, paused: @session.paused?,
audio_muted: @session.audio_muted,
stream_closed: closed, stream_closed: closed,
live: !closed && (@session.live? || @session.paused? || publisher_online), live: !closed && (@session.live? || @session.paused? || publisher_online),
on_air: playable || (!closed && @session.status.in?(%w[connecting live reconnecting paused])), on_air: playable || (!closed && @session.status.in?(%w[connecting live reconnecting paused])),
@@ -1,5 +1,9 @@
module Public module Public
class SessionsController < WebBaseController class SessionsController < WebBaseController
# Logout must succeed even with a stale CSRF token (old tab after deploy /
# session rotate). SameSite=Lax already blocks cross-site cookie POSTs.
skip_before_action :verify_authenticity_token, only: :destroy
def new def new
if logged_in? && current_user.primary_club if logged_in? && current_user.primary_club
redirect_to public_club_path(current_user.primary_club) redirect_to public_club_path(current_user.primary_club)
@@ -9,6 +9,7 @@ module Public
protect_from_forgery with: :exception protect_from_forgery with: :exception
before_action :set_site_locale before_action :set_site_locale
before_action :load_site_announcements
private private
@@ -40,5 +41,27 @@ module Public
def logged_in? def logged_in?
current_user.present? current_user.present?
end 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
end end
@@ -8,6 +8,7 @@ module Public
{ loc: "#{base}/prezzi", changefreq: "monthly", priority: "0.9" }, { loc: "#{base}/prezzi", changefreq: "monthly", priority: "0.9" },
{ loc: "#{base}/pallavolo-giovanile", changefreq: "monthly", priority: "0.85" }, { loc: "#{base}/pallavolo-giovanile", changefreq: "monthly", priority: "0.85" },
{ loc: "#{base}/faq", changefreq: "monthly", priority: "0.8" }, { 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}/live", changefreq: "hourly", priority: "0.85" },
{ loc: "#{base}/squadre", changefreq: "daily", priority: "0.85" }, { loc: "#{base}/squadre", changefreq: "daily", priority: "0.85" },
{ loc: "#{base}/privacy", changefreq: "yearly", priority: "0.3" }, { loc: "#{base}/privacy", changefreq: "yearly", priority: "0.3" },
@@ -19,6 +19,7 @@ module Public
require_club_owner!(@club) require_club_owner!(@club)
team = @club.teams.create!(team_params) team = @club.teams.create!(team_params)
attach_branding_logo(team) 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) redirect_to public_club_path(@club), notice: t("flash.teams.added", name: team.name)
rescue ActiveRecord::RecordInvalid => e rescue ActiveRecord::RecordInvalid => e
flash.now[:alert] = e.record.errors.full_messages.join(", ") flash.now[:alert] = e.record.errors.full_messages.join(", ")
@@ -41,24 +42,33 @@ module Public
def edit def edit
require_club_owner_for_team!(@team) require_club_owner_for_team!(@team)
@club = @team.club @club = @team.club
@entitlements = @team.entitlements
end end
def update def update
require_club_owner_for_team!(@team) require_club_owner_for_team!(@team)
@team.assign_attributes(team_params) @team.assign_attributes(team_params)
attach_branding_logo(@team) attach_branding_logo(@team)
attach_cover_image(@team, entitlements: @team.entitlements)
attach_team_photo(@team) attach_team_photo(@team)
@team.save! @team.save!
enqueue_cover_slate_generation(@team)
redirect_to public_team_details_path(@team), notice: t("flash.teams.updated") redirect_to public_team_details_path(@team), notice: t("flash.teams.updated")
rescue ActiveRecord::RecordInvalid => e rescue ActiveRecord::RecordInvalid => e
@club = @team.club @club = @team.club
@entitlements = @team.entitlements
flash.now[:alert] = e.record.errors.full_messages.join(", ") flash.now[:alert] = e.record.errors.full_messages.join(", ")
render :edit, status: :unprocessable_entity 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 end
def invite def invite
require_club_owner_for_team!(@team) require_club_owner_for_team!(@team)
@entitlements = @team.entitlements redirect_to public_team_details_path(@team, anchor: "invita-trasmissione")
end end
def assign_self_staff def assign_self_staff
@@ -87,18 +97,32 @@ module Public
email = params[:email]&.downcase&.strip email = params[:email]&.downcase&.strip
Teams::StaffEmailValidator.assert_available!(team: @team, email: email, staff_kind: staff_kind) Teams::StaffEmailValidator.assert_available!(team: @team, email: email, staff_kind: staff_kind)
token = TeamInvitation.generate_token token = TeamInvitation.generate_token
@team.team_invitations.create!( invitation = @team.team_invitations.create!(
email: email, email: email,
token_digest: Digest::SHA256.hexdigest(token), token_digest: Digest::SHA256.hexdigest(token),
role: "member", role: "member",
staff_kind: staff_kind, staff_kind: staff_kind,
expires_at: 7.days.from_now expires_at: 7.days.from_now
) )
@invite_url = join_public_invitation_url(token: token) invite_url = public_invitation_url(token: token)
flash.now[:notice] = t("flash.teams.invite_link_generated") flash[:invite_url] = invite_url
render :invite
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 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 end
def remove_member def remove_member
+151
View File
@@ -1,4 +1,21 @@
module AdminHelper module AdminHelper
def format_euros(cents, precision: 2)
return I18n.t("admin.common.dash") if cents.nil?
format("%.*f €", precision, cents.to_f / 100.0)
end
def format_hours(hours)
return I18n.t("admin.common.dash") if hours.nil? || hours.to_f <= 0
total_minutes = (hours.to_f * 60).round
format_duration_minutes(total_minutes)
end
def admin_month_label(month)
I18n.l(month, format: "%B %Y")
end
def format_bytes(bytes) def format_bytes(bytes)
return "" if bytes.nil? return "" if bytes.nil?
@@ -30,6 +47,113 @@ module AdminHelper
links links
end end
def admin_session_ingest_badge_class(role_or_node)
role = role_or_node.respond_to?(:role) ? role_or_node.role : role_or_node
case role.to_s
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(role_or_node)
role = role_or_node.respond_to?(:role) ? role_or_node.role : role_or_node
I18n.t("admin.sessions.ingest.role.#{role}", default: role.to_s.humanize)
end
def admin_session_status_badge_class(status)
case status.to_s
when "live" then "badge--live"
when "connecting", "reconnecting" then "badge--connecting"
when "paused", "idle" then "badge--paused"
when "ended" then "badge--ended"
when "error" then "badge--error"
else "badge--paused"
end
end
def admin_datetime(time, with_seconds: false)
return I18n.t("admin.common.dash") if time.blank?
format = with_seconds ? "%d/%m/%Y %H:%M:%S" : "%d/%m/%Y %H:%M"
time.in_time_zone.strftime(format)
end
def admin_session_duration_label(session)
secs = session.total_duration_secs.to_i
if secs <= 0 && session.started_at
end_time = session.ended_at || (session.terminal? ? session.updated_at : Time.current)
secs = (end_time - session.started_at).to_i if end_time
end
return I18n.t("admin.common.dash") if secs <= 0
hours = secs / 3600
mins = (secs % 3600) / 60
rem = secs % 60
if hours.positive?
format("%dh %02dm %02ds", hours, mins, rem)
elsif mins.positive?
format("%dm %02ds", mins, rem)
else
"#{rem}s"
end
end
def admin_session_client_os_label(session)
case session.client_os.to_s
when "android" then "Android"
when "ios" then "iOS"
else I18n.t("admin.common.dash")
end
end
def admin_session_client_summary(session)
parts = []
parts << admin_session_client_os_label(session) if session.client_os.present?
if session.app_version.present?
ver = session.app_version
ver = "#{ver} (#{session.app_build})" if session.app_build.present?
parts << "app #{ver}"
end
device = [session.device_manufacturer, session.device_model].compact_blank.join(" ")
parts << device if device.present?
parts << "OS #{session.os_version}" if session.os_version.present?
parts << session.carrier if session.carrier.present?
parts.presence&.join(" · ") || I18n.t("admin.common.dash")
end
def admin_format_event_meta(metadata)
return content_tag(:span, I18n.t("admin.common.dash"), class: "muted") if metadata.blank?
pairs = metadata.to_h
return content_tag(:span, I18n.t("admin.common.dash"), class: "muted") if pairs.empty?
content_tag(:dl, class: "admin-event-meta") do
safe_join(
pairs.map do |key, value|
content_tag(:div, class: "admin-event-meta__row") do
content_tag(:dt, key.to_s) +
content_tag(:dd, admin_event_meta_value(value))
end
end
)
end
end
def admin_event_meta_value(value)
case value
when Hash, Array
content_tag(:code, value.to_json)
when TrueClass, FalseClass
value ? "true" : "false"
when NilClass
I18n.t("admin.common.dash")
else
value.to_s
end
end
def admin_regia_expires_label(iso_time) def admin_regia_expires_label(iso_time)
return if iso_time.blank? return if iso_time.blank?
@@ -37,4 +161,31 @@ module AdminHelper
rescue ArgumentError, TypeError rescue ArgumentError, TypeError
nil nil
end 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 end
+1 -1
View File
@@ -1,6 +1,6 @@
module LegalHelper module LegalHelper
def legal_last_updated def legal_last_updated
"3 giugno 2026" "20 agosto 2026"
end end
def cookie_policy_last_updated def cookie_policy_last_updated
+83 -6
View File
@@ -1,6 +1,6 @@
module Public module Public
module BillingHelper 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" if target_plan.slug == "free"
return { kind: :current, label: I18n.t("billing.actions.current_plan") } if current_slug == "free" return { kind: :current, label: I18n.t("billing.actions.current_plan") } if current_slug == "free"
return { kind: :none } if stripe_subscription_active return { kind: :none } if stripe_subscription_active
@@ -8,17 +8,50 @@ module Public
return { kind: :contact, label: I18n.t("billing.actions.contact_for_free") } return { kind: :contact, label: I18n.t("billing.actions.contact_for_free") }
end 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? 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") } return { kind: :disabled, label: I18n.t("billing.actions.stripe_not_configured") }
end end
intervals = Billing::Stripe::PriceCatalog.available_intervals(plan_slug: target_plan.slug) stripe_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? 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 active_interval = current_interval.presence || Billing::Stripe::PriceCatalog::DEFAULT_INTERVAL
if current_slug == target_plan.slug && stripe_subscription_active && !subscription&.plan_change_pending? 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? if other_intervals.empty?
label = "#{I18n.t('billing.actions.current_plan')}#{Billing::Stripe::PriceCatalog.label(plan_slug: target_plan.slug, interval: active_interval)}" label = "#{I18n.t('billing.actions.current_plan')}#{Billing::Stripe::PriceCatalog.label(plan_slug: target_plan.slug, interval: active_interval)}"
return { kind: :current, label: label } return { kind: :current, label: label }
@@ -34,13 +67,17 @@ module Public
} }
end 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 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 else
{ {
kind: :change_options, kind: :change_options,
plan: target_plan, plan: target_plan,
intervals: intervals, intervals: stripe_intervals,
current_slug: current_slug, current_slug: current_slug,
current_interval: active_interval, current_interval: active_interval,
subscription: subscription subscription: subscription
@@ -57,6 +94,10 @@ module Public
Billing::Stripe::PlanChangeMessages.button_label(plan: plan, interval: interval, kind: btn_kind) Billing::Stripe::PlanChangeMessages.button_label(plan: plan, interval: interval, kind: btn_kind)
end 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) def billing_plan_change_info_lines(subscription: nil)
Billing::Stripe::PlanChangeMessages.billing_info_lines(subscription: subscription) Billing::Stripe::PlanChangeMessages.billing_info_lines(subscription: subscription)
end end
@@ -92,5 +133,41 @@ module Public
I18n.t("billing.actions.profile_missing", fields: missing.join(", ")) I18n.t("billing.actions.profile_missing", fields: missing.join(", "))
end 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
end end
@@ -44,11 +44,26 @@ module Public
partialsPrefix: t("regia.board.partials_prefix"), partialsPrefix: t("regia.board.partials_prefix"),
resumed: t("regia.js.resumed"), resumed: t("regia.js.resumed"),
pausedCover: t("regia.js.paused_cover"), pausedCover: t("regia.js.paused_cover"),
muted: t("regia.js.muted"),
unmuted: t("regia.js.unmuted"),
closed: t("regia.js.closed"), closed: t("regia.js.closed"),
resumeError: t("regia.js.resume_error"), resumeError: t("regia.js.resume_error"),
pauseError: t("regia.js.pause_error"), pauseError: t("regia.js.pause_error"),
muteError: t("regia.js.mute_error"),
closeError: t("regia.js.close_error"), closeError: t("regia.js.close_error"),
closeConfirm: t("regia.js.close_confirm"), 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"), linkUnavailable: t("regia.js.link_unavailable"),
linkCopied: t("regia.js.link_copied"), linkCopied: t("regia.js.link_copied"),
copyPrompt: t("regia.js.copy_prompt"), copyPrompt: t("regia.js.copy_prompt"),
@@ -64,6 +79,8 @@ module Public
previewOffHint: t("regia.preview_off_hint"), previewOffHint: t("regia.preview_off_hint"),
resumeLabel: t("regia.resume"), resumeLabel: t("regia.resume"),
pauseLabel: t("regia.pause"), pauseLabel: t("regia.pause"),
muteLabel: t("regia.mute"),
unmuteLabel: t("regia.unmute"),
endedBadge: t("regia.status.ended"), endedBadge: t("regia.status.ended"),
pausedBadge: t("regia.status.paused"), pausedBadge: t("regia.status.paused"),
liveBadge: t("regia.status.live"), liveBadge: t("regia.status.live"),
@@ -0,0 +1,25 @@
# frozen_string_literal: true
module Analytics
class AggregateJob < ApplicationJob
queue_as :default
def perform(purge: false)
redis = Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0"))
lock_key = "analytics:aggregate:lock"
acquired = redis.set(lock_key, "1", nx: true, ex: 120)
return unless acquired
begin
service = Analytics::Aggregate.new
service.call
service.purge_old! if purge
ensure
redis.del(lock_key)
end
rescue StandardError => e
Rails.logger.error("[analytics] aggregate failed: #{e.class} #{e.message}")
raise
end
end
end
@@ -6,7 +6,7 @@ class CleanupExpiredSessionsJob
.where("updated_at < ?", 6.hours.ago) .where("updated_at < ?", 6.hours.ago)
.find_each do |session| .find_each do |session|
session.fail! if session.may_fail? session.fail! if session.may_fail?
Mediamtx::Client.new.delete_path(session) Mediamtx::Client.for_session(session).delete_path(session)
end end
end end
end end
@@ -6,12 +6,37 @@ class ExpireEndedSubscriptionsJob
free_plan = Plan["free"] free_plan = Plan["free"]
Subscription.where(cancel_at_period_end: true) Subscription.where(cancel_at_period_end: true)
.where("current_period_end <= ?", Time.current) .where("current_period_end <= ?", Time.current)
.where.not(plan_id: free_plan.id) .where.not(plan_id: free_plan.id)
.find_each do |sub| .where(admin_comped: false)
Billing::Stripe::FinalizeSubscription.call(club: sub.club) .find_each do |sub|
expire!(sub)
rescue StandardError => e rescue StandardError => e
Rails.logger.warn("[ExpireEndedSubscriptions] club=#{sub.club_id} #{e.message}") Rails.logger.warn("[ExpireEndedSubscriptions] club=#{sub.club_id} #{e.message}")
end end
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 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,19 @@
# frozen_string_literal: true
module Recordings
class ClearTemporaryMediaJob
include Sidekiq::Job
sidekiq_options retry: 3, queue: "default"
def perform(recording_id, reason = "verified")
recording = Recording.find_by(id: recording_id)
return unless recording
Recordings::ClearTemporaryMedia.new(
recording,
reason: reason.to_sym,
force: reason.to_s == "verified"
).call
end
end
end
@@ -1,3 +1,5 @@
# frozen_string_literal: true
module Recordings module Recordings
class PostProcessJob class PostProcessJob
include Sidekiq::Job include Sidekiq::Job
@@ -7,10 +9,30 @@ module Recordings
recording = Recording.find_by(id: recording_id) recording = Recording.find_by(id: recording_id)
return unless recording&.ready? return unless recording&.ready?
Recordings::NotifyReady.new(recording).call begin
Recordings::NotifyReady.new(recording).call
rescue StandardError => e
# SMTP/ntfy giù non deve bloccare VerifyYoutubeReplay né il resto del post-process.
Rails.logger.warn(
"[Recordings::PostProcessJob] notify_failed recording=#{recording.id} " \
"#{e.class}: #{e.message}"
)
end
if recording.temporary_storage?
grace = MatchLiveTv.youtube_replay_verify_grace_secs
Rails.logger.info(
"[Recordings::PostProcessJob] schedule VerifyYoutubeReplay " \
"recording=#{recording.id} grace=#{grace}s"
)
Recordings::VerifyYoutubeReplayJob.perform_in(grace.seconds, recording.id, 0)
return
end
# Solo MatchLiveTV-only: eventuale re-upload manuale/flag legacy (non per live YouTube).
return unless recording.auto_publish_youtube? return unless recording.auto_publish_youtube?
return if recording.youtube_video_id.present? return if recording.youtube_video_id.present?
return if recording.source_platform == "youtube"
Recordings::PublishToYoutubeJob.perform_async(recording.id) Recordings::PublishToYoutubeJob.perform_async(recording.id)
end end
@@ -0,0 +1,35 @@
# frozen_string_literal: true
module Recordings
# Safety net: elimina copie temp scadute (anche se verify non è mai riuscito).
class PurgeTemporaryMediaJob
include Sidekiq::Job
sidekiq_options retry: 1, queue: "default"
def perform
scope = Recording.temporary_media_pending_purge
count = 0
scope.find_each do |recording|
unless recording.youtube_verified_at.present?
Rails.logger.warn(
"[Recordings::PurgeTemporaryMediaJob] anomaly_unverified_expiry " \
"recording=#{recording.id} session=#{recording.stream_session_id} " \
"youtube_video_id=#{recording.youtube_video_id.inspect} " \
"temp_expires_at=#{recording.temp_expires_at.inspect}"
)
end
Recordings::ClearTemporaryMedia.new(
recording,
reason: :max_retention,
force: true
).call
count += 1
end
Rails.logger.info("[Recordings::PurgeTemporaryMediaJob] processed=#{count}")
count
end
end
end
@@ -0,0 +1,56 @@
# frozen_string_literal: true
module Recordings
class VerifyYoutubeReplayJob
include Sidekiq::Job
sidekiq_options retry: 0, queue: "default"
def perform(recording_id, attempt = 0)
recording = Recording.find_by(id: recording_id)
return unless recording&.temporary_storage?
return if recording.deleted?
attempt = attempt.to_i
max = MatchLiveTv.youtube_replay_verify_max_attempts
result = Recordings::VerifyYoutubeReplay.new(recording).call
if result.ok?
recording.reload
embeddable = recording.metadata.is_a?(Hash) && recording.metadata.dig("youtube", "embeddable") != false
if embeddable
Recordings::ClearTemporaryMediaJob.perform_async(recording.id, "verified")
else
Rails.logger.warn(
"[Recordings::VerifyYoutubeReplayJob] keep_temp recording=#{recording.id} " \
"youtube_not_embeddable"
)
end
return
end
if result.retriable? && attempt + 1 < max
delay = backoff_secs(attempt)
Rails.logger.info(
"[Recordings::VerifyYoutubeReplayJob] retry recording=#{recording.id} " \
"attempt=#{attempt + 1}/#{max} in=#{delay}s msg=#{result.message}"
)
self.class.perform_in(delay.seconds, recording.id, attempt + 1)
return
end
Rails.logger.warn(
"[Recordings::VerifyYoutubeReplayJob] give_up recording=#{recording.id} " \
"attempts=#{attempt + 1} msg=#{result.message} " \
"temp_expires_at=#{recording.temp_expires_at.inspect}"
)
end
private
def backoff_secs(attempt)
base = MatchLiveTv.youtube_replay_verify_base_interval_secs
# 300, 600, 1200, ... capped at 1h
[base * (2**attempt), 3600].min
end
end
end
@@ -0,0 +1,45 @@
# 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"
KICK_DEBOUNCE_KEY = "streams:autoscaler:kick"
KICK_DEBOUNCE_SECS = ENV.fetch("STREAM_AUTOSCALE_KICK_DEBOUNCE_SECS", "5").to_i
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
# Kick immediato (es. NoCapacity / soft-free bassi). Debounce anti-flood Sidekiq.
def self.kick!
return false unless Streams::Autoscaler.enabled?
return false unless redis
return false unless redis.set(KICK_DEBOUNCE_KEY, "1", nx: true, ex: KICK_DEBOUNCE_SECS)
perform_async
true
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
+2 -2
View File
@@ -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 class YoutubeRelayEnsureJob < ApplicationJob
queue_as :default queue_as Streams::YoutubeRelay::QUEUE
def perform(session_id) def perform(session_id)
session = StreamSession.find_by(id: session_id) session = StreamSession.find_by(id: session_id)
+16 -3
View File
@@ -1,10 +1,23 @@
# Ferma il relay solo sull'owner. Se il job gira su un altro host, requeue breve.
class YoutubeRelayStopJob < ApplicationJob 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) session = StreamSession.find_by(id: session_id)
return unless session 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
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
+26 -6
View File
@@ -5,17 +5,37 @@ module Billing
@club = @invoice.club @club = @invoice.club
I18n.with_locale(I18n.locale) do I18n.with_locale(I18n.locale) do
prefix = t("mailers.invoice.attachment_prefix") attach_invoice_pdf!
attachments["#{prefix}-#{@invoice.number.parameterize}.pdf"] = {
mime_type: "application/pdf",
content: @invoice.pdf.download
}
mail( mail(
to: @club.billing_email, to: @club.billing_email,
subject: t("mailers.invoice.subject", number: @invoice.display_number) subject: t("mailers.invoice.subject", number: @invoice.display_number)
) )
end end
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
end end
+25
View File
@@ -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
+1
View File
@@ -2,6 +2,7 @@ class AdminAccount < ApplicationRecord
include PasswordComplexity include PasswordComplexity
has_secure_password 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 validates :username, presence: true, uniqueness: true
end end
+14
View File
@@ -0,0 +1,14 @@
# frozen_string_literal: true
class AnalyticsEvent < ApplicationRecord
self.table_name = "analytics_events"
EVENT_TYPES = %w[click scroll pageview move].freeze
DEVICES = %w[mobile tablet desktop].freeze
validates :page_path, presence: true
validates :device, inclusion: { in: DEVICES }
validates :event_type, inclusion: { in: EVENT_TYPES }
validates :occurred_at, presence: true
validates :weight, numericality: { greater_than: 0, less_than_or_equal_to: 500 }
end
+13
View File
@@ -0,0 +1,13 @@
# frozen_string_literal: true
class AnalyticsPageCell < ApplicationRecord
self.table_name = "analytics_page_cells"
GRID_SIZE = 40
validates :day, :page_path, :device, :cell_x, :cell_y, presence: true
validates :device, inclusion: { in: AnalyticsEvent::DEVICES }
validates :cell_x, :cell_y, inclusion: { in: 0...(GRID_SIZE) }
validates :click_count, numericality: { greater_than_or_equal_to: 0 }
validates :move_count, numericality: { greater_than_or_equal_to: 0 }
end
@@ -0,0 +1,12 @@
# frozen_string_literal: true
class AnalyticsPageSnapshot < ApplicationRecord
self.table_name = "analytics_page_snapshots"
has_one_attached :image
validates :page_path, presence: true
validates :device, inclusion: { in: AnalyticsEvent::DEVICES }
validates :width, :height, numericality: { greater_than: 0 }
validates :captured_at, presence: true
end
+14
View File
@@ -0,0 +1,14 @@
# frozen_string_literal: true
class AnalyticsPageStat < ApplicationRecord
self.table_name = "analytics_page_stats"
validates :day, :page_path, :device, presence: true
validates :device, inclusion: { in: AnalyticsEvent::DEVICES }
def avg_scroll_pct
return 0 if scroll_samples.to_i <= 0
(scroll_sum_pct.to_f / scroll_samples).round
end
end
+240
View File
@@ -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
+34
View File
@@ -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
+10 -2
View File
@@ -2,14 +2,17 @@ module Billing
class Payment < ApplicationRecord class Payment < ApplicationRecord
self.table_name = "billing_payments" 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 belongs_to :club
has_one :invoice, class_name: "Billing::Invoice", foreign_key: :billing_payment_id, dependent: :nullify 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 :amount_cents, numericality: { greater_than: 0 }
validates :currency, presence: true validates :currency, presence: true
validates :status, inclusion: { in: STATUSES } validates :status, inclusion: { in: STATUSES }
validates :provider, inclusion: { in: PROVIDERS }
validates :stripe_invoice_id, uniqueness: true, allow_nil: true validates :stripe_invoice_id, uniqueness: true, allow_nil: true
scope :recent, -> { order(paid_at: :desc, created_at: :desc) } scope :recent, -> { order(paid_at: :desc, created_at: :desc) }
@@ -62,7 +65,12 @@ module Billing
end end
def display_status 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 end
def invoice_for_display 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
+12
View File
@@ -1,5 +1,6 @@
class Club < ApplicationRecord class Club < ApplicationRecord
include Brandable include Brandable
include Coverable
include ClubBillingProfile include ClubBillingProfile
has_many :club_memberships, dependent: :destroy has_many :club_memberships, dependent: :destroy
@@ -7,8 +8,11 @@ class Club < ApplicationRecord
has_many :teams, dependent: :destroy has_many :teams, dependent: :destroy
has_one :youtube_credential, dependent: :destroy has_one :youtube_credential, dependent: :destroy
has_one :subscription, 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_payments, class_name: "Billing::Payment", dependent: :destroy
has_many :billing_invoices, class_name: "Billing::Invoice", 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 :name, presence: true
validates :sport, presence: true validates :sport, presence: true
@@ -23,6 +27,14 @@ class Club < ApplicationRecord
club_memberships.exists?(user: user, role: "owner") club_memberships.exists?(user: user, role: "owner")
end end
def active_billing_quote
billing_quote
end
def pending_transfer_order
billing_transfer_orders.awaiting_payment.first
end
private private
def branding_parent def branding_parent
+50
View File
@@ -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
+6
View File
@@ -1,4 +1,6 @@
class Match < ApplicationRecord class Match < ApplicationRecord
include Coverable
belongs_to :team belongs_to :team
has_many :stream_sessions, dependent: :destroy has_many :stream_sessions, dependent: :destroy
@@ -174,4 +176,8 @@ class Match < ApplicationRecord
rescue ArgumentError => e rescue ArgumentError => e
errors.add(:scoring_rules, e.message) errors.add(:scoring_rules, e.message)
end end
def cover_parent
nil
end
end end
+1 -1
View File
@@ -4,7 +4,7 @@ module Ops
KINDS = %w[ KINDS = %w[
disk_space recordings_size service_down http_public http_rails rails_latency 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 ].freeze
SEVERITIES = %w[critical warning info].freeze SEVERITIES = %w[critical warning info].freeze
STATUSES = %w[open acknowledged resolved].freeze STATUSES = %w[open acknowledged resolved].freeze
+4
View File
@@ -84,6 +84,10 @@ class Plan < ApplicationRecord
feature("youtube_enabled") == true feature("youtube_enabled") == true
end end
def custom_cover_enabled?
feature("custom_cover") == true
end
def premium? def premium?
slug.in?(PREMIUM_SLUGS) slug.in?(PREMIUM_SLUGS)
end end
+51
View File
@@ -0,0 +1,51 @@
# frozen_string_literal: true
class PlatformCostEntry < ApplicationRecord
DEFAULT_LABEL = "Piattaforma produzione"
validates :month, presence: true
validates :label, presence: true, length: { maximum: 120 }
validates :amount_cents, numericality: { only_integer: true, greater_than: 0 }
before_validation :normalize_month
scope :for_month, ->(date) { where(month: date.to_date.beginning_of_month) }
scope :ordered, -> { order(month: :desc, created_at: :desc) }
def month=(value)
if value.is_a?(String) && value.match?(/\A\d{4}-\d{2}\z/)
super(Date.strptime(value, "%Y-%m"))
else
super(value)
end
end
def amount_euros
amount_cents.to_f / 100.0
end
private
def normalize_month
return if month.blank?
parsed =
case month
when Date
month
when Time, ActiveSupport::TimeWithZone
month.to_date
when String
if month.match?(/\A\d{4}-\d{2}\z/)
Date.strptime(month, "%Y-%m")
else
Date.parse(month)
end
else
Date.parse(month.to_s)
end
self.month = parsed.beginning_of_month
rescue ArgumentError, TypeError
errors.add(:month, :invalid)
end
end
+51 -4
View File
@@ -2,6 +2,8 @@ class Recording < ApplicationRecord
STATUSES = %w[processing ready expired failed].freeze STATUSES = %w[processing ready expired failed].freeze
PRIVACY_STATUSES = %w[public unlisted].freeze PRIVACY_STATUSES = %w[public unlisted].freeze
STORAGE_BACKENDS = %w[local s3].freeze STORAGE_BACKENDS = %w[local s3].freeze
STORAGE_POLICIES = %w[temporary retained none].freeze
REPLAY_SOURCES = %w[youtube matchlivetv none].freeze
belongs_to :stream_session belongs_to :stream_session
belongs_to :team belongs_to :team
@@ -9,6 +11,7 @@ class Recording < ApplicationRecord
validates :status, inclusion: { in: STATUSES } validates :status, inclusion: { in: STATUSES }
validates :privacy_status, inclusion: { in: PRIVACY_STATUSES } validates :privacy_status, inclusion: { in: PRIVACY_STATUSES }
validates :storage_backend, inclusion: { in: STORAGE_BACKENDS } validates :storage_backend, inclusion: { in: STORAGE_BACKENDS }
validates :storage_policy, inclusion: { in: STORAGE_POLICIES }
scope :not_deleted, -> { where(deleted_at: nil) } scope :not_deleted, -> { where(deleted_at: nil) }
scope :ready, lambda { scope :ready, lambda {
@@ -23,7 +26,16 @@ class Recording < ApplicationRecord
ready.where(expires_at: ..days.days.from_now) ready.where(expires_at: ..days.days.from_now)
} }
scope :expired_pending_purge, lambda { scope :expired_pending_purge, lambda {
not_deleted.where(status: %w[ready failed]).where("expires_at IS NOT NULL AND expires_at <= ?", Time.current) not_deleted
.where(storage_policy: "retained")
.where(status: %w[ready failed])
.where("expires_at IS NOT NULL AND expires_at <= ?", Time.current)
}
scope :temporary_media_pending_purge, lambda {
not_deleted
.where(storage_policy: "temporary")
.where(local_media_purged_at: nil)
.where("temp_expires_at IS NOT NULL AND temp_expires_at <= ?", Time.current)
} }
scope :search_replays, lambda { |query| scope :search_replays, lambda { |query|
q = query.to_s.strip q = query.to_s.strip
@@ -46,19 +58,30 @@ class Recording < ApplicationRecord
end end
def playback_stream_url def playback_stream_url
return nil unless ready? && stream_session_id.present? return nil unless ready? && storage_key.present? && stream_session_id.present?
"#{MatchLiveTv.app_public_url.chomp('/')}/replay/#{stream_session_id}/stream" "#{MatchLiveTv.app_public_url.chomp('/')}/replay/#{stream_session_id}/stream"
end end
def thumbnail_url def thumbnail_url
return youtube_thumbnail_url if thumbnail_storage_key.blank? && youtube_video_id.present?
return nil unless thumbnail_storage_key.present? && stream_session_id.present? return nil unless thumbnail_storage_key.present? && stream_session_id.present?
"#{MatchLiveTv.app_public_url.chomp('/')}/replay/#{stream_session_id}/thumbnail" "#{MatchLiveTv.app_public_url.chomp('/')}/replay/#{stream_session_id}/thumbnail"
end end
def youtube_thumbnail_url
return nil if youtube_video_id.blank?
return nil if youtube_video_id.to_s.start_with?("mock_")
meta_url = metadata.is_a?(Hash) ? metadata.dig("youtube", "thumbnail_url") : nil
return meta_url if meta_url.present?
"https://i.ytimg.com/vi/#{youtube_video_id}/hqdefault.jpg"
end
def download_api_path def download_api_path
return nil unless ready? return nil unless ready? && storage_key.present?
"/api/v1/recordings/#{id}/download" "/api/v1/recordings/#{id}/download"
end end
@@ -70,6 +93,30 @@ class Recording < ApplicationRecord
"https://www.youtube.com/watch?v=#{youtube_video_id}" "https://www.youtube.com/watch?v=#{youtube_video_id}"
end end
def temporary_storage?
storage_policy == "temporary"
end
def retained_storage?
storage_policy == "retained"
end
def local_media_purged?
local_media_purged_at.present?
end
# Sorgente fisica del player in archivio (non lo storage policy).
def replay_source
return "youtube" if youtube_watch_url.present? || (ready? && youtube_video_id.present?)
return "matchlivetv" if ready? && storage_key.present?
"none"
end
def available_in_archive?
ready? && (youtube_watch_url.present? || youtube_video_id.present? || storage_key.present?)
end
def ready? def ready?
status == "ready" && !deleted? && (expires_at.nil? || expires_at.future?) status == "ready" && !deleted? && (expires_at.nil? || expires_at.future?)
end end
@@ -154,7 +201,7 @@ class Recording < ApplicationRecord
end end
def playable_on_site? def playable_on_site?
ready? && storage_key.present? available_in_archive?
end end
def days_until_expiry def days_until_expiry
+1 -1
View File
@@ -1,7 +1,7 @@
class StreamEvent < ApplicationRecord class StreamEvent < ApplicationRecord
EVENT_TYPES = %w[ EVENT_TYPES = %w[
connected disconnected reconnected quality_changed error paused resumed 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 ].freeze
belongs_to :stream_session belongs_to :stream_session
+54
View File
@@ -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
+42 -4
View File
@@ -1,14 +1,17 @@
class StreamSession < ApplicationRecord class StreamSession < ApplicationRecord
include AASM include AASM
PLATFORMS = %w[matchlivetv youtube facebook twitch].freeze LIVE_PLATFORMS = %w[matchlivetv youtube].freeze
PLATFORMS = (LIVE_PLATFORMS + %w[facebook twitch]).freeze
STATUSES = %w[idle connecting live reconnecting paused ended error].freeze STATUSES = %w[idle connecting live reconnecting paused ended error].freeze
PRIVACY_STATUSES = %w[public unlisted private].freeze PRIVACY_STATUSES = %w[public unlisted private].freeze
belongs_to :match belongs_to :match
belongs_to :user belongs_to :user
belongs_to :stream_node, optional: true
has_many :stream_events, dependent: :destroy has_many :stream_events, dependent: :destroy
has_one :score_state, dependent: :destroy has_one :score_state, dependent: :destroy
has_one :recording
has_many :device_states, dependent: :destroy has_many :device_states, dependent: :destroy
validates :platform, inclusion: { in: PLATFORMS } validates :platform, inclusion: { in: PLATFORMS }
@@ -18,6 +21,7 @@ class StreamSession < ApplicationRecord
before_validation :normalize_privacy_status before_validation :normalize_privacy_status
before_validation :ensure_publish_token, on: :create before_validation :ensure_publish_token, on: :create
before_validation :snapshot_ingest_from_node, if: -> { stream_node.present? }
scope :broadcasting, -> { where(status: %w[live connecting reconnecting paused]) } scope :broadcasting, -> { where(status: %w[live connecting reconnecting paused]) }
scope :publicly_listed, -> { where(privacy_status: "public") } scope :publicly_listed, -> { where(privacy_status: "public") }
@@ -72,10 +76,18 @@ class StreamSession < ApplicationRecord
end end
end end
def ingest_slug_display
stream_node&.slug.presence || ingest_slug
end
def ingest_role_display
stream_node&.role.presence || ingest_role
end
def rtmp_ingest_url def rtmp_ingest_url
# RootEncoder richiede rtmp://host:port/app/stream (due segmenti). # RootEncoder richiede rtmp://host:port/app/stream (due segmenti).
# MediaMTX path = live/match_{uuid} (no ?token= nel path). # MediaMTX path = live/match_{uuid} (no ?token= nel path).
"#{MatchLiveTv.mediamtx_rtmp_url}/#{mediamtx_path_name}" "#{rtmp_base_url.chomp('/')}/#{mediamtx_path_name}"
end end
def mediamtx_path_name def mediamtx_path_name
@@ -91,8 +103,29 @@ class StreamSession < ApplicationRecord
end end
def hls_playback_url def hls_playback_url
base = MatchLiveTv.hls_public_url.chomp("/") "#{hls_base_url.chomp('/')}/#{effective_hls_path_name}/index.m3u8"
"#{base}/#{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 end
def effective_hls_path_name def effective_hls_path_name
@@ -185,6 +218,11 @@ class StreamSession < ApplicationRecord
self.privacy_status = "unlisted" if privacy_status == "private" self.privacy_status = "unlisted" if privacy_status == "private"
end end
def snapshot_ingest_from_node
self.ingest_slug = stream_node.slug
self.ingest_role = stream_node.role
end
def record_ended_timestamps! def record_ended_timestamps!
now = Time.current now = Time.current
update!(ended_at: now) if ended_at.nil? update!(ended_at: now) if ended_at.nil?
+4
View File
@@ -35,4 +35,8 @@ class Subscription < ApplicationRecord
def admin_comped? def admin_comped?
admin_comped admin_comped
end end
def bank_transfer?
premium? && stripe_subscription_id.blank? && !admin_comped? && current_period_end.present?
end
end end
+5
View File
@@ -1,5 +1,6 @@
class Team < ApplicationRecord class Team < ApplicationRecord
include Brandable include Brandable
include Coverable
belongs_to :club belongs_to :club
has_many :user_teams, dependent: :destroy has_many :user_teams, dependent: :destroy
@@ -99,6 +100,10 @@ class Team < ApplicationRecord
club club
end end
def cover_parent
club
end
def assign_slug def assign_slug
self.slug = Teams::GenerateSlug.call(self) if slug.blank? self.slug = Teams::GenerateSlug.call(self) if slug.blank?
end end
+16 -2
View File
@@ -4,19 +4,33 @@ class YoutubeCredential < ApplicationRecord
attr_encrypted :access_token, attr_encrypted :access_token,
key: :encryption_key, key: :encryption_key,
attribute: "access_token_encrypted", attribute: "access_token_encrypted",
mode: :single_iv_salt mode: :single_iv_and_salt,
algorithm: "aes-256-cbc",
iv: :encryption_iv
attr_encrypted :refresh_token, attr_encrypted :refresh_token,
key: :encryption_key, key: :encryption_key,
attribute: "refresh_token_encrypted", attribute: "refresh_token_encrypted",
mode: :single_iv_salt mode: :single_iv_and_salt,
algorithm: "aes-256-cbc",
iv: :encryption_iv
def expired? def expired?
expires_at.present? && expires_at < Time.current expires_at.present? && expires_at < Time.current
end end
def usable?
refresh_token.present?
rescue ArgumentError, OpenSSL::Cipher::CipherError, NoMethodError
false
end
private private
def encryption_key def encryption_key
Rails.application.secret_key_base[0, 32] Rails.application.secret_key_base[0, 32]
end end
def encryption_iv
encryption_key[0, 16]
end
end end
@@ -0,0 +1,150 @@
# frozen_string_literal: true
module Admin
class CostAnalytics
TREND_MONTHS = 12
def initialize(month:)
@month = month.to_date.beginning_of_month
end
def summary
build_period(@month)
end
def trend(months: TREND_MONTHS)
start_month = @month - (months - 1).months
months_list = (0...months).map { |i| start_month + i.months }
months_list.map { |m| build_period(m) }
end
def club_breakdown
range = month_range(@month)
rows = session_rows(range)
total_secs = rows.sum { |row| row.total_secs.to_i }
total_cost_cents = PlatformCostEntry.for_month(@month).sum(:amount_cents)
revenue_by_club = revenue_by_club(range)
storage_by_club = storage_by_club_index
club_ids = rows.map(&:club_id)
clubs = Club.where(id: club_ids).includes(subscription: :plan).index_by(&:id)
rows.map do |row|
club = clubs[row.club_id]
secs = row.total_secs.to_i
hours = secs / 3600.0
sessions_count = row.sessions_count.to_i
share = total_secs.positive? ? secs.to_f / total_secs : 0.0
allocated_cents = (total_cost_cents * share).round
revenue_cents = revenue_by_club[row.club_id].to_i
storage_bytes = storage_by_club[row.club_id].to_i
{
club_id: row.club_id,
club_name: row.club_name,
plan_slug: club&.subscription&.plan&.slug,
sessions: sessions_count,
hours: hours.round(2),
hours_share_pct: (share * 100).round(1),
allocated_cost_cents: allocated_cents,
revenue_cents: revenue_cents,
margin_cents: revenue_cents - allocated_cents,
cost_per_hour_cents: hours.positive? ? (allocated_cents / hours).round : nil,
cost_per_session_cents: sessions_count.positive? ? (allocated_cents / sessions_count) : nil,
storage_bytes: storage_bytes
}
end.sort_by { |row| [-row[:hours], row[:club_name]] }
end
private
def build_period(month)
range = month_range(month)
cost_cents = PlatformCostEntry.for_month(month).sum(:amount_cents)
sessions_scope = ended_sessions.where(ended_at: range)
total_secs = sessions_scope.sum(:total_duration_secs).to_i
hours = total_secs / 3600.0
sessions = sessions_scope.count
clubs_active = distinct_active_clubs(range)
revenue_cents = Billing::Payment.where(status: "paid", paid_at: range).sum(:amount_cents)
storage_bytes = Recording.ready.sum(:byte_size).to_i
storage_gb = storage_bytes.positive? ? storage_bytes / (1024.0**3) : 0.0
kpis = compute_kpis(
cost_cents: cost_cents,
hours: hours,
sessions: sessions,
clubs_active: clubs_active,
revenue_cents: revenue_cents,
storage_gb: storage_gb
)
{
month: month,
cost_cents: cost_cents,
sessions: sessions,
hours: hours.round(2),
clubs_active: clubs_active,
revenue_cents: revenue_cents,
storage_bytes: storage_bytes,
**kpis
}
end
def compute_kpis(cost_cents:, hours:, sessions:, clubs_active:, revenue_cents:, storage_gb:)
margin_cents = revenue_cents.to_i - cost_cents.to_i
margin_pct = revenue_cents.to_i.positive? ? ((margin_cents.to_f / revenue_cents.to_i) * 100).round(1) : nil
hours_f = hours.to_f
sessions_i = sessions.to_i
clubs_i = clubs_active.to_i
storage_f = storage_gb.to_f
cost_i = cost_cents.to_i
{
cost_per_hour_cents: hours_f.positive? ? (cost_i / hours_f).round : nil,
cost_per_session_cents: sessions_i.positive? ? (cost_i / sessions_i) : nil,
cost_per_club_cents: clubs_i.positive? ? (cost_i / clubs_i) : nil,
revenue_per_hour_cents: hours_f.positive? ? (revenue_cents.to_i / hours_f).round : nil,
cost_per_gb_cents: storage_f.positive? ? (cost_i / storage_f).round : nil,
margin_cents: margin_cents,
margin_pct: margin_pct
}
end
def month_range(month)
month.beginning_of_month.beginning_of_day..month.end_of_month.end_of_day
end
def ended_sessions
StreamSession.where(status: "ended")
end
def distinct_active_clubs(range)
ended_sessions
.where(ended_at: range)
.joins(match: { team: :club })
.distinct
.count("clubs.id")
end
def session_rows(range)
ended_sessions
.where(ended_at: range)
.joins(match: { team: :club })
.group("clubs.id", "clubs.name")
.select(
"clubs.id AS club_id",
"clubs.name AS club_name",
"COUNT(stream_sessions.id) AS sessions_count",
"SUM(stream_sessions.total_duration_secs) AS total_secs"
)
end
def revenue_by_club(range)
Billing::Payment.where(status: "paid", paid_at: range).group(:club_id).sum(:amount_cents)
end
def storage_by_club_index
Recording.ready.joins(:team).group("teams.club_id").sum(:byte_size)
end
end
end
+107
View File
@@ -0,0 +1,107 @@
# frozen_string_literal: true
module Analytics
class Aggregate
BATCH = 500
UPSERT_RETRIES = 3
def call
loop do
ids = AnalyticsEvent.order(:occurred_at).limit(BATCH).pluck(:id)
break if ids.empty?
events = AnalyticsEvent.where(id: ids).to_a
apply_points(events.select { |e| e.event_type == "click" }, :click_count)
apply_points(events.select { |e| e.event_type == "move" }, :move_count)
apply_stats(events.select { |e| %w[pageview scroll].include?(e.event_type) })
AnalyticsEvent.where(id: ids).delete_all
end
end
def purge_old!(retention: 30.days)
AnalyticsEvent.where("occurred_at < ?", retention.ago).delete_all
AnalyticsPageCell.where("day < ?", retention.ago.to_date).delete_all
AnalyticsPageStat.where("day < ?", retention.ago.to_date).delete_all
AnalyticsPageSnapshot.where("captured_at < ?", retention.ago).find_each do |snap|
snap.image.purge if snap.image.attached?
snap.destroy!
end
end
private
def apply_points(events, counter_attr)
return if events.empty?
grid = AnalyticsPageCell::GRID_SIZE
grouped = Hash.new(0)
events.each do |e|
next if e.x_pct.nil? || e.y_pct.nil?
cell_x = [[(e.x_pct.to_f / 100 * grid).floor, grid - 1].min, 0].max
cell_y = [[(e.y_pct.to_f / 100 * grid).floor, grid - 1].min, 0].max
key = [e.occurred_at.in_time_zone.to_date, e.page_path, e.device, cell_x, cell_y]
grouped[key] += e.weight.to_i.clamp(1, 500)
end
now = Time.current
grouped.each do |(day, page_path, device, cell_x, cell_y), count|
with_unique_retry do
cell = AnalyticsPageCell.find_or_initialize_by(
day: day, page_path: page_path, device: device, cell_x: cell_x, cell_y: cell_y
)
cell[counter_attr] = cell[counter_attr].to_i + count
cell.created_at ||= now
cell.updated_at = now
cell.save!
end
end
end
def apply_stats(events)
return if events.empty?
grouped = {}
events.each do |e|
day = e.occurred_at.in_time_zone.to_date
key = [day, e.page_path, e.device]
bucket = grouped[key] ||= { pageviews: 0, scroll_samples: 0, scroll_sum: 0, max_scroll: 0 }
case e.event_type
when "pageview"
bucket[:pageviews] += e.weight.to_i.clamp(1, 500)
when "scroll"
pct = e.scroll_pct.to_f.round
bucket[:scroll_samples] += 1
bucket[:scroll_sum] += pct
bucket[:max_scroll] = [bucket[:max_scroll], pct].max
end
end
now = Time.current
grouped.each do |(day, page_path, device), vals|
with_unique_retry do
stat = AnalyticsPageStat.find_or_initialize_by(day: day, page_path: page_path, device: device)
stat.pageview_count = stat.pageview_count.to_i + vals[:pageviews]
stat.scroll_samples = stat.scroll_samples.to_i + vals[:scroll_samples]
stat.scroll_sum_pct = stat.scroll_sum_pct.to_i + vals[:scroll_sum]
stat.max_scroll_pct = [stat.max_scroll_pct.to_i, vals[:max_scroll]].max
stat.created_at ||= now
stat.updated_at = now
stat.save!
end
end
end
def with_unique_retry
attempts = 0
begin
attempts += 1
yield
rescue ActiveRecord::RecordNotUnique
raise if attempts >= UPSERT_RETRIES
retry
end
end
end
end
+134
View File
@@ -0,0 +1,134 @@
# frozen_string_literal: true
module Analytics
class Ingest
MAX_BATCH = 80
RATE_LIMIT_PER_MINUTE = 90
Result = Struct.new(:accepted, :rejected, :rate_limited, keyword_init: true)
def initialize(events:, remote_ip:)
@events = Array(events)
@remote_ip = remote_ip.to_s.presence || "unknown"
end
def call
return Result.new(accepted: 0, rejected: 0, rate_limited: true) if rate_limited?
slice = @events.first(MAX_BATCH)
accepted = 0
rejected = 0
rows = []
slice.each do |raw|
attrs = build_attrs(raw)
if attrs
rows << attrs
accepted += 1
else
rejected += 1
end
end
AnalyticsEvent.insert_all(rows) if rows.any?
if rows.any?
begin
Analytics::Aggregate.new.call
rescue ActiveRecord::RecordNotUnique => e
# Dopo i retry interni: non far fallire la richiesta analytics.
Rails.logger.warn("[analytics] aggregate race after retries, enqueue job: #{e.message}")
Analytics::AggregateJob.perform_later
end
end
Result.new(accepted: accepted, rejected: rejected + (@events.size - slice.size), rate_limited: false)
end
private
def rate_limited?
key = "analytics:ingest:#{@remote_ip}"
redis = Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0"))
count = redis.incr(key)
redis.expire(key, 60) if count == 1
count > RATE_LIMIT_PER_MINUTE
rescue StandardError => e
Rails.logger.warn("[analytics] rate limit unavailable: #{e.class}")
false
end
def build_attrs(raw)
data = raw.is_a?(Hash) ? raw.with_indifferent_access : {}
event_type = data[:type].to_s.presence || data[:event_type].to_s
return nil unless AnalyticsEvent::EVENT_TYPES.include?(event_type)
device = data[:device].to_s
return nil unless AnalyticsEvent::DEVICES.include?(device)
page_path = PathNormalizer.normalize(data[:path] || data[:page_path])
return nil if PathNormalizer.excluded?(page_path)
return nil if page_path.length > 200
occurred_at = parse_time(data[:ts] || data[:occurred_at]) || Time.current
weight = data[:n].presence || data[:weight].presence || 1
weight = Integer(weight)
weight = 1 if weight < 1
weight = 500 if weight > 500
attrs = {
page_path: page_path,
device: device,
event_type: event_type,
occurred_at: occurred_at,
weight: weight,
x_pct: nil,
y_pct: nil,
scroll_pct: nil,
created_at: Time.current,
updated_at: Time.current
}
case event_type
when "click", "move"
x = clamp_pct(data[:x] || data[:x_pct])
y = clamp_pct(data[:y] || data[:y_pct])
return nil if x.nil? || y.nil?
attrs[:x_pct] = x
attrs[:y_pct] = y
when "scroll"
scroll = clamp_pct(data[:scroll] || data[:scroll_pct])
return nil if scroll.nil?
attrs[:scroll_pct] = scroll
when "pageview"
# no extra fields
end
attrs
rescue ArgumentError, TypeError
nil
end
def clamp_pct(value)
return nil if value.nil?
n = Float(value)
return nil if n.nan? || n.infinite?
return nil if n.negative? || n > 100
n.round(2)
rescue ArgumentError, TypeError
nil
end
def parse_time(value)
return value if value.is_a?(Time) || value.is_a?(ActiveSupport::TimeWithZone)
return Time.zone.at(value.to_i / 1000.0) if value.is_a?(Numeric) || value.to_s.match?(/\A\d{10,13}\z/)
Time.zone.parse(value.to_s)
rescue ArgumentError, TypeError
nil
end
end
end
@@ -0,0 +1,42 @@
# frozen_string_literal: true
module Analytics
class PathNormalizer
UUID_RE = /\A[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}\z/i
NUMERIC_RE = /\A\d{6,}\z/
TOKENISH_RE = /\A[A-Za-z0-9_\-]{22,}\z/
EXCLUDED_PREFIXES = %w[/admin /regia /analytics /cable /api /webhooks /internal].freeze
class << self
def normalize(raw_path)
path = raw_path.to_s.split("?", 2).first.to_s
path = "/" if path.blank?
path = "/#{path}" unless path.start_with?("/")
path = path.gsub(%r{/{2,}}, "/").chomp("/")
path = "/" if path.blank?
parts = path.split("/").map do |segment|
next segment if segment.blank?
next ":id" if UUID_RE.match?(segment)
next ":id" if NUMERIC_RE.match?(segment)
next ":token" if TOKENISH_RE.match?(segment) && segment.match?(/[A-Z]/)
segment
end
normalized = parts.join("/")
normalized = "/" if normalized.blank?
normalized
end
def excluded?(raw_or_normalized)
path = normalize(raw_or_normalized)
return true if EXCLUDED_PREFIXES.any? { |p| path == p || path.start_with?("#{p}/") }
return true if path.match?(%r{\A/live/[^/]+})
return true if path.match?(%r{\A/replay/[^/]+})
false
end
end
end
end
@@ -0,0 +1,84 @@
# frozen_string_literal: true
module Analytics
class SnapshotIngest
MAX_BYTES = 1_800_000
MIN_REFRESH = 12.hours
RATE_LIMIT_PER_HOUR = 20
Result = Struct.new(:ok, :skipped, :rate_limited, :error, keyword_init: true)
def initialize(path:, device:, width:, height:, image:, remote_ip:)
@path = path
@device = device.to_s
@width = width
@height = height
@image = image
@remote_ip = remote_ip.to_s.presence || "unknown"
end
def call
return Result.new(ok: false, skipped: false, rate_limited: true) if rate_limited?
page_path = PathNormalizer.normalize(@path)
return Result.new(ok: false, skipped: true, error: "path") if PathNormalizer.excluded?(page_path)
return Result.new(ok: false, skipped: true, error: "device") unless AnalyticsEvent::DEVICES.include?(@device)
return Result.new(ok: false, skipped: true, error: "image") unless valid_image?
width = Integer(@width)
height = Integer(@height)
return Result.new(ok: false, skipped: true, error: "size") if width < 200 || height < 200 || width > 6000 || height > 20000
snap = AnalyticsPageSnapshot.find_or_initialize_by(page_path: page_path, device: @device)
if snap.persisted? && snap.captured_at.present? && snap.captured_at > MIN_REFRESH.ago && snap.image.attached?
return Result.new(ok: true, skipped: true)
end
snap.width = width
snap.height = height
snap.captured_at = Time.current
snap.save!
snap.image.purge if snap.image.attached?
snap.image.attach(
io: image_io,
filename: "analytics-#{Digest::SHA1.hexdigest(page_path)[0, 12]}-#{@device}.jpg",
content_type: "image/jpeg"
)
Result.new(ok: true, skipped: false)
rescue ArgumentError, TypeError, ActiveRecord::RecordInvalid => e
Result.new(ok: false, skipped: true, error: e.class.name)
end
private
def rate_limited?
key = "analytics:snapshot:#{@remote_ip}"
redis = Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0"))
count = redis.incr(key)
redis.expire(key, 3600) if count == 1
count > RATE_LIMIT_PER_HOUR
rescue StandardError => e
Rails.logger.warn("[analytics] snapshot rate limit unavailable: #{e.class}")
false
end
def valid_image?
return false if @image.blank?
return false unless @image.respond_to?(:tempfile) || @image.respond_to?(:read)
return false if @image.respond_to?(:size) && @image.size.to_i > MAX_BYTES
content_type = @image.content_type.to_s if @image.respond_to?(:content_type)
content_type.blank? || content_type.start_with?("image/")
end
def image_io
if @image.respond_to?(:tempfile)
@image.tempfile.rewind
@image.tempfile
else
StringIO.new(@image.read)
end
end
end
end
@@ -0,0 +1,42 @@
# frozen_string_literal: true
module Analytics
module Suppress
COOKIE_NAME = "mltv_analytics_suppress"
COOKIE_MAX_AGE = 7 * 24 * 60 * 60
module_function
def active?(cookie_value)
cookie_value.to_s == "1"
end
def enable!(cookie_jar)
cookie_jar[COOKIE_NAME] = cookie_options(value: "1", expires: COOKIE_MAX_AGE.seconds.from_now)
end
def disable!(cookie_jar)
cookie_jar.delete(
COOKIE_NAME,
path: "/",
same_site: :lax,
secure: cookie_secure?
)
end
def cookie_options(value:, expires:)
{
value: value,
expires: expires,
path: "/",
httponly: true,
same_site: :lax,
secure: cookie_secure?
}
end
def cookie_secure?
Rails.application.config.force_ssl || Rails.env.production?
end
end
end
@@ -23,6 +23,7 @@ module Billing
raise Error, "Piano non valido" unless @plan_slug.in?(VALID_PLANS) raise Error, "Piano non valido" unless @plan_slug.in?(VALID_PLANS)
release_stripe_schedule! release_stripe_schedule!
cancel_awaiting_transfers!
AssignPlan.call( AssignPlan.call(
club: @club, club: @club,
@@ -97,5 +98,9 @@ module Billing
ensure ensure
sub&.update!(stripe_schedule_id: nil, pending_plan_id: nil, pending_billing_interval: nil) sub&.update!(stripe_schedule_id: nil, pending_plan_id: nil, pending_billing_interval: nil)
end end
def cancel_awaiting_transfers!
@club.billing_transfer_orders.awaiting_payment.find_each(&:cancel!)
end
end end
end end
@@ -2,13 +2,14 @@ module Billing
class AttachPaymentInvoice class AttachPaymentInvoice
class Error < StandardError; end class Error < StandardError; end
def self.call(payment:, pdf:) def self.call(payment:, pdf:, mailer_action: :invoice_pdf)
new(payment: payment, pdf: pdf).call new(payment: payment, pdf: pdf, mailer_action: mailer_action).call
end end
def initialize(payment:, pdf:) def initialize(payment:, pdf:, mailer_action: :invoice_pdf)
@payment = payment @payment = payment
@pdf = pdf @pdf = pdf
@mailer_action = mailer_action
end end
def call def call
@@ -17,7 +18,7 @@ module Billing
club = @payment.club club = @payment.club
invoice = @payment.invoice || build_invoice!(club) invoice = @payment.invoice || build_invoice!(club)
IssueInvoice.call(invoice: invoice, pdf: @pdf) IssueInvoice.call(invoice: invoice, pdf: @pdf, mailer_action: @mailer_action)
end end
private 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
+11 -5
View File
@@ -2,13 +2,17 @@ module Billing
class IssueInvoice class IssueInvoice
class Error < StandardError; end class Error < StandardError; end
def self.call(invoice:, pdf: nil) MAILER_ACTIONS = %i[invoice_pdf plan_activated_with_invoice].freeze
new(invoice: invoice, pdf: pdf).call
def self.call(invoice:, pdf: nil, mailer_action: :invoice_pdf)
new(invoice: invoice, pdf: pdf, mailer_action: mailer_action).call
end end
def initialize(invoice:, pdf: nil) def initialize(invoice:, pdf: nil, mailer_action: :invoice_pdf)
@invoice = invoice @invoice = invoice
@pdf = pdf @pdf = pdf
@mailer_action = mailer_action.to_sym
raise Error, "Azione email non valida" unless MAILER_ACTIONS.include?(@mailer_action)
end end
def call def call
@@ -22,8 +26,10 @@ module Billing
@invoice.update!(status: "issued") @invoice.update!(status: "issued")
Billing::InvoiceMailer.with(invoice: @invoice).invoice_pdf.deliver_now mail = Billing::InvoiceMailer.with(invoice: @invoice).public_send(@mailer_action)
@invoice.update!(status: "sent", emailed_at: Time.current) if MatchLiveTv.deliver_mail(mail)
@invoice.update!(status: "sent", emailed_at: Time.current)
end
@invoice @invoice
end end
@@ -5,9 +5,13 @@ module Billing
AMOUNT_HINTS = { AMOUNT_HINTS = {
500 => %w[premium_light monthly], 500 => %w[premium_light monthly],
790 => %w[premium_light monthly],
4000 => %w[premium_light yearly], 4000 => %w[premium_light yearly],
5900 => %w[premium_light yearly],
2000 => %w[premium_full monthly], 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 }.freeze
class << self 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:) def button_label(plan:, interval:, kind:)
price = PriceCatalog.label(plan_slug: plan.slug, interval: interval) price = PriceCatalog.label(plan_slug: plan.slug, interval: interval)
name = I18n.t("pages.plans.cta_name.#{plan.slug}", default: plan.name)
if kind == :checkout 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 else
I18n.t("billing.messages.switch_button", plan: plan.name, price: price) I18n.t("billing.messages.switch_button", plan: name, price: price)
end end
end end
@@ -4,9 +4,16 @@ module Billing
INTERVALS = %w[monthly yearly].freeze INTERVALS = %w[monthly yearly].freeze
DEFAULT_INTERVAL = "yearly" DEFAULT_INTERVAL = "yearly"
PRICING = { # Importi in centesimi (EUR). Il yearly è il prezzo lancio addebitato.
"premium_light" => { "yearly" => "€40/anno", "monthly" => "€5/mese" }, AMOUNTS = {
"premium_full" => { "yearly" => "€200/anno", "monthly" => "€20/mese" } "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 }.freeze
class << self class << self
@@ -30,7 +37,23 @@ module Billing
end end
def label(plan_slug:, interval:) 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 end
def available_intervals(plan_slug:) def available_intervals(plan_slug:)
@@ -52,8 +75,45 @@ module Billing
nil nil
end 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 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) def price_id_for(plan_slug, interval)
case [plan_slug, interval] case [plan_slug, interval]
when %w[premium_light monthly] 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
+106 -8
View File
@@ -4,7 +4,15 @@ module Mediamtx
class Client class Client
class Error < StandardError; end class Error < StandardError; end
CREATE_PATH_RETRIES = -> { ENV.fetch("MEDIAMTX_CREATE_RETRIES", "5").to_i }
CREATE_PATH_RETRY_BASE_SECS = -> { ENV.fetch("MEDIAMTX_CREATE_RETRY_BASE_SECS", "0.4").to_f }
def self.for_session(session)
new(base_url: session.mediamtx_api_base_url)
end
def initialize(base_url: MatchLiveTv.mediamtx_api_url) def initialize(base_url: MatchLiveTv.mediamtx_api_url)
@base_url = base_url
@conn = Faraday.new(url: base_url) do |f| @conn = Faraday.new(url: base_url) do |f|
f.request :json f.request :json
f.response :json f.response :json
@@ -12,19 +20,35 @@ module Mediamtx
end end
end end
attr_reader :base_url
# Health probe for CPX readiness (GET /v3/paths/list).
def reachable?(timeout: 2)
conn = Faraday.new(url: @base_url) do |f|
f.adapter Faraday.default_adapter
f.options.open_timeout = timeout
f.options.timeout = timeout
end
response = conn.get("/v3/paths/list")
response.success?
rescue Faraday::Error
false
end
def create_path(session) def create_path(session)
path = session.mediamtx_path_name path = session.mediamtx_path_name
# record: false finché non c'è publisher — con alwaysAvailable MediaMTX registrerebbe # record: false finché non c'è publisher — con alwaysAvailable MediaMTX registrerebbe
# solo la slate (schermo nero) in pausa/attesa. # solo la slate in pausa/attesa. Slate resta accesa anche su YouTube (copertina se l'app cade).
# YouTube: niente slate sul path camera (maschera il video al relay ffmpeg).
body = recording_body(session, enabled: false).merge( body = recording_body(session, enabled: false).merge(
source: "publisher", source: "publisher",
overridePublisher: true overridePublisher: true
) )
body[:alwaysAvailable] = true body[:alwaysAvailable] = true
body[:alwaysAvailableFile] = slate_file_path body[:alwaysAvailableFile] = slate_file_path(session)
# YouTube: telefono → MediaMTX; relay copy verso RTMPS in sidekiq. # YouTube: telefono → MediaMTX; relay copy verso RTMPS in sidekiq.
response = @conn.post("/v3/config/paths/add/#{CGI.escape(path)}", body) response = with_connection_retries("create_path #{path}") do
@conn.post("/v3/config/paths/add/#{CGI.escape(path)}", body)
end
unless response.success? unless response.success?
err = response.body.is_a?(Hash) ? response.body["error"] : response.body err = response.body.is_a?(Hash) ? response.body["error"] : response.body
raise Error, "MediaMTX path create failed: #{response.status} #{err}" raise Error, "MediaMTX path create failed: #{response.status} #{err}"
@@ -68,14 +92,55 @@ module Mediamtx
set_always_available(session, enabled: true) set_always_available(session, enabled: true)
end end
# Slate alwaysAvailable: copertina sullo stesso path quando il telefono è offline. # Garantisce path config + slate: se MediaMTX ha perso la config (restart → all_others),
# Disattivare quando il publisher è in onda; riattivare in pausa/disconnessione. # 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:) def set_always_available(session, enabled:)
path = session.mediamtx_path_name path = session.mediamtx_path_name
return true if always_available_remembered?(path, enabled: enabled) return true if always_available_remembered?(path, enabled: enabled)
body = if enabled body = if enabled
{ alwaysAvailable: true, alwaysAvailableFile: slate_file_path } { alwaysAvailable: true, alwaysAvailableFile: slate_file_path(session) }
else else
{ alwaysAvailable: false } { alwaysAvailable: false }
end end
@@ -98,12 +163,39 @@ module Mediamtx
[] []
end 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 def online_path_names
Set.new(list_paths.filter_map { |item| item["name"] if item["online"] }) Set.new(list_paths.filter_map { |item| item["name"] if item["online"] })
end end
private private
def with_connection_retries(label)
attempts = [CREATE_PATH_RETRIES.call, 1].max
base = CREATE_PATH_RETRY_BASE_SECS.call
try = 0
begin
try += 1
yield
rescue Faraday::ConnectionFailed, Faraday::TimeoutError => e
raise if try >= attempts
sleep_secs = base * (2**(try - 1))
Rails.logger.warn("[Mediamtx::Client] #{label} retry #{try}/#{attempts} after #{e.class}: #{e.message} (sleep #{sleep_secs}s)")
sleep(sleep_secs)
retry
end
end
def recording_body(session, enabled:) def recording_body(session, enabled:)
ent = session.match.team.entitlements ent = session.match.team.entitlements
can_record = ent.recording_enabled_for_mediamtx? can_record = ent.recording_enabled_for_mediamtx?
@@ -157,7 +249,13 @@ module Mediamtx
SCRIPT SCRIPT
end 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") ENV.fetch("MEDIAMTX_SLATE_FILE", "/slates/offline.mp4")
end end
@@ -4,17 +4,37 @@ module Mediamtx
module_function module_function
def active?(session) def active?(session)
return true if rtmp_publisher?(session)
active_path?(path_info(session)) active_path?(path_info(session))
end end
def path_info(session) 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 end
# MediaMTX <=1.19: online + source.type=rtmpConn.
# MediaMTX 1.20+: online/source spesso null anche con publisher; usare rtmpconns.
def active_path?(info) def active_path?(info)
return false unless 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 end
def h264_video?(info) def h264_video?(info)
@@ -26,12 +46,14 @@ module Mediamtx
end end
def video_publishing?(session) def video_publishing?(session)
info = path_info(session) return false unless active?(session)
return false unless active_path?(info)
# Slate alwaysAvailable ha H264 ma non è il telefono.
return false unless info.dig("source", "type") == "rtmpConn"
info = path_info(session)
h264_video?(info) h264_video?(info)
end end
def rtmp_source?(type)
type.to_s.match?(/\Artmps?Conn\z/)
end
end end
end end
+13 -20
View File
@@ -12,14 +12,13 @@ module Mediamtx
return @session if @session.terminal? return @session if @session.terminal?
path_info = Mediamtx::PublisherOnline.path_info(@session) path_info = Mediamtx::PublisherOnline.path_info(@session)
publisher_online = Mediamtx::PublisherOnline.active_path?(path_info) publisher_online = Mediamtx::PublisherOnline.active?(@session)
if publisher_online if publisher_online
clear_publisher_misses!(@session.id) clear_publisher_misses!(@session.id)
if @session.paused? if @session.paused?
# RTMP ancora connesso in pausa: non forzare live/reconnect. # RTMP ancora connesso in pausa: non forzare live/reconnect.
else else
enable_live_path_once!(@session)
if @session.may_go_live? if @session.may_go_live?
@session.go_live! @session.go_live!
@session.reload @session.reload
@@ -75,27 +74,13 @@ module Mediamtx
@session.update!(timeout_job_id: job) @session.update!(timeout_job_id: job)
end end
# Compatibilità con webhook / controller.
def self.schedule_youtube_pipeline!(session, force: false) def self.schedule_youtube_pipeline!(session, force: false)
Youtube::LivePipeline.schedule!(session, force: force) Youtube::LivePipeline.schedule!(session, force: force)
end 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) def restore_slate_path!(session)
return if session.platform == "matchlivetv" # Anche su matchlivetv: dopo restart MediaMTX la path può cadere su all_others senza slate.
Client.for_session(session).ensure_path_with_slate!(session)
Client.new.set_always_available(session, enabled: true)
rescue Client::Error => e rescue Client::Error => e
Rails.logger.warn("[PublisherSync] enable slate session=#{session.id}: #{e.message}") Rails.logger.warn("[PublisherSync] enable slate session=#{session.id}: #{e.message}")
end end
@@ -128,7 +113,7 @@ module Mediamtx
return return
end 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) redis.set(key, desired, ex: 48.hours.to_i)
mark_recording_patch!(session.id) mark_recording_patch!(session.id)
rescue Client::Error => e rescue Client::Error => e
@@ -160,6 +145,14 @@ module Mediamtx
nil nil
end 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) def recording_state_key(session_id)
"mediamtx:recording:#{session_id}" "mediamtx:recording:#{session_id}"
end end
@@ -183,7 +176,7 @@ module Mediamtx
key = format("youtube_relay:sched:%s", session.id) key = format("youtube_relay:sched:%s", session.id)
return unless redis.set(key, "1", nx: true, ex: 10) return unless redis.set(key, "1", nx: true, ex: 10)
YoutubeRelayEnsureJob.perform_later(session.id) Streams::YoutubeRelay.enqueue_ensure!(session)
end end
def redis def redis
+56 -1
View File
@@ -44,7 +44,8 @@ module Ops
check_sidekiq_heartbeat, check_sidekiq_heartbeat,
check_sidekiq_dead, check_sidekiq_dead,
check_http_rails, check_http_rails,
check_rails_latency check_rails_latency,
check_stream_overflow
] ]
findings << check_http_public if public_check_due? findings << check_http_public if public_check_due?
findings findings
@@ -53,6 +54,8 @@ module Ops
def process_finding(finding) def process_finding(finding)
if finding.healthy if finding.healthy
Ops::IncidentRecorder.resolve(fingerprint: finding.fingerprint) Ops::IncidentRecorder.resolve(fingerprint: finding.fingerprint)
# overflow usa fingerprint diverse (at_max / orphan_idle / budget) rispetto al check sano
Ops::IncidentRecorder.resolve_kind(finding.kind) if finding.kind == "stream_overflow"
else else
Ops::IncidentRecorder.record( Ops::IncidentRecorder.record(
kind: finding.kind, kind: finding.kind,
@@ -161,6 +164,58 @@ module Ops
fail_finding("garage_storage", "warning", "garage_storage:head", "Garage storage non raggiungibile", e.message) fail_finding("garage_storage", "warning", "garage_storage:head", "Garage storage non raggiungibile", e.message)
end 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 def check_sidekiq_heartbeat
redis = Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0")) redis = Redis.new(url: ENV.fetch("REDIS_URL", "redis://localhost:6379/0"))
last = redis.get(Ops::HealthMonitorJob::HEARTBEAT_KEY).to_i last = redis.get(Ops::HealthMonitorJob::HEARTBEAT_KEY).to_i
@@ -10,6 +10,10 @@ module Ops
def resolve(fingerprint:) def resolve(fingerprint:)
new.resolve(fingerprint: fingerprint) new.resolve(fingerprint: fingerprint)
end end
def resolve_kind(kind)
new.resolve_kind(kind)
end
end end
def record(finding) def record(finding)
@@ -49,6 +53,10 @@ module Ops
Ops::Incident.open.where(fingerprint: fingerprint).find_each(&:resolve!) Ops::Incident.open.where(fingerprint: fingerprint).find_each(&:resolve!)
end end
def resolve_kind(kind)
Ops::Incident.open.where(kind: kind).find_each(&:resolve!)
end
private private
def fingerprint_for(finding) def fingerprint_for(finding)
@@ -95,7 +95,7 @@ module Recordings
def cleanup_mediamtx_path(session) def cleanup_mediamtx_path(session)
return unless session.status.in?(%w[ended error]) 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 rescue Mediamtx::Client::Error => e
@logger.warn("[Recordings::CleanupLocal] delete_path #{session.id}: #{e.message}") @logger.warn("[Recordings::CleanupLocal] delete_path #{session.id}: #{e.message}")
end end
@@ -0,0 +1,72 @@
# frozen_string_literal: true
module Recordings
# Rimuove solo l'MP4 temporaneo su object storage. Non soft-delete, non tocca YouTube.
class ClearTemporaryMedia
def initialize(recording, reason: :verified, force: false)
@recording = recording
@reason = reason
@force = force
end
def call
unless @recording.temporary_storage?
Rails.logger.info(
"[Recordings::ClearTemporaryMedia] skip recording=#{@recording.id} not_temporary"
)
return @recording
end
if @recording.local_media_purged?
Rails.logger.info(
"[Recordings::ClearTemporaryMedia] already_purged recording=#{@recording.id}"
)
return @recording
end
unless allowed?
Rails.logger.info(
"[Recordings::ClearTemporaryMedia] skip recording=#{@recording.id} " \
"reason=not_verified_and_not_expired force=#{@force}"
)
return @recording
end
delete_video_object!
@recording.update!(
storage_key: nil,
byte_size: nil,
local_media_purged_at: Time.current
)
Rails.logger.info(
"[Recordings::ClearTemporaryMedia] purged recording=#{@recording.id} " \
"reason=#{@reason} youtube_verified=#{@recording.youtube_verified_at.present?} " \
"youtube_video_id=#{@recording.youtube_video_id.inspect}"
)
@recording
end
private
def allowed?
return true if @force
return true if @recording.youtube_verified_at.present?
return true if @recording.temp_expires_at.present? && @recording.temp_expires_at <= Time.current
false
end
def delete_video_object!
key = @recording.storage_key
return if key.blank?
Recordings::Storage.new.delete(key: key)
rescue Recordings::Storage::Error => e
# File già assente → ok (idempotente)
Rails.logger.info(
"[Recordings::ClearTemporaryMedia] storage delete recording=#{@recording.id}: #{e.message}"
)
end
end
end
@@ -1,3 +1,5 @@
# frozen_string_literal: true
module Recordings module Recordings
class FinalizeSession class FinalizeSession
def initialize(session) def initialize(session)
@@ -5,29 +7,69 @@ module Recordings
end end
def call def call
team = @session.match.team policy = Recordings::StoragePolicy.call(@session)
return unless team.entitlements.can_create_recordings? return if policy == Recordings::StoragePolicy::NONE
retention_days = team.entitlements.recording_retention_days team = @session.match.team
expires_at = retention_days.positive? ? retention_days.days.from_now : nil attrs = attributes_for(policy, team)
recording = Recording.find_or_initialize_by(stream_session: @session) recording = Recording.find_or_initialize_by(stream_session: @session)
recording.assign_attributes( if skip_reinitialize?(recording)
Rails.logger.info(
"[Recordings::FinalizeSession] skip already-finalized " \
"recording=#{recording.id} status=#{recording.status}"
)
return recording
end
recording.assign_attributes(attrs)
recording.save!
Rails.logger.info(
"[Recordings::FinalizeSession] session=#{@session.id} recording=#{recording.id} " \
"storage_policy=#{policy} expires_at=#{recording.expires_at.inspect} " \
"temp_expires_at=#{recording.temp_expires_at.inspect}"
)
recording
end
private
def skip_reinitialize?(recording)
return false unless recording.persisted?
recording.ready? ||
recording.status == "processing" ||
recording.storage_key.present?
end
def attributes_for(policy, team)
{
team: team, team: team,
status: "processing", status: "processing",
title: default_title, title: default_title,
privacy_status: privacy_from_session, privacy_status: privacy_from_session,
storage_path: @session.mediamtx_path_name, storage_path: @session.mediamtx_path_name,
recorded_at: @session.ended_at || Time.current, recorded_at: @session.ended_at || Time.current,
expires_at: expires_at, storage_policy: policy,
expires_at: archive_expires_at(policy, team),
temp_expires_at: temp_expires_at_for(policy),
error_message: nil, error_message: nil,
metadata: initial_metadata(team) metadata: initial_metadata(policy, team)
) }
recording.save!
recording
end end
private def archive_expires_at(policy, team)
return nil if policy == Recordings::StoragePolicy::TEMPORARY
retention_days = team.entitlements.recording_retention_days
retention_days.positive? ? retention_days.days.from_now : nil
end
def temp_expires_at_for(policy)
return nil unless policy == Recordings::StoragePolicy::TEMPORARY
MatchLiveTv.youtube_temp_replay_retention_hours.hours.from_now
end
def default_title def default_title
match = @session.match match = @session.match
@@ -38,12 +80,13 @@ module Recordings
@session.privacy_status == "public" ? "public" : "unlisted" @session.privacy_status == "public" ? "public" : "unlisted"
end end
def initial_metadata(team) def initial_metadata(policy, team)
ent = team.entitlements
{ {
"source_platform" => @session.platform, "source_platform" => @session.platform,
"session_privacy" => @session.privacy_status, "session_privacy" => @session.privacy_status,
"auto_publish_youtube" => ent.premium_full? && ent.youtube_enabled? && @session.platform == "youtube", # Re-upload automatico disabilitato: le live YouTube usano VerifyYoutubeReplay.
"auto_publish_youtube" => false,
"storage_policy" => policy,
"ai" => {} "ai" => {}
} }
end end
@@ -23,7 +23,7 @@ module Recordings
private private
def deliver_expiring_soon(user) def deliver_expiring_soon(user)
Recordings::ReplayMailer.replay_expiring_soon(recording: @recording, recipient: user).deliver_now MatchLiveTv.deliver_mail(Recordings::ReplayMailer.replay_expiring_soon(recording: @recording, recipient: user))
rescue EOFError => e rescue EOFError => e
Rails.logger.warn("[Recordings::NotifyExpiring] SMTP EOF on close for #{user.email}: #{e.message}") Rails.logger.warn("[Recordings::NotifyExpiring] SMTP EOF on close for #{user.email}: #{e.message}")
end end

Some files were not shown because too many files have changed in this diff Show More