Implemented streaming during indexing

This commit is contained in:
Kieran Eglin 2024-03-03 12:09:33 -08:00
parent a46dd19ec7
commit d2865130ff
No known key found for this signature in database
GPG key ID: 193984967FCF432D
18 changed files with 437 additions and 74 deletions

View file

@ -15,6 +15,8 @@ alias Pinchflat.Sources
alias Pinchflat.MediaClient.{SourceDetails, VideoDownloader} alias Pinchflat.MediaClient.{SourceDetails, VideoDownloader}
alias Pinchflat.Metadata.{Zipper, ThumbnailFetcher} alias Pinchflat.Metadata.{Zipper, ThumbnailFetcher}
alias Pinchflat.Utils.FilesystemUtils.FileFollowerServer
defmodule IexHelpers do defmodule IexHelpers do
def playlist_url do def playlist_url do
"https://www.youtube.com/playlist?list=PLmqC3wPkeL8kSlTCcSMDD63gmSi7evcXS" "https://www.youtube.com/playlist?list=PLmqC3wPkeL8kSlTCcSMDD63gmSi7evcXS"

View file

@ -20,7 +20,8 @@ config :pinchflat,
# Setting AUTH_USERNAME and AUTH_PASSWORD implies you want to use basic auth. # Setting AUTH_USERNAME and AUTH_PASSWORD implies you want to use basic auth.
# If either is unset, basic auth will not be used. # If either is unset, basic auth will not be used.
basic_auth_username: System.get_env("AUTH_USERNAME"), basic_auth_username: System.get_env("AUTH_USERNAME"),
basic_auth_password: System.get_env("AUTH_PASSWORD") basic_auth_password: System.get_env("AUTH_PASSWORD"),
file_watcher_poll_interval: 1000
# Configures the endpoint # Configures the endpoint
config :pinchflat, PinchflatWeb.Endpoint, config :pinchflat, PinchflatWeb.Endpoint,

View file

@ -5,7 +5,8 @@ config :pinchflat,
yt_dlp_executable: Path.join([File.cwd!(), "/test/support/scripts/yt-dlp-mocks/repeater.sh"]), yt_dlp_executable: Path.join([File.cwd!(), "/test/support/scripts/yt-dlp-mocks/repeater.sh"]),
media_directory: Path.join([System.tmp_dir!(), "test", "videos"]), media_directory: Path.join([System.tmp_dir!(), "test", "videos"]),
metadata_directory: Path.join([System.tmp_dir!(), "test", "metadata"]), metadata_directory: Path.join([System.tmp_dir!(), "test", "metadata"]),
tmpfile_directory: Path.join([System.tmp_dir!(), "test", "tmpfiles"]) tmpfile_directory: Path.join([System.tmp_dir!(), "test", "tmpfiles"]),
file_watcher_poll_interval: 50
config :pinchflat, Oban, testing: :manual config :pinchflat, Oban, testing: :manual

View file

@ -4,4 +4,5 @@ defmodule Pinchflat.MediaClient.Backends.BackendCommandRunner do
""" """
@callback run(binary(), keyword(), binary()) :: {:ok, binary()} | {:error, binary(), integer()} @callback run(binary(), keyword(), binary()) :: {:ok, binary()} | {:error, binary(), integer()}
@callback run(binary(), keyword(), binary(), keyword()) :: {:ok, binary()} | {:error, binary(), integer()}
end end

View file

@ -6,6 +6,7 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.CommandRunner do
require Logger require Logger
alias Pinchflat.Utils.StringUtils alias Pinchflat.Utils.StringUtils
alias Pinchflat.Utils.FilesystemUtils, as: FSUtils
alias Pinchflat.MediaClient.Backends.BackendCommandRunner alias Pinchflat.MediaClient.Backends.BackendCommandRunner
@behaviour BackendCommandRunner @behaviour BackendCommandRunner
@ -15,19 +16,20 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.CommandRunner do
a file and then returns its contents because yt-dlp will return warnings a file and then returns its contents because yt-dlp will return warnings
to stdout even if the command is successful, but these will break JSON parsing. to stdout even if the command is successful, but these will break JSON parsing.
Returns {:ok, binary()} | {:error, output, status}. Additional Opts:
- :output_filepath - the path to save the output to. If not provided, a temporary
file will be created and used. Useful for if you need a reference to the file
for a file watcher.
IDEA: Indexing takes a long time, but the output is actually streamed to stdout. Returns {:ok, binary()} | {:error, output, status}.
Maybe we could listen to that stream instead so we can index videos as they're discovered.
See: https://stackoverflow.com/a/49061086/5665799
""" """
@impl BackendCommandRunner @impl BackendCommandRunner
def run(url, command_opts, output_template) do def run(url, command_opts, output_template, addl_opts \\ []) do
command = backend_executable() command = backend_executable()
# These must stay in exactly this order, hence why I'm giving it its own variable. # These must stay in exactly this order, hence why I'm giving it its own variable.
# Also, can't use RAM file since yt-dlp needs a concrete filepath. # Also, can't use RAM file since yt-dlp needs a concrete filepath.
json_output_path = generate_json_output_path() output_filepath = Keyword.get(addl_opts, :output_filepath, FSUtils.generate_metadata_tmpfile(:json))
print_to_file_opts = [{:print_to_file, output_template}, json_output_path] print_to_file_opts = [{:print_to_file, output_template}, output_filepath]
formatted_command_opts = [url] ++ parse_options(command_opts ++ print_to_file_opts) formatted_command_opts = [url] ++ parse_options(command_opts ++ print_to_file_opts)
Logger.info("[yt-dlp] called with: #{Enum.join(formatted_command_opts, " ")}") Logger.info("[yt-dlp] called with: #{Enum.join(formatted_command_opts, " ")}")
@ -36,24 +38,13 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.CommandRunner do
{_, 0} -> {_, 0} ->
# IDEA: consider deleting the file after reading it # IDEA: consider deleting the file after reading it
# (even on error? especially on error?) # (even on error? especially on error?)
File.read(json_output_path) File.read(output_filepath)
{output, status} -> {output, status} ->
{:error, output, status} {:error, output, status}
end end
end end
defp generate_json_output_path do
tmpfile_directory = Application.get_env(:pinchflat, :tmpfile_directory)
filepath = Path.join([tmpfile_directory, "#{StringUtils.random_string(64)}.json"])
# Ensure the file can be created and written to BEFORE we run the `yt-dlp` command
:ok = File.mkdir_p!(Path.dirname(filepath))
:ok = File.write(filepath, "")
filepath
end
# We want to satisfy the following behaviours: # We want to satisfy the following behaviours:
# #
# 1. If the key is an atom, convert it to a string and convert it to kebab case (for convenience) # 1. If the key is an atom, convert it to a string and convert it to kebab case (for convenience)

View file

@ -4,19 +4,33 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.VideoCollection do
videos (aka: a source [ie: channels, playlists]). videos (aka: a source [ie: channels, playlists]).
""" """
require Logger
alias Pinchflat.Utils.FunctionUtils alias Pinchflat.Utils.FunctionUtils
alias Pinchflat.Utils.FilesystemUtils
@doc """ @doc """
Returns a list of maps representing the videos in the collection. Returns a list of maps representing the videos in the collection.
Options:
- :file_listener_handler - a function that will be called with the path to the
file that will be written to when yt-dlp is done. This is useful for
setting up a file watcher to know when the file is ready to be read.
Returns {:ok, [map()]} | {:error, any, ...}. Returns {:ok, [map()]} | {:error, any, ...}.
""" """
def get_media_attributes(url, command_opts \\ []) do def get_media_attributes(url, addl_opts \\ []) do
runner = Application.get_env(:pinchflat, :yt_dlp_runner) runner = Application.get_env(:pinchflat, :yt_dlp_runner)
opts = command_opts ++ [:simulate, :skip_download] command_opts = [:simulate, :skip_download]
output_template = "%(.{id,title,was_live,original_url,description})j" output_template = "%(.{id,title,was_live,original_url,description})j"
output_filepath = FilesystemUtils.generate_metadata_tmpfile(:json)
file_listener_handler = Keyword.get(addl_opts, :file_listener_handler, false)
case runner.run(url, opts, output_template) do if file_listener_handler do
file_listener_handler.(output_filepath)
end
case runner.run(url, command_opts, output_template, output_filepath: output_filepath) do
{:ok, output} -> {:ok, output} ->
output output
|> String.split("\n", trim: true) |> String.split("\n", trim: true)

View file

@ -22,16 +22,21 @@ defmodule Pinchflat.MediaClient.SourceDetails do
Returns a list of basic video data mapsfor the given source URL OR Returns a list of basic video data mapsfor the given source URL OR
source record using the given backend. source record using the given backend.
Options:
- :file_listener_handler - a function that will be called with the path to the
file that will be written to by yt-dlp. This is useful for
setting up a file watcher to read the file as it gets written to.
Returns {:ok, [map()]} | {:error, any, ...}. Returns {:ok, [map()]} | {:error, any, ...}.
""" """
def get_media_attributes(sourceable, backend \\ :yt_dlp) def get_media_attributes(sourceable, opts \\ [], backend \\ :yt_dlp)
def get_media_attributes(%Source{} = source, backend) do def get_media_attributes(%Source{} = source, opts, backend) do
source_module(backend).get_media_attributes(source.collection_id) get_media_attributes(source.collection_id, opts, backend)
end end
def get_media_attributes(source_url, backend) when is_binary(source_url) do def get_media_attributes(source_url, opts, backend) when is_binary(source_url) do
source_module(backend).get_media_attributes(source_url) source_module(backend).get_media_attributes(source_url, opts)
end end
defp source_module(backend) do defp source_module(backend) do

View file

@ -1,7 +1,9 @@
defmodule Pinchflat.Tasks.MediaItemTasks do defmodule Pinchflat.Tasks.MediaItemTasks do
@moduledoc """ @moduledoc """
This module contains methods used by or used to control tasks (aka workers) Contains methods used by OR used to create/manage tasks for media items.
related to media items.
Tasks/workers are meant to be thin wrappers so most of the actual work they
do is also defined here. Essentially, a one-stop-shop for media-related tasks/workers.
""" """
alias Pinchflat.Media alias Pinchflat.Media

View file

@ -1,9 +1,13 @@
defmodule Pinchflat.Tasks.SourceTasks do defmodule Pinchflat.Tasks.SourceTasks do
@moduledoc """ @moduledoc """
This module contains methods used by or used to control tasks (aka workers) Contains methods used by OR used to create/manage tasks for sources.
related to sources.
Tasks/workers are meant to be thin wrappers so most of the actual work they
do is also defined here. Essentially, a one-stop-shop for source-related tasks/workers.
""" """
require Logger
alias Pinchflat.Media alias Pinchflat.Media
alias Pinchflat.Tasks alias Pinchflat.Tasks
alias Pinchflat.Sources alias Pinchflat.Sources
@ -11,6 +15,7 @@ defmodule Pinchflat.Tasks.SourceTasks do
alias Pinchflat.MediaClient.SourceDetails alias Pinchflat.MediaClient.SourceDetails
alias Pinchflat.Workers.MediaIndexingWorker alias Pinchflat.Workers.MediaIndexingWorker
alias Pinchflat.Workers.VideoDownloadWorker alias Pinchflat.Workers.VideoDownloadWorker
alias Pinchflat.Utils.FilesystemUtils.FileFollowerServer
@doc """ @doc """
Starts tasks for indexing a source's media regardless of the source's indexing Starts tasks for indexing a source's media regardless of the source's indexing
@ -37,27 +42,28 @@ defmodule Pinchflat.Tasks.SourceTasks do
Given a media source, creates (indexes) the media by creating media_items for each Given a media source, creates (indexes) the media by creating media_items for each
media ID in the source. media ID in the source.
Indexing is slow and usually returns a list of all media data at once for record creation.
To help with this, we use a file follower to watch the file that yt-dlp writes to
so we can create media items as they come in. This parallelizes the process and adds
clarity to the user experience. This has a few things to be aware of which are documented
below in the file watcher setup method.
Since indexing returns all media data EVERY TIME, we rely on the unique index of the
media_id to prevent duplicates. Due to both the file follower and the fact that future
indexing will index a lot of existing data, this method will MOSTLY return error
changesets (from the unique index violation) and not media items. This is intended.
Returns [%MediaItem{}, ...] | [%Ecto.Changeset{}, ...] Returns [%MediaItem{}, ...] | [%Ecto.Changeset{}, ...]
""" """
def index_media_items(%Source{} = source) do def index_media_items(%Source{} = source) do
{:ok, media_attributes} = SourceDetails.get_media_attributes(source.original_url) # See the method definition below for more info on how file watchers work
# (important reading if you're not familiar with it)
{:ok, media_attributes} = get_media_attributes_and_setup_file_watcher(source)
Sources.update_source(source, %{last_indexed_at: DateTime.utc_now()}) Sources.update_source(source, %{last_indexed_at: DateTime.utc_now()})
media_attributes Enum.map(media_attributes, fn media_attrs ->
|> Enum.map(fn media_attrs -> create_media_item_from_attributes(source, media_attrs)
attrs = %{
source_id: source.id,
title: media_attrs["title"],
media_id: media_attrs["id"],
original_url: media_attrs["original_url"],
livestream: media_attrs["was_live"],
description: media_attrs["description"]
}
case Media.create_media_item(attrs) do
{:ok, media_item} -> media_item
{:error, changeset} -> changeset
end
end) end)
end end
@ -99,4 +105,60 @@ defmodule Pinchflat.Tasks.SourceTasks do
|> Media.list_pending_media_items_for() |> Media.list_pending_media_items_for()
|> Enum.each(&Tasks.delete_pending_tasks_for/1) |> Enum.each(&Tasks.delete_pending_tasks_for/1)
end end
# The file follower is a GenServer that watches a file for new lines and
# processes them. This works well, but we have to be resilliant to partially-written
# lines (ie: you should gracefully fail if you can't parse a line).
#
# This works in-tandem with the normal (blocking) media indexing behaviour. When
# the `get_media_attributes` method completes it'll return the FULL result to
# the caller for parsing. Ideally, every item in the list will have already
# been processed by the file follower, but if not, the caller handles creation
# of any media items that were missed/initially failed.
#
# It attempts a graceful shutdown of the file follower after the indexing is done,
# but the FileFollowerServer will also stop itself if it doesn't see any activity
# for a sufficiently long time.
defp get_media_attributes_and_setup_file_watcher(source) do
{:ok, pid} = FileFollowerServer.start_link()
handler = fn filepath -> setup_file_follower_watcher(pid, filepath, source) end
result = SourceDetails.get_media_attributes(source.original_url, file_listener_handler: handler)
FileFollowerServer.stop(pid)
result
end
defp setup_file_follower_watcher(pid, filepath, source) do
FileFollowerServer.watch_file(pid, filepath, fn line ->
case Phoenix.json_library().decode(line) do
{:ok, media_attrs} ->
Logger.debug("FileFollowerServer Handler: Got media attributes: #{inspect(media_attrs)}")
create_media_item_from_attributes(source, media_attrs)
err ->
Logger.debug("FileFollowerServer Handler: Error decoding JSON: #{inspect(err)}")
err
end
end)
end
defp create_media_item_from_attributes(source, media_attrs) do
attrs = %{
source_id: source.id,
title: media_attrs["title"],
media_id: media_attrs["id"],
original_url: media_attrs["original_url"],
livestream: media_attrs["was_live"],
description: media_attrs["description"]
}
case Media.create_media_item(attrs) do
{:ok, media_item} -> media_item
{:error, changeset} -> changeset
end
end
end end

View file

@ -0,0 +1,23 @@
defmodule Pinchflat.Utils.FilesystemUtils do
@moduledoc """
Utility methods for working with the filesystem
"""
alias Pinchflat.Utils.StringUtils
@doc """
Generates a temporary file and returns its path. The file is empty and has the given type.
Generates all the directories in the path if they don't exist.
Returns binary()
"""
def generate_metadata_tmpfile(type) do
tmpfile_directory = Application.get_env(:pinchflat, :tmpfile_directory)
filepath = Path.join([tmpfile_directory, "#{StringUtils.random_string(64)}.#{type}"])
:ok = File.mkdir_p!(Path.dirname(filepath))
:ok = File.write(filepath, "")
filepath
end
end

View file

@ -0,0 +1,121 @@
defmodule Pinchflat.Utils.FilesystemUtils.FileFollowerServer do
@moduledoc """
A GenServer that watches a file for new lines and processes them as they come in.
This is useful for tailing log files and other similar tasks. If there's no activity
for a certain amount of time, the server will stop itself.
"""
use GenServer
require Logger
@poll_interval_ms Application.compile_env(:pinchflat, :file_watcher_poll_interval)
@activity_timeout_ms 10_000
# Client API
@doc """
Starts the file follower server
Returns {:ok, pid} or {:error, reason}
"""
def start_link() do
GenServer.start_link(__MODULE__, [])
end
@doc """
Starts the file watcher for a given filepath and handler function.
Returns :ok
"""
def watch_file(process, filepath, handler) do
GenServer.cast(process, {:watch_file, filepath, handler})
end
@doc """
Stops the file watcher and closes the file.
Returns :ok
"""
def stop(process) do
GenServer.cast(process, :stop)
end
# Server Callbacks
@impl true
def init(_opts) do
# Start with a blank state because, based on the common calling
# pattern for this module, we'll need a reference to the server's
# PID before we start watching any files so we can later stop the
# server gracefully.
{:ok, %{}}
end
@impl true
def handle_cast({:watch_file, filepath, handler}, _old_state) do
{:ok, io_device} = :file.open(filepath, [:raw, :read_ahead, :binary])
state = %{
io_device: io_device,
last_activity: DateTime.utc_now(),
handler: handler
}
Process.send(self(), :read_new_lines, [])
{:noreply, state}
end
@impl true
def handle_cast(:stop, state) do
Logger.debug("Gracefully stopping file follower")
:file.close(state.io_device)
{:stop, :normal, state}
end
@impl true
def handle_info(:read_new_lines, state) do
last_activity = state.last_activity
# If there's no new lines written for a certain amount of time, stop the server
if DateTime.diff(DateTime.utc_now(), last_activity, :millisecond) > @activity_timeout_ms do
Logger.debug("No activity for #{@activity_timeout_ms}ms. Requesting stop.")
stop(self())
{:noreply, state}
else
attempt_process_new_lines(state)
end
end
defp attempt_process_new_lines(state) do
io_device = state.io_device
# This reads one line at a time. If a line is found, it
# will be passed to the handler, we'll note the time of
# the last activity, and then we'll immediately call this
# again to read the next line.
#
# If there are no lines, it waits for the poll interval
# before trying again.
case :file.read_line(io_device) do
{:ok, line} ->
state.handler.(line)
Process.send(self(), :read_new_lines, [])
{:noreply, %{state | last_activity: DateTime.utc_now()}}
:eof ->
Logger.debug("EOF reached, waiting before trying to read new lines")
Process.send_after(self(), :read_new_lines, @poll_interval_ms)
{:noreply, state}
{:error, reason} ->
Logger.error("Error reading file: #{reason}")
stop(self())
{:noreply, state}
end
end
end

View file

@ -10,7 +10,7 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.CommandRunnerTest do
on_exit(&reset_executable/0) on_exit(&reset_executable/0)
end end
describe "run/2" do describe "run/4" do
test "it returns the output and status when the command succeeds" do test "it returns the output and status when the command succeeds" do
assert {:ok, _output} = Runner.run(@video_url, [], "") assert {:ok, _output} = Runner.run(@video_url, [], "")
end end
@ -57,6 +57,12 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.CommandRunnerTest do
assert {:error, "", 1} = Runner.run(@video_url, [], "") assert {:error, "", 1} = Runner.run(@video_url, [], "")
end) end)
end end
test "optionally lets you specify an output_filepath" do
assert {:ok, output} = Runner.run(@video_url, [], "%(id)s", output_filepath: "/tmp/yt-dlp-output.json")
assert String.contains?(output, "--print-to-file %(id)s /tmp/yt-dlp-output.json")
end
end end
defp wrap_executable(new_executable, fun) do defp wrap_executable(new_executable, fun) do

View file

@ -11,14 +11,16 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.VideoCollectionTest do
describe "get_media_attributes/2" do describe "get_media_attributes/2" do
test "returns a list of video attributes with no blank elements" do test "returns a list of video attributes with no blank elements" do
expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, source_attributes_return_fixture() <> "\n\n"} end) expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts ->
{:ok, source_attributes_return_fixture() <> "\n\n"}
end)
assert {:ok, [%{"id" => "video1"}, %{"id" => "video2"}, %{"id" => "video3"}]} = assert {:ok, [%{"id" => "video1"}, %{"id" => "video2"}, %{"id" => "video3"}]} =
VideoCollection.get_media_attributes(@channel_url) VideoCollection.get_media_attributes(@channel_url)
end end
test "it passes the expected default args" do test "it passes the expected default args" do
expect(YtDlpRunnerMock, :run, fn _url, opts, ot -> expect(YtDlpRunnerMock, :run, fn _url, opts, ot, _addl_opts ->
assert opts == [:simulate, :skip_download] assert opts == [:simulate, :skip_download]
assert ot == "%(.{id,title,was_live,original_url,description})j" assert ot == "%(.{id,title,was_live,original_url,description})j"
@ -28,20 +30,35 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.VideoCollectionTest do
assert {:ok, _} = VideoCollection.get_media_attributes(@channel_url) assert {:ok, _} = VideoCollection.get_media_attributes(@channel_url)
end end
test "it passes the expected custom args" do test "returns the error straight through when the command fails" do
expect(YtDlpRunnerMock, :run, fn _url, opts, _ot -> expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts -> {:error, "Big issue", 1} end)
assert opts == [:custom_arg, :simulate, :skip_download]
assert {:error, "Big issue", 1} = VideoCollection.get_media_attributes(@channel_url)
end
test "passes the explict tmpfile path to runner" do
expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, addl_opts ->
assert [{:output_filepath, filepath}] = addl_opts
assert String.ends_with?(filepath, ".json")
{:ok, ""} {:ok, ""}
end) end)
assert {:ok, _} = VideoCollection.get_media_attributes(@channel_url, [:custom_arg]) assert {:ok, _} = VideoCollection.get_media_attributes(@channel_url)
end end
test "returns the error straight through when the command fails" do test "supports an optional file_listener_handler that gets passed a filename" do
expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:error, "Big issue", 1} end) expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts -> {:ok, ""} end)
current_self = self()
assert {:error, "Big issue", 1} = VideoCollection.get_media_attributes(@channel_url) handler = fn filename ->
send(current_self, {:handler, filename})
end
assert {:ok, _} = VideoCollection.get_media_attributes(@channel_url, file_listener_handler: handler)
assert_receive {:handler, filename}
assert String.ends_with?(filename, ".json")
end end
end end

View file

@ -45,7 +45,7 @@ defmodule Pinchflat.MediaClient.SourceDetailsTest do
describe "get_media_attributes/2 when passed a string" do describe "get_media_attributes/2 when passed a string" do
test "it passes the expected arguments to the backend" do test "it passes the expected arguments to the backend" do
expect(YtDlpRunnerMock, :run, fn @channel_url, opts, ot -> expect(YtDlpRunnerMock, :run, fn @channel_url, opts, ot, _addl_opts ->
assert opts == [:simulate, :skip_download] assert opts == [:simulate, :skip_download]
assert ot == "%(.{id,title,was_live,original_url,description})j" assert ot == "%(.{id,title,was_live,original_url,description})j"
@ -56,7 +56,7 @@ defmodule Pinchflat.MediaClient.SourceDetailsTest do
end end
test "it returns a list of maps" do test "it returns a list of maps" do
expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts ->
{:ok, source_attributes_return_fixture()} {:ok, source_attributes_return_fixture()}
end) end)
@ -68,7 +68,7 @@ defmodule Pinchflat.MediaClient.SourceDetailsTest do
test "it calls the backend with the source's collection ID" do test "it calls the backend with the source's collection ID" do
source = source_fixture() source = source_fixture()
expect(YtDlpRunnerMock, :run, fn url, _opts, _ot -> expect(YtDlpRunnerMock, :run, fn url, _opts, _ot, _addl_opts ->
assert source.collection_id == url assert source.collection_id == url
{:ok, source_attributes_return_fixture()} {:ok, source_attributes_return_fixture()}
end) end)
@ -77,7 +77,7 @@ defmodule Pinchflat.MediaClient.SourceDetailsTest do
end end
test "it builds options based on the source's media profile" do test "it builds options based on the source's media profile" do
expect(YtDlpRunnerMock, :run, fn _url, opts, _ot -> expect(YtDlpRunnerMock, :run, fn _url, opts, _ot, _addl_opts ->
assert opts == [:simulate, :skip_download] assert opts == [:simulate, :skip_download]
{:ok, ""} {:ok, ""}
end) end)
@ -91,5 +91,22 @@ defmodule Pinchflat.MediaClient.SourceDetailsTest do
source = source_fixture(media_profile_id: media_profile.id) source = source_fixture(media_profile_id: media_profile.id)
assert {:ok, _} = SourceDetails.get_media_attributes(source) assert {:ok, _} = SourceDetails.get_media_attributes(source)
end end
test "lets you pass through an optional file_listener_handler" do
expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts ->
{:ok, source_attributes_return_fixture()}
end)
source = source_fixture()
current_self = self()
handler = fn filename ->
send(current_self, {:handler, filename})
end
assert {:ok, _} = SourceDetails.get_media_attributes(source, file_listener_handler: handler)
assert_receive {:handler, _}
end
end end
end end

View file

@ -45,7 +45,9 @@ defmodule Pinchflat.Tasks.SourceTasksTest do
describe "index_media_items/1" do describe "index_media_items/1" do
setup do setup do
stub(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, source_attributes_return_fixture()} end) stub(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts ->
{:ok, source_attributes_return_fixture()}
end)
{:ok, [source: source_fixture()]} {:ok, [source: source_fixture()]}
end end
@ -108,6 +110,30 @@ defmodule Pinchflat.Tasks.SourceTasksTest do
end end
end end
describe "index_media_items/1 when testing file watcher" do
setup do
{:ok, [source: source_fixture()]}
end
test "it creates a new media item for everything already in the file", %{source: source} do
watcher_poll_interval = Application.get_env(:pinchflat, :file_watcher_poll_interval)
stub(YtDlpRunnerMock, :run, fn _url, _opts, _ot, addl_opts ->
filepath = Keyword.get(addl_opts, :output_filepath)
File.write(filepath, source_attributes_return_fixture())
# Need to add a delay to ensure the file watcher has time to read the file
:timer.sleep(watcher_poll_interval * 2)
{:ok, ""}
end)
assert Repo.aggregate(MediaItem, :count, :id) == 0
SourceTasks.index_media_items(source)
assert Repo.aggregate(MediaItem, :count, :id) == 3
end
end
describe "enqueue_pending_media_tasks/1" do describe "enqueue_pending_media_tasks/1" do
test "it enqueues a job for each pending media item" do test "it enqueues a job for each pending media item" do
source = source_fixture() source = source_fixture()

View file

@ -0,0 +1,52 @@
defmodule Pinchflat.Utils.FilesystemUtils.FileFollowerServerTest do
use ExUnit.Case, async: true
alias alias Pinchflat.Utils.FilesystemUtils
alias Pinchflat.Utils.FilesystemUtils.FileFollowerServer
setup do
{:ok, pid} = FileFollowerServer.start_link()
tmpfile = FilesystemUtils.generate_metadata_tmpfile(:txt)
{:ok, %{pid: pid, tmpfile: tmpfile}}
end
describe "watch_file" do
test "calls the handler for each existing line in the file", %{pid: pid, tmpfile: tmpfile} do
File.write!(tmpfile, "line1\nline2")
parent = self()
handler = fn line -> send(parent, line) end
FileFollowerServer.watch_file(pid, tmpfile, handler)
assert_receive "line1\n"
assert_receive "line2"
end
test "calls the handler for each new line in the file", %{pid: pid, tmpfile: tmpfile} do
parent = self()
file = File.open!(tmpfile, [:append])
handler = fn line -> send(parent, line) end
FileFollowerServer.watch_file(pid, tmpfile, handler)
IO.binwrite(file, "line1\n")
assert_receive "line1\n"
IO.binwrite(file, "line2")
assert_receive "line2"
end
end
describe "stop" do
test "stops the watcher", %{pid: pid, tmpfile: tmpfile} do
handler = fn _line -> :noop end
FileFollowerServer.watch_file(pid, tmpfile, handler)
refute is_nil(Process.info(pid))
FileFollowerServer.stop(pid)
# Gotta wait for the server to stop async
:timer.sleep(10)
assert is_nil(Process.info(pid))
end
end
end

View file

@ -0,0 +1,16 @@
defmodule Pinchflat.Utils.FilesystemUtilsTest do
use ExUnit.Case, async: true
alias Pinchflat.Utils.FilesystemUtils
describe "generate_metadata_tmpfile/1" do
test "creates a tmpfile and returns its path" do
res = FilesystemUtils.generate_metadata_tmpfile(:json)
assert String.ends_with?(res, ".json")
assert File.exists?(res)
File.rm!(res)
end
end
end

View file

@ -13,7 +13,7 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do
describe "perform/1" do describe "perform/1" do
test "it indexes the source if it should be indexed" do test "it indexes the source if it should be indexed" do
expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, ""} end) expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts -> {:ok, ""} end)
source = source_fixture(index_frequency_minutes: 10) source = source_fixture(index_frequency_minutes: 10)
@ -21,7 +21,7 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do
end end
test "it indexes the source no matter what if the source has never been indexed before" do test "it indexes the source no matter what if the source has never been indexed before" do
expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, ""} end) expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts -> {:ok, ""} end)
source = source_fixture(index_frequency_minutes: 0, last_indexed_at: nil) source = source_fixture(index_frequency_minutes: 0, last_indexed_at: nil)
@ -29,7 +29,7 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do
end end
test "it does not do any indexing if the source has been indexed and shouldn't be rescheduled" do test "it does not do any indexing if the source has been indexed and shouldn't be rescheduled" do
expect(YtDlpRunnerMock, :run, 0, fn _url, _opts, _ot -> {:ok, ""} end) expect(YtDlpRunnerMock, :run, 0, fn _url, _opts, _ot, _addl_opts -> {:ok, ""} end)
source = source_fixture(index_frequency_minutes: -1, last_indexed_at: DateTime.utc_now()) source = source_fixture(index_frequency_minutes: -1, last_indexed_at: DateTime.utc_now())
@ -37,7 +37,7 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do
end end
test "it does not reschedule if the source shouldn't be indexed" do test "it does not reschedule if the source shouldn't be indexed" do
stub(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, ""} end) stub(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts -> {:ok, ""} end)
source = source_fixture(index_frequency_minutes: -1) source = source_fixture(index_frequency_minutes: -1)
perform_job(MediaIndexingWorker, %{id: source.id}) perform_job(MediaIndexingWorker, %{id: source.id})
@ -46,7 +46,9 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do
end end
test "it kicks off a download job for each pending media item" do test "it kicks off a download job for each pending media item" do
expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, source_attributes_return_fixture()} end) expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts ->
{:ok, source_attributes_return_fixture()}
end)
source = source_fixture(index_frequency_minutes: 10) source = source_fixture(index_frequency_minutes: 10)
perform_job(MediaIndexingWorker, %{id: source.id}) perform_job(MediaIndexingWorker, %{id: source.id})
@ -55,7 +57,9 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do
end end
test "it starts a job for any pending media item even if it's from another run" do test "it starts a job for any pending media item even if it's from another run" do
expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, source_attributes_return_fixture()} end) expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts ->
{:ok, source_attributes_return_fixture()}
end)
source = source_fixture(index_frequency_minutes: 10) source = source_fixture(index_frequency_minutes: 10)
media_item_fixture(%{source_id: source.id, media_filepath: nil}) media_item_fixture(%{source_id: source.id, media_filepath: nil})
@ -65,7 +69,9 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do
end end
test "it does not kick off a job for media items that could not be saved" do test "it does not kick off a job for media items that could not be saved" do
expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, source_attributes_return_fixture()} end) expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts ->
{:ok, source_attributes_return_fixture()}
end)
source = source_fixture(index_frequency_minutes: 10) source = source_fixture(index_frequency_minutes: 10)
media_item_fixture(%{source_id: source.id, media_filepath: nil, media_id: "video1"}) media_item_fixture(%{source_id: source.id, media_filepath: nil, media_id: "video1"})
@ -76,7 +82,7 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do
end end
test "it reschedules the job based on the index frequency" do test "it reschedules the job based on the index frequency" do
expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, ""} end) expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts -> {:ok, ""} end)
source = source_fixture(index_frequency_minutes: 10) source = source_fixture(index_frequency_minutes: 10)
perform_job(MediaIndexingWorker, %{id: source.id}) perform_job(MediaIndexingWorker, %{id: source.id})
@ -89,7 +95,7 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do
end end
test "it creates a task for the rescheduled job" do test "it creates a task for the rescheduled job" do
expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, ""} end) expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts -> {:ok, ""} end)
source = source_fixture(index_frequency_minutes: 10) source = source_fixture(index_frequency_minutes: 10)
task_count_fetcher = fn -> Enum.count(Tasks.list_tasks()) end task_count_fetcher = fn -> Enum.count(Tasks.list_tasks()) end
@ -100,7 +106,7 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do
end end
test "it creates the basic media_item records" do test "it creates the basic media_item records" do
expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, source_attributes_return_fixture()} end) expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot, _addl_opts -> {:ok, source_attributes_return_fixture()} end)
source = source_fixture(index_frequency_minutes: 10) source = source_fixture(index_frequency_minutes: 10)