Refactor modules into contexts (#78)
* [WIP] break out a few contexts, start refactoring fast index modules * [WIP] more contexts, this time around slow indexing and downloads * [WIP] got all tests passing * [WIP] Added moduledocs * Built a genserver to rename old jobs on boot * Added a module naming check; moved things around * Fixed specs
This commit is contained in:
@@ -0,0 +1,121 @@
|
||||
defmodule Pinchflat.SlowIndexing.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 60_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
|
||||
@@ -0,0 +1,107 @@
|
||||
defmodule Pinchflat.SlowIndexing.MediaCollectionIndexingWorker do
|
||||
@moduledoc false
|
||||
|
||||
use Oban.Worker,
|
||||
queue: :media_collection_indexing,
|
||||
unique: [period: :infinity, states: [:available, :scheduled, :retryable]],
|
||||
tags: ["media_source", "media_collection_indexing"]
|
||||
|
||||
alias __MODULE__
|
||||
alias Pinchflat.Tasks
|
||||
alias Pinchflat.Sources
|
||||
alias Pinchflat.Sources.Source
|
||||
alias Pinchflat.FastIndexing.FastIndexingWorker
|
||||
alias Pinchflat.SlowIndexing.SlowIndexingHelpers
|
||||
|
||||
@impl Oban.Worker
|
||||
@doc """
|
||||
The ID is that of a source _record_, not a YouTube channel/playlist ID. Indexes
|
||||
the provided source, kicks off downloads for each new MediaItem, and
|
||||
reschedules the job to run again in the future. It will ALWAYS index a source
|
||||
if it's never been indexed before, but rescheduling is determined by the
|
||||
`index_frequency_minutes` field.
|
||||
|
||||
README: Re-scheduling here works a little different than you may expect.
|
||||
The reschedule time is relative to the time the job has actually _completed_.
|
||||
This has some benefits but also side effects to be aware of:
|
||||
|
||||
- Benefit: No chance for jobs to overlap if a job takes longer than the
|
||||
scheduled interval. Less likely to hit API rate limits.
|
||||
- Side effect: Intervals are "soft" and _always_ walk forward. This may cause
|
||||
user confusion since a 30-minute job scheduled for every hour will
|
||||
actually run every 1 hour and 30 minutes. The tradeoff of not inundating
|
||||
the API with requests and also not overlapping jobs is worth it, IMO.
|
||||
|
||||
Order of operations:
|
||||
1. The user saves a source
|
||||
2. This job is automatically scheduled immediately. This happens in all cases.
|
||||
3. This job indexes all content for the given source. A download job is
|
||||
enqueued for each media item that should be downloaded. This can be impacted
|
||||
by the `download_media` field on the source as well as the profile's
|
||||
shorts/livestream behaviour. At this step we also attach a file reader
|
||||
to the `yt-dlp` output file so we can create media items as they come in
|
||||
for a little speedup (see {Fast,Slow}IndexingHelpers comments for more)
|
||||
4. If this job is meant to reschedule (ie: has an index frequency > 0),
|
||||
it reschedules itself. If not, it runs once and does not reschedule
|
||||
5. If the source uses fast indexing, that job is kicked off as well. It
|
||||
uses RSS to run a smaller, faster, and more frequent index. That job
|
||||
handles rescheduling itself but largely has a similar behaviour to this
|
||||
job in that it kicks off index and maybe download jobs. The biggest difference
|
||||
is that an index job is kicked off _for each new media item_ as opposed
|
||||
to one larger index job. Check out `MediaIndexingWorker` comments for more.
|
||||
6. If the job reschedules, the cycle from step 3 repeats until the heat death
|
||||
of the universe. The user changing things like the index frequency can
|
||||
dequeue or reschedule jobs as well
|
||||
|
||||
NOTE: Since indexing can take a LONG time, I should check what happens if an
|
||||
application restart occurs while a job is running. Will the job be lost?
|
||||
|
||||
Returns :ok | {:ok, %Task{}}
|
||||
"""
|
||||
def perform(%Oban.Job{args: %{"id" => source_id}}) do
|
||||
source = Sources.get_source!(source_id)
|
||||
|
||||
case {source.index_frequency_minutes, source.last_indexed_at} do
|
||||
{index_freq, _} when index_freq > 0 ->
|
||||
# If the indexing is on a schedule simply run indexing and reschedule
|
||||
SlowIndexingHelpers.index_and_enqueue_download_for_media_items(source)
|
||||
maybe_enqueue_fast_indexing_task(source)
|
||||
reschedule_indexing(source)
|
||||
|
||||
{_, nil} ->
|
||||
# If the source has never been indexed, index it once
|
||||
# even if it's not meant to reschedule
|
||||
SlowIndexingHelpers.index_and_enqueue_download_for_media_items(source)
|
||||
:ok
|
||||
|
||||
_ ->
|
||||
# If the source HAS been indexed and is not meant to reschedule,
|
||||
# perform a no-op
|
||||
:ok
|
||||
end
|
||||
end
|
||||
|
||||
defp reschedule_indexing(source) do
|
||||
next_run_in = source.index_frequency_minutes * 60
|
||||
|
||||
%{id: source.id}
|
||||
|> MediaCollectionIndexingWorker.new(schedule_in: next_run_in)
|
||||
|> Tasks.create_job_with_task(source)
|
||||
|> case do
|
||||
{:ok, task} -> {:ok, task}
|
||||
{:error, :duplicate_job} -> {:ok, :job_exists}
|
||||
end
|
||||
end
|
||||
|
||||
defp maybe_enqueue_fast_indexing_task(source) do
|
||||
if source.fast_index do
|
||||
Tasks.delete_pending_tasks_for(source, "FastIndexingWorker")
|
||||
|
||||
next_run_in = Source.fast_index_frequency() * 60
|
||||
|
||||
%{id: source.id}
|
||||
|> FastIndexingWorker.new(schedule_in: next_run_in)
|
||||
|> Tasks.create_job_with_task(source)
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -0,0 +1,136 @@
|
||||
defmodule Pinchflat.SlowIndexing.SlowIndexingHelpers do
|
||||
@moduledoc """
|
||||
Methods for performing slow indexing tasks and managing the indexing process.
|
||||
|
||||
Many of these methods are made to be kickoff or be consumed by workers.
|
||||
"""
|
||||
|
||||
require Logger
|
||||
|
||||
alias Pinchflat.Media
|
||||
alias Pinchflat.Tasks
|
||||
alias Pinchflat.Sources
|
||||
alias Pinchflat.Sources.Source
|
||||
alias Pinchflat.Media.MediaItem
|
||||
alias Pinchflat.YtDlp.MediaCollection
|
||||
alias Pinchflat.Downloading.DownloadingHelpers
|
||||
alias Pinchflat.SlowIndexing.FileFollowerServer
|
||||
alias Pinchflat.Downloading.MediaDownloadWorker
|
||||
alias Pinchflat.SlowIndexing.MediaCollectionIndexingWorker
|
||||
|
||||
alias Pinchflat.YtDlp.Media, as: YtDlpMedia
|
||||
|
||||
@doc """
|
||||
Starts tasks for indexing a source's media regardless of the source's indexing
|
||||
frequency. It's assumed the caller will check for indexing frequency.
|
||||
|
||||
Returns {:ok, %Task{}}.
|
||||
"""
|
||||
def kickoff_indexing_task(%Source{} = source) do
|
||||
Tasks.delete_pending_tasks_for(source, "FastIndexingWorker")
|
||||
Tasks.delete_pending_tasks_for(source, "MediaIndexingWorker")
|
||||
Tasks.delete_pending_tasks_for(source, "MediaCollectionIndexingWorker")
|
||||
|
||||
%{id: source.id}
|
||||
# Schedule this one immediately, but future ones will be on an interval
|
||||
|> MediaCollectionIndexingWorker.new()
|
||||
|> Tasks.create_job_with_task(source)
|
||||
end
|
||||
|
||||
@doc """
|
||||
Given a media source, creates (indexes) the media by creating media_items for each
|
||||
media ID in the source. Afterward, kicks off a download task for each pending media
|
||||
item belonging to the source. You can't tell me the method name isn't descriptive!
|
||||
|
||||
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.
|
||||
|
||||
NOTE: downloads are only enqueued if the source is set to download media. Downloads are
|
||||
also enqueued for ALL pending media items, not just the ones that were indexed in this
|
||||
job run. This should ensure that any stragglers are caught if, for some reason, they
|
||||
weren't enqueued or somehow got de-queued.
|
||||
|
||||
Since indexing returns all media data EVERY TIME, we that that opportunity to update
|
||||
indexing metadata for media items that have already been created.
|
||||
|
||||
Returns [%MediaItem{}, ...]
|
||||
"""
|
||||
def index_and_enqueue_download_for_media_items(%Source{} = source) do
|
||||
# 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_for_collection_and_setup_file_watcher(source)
|
||||
|
||||
result =
|
||||
Enum.map(media_attributes, fn media_attrs ->
|
||||
case Media.create_media_item_from_backend_attrs(source, media_attrs) do
|
||||
{:ok, media_item} -> media_item
|
||||
{:error, changeset} -> changeset
|
||||
end
|
||||
end)
|
||||
|
||||
Sources.update_source(source, %{last_indexed_at: DateTime.utc_now()})
|
||||
DownloadingHelpers.enqueue_pending_download_tasks(source)
|
||||
|
||||
result
|
||||
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_for_collection` 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_for_collection_and_setup_file_watcher(source) do
|
||||
{:ok, pid} = FileFollowerServer.start_link()
|
||||
|
||||
handler = fn filepath -> setup_file_follower_watcher(pid, filepath, source) end
|
||||
result = MediaCollection.get_media_attributes_for_collection(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)}")
|
||||
|
||||
media_struct = YtDlpMedia.response_to_struct(media_attrs)
|
||||
create_media_item_and_enqueue_download(source, media_struct)
|
||||
|
||||
err ->
|
||||
Logger.debug("FileFollowerServer Handler: Error decoding JSON: #{inspect(err)}")
|
||||
|
||||
err
|
||||
end
|
||||
end)
|
||||
end
|
||||
|
||||
defp create_media_item_and_enqueue_download(source, media_attrs) do
|
||||
case Media.create_media_item_from_backend_attrs(source, media_attrs) do
|
||||
{:ok, %MediaItem{} = media_item} ->
|
||||
if source.download_media && Media.pending_download?(media_item) do
|
||||
Logger.debug("FileFollowerServer Handler: Enqueuing download task for #{inspect(media_attrs)}")
|
||||
|
||||
%{id: media_item.id}
|
||||
|> MediaDownloadWorker.new()
|
||||
|> Tasks.create_job_with_task(media_item)
|
||||
end
|
||||
|
||||
{:error, changeset} ->
|
||||
changeset
|
||||
end
|
||||
end
|
||||
end
|
||||
Reference in New Issue
Block a user