Fix to use ETS based transaction_id counter
This commit is contained in:
parent
25b0271ab1
commit
57611f9dc0
3 changed files with 28 additions and 17 deletions
|
|
@ -11,7 +11,7 @@ defmodule Tres.SecureChannel do
|
|||
|
||||
@process_flags [
|
||||
trap_exit: true,
|
||||
message_queue_data: :on_heap
|
||||
message_queue_data: :off_heap
|
||||
]
|
||||
|
||||
@supported_version 4
|
||||
|
|
@ -379,7 +379,7 @@ defmodule Tres.SecureChannel do
|
|||
|
||||
defp initiate_hello_handshake(state_data) do
|
||||
send_hello(state_data)
|
||||
ref = :erlang.send_after(@hello_handshake_timeout, self(), :hello_timeout)
|
||||
ref = Process.send_after(self(), :hello_timeout, @hello_handshake_timeout)
|
||||
{:keep_state, %{state_data | timer_ref: ref}}
|
||||
end
|
||||
|
||||
|
|
@ -397,7 +397,7 @@ defmodule Tres.SecureChannel do
|
|||
defp initiate_features_handshake(state_data) do
|
||||
new_xid = State.increment_transaction_id(state_data.xid)
|
||||
send_features(new_xid, state_data)
|
||||
ref = :erlang.send_after(@features_handshake_timeout, self(), :features_timeout)
|
||||
ref = Process.send_after(self(), :features_timeout, @features_handshake_timeout)
|
||||
{:keep_state, %{state_data | timer_ref: ref}}
|
||||
end
|
||||
|
||||
|
|
@ -445,7 +445,7 @@ defmodule Tres.SecureChannel do
|
|||
end
|
||||
|
||||
defp start_periodic_idle_check do
|
||||
:erlang.send_after(@ping_interval, self(), :idle_check)
|
||||
Process.send_after(self(), :idle_check, @ping_interval)
|
||||
end
|
||||
|
||||
defp maybe_ping(state_data) do
|
||||
|
|
@ -463,10 +463,10 @@ defmodule Tres.SecureChannel do
|
|||
false
|
||||
end
|
||||
|
||||
defp send_ping(%State{xid: x_agent} = state_data) do
|
||||
xid = State.increment_transaction_id(x_agent)
|
||||
defp send_ping(%State{xid: table_ref} = state_data) do
|
||||
xid = State.increment_transaction_id(table_ref)
|
||||
send_echo_request(xid, "", state_data)
|
||||
ping_ref = :erlang.send_after(@ping_timeout, self(), :ping_timeout)
|
||||
ping_ref = Process.send_after(self(), :ping_timeout, @ping_timeout)
|
||||
%{state_data | ping_timer_ref: ping_ref, ping_xid: xid}
|
||||
end
|
||||
|
||||
|
|
@ -516,7 +516,7 @@ defmodule Tres.SecureChannel do
|
|||
defp maybe_cancel_timer(ref) when not is_reference(ref), do: :ok
|
||||
|
||||
defp maybe_cancel_timer(ref) do
|
||||
:erlang.cancel_timer(ref)
|
||||
Process.cancel_timer(ref)
|
||||
:ok
|
||||
end
|
||||
|
||||
|
|
|
|||
|
|
@ -29,7 +29,7 @@ defmodule Tres.SecureChannelState do
|
|||
socket = Keyword.get(options, :socket)
|
||||
transport = Keyword.get(options, :transport)
|
||||
{:ok, {ip_addr, port}} = :inet.peername(socket)
|
||||
{:ok, xid_agent} = Agent.start_link(fn -> 0 end)
|
||||
{:ok, table_ref} = create_counter()
|
||||
kv_ref = XACT_KV.create()
|
||||
|
||||
%SecureChannelState{
|
||||
|
|
@ -38,20 +38,31 @@ defmodule Tres.SecureChannelState do
|
|||
transport: transport,
|
||||
ip_addr: :inet.ntoa(ip_addr),
|
||||
port: port,
|
||||
xid: xid_agent,
|
||||
xid: table_ref,
|
||||
xact_kv_ref: kv_ref
|
||||
}
|
||||
end
|
||||
|
||||
def increment_transaction_id(xid_agent) do
|
||||
Agent.get_and_update(xid_agent, &{&1 + 1, &1 + 1})
|
||||
def increment_transaction_id(table_ref) do
|
||||
:ets.update_counter(table_ref, :datapath_xid, {2, 1, 0xffffffff, 0})
|
||||
end
|
||||
|
||||
def set_transaction_id(xid_agent, xid) do
|
||||
Agent.update(xid_agent, fn _ -> xid end)
|
||||
def set_transaction_id(table_ref, xid) do
|
||||
:ets.insert(table_ref, {:datapath_xid, xid})
|
||||
end
|
||||
|
||||
def get_transaction_id(xid_agent) do
|
||||
Agent.get(xid_agent, & &1)
|
||||
def get_transaction_id(table_ref) do
|
||||
case :ets.lookup(table_ref, :datapath_xid) do
|
||||
[{_, xid} | _] -> xid
|
||||
end
|
||||
end
|
||||
|
||||
# private functions
|
||||
|
||||
@spec create_counter() :: reference()
|
||||
defp create_counter do
|
||||
table_ref = :ets.new(:xid_counter, [:set, :private])
|
||||
_ = :ets.insert(table_ref, {:datapath_xid, 0})
|
||||
{:ok, table_ref}
|
||||
end
|
||||
end
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue