Packages
upstream
1.5.5
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/store/store.ex
defmodule Upstream.Store do
@moduledoc """
The Store module is used to store the state of uploads.
If you are using the Upstream module in a distributed system you will need to,
set the redis_url: option for upstream, as uploads can happen from any of your node.
"""
use GenServer
alias Upstream.Store.{
Redis, Ets
}
def start_link(_) do
GenServer.start_link(__MODULE__, :ok, name: __MODULE__)
end
def exist?(key) do
GenServer.call(__MODULE__, {:exist?, key})
end
def add_member(key, value) do
GenServer.call(__MODULE__, {:add_member, key, value})
end
def is_member?(key, value) do
GenServer.call(__MODULE__, {:is_member?, key, value})
end
def remove_member(key, value) do
GenServer.call(__MODULE__, {:remove_member, key, value})
end
def move_member(from, to, value) do
GenServer.call(__MODULE__, {:move_member, from, to, value})
end
def get(key) do
GenServer.call(__MODULE__, {:get, key})
end
def set(key, value) do
GenServer.call(__MODULE__, {:set, key, value})
end
def remove(key) do
GenServer.call(__MODULE__, {:remove, key})
end
# Callbacks
def init(:ok) do
with {:ok, conn, type} <- create_store() do
{:ok, {conn, type}}
end
end
def handle_call({:exist?, key}, _from, {conn, :ets}) do
case :ets.lookup(conn, key) do
[{_k, _value}] -> {:reply, true, {conn, :ets}}
[] -> {:reply, false, {conn, :ets}}
end
end
def handle_call({:get, key}, _from, {conn, :ets}) do
case :ets.lookup(conn, key) do
[{_k, value}] -> {:reply, value, {conn, :ets}}
[] -> {:reply, nil, {conn, :ets}}
end
end
def handle_call({:remove, key}, _from, {conn, :ets}) do
:ets.delete(conn, key)
{:reply, :ok, {conn, :redis}}
end
def handle_call({:set, key, value}, _from, {conn, :ets}) do
case :ets.insert_new(conn, {key, value}) do
true -> {:reply, {:ok, value}, {conn, :ets}}
false -> {:reply, {:error, :already_set}, {conn, :ets}}
end
end
def handle_call({:add_member, key, value}, _from, {conn, :ets}) do
case :ets.lookup(conn, key) do
[{_k, existing}] ->
:ets.insert(conn, {key, [value | existing]})
{:reply, {:ok, value}, {conn, :ets}}
[] ->
:ets.insert_new(conn, {key, [value]})
{:reply, {:ok, value}, {conn, :ets}}
end
end
def handle_call({:move_member, from, to, value}, _from, {conn, :ets}) do
case {:ets.lookup(conn, from), :ets.lookup(conn, to)} do
{[{from_key, from_value}], [{to_key, to_existing}]} ->
Ets.remove_member(conn, from_key, from_value, value)
:ets.insert(conn, {to_key, [value | to_existing]})
{:reply, :ok, {conn, :ets}}
{[{from_key, from_value}], []} ->
Ets.remove_member(conn, from_key, from_value, value)
:ets.insert_new(conn, {to, [value]})
{:reply, :ok, {conn, :ets}}
_ ->
{:reply, :error, {conn, :ets}}
end
end
def handle_call({:is_member?, key, value}, _from, {conn, :ets}) do
case :ets.lookup(conn, key) do
[{_, existing}] -> {:reply, Enum.member?(existing, value), {conn, :ets}}
[] -> {:reply, false, {conn, :ets}}
end
end
def handle_call({:remove_member, key, value}, _from, {conn, :ets}) do
case :ets.lookup(conn, key) do
[{_k, existing}] ->
Ets.remove_member(conn, key, existing, value)
{:reply, :ok, {conn, :ets}}
[] ->
{:reply, :error, {conn, :ets}}
end
end
def handle_call({:exist?, key}, _from, {conn, :redis}) do
case Redix.command(conn, ["EXISTS", Redis.namespace(key)]) do
{:ok, 0} -> {:reply, false, {conn, :redis}}
{:ok, 1} -> {:reply, true, {conn, :redis}}
end
end
def handle_call({:is_member?, key, value}, _from, {conn, :redis}) do
case Redix.command(conn, ["SISMEMBER", Redis.namespace(key), value]) do
{:ok, 1} -> {:reply, true, {conn, :redis}}
{:ok, 0} -> {:reply, false, {conn, :redis}}
end
end
def handle_call({:move_member, from, to, value}, _from, {conn, :redis}) do
case Redix.command(conn, ["SMOVE", Redis.namespace(from), Redis.namespace(to), value]) do
{:ok, 1} -> {:reply, :ok, {conn, :redis}}
{:ok, 0} -> {:reply, :error, {conn, :redis}}
end
end
def handle_call({:add_member, key, value}, _from, {conn, :redis}) do
case Redix.command(conn, ["SADD", Redis.namespace(key), value]) do
{:ok, 1} -> {:reply, {:ok, value}, {conn, :redis}}
{:ok, 0} -> {:reply, {:error, :already_exists}, {conn, :redis}}
end
end
def handle_call({:remove_member, key, value}, _from, {conn, :redis}) do
case Redix.command(conn, ["SREM", Redis.namespace(key), value]) do
{:ok, 1} -> {:reply, :ok, {conn, :redis}}
{:ok, 0} -> {:reply, :error, {conn, :redis}}
end
end
def handle_call({:get, key}, _from, {conn, :redis}) do
with {:ok, type} <- Redix.command(conn, ["TYPE", Redis.namespace(key)]),
do: Redis.get(type, conn, key)
end
def handle_call({:remove, key}, _from, {conn, :redis}) do
{:ok, _} = Redix.command(conn, ["DEL", Redis.namespace(key)])
{:reply, :ok, {conn, :redis}}
end
def handle_call({:set, key, value}, _from, {conn, :redis}) when is_map(value) do
command = Enum.reduce(value, [Redis.namespace(key), "HMSET"], fn {k, v}, acc -> [[v, k] | acc] end)
case Redix.command(conn, command |> List.flatten() |> Enum.reverse()) do
{:ok, "OK"} -> {:reply, {:ok, value}, {conn, :redis}}
end
end
def handle_call({:set, key, value}, _from, {conn, :redis}) when is_binary(value) do
case Redix.command(conn, ["SETNX", Redis.namespace(key), value]) do
{:ok, 1} -> {:reply, {:ok, value}, {conn, :redis}}
{:ok, 0} -> {:reply, {:error, :already_set}, {conn, :redis}}
end
end
defp create_store do
if is_nil(Upstream.config(:redis_url)) do
{:ok, :ets.new(__MODULE__, [:set, :private, :named_table]), :ets}
else
{:ok, conn} = Redix.start_link(Upstream.config(:redis_url))
{:ok, conn, :redis}
end
end
end