Skip to content

Commit ac33b5a

Browse files
committed
Make anti_swap lock correctly exclusive
1 parent 16ee77d commit ac33b5a

1 file changed

Lines changed: 35 additions & 22 deletions

File tree

src/ets_ring_buffer_ro.erl

Lines changed: 35 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -70,30 +70,30 @@
7070
anti_swap_lock = 0 :: non_neg_integer() | '_',
7171
swapping_lock = 0 :: non_neg_integer() | '_',
7272
ring_size = 0 :: ring_size() | '_',
73-
read_loc = 0 :: ring_loc() | 0 | '_' % 0 signifies a buffer that has never been read
73+
read_loc = 0 :: ring_loc() | 0 | '_' % 0: buffer has never been read
7474
}).
7575

7676
-define(RING_SIZE, {#ring_ro_metadata.ring_size, 0}).
7777
-define(RING_BUFFER, {#ring_ro_metadata.buffer, 0}).
7878
-define(LAST_READ_LOC, {#ring_ro_metadata.read_loc, 0}).
7979
-define(BUMP_GENERATION, {#ring_ro_metadata.generation, 1}).
8080

81-
-define(READ_ANTI_SWAP, {#ring_ro_metadata.anti_swap_lock, 0}).
8281
-define(LOCK_ANTI_SWAP, {#ring_ro_metadata.anti_swap_lock, 1}).
83-
-define(UNLOCK_ANTI_SWAP, {#ring_ro_metadata.anti_swap_lock, -1}).
82+
-define(UNLOCK_ANTI_SWAP, {#ring_ro_metadata.anti_swap_lock, 1, 0, 0}).
8483

8584
-define(RESERVE_READ_LOC(__Size), [?RING_BUFFER, {#ring_ro_metadata.read_loc, 1, __Size, 1}]).
8685
-define(RESERVE_READ_ALL_LOCS, [?RING_BUFFER, ?RING_SIZE, ?LAST_READ_LOC]).
8786

8887
meta_key(Ring_Name) -> {meta, Ring_Name}.
8988
make_meta(Ring_Name, Ring_Table_Id, Size) ->
90-
#ring_ro_metadata{name=meta_key(Ring_Name), buffer=Ring_Table_Id, created=os:timestamp(), ring_size=Size}.
89+
#ring_ro_metadata{name = meta_key(Ring_Name), buffer = Ring_Table_Id,
90+
created = os:timestamp(), ring_size = Size}.
9191

9292
%% Convert record to proplist.
9393
make_ring_proplist(#ring_ro_metadata{name={meta, Name}, generation=Gen, buffer=Buffer,
9494
created=Created, ring_size=Size, read_loc=Read_Loc}) ->
95-
[{name, Name}, {generation, Gen}, {buffer, Buffer},
96-
{created, Created}, {ring_size, Size}, {read_loc, Read_Loc}].
95+
[{name, Name}, {generation, Gen}, {buffer, Buffer},
96+
{created, Created}, {ring_size, Size}, {read_loc, Read_Loc}].
9797

9898
%% Match specs for buffers.
9999
all_rings() ->
@@ -111,8 +111,11 @@ get_ring_size (Ring_Name) -> get_ring_metadata_field(Ring_Name, ?RING_SIZE)
111111
get_ring_last_read (Ring_Name) -> get_ring_metadata_field(Ring_Name, ?LAST_READ_LOC).
112112
get_ring_read_all (Ring_Name) -> get_ring_metadata_field(Ring_Name, ?RESERVE_READ_ALL_LOCS).
113113

114-
get_ring_size_lock (Ring_Name) -> get_ring_metadata_field(Ring_Name, [?RING_SIZE, ?LOCK_ANTI_SWAP]).
115114
unlock_anti_swap (Ring_Name) -> get_ring_metadata_field(Ring_Name, ?UNLOCK_ANTI_SWAP).
115+
get_ring_size_lock (Ring_Name) ->
116+
{Size, Lock} = get_ring_metadata_field(Ring_Name, [?RING_SIZE, ?LOCK_ANTI_SWAP]),
117+
{Size, Lock =:= 1}.
118+
116119

117120
%% Reserve the next read location for the calling process to read.
118121
%% Since the pointer starts at 0, the reserved location is after increment
@@ -126,15 +129,23 @@ unlock_anti_swap (Ring_Name) -> get_ring_metadata_field(Ring_Name, ?UNLOCK_ANT
126129
%% have obtained a read location from the single metadata record which
127130
%% will be highly contended itself.
128131
get_ring_next_read(Ring_Name) ->
129-
%% Unfortunately, must fetch size before increment so we can properly wrap around the read pointer.
130-
%% This creates a race that could fail when a ring buffer is replaced, so we lock against swaps.
131-
try get_ring_size_lock(Ring_Name) of
132-
[0, _] -> {undefined, 0, 0};
133-
[Size, _] -> [Ring_Buffer, Read_Loc] = get_ring_metadata_field(Ring_Name, ?RESERVE_READ_LOC(Size)),
134-
{Ring_Buffer, Size, Read_Loc}
132+
%% Unfortunately, must fetch size before increment so we can properly
133+
%% wrap around the read pointer. This creates a race that could fail
134+
%% when a ring buffer is replaced, so we lock against swaps.
135+
try obtain_ring_size_lock(Ring_Name)
135136
after _ = unlock_anti_swap(Ring_Name)
136137
end.
137138

139+
obtain_ring_size_lock(Ring_Name) ->
140+
case get_ring_size_lock(Ring_Name) of
141+
{0, _} -> {undefined, 0, 0};
142+
{ _, false} -> erlang:yield(),
143+
obtain_ring_size_lock(Ring_Name);
144+
{Size, true} -> [Ring_Buffer, Read_Loc]
145+
= get_ring_metadata_field(Ring_Name, ?RESERVE_READ_LOC(Size)),
146+
{Ring_Buffer, Size, Read_Loc}
147+
end.
148+
138149
%% Uses update_counter to maintain the write_concurrency lock.
139150
get_ring_metadata_field(Ring_Name, Update_Cmd) ->
140151
try ets:update_counter(?RING_RO_TABLE, meta_key(Ring_Name), Update_Cmd)
@@ -160,8 +171,8 @@ get_ring_metadata(Ring_Name) ->
160171
data :: ring_data()
161172
}).
162173

163-
ring_key(Name, Loc) -> {Name, Loc}.
164-
make_ring_data(Name, Loc, Data) -> #ring_ro_data{key=ring_key(Name, Loc), data=Data}.
174+
ring_key (Name, Loc) -> {Name, Loc}.
175+
make_ring_data (Name, Loc, Data) -> #ring_ro_data{key=ring_key(Name, Loc), data=Data}.
165176

166177

167178
%%%------------------------------------------------------------------------------
@@ -246,7 +257,7 @@ clear(Ring_Name) when is_atom(Ring_Name) ->
246257
clear(_Ring_Name, false) -> false;
247258
clear(_Ring_Name, missing) -> false;
248259
clear( Ring_Name, #ring_ro_metadata{anti_swap_lock=0, buffer=Ring_Buffer}) ->
249-
%% Clear the pointers in the metadata first, for mmediate effect...
260+
%% Clear the pointers in the metadata first, for immediate effect...
250261
try reset_metadata(Ring_Name, undefined, [])
251262
after epocxy_ets_fsm:delete_ets_table(Ring_Buffer)
252263
end;
@@ -256,7 +267,7 @@ clear( Ring_Name, #ring_ro_metadata{}) ->
256267
clear(Ring_Name).
257268

258269
%% @doc
259-
%% Delete the ring metadata, then delete the ring data buffer. Thie function
270+
%% Delete the ring metadata, then delete the ring data buffer. This function
260271
%% does a read on the ring metadata table so it will incur a slow lock
261272
%% penalty, but it is used infrequently.
262273
%% @end
@@ -283,8 +294,8 @@ delete( Ring_Name, #ring_ro_metadata{}) ->
283294
read(Ring_Name) when is_atom(Ring_Name) ->
284295
read(Ring_Name, get_ring_next_read(Ring_Name)).
285296

286-
read(Ring_Name, false ) -> {error, {buffer_is_empty, Ring_Name}};
287-
read(Ring_Name, {undefined, 0, 0} ) -> {error, {buffer_is_empty, Ring_Name}};
297+
read(Ring_Name, false ) -> {error, {buffer_is_empty, Ring_Name}};
298+
read(Ring_Name, {undefined, 0, 0}) -> {error, {buffer_is_empty, Ring_Name}};
288299
read(Ring_Name, {Ring_Buffer, _Size, Location})
289300
when is_integer(Location), Location > 0 ->
290301
read_value(Ring_Name, Ring_Buffer, Location).
@@ -316,13 +327,14 @@ ring_size(Ring_Name) when is_atom(Ring_Name) ->
316327

317328
create_ring(Ring_Name, Ring_Buffer, Ring_Values) ->
318329
insert_values(Ring_Name, Ring_Buffer, Ring_Values, 1),
319-
case ets:insert_new(?RING_RO_TABLE, make_meta(Ring_Name, Ring_Buffer, length(Ring_Values))) of
330+
Num_Values = length(Ring_Values),
331+
case ets:insert_new(?RING_RO_TABLE, make_meta(Ring_Name, Ring_Buffer, Num_Values)) of
320332
false -> epocxy_ets_fsm:delete_ets_table(Ring_Buffer),
321333
false;
322334
true -> true
323335
end.
324336

325-
replace_ring(Ring_Name, New_Ring_Buffer, Ring_Values) ->
337+
replace_ring( Ring_Name, New_Ring_Buffer, Ring_Values) ->
326338
true = insert_values(Ring_Name, New_Ring_Buffer, Ring_Values, 1),
327339
replace_ring(Ring_Name, New_Ring_Buffer, Ring_Values, get_ring_metadata(Ring_Name)).
328340

@@ -331,12 +343,13 @@ replace_ring(_Ring_Name, New_Ring_Buffer, _Ring_Values, missing) ->
331343
false;
332344
replace_ring( Ring_Name, New_Ring_Buffer, Ring_Values,
333345
#ring_ro_metadata{anti_swap_lock=0, buffer=Old_Ring_Buffer}) ->
334-
try reset_metadata (Ring_Name, New_Ring_Buffer, Ring_Values)
346+
try reset_metadata(Ring_Name, New_Ring_Buffer, Ring_Values)
335347
after Old_Ring_Buffer =:= undefined
336348
orelse epocxy_ets_fsm:delete_ets_table(Old_Ring_Buffer)
337349
end;
338350
%% Anti-swap lock is set, wait for it to clear.
339351
replace_ring( Ring_Name, New_Ring_Buffer, Ring_Values, #ring_ro_metadata{}) ->
352+
erlang:yield(),
340353
replace_ring(Ring_Name, New_Ring_Buffer, Ring_Values, get_ring_metadata(Ring_Name)).
341354

342355
reset_metadata(Ring_Name, New_Ring_Buffer, Ring_Values) ->

0 commit comments

Comments
 (0)