Current section

2 Versions

Jump to

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…