Current section
2 Versions
Jump to
Current section
2 Versions
Compare versions
9
files changed
+459
additions
-16
deletions
| @@ -3,11 +3,13 @@ | |
| 3 3 | {<<"description">>,<<"A simple persistent queues">>}. |
| 4 4 | {<<"elixir">>,<<"~> 1.6">>}. |
| 5 5 | {<<"files">>, |
| 6 | - [<<"lib">>,<<"lib/simple_queue.ex">>,<<".formatter.exs">>,<<"mix.exs">>, |
| 7 | - <<"README.md">>]}. |
| 6 | + [<<"lib">>,<<"lib/queue.ex">>,<<"lib/reader.ex">>,<<"lib/simple_queue.ex">>, |
| 7 | + <<"lib/store.ex">>,<<"lib/unacked.ex">>,<<"lib/utils">>, |
| 8 | + <<"lib/utils/id.ex">>,<<"lib/writer.ex">>,<<".formatter.exs">>, |
| 9 | + <<"mix.exs">>,<<"README.md">>]}. |
| 8 10 | {<<"licenses">>,[<<"MIT">>]}. |
| 9 11 | {<<"links">>,[{<<"GitHub">>,<<"https://github.com/vtm9/simple_queue">>}]}. |
| 10 12 | {<<"maintainers">>,[<<"vtm">>]}. |
| 11 13 | {<<"name">>,<<"simple_queue">>}. |
| 12 14 | {<<"requirements">>,[]}. |
| 13 | - {<<"version">>,<<"0.1.0">>}. |
| 15 | + {<<"version">>,<<"0.1.1">>}. |
| @@ -0,0 +1,158 @@ | |
| 1 | + defmodule SimpleQueue.Queue do |
| 2 | + use GenServer |
| 3 | + |
| 4 | + alias __MODULE__ |
| 5 | + alias SimpleQueue.{Store, Unacked} |
| 6 | + alias SimpleQueue.Utils.Id |
| 7 | + |
| 8 | + defstruct( |
| 9 | + # queue data structure |
| 10 | + # in-memory queue head |
| 11 | + head: :queue.new(), |
| 12 | + # on-disk overflow queue tail |
| 13 | + store: nil, |
| 14 | + # in-flight map |
| 15 | + unacked: Unacked.new(), |
| 16 | + # number of message to keep in-memory |
| 17 | + capacity: 10, |
| 18 | + # message visibility timeout, time to keep message in-flight queue |
| 19 | + ttf: 1_000, |
| 20 | + # message time-to-sync |
| 21 | + tts: 100, |
| 22 | + expire_at: System.os_time() |
| 23 | + ) |
| 24 | + |
| 25 | + @timeout Application.get_env(:simple_queue, :timeout) |
| 26 | + |
| 27 | + # API |
| 28 | + def start_link(dir, params) do |
| 29 | + GenServer.start_link(__MODULE__, [dir, params]) |
| 30 | + end |
| 31 | + |
| 32 | + def add(queue_pid, payload) do |
| 33 | + GenServer.call(queue_pid, {:add, payload}, @timeout) |
| 34 | + end |
| 35 | + |
| 36 | + def get(queue_pid) do |
| 37 | + GenServer.call(queue_pid, :get, @timeout) |
| 38 | + end |
| 39 | + |
| 40 | + def ack(queue_pid, message_id) do |
| 41 | + GenServer.call(queue_pid, {:ack, message_id}, @timeout) |
| 42 | + end |
| 43 | + |
| 44 | + def reject(queue_pid, message_id) do |
| 45 | + GenServer.call(queue_pid, {:reject, message_id}, @timeout) |
| 46 | + end |
| 47 | + |
| 48 | + # SERVER |
| 49 | + |
| 50 | + def init([dir, params]) do |
| 51 | + {:ok, |
| 52 | + %Queue{ |
| 53 | + store: Store.new(dir) |
| 54 | + }} |
| 55 | + end |
| 56 | + |
| 57 | + def handle_call({:add, payload}, _from, state) do |
| 58 | + new_state = add_(payload, state) |
| 59 | + {:reply, :ok, new_state} |
| 60 | + end |
| 61 | + |
| 62 | + def handle_call(:get, _from, state) do |
| 63 | + {message, new_state} = get_(state) |
| 64 | + {:reply, message, new_state} |
| 65 | + end |
| 66 | + |
| 67 | + def handle_call({:ack, message_id}, _from, state) do |
| 68 | + new_state = ack_(message_id, state) |
| 69 | + {:reply, :ok, new_state} |
| 70 | + end |
| 71 | + |
| 72 | + def handle_call({:reject, message_id}, _from, state) do |
| 73 | + new_state = reject_(message_id, state) |
| 74 | + {:reply, :ok, new_state} |
| 75 | + end |
| 76 | + |
| 77 | + # def handle_call(:drop, _from, state) do |
| 78 | + # # {:reply, {:ok}, state} |
| 79 | + |
| 80 | + # {:stop, :normal, state} |
| 81 | + # end |
| 82 | + |
| 83 | + # PRIVATE |
| 84 | + |
| 85 | + defp add_(payload, state) do |
| 86 | + id = Id.new() |
| 87 | + |
| 88 | + new_state = |
| 89 | + state |
| 90 | + |> maybe_sync_on_disk_store() |
| 91 | + |
| 92 | + new_store = Store.add(pack(id, payload), new_state.store) |
| 93 | + %Queue{new_state | store: new_store} |
| 94 | + end |
| 95 | + |
| 96 | + defp get_(state) do |
| 97 | + new_state = |
| 98 | + state |
| 99 | + |> maybe_sync_on_disk_store |
| 100 | + |> maybe_shift_on_disk_store |
| 101 | + |
| 102 | + %Queue{head: head} = new_state |
| 103 | + |
| 104 | + case head do |
| 105 | + {[], []} -> |
| 106 | + {:empty, new_state} |
| 107 | + |
| 108 | + queue -> |
| 109 | + {{:value, message}, new_head} = :queue.out(head) |
| 110 | + new_state = add_to_unacked(message, %Queue{new_state | head: new_head}) |
| 111 | + {message, new_state} |
| 112 | + end |
| 113 | + end |
| 114 | + |
| 115 | + defp ack_(message_id, %Queue{unacked: unacked} = state) do |
| 116 | + new_unacked = Unacked.ack(message_id, unacked) |
| 117 | + %Queue{state | unacked: new_unacked} |
| 118 | + end |
| 119 | + |
| 120 | + defp reject_(message_id, %Queue{unacked: unacked} = state) do |
| 121 | + case Unacked.get(message_id, unacked) do |
| 122 | + nil -> |
| 123 | + state |
| 124 | + message -> |
| 125 | + new_state = add_(message.payload, state) |
| 126 | + new_unacked = Unacked.reject(message_id, new_state.unacked) |
| 127 | + %Queue{new_state | unacked: new_unacked} |
| 128 | + end |
| 129 | + end |
| 130 | + |
| 131 | + defp maybe_shift_on_disk_store(%Queue{head: {[], []}, store: store, capacity: capacity} = state) do |
| 132 | + {head, new_store} = Store.get(capacity + 1, store) |
| 133 | + %Queue{state | head: head, store: new_store} |
| 134 | + end |
| 135 | + |
| 136 | + defp maybe_shift_on_disk_store(state) do |
| 137 | + state |
| 138 | + end |
| 139 | + |
| 140 | + def maybe_sync_on_disk_store(%Queue{store: store, tts: tts, expire_at: expire_at} = state) do |
| 141 | + case System.os_time() do |
| 142 | + now when now > expire_at -> |
| 143 | + %Queue{state | store: Store.sync(store), expire_at: System.os_time() + tts} |
| 144 | + |
| 145 | + _ -> |
| 146 | + state |
| 147 | + end |
| 148 | + end |
| 149 | + |
| 150 | + defp add_to_unacked(message, %Queue{unacked: unacked} = state) do |
| 151 | + new_unacked = Unacked.add(message, unacked) |
| 152 | + %Queue{state | unacked: new_unacked} |
| 153 | + end |
| 154 | + |
| 155 | + defp pack(id, payload) do |
| 156 | + %{id: id, payload: payload} |
| 157 | + end |
| 158 | + end |
| @@ -0,0 +1,111 @@ | |
| 1 | + defmodule SimpleQueue.Reader do |
| 2 | + alias __MODULE__ |
| 3 | + |
| 4 | + @extension Application.get_env(:simple_queue, :timed_extension) |
| 5 | + @chunk Application.get_env(:simple_queue, :chunk) |
| 6 | + |
| 7 | + defstruct [:descriptor, :dir, :file, chunk: <<>>] |
| 8 | + |
| 9 | + def new(dir) do |
| 10 | + %Reader{dir: dir} |
| 11 | + end |
| 12 | + |
| 13 | + def get(%Reader{chunk: <<>>} = reader) do |
| 14 | + case open(reader) do |
| 15 | + :eof -> |
| 16 | + {:eof, reader} |
| 17 | + |
| 18 | + new_reader -> |
| 19 | + get(read(new_reader)) |
| 20 | + end |
| 21 | + end |
| 22 | + |
| 23 | + def get(%Reader{chunk: chunk} = reader) do |
| 24 | + case decode(chunk) do |
| 25 | + :noent -> |
| 26 | + get(read(reader)) |
| 27 | + |
| 28 | + {<<>>, new_chunk} -> |
| 29 | + get(%Reader{reader | chunk: new_chunk}) |
| 30 | + |
| 31 | + {message, new_chunk} -> |
| 32 | + {message, %Reader{reader | chunk: new_chunk}} |
| 33 | + end |
| 34 | + end |
| 35 | + |
| 36 | + # utility function to check length of file segments |
| 37 | + def length(%Reader{dir: dir}) do |
| 38 | + Reader.length(dir) |
| 39 | + end |
| 40 | + |
| 41 | + def length(dir) do |
| 42 | + pattern = Path.join([dir, "*", "sq", @extension]) |
| 43 | + |
| 44 | + case Path.wildcard(pattern) do |
| 45 | + [] -> |
| 46 | + 0 |
| 47 | + |
| 48 | + _ -> |
| 49 | + :infinity |
| 50 | + end |
| 51 | + end |
| 52 | + |
| 53 | + def open(%Reader{descriptor: nil, dir: dir} = reader) do |
| 54 | + pattern = Path.join([dir, "*", ["sq", @extension]]) |
| 55 | + |
| 56 | + case Path.wildcard(pattern) do |
| 57 | + [] -> |
| 58 | + :eof |
| 59 | + |
| 60 | + [head | _] -> |
| 61 | + {:ok, descriptor} = :file.open(head, [:raw, :binary, :read, {:read_ahead, @chunk}]) |
| 62 | + %Reader{reader | descriptor: descriptor, file: head} |
| 63 | + end |
| 64 | + end |
| 65 | + |
| 66 | + def open(reader), do: reader |
| 67 | + |
| 68 | + defp read(%Reader{descriptor: descriptor, chunk: head_chunk} = reader) do |
| 69 | + case :file.read(descriptor, @chunk) do |
| 70 | + :eof -> |
| 71 | + close(reader) |
| 72 | + |
| 73 | + {:ok, chunk} -> |
| 74 | + %Reader{reader | chunk: <<head_chunk::binary, chunk::binary>>} |
| 75 | + end |
| 76 | + end |
| 77 | + |
| 78 | + # close any open file and rotate active head |
| 79 | + defp close(%Reader{descriptor: nil} = reader), do: reader |
| 80 | + |
| 81 | + defp close(%Reader{descriptor: descriptor, file: file} = reader) do |
| 82 | + :ok = :file.close(descriptor) |
| 83 | + :ok = :file.delete(file) |
| 84 | + :file.del_dir(Path.dirname(file)) |
| 85 | + %Reader{reader | descriptor: nil, file: nil, chunk: <<>>} |
| 86 | + end |
| 87 | + |
| 88 | + # decode message from memory buffer |
| 89 | + defp decode(<<0::16, len::32, hash::32, tail::binary>>) do |
| 90 | + case byte_size(tail) do |
| 91 | + x when x < len -> |
| 92 | + :noent |
| 93 | + |
| 94 | + _ -> |
| 95 | + <<binary_message::binary-size(len), rest::binary>> = tail |
| 96 | + |
| 97 | + case :erlang.crc32(binary_message) do |
| 98 | + hash -> {:erlang.binary_to_term(binary_message), rest} |
| 99 | + _ -> {<<>>, rest} |
| 100 | + end |
| 101 | + end |
| 102 | + end |
| 103 | + |
| 104 | + defp decode(x) when byte_size(x) < 64 do |
| 105 | + :noent |
| 106 | + end |
| 107 | + |
| 108 | + defp decode(<<_::8, tail::binary>>) do |
| 109 | + decode(tail) |
| 110 | + end |
| 111 | + end |
| @@ -1,18 +1,27 @@ | |
| 1 1 | defmodule SimpleQueue do |
| 2 | - @moduledoc """ |
| 3 | - Documentation for SimpleQueue. |
| 4 | - """ |
| 2 | + alias SimpleQueue.Queue |
| 5 3 | |
| 6 | - @doc """ |
| 7 | - Hello world. |
| 4 | + def new(dir) do |
| 5 | + new(dir, %{}) |
| 6 | + end |
| 8 7 | |
| 9 | - ## Examples |
| 8 | + def new(dir, params) do |
| 9 | + Queue.start_link(dir, params) |
| 10 | + end |
| 10 11 | |
| 11 | - iex> SimpleQueue.hello |
| 12 | - :world |
| 12 | + def add(queue_pid, message) do |
| 13 | + Queue.add(queue_pid, message) |
| 14 | + end |
| 13 15 | |
| 14 | - """ |
| 15 | - def hello do |
| 16 | - :world |
| 16 | + def get(queue_pid) do |
| 17 | + Queue.get(queue_pid) |
| 18 | + end |
| 19 | + |
| 20 | + def ack(queue_pid, msg_id) do |
| 21 | + Queue.ack(queue_pid, msg_id) |
| 22 | + end |
| 23 | + |
| 24 | + def reject(queue_pid, msg_id) do |
| 25 | + Queue.reject(queue_pid, msg_id) |
| 17 26 | end |
| 18 27 | end |
| @@ -0,0 +1,52 @@ | |
| 1 | + defmodule SimpleQueue.Store do |
| 2 | + alias __MODULE__ |
| 3 | + alias SimpleQueue.{Reader, Writer} |
| 4 | + |
| 5 | + defstruct [:writer, :reader, :len] |
| 6 | + |
| 7 | + def new(dir) do |
| 8 | + %Store{ |
| 9 | + writer: Writer.new(dir), |
| 10 | + reader: Reader.new(dir), |
| 11 | + len: Reader.length(dir) |
| 12 | + } |
| 13 | + end |
| 14 | + |
| 15 | + def add(message, %Store{writer: writer, len: len} = store) do |
| 16 | + %Store{store | writer: Writer.add(message, writer), len: inc(len)} |
| 17 | + end |
| 18 | + |
| 19 | + def get(n, store) do |
| 20 | + get(n, :queue.new(), store) |
| 21 | + end |
| 22 | + |
| 23 | + def get(0, acc, store) do |
| 24 | + {acc, store} |
| 25 | + end |
| 26 | + |
| 27 | + def get(n, acc, %Store{reader: reader, writer: writer, len: len} = store) do |
| 28 | + case Reader.get(reader) do |
| 29 | + {:eof, new_reader} -> |
| 30 | + case Writer.length(writer) do |
| 31 | + 0 -> |
| 32 | + {acc, %Store{store | reader: new_reader, len: 0}} |
| 33 | + |
| 34 | + _ -> |
| 35 | + {acc, %Store{store | reader: new_reader, len: dec(len)}} |
| 36 | + end |
| 37 | + |
| 38 | + {message, new_reader} -> |
| 39 | + get(n - 1, :queue.in(message, acc), %Store{store | reader: new_reader, len: dec(len)}) |
| 40 | + end |
| 41 | + end |
| 42 | + |
| 43 | + def sync(%Store{writer: writer} = store) do |
| 44 | + %Store{store | writer: Writer.close_and_rotate(writer)} |
| 45 | + end |
| 46 | + |
| 47 | + defp inc(:infinity), do: :infinity |
| 48 | + defp inc(x), do: x + 1 |
| 49 | + |
| 50 | + defp dec(:infinity), do: :infinity |
| 51 | + defp dec(x), do: x - 1 |
| 52 | + end |
Loading more files…