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
8887meta_key (Ring_Name ) -> {meta , Ring_Name }.
8988make_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.
9393make_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.
9999all_rings () ->
@@ -111,8 +111,11 @@ get_ring_size (Ring_Name) -> get_ring_metadata_field(Ring_Name, ?RING_SIZE)
111111get_ring_last_read (Ring_Name ) -> get_ring_metadata_field (Ring_Name , ? LAST_READ_LOC ).
112112get_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 ]).
115114unlock_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.
128131get_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.
139150get_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) ->
246257clear (_Ring_Name , false ) -> false ;
247258clear (_Ring_Name , missing ) -> false ;
248259clear ( 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{}) ->
283294read (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 }};
288299read (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
317328create_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 ;
332344replace_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.
339351replace_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
342355reset_metadata (Ring_Name , New_Ring_Buffer , Ring_Values ) ->
0 commit comments