Packages

Grow-only and positive-negative counters for Gleam — conflict-free replicated counter types

Current section

Files

Jump to
lattice_counters src lattice_counters@pn_counter.erl
Raw

src/lattice_counters@pn_counter.erl

-module(lattice_counters@pn_counter).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/lattice_counters/pn_counter.gleam").
-export([new/1, try_increment/2, increment/2, try_decrement/2, decrement/2, value/1, merge/2, to_json/1, from_json/1]).
-export_type([p_n_counter/0, update_error/0]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
?MODULEDOC(
" A positive-negative counter (PN-Counter) CRDT.\n"
"\n"
" Supports both increment and decrement operations by pairing two G-Counters:\n"
" one tracking increments and one tracking decrements. The value is the\n"
" difference between the two totals. Merge delegates to G-Counter merge on\n"
" each half independently.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" import lattice_core/replica_id\n"
" import lattice_counters/pn_counter\n"
"\n"
" let counter = pn_counter.new(replica_id.new(\"node-a\"))\n"
" |> pn_counter.increment(10)\n"
" |> pn_counter.decrement(3)\n"
" pn_counter.value(counter) // -> 7\n"
" ```\n"
).
-opaque p_n_counter() :: {p_n_counter,
lattice_counters@g_counter:g_counter(),
lattice_counters@g_counter:g_counter()}.
-type update_error() :: {negative_delta, integer()}.
-file("src/lattice_counters/pn_counter.gleam", 42).
?DOC(
" Create a new PN-Counter for the given replica.\n"
"\n"
" Returns a fresh counter with a zero value. Both inner G-Counters are\n"
" initialized with `replica_id` as their node identifier.\n"
).
-spec new(lattice_core@replica_id:replica_id()) -> p_n_counter().
new(Replica_id) ->
{p_n_counter,
lattice_counters@g_counter:new(Replica_id),
lattice_counters@g_counter:new(Replica_id)}.
-file("src/lattice_counters/pn_counter.gleam", 62).
?DOC(
" Safely increment the counter by `delta`.\n"
"\n"
" Returns `Error(NegativeDelta(delta))` if `delta` is negative.\n"
).
-spec try_increment(p_n_counter(), integer()) -> {ok, p_n_counter()} |
{error, update_error()}.
try_increment(Counter, Delta) ->
{p_n_counter, Positive, Negative} = Counter,
case lattice_counters@g_counter:try_increment(Positive, Delta) of
{ok, Updated_positive} ->
{ok, {p_n_counter, Updated_positive, Negative}};
{error, {negative_delta, D}} ->
{error, {negative_delta, D}}
end.
-file("src/lattice_counters/pn_counter.gleam", 54).
?DOC(
" Increment the counter by `delta`.\n"
"\n"
" Adds `delta` to the positive G-Counter. `delta` should be a non-negative\n"
" integer; the positive G-Counter is grow-only so passing a negative value\n"
" violates the invariant.\n"
).
-spec increment(p_n_counter(), integer()) -> p_n_counter().
increment(Counter, Delta) ->
Updated@1 = case try_increment(Counter, Delta) of
{ok, Updated} -> Updated;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"lattice_counters/pn_counter"/utf8>>,
function => <<"increment"/utf8>>,
line => 55,
value => _assert_fail,
start => 1766,
'end' => 1820,
pattern_start => 1777,
pattern_end => 1788})
end,
Updated@1.
-file("src/lattice_counters/pn_counter.gleam", 87).
?DOC(
" Safely decrement the counter by `delta`.\n"
"\n"
" Returns `Error(NegativeDelta(delta))` if `delta` is negative.\n"
).
-spec try_decrement(p_n_counter(), integer()) -> {ok, p_n_counter()} |
{error, update_error()}.
try_decrement(Counter, Delta) ->
{p_n_counter, Positive, Negative} = Counter,
case lattice_counters@g_counter:try_increment(Negative, Delta) of
{ok, Updated_negative} ->
{ok, {p_n_counter, Positive, Updated_negative}};
{error, {negative_delta, D}} ->
{error, {negative_delta, D}}
end.
-file("src/lattice_counters/pn_counter.gleam", 79).
?DOC(
" Decrement the counter by `delta`.\n"
"\n"
" Adds `delta` to the negative G-Counter (which reduces the visible value).\n"
" `delta` should be a non-negative integer; the negative G-Counter is\n"
" grow-only so passing a negative value violates the invariant.\n"
).
-spec decrement(p_n_counter(), integer()) -> p_n_counter().
decrement(Counter, Delta) ->
Updated@1 = case try_decrement(Counter, Delta) of
{ok, Updated} -> Updated;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"lattice_counters/pn_counter"/utf8>>,
function => <<"decrement"/utf8>>,
line => 80,
value => _assert_fail,
start => 2633,
'end' => 2687,
pattern_start => 2644,
pattern_end => 2655})
end,
Updated@1.
-file("src/lattice_counters/pn_counter.gleam", 103).
?DOC(
" Get the current value of the counter.\n"
"\n"
" Returns the sum of positive increments minus the sum of negative\n"
" decrements observed across all replicas.\n"
).
-spec value(p_n_counter()) -> integer().
value(Counter) ->
{p_n_counter, Positive, Negative} = Counter,
lattice_counters@g_counter:value(Positive) - lattice_counters@g_counter:value(
Negative
).
-file("src/lattice_counters/pn_counter.gleam", 115).
?DOC(
" Merge two PN-Counters.\n"
"\n"
" Merges the positive G-Counters and negative G-Counters independently using\n"
" pairwise maximum. The result's `self_id` is taken from `a`'s positive\n"
" G-Counter.\n"
"\n"
" This operation is commutative, associative, and idempotent.\n"
).
-spec merge(p_n_counter(), p_n_counter()) -> p_n_counter().
merge(A, B) ->
{p_n_counter, Positive_a, Negative_a} = A,
{p_n_counter, Positive_b, Negative_b} = B,
{p_n_counter,
lattice_counters@g_counter:merge(Positive_a, Positive_b),
lattice_counters@g_counter:merge(Negative_a, Negative_b)}.
-file("src/lattice_counters/pn_counter.gleam", 131).
?DOC(
" Encode a PN-Counter as a self-describing JSON value.\n"
"\n"
" Produces an envelope with `type`, `v` (schema version), and `state`.\n"
" Format: `{\"type\": \"pn_counter\", \"v\": 1, \"state\": {\"positive\": {...}, \"negative\": {...}}}`\n"
"\n"
" Use `from_json` to decode the result back into a `PNCounter`.\n"
).
-spec to_json(p_n_counter()) -> gleam@json:json().
to_json(Counter) ->
{p_n_counter, Positive, Negative} = Counter,
{Pos_dict, Pos_id} = lattice_counters@g_counter:to_parts(Positive),
{Neg_dict, Neg_id} = lattice_counters@g_counter:to_parts(Negative),
gleam@json:object(
[{<<"type"/utf8>>, gleam@json:string(<<"pn_counter"/utf8>>)},
{<<"v"/utf8>>, gleam@json:int(1)},
{<<"state"/utf8>>,
gleam@json:object(
[{<<"positive"/utf8>>,
gleam@json:object(
[{<<"self_id"/utf8>>,
lattice_core@replica_id:to_json(Pos_id)},
{<<"counts"/utf8>>,
gleam@json:dict(
Pos_dict,
fun(K) ->
lattice_core@replica_id:to_string(
K
)
end,
fun gleam@json:int/1
)}]
)},
{<<"negative"/utf8>>,
gleam@json:object(
[{<<"self_id"/utf8>>,
lattice_core@replica_id:to_json(Neg_id)},
{<<"counts"/utf8>>,
gleam@json:dict(
Neg_dict,
fun(K@1) ->
lattice_core@replica_id:to_string(
K@1
)
end,
fun gleam@json:int/1
)}]
)}]
)}]
).
-file("src/lattice_counters/pn_counter.gleam", 170).
?DOC(
" Decode a PN-Counter from a JSON string produced by `to_json`.\n"
"\n"
" Returns `Ok(PNCounter)` on success, or `Error(json.DecodeError)` if the\n"
" input is not a valid PN-Counter JSON envelope.\n"
).
-spec from_json(binary()) -> {ok, p_n_counter()} |
{error, gleam@json:decode_error()}.
from_json(Json_string) ->
G_counter_state_decoder = begin
gleam@dynamic@decode:field(
<<"self_id"/utf8>>,
lattice_core@replica_id:decoder(),
fun(Self_id) ->
Non_negative_int = begin
_pipe = {decoder, fun gleam@dynamic@decode:decode_int/1},
gleam@dynamic@decode:then(
_pipe,
fun(Val) -> case Val >= 0 of
true ->
gleam@dynamic@decode:success(Val);
false ->
gleam@dynamic@decode:failure(
Val,
<<"a non-negative integer"/utf8>>
)
end end
)
end,
gleam@dynamic@decode:field(
<<"counts"/utf8>>,
gleam@dynamic@decode:dict(
lattice_core@replica_id:decoder(),
Non_negative_int
),
fun(Counts) ->
gleam@dynamic@decode:success(
lattice_counters@g_counter:from_parts(
Counts,
Self_id
)
)
end
)
end
)
end,
State_decoder = begin
gleam@dynamic@decode:field(
<<"state"/utf8>>,
begin
gleam@dynamic@decode:field(
<<"positive"/utf8>>,
G_counter_state_decoder,
fun(Positive) ->
gleam@dynamic@decode:field(
<<"negative"/utf8>>,
G_counter_state_decoder,
fun(Negative) ->
gleam@dynamic@decode:success(
{p_n_counter, Positive, Negative}
)
end
)
end
)
end,
fun(State) -> gleam@dynamic@decode:success(State) end
)
end,
Envelope_decoder = begin
gleam@dynamic@decode:field(
<<"type"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Type_tag) ->
gleam@dynamic@decode:field(
<<"v"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_int/1},
fun(Version) ->
gleam@dynamic@decode:success({Type_tag, Version})
end
)
end
)
end,
case gleam@json:parse(Json_string, Envelope_decoder) of
{error, E} ->
{error, E};
{ok, {Type_tag@1, Version@1}} ->
case (Type_tag@1 =:= <<"pn_counter"/utf8>>) andalso (Version@1 =:= 1) of
true ->
gleam@json:parse(Json_string, State_decoder);
false ->
{error,
{unable_to_decode,
[{decode_error,
<<"type=pn_counter and v=1"/utf8>>,
<<<<Type_tag@1/binary, " v="/utf8>>/binary,
(erlang:integer_to_binary(Version@1))/binary>>,
[]}]}}
end
end.