Packages
upstream
1.6.6
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 20_000
@behaviour unquote(__MODULE__)
alias Upstream.Job
alias Upstream.B2.Upload
alias Upstream.Uploader.{
TaskSupervisor,
Checksum,
Flow
}
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
def init(job) do
Job.start(job)
{:ok, handle_setup(%{job: job, uid: job.uid, current_state: :started})}
end
def handle_call(:upload, _from, state) do
case task(state) do
{:ok, result} ->
Job.complete(state, result)
{:stop, :normal, {:ok, result},
Map.merge(state, %{
current_state: :uploaded
})}
{:error, reason} ->
Job.error(state, reason)
{:stop, {:error, reason}, {:error, reason},
Map.merge(state, %{
current_state: :upload_failed
})}
end
end
def terminate(reason, state) do
handle_stop(state)
cond do
Job.completed?(state) ->
Logger.info("[Upstream] Completed #{state.uid.name}")
Job.errored?(state) ->
Logger.info("[Upstream] Errored #{state.uid.name}")
true ->
Job.error(state, %{error: reason})
end
reason
end
# Private functions
defp handle_stop(state), do: nil
defp handle_setup(state), do: state
defoverridable handle_stop: 1, handle_setup: 1
end
end
@callback task(map) :: {:ok, any} | {:error, any}
end