Packages
upstream
2.0.4
2.1.4
2.1.3
2.1.2
2.1.1
2.1.0
2.0.7
2.0.6
2.0.5
2.0.4
2.0.3
2.0.2
2.0.0
1.8.2
1.8.1
1.7.2
1.7.1
1.7.0
1.6.14
1.6.13
1.6.12
1.6.11
1.6.10
1.6.9
1.6.8
1.6.7
1.6.6
1.6.5
1.6.4
1.6.3
1.6.2
1.6.1
1.6.0
1.5.12
1.5.11
1.5.11-dev-3
1.5.11-dev-2
1.5.11-dev-1
1.5.11-dev
1.5.10
1.5.9
1.5.8
1.5.7
1.5.6
1.5.5
1.5.2
1.5.1
1.5.0
1.4.12
1.4.9
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.7
1.3.5
1.3.4
1.3.3
1.3.2
1.3.1
1.3.0
1.2.3
1.2.1
Upstream is for integrating into projects that need to do large file uploads to B2 service. It integrates tightly with Backblaze B2 for now, with plans to support Amazon S3.
Current section
Files
Jump to
Current section
Files
lib/upstream/worker/base.ex
defmodule Upstream.Worker.Base do
@moduledoc """
Simple Worker for single threaded uploading
"""
defmacro __using__(_) do
quote do
use GenServer
@upload_timeout Application.get_env(:upstream, :upload)[:timeout] || 20_000
@behaviour unquote(__MODULE__)
alias Upstream.Job
alias Upstream.B2.Upload
alias Upstream.Worker.{
Checksum,
Flow
}
alias Upstream.B2.Account
require Logger
# Client API
def start_link(job) do
GenServer.start_link(__MODULE__, job)
end
def upload(pid) do
GenServer.call(pid, :upload, @upload_timeout)
end
# Server Callbacks
@impl true
def init(job) do
Job.State.start(job)
{:ok, handle_setup(%{job: job, current_state: :started})}
end
@impl true
def handle_call(:upload, _from, %{job: job} = state) do
case task(state) do
{:ok, result} ->
Job.State.complete(job, result)
{:stop, :normal, {:ok, result},
Map.merge(state, %{
current_state: :uploaded
})}
{:error, reason} ->
Job.State.error(job, reason)
{:stop, {:error, reason}, {:error, reason},
Map.merge(state, %{
current_state: :upload_failed
})}
end
end
@impl true
def terminate(reason, %{job: job} = state) do
handle_stop(state)
cond do
Job.State.completed?(job) ->
Logger.info("[Upstream] Completed #{job.uid.name}")
Job.State.errored?(job) ->
Logger.info("[Upstream] Errored #{job.uid.name}")
true ->
Job.State.error(job, %{error: reason})
end
reason
end
defp handle_setup(state), do: state
defp handle_stop(state), do: {:ok, state.current_state}
defoverridable handle_stop: 1, handle_setup: 1
end
end
@callback task(map) :: {:ok, any} | {:error, any}
end