Current section
Files
Jump to
Current section
Files
lib/storage/s3.ex
defmodule Membrane.RTC.Engine.Endpoint.Recording.Storage.S3 do
@moduledoc """
`Membrane.RTC.Engine.Endpoint.Recording.Storage` implementation that saves the stream to the pointed AWS S3 bucket.
"""
@behaviour Membrane.RTC.Engine.Endpoint.Recording.Storage
require Membrane.Logger
alias Membrane.RTC.Engine.Endpoint.Recording.Storage
# minimal chunk size based on aws specification (in bytes)
@chunk_size 5_242_880
@type credentials_t :: %{
access_key_id: String.t(),
secret_access_key: String.t(),
region: String.t(),
bucket: String.t()
}
@type storage_opts :: %{:credentials => credentials_t(), optional(:path_prefix) => Path.t()}
@impl true
@spec get_sink(Storage.recording_config(), storage_opts()) :: struct()
def get_sink(config, storage_opts) do
path = s3_path(config, storage_opts)
%__MODULE__.Sink{
path: path,
credentials: storage_opts.credentials,
chunk_size: @chunk_size
}
end
@impl true
def save_object(config, %{credentials: credentials} = storage_opts) do
path = s3_path(config, storage_opts)
aws_config = create_aws_config(credentials)
result =
credentials.bucket
|> ExAws.S3.put_object(path, config.object, [])
|> ExAws.request(aws_config)
error_msg = "Couldn't save object on S3 bucket, recording id: #{config.recording_id}"
handle_upload_result(result, error_msg)
end
@impl true
def on_close(files, recording_id, storage_opts) do
list_objects_result = list_objects(recording_id, storage_opts)
with {:ok, objects} <- list_objects_result,
objects_to_fix = objects_to_fix(objects, files),
:ok <- fix_objects(files, recording_id, storage_opts, objects_to_fix) do
:ok
else
{:error, :list_objects} ->
:error
{:error, :not_fixed} ->
{:ok, objects} = list_objects_result
objects
|> Enum.map(fn {filename, _size} -> filename end)
|> clean_objects(recording_id, storage_opts)
:error
end
end
@spec create_aws_config(credentials_t()) :: list()
def create_aws_config(credentials) do
credentials
|> Enum.reject(fn {key, _value} -> key == :bucket end)
|> then(&ExAws.Config.new(:s3, &1))
|> Map.to_list()
end
defp objects_to_fix(objects, files) do
Enum.reject(files, fn {filename, {_file_path, file_size}} ->
s3_size_result = Map.get(objects, filename, -1)
s3_size_result >= file_size
end)
end
defp fix_objects(files, recording_id, storage_opts, objects) do
all_fixed? =
Enum.all?(objects, fn {filename, _size} ->
{file_path, _size} = Map.fetch!(files, filename)
fix_object(file_path, filename, recording_id, storage_opts)
end)
if all_fixed?, do: :ok, else: {:error, :not_fixed}
end
defp list_objects(recording_id, storage_opts) do
path_prefix =
storage_opts
|> Map.get(:path_prefix, "")
|> Path.join(recording_id)
credentials = storage_opts.credentials
config = create_aws_config(credentials)
response =
credentials.bucket
|> ExAws.S3.list_objects(prefix: path_prefix)
|> ExAws.request(config)
case response do
{:ok, %{body: %{contents: contents}}} ->
{:ok, Map.new(contents, &parse_stats/1)}
_else ->
Membrane.Logger.error("Couldn't list objects on S3 bucket, recording id: #{recording_id}")
{:error, :list_objects}
end
end
defp clean_objects(objects, recording_id, %{credentials: credentials}) do
config = create_aws_config(credentials)
result =
credentials.bucket
|> ExAws.S3.delete_all_objects(objects)
|> ExAws.request(config)
case result do
{:ok, _term} ->
:ok
{:error, _reason} ->
Membrane.Logger.error(
"Couldn't clean objects on S3 bucket, recording id: #{recording_id}"
)
:error
end
end
defp fix_object(file_path, filename, recording_id, storage_opts) do
config = save_object_config(nil, recording_id, filename)
case stream_object(file_path, config, storage_opts) do
:ok -> true
{:error, _response} -> false
end
end
defp stream_object(file_path, config, %{credentials: credentials} = storage_opts) do
path = s3_path(config, storage_opts)
aws_config = create_aws_config(credentials)
result =
file_path
|> ExAws.S3.Upload.stream_file()
|> ExAws.S3.upload(credentials.bucket, path)
|> ExAws.request(aws_config)
error_msg = "Couldn't stream object on S3 bucket, recording id: #{config.recording_id}"
handle_upload_result(result, error_msg)
end
defp save_object_config(object, recording_id, filename) do
%{
object: object,
recording_id: recording_id,
filename: filename
}
end
defp handle_upload_result({:ok, %{status_code: 200}}, _error_msg), do: :ok
defp handle_upload_result({:error, response}, error_msg) do
Membrane.Logger.error(error_msg)
{:error, response}
end
defp parse_stats(stats) do
filename = stats.key |> String.split("/") |> List.last()
size = String.to_integer(stats.size)
{filename, size}
end
defp s3_path(config, storage_opts) do
storage_opts
|> Map.get(:path_prefix, "")
|> Path.join(config.filename)
end
end