Packages
macula
0.21.2
7.0.0
6.0.0
5.2.2
5.2.1
5.2.0
5.1.0
5.0.0
4.8.0
4.7.1
4.7.0
4.6.0
4.5.0
4.4.10
4.4.9
4.4.8
4.4.7
4.4.6
4.4.5
4.4.4
4.4.3
4.4.2
4.4.1
4.4.0
4.3.1
4.3.0
4.2.9
4.2.8
4.2.7
4.2.6
4.2.5
4.2.4
4.2.3
4.2.2
4.2.1
4.2.0
4.1.1
4.1.0
4.0.0
3.16.0
3.15.3
3.15.2
3.15.1
3.14.0
3.13.0
3.12.1
3.12.0
3.11.1
3.11.0
3.10.3
3.10.2
3.10.1
3.9.0
3.8.0
3.7.0
3.5.0
3.4.0
3.3.0
3.2.0
3.1.0
3.0.0
2.1.1
2.1.0
2.0.0
1.5.2
1.5.1
1.4.30
1.4.29
1.4.28
1.4.27
1.4.26
1.4.25
1.4.24
1.4.23
1.4.22
1.4.21
1.4.20
1.4.19
1.4.18
1.4.17
1.4.16
1.4.15
1.4.14
1.4.13
1.4.11
1.4.10
1.4.9
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.1
1.3.0
1.2.0
1.1.0
1.0.10
1.0.9
1.0.8
1.0.7
1.0.6
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
0.48.6
0.48.5
0.48.4
0.48.3
0.48.2
0.48.1
0.48.0
0.47.1
0.47.0
0.46.3
0.46.1
0.46.0
0.45.3
0.45.2
0.45.1
0.45.0
0.44.2
0.44.1
0.44.0
0.43.3
0.43.2
0.43.1
0.43.0
0.42.9
0.42.8
0.42.7
0.42.6
0.42.5
0.42.4
0.42.3
0.42.2
0.42.1
0.42.0
0.41.1
0.41.0
0.40.1
0.40.0
0.39.9
0.39.8
0.39.7
0.39.6
0.39.5
0.39.4
0.39.3
0.39.2
0.39.1
0.39.0
0.38.8
0.38.7
0.38.6
0.38.5
0.38.4
0.38.3
0.38.2
0.38.1
0.38.0
0.37.7
0.37.6
0.37.5
0.37.4
0.37.3
0.37.2
0.37.1
0.37.0
0.36.6
0.36.5
0.36.4
0.36.3
0.36.2
0.36.1
0.36.0
0.35.4
0.35.3
0.35.2
0.35.1
0.35.0
0.34.1
0.34.0
0.33.1
0.33.0
0.32.5
0.32.4
0.32.3
0.32.2
0.32.1
0.32.0
0.31.9
0.31.8
0.31.7
0.31.6
0.31.5
0.31.4
0.31.3
0.31.2
0.31.1
0.31.0
0.30.10
0.30.9
0.30.8
0.30.7
0.30.6
0.30.5
0.30.4
0.30.3
0.30.2
0.30.1
0.30.0
0.29.0
0.28.3
0.28.2
0.28.1
0.28.0
0.27.1
0.27.0
0.26.1
0.26.0
0.25.6
0.25.5
0.25.4
0.25.3
0.25.2
0.25.1
0.25.0
0.24.6
0.24.5
0.24.4
0.24.3
0.24.2
0.24.1
0.24.0
0.23.3
0.23.2
0.23.1
0.23.0
0.22.12
0.22.11
0.22.10
0.22.9
0.22.8
0.22.7
0.22.6
0.22.5
0.22.4
0.22.3
0.22.2
0.22.1
0.22.0
0.21.7
0.21.6
0.21.5
0.21.4
0.21.2
0.21.1
0.21.0
0.20.25
0.20.24
0.20.23
0.20.22
0.20.21
0.20.20
0.20.19
0.20.18
0.20.17
0.20.16
0.20.15
0.20.14
0.20.13
0.20.12
0.20.11
0.20.10
0.20.9
0.20.8
0.20.7
0.20.6
0.20.5
0.20.3
0.20.2
0.20.1
0.20.0
0.19.2
0.19.1
0.19.0
0.18.1
0.18.0
0.17.4
0.17.3
0.17.2
0.17.1
0.17.0
0.16.6
0.16.5
0.16.4
0.16.3
0.16.2
0.16.1
0.16.0
0.15.1
0.15.0
0.14.3
0.14.2
0.14.1
0.14.0
0.12.6
0.12.5
0.12.3
0.11.3
0.10.2
0.10.1
0.10.0
0.9.2
0.9.1
0.9.0
0.8.25
0.8.24
0.8.23
0.8.22
0.8.21
0.8.20
0.8.19
0.8.18
0.8.17
0.8.16
0.8.15
0.8.14
0.8.13
0.8.12
0.8.11
0.8.10
0.8.9
0.8.8
0.8.7
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
0.7.30
0.7.29
0.7.28
0.7.27
0.7.26
0.7.25
0.7.24
0.7.23
0.7.22
0.7.21
0.7.20
0.7.19
0.7.18
0.7.17
0.7.16
0.7.15
0.7.14
0.7.13
0.7.12
0.7.11
0.7.10
0.7.9
0.7.8
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.7
0.6.6
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.0
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.4
0.3.3
0.3.2
0.3.1
Macula HTTP/3 Mesh SDK — connect, subscribe, publish, call, advertise
Current section
Files
Jump to
Current section
Files
src/macula_platform_system/macula_crdt.erl
%%%-------------------------------------------------------------------
%%% @doc Conflict-Free Replicated Data Types (CRDTs) for shared state.
%%%
%%% Implements CRDTs for eventually-consistent distributed state management:
%%% - LWW-Register (Last-Write-Wins Register) - single value with timestamp
%%% - OR-Set (Observed-Remove Set) - set with add/remove semantics
%%% - G-Counter (Grow-only Counter) - monotonically increasing counter
%%% - PN-Counter (Positive-Negative Counter) - increment/decrement counter
%%%
%%% These CRDTs replace Ra/Raft consensus for Macula's masterless architecture.
%%% State is synchronized via gossip protocol between nodes.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_crdt).
%% LWW-Register API
-export([
new_lww_register/0,
new_lww_register/1,
lww_set/3,
lww_get/1,
lww_merge/2,
lww_timestamp/1,
lww_value/1
]).
%% OR-Set API
-export([
new_or_set/0,
or_add/2,
or_add/3,
or_remove/2,
or_contains/2,
or_elements/1,
or_size/1,
or_merge/2
]).
%% G-Counter API
-export([
new_gcounter/0,
gcounter_increment/1,
gcounter_increment/2,
gcounter_value/1,
gcounter_merge/2
]).
%% PN-Counter API
-export([
new_pncounter/0,
pncounter_increment/1,
pncounter_increment/2,
pncounter_decrement/1,
pncounter_decrement/2,
pncounter_value/1,
pncounter_merge/2
]).
%% Type definitions
-type timestamp() :: non_neg_integer().
-type unique_tag() :: binary().
-type lww_register() :: #{
value => term(),
timestamp => timestamp(),
node => node()
}.
%% OR-Set: each element has unique tags for add operations
%% An element is in the set if it has at least one tag not in tombstones
-type or_set() :: #{
elements => #{term() => sets:set(unique_tag())}, % element -> set of add tags
tombstones => sets:set(unique_tag()) % removed tags
}.
%% G-Counter: per-node counts that only grow
-type gcounter() :: #{node() => non_neg_integer()}.
%% PN-Counter: pair of G-Counters for increments and decrements
-type pncounter() :: #{
p => gcounter(), % positive (increments)
n => gcounter() % negative (decrements)
}.
-export_type([
lww_register/0,
timestamp/0,
or_set/0,
gcounter/0,
pncounter/0
]).
%%==============================================================================
%% LWW-Register (Last-Write-Wins Register)
%%==============================================================================
%% A simple CRDT that resolves conflicts by keeping the value with
%% the highest timestamp. Ties are broken by node name (lexicographic order).
%% @doc Create a new empty LWW-Register
-spec new_lww_register() -> lww_register().
new_lww_register() ->
#{
value => undefined,
timestamp => 0,
node => node()
}.
%% @doc Create a new LWW-Register with initial value
-spec new_lww_register(term()) -> lww_register().
new_lww_register(Value) ->
#{
value => Value,
timestamp => erlang:system_time(microsecond),
node => node()
}.
%% @doc Set a value in the LWW-Register with timestamp
-spec lww_set(lww_register(), term(), timestamp()) -> lww_register().
lww_set(_Register, Value, Timestamp) when is_integer(Timestamp), Timestamp >= 0 ->
#{
value => Value,
timestamp => Timestamp,
node => node()
}.
%% @doc Get the current value of the LWW-Register
-spec lww_get(lww_register()) -> term().
lww_get(#{value := Value}) ->
Value.
%% @doc Get the timestamp of the LWW-Register
-spec lww_timestamp(lww_register()) -> timestamp().
lww_timestamp(#{timestamp := Timestamp}) ->
Timestamp.
%% @doc Get the value of the LWW-Register (alias for lww_get)
-spec lww_value(lww_register()) -> term().
lww_value(Register) ->
lww_get(Register).
%% @doc Merge two LWW-Registers
%% Keeps the value with the highest timestamp
%% Ties broken by node name (lexicographic order)
-spec lww_merge(lww_register(), lww_register()) -> lww_register().
lww_merge(
#{timestamp := T1} = R1,
#{timestamp := T2}
) when T1 > T2 ->
R1;
lww_merge(
#{timestamp := T1},
#{timestamp := T2} = R2
) when T1 < T2 ->
R2;
lww_merge(
#{timestamp := T, node := N1} = R1,
#{timestamp := T, node := N2} = R2
) ->
%% Same timestamp - use lexicographic order of node names
case N1 < N2 of
true -> R1;
false -> R2
end.
%%==============================================================================
%% OR-Set (Observed-Remove Set)
%%==============================================================================
%% A set CRDT that supports both add and remove operations.
%% Each add creates a unique tag. Remove adds tags to tombstone set.
%% An element is present if it has at least one tag not tombstoned.
%% @doc Create a new empty OR-Set
-spec new_or_set() -> or_set().
new_or_set() ->
#{
elements => #{},
tombstones => sets:new()
}.
%% @doc Add an element to the OR-Set (auto-generates unique tag)
-spec or_add(or_set(), term()) -> or_set().
or_add(Set, Element) ->
Tag = generate_unique_tag(),
or_add(Set, Element, Tag).
%% @doc Add an element to the OR-Set with a specific tag
-spec or_add(or_set(), term(), unique_tag()) -> or_set().
or_add(#{elements := Elements, tombstones := Tombstones}, Element, Tag) ->
CurrentTags = maps:get(Element, Elements, sets:new()),
NewTags = sets:add_element(Tag, CurrentTags),
#{
elements => maps:put(Element, NewTags, Elements),
tombstones => Tombstones
}.
%% @doc Remove an element from the OR-Set
%% Removes all current tags for that element (add-wins semantics)
-spec or_remove(or_set(), term()) -> or_set().
or_remove(#{elements := Elements, tombstones := Tombstones} = Set, Element) ->
case maps:find(Element, Elements) of
{ok, Tags} ->
%% Move all tags to tombstones
NewTombstones = sets:union(Tombstones, Tags),
#{
elements => maps:remove(Element, Elements),
tombstones => NewTombstones
};
error ->
%% Element not in set - no-op
Set
end.
%% @doc Check if element is in the OR-Set
-spec or_contains(or_set(), term()) -> boolean().
or_contains(#{elements := Elements, tombstones := Tombstones}, Element) ->
case maps:find(Element, Elements) of
{ok, Tags} ->
%% Element is present if it has any non-tombstoned tags
ActiveTags = sets:subtract(Tags, Tombstones),
sets:size(ActiveTags) > 0;
error ->
false
end.
%% @doc Get all elements in the OR-Set
-spec or_elements(or_set()) -> [term()].
or_elements(#{elements := Elements, tombstones := Tombstones}) ->
maps:fold(
fun(Element, Tags, Acc) ->
ActiveTags = sets:subtract(Tags, Tombstones),
case sets:size(ActiveTags) > 0 of
true -> [Element | Acc];
false -> Acc
end
end,
[],
Elements
).
%% @doc Get the number of elements in the OR-Set
-spec or_size(or_set()) -> non_neg_integer().
or_size(Set) ->
length(or_elements(Set)).
%% @doc Merge two OR-Sets
%% Union of elements, union of tombstones
-spec or_merge(or_set(), or_set()) -> or_set().
or_merge(
#{elements := E1, tombstones := T1},
#{elements := E2, tombstones := T2}
) ->
%% Merge elements: for each key, union the tag sets
MergedElements = maps:fold(
fun(Element, Tags2, Acc) ->
Tags1 = maps:get(Element, Acc, sets:new()),
maps:put(Element, sets:union(Tags1, Tags2), Acc)
end,
E1,
E2
),
%% Merge tombstones: union
MergedTombstones = sets:union(T1, T2),
#{
elements => MergedElements,
tombstones => MergedTombstones
}.
%%==============================================================================
%% G-Counter (Grow-only Counter)
%%==============================================================================
%% A counter that can only be incremented.
%% Each node maintains its own count; total is sum of all nodes.
%% @doc Create a new G-Counter
-spec new_gcounter() -> gcounter().
new_gcounter() ->
#{}.
%% @doc Increment the G-Counter by 1 for current node
-spec gcounter_increment(gcounter()) -> gcounter().
gcounter_increment(Counter) ->
gcounter_increment(Counter, 1).
%% @doc Increment the G-Counter by N for current node
-spec gcounter_increment(gcounter(), pos_integer()) -> gcounter().
gcounter_increment(Counter, N) when is_integer(N), N > 0 ->
Node = node(),
Current = maps:get(Node, Counter, 0),
maps:put(Node, Current + N, Counter).
%% @doc Get the total value of the G-Counter
-spec gcounter_value(gcounter()) -> non_neg_integer().
gcounter_value(Counter) ->
maps:fold(fun(_Node, Count, Acc) -> Acc + Count end, 0, Counter).
%% @doc Merge two G-Counters (take max of each node's count)
-spec gcounter_merge(gcounter(), gcounter()) -> gcounter().
gcounter_merge(C1, C2) ->
maps:fold(
fun(Node, Count2, Acc) ->
Count1 = maps:get(Node, Acc, 0),
maps:put(Node, max(Count1, Count2), Acc)
end,
C1,
C2
).
%%==============================================================================
%% PN-Counter (Positive-Negative Counter)
%%==============================================================================
%% A counter that supports both increment and decrement.
%% Implemented as pair of G-Counters (positive - negative).
%% @doc Create a new PN-Counter
-spec new_pncounter() -> pncounter().
new_pncounter() ->
#{
p => new_gcounter(),
n => new_gcounter()
}.
%% @doc Increment the PN-Counter by 1
-spec pncounter_increment(pncounter()) -> pncounter().
pncounter_increment(Counter) ->
pncounter_increment(Counter, 1).
%% @doc Increment the PN-Counter by N
-spec pncounter_increment(pncounter(), pos_integer()) -> pncounter().
pncounter_increment(#{p := P, n := N}, Amount) when is_integer(Amount), Amount > 0 ->
#{
p => gcounter_increment(P, Amount),
n => N
}.
%% @doc Decrement the PN-Counter by 1
-spec pncounter_decrement(pncounter()) -> pncounter().
pncounter_decrement(Counter) ->
pncounter_decrement(Counter, 1).
%% @doc Decrement the PN-Counter by N
-spec pncounter_decrement(pncounter(), pos_integer()) -> pncounter().
pncounter_decrement(#{p := P, n := N}, Amount) when is_integer(Amount), Amount > 0 ->
#{
p => P,
n => gcounter_increment(N, Amount)
}.
%% @doc Get the value of the PN-Counter (positive - negative)
-spec pncounter_value(pncounter()) -> integer().
pncounter_value(#{p := P, n := N}) ->
gcounter_value(P) - gcounter_value(N).
%% @doc Merge two PN-Counters
-spec pncounter_merge(pncounter(), pncounter()) -> pncounter().
pncounter_merge(#{p := P1, n := N1}, #{p := P2, n := N2}) ->
#{
p => gcounter_merge(P1, P2),
n => gcounter_merge(N1, N2)
}.
%%==============================================================================
%% Internal Functions
%%==============================================================================
%% @private Generate a unique tag for OR-Set operations
-spec generate_unique_tag() -> unique_tag().
generate_unique_tag() ->
%% Combine timestamp, node, and random bytes for uniqueness
Timestamp = erlang:system_time(nanosecond),
NodeBin = atom_to_binary(node(), utf8),
Random = crypto:strong_rand_bytes(8),
TimestampBin = integer_to_binary(Timestamp),
<<TimestampBin/binary, "-", NodeBin/binary, "-", Random/binary>>.