488 lines
16 KiB
Ruby
488 lines
16 KiB
Ruby
require "digest"
|
|
require "fileutils"
|
|
require "ipaddr"
|
|
require "net/http"
|
|
require "uri"
|
|
|
|
module MeowClient
|
|
module Media
|
|
PlayRequest = Struct.new(
|
|
:name, :url, :kind, :tag, :volume, :fade_in, :fade_out, :start_at,
|
|
:finish_at, :loops, :priority, :continue_existing, :key, :caption,
|
|
keyword_init: true
|
|
)
|
|
StopRequest = Struct.new(
|
|
:name, :kind, :tag, :priority, :key, :fade_away, :fade_out,
|
|
keyword_init: true
|
|
)
|
|
PreloadRequest = Struct.new(:name, :url, keyword_init: true)
|
|
DownloadResult = Struct.new(:id, :path, :url, :error, keyword_init: true)
|
|
Track = Struct.new(:request, :handle, :url, :remaining_loops, :started_at, keyword_init: true)
|
|
|
|
class URLPolicy
|
|
AUDIO_EXTENSIONS = %w[.mp3 .ogg .opus .wav .flac .m4a .aac].freeze
|
|
|
|
def resolve(base_url, name, explicit_url = nil)
|
|
clean_name = validate_name(name)
|
|
base = explicit_url.to_s.strip == "" ? base_url : explicit_url
|
|
raise ArgumentError, "No media base URL was supplied" if base.to_s.strip == ""
|
|
base_uri = validate_uri(URI.parse(base.to_s))
|
|
value = clean_name == "" ? base_uri : URI.join(ensure_directory(base_uri.to_s), clean_name)
|
|
validate_uri(value)
|
|
rescue URI::InvalidURIError => error
|
|
raise ArgumentError, "Invalid media URL: #{error.message}"
|
|
end
|
|
|
|
def validate_uri(uri)
|
|
raise ArgumentError, "Only HTTPS media URLs are allowed" unless uri.scheme.to_s.downcase == "https"
|
|
raise ArgumentError, "Media URL must include a host" if uri.host.to_s == ""
|
|
raise ArgumentError, "Media URL credentials are not allowed" if uri.userinfo != nil
|
|
raise ArgumentError, "Media URL contains control characters" if uri.to_s.match?(/[\x00-\x1f\x7f]/)
|
|
reject_private_literal(uri.host)
|
|
uri
|
|
end
|
|
|
|
def validate_name(name)
|
|
value = name.to_s.tr("\\", "/")
|
|
raise ArgumentError, "Media name contains control characters" if value.match?(/[\x00-\x1f\x7f]/)
|
|
raise ArgumentError, "Media name must be relative" if value.start_with?("/", "//") || value.match?(/\A[A-Za-z]:/)
|
|
segments = value.split("/")
|
|
raise ArgumentError, "Media name may not traverse parent directories" if segments.include?("..")
|
|
value
|
|
end
|
|
|
|
def supported_audio?(uri)
|
|
AUDIO_EXTENSIONS.include?(File.extname(uri.path.to_s).downcase)
|
|
end
|
|
|
|
private
|
|
|
|
def ensure_directory(value)
|
|
value.end_with?("/") ? value : value + "/"
|
|
end
|
|
|
|
def reject_private_literal(host)
|
|
ip = IPAddr.new(host)
|
|
private_ip = ip.loopback? || ip.link_local? || ip.private? || ip.multicast? || ip.unspecified?
|
|
raise ArgumentError, "Private or reserved media addresses are not allowed" if private_ip
|
|
rescue IPAddr::InvalidAddressError
|
|
end
|
|
end
|
|
|
|
class Downloader
|
|
MAX_FILE_BYTES = 25 * 1_048_576
|
|
MAX_CACHE_BYTES = 100 * 1_048_576
|
|
MAX_QUEUE = 32
|
|
MAX_REDIRECTS = 5
|
|
OPEN_TIMEOUT = 10
|
|
READ_TIMEOUT = 15
|
|
|
|
attr_reader :results
|
|
|
|
def initialize(cache_directory, policy = URLPolicy.new)
|
|
@cache_directory = cache_directory.to_s
|
|
@policy = policy
|
|
@queue = SizedQueue.new(MAX_QUEUE)
|
|
@results = Queue.new
|
|
@closed = false
|
|
FileUtils.mkdir_p(@cache_directory)
|
|
@worker = Thread.new { work }
|
|
@worker.report_on_exception = false
|
|
end
|
|
|
|
def enqueue(id, uri)
|
|
return false if @closed
|
|
cached = cache_path(uri)
|
|
if File.file?(cached)
|
|
File.utime(Time.now, Time.now, cached) rescue nil
|
|
@results << DownloadResult.new(:id => id, :path => cached, :url => uri.to_s)
|
|
return true
|
|
end
|
|
@queue.push([id, uri], true)
|
|
true
|
|
rescue ThreadError
|
|
false
|
|
end
|
|
|
|
def close
|
|
return if @closed
|
|
@closed = true
|
|
@queue.push(nil, true) rescue nil
|
|
@worker.join(1.0) if @worker != nil && @worker != Thread.current
|
|
@worker.kill if @worker != nil && @worker.alive?
|
|
nil
|
|
end
|
|
|
|
private
|
|
|
|
def work
|
|
loop do
|
|
job = @queue.pop
|
|
break if job == nil || @closed
|
|
id, uri = job
|
|
begin
|
|
path, final_uri = download(uri)
|
|
@results << DownloadResult.new(:id => id, :path => path, :url => final_uri.to_s)
|
|
rescue Exception => error
|
|
@results << DownloadResult.new(:id => id, :url => uri.to_s, :error => error.message)
|
|
end
|
|
end
|
|
end
|
|
|
|
def download(uri, redirects = 0)
|
|
raise "Too many media redirects" if redirects > MAX_REDIRECTS
|
|
current = @policy.validate_uri(uri)
|
|
response = nil
|
|
Net::HTTP.start(
|
|
current.host,
|
|
current.port,
|
|
:use_ssl => true,
|
|
:open_timeout => OPEN_TIMEOUT,
|
|
:read_timeout => READ_TIMEOUT
|
|
) do |http|
|
|
request = Net::HTTP::Get.new(current.request_uri)
|
|
http.request(request) do |incoming|
|
|
response = incoming
|
|
if incoming.is_a?(Net::HTTPRedirection)
|
|
location = incoming["location"]
|
|
raise "Media redirect did not include a location" if location.to_s == ""
|
|
redirected = @policy.validate_uri(URI.join(current.to_s, location))
|
|
return download(redirected, redirects + 1)
|
|
end
|
|
raise "Media download failed with HTTP #{incoming.code}" unless incoming.is_a?(Net::HTTPSuccess)
|
|
length = Integer(incoming["content-length"], :exception => false)
|
|
raise "Media file exceeds #{MAX_FILE_BYTES} bytes" if length != nil && length > MAX_FILE_BYTES
|
|
target = cache_path(current)
|
|
temporary = target + ".part-#{Thread.current.object_id}"
|
|
total = 0
|
|
begin
|
|
File.open(temporary, "wb") do |file|
|
|
incoming.read_body do |chunk|
|
|
total += chunk.bytesize
|
|
raise "Media file exceeds #{MAX_FILE_BYTES} bytes" if total > MAX_FILE_BYTES
|
|
file.write(chunk)
|
|
end
|
|
end
|
|
if File.file?(target)
|
|
File.delete(temporary)
|
|
else
|
|
File.rename(temporary, target)
|
|
end
|
|
ensure
|
|
File.delete(temporary) if File.exist?(temporary)
|
|
end
|
|
prune_cache(target)
|
|
return [target, current]
|
|
end
|
|
end
|
|
raise "Media download failed" if response == nil
|
|
end
|
|
|
|
def cache_path(uri)
|
|
extension = File.extname(uri.path.to_s).downcase
|
|
extension = ".media" unless URLPolicy::AUDIO_EXTENSIONS.include?(extension)
|
|
File.join(@cache_directory, Digest::SHA256.hexdigest(uri.to_s) + extension)
|
|
end
|
|
|
|
def prune_cache(protected_path)
|
|
rows = Dir.glob(File.join(@cache_directory, "*")).filter_map do |file|
|
|
next unless File.file?(file) && !File.basename(file).include?(".part-")
|
|
stat = File.stat(file)
|
|
[file, stat.size, stat.mtime]
|
|
rescue SystemCallError
|
|
nil
|
|
end
|
|
total = rows.sum { |row| row[1] }
|
|
rows.sort_by { |row| row[2] }.each do |file, size, _time|
|
|
break if total <= MAX_CACHE_BYTES
|
|
next if file == protected_path
|
|
File.delete(file) rescue next
|
|
total -= size
|
|
end
|
|
end
|
|
end
|
|
|
|
class EltenSoundHandle
|
|
def initialize(path, request, category_volume)
|
|
@sound = Sound.new(path, :loop => request.loops.to_i == -1)
|
|
@target_volume = [[request.volume.to_i, 0].max, 100].min / 100.0 * category_volume
|
|
@sound.position = request.start_at.to_f / 1000.0 if request.start_at.to_i > 0
|
|
volume = @sound.attribute(:volume)
|
|
volume.value = request.fade_in.to_i > 0 ? 0.0 : @target_volume
|
|
@sound.play
|
|
volume.slide(@target_volume, :duration => request.fade_in.to_f / 1000.0) if request.fade_in.to_i > 0
|
|
end
|
|
|
|
def finished?
|
|
@sound.finished?
|
|
end
|
|
|
|
def position_ms
|
|
(@sound.position.to_f * 1000).to_i
|
|
end
|
|
|
|
def restart(position_ms)
|
|
@sound.position = position_ms.to_f / 1000.0
|
|
@sound.play
|
|
end
|
|
|
|
def stop(fade_ms = 0)
|
|
if fade_ms.to_i > 0 && !finished?
|
|
@sound.attribute(:volume).slide(0.0, :duration => fade_ms.to_f / 1000.0)
|
|
Thread.new do
|
|
sleep(fade_ms.to_f / 1000.0)
|
|
close
|
|
end
|
|
else
|
|
close
|
|
end
|
|
end
|
|
|
|
def close
|
|
@sound.stop rescue nil
|
|
@sound.close rescue nil
|
|
end
|
|
end
|
|
|
|
class EltenBackend
|
|
def start(path, request, category_volume)
|
|
EltenSoundHandle.new(path, request, category_volume)
|
|
end
|
|
end
|
|
|
|
class Manager
|
|
MAX_TRACKS = 16
|
|
|
|
attr_reader :captions, :errors, :tracks
|
|
attr_reader :muted
|
|
|
|
def initialize(cache_directory, sound_volume: 1.0, music_volume: 1.0, backend: EltenBackend.new, downloader: nil)
|
|
@policy = URLPolicy.new
|
|
@downloader = downloader || Downloader.new(cache_directory, @policy)
|
|
@backend = backend
|
|
@sound_volume = sound_volume.to_f
|
|
@music_volume = music_volume.to_f
|
|
@tracks = []
|
|
@pending = {}
|
|
@sequence = 0
|
|
@captions = []
|
|
@pending_captions = []
|
|
@errors = []
|
|
@default_url = nil
|
|
@loaded_urls = {}
|
|
@muted = false
|
|
end
|
|
|
|
def default_url=(value)
|
|
@default_url = @policy.validate_uri(URI.parse(value.to_s)).to_s
|
|
rescue ArgumentError, URI::InvalidURIError => error
|
|
add_error(error.message)
|
|
end
|
|
|
|
def preload(request)
|
|
queue(request, false)
|
|
end
|
|
|
|
def play(request)
|
|
request.kind = normalize_kind(request.kind)
|
|
request.volume = 50 if request.volume == nil
|
|
request.loops = 1 if request.loops == nil || request.loops.to_i == 0
|
|
request.priority = 50 if request.priority == nil
|
|
remember_caption(request.caption)
|
|
return nil if @muted
|
|
existing = matching_identity(request)
|
|
return existing if request.continue_existing != false && existing != nil
|
|
stop(StopRequest.new(:key => request.key)) if request.key.to_s != ""
|
|
queue(request, true)
|
|
rescue Exception => error
|
|
add_error(error.message)
|
|
nil
|
|
end
|
|
|
|
def muted=(value)
|
|
@muted = value == true
|
|
stop if @muted
|
|
@muted
|
|
end
|
|
|
|
def stop(request = StopRequest.new)
|
|
pending_ids = @pending.filter_map do |id, (pending_request, play_after)|
|
|
id if play_after && stop_match?(pending_request, request)
|
|
end
|
|
pending_ids.each { |id| @pending.delete(id) }
|
|
|
|
selected = @tracks.select { |track| stop_match?(track.request, request) }
|
|
selected.each do |track|
|
|
fade = request.fade_away ? (track.request.fade_out || request.fade_out || 5_000) : (request.fade_out || 0)
|
|
track.handle.stop(fade)
|
|
@tracks.delete(track)
|
|
end
|
|
selected.size + pending_ids.size
|
|
end
|
|
|
|
def drain_captions
|
|
values = @pending_captions
|
|
@pending_captions = []
|
|
values
|
|
end
|
|
|
|
def tick
|
|
drain_downloads
|
|
@tracks.dup.each do |track|
|
|
request = track.request
|
|
if request.finish_at.to_i > 0 && track.handle.position_ms >= request.finish_at.to_i
|
|
finish_iteration(track)
|
|
elsif track.handle.finished?
|
|
finish_iteration(track)
|
|
end
|
|
rescue Exception => error
|
|
add_error(error.message)
|
|
remove_track(track)
|
|
end
|
|
end
|
|
|
|
def close
|
|
@pending.clear
|
|
@tracks.dup.each { |track| remove_track(track) }
|
|
@downloader.close
|
|
nil
|
|
end
|
|
|
|
private
|
|
|
|
def queue(request, play_after)
|
|
loaded = @loaded_urls[request.name.to_s] if request.url.to_s == ""
|
|
if loaded != nil
|
|
uri = @policy.validate_uri(URI.parse(loaded))
|
|
else
|
|
base = request.url.to_s == "" ? @default_url : request.url
|
|
uri = @policy.resolve(base, request.name, request.url)
|
|
end
|
|
raise "Unsupported audio file type" unless @policy.supported_audio?(uri)
|
|
@sequence += 1
|
|
@pending[@sequence] = [request, play_after]
|
|
unless @downloader.enqueue(@sequence, uri)
|
|
@pending.delete(@sequence)
|
|
raise "Media download queue is full"
|
|
end
|
|
@loaded_urls[request.name.to_s] = uri.to_s
|
|
@sequence
|
|
end
|
|
|
|
def drain_downloads
|
|
loop do
|
|
result = @downloader.results.pop(true)
|
|
request, play_after = @pending.delete(result.id)
|
|
next if request == nil
|
|
if result.error != nil
|
|
add_error(result.error)
|
|
elsif play_after
|
|
start_track(request, result.path, result.url)
|
|
end
|
|
end
|
|
rescue ThreadError
|
|
end
|
|
|
|
def start_track(request, path, url)
|
|
priority = request.priority.to_i
|
|
return if @tracks.any? { |track| track.request.priority.to_i > priority }
|
|
@tracks.select { |track| track.request.priority.to_i < priority }.each { |track| remove_track(track) }
|
|
if @tracks.size >= MAX_TRACKS
|
|
victim = @tracks.min_by { |track| [track.request.priority.to_i, track.started_at] }
|
|
return if victim.request.priority.to_i > priority
|
|
remove_track(victim)
|
|
end
|
|
volume = request.kind == "music" ? @music_volume : @sound_volume
|
|
handle = @backend.start(path, request, volume)
|
|
loops = request.loops.to_i
|
|
@tracks << Track.new(
|
|
:request => request,
|
|
:handle => handle,
|
|
:url => url,
|
|
:remaining_loops => loops < 0 ? -1 : loops,
|
|
:started_at => Process.clock_gettime(Process::CLOCK_MONOTONIC)
|
|
)
|
|
end
|
|
|
|
def finish_iteration(track)
|
|
if track.remaining_loops == -1
|
|
track.handle.restart(track.request.start_at.to_i)
|
|
elsif track.remaining_loops > 1
|
|
track.remaining_loops -= 1
|
|
track.handle.restart(track.request.start_at.to_i)
|
|
else
|
|
remove_track(track)
|
|
end
|
|
end
|
|
|
|
def remove_track(track)
|
|
@tracks.delete(track)
|
|
track.handle.close rescue nil
|
|
end
|
|
|
|
def matching_identity(request)
|
|
@tracks.find do |track|
|
|
(request.key.to_s != "" && track.request.key.to_s == request.key.to_s) ||
|
|
track.request.name.to_s == request.name.to_s
|
|
end
|
|
end
|
|
|
|
def stop_match?(playing, stop)
|
|
filters = []
|
|
filters << playing.name.to_s == stop.name.to_s if stop.name.to_s != ""
|
|
filters << playing.kind.to_s == stop.kind.to_s if stop.kind.to_s != ""
|
|
filters << playing.tag.to_s == stop.tag.to_s if stop.tag.to_s != ""
|
|
filters << playing.key.to_s == stop.key.to_s if stop.key.to_s != ""
|
|
filters << playing.priority.to_i <= stop.priority.to_i if stop.priority != nil
|
|
filters.empty? || filters.all?
|
|
end
|
|
|
|
def normalize_kind(value)
|
|
%w[sound music video].include?(value.to_s.downcase) ? value.to_s.downcase : "sound"
|
|
end
|
|
|
|
def remember_caption(value)
|
|
return if value.to_s.strip == ""
|
|
@captions << value.to_s
|
|
@pending_captions << value.to_s
|
|
@captions.shift while @captions.size > 100
|
|
end
|
|
|
|
def add_error(value)
|
|
@errors << value.to_s
|
|
@errors.shift while @errors.size > 50
|
|
end
|
|
end
|
|
|
|
module ClientMedia
|
|
module_function
|
|
|
|
def play_request(data)
|
|
value = data.is_a?(Hash) ? data : {}
|
|
PlayRequest.new(
|
|
:name => value["name"], :url => value["url"], :kind => value["type"],
|
|
:tag => value["tag"], :volume => value["volume"], :fade_in => value["fadein"],
|
|
:fade_out => value["fadeout"], :start_at => value["start"], :finish_at => value["finish"],
|
|
:loops => value["loops"], :priority => value["priority"],
|
|
:continue_existing => value.fetch("continue", true), :key => value["key"],
|
|
:caption => value["caption"]
|
|
)
|
|
end
|
|
|
|
def stop_request(data)
|
|
value = data.is_a?(Hash) ? data : {}
|
|
StopRequest.new(
|
|
:name => value["name"], :kind => value["type"], :tag => value["tag"],
|
|
:priority => value["priority"], :key => value["key"],
|
|
:fade_away => value["fadeaway"] == true, :fade_out => value["fadeout"]
|
|
)
|
|
end
|
|
|
|
def preload_request(data)
|
|
value = data.is_a?(Hash) ? data : {}
|
|
PreloadRequest.new(:name => value["name"], :url => value["url"])
|
|
end
|
|
end
|
|
end
|
|
end
|