defmodule SuperCache.Queue do @moduledoc """ A queue based on cache. Data in queue stored in cache. Can handle multiple queue with different name. A queue is a FIFO data structure. This is global data, any process can access to queue data. Need to start SuperCache.start!/1 before using this module. Ex: ```elixir alias SuperCache.Queue SuperCache.start!() Queue.add("my_queue", "Hello") Queue.out("my_queue") # => "Hello" ``` """ alias SuperCache.Storage alias SuperCache.Partition require Logger ### Api ### @doc """ Add value to queue with name is queue_name. Queue is a FIFO data structure. If queue_name is not existed, it will be created. """ @spec add(any, any) :: true def add(queue_name, value) do part = Partition.get_partition(queue_name) queue_in(part, queue_name, value) end @doc """ Pop value from queue with name is queue_name. If queue_name is not existed or no data, it will return default value. """ @spec out(any, any) :: any def out(queue_name, default \\ nil) do part = Partition.get_partition(queue_name) queue_out(part, queue_name, default) end @doc """ Peak value from queue with name is queue_name. Data still in queue after peak. If queue_name is not existed or no data, it will return default value. """ @spec peak(any, any) :: any def peak(queue_name, default \\ nil) do part = Partition.get_partition(queue_name) queue_peak(part, queue_name, default) end @doc """ Get all items in queue. If queue_name is not existed, it will return []. If queue is empty, it will return []. """ @spec get_all(any()) :: list() def get_all(queue_name) do partition = Partition.get_partition(queue_name) to_list(partition, queue_name) end @doc """ Count number of data in queue with name is queue_name. If queue_name is not existed, it will return 0. """ @spec count(any()) :: number() def count(queue_name) do partition = Partition.get_partition(queue_name) case Storage.get({:queue, :tail, queue_name}, partition) do # stack is not initialized [] -> 0 [{_, tail}] -> case Storage.get({:queue, :head, queue_name}, partition) do # in updating state [] -> count(queue_name) [{_, head}] -> tail - head + 1 end end end ### Internal functions ### defp queue_in(partition, queue_name, value) do case Storage.take({:queue, :tail, queue_name}, partition) do # stack is not initialized [] -> case Storage.get({:queue, :updating, queue_name}, partition) do # stack is not initialized [] -> queue_init(queue_name) queue_in(partition, queue_name, value) # queue is updating _ -> # wait for stack is ready :erlang.yield() queue_in(partition, queue_name, value) end [{_, 0}] -> Storage.put({{:queue, :updating, queue_name}, true}, partition) Storage.delete({:queue, :tail, queue_name}, partition) Storage.delete({:queue, :head, queue_name}, partition) Storage.put({{:queue, queue_name, 1}, value}, partition) Storage.put({{:queue, :head, queue_name}, 1}, partition) Storage.put({{:queue, :tail, queue_name}, 1}, partition) Storage.delete({:queue, :updating, queue_name}, partition) Logger.debug( "super_cache, queue, push value: #{inspect(value)} to queue #{inspect(queue_name)}" ) true [{_, counter}] -> next_counter = counter + 1 Storage.put({{:queue, :updating, queue_name}, true}, partition) Storage.delete({:queue, :tail, queue_name}, partition) Storage.put({{:queue, queue_name, next_counter}, value}, partition) Storage.put({{:queue, :tail, queue_name}, next_counter}, partition) Storage.delete({:queue, :updating, queue_name}, partition) Logger.debug( "super_cache, queue, push value: #{inspect(value)} to queue #{inspect(queue_name)}" ) true end end defp queue_out(partition, queue_name, default) do case Storage.take({:queue, :head, queue_name}, partition) do # maybe stack is not initialized/updating process [] -> case Storage.get({:queue, :updating, queue_name}, partition) do # stack is not initialized [] -> default # queue is updating _ -> # wait for stack is ready :erlang.yield() queue_out(partition, queue_name, default) end # queue is empty [{_, 0}] -> default [{_, counter}] -> next_counter = counter + 1 Storage.put({{:queue, :updating, queue_name}, true}, partition) value = case Storage.take({:queue, queue_name, counter}, partition) do [] -> # no more data in queue, reset queue Storage.delete({:queue, :tail, queue_name}, partition) Storage.delete({:queue, :head, queue_name}, partition) Storage.put({{:queue, :head, queue_name}, 0}, partition) Storage.put({{:queue, :tail, queue_name}, 0}, partition) default [{_, value}] -> Storage.delete({:queue, :head, queue_name}, partition) Storage.put({{:queue, :head, queue_name}, next_counter}, partition) value end Storage.delete({:queue, :updating, queue_name}, partition) Logger.debug(fn -> "super_cache, queue, get value: #{inspect(value)} from queue #{inspect(queue_name)}" end) value end end defp queue_peak(partition, queue_name, default) do case Storage.get({:queue, :head, queue_name}, partition) do # stack is not initialized [] -> case Storage.get({:queue, :updating, queue_name}, partition) do # stack is not initialized [] -> default # queue is updating _ -> # wait for stack is ready :erlang.yield() queue_peak(partition, queue_name, default) end # queue is empty [{_, 0}] -> default [{_, counter}] -> case Storage.get({:queue, queue_name, counter}, partition) do [] -> default [{_, value}] -> value end end end defp queue_init(queue_name) do Logger.debug("super_cache, stack, init stack: #{inspect(queue_name)}") partition = Partition.get_partition(queue_name) Storage.put({{:queue, :updating, queue_name}, true}, partition) Storage.put({{:queue, :head, queue_name}, 0}, partition) Storage.put({{:queue, :tail, queue_name}, 0}, partition) Storage.delete({:queue, :updating, queue_name}, partition) end defp to_list(partition, queue_name) do case Storage.take({:queue, :head, queue_name}, partition) do # stack is not initialized [] -> case Storage.get({:queue, :updating, queue_name}, partition) do # stack is not initialized [] -> [] # queue is updating _ -> # wait for stack is ready :erlang.yield() to_list(partition, queue_name) end # queue is empty [{_, 0}] -> [] [{_, first}] -> Storage.put({{:queue, :updating, queue_name}, true}, partition) [{_, last}] = Storage.take({:queue, :tail, queue_name}, partition) value = Enum.reduce(first..last, [], fn i, acc -> case Storage.take({:queue, queue_name, i}, partition) do [] -> acc [{_, value}] -> [value | acc] end end) # reset queue Storage.delete({:queue, :head, queue_name}, partition) Storage.delete({:queue, :tail, queue_name}, partition) Storage.put({{:queue, :head, queue_name}, 0}, partition) Storage.put({{:queue, :tail, queue_name}, 0}, partition) Storage.delete({:queue, :updating, queue_name}, partition) Logger.debug(fn -> "super_cache, queue, get all length: #{inspect(length(value))} from queue #{inspect(queue_name)}" end) Enum.reverse(value) end end end