Packages
Grow-only, two-phase, and observed-remove sets for Gleam — conflict-free replicated set types
Current section
Files
Jump to
Current section
Files
src/lattice_sets@or_set.erl
-module(lattice_sets@or_set).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/lattice_sets/or_set.gleam").
-export([new/1, add/2, remove/2, contains/2, value/1, merge/2, remove_with_bound/2, pruned_vv/1, prune/2, from_json/1, to_json/1]).
-export_type([tag/0, o_r_set/1]).
-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(
" An observed-remove set (OR-Set) CRDT.\n"
"\n"
" The most flexible set CRDT: supports add, remove, and re-add. Each add\n"
" creates a unique tag. Remove only deletes tags observed locally, so a\n"
" concurrent add on another replica survives (add-wins semantics). This makes\n"
" OR-Set suitable for collaborative data where elements may be toggled.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" import lattice_core/replica_id\n"
" import lattice_sets/or_set\n"
"\n"
" let a = or_set.new(replica_id.new(\"node-a\")) |> or_set.add(\"item\")\n"
" let b = or_set.new(replica_id.new(\"node-b\")) |> or_set.add(\"item\") |> or_set.remove(\"item\")\n"
" let merged = or_set.merge(a, b)\n"
" or_set.contains(merged, \"item\") // -> True (concurrent add wins)\n"
" ```\n"
).
-opaque tag() :: {tag, lattice_core@replica_id:replica_id(), integer()}.
-opaque o_r_set(PMP) :: {o_r_set,
lattice_core@replica_id:replica_id(),
integer(),
gleam@dict:dict(PMP, gleam@set:set(tag())),
gleam@set:set(tag()),
lattice_core@version_vector:version_vector()}.
-file("src/lattice_sets/or_set.gleam", 66).
?DOC(
" Create a new empty OR-Set for the given replica.\n"
"\n"
" Each replica should have a unique `replica_id` to ensure that tags\n"
" generated on different replicas never collide.\n"
).
-spec new(lattice_core@replica_id:replica_id()) -> o_r_set(any()).
new(Replica_id) ->
{o_r_set,
Replica_id,
0,
maps:new(),
gleam@set:new(),
lattice_core@version_vector:new()}.
-file("src/lattice_sets/or_set.gleam", 81).
?DOC(
" Add an element to the set.\n"
"\n"
" Creates a fresh unique tag for this add operation using the replica's\n"
" monotonically-increasing counter. The element may already be present;\n"
" in that case a new tag is added alongside existing ones.\n"
).
-spec add(o_r_set(PMS), PMS) -> o_r_set(PMS).
add(Orset, Element) ->
New_counter = erlang:element(3, Orset) + 1,
Tag = {tag, erlang:element(2, Orset), New_counter},
Existing_tags = gleam@result:unwrap(
gleam_stdlib:map_get(erlang:element(4, Orset), Element),
gleam@set:new()
),
New_tags = gleam@set:insert(Existing_tags, Tag),
{o_r_set,
erlang:element(2, Orset),
New_counter,
gleam@dict:insert(erlang:element(4, Orset), Element, New_tags),
erlang:element(5, Orset),
erlang:element(6, Orset)}.
-file("src/lattice_sets/or_set.gleam", 100).
?DOC(
" Remove an element from the set.\n"
"\n"
" Removes all currently observed tags for the element (observed-remove\n"
" semantics). Any concurrent add on another replica that created a new tag\n"
" not yet observed here will survive this remove after merging.\n"
).
-spec remove(o_r_set(PMV), PMV) -> o_r_set(PMV).
remove(Orset, Element) ->
Removed_tags = gleam@result:unwrap(
gleam_stdlib:map_get(erlang:element(4, Orset), Element),
gleam@set:new()
),
{o_r_set,
erlang:element(2, Orset),
erlang:element(3, Orset),
gleam@dict:delete(erlang:element(4, Orset), Element),
gleam@set:union(erlang:element(5, Orset), Removed_tags),
erlang:element(6, Orset)}.
-file("src/lattice_sets/or_set.gleam", 116).
?DOC(
" Check if the set contains the given element.\n"
"\n"
" Returns `True` if the element has at least one live tag (i.e., it has\n"
" been added and not yet removed on this replica, or a concurrent add\n"
" survived a remove after merging).\n"
).
-spec contains(o_r_set(PMY), PMY) -> boolean().
contains(Orset, Element) ->
case gleam_stdlib:map_get(erlang:element(4, Orset), Element) of
{error, _} ->
false;
{ok, Tags} ->
not gleam@set:is_empty(Tags)
end.
-file("src/lattice_sets/or_set.gleam", 126).
?DOC(
" Return the set of all elements currently in the OR-Set.\n"
"\n"
" An element is included only when its tag set is non-empty.\n"
).
-spec value(o_r_set(PNA)) -> gleam@set:set(PNA).
value(Orset) ->
_pipe = maps:keys(erlang:element(4, Orset)),
gleam@set:from_list(_pipe).
-file("src/lattice_sets/or_set.gleam", 182).
-spec not_dominated(tag(), lattice_core@version_vector:version_vector()) -> boolean().
not_dominated(Tag, Pruned) ->
{tag, Rid, C} = Tag,
lattice_core@version_vector:get(Pruned, Rid) < C.
-file("src/lattice_sets/or_set.gleam", 198).
-spec pruned_on_side_without_live_tag(
tag(),
gleam@set:set(tag()),
lattice_core@version_vector:version_vector()
) -> boolean().
pruned_on_side_without_live_tag(Tag, Live_tags, Pruned) ->
{tag, Rid, C} = Tag,
(lattice_core@version_vector:get(Pruned, Rid) >= C) andalso not gleam@set:contains(
Live_tags,
Tag
).
-file("src/lattice_sets/or_set.gleam", 187).
-spec is_pruned_zombie(
tag(),
gleam@set:set(tag()),
lattice_core@version_vector:version_vector(),
gleam@set:set(tag()),
lattice_core@version_vector:version_vector()
) -> boolean().
is_pruned_zombie(Tag, A_tags, A_pruned, B_tags, B_pruned) ->
pruned_on_side_without_live_tag(Tag, A_tags, A_pruned) orelse pruned_on_side_without_live_tag(
Tag,
B_tags,
B_pruned
).
-file("src/lattice_sets/or_set.gleam", 144).
?DOC(
" Merge two OR-Sets.\n"
"\n"
" For each element, the merged tag set is the union of both sides' tags,\n"
" minus merged tombstones, and minus any tags dominated by the merged\n"
" pruned vector that are not live on the side that pruned them (zombie\n"
" detection). An element is present if it has at least one surviving tag.\n"
"\n"
" The merged counter is the maximum of both sides, ensuring future adds on\n"
" either replica generate unique tags.\n"
"\n"
" Merge is commutative, associative, and idempotent (a valid CRDT join).\n"
).
-spec merge(o_r_set(PND), o_r_set(PND)) -> o_r_set(PND).
merge(A, B) ->
Merged_pruned = lattice_core@version_vector:merge(
erlang:element(6, A),
erlang:element(6, B)
),
Merged_tombstones = begin
_pipe = gleam@set:union(erlang:element(5, A), erlang:element(5, B)),
gleam@set:filter(
_pipe,
fun(Tag) -> not_dominated(Tag, Merged_pruned) end
)
end,
Merged_counter = gleam@int:max(erlang:element(3, A), erlang:element(3, B)),
A_keys = maps:keys(erlang:element(4, A)),
B_keys = maps:keys(erlang:element(4, B)),
All_keys = gleam@list:unique(lists:append(A_keys, B_keys)),
Merged_entries = gleam@list:fold(
All_keys,
maps:new(),
fun(Acc, Element) ->
A_tags = gleam@result:unwrap(
gleam_stdlib:map_get(erlang:element(4, A), Element),
gleam@set:new()
),
B_tags = gleam@result:unwrap(
gleam_stdlib:map_get(erlang:element(4, B), Element),
gleam@set:new()
),
Combined = begin
_pipe@1 = gleam@set:union(A_tags, B_tags),
gleam@set:filter(
_pipe@1,
fun(Tag@1) ->
not gleam@set:contains(Merged_tombstones, Tag@1) andalso not is_pruned_zombie(
Tag@1,
A_tags,
erlang:element(6, A),
B_tags,
erlang:element(6, B)
)
end
)
end,
case gleam@set:is_empty(Combined) of
true ->
Acc;
false ->
gleam@dict:insert(Acc, Element, Combined)
end
end
),
{o_r_set,
erlang:element(2, A),
Merged_counter,
Merged_entries,
Merged_tombstones,
Merged_pruned}.
-file("src/lattice_sets/or_set.gleam", 230).
-spec tags_to_bound(gleam@set:set(tag())) -> lattice_core@version_vector:version_vector().
tags_to_bound(Tags) ->
gleam@set:fold(
Tags,
lattice_core@version_vector:new(),
fun(Vv, Tag) ->
{tag, Rid, C} = Tag,
lattice_core@version_vector:set_max(Vv, Rid, C)
end
).
-file("src/lattice_sets/or_set.gleam", 215).
?DOC(
" Remove an element and return a causal bound for the removed tags.\n"
"\n"
" Behaves identically to `remove` but also returns a `VersionVector`\n"
" representing the maximum counter per replica across all tags that were\n"
" live for the element. This bound can be compared against a pruned vector\n"
" to determine when the removal is causally stable.\n"
"\n"
" Returns an empty `VersionVector` if the element had no live tags.\n"
).
-spec remove_with_bound(o_r_set(PNK), PNK) -> {o_r_set(PNK),
lattice_core@version_vector:version_vector()}.
remove_with_bound(Orset, Element) ->
Removed_tags = gleam@result:unwrap(
gleam_stdlib:map_get(erlang:element(4, Orset), Element),
gleam@set:new()
),
Bound = tags_to_bound(Removed_tags),
Updated = {o_r_set,
erlang:element(2, Orset),
erlang:element(3, Orset),
gleam@dict:delete(erlang:element(4, Orset), Element),
gleam@set:union(erlang:element(5, Orset), Removed_tags),
erlang:element(6, Orset)},
{Updated, Bound}.
-file("src/lattice_sets/or_set.gleam", 242).
?DOC(
" Return the pruned version vector.\n"
"\n"
" This is the causal horizon below which tombstones have been garbage\n"
" collected. Useful for determining whether a remove bound is fully\n"
" dominated (causally stable).\n"
).
-spec pruned_vv(o_r_set(any())) -> lattice_core@version_vector:version_vector().
pruned_vv(Orset) ->
erlang:element(6, Orset).
-file("src/lattice_sets/or_set.gleam", 253).
?DOC(
" Prune tombstones based on a stable version vector.\n"
"\n"
" Updates the `pruned` vector by merging it with `stable_vv`. Any tombstones\n"
" dominated by the new `pruned` vector are removed. This function should only\n"
" be called with a version vector representing events that have been seen by\n"
" all replicas (causally stable), otherwise \"zombie\" updates might be\n"
" incorrectly ignored.\n"
).
-spec prune(o_r_set(PNQ), lattice_core@version_vector:version_vector()) -> o_r_set(PNQ).
prune(Orset, Stable_vv) ->
New_pruned = lattice_core@version_vector:merge(
erlang:element(6, Orset),
Stable_vv
),
Pruned_tombstones = gleam@set:filter(
erlang:element(5, Orset),
fun(Tag) -> not_dominated(Tag, New_pruned) end
),
{o_r_set,
erlang:element(2, Orset),
erlang:element(3, Orset),
erlang:element(4, Orset),
Pruned_tombstones,
New_pruned}.
-file("src/lattice_sets/or_set.gleam", 297).
?DOC(
" Decode an `ORSet(String)` from a JSON string produced by `to_json`.\n"
"\n"
" Supports both v1 (no pruned field) and v2 formats. Returns `Error` if the\n"
" string is not valid JSON or does not match the expected format.\n"
).
-spec from_json(binary()) -> {ok, o_r_set(binary())} |
{error, gleam@json:decode_error()}.
from_json(Json_string) ->
Tag_decoder = begin
gleam@dynamic@decode:field(
<<"r"/utf8>>,
lattice_core@replica_id:decoder(),
fun(R) ->
gleam@dynamic@decode:field(
<<"c"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_int/1},
fun(C) -> gleam@dynamic@decode:success({tag, R, C}) end
)
end
)
end,
Tag_set_decoder = gleam@dynamic@decode:map(
gleam@dynamic@decode:list(Tag_decoder),
fun gleam@set:from_list/1
),
V1_state_decoder = begin
gleam@dynamic@decode:field(
<<"state"/utf8>>,
begin
gleam@dynamic@decode:field(
<<"replica_id"/utf8>>,
lattice_core@replica_id:decoder(),
fun(Replica_id) ->
gleam@dynamic@decode:field(
<<"counter"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_int/1},
fun(Counter) ->
gleam@dynamic@decode:field(
<<"entries"/utf8>>,
gleam@dynamic@decode:dict(
{decoder,
fun gleam@dynamic@decode:decode_string/1},
Tag_set_decoder
),
fun(Entries) ->
gleam@dynamic@decode:optional_field(
<<"tombstones"/utf8>>,
[],
gleam@dynamic@decode:list(
Tag_decoder
),
fun(Tombstones) ->
gleam@dynamic@decode:success(
{o_r_set,
Replica_id,
Counter,
Entries,
gleam@set:from_list(
Tombstones
),
lattice_core@version_vector:new(
)}
)
end
)
end
)
end
)
end
)
end,
fun(State) -> gleam@dynamic@decode:success(State) end
)
end,
V2_state_decoder = begin
gleam@dynamic@decode:field(
<<"state"/utf8>>,
begin
gleam@dynamic@decode:field(
<<"replica_id"/utf8>>,
lattice_core@replica_id:decoder(),
fun(Replica_id@1) ->
gleam@dynamic@decode:field(
<<"counter"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_int/1},
fun(Counter@1) ->
gleam@dynamic@decode:field(
<<"entries"/utf8>>,
gleam@dynamic@decode:dict(
{decoder,
fun gleam@dynamic@decode:decode_string/1},
Tag_set_decoder
),
fun(Entries@1) ->
gleam@dynamic@decode:field(
<<"tombstones"/utf8>>,
Tag_set_decoder,
fun(Tombstones@1) ->
gleam@dynamic@decode:field(
<<"pruned"/utf8>>,
lattice_core@version_vector:decoder(
),
fun(Pruned) ->
gleam@dynamic@decode:success(
{o_r_set,
Replica_id@1,
Counter@1,
Entries@1,
Tombstones@1,
Pruned}
)
end
)
end
)
end
)
end
)
end
)
end,
fun(State@1) -> gleam@dynamic@decode:success(State@1) 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 =:= <<"or_set"/utf8>> of
false ->
{error,
{unable_to_decode,
[{decode_error,
<<"type=or_set"/utf8>>,
Type_tag@1,
[]}]}};
true ->
case Version@1 of
1 ->
gleam@json:parse(Json_string, V1_state_decoder);
2 ->
gleam@json:parse(Json_string, V2_state_decoder);
_ ->
{error,
{unable_to_decode,
[{decode_error,
<<"v=1 or v=2"/utf8>>,
erlang:integer_to_binary(Version@1),
[<<"v"/utf8>>]}]}}
end
end
end.
-file("src/lattice_sets/or_set.gleam", 388).
-spec encode_tag(tag()) -> gleam@json:json().
encode_tag(Tag) ->
{tag, Rid, C} = Tag,
gleam@json:object(
[{<<"r"/utf8>>,
gleam@json:string(lattice_core@replica_id:to_string(Rid))},
{<<"c"/utf8>>, gleam@json:int(C)}]
).
-file("src/lattice_sets/or_set.gleam", 271).
?DOC(
" Encode an `ORSet(String)` as a self-describing JSON value.\n"
"\n"
" Entries are encoded as a JSON dict where values are arrays of tag objects\n"
" `{\"r\": replica_id, \"c\": counter}`. Removed tags are encoded separately in\n"
" `tombstones`. The `pruned` version vector tracks garbage-collected causal\n"
" history.\n"
"\n"
" Format: `{\"type\": \"or_set\", \"v\": 2, \"state\": {\"replica_id\": \"...\", \"counter\": N, \"entries\": {...}, \"tombstones\": [...], \"pruned\": {...}}}`\n"
"\n"
" The encoded value can be restored with `from_json`.\n"
).
-spec to_json(o_r_set(binary())) -> gleam@json:json().
to_json(Orset) ->
gleam@json:object(
[{<<"type"/utf8>>, gleam@json:string(<<"or_set"/utf8>>)},
{<<"v"/utf8>>, gleam@json:int(2)},
{<<"state"/utf8>>,
gleam@json:object(
[{<<"replica_id"/utf8>>,
lattice_core@replica_id:to_json(
erlang:element(2, Orset)
)},
{<<"counter"/utf8>>,
gleam@json:int(erlang:element(3, Orset))},
{<<"entries"/utf8>>,
gleam@json:dict(
erlang:element(4, Orset),
fun(K) -> K end,
fun(Tag_set) ->
gleam@json:array(
gleam@set:to_list(Tag_set),
fun encode_tag/1
)
end
)},
{<<"tombstones"/utf8>>,
gleam@json:array(
gleam@set:to_list(erlang:element(5, Orset)),
fun encode_tag/1
)},
{<<"pruned"/utf8>>,
lattice_core@version_vector:to_json(
erlang:element(6, Orset)
)}]
)}]
).