Compare commits

..
Author SHA1 Message Date
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
67 changed files with 2251 additions and 219 deletions
@@ -79,14 +79,21 @@ module Admin
if @filters[:status].present? && StreamSession::STATUSES.include?(@filters[:status])
scope = scope.where(stream_sessions: { status: @filters[:status] })
end
if @filters[:platform].present? && StreamSession::PLATFORMS.include?(@filters[:platform])
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?
scope = scope.where(stream_sessions: { stream_node_id: @filters[:stream_node_id] })
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(
@@ -97,10 +97,11 @@ module Api
replay_url: recording.replay_url,
playback_url: recording.playback_stream_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_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_label: recording.source_platform_label,
expires_at: recording.expires_at,
@@ -77,6 +77,7 @@ module Api
thermal_state: sanitized_thermal_state,
last_seen_at: Time.current
)
Sessions::ApplyClientInfo.call(@session, params[:client]) if params[:client].present?
sync_publisher_when_streaming!(params[:fps].to_f)
SessionChannel.broadcast_message(@session, state.as_cable_payload)
head :no_content
@@ -180,7 +181,12 @@ module Api
end
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
def score_sync_params
@@ -171,12 +171,13 @@ module Api
replay_url: recording.replay_url,
playback_url: recording.playback_stream_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,
views_label: recording.views_label,
youtube_video_id: recording.youtube_video_id,
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_label: recording.source_platform_label,
expires_at: recording.expires_at,
+29 -4
View File
@@ -30,8 +30,9 @@ module AdminHelper
links
end
def admin_session_ingest_badge_class(node)
case node.role
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"
@@ -39,8 +40,9 @@ module AdminHelper
end
end
def admin_session_ingest_role_label(node)
I18n.t("admin.sessions.ingest.role.#{node.role}", default: node.role.to_s.humanize)
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)
@@ -81,6 +83,29 @@ module AdminHelper
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?
@@ -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
class PostProcessJob
include Sidekiq::Job
@@ -7,10 +9,30 @@ module Recordings
recording = Recording.find_by(id: recording_id)
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 if recording.youtube_video_id.present?
return if recording.source_platform == "youtube"
Recordings::PublishToYoutubeJob.perform_async(recording.id)
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,47 @@
# 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?
Recordings::ClearTemporaryMediaJob.perform_async(recording.id, "verified")
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
@@ -8,6 +8,8 @@ module Streams
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
@@ -17,6 +19,16 @@ module Streams
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
+51 -4
View File
@@ -2,6 +2,8 @@ class Recording < ApplicationRecord
STATUSES = %w[processing ready expired failed].freeze
PRIVACY_STATUSES = %w[public unlisted].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 :team
@@ -9,6 +11,7 @@ class Recording < ApplicationRecord
validates :status, inclusion: { in: STATUSES }
validates :privacy_status, inclusion: { in: PRIVACY_STATUSES }
validates :storage_backend, inclusion: { in: STORAGE_BACKENDS }
validates :storage_policy, inclusion: { in: STORAGE_POLICIES }
scope :not_deleted, -> { where(deleted_at: nil) }
scope :ready, lambda {
@@ -23,7 +26,16 @@ class Recording < ApplicationRecord
ready.where(expires_at: ..days.days.from_now)
}
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|
q = query.to_s.strip
@@ -46,19 +58,30 @@ class Recording < ApplicationRecord
end
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"
end
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?
"#{MatchLiveTv.app_public_url.chomp('/')}/replay/#{stream_session_id}/thumbnail"
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
return nil unless ready?
return nil unless ready? && storage_key.present?
"/api/v1/recordings/#{id}/download"
end
@@ -70,6 +93,30 @@ class Recording < ApplicationRecord
"https://www.youtube.com/watch?v=#{youtube_video_id}"
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?
status == "ready" && !deleted? && (expires_at.nil? || expires_at.future?)
end
@@ -154,7 +201,7 @@ class Recording < ApplicationRecord
end
def playable_on_site?
ready? && storage_key.present?
available_in_archive?
end
def days_until_expiry
+16 -1
View File
@@ -1,7 +1,8 @@
class StreamSession < ApplicationRecord
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
PRIVACY_STATUSES = %w[public unlisted private].freeze
@@ -20,6 +21,7 @@ class StreamSession < ApplicationRecord
before_validation :normalize_privacy_status
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 :publicly_listed, -> { where(privacy_status: "public") }
@@ -74,6 +76,14 @@ class StreamSession < ApplicationRecord
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
# RootEncoder richiede rtmp://host:port/app/stream (due segmenti).
# MediaMTX path = live/match_{uuid} (no ?token= nel path).
@@ -208,6 +218,11 @@ class StreamSession < ApplicationRecord
self.privacy_status = "unlisted" if privacy_status == "private"
end
def snapshot_ingest_from_node
self.ingest_slug = stream_node.slug
self.ingest_role = stream_node.role
end
def record_ended_timestamps!
now = Time.current
update!(ended_at: now) if ended_at.nil?
@@ -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
class FinalizeSession
def initialize(session)
@@ -5,29 +7,69 @@ module Recordings
end
def call
team = @session.match.team
return unless team.entitlements.can_create_recordings?
policy = Recordings::StoragePolicy.call(@session)
return if policy == Recordings::StoragePolicy::NONE
retention_days = team.entitlements.recording_retention_days
expires_at = retention_days.positive? ? retention_days.days.from_now : nil
team = @session.match.team
attrs = attributes_for(policy, team)
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,
status: "processing",
title: default_title,
privacy_status: privacy_from_session,
storage_path: @session.mediamtx_path_name,
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,
metadata: initial_metadata(team)
)
recording.save!
recording
metadata: initial_metadata(policy, team)
}
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
match = @session.match
@@ -38,12 +80,13 @@ module Recordings
@session.privacy_status == "public" ? "public" : "unlisted"
end
def initial_metadata(team)
ent = team.entitlements
def initial_metadata(policy, team)
{
"source_platform" => @session.platform,
"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" => {}
}
end
@@ -0,0 +1,92 @@
# frozen_string_literal: true
require "cgi"
require "fileutils"
require "net/http"
require "uri"
module Recordings
# Recupera i segmenti MediaMTX dal disco del CPX (agent :9100) verso un tmpdir locale.
class PullFromCloudNode
class Error < StandardError; end
def initialize(session)
@session = session
end
def applicable?
node = @session.stream_node
node.present? && node.role == "cloud" && agent_url.present?
end
# @return [String, nil] directory con i file, o nil se il nodo non è cloud / 404
def fetch
return unless applicable?
dest = Dir.mktmpdir("mltv-cpx-rec-")
uri = recordings_uri
http = Net::HTTP.new(uri.host, uri.port)
http.open_timeout = 5
http.read_timeout = ENV.fetch("STREAM_NODE_RECORDINGS_PULL_TIMEOUT", "180").to_i
req = Net::HTTP::Get.new(uri)
req["Authorization"] = "Bearer #{agent_secret}" if agent_secret.present?
res = http.request(req)
if res.is_a?(Net::HTTPNotFound)
FileUtils.remove_entry(dest)
return nil
end
unless res.is_a?(Net::HTTPSuccess) && res.body.present?
FileUtils.remove_entry(dest)
raise Error, "agent GET recordings HTTP #{res.code} #{res.body.to_s.truncate(200)}"
end
tar_path = File.join(dest, "recordings.tar.gz")
File.binwrite(tar_path, res.body)
unpack!(tar_path, dest)
FileUtils.rm_f(tar_path)
dest
rescue StandardError
FileUtils.remove_entry(dest) if dest && Dir.exist?(dest)
raise
end
def cleanup_remote!
return unless applicable?
uri = recordings_uri
http = Net::HTTP.new(uri.host, uri.port)
http.open_timeout = 5
http.read_timeout = 15
req = Net::HTTP::Delete.new(uri)
req["Authorization"] = "Bearer #{agent_secret}" if agent_secret.present?
res = http.request(req)
return if res.is_a?(Net::HTTPSuccess) || res.is_a?(Net::HTTPNotFound)
Rails.logger.warn(
"[Recordings::PullFromCloudNode] delete HTTP #{res.code} session=#{@session.id}"
)
rescue StandardError => e
Rails.logger.warn("[Recordings::PullFromCloudNode] delete #{e.class}: #{e.message}")
end
private
def recordings_uri
path = @session.mediamtx_path_name.to_s
URI.parse("#{agent_url.chomp('/')}/recordings/#{CGI.escape(path)}")
end
def agent_url
ENV["STREAM_NODE_RELAY_AGENT_URL"].presence || @session.stream_node&.relay_agent_url
end
def agent_secret
ENV["STREAM_NODE_AGENT_SECRET"].presence || "mediamtx_webhook_dev_secret"
end
def unpack!(tar_path, dest)
ok = system("tar", "-xzf", tar_path, "-C", dest, out: File::NULL, err: File::NULL)
raise Error, "tar extract failed" unless ok
end
end
end
@@ -0,0 +1,36 @@
# frozen_string_literal: true
module Recordings
# Decisione centralizzata: temporary (YouTube) | retained (HLS) | none.
class StoragePolicy
TEMPORARY = "temporary"
RETAINED = "retained"
NONE = "none"
POLICIES = [TEMPORARY, RETAINED, NONE].freeze
def self.call(session)
new(session).call
end
def initialize(session)
@session = session
end
def call
team = @session.match.team
ent = team.entitlements
return NONE unless ent.can_create_recordings?
if youtube_destination?
TEMPORARY
else
RETAINED
end
end
def youtube_destination?
@session.platform.to_s == "youtube"
end
end
end
@@ -9,8 +9,9 @@ module Recordings
def call
recording = Recording.find_by(stream_session: @session)
return unless recording&.status == "processing"
return if recording.storage_key.present?
source_files = local_source_files
source_files = collect_source_files
if source_files.empty?
fail_recording!(recording, "Nessun file di registrazione trovato")
cleanup_mediamtx_path
@@ -35,23 +36,52 @@ module Recordings
error_message: nil
)
Rails.logger.info(
"[Recordings::UploadFromSession] ready recording=#{recording.id} " \
"storage_policy=#{recording.storage_policy} key=#{storage_key} bytes=#{byte_size}"
)
cleanup_local_sources(source_files, merged_path)
cleanup_cloud_pull!
cleanup_remote_recordings!
cleanup_mediamtx_path
Recordings::PostProcessJob.perform_async(recording.id)
recording
rescue Error, Recordings::Storage::Error => e
rescue Error, Recordings::Storage::Error, Recordings::PullFromCloudNode::Error => e
recording = Recording.find_by(stream_session: @session)
fail_recording!(recording, e.message) if recording
raise
ensure
FileUtils.rm_f(@merged_temp_path) if @merged_temp_path && File.exist?(@merged_temp_path)
cleanup_cloud_pull! unless recording_ready?
end
private
def collect_source_files
files = local_source_files
return files if files.any?
puller = Recordings::PullFromCloudNode.new(@session)
return [] unless puller.applicable?
3.times do |i|
cleanup_cloud_pull!
@cloud_pull_dir = puller.fetch
files = scan_recording_files(@cloud_pull_dir)
return files if files.any?
sleep 2 if i < 2
end
[]
end
def local_source_files
base = File.join(MatchLiveTv.recordings_local_path, @session.mediamtx_path_name)
return [] unless Dir.exist?(base)
scan_recording_files(File.join(MatchLiveTv.recordings_local_path, @session.mediamtx_path_name))
end
def scan_recording_files(base)
return [] if base.blank? || !Dir.exist?(base)
Dir.glob(File.join(base, "**", "*"))
.select { |path| File.file?(path) && path.match?(/\.(mp4|fmp4|m4s|ts)$/i) }
@@ -83,7 +113,11 @@ module Recordings
end
def object_key(recording)
"teams/#{recording.team_id}/sessions/#{@session.id}/replay.mp4"
if recording.temporary_storage?
"temporary_replays/teams/#{recording.team_id}/sessions/#{@session.id}/replay.mp4"
else
"teams/#{recording.team_id}/sessions/#{@session.id}/replay.mp4"
end
end
def cleanup_local_sources(source_files, merged_path)
@@ -100,6 +134,24 @@ module Recordings
Rails.logger.warn("[Recordings::UploadFromSession] delete_path: #{e.message}")
end
def cleanup_cloud_pull!
return if @cloud_pull_dir.blank? || !Dir.exist?(@cloud_pull_dir)
FileUtils.remove_entry(@cloud_pull_dir)
@cloud_pull_dir = nil
rescue StandardError
nil
end
def cleanup_remote_recordings!
Recordings::PullFromCloudNode.new(@session).cleanup_remote!
end
def recording_ready?
rec = Recording.find_by(stream_session: @session)
rec&.status == "ready"
end
def fail_recording!(recording, message)
recording.update!(status: "failed", error_message: message)
cleanup_mediamtx_path
@@ -0,0 +1,86 @@
# frozen_string_literal: true
module Recordings
# Collega la live YouTube (broadcast_id) al VOD e aggiorna metadata recording.
class VerifyYoutubeReplay
class Error < StandardError; end
Result = Struct.new(:status, :message, keyword_init: true) do
def ok?
status == :ok
end
def retriable?
status == :pending
end
end
def initialize(recording)
@recording = recording
end
def call
unless @recording.temporary_storage?
return Result.new(status: :skipped, message: "not_temporary")
end
if @recording.youtube_verified_at.present? && @recording.youtube_video_id.present?
return Result.new(status: :ok, message: "already_verified")
end
session = @recording.stream_session
broadcast_id = session&.youtube_broadcast_id
if broadcast_id.blank?
Rails.logger.warn(
"[Recordings::VerifyYoutubeReplay] missing broadcast_id recording=#{@recording.id}"
)
return Result.new(status: :failed, message: "missing_broadcast_id")
end
info = Youtube::VodStatus.new(@recording.team, channel: "team").fetch(broadcast_id)
unless info.ready
Rails.logger.info(
"[Recordings::VerifyYoutubeReplay] not_ready recording=#{@recording.id} " \
"video=#{broadcast_id} upload_status=#{info.upload_status.inspect}"
)
return Result.new(status: :pending, message: "vod_not_ready")
end
apply_verified!(info)
Rails.logger.info(
"[Recordings::VerifyYoutubeReplay] ok recording=#{@recording.id} " \
"youtube_video_id=#{info.video_id}"
)
Result.new(status: :ok, message: "verified")
rescue Youtube::VodStatus::Error => e
Rails.logger.warn(
"[Recordings::VerifyYoutubeReplay] api_error recording=#{@recording.id}: #{e.message}"
)
Result.new(status: :pending, message: e.message)
end
private
def apply_verified!(info)
meta = @recording.metadata.is_a?(Hash) ? @recording.metadata.deep_dup : {}
yt = meta.fetch("youtube", {}).merge(
"broadcast_id" => @recording.stream_session.youtube_broadcast_id,
"thumbnail_url" => info.thumbnail_url,
"privacy_status" => info.privacy_status,
"upload_status" => info.upload_status,
"verified_via" => "live_broadcast"
).compact
attrs = {
youtube_video_id: info.video_id,
youtube_verified_at: Time.current,
youtube_published_at: @recording.youtube_published_at || Time.current,
metadata: meta.merge("youtube" => yt)
}
attrs[:duration_secs] = info.duration_secs if info.duration_secs.to_i.positive?
attrs[:title] = info.title if info.title.present? && @recording.title.blank?
@recording.update!(attrs)
end
end
end
@@ -0,0 +1,72 @@
# frozen_string_literal: true
module Sessions
# Normalizza e applica fingerprint del client (OS, app, device, operatore)
# sulla sessione, a create e/o a ogni telemetry.
class ApplyClientInfo
OS_VALUES = %w[android ios].freeze
MAX_LEN = 80
ATTRS = %i[
client_os app_version app_build device_manufacturer device_model os_version carrier
].freeze
def self.call(session, raw)
new(session, raw).call
end
def initialize(session, raw)
@session = session
@raw = normalize_hash(raw)
end
def call
attrs = extract_attrs
return @session if attrs.empty?
@session.assign_attributes(attrs)
@session.save! if @session.persisted? && @session.changed?
@session
end
private
def normalize_hash(raw)
return {} if raw.blank?
data = raw.respond_to?(:to_unsafe_h) ? raw.to_unsafe_h : raw
data = data.to_h if data.respond_to?(:to_h)
data.with_indifferent_access
rescue StandardError
{}
end
def extract_attrs
attrs = {}
os = @raw[:os].presence || @raw[:client_os].presence
os = os.to_s.downcase.strip
attrs[:client_os] = os if OS_VALUES.include?(os)
{
app_version: %i[app_version version],
app_build: %i[app_build build build_number],
device_manufacturer: %i[device_manufacturer manufacturer],
device_model: %i[device_model model],
os_version: %i[os_version system_version],
carrier: %i[carrier network_operator operator]
}.each do |column, keys|
value = keys.map { |k| @raw[k] }.find(&:present?)
next if value.blank?
attrs[column] = truncate(value.to_s.strip)
end
attrs
end
def truncate(value)
value.bytesize <= MAX_LEN ? value : value.byteslice(0, MAX_LEN)
end
end
end
+31 -3
View File
@@ -26,6 +26,7 @@ module Sessions
target_fps: @params[:target_fps] || 30,
status: "idle"
)
Sessions::ApplyClientInfo.call(session, @params[:client])
youtube_channel = nil
if session.platform == "youtube"
@@ -46,18 +47,23 @@ module Sessions
{
created: true,
platform: session.platform,
stream_node: session.stream_node&.slug
}
stream_node: session.stream_node&.slug,
client_os: session.client_os,
app_version: session.app_version,
device_model: session.device_model
}.compact
)
end
kick_autoscaler_if_soft_limit!
if session.platform == "youtube"
YoutubeBroadcastSetupJob.perform_later(session.id, youtube_channel)
end
session
rescue Streams::NodeRegistry::NoCapacityError => e
raise Teams::EntitlementError.new(e.message, code: "stream_capacity_exhausted")
raise_no_capacity!(e)
rescue Faraday::ConnectionFailed, Faraday::TimeoutError => e
raise Streams::IngestUnavailableError, "Ingest temporaneamente non disponibile (#{e.class})"
rescue Mediamtx::Client::Error => e
@@ -69,6 +75,28 @@ module Sessions
private
def raise_no_capacity!(error)
if Streams::Autoscaler.enabled?
Streams::AutoscalerJob.kick!
raise Teams::EntitlementError.new(
"Capacità streaming in espansione. Riprova tra poco.",
code: "stream_capacity_scaling"
)
end
raise Teams::EntitlementError.new(error.message, code: "stream_capacity_exhausted")
end
def kick_autoscaler_if_soft_limit!
return unless Streams::Autoscaler.enabled?
free = Streams::Autoscaler.metrics[:free_slots].to_i
return if free > Streams::Autoscaler.soft_free_slots
Streams::AutoscalerJob.kick!
end
def assert_youtube_channel!(youtube_channel)
resolver = Youtube::CredentialResolver.new(@match.team, channel: youtube_channel)
if resolver.resolve.blank?
+1 -1
View File
@@ -12,7 +12,7 @@ module Sessions
complete_youtube_broadcast! if @session.youtube_broadcast_id.present?
recording = Recordings::FinalizeSession.new(@session).call
Recordings::UploadJob.perform_async(@session.id) if recording
Recordings::UploadJob.perform_async(@session.id) if recording&.status == "processing"
remove_mediamtx_paths!
log_event("ended")
@@ -0,0 +1,85 @@
# frozen_string_literal: true
module Youtube
# Stato VOD dopo una live (broadcast_id ≈ video_id su YouTube).
class VodStatus
class Error < StandardError; end
Result = Struct.new(
:ready,
:video_id,
:title,
:duration_secs,
:thumbnail_url,
:privacy_status,
:upload_status,
keyword_init: true
)
def initialize(team, channel: "team")
@team = team
@channel = channel
end
def fetch(video_id)
raise Error, "video_id mancante" if video_id.blank?
if mock_or_unconfigured?(video_id)
return Result.new(
ready: true,
video_id: video_id,
title: nil,
duration_secs: nil,
thumbnail_url: nil,
privacy_status: "unlisted",
upload_status: "processed"
)
end
client = authorized_client
item = client.list_videos("snippet,contentDetails,status", id: video_id).items&.first
return Result.new(ready: false, video_id: video_id) if item.blank?
upload_status = item.status&.upload_status.to_s
ready = upload_status.in?(%w[processed uploaded]) ||
(item.snippet.present? && upload_status != "deleted" && upload_status != "rejected" && upload_status != "failed")
Result.new(
ready: ready,
video_id: item.id,
title: item.snippet&.title,
duration_secs: parse_duration(item.content_details&.duration),
thumbnail_url: item.snippet&.thumbnails&.high&.url || item.snippet&.thumbnails&.default&.url,
privacy_status: item.status&.privacy_status,
upload_status: upload_status
)
rescue Google::Apis::Error => e
raise Error, e.message
end
private
def mock_or_unconfigured?(video_id)
video_id.to_s.start_with?("mock_") ||
ENV["YOUTUBE_CLIENT_ID"].blank? ||
CredentialResolver.new(@team, channel: @channel).resolve.blank?
end
def authorized_client
credential = CredentialResolver.new(@team, channel: @channel).resolve
raise Error, "Credenziali YouTube non disponibili" if credential.blank?
OauthRefresh.new(credential).apply!(Google::Apis::YoutubeV3::YouTubeService.new)
end
def parse_duration(iso)
return nil if iso.blank?
# PT1H2M3S
match = iso.match(/\APT(?:(\d+)H)?(?:(\d+)M)?(?:(\d+)S)?\z/)
return nil unless match
match[1].to_i * 3600 + match[2].to_i * 60 + match[3].to_i
end
end
end
@@ -107,6 +107,7 @@
<tr>
<th><%= t("admin.dashboard.sessions.table.match") %></th>
<th><%= t("admin.dashboard.sessions.table.status") %></th>
<th><%= t("admin.dashboard.sessions.table.client") %></th>
<th><%= t("admin.dashboard.sessions.table.ingest") %></th>
<th><%= t("admin.dashboard.sessions.table.start") %></th>
<th><%= t("admin.dashboard.sessions.table.link") %></th>
@@ -118,6 +119,7 @@
<tr>
<td><%= s.match.team.name %> vs <%= s.match.opponent_name %></td>
<td><span class="badge badge--<%= s.status == 'live' ? 'live' : (s.status == 'paused' ? 'paused' : 'connecting') %>"><%= s.status %></span></td>
<td class="muted"><%= admin_session_client_summary(s) %></td>
<td><%= render "admin/sessions/ingest_cell", session: s %></td>
<td class="muted"><%= s.started_at&.strftime("%d/%m %H:%M") || t("admin.common.dash") %></td>
<td>
@@ -1,8 +1,11 @@
<% node = session.stream_node %>
<% if node %>
<% slug = session.ingest_slug_display %>
<% role = session.ingest_role_display %>
<% if slug.present? %>
<div class="admin-ingest">
<code class="admin-ingest__slug"><%= node.slug %></code>
<span class="badge <%= admin_session_ingest_badge_class(node) %>"><%= admin_session_ingest_role_label(node) %></span>
<code class="admin-ingest__slug"><%= slug %></code>
<% if role.present? %>
<span class="badge <%= admin_session_ingest_badge_class(role) %>"><%= admin_session_ingest_role_label(role) %></span>
<% end %>
</div>
<% else %>
<span class="muted"><%= t("admin.sessions.ingest.none") %></span>
@@ -26,7 +26,7 @@
<span><%= t("admin.sessions.index.filters.platform") %></span>
<%= select_tag :platform,
options_for_select(
[[t("admin.sessions.index.filters.any"), ""]] + StreamSession::PLATFORMS.map { |p| [p, p] },
[[t("admin.sessions.index.filters.any"), ""]] + StreamSession::LIVE_PLATFORMS.map { |p| [p, p] },
@filters[:platform]
) %>
</label>
@@ -83,6 +83,7 @@
<th><%= t("admin.sessions.index.table.ended") %></th>
<th><%= t("admin.sessions.index.table.duration") %></th>
<th><%= t("admin.sessions.index.table.ingest") %></th>
<th><%= t("admin.sessions.index.table.client") %></th>
<th><%= t("admin.sessions.index.table.disconnects") %></th>
<th><%= t("admin.sessions.index.table.link") %></th>
<th></th>
@@ -110,6 +111,7 @@
<td class="muted"><%= s.ended_at ? admin_datetime(s.ended_at) : t("admin.common.dash") %></td>
<td class="muted"><%= admin_session_duration_label(s) %></td>
<td><%= render "admin/sessions/ingest_cell", session: s %></td>
<td class="muted admin-table-sub"><%= admin_session_client_summary(s) %></td>
<td><%= s.disconnection_count %></td>
<td>
<div class="admin-link-compact">
+46 -6
View File
@@ -67,6 +67,38 @@
<dt><%= t("admin.sessions.show.fields.platform") %></dt>
<dd><%= @session.platform %></dd>
</div>
<div>
<dt><%= t("admin.sessions.show.fields.client_os") %></dt>
<dd><%= admin_session_client_os_label(@session) %></dd>
</div>
<div>
<dt><%= t("admin.sessions.show.fields.app_version") %></dt>
<dd>
<% if @session.app_version.present? %>
<%= @session.app_version %>
<% if @session.app_build.present? %>
<span class="muted">(<%= @session.app_build %>)</span>
<% end %>
<% else %>
<%= t("admin.common.dash") %>
<% end %>
</dd>
</div>
<div>
<dt><%= t("admin.sessions.show.fields.device") %></dt>
<dd>
<% device = [@session.device_manufacturer, @session.device_model].compact_blank.join(" ") %>
<%= device.presence || t("admin.common.dash") %>
</dd>
</div>
<div>
<dt><%= t("admin.sessions.show.fields.os_version") %></dt>
<dd><%= @session.os_version.presence || t("admin.common.dash") %></dd>
</div>
<div>
<dt><%= t("admin.sessions.show.fields.carrier") %></dt>
<dd><%= @session.carrier.presence || t("admin.common.dash") %></dd>
</div>
<div>
<dt><%= t("admin.sessions.show.fields.privacy") %></dt>
<dd><%= @session.privacy_status %></dd>
@@ -130,12 +162,20 @@
<div>
<dt><%= t("admin.sessions.show.fields.node") %></dt>
<dd>
<% if @session.stream_node %>
<code><%= @session.stream_node.slug %></code>
<span class="badge <%= admin_session_ingest_badge_class(@session.stream_node) %>">
<%= admin_session_ingest_role_label(@session.stream_node) %>
</span>
<span class="muted">(<%= @session.stream_node.provider %>)</span>
<% slug = @session.ingest_slug_display %>
<% role = @session.ingest_role_display %>
<% if slug.present? %>
<code><%= slug %></code>
<% if role.present? %>
<span class="badge <%= admin_session_ingest_badge_class(role) %>">
<%= admin_session_ingest_role_label(role) %>
</span>
<% end %>
<% if @session.stream_node %>
<span class="muted">(<%= @session.stream_node.provider %>)</span>
<% else %>
<span class="muted"><%= t("admin.sessions.ingest.decommissioned") %></span>
<% end %>
<% else %>
<span class="muted"><%= t("admin.sessions.ingest.none") %></span>
<% end %>
+29 -21
View File
@@ -14,6 +14,33 @@
<h2><%= t("replay.show.processing_title") %></h2>
<p><%= t("replay.show.processing_body") %></p>
</div>
<% elsif @recording.ready? && @recording.replay_source == "youtube" && @recording.youtube_video_id.present? && !@recording.youtube_video_id.to_s.start_with?("mock_") %>
<div class="live-player-wrap live-player-wrap--embed">
<iframe src="https://www.youtube.com/embed/<%= @recording.youtube_video_id %>"
title="<%= @recording.title_or_default %>"
class="live-player-embed"
allow="accelerometer; autoplay; clipboard-write; encrypted-media; gyroscope; picture-in-picture"
allowfullscreen></iframe>
<%= render "public/live/player_overlays",
match: @match,
session: @session,
stream_closed: true,
on_air: false,
badge_label: t("replay.show.badge_label"),
badge_class: "badge-ended" %>
</div>
<p class="replay-show__meta-line">
<%= l_local(@recording.recorded_at_or_fallback) %>
· <%= t("replay.show.meta_line_duration", value: @recording.duration_label) %>
· <%= @recording.views_label %>
<% if @recording.source_platform_label != "—" %>
· <%= @recording.source_platform_label %>
<% end %>
</p>
<p class="replay-show__meta-line">
<%= t("replay.show.youtube_only_body") %>
<%= link_to t("replay.show.youtube_link"), @recording.youtube_watch_url, target: "_blank", rel: "noopener" %>
</p>
<% elsif @recording.ready? && @recording.storage_key.present? %>
<div class="live-player-wrap">
<video id="replay-player" controls playsinline preload="metadata"
@@ -43,9 +70,9 @@
</p>
<% ent = @recording.team.entitlements %>
<% if (ent.phone_download_enabled? && (logged_in? || @recording.unlisted?)) || @recording.youtube_watch_url %>
<% if (ent.phone_download_enabled? && @recording.storage_key.present? && (logged_in? || @recording.unlisted?)) || @recording.youtube_watch_url %>
<div class="replay-show__actions">
<% if ent.phone_download_enabled? && (logged_in? || @recording.unlisted?) %>
<% if ent.phone_download_enabled? && @recording.storage_key.present? && (logged_in? || @recording.unlisted?) %>
<%= link_to t("replay.show.download_link"), public_replay_download_path(@session), class: "btn btn-primary" %>
<% end %>
<% if @recording.youtube_watch_url %>
@@ -53,25 +80,6 @@
<% end %>
</div>
<% end %>
<% elsif @recording.ready? && @recording.youtube_watch_url.present? %>
<div class="live-player-wrap live-player-wrap--embed">
<iframe src="https://www.youtube.com/embed/<%= @recording.youtube_video_id %>"
title="<%= @recording.title_or_default %>"
class="live-player-embed"
allow="accelerometer; autoplay; clipboard-write; encrypted-media; gyroscope; picture-in-picture"
allowfullscreen></iframe>
<%= render "public/live/player_overlays",
match: @match,
session: @session,
stream_closed: true,
on_air: false,
badge_label: t("replay.show.badge_label"),
badge_class: "badge-ended" %>
</div>
<p class="replay-show__meta-line">
<%= t("replay.show.youtube_only_body") %>
<%= link_to t("replay.show.youtube_link"), @recording.youtube_watch_url, target: "_blank", rel: "noopener" %>
</p>
<% elsif @recording.ready? %>
<div class="stream-ended" role="status">
<h2><%= t("replay.show.file_missing_title") %></h2>
@@ -158,16 +158,16 @@
<% end %>
<span class="visually-hidden"><%= privacy_label %></span>
<% end %>
<% if ent.phone_download_enabled? && rec.ready? %>
<% if ent.phone_download_enabled? && rec.ready? && rec.storage_key.present? %>
<%= link_to "MP4", public_replay_download_path(rec.stream_session_id), class: "replay-archive__action replay-archive__action--secondary", title: t("recordings.archive.download_mp4_title") %>
<% end %>
<% if ent.premium_full? && ent.youtube_enabled? && rec.ready? && rec.youtube_video_id.blank? %>
<% if ent.premium_full? && ent.youtube_enabled? && rec.ready? && rec.storage_key.present? && rec.youtube_video_id.blank? %>
<%= button_to "YT", paths.publish_youtube.call(rec), method: :post, class: "replay-archive__action replay-archive__action--secondary", title: t("recordings.archive.publish_youtube_title"), form: { class: "replay-archive__action-form" } %>
<% elsif rec.youtube_watch_url %>
<%= link_to "YT", rec.youtube_watch_url, class: "replay-archive__action replay-archive__action--secondary", target: "_blank", rel: "noopener", title: t("recordings.archive.open_youtube_title") %>
<% end %>
<% has_youtube = rec.youtube_video_id.present? && !rec.youtube_video_id.to_s.start_with?("mock_") %>
<% has_site = rec.storage_key.present? || %w[ready processing failed].include?(rec.status) %>
<% has_site = rec.available_in_archive? || %w[processing failed].include?(rec.status) %>
<% delete_confirm = t("recordings.archive.delete_confirm") %>
<%= button_to paths.destroy.call(rec), method: :delete,
params: filter_params,
@@ -196,6 +196,24 @@ module MatchLiveTv
ENV.fetch("REPLAY_MEDIA_REDIRECT", "true") == "true"
end
# Retention copia temporanea post-live YouTube (ore). Clamp 2472, default 48.
def youtube_temp_replay_retention_hours
raw = ENV.fetch("YOUTUBE_TEMP_REPLAY_RETENTION_HOURS", "48").to_i
raw.clamp(24, 72)
end
def youtube_replay_verify_max_attempts
ENV.fetch("YOUTUBE_REPLAY_VERIFY_MAX_ATTEMPTS", "12").to_i
end
def youtube_replay_verify_base_interval_secs
ENV.fetch("YOUTUBE_REPLAY_VERIFY_BASE_INTERVAL_SECS", "300").to_i
end
def youtube_replay_verify_grace_secs
ENV.fetch("YOUTUBE_REPLAY_VERIFY_GRACE_SECS", "120").to_i
end
def ops_http_rails_url
ENV.fetch("OPS_HTTP_RAILS_URL", "http://edge/up")
end
+8
View File
@@ -111,6 +111,7 @@ de:
table:
match: Spiel
status: Status
client: Client
ingest: Ingest
start: Start
link: Link
@@ -276,6 +277,7 @@ de:
duration: Dauer
ingest: Ingest
disconnects: Verbindungsabbrüche
client: Client
link: Link
detail: Details
regia: Regie
@@ -308,6 +310,11 @@ de:
opponent: Gegner
operator: Operator
platform: Plattform
client_os: System
app_version: App-Version
device: Gerät
os_version: OS-Version
carrier: Mobilfunkanbieter
privacy: Privacy
quality: Qualität
min_quality: Min. Qualität
@@ -346,6 +353,7 @@ de:
generate_button: Regie-Link erzeugen
ingest:
none: "—"
decommissioned: "(Knoten entfernt)"
role:
home: Home-lab
lab: Lab
+8
View File
@@ -111,6 +111,7 @@ en:
table:
match: Match
status: Status
client: Client
ingest: Ingest
start: Start
link: Link
@@ -276,6 +277,7 @@ en:
duration: Duration
ingest: Ingest
disconnects: Disconnects
client: Client
link: Link
detail: Details
regia: Control
@@ -308,6 +310,11 @@ en:
opponent: Opponent
operator: Operator
platform: Platform
client_os: OS
app_version: App version
device: Device
os_version: OS version
carrier: Carrier
privacy: Privacy
quality: Quality
min_quality: Min quality
@@ -346,6 +353,7 @@ en:
generate_button: Generate control link
ingest:
none: "—"
decommissioned: "(node removed)"
role:
home: Home-lab
lab: Lab
+8
View File
@@ -111,6 +111,7 @@ es:
table:
match: Partido
status: Estado
client: Cliente
ingest: Ingest
start: Inicio
link: Enlace
@@ -276,6 +277,7 @@ es:
duration: Duración
ingest: Ingest
disconnects: Desconexiones
client: Cliente
link: Enlace
detail: Detalle
regia: Regie
@@ -308,6 +310,11 @@ es:
opponent: Rival
operator: Operador
platform: Plataforma
client_os: Sistema
app_version: Versión app
device: Dispositivo
os_version: Versión OS
carrier: Operador móvil
privacy: Privacidad
quality: Calidad
min_quality: Calidad mínima
@@ -346,6 +353,7 @@ es:
generate_button: Generar enlace de regie
ingest:
none: "—"
decommissioned: "(nodo eliminado)"
role:
home: Home-lab
lab: Lab
+8
View File
@@ -111,6 +111,7 @@ fr:
table:
match: Match
status: Statut
client: Client
ingest: Ingest
start: Début
link: Lien
@@ -276,6 +277,7 @@ fr:
duration: Durée
ingest: Ingest
disconnects: Déconnexions
client: Client
link: Lien
detail: Détail
regia: Régie
@@ -308,6 +310,11 @@ fr:
opponent: Adversaire
operator: Opérateur
platform: Plateforme
client_os: Système
app_version: Version app
device: Appareil
os_version: Version OS
carrier: Opérateur mobile
privacy: Confidentialité
quality: Qualité
min_quality: Qualité mini
@@ -346,6 +353,7 @@ fr:
generate_button: Générer le lien de régie
ingest:
none: "—"
decommissioned: "(nœud retiré)"
role:
home: Home-lab
lab: Lab
+8
View File
@@ -115,6 +115,7 @@ it:
table:
match: Partita
status: Stato
client: Client
ingest: Ingest
start: Inizio
link: Link
@@ -297,6 +298,7 @@ it:
duration: Durata
ingest: Ingest
disconnects: Disconnessioni
client: Client
link: Link
detail: Dettaglio
regia: Regia
@@ -329,6 +331,11 @@ it:
opponent: Avversario
operator: Operatore
platform: Piattaforma
client_os: Sistema
app_version: Versione app
device: Dispositivo
os_version: Versione OS
carrier: Operatore telefonico
privacy: Privacy
quality: Qualità
min_quality: Qualità minima
@@ -367,6 +374,7 @@ it:
generate_button: Genera link regia
ingest:
none: "—"
decommissioned: "(nodo rimosso)"
role:
home: Home-lab
lab: Lab
@@ -0,0 +1,15 @@
# frozen_string_literal: true
class AddClientTelemetryToStreamSessions < ActiveRecord::Migration[7.2]
def change
change_table :stream_sessions, bulk: true do |t|
t.string :client_os
t.string :app_version
t.string :app_build
t.string :device_manufacturer
t.string :device_model
t.string :os_version
t.string :carrier
end
end
end
@@ -0,0 +1,15 @@
# frozen_string_literal: true
class AddTemporaryReplayFieldsToRecordings < ActiveRecord::Migration[7.2]
def change
change_table :recordings, bulk: true do |t|
t.string :storage_policy, null: false, default: "retained"
t.datetime :temp_expires_at
t.datetime :youtube_verified_at
t.datetime :local_media_purged_at
end
add_index :recordings, %i[storage_policy temp_expires_at],
name: "index_recordings_on_storage_policy_and_temp_expires_at"
end
end
@@ -0,0 +1,51 @@
# frozen_string_literal: true
class AddIngestSnapshotToStreamSessions < ActiveRecord::Migration[7.2]
def up
add_column :stream_sessions, :ingest_slug, :string
add_column :stream_sessions, :ingest_role, :string
add_index :stream_sessions, :ingest_slug
execute <<~SQL
UPDATE stream_sessions ss
SET ingest_slug = sn.slug,
ingest_role = sn.role
FROM stream_nodes sn
WHERE ss.stream_node_id = sn.id
AND ss.ingest_slug IS NULL
SQL
execute <<~SQL
UPDATE stream_sessions ss
SET ingest_slug = ev.slug
FROM (
SELECT DISTINCT ON (stream_session_id)
stream_session_id,
metadata->>'stream_node' AS slug
FROM stream_events
WHERE event_type = 'pairing'
AND COALESCE(metadata->>'stream_node', '') <> ''
ORDER BY stream_session_id, occurred_at ASC
) ev
WHERE ss.id = ev.stream_session_id
AND ss.ingest_slug IS NULL
SQL
execute <<~SQL
UPDATE stream_sessions
SET ingest_role = CASE
WHEN ingest_slug = 'home' THEN 'home'
WHEN ingest_slug LIKE 'ingest-lab%' THEN 'lab'
WHEN ingest_slug LIKE 'ingest-%' THEN 'cloud'
ELSE ingest_role
END
WHERE ingest_slug IS NOT NULL AND ingest_role IS NULL
SQL
end
def down
remove_index :stream_sessions, :ingest_slug
remove_column :stream_sessions, :ingest_role
remove_column :stream_sessions, :ingest_slug
end
end
+16 -1
View File
@@ -10,7 +10,7 @@
#
# It's strongly recommended that you check this file into your version control system.
ActiveRecord::Schema[7.2].define(version: 2026_08_20_220000) do
ActiveRecord::Schema[7.2].define(version: 2026_08_28_124000) do
# These are extensions that must be enabled in order to support this database
enable_extension "pgcrypto"
enable_extension "plpgsql"
@@ -338,10 +338,15 @@ ActiveRecord::Schema[7.2].define(version: 2026_08_20_220000) do
t.datetime "expiry_warning_sent_at"
t.string "youtube_video_id"
t.datetime "youtube_published_at"
t.string "storage_policy", default: "retained", null: false
t.datetime "temp_expires_at"
t.datetime "youtube_verified_at"
t.datetime "local_media_purged_at"
t.index ["deleted_at"], name: "index_recordings_on_deleted_at"
t.index ["expires_at"], name: "index_recordings_on_expires_at"
t.index ["privacy_status"], name: "index_recordings_on_privacy_status"
t.index ["storage_key"], name: "index_recordings_on_storage_key", unique: true, where: "(storage_key IS NOT NULL)"
t.index ["storage_policy", "temp_expires_at"], name: "index_recordings_on_storage_policy_and_temp_expires_at"
t.index ["stream_session_id"], name: "index_recordings_on_stream_session_id", unique: true
t.index ["team_id", "status"], name: "index_recordings_on_team_id_and_status"
t.index ["team_id"], name: "index_recordings_on_team_id"
@@ -427,6 +432,16 @@ ActiveRecord::Schema[7.2].define(version: 2026_08_20_220000) do
t.uuid "stream_node_id"
t.string "min_quality_preset", default: "auto", null: false
t.boolean "audio_muted", default: false, null: false
t.string "client_os"
t.string "app_version"
t.string "app_build"
t.string "device_manufacturer"
t.string "device_model"
t.string "os_version"
t.string "carrier"
t.string "ingest_slug"
t.string "ingest_role"
t.index ["ingest_slug"], name: "index_stream_sessions_on_ingest_slug"
t.index ["match_id"], name: "index_stream_sessions_on_match_id"
t.index ["publish_token"], name: "index_stream_sessions_on_publish_token", unique: true
t.index ["regia_token_digest"], name: "index_stream_sessions_on_regia_token_digest", unique: true
+33
View File
@@ -20,4 +20,37 @@ namespace :recordings do
orphan = Mediamtx::CleanupOrphanPaths.new.call
puts "Path MediaMTX orfani rimossi: #{orphan.removed} (saltati: #{orphan.skipped})"
end
desc "Elimina copie temporanee YouTube scadute (solo storage, non soft-delete)"
task purge_temporary: :environment do
count = Recordings::PurgeTemporaryMediaJob.new.perform
puts "Purge temporary replays: processed=#{count}"
end
desc "Audit dry-run: recording YouTube con ancora MP4 in Garage (duplicati storici). " \
"Cleanup solo se APPLY=1 (non implementato di default — solo report)."
task audit_youtube_duplicates: :environment do
apply = ENV["APPLY"].to_s == "1"
scope = Recording.not_deleted
.where.not(youtube_video_id: [nil, ""])
.where.not(storage_key: [nil, ""])
.where("youtube_video_id NOT LIKE 'mock_%'")
puts "Audit YouTube duplicates (dry-run=#{!apply}): count=#{scope.count}"
scope.find_each do |rec|
puts [
"id=#{rec.id}",
"policy=#{rec.storage_policy}",
"key=#{rec.storage_key}",
"yt=#{rec.youtube_video_id}",
"purged_at=#{rec.local_media_purged_at.inspect}",
"temp_expires=#{rec.temp_expires_at.inspect}"
].join(" ")
end
if apply
puts "APPLY=1 richiesto ma cleanup storico automatico NON eseguito " \
"(fuori scope; solo dry-run supportato)."
end
end
end
@@ -34,4 +34,35 @@ RSpec.describe StreamSession do
expect(session.mediamtx_api_base_url).to eq("http://10.0.0.2:9997")
expect(session.mediamtx_internal_rtmp_url).to eq("rtmp://10.0.0.2:1935")
end
it "keeps ingest snapshot after the stream node is destroyed" do
node = StreamNode.create!(
slug: "ingest-01",
hostname: "ingest-01.mltv-stream.net",
role: "cloud",
status: "ready",
provider: "hetzner",
rtmp_base_url: "rtmp://ingest-01.mltv-stream.net:1935",
hls_base_url: "https://ingest-01.mltv-stream.net/hls",
api_base_url: "http://10.0.0.2:9997",
internal_rtmp_url: "rtmp://10.0.0.2:1935",
max_publishers: 4,
max_relays: 4
)
user = User.create!(email: "snap@example.com", name: "S", password: "Password123", role: "coach")
club = Club.create!(name: "Club Snap", sport: "volleyball")
team = club.teams.create!(name: "Team", sport: "volleyball", slug: "team-snap")
match = team.matches.create!(opponent_name: "Opp", scheduled_at: 1.hour.from_now)
session = StreamSession.create!(
match: match, user: user, platform: "matchlivetv", status: "ended", stream_node: node
)
expect(session.ingest_slug).to eq("ingest-01")
expect(session.ingest_role).to eq("cloud")
node.destroy!
session.reload
expect(session.stream_node).to be_nil
expect(session.ingest_slug_display).to eq("ingest-01")
expect(session.ingest_role_display).to eq("cloud")
end
end
@@ -60,4 +60,31 @@ RSpec.describe "Admin sessions index", type: :request do
expect(response.body).to include("Dettaglio sessione").or include("Session details")
expect(response.body).to include(session_a.id)
end
it "mostra lo snapshot ingest anche se il nodo è stato decommissionato" do
node = StreamNode.create!(
slug: "ingest-09",
hostname: "ingest-09.mltv-stream.net",
role: "cloud",
status: "ready",
provider: "hetzner",
rtmp_base_url: "rtmp://ingest-09.mltv-stream.net:1935",
hls_base_url: "https://collaudo.example/hls",
api_base_url: "http://10.0.0.9:9997",
max_publishers: 4,
max_relays: 4
)
session_a.update!(stream_node: node)
node.destroy!
session_a.reload
get admin_sessions_path
expect(response).to have_http_status(:ok)
expect(response.body).to include("ingest-09")
expect(response.body).to include("Hetzner")
get admin_session_path(session_a)
expect(response.body).to include("ingest-09")
expect(response.body).to include("nodo rimosso").or include("node removed")
end
end
@@ -0,0 +1,60 @@
# frozen_string_literal: true
require "rails_helper"
RSpec.describe "Public replay show youtube-only", type: :request do
let!(:club) { Club.create!(name: "Club", sport: "volleyball", primary_color: "#e53935", secondary_color: "#ffffff") }
let!(:team) { club.teams.create!(name: "Team", sport: "volleyball") }
let!(:user) { User.create!(email: "pub-replay@test.com", name: "Coach", password: "Password123", role: "coach") }
let!(:match) { team.matches.create!(opponent_name: "Rival", scheduled_at: 1.day.ago) }
let!(:session) do
StreamSession.create!(
match: match, user: user, platform: "youtube", status: "ended", ended_at: 1.hour.ago,
youtube_broadcast_id: "bcast"
)
end
it "mostra embed YouTube senza storage_key" do
Recording.create!(
stream_session: session,
team: team,
status: "ready",
storage_policy: "temporary",
storage_key: nil,
youtube_video_id: "publicYtVid01",
youtube_verified_at: Time.current,
local_media_purged_at: Time.current,
privacy_status: "public",
title: "Derby",
expires_at: nil,
metadata: { "source_platform" => "youtube" }
)
get "/replay/#{session.id}"
expect(response).to have_http_status(:ok)
expect(response.body).to include("youtube.com/embed/publicYtVid01")
expect(response.body).not_to include("id=\"replay-player\"")
end
it "mostra player MP4 per retained con storage_key" do
mltv = StreamSession.create!(
match: match, user: user, platform: "matchlivetv", status: "ended", ended_at: 1.hour.ago
)
Recording.create!(
stream_session: mltv,
team: team,
status: "ready",
storage_policy: "retained",
storage_key: "teams/#{team.id}/sessions/#{mltv.id}/replay.mp4",
privacy_status: "public",
title: "Casa",
expires_at: 30.days.from_now,
metadata: { "source_platform" => "matchlivetv" }
)
get "/replay/#{mltv.id}"
expect(response).to have_http_status(:ok)
expect(response.body).to include("id=\"replay-player\"")
expect(response.body).not_to include("youtube.com/embed/")
end
end
@@ -0,0 +1,94 @@
# frozen_string_literal: true
require "rails_helper"
RSpec.describe "Api::V1 team recordings replay_source", type: :request do
let!(:user) { User.create!(email: "rec-api@test.it", name: "U", password: "Password123", role: "coach") }
let!(:club) { Club.create!(name: "C", sport: "volleyball", primary_color: "#e53935", secondary_color: "#ffffff") }
let!(:owner) { club.club_memberships.create!(user: user, role: "owner") }
let!(:team) { club.teams.create!(name: "T", sport: "volleyball") }
let!(:match) { team.matches.create!(opponent_name: "Opp", sport: "volleyball", scheduled_at: 1.day.ago) }
def auth_headers
post "/api/v1/auth/login", params: { email: user.email, password: "Password123" }
token = response.parsed_body["access_token"]
{ "Authorization" => "Bearer #{token}" }
end
before do
load Rails.root.join("db/seeds/plans.rb")
Billing::AssignPlan.call(club: club, plan_slug: "premium_full")
end
it "youtube-only: replay_source youtube, no playback/download" do
session = StreamSession.create!(
match: match, user: user, platform: "youtube", status: "ended", ended_at: 1.hour.ago,
youtube_broadcast_id: "bcast1"
)
Recording.create!(
stream_session: session,
team: team,
status: "ready",
storage_policy: "temporary",
storage_key: nil,
youtube_video_id: "ytVideoId123",
youtube_verified_at: Time.current,
local_media_purged_at: Time.current,
privacy_status: "public",
expires_at: nil,
metadata: { "source_platform" => "youtube" }
)
get "/api/v1/teams/#{team.id}/recordings", headers: auth_headers
expect(response).to have_http_status(:ok)
item = response.parsed_body.find { |r| r["session_id"] == session.id } ||
response.parsed_body.find { |r| r["youtube_video_id"] == "ytVideoId123" }
expect(item).to be_present
expect(item["replay_source"]).to eq("youtube")
expect(item["youtube_watch_url"]).to include("ytVideoId123")
expect(item["playback_url"]).to be_nil
expect(item["download_enabled"]).to eq(false)
end
it "matchlivetv retained: replay_source matchlivetv with playback" do
session = StreamSession.create!(
match: match, user: user, platform: "matchlivetv", status: "ended", ended_at: 1.hour.ago
)
Recording.create!(
stream_session: session,
team: team,
status: "ready",
storage_policy: "retained",
storage_key: "teams/#{team.id}/sessions/#{session.id}/replay.mp4",
privacy_status: "public",
expires_at: 30.days.from_now,
metadata: { "source_platform" => "matchlivetv" }
)
get "/api/v1/teams/#{team.id}/recordings", headers: auth_headers
expect(response).to have_http_status(:ok)
item = response.parsed_body.find { |r| r["session_id"] == session.id }
expect(item["replay_source"]).to eq("matchlivetv")
expect(item["playback_url"]).to be_present
expect(item["download_enabled"]).to eq(true)
end
it "download API rejects youtube-only recording" do
session = StreamSession.create!(
match: match, user: user, platform: "youtube", status: "ended", ended_at: 1.hour.ago
)
recording = Recording.create!(
stream_session: session,
team: team,
status: "ready",
storage_policy: "temporary",
storage_key: nil,
youtube_video_id: "ytVideoId999",
youtube_verified_at: Time.current,
privacy_status: "public"
)
get "/api/v1/recordings/#{recording.id}/download", headers: auth_headers
expect(response).to have_http_status(:unprocessable_entity)
end
end
@@ -25,7 +25,9 @@ RSpec.describe Recordings::FinalizeSession do
rec = described_class.new(session).call
expect(rec.status).to eq("processing")
expect(rec.privacy_status).to eq("public")
expect(rec.storage_policy).to eq("retained")
expect(rec.expires_at).to be > 29.days.from_now
expect(rec.metadata["auto_publish_youtube"]).to eq(false)
end
it "skips free plan" do
@@ -39,6 +41,18 @@ RSpec.describe Recordings::FinalizeSession do
expect(rec.expires_at).to be > 89.days.from_now
end
it "non riazzerare un recording già ready su stop ripetuto" do
rec = described_class.new(session).call
rec.update!(status: "ready", storage_key: "teams/x/sessions/#{session.id}/replay.mp4", byte_size: 12_345)
again = described_class.new(session).call
expect(again.id).to eq(rec.id)
expect(again.status).to eq("ready")
expect(again.storage_key).to eq("teams/x/sessions/#{session.id}/replay.mp4")
expect(again.byte_size).to eq(12_345)
expect(again.error_message).to be_nil
end
it "non crea recording se abbonamento scaduto" do
Billing::AssignPlan.call(club: club, plan_slug: "premium_light")
club.subscription.update!(status: "canceled")
@@ -0,0 +1,79 @@
require "rails_helper"
require "zlib"
require "rubygems/package"
RSpec.describe Recordings::PullFromCloudNode do
let!(:club) { Club.create!(name: "Club", sport: "volleyball", primary_color: "#e53935", secondary_color: "#ffffff") }
let!(:team) { club.teams.create!(name: "Team", sport: "volleyball") }
let!(:user) { User.create!(email: "cpx-rec@test.com", name: "Coach", password: "Password123", role: "coach") }
let!(:match) { team.matches.create!(opponent_name: "Rival", scheduled_at: 1.day.from_now) }
let!(:home) do
StreamNode.create!(
slug: "spec-home-cpx-rec",
hostname: "home.local",
role: "home",
status: "ready",
provider: "local",
rtmp_base_url: "rtmp://home/live",
hls_base_url: "http://home/hls",
api_base_url: "http://home:9997",
max_publishers: 2,
max_relays: 2
)
end
let!(:cloud) do
StreamNode.create!(
slug: "spec-cpx-rec",
hostname: "cpx.example",
role: "cloud",
status: "ready",
provider: "hetzner",
rtmp_base_url: "rtmp://cpx/live",
hls_base_url: "http://cpx/hls",
api_base_url: "http://cpx:9997",
max_publishers: 4,
max_relays: 4,
metadata: { "public_ip" => "203.0.113.10" }
)
end
def session_on(node)
StreamSession.create!(match: match, user: user, platform: "youtube", status: "ended", stream_node: node)
end
it "non è applicable sul nodo home" do
expect(described_class.new(session_on(home)).applicable?).to eq(false)
end
it "è applicable sul CPX con agent URL" do
expect(described_class.new(session_on(cloud)).applicable?).to eq(true)
end
it "estrae i segmenti dal tar dell'agent" do
session = session_on(cloud)
raw = "fake-mp4-bytes"
tar_io = StringIO.new
Zlib::GzipWriter.wrap(tar_io) do |gz|
Gem::Package::TarWriter.new(gz) do |tar|
tar.add_file_simple("clip.mp4", 0o644, raw.bytesize) { |io| io.write(raw) }
end
end
gz_bytes = tar_io.string
http = instance_double(Net::HTTP)
allow(Net::HTTP).to receive(:new).and_return(http)
allow(http).to receive(:open_timeout=)
allow(http).to receive(:read_timeout=)
success = Net::HTTPOK.new("1.1", "200", "OK")
allow(success).to receive(:body).and_return(gz_bytes)
allow(http).to receive(:request).and_return(success)
dest = described_class.new(session).fetch
expect(dest).to be_present
files = Dir.glob(File.join(dest, "**", "*.mp4"))
expect(files.size).to eq(1)
expect(File.binread(files.first)).to eq(raw)
ensure
FileUtils.remove_entry(dest) if dest && Dir.exist?(dest)
end
end
@@ -0,0 +1,40 @@
# frozen_string_literal: true
require "rails_helper"
RSpec.describe Recordings::StoragePolicy do
let!(:club) { Club.create!(name: "Club", sport: "volleyball", primary_color: "#e53935", secondary_color: "#ffffff") }
let!(:team) { club.teams.create!(name: "Team", sport: "volleyball") }
let!(:user) { User.create!(email: "coach@test.com", name: "Coach", password: "Password123", role: "coach") }
let!(:match) { team.matches.create!(opponent_name: "Rival", scheduled_at: 1.day.from_now) }
def session_for(platform)
StreamSession.create!(
match: match,
user: user,
platform: platform,
status: "ended",
privacy_status: "public",
ended_at: Time.current
)
end
before do
load Rails.root.join("db/seeds/plans.rb")
end
it "returns temporary for youtube with premium" do
Billing::AssignPlan.call(club: club, plan_slug: "premium_full")
expect(described_class.call(session_for("youtube"))).to eq("temporary")
end
it "returns retained for matchlivetv with premium" do
Billing::AssignPlan.call(club: club, plan_slug: "premium_light")
expect(described_class.call(session_for("matchlivetv"))).to eq("retained")
end
it "returns none without recordings entitlement" do
Billing::AssignPlan.call(club: club, plan_slug: "free")
expect(described_class.call(session_for("youtube"))).to eq("none")
end
end
@@ -0,0 +1,243 @@
# frozen_string_literal: true
require "rails_helper"
RSpec.describe "YouTube temporary replay retention" do
let!(:club) { Club.create!(name: "Club", sport: "volleyball", primary_color: "#e53935", secondary_color: "#ffffff") }
let!(:team) { club.teams.create!(name: "Team", sport: "volleyball") }
let!(:user) { User.create!(email: "coach@test.com", name: "Coach", password: "Password123", role: "coach") }
let!(:match) { team.matches.create!(opponent_name: "Rival", scheduled_at: 1.day.from_now) }
before do
load Rails.root.join("db/seeds/plans.rb")
Billing::AssignPlan.call(club: club, plan_slug: "premium_full")
end
def youtube_session!(broadcast_id: "mock_broadcast_abc")
StreamSession.create!(
match: match,
user: user,
platform: "youtube",
status: "ended",
privacy_status: "public",
ended_at: Time.current,
youtube_broadcast_id: broadcast_id,
stream_key: "sk",
rtmp_url: "rtmp://a.rtmp.youtube.com/live2"
)
end
def mltv_session!
StreamSession.create!(
match: match,
user: user,
platform: "matchlivetv",
status: "ended",
privacy_status: "public",
ended_at: Time.current
)
end
def create_temp_ready!(session, storage_key: nil)
key = storage_key || "temporary_replays/teams/#{team.id}/sessions/#{session.id}/replay.mp4"
Recording.create!(
stream_session: session,
team: team,
status: "ready",
storage_policy: "temporary",
storage_key: key,
byte_size: 1_024,
duration_secs: 120,
temp_expires_at: 48.hours.from_now,
expires_at: nil,
privacy_status: "public",
metadata: { "source_platform" => "youtube", "auto_publish_youtube" => false }
)
end
# A. YT → temp → verify OK → metadata → temp eliminata
it "A: verify OK sets youtube metadata and clears temporary media" do
session = youtube_session!
recording = create_temp_ready!(session)
storage = instance_double(Recordings::Storage, delete: true)
allow(Recordings::Storage).to receive(:new).and_return(storage)
result = Recordings::VerifyYoutubeReplay.new(recording).call
expect(result).to be_ok
recording.reload
expect(recording.youtube_video_id).to eq("mock_broadcast_abc")
expect(recording.youtube_verified_at).to be_present
expect(recording).to be_ready
expect(recording.available_in_archive?).to eq(true)
expect(recording.replay_source).to eq("youtube")
Recordings::ClearTemporaryMedia.new(recording, reason: :verified, force: true).call
recording.reload
expect(recording.storage_key).to be_nil
expect(recording.local_media_purged_at).to be_present
expect(recording.youtube_video_id).to eq("mock_broadcast_abc")
expect(recording).to be_ready
expect(recording.available_in_archive?).to eq(true)
expect(storage).to have_received(:delete).with(key: a_string_including("temporary_replays/"))
end
# B. YT → YT non pronto → temp non eliminata
it "B: does not clear temp when VOD is not ready" do
session = youtube_session!(broadcast_id: "real_broadcast_id")
recording = create_temp_ready!(session)
allow_any_instance_of(Youtube::VodStatus).to receive(:fetch).and_return(
Youtube::VodStatus::Result.new(ready: false, video_id: "real_broadcast_id", upload_status: "uploaded")
)
result = Recordings::VerifyYoutubeReplay.new(recording).call
expect(result).to be_retriable
expect(recording.reload.youtube_verified_at).to be_nil
Recordings::ClearTemporaryMedia.new(recording, reason: :verified, force: false).call
expect(recording.reload.storage_key).to be_present
expect(recording.local_media_purged_at).to be_nil
end
# C. YT → verify fail → max retention → safety + log
it "C: safety net purges at temp_expires_at even if unverified" do
session = youtube_session!
recording = create_temp_ready!(session)
recording.update!(temp_expires_at: 1.hour.ago, youtube_verified_at: nil)
storage = instance_double(Recordings::Storage, delete: true)
allow(Recordings::Storage).to receive(:new).and_return(storage)
expect(Rails.logger).to receive(:warn).with(/anomaly_unverified_expiry/).at_least(:once)
Recordings::PurgeTemporaryMediaJob.new.perform
recording.reload
expect(recording.local_media_purged_at).to be_present
expect(recording.storage_key).to be_nil
expect(recording.deleted_at).to be_nil
expect(recording.status).to eq("ready")
end
# D. solo MLTV → retained + retention piano
it "D: matchlivetv finalize uses retained and plan retention" do
session = mltv_session!
rec = Recordings::FinalizeSession.new(session).call
expect(rec.storage_policy).to eq("retained")
expect(rec.expires_at).to be > 89.days.from_now
expect(rec.temp_expires_at).to be_nil
expect(rec.metadata["auto_publish_youtube"]).to eq(false)
end
# E. no replay piano → none
it "E: free plan skips finalize (none)" do
Billing::AssignPlan.call(club: club, plan_slug: "free")
session = youtube_session!
expect(Recordings::StoragePolicy.call(session)).to eq("none")
expect(Recordings::FinalizeSession.new(session).call).to be_nil
end
# F. job doppio → idempotente
it "F: clear temporary media is idempotent" do
session = youtube_session!
recording = create_temp_ready!(session)
recording.update!(youtube_video_id: "mock_broadcast_abc", youtube_verified_at: Time.current)
storage = instance_double(Recordings::Storage, delete: true)
allow(Recordings::Storage).to receive(:new).and_return(storage)
Recordings::ClearTemporaryMedia.new(recording, reason: :verified, force: true).call
Recordings::ClearTemporaryMedia.new(recording.reload, reason: :verified, force: true).call
expect(storage).to have_received(:delete).once
expect(recording.reload.local_media_purged_at).to be_present
end
# G. file assente → cleanup ok
it "G: clear succeeds when storage delete raises (missing object)" do
session = youtube_session!
recording = create_temp_ready!(session)
recording.update!(youtube_verified_at: Time.current, youtube_video_id: "mock_broadcast_abc")
storage = instance_double(Recordings::Storage)
allow(storage).to receive(:delete).and_raise(Recordings::Storage::Error, "NoSuchKey")
allow(Recordings::Storage).to receive(:new).and_return(storage)
expect do
Recordings::ClearTemporaryMedia.new(recording, reason: :verified, force: true).call
end.not_to raise_error
expect(recording.reload.local_media_purged_at).to be_present
expect(recording.storage_key).to be_nil
end
# H. archivio ok senza S3
it "H: available_in_archive without storage_key when youtube linked" do
session = youtube_session!
recording = create_temp_ready!(session)
recording.update!(
storage_key: nil,
local_media_purged_at: Time.current,
youtube_video_id: "abc123XYZ01",
youtube_verified_at: Time.current
)
expect(recording.available_in_archive?).to eq(true)
expect(recording.playable_on_site?).to eq(true)
expect(recording.replay_source).to eq("youtube")
expect(recording.playback_stream_url).to be_nil
expect(recording.download_api_path).to be_nil
end
# I. API replay_source
it "I: exposes replay_source youtube for youtube-only ready recording" do
session = youtube_session!
recording = create_temp_ready!(session)
recording.update!(
storage_key: nil,
youtube_video_id: "abc123XYZ01",
youtube_verified_at: Time.current,
local_media_purged_at: Time.current
)
expect(recording.replay_source).to eq("youtube")
expect(recording.youtube_watch_url).to include("abc123XYZ01")
end
# J. download assente per youtube-only
it "J: DownloadUrl rejects youtube-only without storage_key" do
session = youtube_session!
recording = create_temp_ready!(session)
recording.update!(
storage_key: nil,
youtube_video_id: "abc123XYZ01",
youtube_verified_at: Time.current
)
expect do
Recordings::DownloadUrl.new(recording, viewer: user).call
end.to raise_error(ArgumentError, /File non disponibile/)
end
it "finalize youtube sets temporary policy and temp_expires_at" do
session = youtube_session!
rec = Recordings::FinalizeSession.new(session).call
expect(rec.storage_policy).to eq("temporary")
expect(rec.expires_at).to be_nil
expect(rec.temp_expires_at).to be_within(2.seconds).of(48.hours.from_now)
expect(rec.metadata["auto_publish_youtube"]).to eq(false)
end
it "PostProcess schedules verify for temporary, not PublishToYoutube" do
session = youtube_session!
recording = create_temp_ready!(session)
expect(Recordings::VerifyYoutubeReplayJob).to receive(:perform_in)
expect(Recordings::PublishToYoutubeJob).not_to receive(:perform_async)
Recordings::PostProcessJob.new.perform(recording.id)
end
it "PostProcess still schedules verify if notify/mailer fails" do
session = youtube_session!
recording = create_temp_ready!(session)
allow_any_instance_of(Recordings::NotifyReady).to receive(:call).and_raise(Errno::ECONNREFUSED)
expect(Recordings::VerifyYoutubeReplayJob).to receive(:perform_in)
expect(Recordings::PublishToYoutubeJob).not_to receive(:perform_async)
Recordings::PostProcessJob.new.perform(recording.id)
end
it "expired_pending_purge ignores temporary with nil expires_at" do
session = youtube_session!
create_temp_ready!(session)
expect(Recording.expired_pending_purge.count).to eq(0)
end
end
@@ -0,0 +1,45 @@
# frozen_string_literal: true
require "rails_helper"
RSpec.describe Sessions::ApplyClientInfo do
let!(:user) { User.create!(email: "client-info@test.it", name: "U", password: "Password123", role: "coach") }
let!(:club) { Club.create!(name: "C", sport: "volleyball", primary_color: "#e53935", secondary_color: "#ffffff") }
let!(:team) { club.teams.create!(name: "T", sport: "volleyball") }
let!(:match) { team.matches.create!(opponent_name: "Opp", sport: "volleyball") }
let!(:session) { StreamSession.create!(match: match, user: user, platform: "matchlivetv", status: "idle") }
it "applica i campi client sulla sessione" do
described_class.call(session, {
os: "android",
app_version: "1.4.0",
app_build: "42",
device_manufacturer: "Samsung",
device_model: "SM-G991B",
os_version: "14",
carrier: "TIM"
})
session.reload
expect(session.client_os).to eq("android")
expect(session.app_version).to eq("1.4.0")
expect(session.app_build).to eq("42")
expect(session.device_manufacturer).to eq("Samsung")
expect(session.device_model).to eq("SM-G991B")
expect(session.os_version).to eq("14")
expect(session.carrier).to eq("TIM")
end
it "ignora os sconosciuti" do
described_class.call(session, { os: "windows" })
expect(session.reload.client_os).to be_nil
end
it "assegna senza salvare su record non persistito" do
draft = StreamSession.new(match: match, user: user, platform: "matchlivetv", status: "idle")
described_class.call(draft, { os: "ios", app_version: "2.0.0" })
expect(draft).not_to be_persisted
expect(draft.client_os).to eq("ios")
expect(draft.app_version).to eq("2.0.0")
end
end
@@ -0,0 +1,50 @@
# frozen_string_literal: true
require "rails_helper"
RSpec.describe Sessions::Create, "capacity / autoscaler kick" do
def with_env(vars)
previous = vars.keys.index_with { |k| ENV[k] }
vars.each { |k, v| ENV[k] = v }
yield
ensure
previous.each { |k, v| v.nil? ? ENV.delete(k) : ENV[k] = v }
end
let(:redis) { Redis.new(url: ENV.fetch("REDIS_URL", "redis://redis:6379/0")) }
let(:user) { User.create!(email: "cap@example.com", name: "C", password: "Password123", role: "coach") }
let(:club) { Club.create!(name: "Cap Club", sport: "volleyball") }
let(:team) { club.teams.create!(name: "Cap Team", sport: "volleyball", slug: "cap-team") }
let(:match) { team.matches.create!(opponent_name: "Opp", scheduled_at: 1.hour.from_now) }
before do
redis.del(Streams::AutoscalerJob::KICK_DEBOUNCE_KEY)
Billing::AssignPlan.call(club: club, plan_slug: "premium_full")
end
it "kicks AutoscalerJob and raises stream_capacity_scaling when autoscaler is on" do
with_env("STREAM_AUTOSCALE_ENABLED" => "1") do
allow(Streams::NodeRegistry).to receive(:allocate!).and_raise(Streams::NodeRegistry::NoCapacityError, "full")
expect(Streams::AutoscalerJob).to receive(:kick!).and_return(true)
expect {
described_class.new(user: user, match: match, params: { platform: "matchlivetv" }).call
}.to raise_error(Teams::EntitlementError) { |e|
expect(e.code).to eq("stream_capacity_scaling")
}
end
end
it "raises stream_capacity_exhausted when autoscaler is off" do
with_env("STREAM_AUTOSCALE_ENABLED" => "0") do
allow(Streams::NodeRegistry).to receive(:allocate!).and_raise(Streams::NodeRegistry::NoCapacityError, "full")
expect(Streams::AutoscalerJob).not_to receive(:kick!)
expect {
described_class.new(user: user, match: match, params: { platform: "matchlivetv" }).call
}.to raise_error(Teams::EntitlementError) { |e|
expect(e.code).to eq("stream_capacity_exhausted")
}
end
end
end
+86 -99
View File
@@ -2,156 +2,143 @@
## Panoramica
Al termine di ogni diretta Premium, MediaMTX registra i segmenti. Sidekiq unisce i segmenti, genera thumbnail, carica su **Garage (S3-compatible)** e attiva i servizi di notifica, statistiche e (opzionale) republicazione YouTube.
Al termine di ogni diretta Premium, MediaMTX registra i segmenti. Sidekiq unisce i segmenti, genera thumbnail e carica su **Garage (S3-compatible)**.
- **Live MatchLiveTV (HLS)**: copia **permanente** (policy `retained`) con retention del piano.
- **Live YouTube**: copia **temporanea** di sicurezza (policy `temporary`, default 48h), poi player YouTube in archivio. **Niente** re-upload automatico del MP4 su YouTube.
## Piani e retention
| Piano | Registrazione | Retention | Download MP4 | YouTube VOD |
|-------|---------------|-----------|--------------|-------------|
| Piano | Registrazione | Retention archivio MLTV | Download MP4 | YouTube |
|-------|---------------|-------------------------|--------------|---------|
| Free | No | — | No | No |
| Premium Light | Sì | 30 giorni | Sì | No |
| Premium Full | Sì | 90 giorni | Sì | Sì (auto se diretta YouTube) |
| Premium Light | Sì | 30 giorni (`expires_at`) | Sì (solo se file MLTV) | No |
| Premium Full | Sì | 90 giorni (`expires_at`) se HLS; **nessuna scadenza archivio** se YouTube (`expires_at` nil) | Sì solo con file Garage | Live → VOD nativo; verify collega `youtube_video_id` |
## Storage policy (`Recordings::StoragePolicy`)
```text
if youtube_destination? → temporary
elsif can_create_recordings? → retained
else → none
```
| Policy | Prefisso storage | `expires_at` (archivio) | `temp_expires_at` |
|--------|------------------|-------------------------|-------------------|
| `temporary` | `temporary_replays/...` | `nil` | now + `YOUTUBE_TEMP_REPLAY_RETENTION_HOURS` (2472, default 48) |
| `retained` | `teams/...` | piano `recording_days` | `nil` |
| `none` | — | — | — |
## Archivio logico vs storage fisico
Dopo cleanup della copia temp, la partita **resta** in archivio MatchLiveTV. Cambia solo la sorgente:
| `replay_source` | Player | File Garage permanente |
|-----------------|--------|------------------------|
| `youtube` | embed / link YouTube | no (MP4 temp già eliminato) |
| `matchlivetv` | player MP4 proprietario | sì |
| `none` | non playable | — |
`available_in_archive?` = `ready` e (YouTube URL/`youtube_video_id` **oppure** `storage_key`).
Clear/purge temporary: **solo** oggetti S3 video; mai soft-delete; mai azzerare `youtube_video_id` / metadata / score. Thumbnail preferibilmente conservata.
## Flusso YouTube
```text
Stop → StoragePolicy(temporary) → Finalize (expires_at nil, temp_expires_at)
→ Upload temporary_replays/... → ready
→ VerifyYoutubeReplayJob (backoff)
→ ok: youtube_video_id + youtube_verified_at → ClearTemporaryMediaJob
→ fail fino a max attempts: aspetta safety net
→ PurgeTemporaryMediaJob (cron): a temp_expires_at elimina file + log se non verificato
```
## Funzionalità implementate
### 1. Registrazione automatica
Pipeline: `Sessions::Stop``FinalizeSession``UploadJob` → storage.
Pipeline: `Sessions::Stop``FinalizeSession``UploadJob` → storage`PostProcessJob`.
MediaMTX registra **solo mentre il telefono pubblica** (RTMP connesso). In pausa o con sola slate `alwaysAvailable` la registrazione è disattivata, così il replay non contiene minuti di schermo «Trasmissione in pausa».
MediaMTX registra **solo mentre il telefono pubblica** (RTMP connesso).
### 2. Retention e purge
`Recordings::PurgeExpiredJob` / `rails recordings:purge_expired`
- Retained: `Recordings::PurgeExpiredJob` / `rails recordings:purge_expired` (soft-delete + S3 + YouTube se collegato)
- Temporary: `rails recordings:purge_temporary` (solo storage)
### 3. Archivio Replay
- Web: `/clubs/:id/replays` (gestione società)
- Pubblico: `/replay` con filtro società/squadra
- App mobile: **nessuna** UI replay (gestione solo da sito web)
- Stessa UX listing per YouTube e MLTV (differisce solo il player)
### 4. Riproduzione
- Player MP4: `/replay/:id` + stream `/replay/:id/stream`
- Delivery: Rails verifica accesso → redirect 302 a `/media/...` (edge → Garage). Puma non proxya il file.
- YouTube: embed su `/replay/:id`
- MLTV: MP4 via `/replay/:id/stream` → redirect `/media/...` → Garage
- Contatore visualizzazioni su ogni play (`view_count`)
### 5. Visibilità
- **Pubblico** (`public`): compare in `/replay`, indicizzabile, YouTube `public`
- **Privato** (`unlisted`): non in catalogo pubblico, accessibile solo con link diretto, `noindex`, YouTube `unlisted`
Modificabile dallarchivio società (`/clubs/:id/replays`). Il cambio privacy sincronizza automaticamente il VOD YouTube se presente (`Recordings::SyncYoutubePrivacyJob`).
Pubblico / privato (`unlisted`) come prima. Sync privacy YouTube se `youtube_video_id` presente.
### 6. Eliminazione anticipata
`Recordings::Delete` — elimina video YouTube collegato, oggetti S3/Garage (MP4 + thumbnail) e soft-delete DB.
`Recordings::Delete` manuale — resta distruttiva sul record (distinta dal cleanup temporary automatico).
Stesso comportamento alla scadenza retention (`PurgeExpiredJob` alle 03:00).
### 78. KPI, email, thumbnail, download, views — invariati dove applicabili.
### 7. Dashboard KPI
Replay disponibili, in scadenza (7 gg), spazio occupato, **visualizzazioni totali**.
### 8. Miglioramenti Premium
| Feature | Descrizione |
|---------|-------------|
| **Email replay pronto** | A owner e membri società quando status → `ready` |
| **Email scadenza** | 7 giorni prima di `expires_at` (`rails recordings:expiry_warnings`) |
| **Thumbnail** | Frame ffmpeg, URL `/replay/:id/thumbnail` |
| **Download MP4** | Premium, link presigned 15 min o `/replay/:id/download` |
| **Statistiche views** | `view_count` su ogni replay |
| **metadata JSONB** | `source_platform`, `auto_publish_youtube`, `ai: {}` per estensioni |
| **YouTube VOD** | Premium Full: upload automatico se diretta era su YouTube; manuale da archivio |
## Storage
## Config ENV
```env
YOUTUBE_TEMP_REPLAY_RETENTION_HOURS=48
YOUTUBE_REPLAY_VERIFY_MAX_ATTEMPTS=12
YOUTUBE_REPLAY_VERIFY_BASE_INTERVAL_SECS=300
YOUTUBE_REPLAY_VERIFY_GRACE_SECS=120
REPLAY_STORAGE_ENDPOINT=http://garage:3900
REPLAY_STORAGE_BUCKET=matchlivetv-replays
REPLAY_STORAGE_ACCESS_KEY_ID=...
REPLAY_STORAGE_SECRET_ACCESS_KEY=...
REPLAY_DOWNLOAD_URL_TTL_SECONDS=900
...
```
Valori letti da `MatchLiveTv` (`config/initializers/match_live_tv.rb`).
## Storage Garage
Senza endpoint S3: storage locale `/recordings/replays/`.
### Dev locale con Garage
```bash
cd infra
docker compose up -d # include garage
bash scripts/setup_garage_replays.sh # layout, bucket, chiave → .env
docker compose up -d rails sidekiq
# Verifica pipeline (segmento finto → Garage)
docker compose exec -T rails bundle exec rails replay:e2e_garage
```
### Produzione (`/opt/matchlivetv`)
```bash
cd infra
bash scripts/setup_garage_production.sh # garage.prod.toml, bucket, chiavi → .env
# Rails/Sidekiq usano REPLAY_STORAGE_ENDPOINT=http://garage:3900 (rete Docker)
```
**Credenziali in `.env`:** non compilarle a mano in `.env.example`. Lo script scrive in `infra/.env`:
- `REPLAY_STORAGE_ACCESS_KEY_ID` = Key ID Garage (es. `GK...`)
- `REPLAY_STORAGE_SECRET_ACCESS_KEY` = Secret mostrato una sola volta alla creazione chiave
### Dev / produzione Garage
Vedi script `infra/scripts/setup_garage_*.sh` come in precedenza.
## Cron (produzione)
Installazione sul server:
```bash
cd /opt/matchlivetv/infra
bash scripts/install_production_cron.sh
```
Job attivi (utente `eminux`):
| Orario | Task |
|--------|------|
| 03:00 | `recordings:purge_expired`elimina replay scaduti (DB + Garage) |
| 03:00 | `recordings:purge_expired`replay retained scaduti |
| 03:30 | `recordings:purge_temporary` — copie temp YouTube scadute |
| 08:00 | `recordings:expiry_warnings` — email avviso scadenza (7 gg) |
Log: `${MATCHLIVETV_VIDEOS_ROOT}/log/cron-replay.log` (es. `/media/videos/matchlivetv/log/cron-replay.log`)
### Capacità Garage
Allinstallazione, `setup_garage_production.sh` assegna ~**(disco libero 20 GB)** al nodo.
Per ridimensionare dopo:
### Audit duplicati storici
```bash
bash scripts/garage_set_capacity.sh 85G # oppure senza argomento: calcolo automatico
rails recordings:audit_youtube_duplicates # dry-run di default
```
## Piattaforme live e replay
Nessun cleanup automatico dei vecchi MP4 YouTube già in Garage.
| Origine diretta | Copia Garage/S3 | Player sito `/replay/:id` | YouTube VOD |
|-----------------|-----------------|----------------------------|-------------|
| `matchlivetv` | Sì (Premium) | MP4 via redirect `/media/` → Garage | Manuale (Premium Full) |
| `youtube` | Sì (Premium) | MP4 via redirect `/media/` → Garage | Auto se Premium Full + diretta YouTube |
## API
Se il file S3 non è disponibile ma esiste `youtube_video_id`, la pagina replay mostra embed YouTube.
- `GET /teams/:id/recordings` — include `replay_source`, `youtube_watch_url`, `playback_url` (null se youtube-only), `download_enabled` solo con `storage_key`
- `GET /recordings/:id/download` — richiede file Garage
- `POST /recordings/:id/publish_youtube` — solo MatchLiveTV-only con file (non per live già YouTube)
## Retention e abbonamento
## Colonne DB (migration)
- `expires_at` viene impostato alla fine diretta: **30 giorni** (Premium Light), **90 giorni** (Premium Full)
- `expires_at` **non cambia** al cambio piano successivo
- **Nuove registrazioni** richiedono abbonamento attivo (`can_create_recordings?`)
- **Archivio esistente** resta accessibile fino a `expires_at` anche se labbonamento scade (`can_access_recordings?`)
- `storage_policy` (default `retained`)
- `temp_expires_at`, `youtube_verified_at`, `local_media_purged_at`
## API mobile
## Rischi
- `GET /teams/:id/recordings` — lista con thumbnail, views, download_enabled
- `GET /recordings/:id/download` — URL download temporaneo
- `POST /recordings/:id/publish_youtube` — coda republicazione YouTube
## Estensioni future (metadata.ai)
Campo `metadata["ai"]` predisposto per:
- highlight automatici
- trascrizioni
- articoli generati
- clip / Shorts
Esempio aggiornamento via API:
```json
{ "recording": { "metadata": { "ai": { "transcript_status": "queued" } } } } }
```
- VOD YouTube in ritardo oltre retention temp → safety net elimina file; archivio resta se id collegato; altrimenti anomaly log.
- Thumb: se manca storage, poster YouTube da metadata / default.
- Record storici restano `storage_policy=retained`.
+6
View File
@@ -64,3 +64,9 @@ REPLAY_STORAGE_FORCE_PATH_STYLE=true
# Opzionale: durata link firmati per play/download (secondi)
# REPLAY_PRESIGNED_URL_TTL_SECONDS=7200
# REPLAY_DOWNLOAD_URL_TTL_SECONDS=900
# Replay YouTube: copia temporanea di sicurezza (ore, clamp 2472, default 48)
# YOUTUBE_TEMP_REPLAY_RETENTION_HOURS=48
# YOUTUBE_REPLAY_VERIFY_MAX_ATTEMPTS=12
# YOUTUBE_REPLAY_VERIFY_BASE_INTERVAL_SECS=300
# YOUTUBE_REPLAY_VERIFY_GRACE_SECS=120
+5
View File
@@ -78,6 +78,11 @@ REPLAY_STORAGE_FORCE_PATH_STYLE=true
# Streaming replay: Rails fa auth, edge serve /media/ → Garage (non passa da Puma)
REPLAY_MEDIA_PUBLIC_BASE_URL=https://www.matchlivetv.it/media
REPLAY_MEDIA_REDIRECT=true
# Copia temporanea post-live YouTube (ore, clamp 2472)
YOUTUBE_TEMP_REPLAY_RETENTION_HOURS=48
YOUTUBE_REPLAY_VERIFY_MAX_ATTEMPTS=12
YOUTUBE_REPLAY_VERIFY_BASE_INTERVAL_SECS=300
YOUTUBE_REPLAY_VERIFY_GRACE_SECS=120
# Puma: 5 thread (no secondo worker senza più RAM)
RAILS_MAX_THREADS=5
+1
View File
@@ -28,6 +28,7 @@ MARKER="# matchlivetv-replay-cron"
CRON_BLOCK="${MARKER}
30 * * * * mkdir -p ${LOG_DIR} && ${RUNNER} recordings:cleanup_local >> ${LOG_FILE} 2>&1
0 3 * * * mkdir -p ${LOG_DIR} && ${RUNNER} recordings:purge_expired >> ${LOG_FILE} 2>&1
30 3 * * * mkdir -p ${LOG_DIR} && ${RUNNER} recordings:purge_temporary >> ${LOG_FILE} 2>&1
0 8 * * * mkdir -p ${LOG_DIR} && ${RUNNER} recordings:expiry_warnings >> ${LOG_FILE} 2>&1
15 * * * * mkdir -p ${LOG_DIR} && ${RUNNER} analytics:aggregate >> ${LOG_DIR}/cron-analytics.log 2>&1
20 4 * * * mkdir -p ${LOG_DIR} && ${RUNNER} analytics:purge >> ${LOG_DIR}/cron-analytics.log 2>&1
+66 -7
View File
@@ -57,6 +57,7 @@ write_files:
STREAM_NODE_LOCAL_HLS_URL=http://127.0.0.1:8888
STREAM_NODE_RELAY_LOG_DIR=/var/log
STREAM_NODE_SLATES_DIR=/slates/custom
STREAM_NODE_RECORDINGS_DIR=/recordings
- path: /opt/stream-node/relay-agent.py
permissions: "0755"
content: |
@@ -68,17 +69,21 @@ write_files:
import json
import os
import shutil
import signal
import subprocess
import tarfile
import tempfile
import threading
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from urllib.parse import urlparse
from urllib.parse import unquote, urlparse
SECRET = os.environ.get("STREAM_NODE_AGENT_SECRET", "")
LISTEN = os.environ.get("STREAM_NODE_AGENT_LISTEN", "0.0.0.0:9100")
HLS_BASE = os.environ.get("STREAM_NODE_LOCAL_HLS_URL", "http://127.0.0.1:8888").rstrip("/")
LOG_DIR = os.environ.get("STREAM_NODE_RELAY_LOG_DIR", "/var/log")
SLATES_DIR = os.environ.get("STREAM_NODE_SLATES_DIR", "/slates/custom")
RECORDINGS_DIR = os.environ.get("STREAM_NODE_RECORDINGS_DIR", "/recordings")
_lock = threading.Lock()
# session_id -> {"pid": int, "path": str, "proc": Popen}
@@ -160,6 +165,27 @@ write_files:
pass
def _recordings_root(path_name):
name = unquote(path_name or "").strip().lstrip("/")
if not name or ".." in name.split("/"):
return None
base = os.path.realpath(RECORDINGS_DIR)
root = os.path.realpath(os.path.join(base, name))
if root != base and not root.startswith(base + os.sep):
return None
return root
def _has_recording_files(root: str) -> bool:
if not os.path.isdir(root):
return False
for dirpath, _dirnames, filenames in os.walk(root):
for name in filenames:
if name.lower().endswith((".mp4", ".fmp4", ".m4s", ".ts")):
return True
return False
def _start_session(session_id: str, path: str, rtmps: str) -> int:
existing = _relays.get(session_id)
if existing:
@@ -214,6 +240,30 @@ write_files:
return
self._json(200, {"running": True, "session_id": session_id, "pid": info["pid"], "path": info["path"]})
return
if parsed.path.startswith("/recordings/"):
name = parsed.path.split("/recordings/", 1)[1].strip("/")
root = _recordings_root(name)
if not root or not _has_recording_files(root):
self._json(404, {"error": "not found"})
return
tmp = tempfile.NamedTemporaryFile(prefix="mltv-rec-", suffix=".tar.gz", delete=False)
tmp.close()
try:
with tarfile.open(tmp.name, "w:gz") as tar:
tar.add(root, arcname=".")
size = os.path.getsize(tmp.name)
self.send_response(200)
self.send_header("Content-Type", "application/gzip")
self.send_header("Content-Length", str(size))
self.end_headers()
with open(tmp.name, "rb") as fh:
shutil.copyfileobj(fh, self.wfile)
finally:
try:
os.unlink(tmp.name)
except OSError:
pass
return
self._json(404, {"error": "not found"})
def do_POST(self) -> None:
@@ -244,13 +294,22 @@ write_files:
self._json(401, {"error": "unauthorized"})
return
parsed = urlparse(self.path)
if not parsed.path.startswith("/relays/"):
self._json(404, {"error": "not found"})
if parsed.path.startswith("/relays/"):
session_id = parsed.path.split("/relays/", 1)[1].strip("/")
with _lock:
_stop_session(session_id)
self._json(200, {"stopped": True, "session_id": session_id})
return
session_id = parsed.path.split("/relays/", 1)[1].strip("/")
with _lock:
_stop_session(session_id)
self._json(200, {"stopped": True, "session_id": session_id})
if parsed.path.startswith("/recordings/"):
name = parsed.path.split("/recordings/", 1)[1].strip("/")
root = _recordings_root(name)
if not root or not os.path.isdir(root):
self._json(404, {"error": "not found"})
return
shutil.rmtree(root, ignore_errors=True)
self._json(200, {"deleted": True, "path": name})
return
self._json(404, {"error": "not found"})
def do_PUT(self) -> None:
if not _authorized(self):
+65 -7
View File
@@ -6,17 +6,21 @@ from __future__ import annotations
import json
import os
import shutil
import signal
import subprocess
import tarfile
import tempfile
import threading
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from urllib.parse import urlparse
from urllib.parse import unquote, urlparse
SECRET = os.environ.get("STREAM_NODE_AGENT_SECRET", "")
LISTEN = os.environ.get("STREAM_NODE_AGENT_LISTEN", "0.0.0.0:9100")
HLS_BASE = os.environ.get("STREAM_NODE_LOCAL_HLS_URL", "http://127.0.0.1:8888").rstrip("/")
LOG_DIR = os.environ.get("STREAM_NODE_RELAY_LOG_DIR", "/var/log")
SLATES_DIR = os.environ.get("STREAM_NODE_SLATES_DIR", "/slates/custom")
RECORDINGS_DIR = os.environ.get("STREAM_NODE_RECORDINGS_DIR", "/recordings")
_lock = threading.Lock()
# session_id -> {"pid": int, "path": str, "proc": Popen}
@@ -98,6 +102,27 @@ def _stop_session(session_id: str) -> None:
pass
def _recordings_root(path_name):
name = unquote(path_name or "").strip().lstrip("/")
if not name or ".." in name.split("/"):
return None
base = os.path.realpath(RECORDINGS_DIR)
root = os.path.realpath(os.path.join(base, name))
if root != base and not root.startswith(base + os.sep):
return None
return root
def _has_recording_files(root: str) -> bool:
if not os.path.isdir(root):
return False
for dirpath, _dirnames, filenames in os.walk(root):
for name in filenames:
if name.lower().endswith((".mp4", ".fmp4", ".m4s", ".ts")):
return True
return False
def _start_session(session_id: str, path: str, rtmps: str) -> int:
existing = _relays.get(session_id)
if existing:
@@ -152,6 +177,30 @@ class Handler(BaseHTTPRequestHandler):
return
self._json(200, {"running": True, "session_id": session_id, "pid": info["pid"], "path": info["path"]})
return
if parsed.path.startswith("/recordings/"):
name = parsed.path.split("/recordings/", 1)[1].strip("/")
root = _recordings_root(name)
if not root or not _has_recording_files(root):
self._json(404, {"error": "not found"})
return
tmp = tempfile.NamedTemporaryFile(prefix="mltv-rec-", suffix=".tar.gz", delete=False)
tmp.close()
try:
with tarfile.open(tmp.name, "w:gz") as tar:
tar.add(root, arcname=".")
size = os.path.getsize(tmp.name)
self.send_response(200)
self.send_header("Content-Type", "application/gzip")
self.send_header("Content-Length", str(size))
self.end_headers()
with open(tmp.name, "rb") as fh:
shutil.copyfileobj(fh, self.wfile)
finally:
try:
os.unlink(tmp.name)
except OSError:
pass
return
self._json(404, {"error": "not found"})
def do_POST(self) -> None:
@@ -182,13 +231,22 @@ class Handler(BaseHTTPRequestHandler):
self._json(401, {"error": "unauthorized"})
return
parsed = urlparse(self.path)
if not parsed.path.startswith("/relays/"):
self._json(404, {"error": "not found"})
if parsed.path.startswith("/relays/"):
session_id = parsed.path.split("/relays/", 1)[1].strip("/")
with _lock:
_stop_session(session_id)
self._json(200, {"stopped": True, "session_id": session_id})
return
session_id = parsed.path.split("/relays/", 1)[1].strip("/")
with _lock:
_stop_session(session_id)
self._json(200, {"stopped": True, "session_id": session_id})
if parsed.path.startswith("/recordings/"):
name = parsed.path.split("/recordings/", 1)[1].strip("/")
root = _recordings_root(name)
if not root or not os.path.isdir(root):
self._json(404, {"error": "not found"})
return
shutil.rmtree(root, ignore_errors=True)
self._json(200, {"deleted": True, "path": name})
return
self._json(404, {"error": "not found"})
def do_PUT(self) -> None:
if not _authorized(self):
@@ -43,7 +43,11 @@ class E2ECoverUploadTest {
openMatchWizardViaQuickStart()
waitForText(s(R.string.wizard_step_title_match), timeoutMs = 45_000)
scrollUntilVisible(s(R.string.wizard_match_cover_title))
scrollUntilVisible(
s(R.string.wizard_match_cover_title),
"Copertina in pausa",
"Pause cover",
)
tapClickableText(s(R.string.wizard_match_cover_change))
pickFirstPhotoFromPicker()
@@ -185,12 +189,12 @@ class E2ECoverUploadTest {
}
}
private fun scrollUntilVisible(text: String) {
repeat(8) {
if (device.hasObject(By.text(text))) return
private fun scrollUntilVisible(vararg texts: String) {
repeat(20) {
if (texts.any { device.hasObject(By.text(it)) }) return
scrollDown(1)
}
waitForText(text, timeoutMs = 10_000)
waitForAnyText(*texts, timeoutMs = 10_000)
}
private fun grantRuntimePermissions() {
@@ -22,8 +22,8 @@ import org.junit.runner.RunWith
@RunWith(AndroidJUnit4::class)
class E2EWizardFlowTest {
private lateinit var device: UiDevice
private val pkg = "com.matchlivetv.match_live_tv"
private val ctx by lazy { InstrumentationRegistry.getInstrumentation().targetContext }
private val pkg by lazy { ctx.packageName }
private fun s(id: Int): String = ctx.getString(id)
private fun su(id: Int): String = s(id).uppercase()
@@ -20,7 +20,7 @@ import org.junit.runner.RunWith
@RunWith(AndroidJUnit4::class)
class LocalAdaptiveBitrateUiTest {
private lateinit var device: UiDevice
private val pkg = "com.matchlivetv.match_live_tv"
private val pkg by lazy { InstrumentationRegistry.getInstrumentation().targetContext.packageName }
@Before
fun setUp() {
@@ -47,8 +47,21 @@ class LocalAdaptiveBitrateUiTest {
tapAny("AVANTI >", "NEXT >")
waitForAny(45_000, "02 · Trasmissione", "02 · Broadcast")
waitForAny(30_000, "Piattaforma", "Platform")
scrollDown()
tapAny("AVANTI >", "NEXT >")
// Come E2EWizardFlowTest: privacy non-in-elenco e AVANTI riprovato con scroll.
waitForAny(10_000, "NON IN ELENCO", "UNLISTED")
runCatching { tapAny("NON IN ELENCO", "UNLISTED") }
val onNetworkStep = {
hasAny("03 · Test rete", "03 · Network test", "AVVIA TEST RETE", "START NETWORK TEST")
}
repeat(4) {
if (onNetworkStep()) return@repeat
scrollDown()
runCatching { tapAny("AVANTI >", "NEXT >") }
SystemClock.sleep(1_200)
if (onNetworkStep()) return@repeat
scrollUp()
SystemClock.sleep(800)
}
waitForAny(45_000, "03 · Test rete", "03 · Network test")
waitForAny(30_000, "AVVIA TEST RETE", "START NETWORK TEST")
tapAny("AVVIA TEST RETE", "START NETWORK TEST")
@@ -262,4 +275,15 @@ class LocalAdaptiveBitrateUiTest {
SystemClock.sleep(300)
}
}
private fun scrollUp(steps: Int = 1) {
val centerX = device.displayWidth / 2
val startY = (device.displayHeight * 0.35).toInt()
val endY = (device.displayHeight * 0.75).toInt()
repeat(steps) {
device.swipe(centerX, startY, centerX, endY, 24)
device.waitForIdle()
SystemClock.sleep(300)
}
}
}
@@ -28,7 +28,7 @@ class ReleaseApiSmokeTest {
fun login_parsesResponse() = runBlocking {
val session = container.authRepository.login(
email = "coach@matchlivetv.test",
password = "password123",
password = "Password123",
)
assertEquals("coach@matchlivetv.test", session.user.email)
assertTrue(session.accessToken.isNotBlank())
@@ -38,7 +38,7 @@ class ReleaseApiSmokeTest {
fun fetchMatches_afterLogin() = runBlocking {
container.authRepository.login(
email = "coach@matchlivetv.test",
password = "password123",
password = "Password123",
)
val matches = container.matchRepository.fetchMatches()
assertTrue(matches.isNotEmpty())
@@ -48,15 +48,22 @@ class ReleaseApiSmokeTest {
fun scheduledMatch_parsesAndIsVisible() = runBlocking {
container.authRepository.login(
email = "coach@matchlivetv.test",
password = "password123",
password = "Password123",
)
val teams = container.matchRepository.fetchTeams()
val tigers = teams.first { it.name == "Tigers Volley" }
val raw = container.api.matches(tigers.id)
val scheduled = raw.first { it.opponentName.contains("Crazy Volley") }
assertNotNull(scheduled.scheduledAt)
assertNotNull(parseApiInstant(scheduled.scheduledAt))
val domain = scheduled.toDomain()
assertTrue(teams.isNotEmpty())
var scheduled: com.matchlivetv.match_live_tv.data.api.MatchDto? = null
for (team in teams) {
val found = container.api.matches(team.id).firstOrNull { !it.scheduledAt.isNullOrBlank() }
if (found != null) {
scheduled = found
break
}
}
val match = checkNotNull(scheduled) { "Nessuna partita con scheduled_at tra i team del coach" }
assertNotNull(match.scheduledAt)
assertNotNull(parseApiInstant(match.scheduledAt))
val domain = match.toDomain()
assertTrue(domain.isCoachHubVisible())
}
}
@@ -3,9 +3,13 @@ package com.matchlivetv.match_live_tv.core
import android.content.Context
import android.content.Intent
import android.content.IntentFilter
import android.content.pm.PackageManager
import android.net.ConnectivityManager
import android.net.NetworkCapabilities
import android.os.BatteryManager
import android.os.Build
import android.telephony.TelephonyManager
import com.matchlivetv.match_live_tv.data.api.ClientInfoPayload
data class DeviceHealthSnapshot(
val batteryPercent: Int,
@@ -27,6 +31,38 @@ object DeviceTelemetry {
}
}.getOrDefault("Sconosciuto")
fun clientInfo(context: Context): ClientInfoPayload {
val packageInfo = runCatching {
if (Build.VERSION.SDK_INT >= 33) {
context.packageManager.getPackageInfo(
context.packageName,
PackageManager.PackageInfoFlags.of(0),
)
} else {
@Suppress("DEPRECATION")
context.packageManager.getPackageInfo(context.packageName, 0)
}
}.getOrNull()
val versionName = packageInfo?.versionName
val versionCode = packageInfo?.let {
if (Build.VERSION.SDK_INT >= 28) it.longVersionCode.toString() else {
@Suppress("DEPRECATION")
it.versionCode.toString()
}
}
return ClientInfoPayload(
os = "android",
appVersion = versionName,
appBuild = versionCode,
deviceManufacturer = Build.MANUFACTURER?.takeIf { it.isNotBlank() },
deviceModel = Build.MODEL?.takeIf { it.isNotBlank() },
osVersion = Build.VERSION.RELEASE,
carrier = carrierName(context),
)
}
fun snapshot(context: Context, thermalState: ThermalState? = null): DeviceHealthSnapshot {
val batteryIntent = context.registerReceiver(null, IntentFilter(Intent.ACTION_BATTERY_CHANGED))
val batteryPercent = readBatteryPercent(batteryIntent)
@@ -37,6 +73,13 @@ object DeviceTelemetry {
)
}
private fun carrierName(context: Context): String? = runCatching {
val tm = context.getSystemService(Context.TELEPHONY_SERVICE) as? TelephonyManager ?: return null
sequenceOf(tm.networkOperatorName, tm.simOperatorName)
.mapNotNull { it?.trim()?.takeIf { name -> name.isNotEmpty() } }
.firstOrNull()
}.getOrNull()
private fun readBatteryPercent(intent: Intent?): Int {
if (intent == null) return 100
val level = intent.getIntExtra(BatteryManager.EXTRA_LEVEL, -1)
@@ -96,7 +96,7 @@ class AppContainer(context: Context) {
.filter { it.id !in dismissed }
}
val sessionRepository = SessionRepository(api)
val sessionRepository = SessionRepository(api, appContext)
val scoreRepository = ScoreRepository(api)
@@ -27,7 +27,10 @@ data class ChangePasswordRequest(
data class MessageResponse(val message: String)
data class ApiErrorResponse(val error: String? = null)
data class ApiErrorResponse(
val error: String? = null,
@Json(name = "error_code") val errorCode: String? = null,
)
data class LoginResponse(
val user: UserDto,
@@ -295,6 +298,17 @@ data class CreateSessionRequest(
@Json(name = "target_bitrate") val targetBitrate: Int = 2_500_000,
@Json(name = "target_fps") val targetFps: Int = 30,
@Json(name = "youtube_channel") val youtubeChannel: String? = null,
val client: ClientInfoPayload? = null,
)
data class ClientInfoPayload(
val os: String,
@Json(name = "app_version") val appVersion: String? = null,
@Json(name = "app_build") val appBuild: String? = null,
@Json(name = "device_manufacturer") val deviceManufacturer: String? = null,
@Json(name = "device_model") val deviceModel: String? = null,
@Json(name = "os_version") val osVersion: String? = null,
val carrier: String? = null,
)
data class MinQualityRequest(
@@ -413,6 +427,7 @@ data class TelemetryRequest(
@Json(name = "target_bitrate") val targetBitrate: Int? = null,
val fps: Int? = null,
@Json(name = "thermal_state") val thermalState: String? = null,
val client: ClientInfoPayload? = null,
)
data class AnnouncementDto(
@@ -1,28 +1,64 @@
package com.matchlivetv.match_live_tv.data.repository
import android.content.Context
import com.matchlivetv.match_live_tv.core.DeviceTelemetry
import com.matchlivetv.match_live_tv.data.api.ApiErrorResponse
import com.matchlivetv.match_live_tv.data.api.CreateSessionRequest
import com.matchlivetv.match_live_tv.data.api.MatchLiveApi
import com.matchlivetv.match_live_tv.data.api.MinQualityRequest
import com.matchlivetv.match_live_tv.data.api.NetworkTestRequest
import com.matchlivetv.match_live_tv.data.api.NetworkTestResponse
import com.matchlivetv.match_live_tv.domain.StreamSession
import com.squareup.moshi.Moshi
import com.squareup.moshi.kotlin.reflect.KotlinJsonAdapterFactory
import kotlinx.coroutines.delay
import retrofit2.HttpException
class SessionRepository(
private val api: MatchLiveApi,
private val appContext: Context,
) {
private val errorAdapter by lazy {
Moshi.Builder()
.add(KotlinJsonAdapterFactory())
.build()
.adapter(ApiErrorResponse::class.java)
}
suspend fun createSession(
matchId: String,
platform: String = "matchlivetv",
privacyStatus: String = "public",
youtubeChannel: String? = null,
): StreamSession = api.createSession(
matchId,
CreateSessionRequest(
platform = platform,
privacyStatus = privacyStatus,
youtubeChannel = youtubeChannel,
),
).toDomain()
scalingRetries: Int = 8,
scalingRetryDelayMs: Long = 15_000L,
): StreamSession {
var attempt = 0
while (true) {
try {
return api.createSession(
matchId,
CreateSessionRequest(
platform = platform,
privacyStatus = privacyStatus,
youtubeChannel = youtubeChannel,
client = DeviceTelemetry.clientInfo(appContext),
),
).toDomain()
} catch (e: HttpException) {
val code = parseErrorCode(e)
if (code != "stream_capacity_scaling" || attempt >= scalingRetries) throw e
attempt += 1
delay(scalingRetryDelayMs)
}
}
}
private fun parseErrorCode(error: HttpException): String? {
val body = error.response()?.errorBody()?.string().orEmpty()
if (body.isBlank()) return null
return runCatching { errorAdapter.fromJson(body)?.errorCode }.getOrNull()
}
suspend fun fetchSession(sessionId: String): StreamSession =
api.session(sessionId).toDomain()
@@ -81,6 +117,7 @@ class SessionRepository(
targetBitrate = targetBitrate,
fps = fps,
thermalState = thermalState,
client = DeviceTelemetry.clientInfo(appContext),
),
)
}
@@ -1,6 +1,7 @@
import Foundation
import UIKit
import Network
import CoreTelephony
struct DeviceHealth: Sendable {
let batteryPercent: Int
@@ -8,6 +9,16 @@ struct DeviceHealth: Sendable {
let networkType: String
}
struct ClientInfoPayload: Encodable, Sendable {
let os: String
let appVersion: String?
let appBuild: String?
let deviceManufacturer: String?
let deviceModel: String?
let osVersion: String?
let carrier: String?
}
enum DeviceTelemetry {
static func snapshot(thermalState: ThermalState? = nil) -> DeviceHealth {
UIDevice.current.isBatteryMonitoringEnabled = true
@@ -20,6 +31,24 @@ enum DeviceTelemetry {
)
}
static func clientInfo() -> ClientInfoPayload {
let bundle = Bundle.main
let version = bundle.object(forInfoDictionaryKey: "CFBundleShortVersionString") as? String
let build = bundle.object(forInfoDictionaryKey: "CFBundleVersion") as? String
let model = UIDevice.current.model
// Prefer machine identifier when available (e.g. iPhone15,2)
let machine = utsnameMachine()
return ClientInfoPayload(
os: "ios",
appVersion: version,
appBuild: build,
deviceManufacturer: "Apple",
deviceModel: machine ?? model,
osVersion: UIDevice.current.systemVersion,
carrier: carrierName()
)
}
private static func currentNetworkType() -> String {
let monitor = NWPathMonitor()
let semaphore = DispatchSemaphore(value: 0)
@@ -40,4 +69,27 @@ enum DeviceTelemetry {
monitor.cancel()
return result
}
private static func carrierName() -> String? {
let info = CTTelephonyNetworkInfo()
if let providers = info.serviceSubscriberCellularProviders {
for carrier in providers.values {
if let name = carrier.carrierName?.trimmingCharacters(in: .whitespacesAndNewlines),
!name.isEmpty {
return name
}
}
}
return nil
}
private static func utsnameMachine() -> String? {
var systemInfo = utsname()
uname(&systemInfo)
return withUnsafePointer(to: &systemInfo.machine) {
$0.withMemoryRebound(to: CChar.self, capacity: 1) {
String(validatingUTF8: $0)
}
}
}
}
@@ -435,6 +435,7 @@ struct CreateSessionRequest: Encodable {
let targetBitrate: Int
let targetFps: Int
let youtubeChannel: String?
let client: ClientInfoPayload?
}
struct AudioMuteRequest: Encodable {
@@ -531,6 +532,7 @@ struct TelemetryRequest: Encodable {
let targetBitrate: Int?
let fps: Int?
let thermalState: String?
let client: ClientInfoPayload?
}
extension ScoringRules {
@@ -25,7 +25,8 @@ final class SessionRepository {
qualityPreset: qualityPreset,
targetBitrate: targetBitrate,
targetFps: targetFps,
youtubeChannel: youtubeChannel
youtubeChannel: youtubeChannel,
client: DeviceTelemetry.clientInfo()
)
).toDomain()
}
@@ -97,7 +98,8 @@ final class SessionRepository {
currentBitrate: currentBitrate,
targetBitrate: targetBitrate,
fps: fps,
thermalState: health.thermalState.apiValue
thermalState: health.thermalState.apiValue,
client: DeviceTelemetry.clientInfo()
)
)
}
+2
View File
@@ -23,6 +23,8 @@ STREAM_NODE_AGENT_SECRET=${SECRET}
STREAM_NODE_AGENT_LISTEN=0.0.0.0:9100
STREAM_NODE_LOCAL_HLS_URL=http://127.0.0.1:8888
STREAM_NODE_RELAY_LOG_DIR=/var/log
STREAM_NODE_SLATES_DIR=/slates/custom
STREAM_NODE_RECORDINGS_DIR=/recordings
EOF
chmod 600 /opt/stream-node/agent.env"
"${SSH[@]}" "cat >/etc/systemd/system/mltv-relay-agent.service <<'EOF'