Packages
Last-writer-wins and observed-remove maps for Gleam — conflict-free replicated map types with generic value support
Current section
Files
Jump to
Current section
Files
src/lattice_maps@or_map.erl
-module(lattice_maps@or_map).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/lattice_maps/or_map.gleam").
-export([new/2, get/2, remove/2, keys/1, values/1, prune/2, internal_value_count/1, to_json/1, update/3, merge/2, from_json/1]).
-export_type([o_r_map/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(
" An observed-remove map (OR-Map) CRDT.\n"
"\n"
" Keys are tracked using an OR-Set with add-wins semantics: concurrent update\n"
" and remove of the same key resolves in favor of the update. Each value is\n"
" itself a CRDT (specified by `CrdtSpec` at construction), enabling nested\n"
" convergent data structures.\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" import lattice_maps/crdt\n"
" import lattice_core/replica_id\n"
" import lattice_counters/g_counter\n"
" import lattice_maps/or_map\n"
"\n"
" let map = or_map.new(replica_id.new(\"node-a\"), crdt.GCounterSpec)\n"
" |> or_map.update(\"score\", fn(c) {\n"
" let assert crdt.CrdtGCounter(gc) = c\n"
" crdt.CrdtGCounter(g_counter.increment(gc, 10))\n"
" })\n"
" ```\n"
).
-opaque o_r_map() :: {o_r_map,
lattice_core@replica_id:replica_id(),
lattice_maps@crdt:crdt_spec(),
lattice_sets@or_set:o_r_set(binary()),
gleam@dict:dict(binary(), lattice_maps@crdt:crdt()),
gleam@dict:dict(binary(), lattice_core@version_vector:version_vector())}.
-file("src/lattice_maps/or_map.gleam", 54).
-spec spec_to_string(lattice_maps@crdt:crdt_spec()) -> binary().
spec_to_string(Spec) ->
case Spec of
g_counter_spec ->
<<"g_counter"/utf8>>;
pn_counter_spec ->
<<"pn_counter"/utf8>>;
lww_register_spec ->
<<"lww_register"/utf8>>;
mv_register_spec ->
<<"mv_register"/utf8>>;
g_set_spec ->
<<"g_set"/utf8>>;
two_p_set_spec ->
<<"two_p_set"/utf8>>;
or_set_spec ->
<<"or_set"/utf8>>
end.
-file("src/lattice_maps/or_map.gleam", 66).
-spec string_to_spec(binary()) -> {ok, lattice_maps@crdt:crdt_spec()} |
{error, nil}.
string_to_spec(S) ->
case S of
<<"g_counter"/utf8>> ->
{ok, g_counter_spec};
<<"pn_counter"/utf8>> ->
{ok, pn_counter_spec};
<<"lww_register"/utf8>> ->
{ok, lww_register_spec};
<<"mv_register"/utf8>> ->
{ok, mv_register_spec};
<<"g_set"/utf8>> ->
{ok, g_set_spec};
<<"two_p_set"/utf8>> ->
{ok, two_p_set_spec};
<<"or_set"/utf8>> ->
{ok, or_set_spec};
_ ->
{error, nil}
end.
-file("src/lattice_maps/or_map.gleam", 83).
?DOC(
" Create a new empty OR-Map for the given replica with the specified CRDT type.\n"
"\n"
" The `crdt_spec` determines what type of CRDT is auto-created when `update`\n"
" is called on a key that does not yet exist in the map.\n"
).
-spec new(lattice_core@replica_id:replica_id(), lattice_maps@crdt:crdt_spec()) -> o_r_map().
new(Replica_id, Crdt_spec) ->
{o_r_map,
Replica_id,
Crdt_spec,
lattice_sets@or_set:new(Replica_id),
maps:new(),
maps:new()}.
-file("src/lattice_maps/or_map.gleam", 124).
?DOC(
" Get the CRDT value at `key`.\n"
"\n"
" Returns `Ok(crdt)` if the key is active in the OR-Set.\n"
" Returns `Error(Nil)` if the key has never been added, or has been removed\n"
" and not re-added.\n"
).
-spec get(o_r_map(), binary()) -> {ok, lattice_maps@crdt:crdt()} | {error, nil}.
get(Map, Key) ->
case lattice_sets@or_set:contains(erlang:element(4, Map), Key) of
true ->
case gleam_stdlib:map_get(erlang:element(5, Map), Key) of
{ok, Val} ->
{ok, Val};
{error, _} ->
{error, nil}
end;
false ->
{error, nil}
end.
-file("src/lattice_maps/or_map.gleam", 141).
?DOC(
" Remove a key from the OR-Map.\n"
"\n"
" Removes the key from the OR-Set (marking it inactive). The underlying\n"
" CRDT value is retained until prune determines the removal is causally\n"
" stable. A causal bound is recorded so prune can later decide when it is\n"
" safe to discard the value.\n"
).
-spec remove(o_r_map(), binary()) -> o_r_map().
remove(Map, Key) ->
{Updated_key_set, Bound} = lattice_sets@or_set:remove_with_bound(
erlang:element(4, Map),
Key
),
Updated_bounds = case lattice_core@version_vector:is_empty(Bound) of
true ->
erlang:element(6, Map);
false ->
gleam@dict:insert(erlang:element(6, Map), Key, Bound)
end,
{o_r_map,
erlang:element(2, Map),
erlang:element(3, Map),
Updated_key_set,
erlang:element(5, Map),
Updated_bounds}.
-file("src/lattice_maps/or_map.gleam", 153).
?DOC(
" Return the list of all active keys (those present in the OR-Set).\n"
"\n"
" Order is not guaranteed.\n"
).
-spec keys(o_r_map()) -> list(binary()).
keys(Map) ->
gleam@set:to_list(lattice_sets@or_set:value(erlang:element(4, Map))).
-file("src/lattice_maps/or_map.gleam", 160).
?DOC(
" Return the CRDT values for all active keys.\n"
"\n"
" Order is not guaranteed and does not correspond to the order of `keys`.\n"
).
-spec values(o_r_map()) -> list(lattice_maps@crdt:crdt()).
values(Map) ->
Active_keys = lattice_sets@or_set:value(erlang:element(4, Map)),
gleam@dict:fold(
erlang:element(5, Map),
[],
fun(Acc, Key, Val) -> case gleam@set:contains(Active_keys, Key) of
true ->
[Val | Acc];
false ->
Acc
end end
).
-file("src/lattice_maps/or_map.gleam", 257).
?DOC(
" Prune tombstones for keys and compact removed values whose removal is\n"
" causally stable.\n"
"\n"
" Delegates to `or_set.prune` to remove tombstones from the internal key\n"
" tracker. Then, for each removed key that has a recorded causal bound, if\n"
" the pruned version vector dominates that bound, the key's CRDT value is\n"
" discarded (the removal is stable and no concurrent re-add can reference\n"
" the old value).\n"
"\n"
" Only call this with a version vector representing events that have been\n"
" seen by all replicas (causally stable), otherwise zombie updates might be\n"
" incorrectly ignored.\n"
).
-spec prune(o_r_map(), lattice_core@version_vector:version_vector()) -> o_r_map().
prune(Map, Stable_vv) ->
Pruned_key_set = lattice_sets@or_set:prune(
erlang:element(4, Map),
Stable_vv
),
Pruned_vv = lattice_sets@or_set:pruned_vv(Pruned_key_set),
Active_keys = lattice_sets@or_set:value(Pruned_key_set),
{Compacted_values, Compacted_bounds} = gleam@dict:fold(
erlang:element(5, Map),
{maps:new(), erlang:element(6, Map)},
fun(Acc, Key, Val) ->
{Vals, Bounds} = Acc,
case gleam@set:contains(Active_keys, Key) of
true ->
{gleam@dict:insert(Vals, Key, Val), Bounds};
false ->
case gleam_stdlib:map_get(erlang:element(6, Map), Key) of
{ok, Bound} ->
case lattice_core@version_vector:dominates(
Pruned_vv,
Bound
) of
true ->
{Vals, gleam@dict:delete(Bounds, Key)};
false ->
{gleam@dict:insert(Vals, Key, Val), Bounds}
end;
{error, _} ->
{gleam@dict:insert(Vals, Key, Val), Bounds}
end
end
end
),
{o_r_map,
erlang:element(2, Map),
erlang:element(3, Map),
Pruned_key_set,
Compacted_values,
Compacted_bounds}.
-file("src/lattice_maps/or_map.gleam", 292).
?DOC(false).
-spec internal_value_count(o_r_map()) -> integer().
internal_value_count(Map) ->
maps:size(erlang:element(5, Map)).
-file("src/lattice_maps/or_map.gleam", 304).
?DOC(
" Encode an `ORMap` as a self-describing JSON value.\n"
"\n"
" The nested OR-Set (`key_set`) and CRDT values are double-encoded as JSON\n"
" strings so they can be decoded using the existing `from_json` APIs.\n"
"\n"
" Format: `{\"type\": \"or_map\", \"v\": 2, \"state\": {\"replica_id\": \"...\", \"crdt_spec\": \"...\", \"key_set\": \"...\", \"values\": [...], \"remove_bounds\": {...}}}`\n"
"\n"
" The encoded value can be restored with `from_json`.\n"
).
-spec to_json(o_r_map()) -> gleam@json:json().
to_json(Map) ->
{o_r_map, Rid, Crdt_spec, Key_set, Values, Remove_bounds} = Map,
Values_json = gleam@json:array(
maps:to_list(Values),
fun(Pair) ->
{Key, Crdt_val} = Pair,
gleam@json:object(
[{<<"key"/utf8>>, gleam@json:string(Key)},
{<<"crdt"/utf8>>,
gleam@json:string(
gleam@json:to_string(
lattice_maps@crdt:to_json(Crdt_val)
)
)}]
)
end
),
Bounds_json = gleam@json:dict(
Remove_bounds,
fun(K) -> K end,
fun(Vv) -> lattice_core@version_vector:to_json(Vv) end
),
gleam@json:object(
[{<<"type"/utf8>>, gleam@json:string(<<"or_map"/utf8>>)},
{<<"v"/utf8>>, gleam@json:int(2)},
{<<"state"/utf8>>,
gleam@json:object(
[{<<"replica_id"/utf8>>,
gleam@json:string(
lattice_core@replica_id:to_string(Rid)
)},
{<<"crdt_spec"/utf8>>,
gleam@json:string(spec_to_string(Crdt_spec))},
{<<"key_set"/utf8>>,
gleam@json:string(
gleam@json:to_string(
lattice_sets@or_set:to_json(Key_set)
)
)},
{<<"values"/utf8>>, Values_json},
{<<"remove_bounds"/utf8>>, Bounds_json}]
)}]
).
-file("src/lattice_maps/or_map.gleam", 481).
-spec matches_spec(lattice_maps@crdt:crdt(), lattice_maps@crdt:crdt_spec()) -> boolean().
matches_spec(Value, Spec) ->
case {Value, Spec} of
{{crdt_g_counter, _}, g_counter_spec} ->
true;
{{crdt_pn_counter, _}, pn_counter_spec} ->
true;
{{crdt_lww_register, _}, lww_register_spec} ->
true;
{{crdt_mv_register, _}, mv_register_spec} ->
true;
{{crdt_g_set, _}, g_set_spec} ->
true;
{{crdt_two_p_set, _}, two_p_set_spec} ->
true;
{{crdt_or_set, _}, or_set_spec} ->
true;
{_, _} ->
false
end.
-file("src/lattice_maps/or_map.gleam", 98).
?DOC(
" Apply a function to the CRDT value at `key`, auto-creating it if absent.\n"
"\n"
" If the key does not exist, a default value is created from `crdt_spec`\n"
" and passed to `f`. The key is added to the OR-Set, marking it active.\n"
" The return value of `f` replaces (or sets) the value for that key.\n"
).
-spec update(
o_r_map(),
binary(),
fun((lattice_maps@crdt:crdt()) -> lattice_maps@crdt:crdt())
) -> o_r_map().
update(Map, Key, F) ->
Current = case {lattice_sets@or_set:contains(erlang:element(4, Map), Key),
gleam_stdlib:map_get(erlang:element(5, Map), Key)} of
{true, {ok, Crdt_val}} ->
Crdt_val;
{_, _} ->
lattice_maps@crdt:default_crdt(
erlang:element(3, Map),
erlang:element(2, Map)
)
end,
Updated = case matches_spec(F(Current), erlang:element(3, Map)) of
true ->
F(Current);
false ->
Current
end,
{o_r_map,
erlang:element(2, Map),
erlang:element(3, Map),
lattice_sets@or_set:add(erlang:element(4, Map), Key),
gleam@dict:insert(erlang:element(5, Map), Key, Updated),
gleam@dict:delete(erlang:element(6, Map), Key)}.
-file("src/lattice_maps/or_map.gleam", 494).
-spec valid_value(o_r_map(), binary()) -> {ok, lattice_maps@crdt:crdt()} |
{error, nil}.
valid_value(Map, Key) ->
case gleam_stdlib:map_get(erlang:element(5, Map), Key) of
{ok, Value} ->
case matches_spec(Value, erlang:element(3, Map)) of
true ->
{ok, Value};
false ->
{ok,
lattice_maps@crdt:default_crdt(
erlang:element(3, Map),
erlang:element(2, Map)
)}
end;
{error, _} ->
{error, nil}
end.
-file("src/lattice_maps/or_map.gleam", 179).
?DOC(
" Merge two OR-Maps.\n"
"\n"
" The OR-Set key trackers are merged with add-wins semantics: if a key was\n"
" concurrently updated on one replica and removed on another, the key\n"
" survives in the merged result. CRDT values are merged per-key using\n"
" `crdt.merge` for type-specific convergence.\n"
"\n"
" Returns `Error(TypeMismatch(...))` if the two maps have different\n"
" `crdt_spec` values (e.g., one holds counters and the other holds sets).\n"
).
-spec merge(o_r_map(), o_r_map()) -> {ok, o_r_map()} |
{error, lattice_maps@crdt:merge_error()}.
merge(A, B) ->
case erlang:element(3, A) =:= erlang:element(3, B) of
false ->
{error,
{type_mismatch,
spec_to_string(erlang:element(3, A)),
spec_to_string(erlang:element(3, B))}};
true ->
Merged_key_set = lattice_sets@or_set:merge(
erlang:element(4, A),
erlang:element(4, B)
),
Active_keys = lattice_sets@or_set:value(Merged_key_set),
All_value_keys = gleam@set:to_list(
gleam@set:union(
gleam@set:from_list(maps:keys(erlang:element(5, A))),
gleam@set:from_list(maps:keys(erlang:element(5, B)))
)
),
Merged_values = gleam@list:fold(
All_value_keys,
maps:new(),
fun(Acc, Key) ->
Merged_crdt = case {valid_value(A, Key),
valid_value(B, Key)} of
{{ok, Ca}, {ok, Cb}} ->
case lattice_maps@crdt:merge(Ca, Cb) of
{ok, Merged} ->
Merged;
{error, _} ->
lattice_maps@crdt:default_crdt(
erlang:element(3, A),
erlang:element(2, A)
)
end;
{{ok, Ca@1}, {error, _}} ->
Ca@1;
{{error, _}, {ok, Cb@1}} ->
Cb@1;
{{error, _}, {error, _}} ->
erlang:error(#{gleam_error => panic,
message => <<"unreachable: key must exist in at least one map"/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"lattice_maps/or_map"/utf8>>,
function => <<"merge"/utf8>>,
line => 205})
end,
gleam@dict:insert(Acc, Key, Merged_crdt)
end
),
All_bound_keys = gleam@set:to_list(
gleam@set:union(
gleam@set:from_list(maps:keys(erlang:element(6, A))),
gleam@set:from_list(maps:keys(erlang:element(6, B)))
)
),
Merged_bounds = gleam@list:fold(
All_bound_keys,
maps:new(),
fun(Acc@1, Key@1) ->
case gleam@set:contains(Active_keys, Key@1) of
true ->
Acc@1;
false ->
case {gleam_stdlib:map_get(
erlang:element(6, A),
Key@1
),
gleam_stdlib:map_get(
erlang:element(6, B),
Key@1
)} of
{{ok, Ba}, {ok, Bb}} ->
gleam@dict:insert(
Acc@1,
Key@1,
lattice_core@version_vector:merge(
Ba,
Bb
)
);
{{ok, Ba@1}, {error, _}} ->
gleam@dict:insert(Acc@1, Key@1, Ba@1);
{{error, _}, {ok, Bb@1}} ->
gleam@dict:insert(Acc@1, Key@1, Bb@1);
{{error, _}, {error, _}} ->
Acc@1
end
end
end
),
{ok,
{o_r_map,
erlang:element(2, A),
erlang:element(3, A),
Merged_key_set,
Merged_values,
Merged_bounds}}
end.
-file("src/lattice_maps/or_map.gleam", 505).
-spec crdt_name(lattice_maps@crdt:crdt()) -> binary().
crdt_name(Value) ->
case Value of
{crdt_g_counter, _} ->
<<"g_counter"/utf8>>;
{crdt_pn_counter, _} ->
<<"pn_counter"/utf8>>;
{crdt_lww_register, _} ->
<<"lww_register"/utf8>>;
{crdt_mv_register, _} ->
<<"mv_register"/utf8>>;
{crdt_g_set, _} ->
<<"g_set"/utf8>>;
{crdt_two_p_set, _} ->
<<"two_p_set"/utf8>>;
{crdt_or_set, _} ->
<<"or_set"/utf8>>;
{crdt_version_vector, _} ->
<<"version_vector"/utf8>>
end.
-file("src/lattice_maps/or_map.gleam", 422).
-spec decode_or_map_state(
binary(),
binary(),
binary(),
list({binary(), binary()}),
gleam@dict:dict(binary(), lattice_core@version_vector:version_vector())
) -> {ok, o_r_map()} | {error, gleam@json:decode_error()}.
decode_or_map_state(
Replica_id_str,
Crdt_spec_str,
Key_set_str,
Values_list,
Remove_bounds
) ->
case string_to_spec(Crdt_spec_str) of
{error, _} ->
{error,
{unable_to_decode,
[{decode_error,
<<"known CrdtSpec"/utf8>>,
Crdt_spec_str,
[<<"state"/utf8>>, <<"crdt_spec"/utf8>>]}]}};
{ok, Crdt_spec} ->
case lattice_sets@or_set:from_json(Key_set_str) of
{error, E} ->
{error, E};
{ok, Key_set} ->
Values_result = gleam@list:try_map(
Values_list,
fun(Pair) ->
{Key, Crdt_str} = Pair,
case lattice_maps@crdt:from_json(Crdt_str) of
{ok, C} ->
case matches_spec(C, Crdt_spec) of
true ->
{ok, {Key, C}};
false ->
{error,
{unable_to_decode,
[{decode_error,
spec_to_string(
Crdt_spec
),
crdt_name(C),
[<<"state"/utf8>>,
<<"values"/utf8>>]}]}}
end;
{error, E@1} ->
{error, E@1}
end
end
),
case Values_result of
{error, E@2} ->
{error, E@2};
{ok, Pairs} ->
{ok,
{o_r_map,
lattice_core@replica_id:new(Replica_id_str),
Crdt_spec,
Key_set,
maps:from_list(Pairs),
Remove_bounds}}
end
end
end.
-file("src/lattice_maps/or_map.gleam", 340).
?DOC(
" Decode an `ORMap` from a JSON string produced by `to_json`.\n"
"\n"
" Supports both v1 (no remove_bounds) and v2 (with remove_bounds) formats.\n"
" v1 maps are decoded with empty remove_bounds, meaning no value compaction\n"
" is possible until new removes are performed.\n"
"\n"
" Returns `Error` if the string is not valid JSON, does not match the\n"
" expected format, or contains an unknown `crdt_spec` string.\n"
).
-spec from_json(binary()) -> {ok, o_r_map()} |
{error, gleam@json:decode_error()}.
from_json(Json_string) ->
Value_pair_decoder = begin
gleam@dynamic@decode:field(
<<"key"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Key) ->
gleam@dynamic@decode:field(
<<"crdt"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Crdt_str) ->
gleam@dynamic@decode:success({Key, Crdt_str})
end
)
end
)
end,
Bounds_decoder = gleam@dynamic@decode:dict(
{decoder, fun gleam@dynamic@decode:decode_string/1},
lattice_core@version_vector:decoder()
),
State_decoder = begin
gleam@dynamic@decode:field(
<<"state"/utf8>>,
begin
gleam@dynamic@decode:field(
<<"replica_id"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Replica_id_str) ->
gleam@dynamic@decode:field(
<<"crdt_spec"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Crdt_spec_str) ->
gleam@dynamic@decode:field(
<<"key_set"/utf8>>,
{decoder,
fun gleam@dynamic@decode:decode_string/1},
fun(Key_set_str) ->
gleam@dynamic@decode:field(
<<"values"/utf8>>,
gleam@dynamic@decode:list(
Value_pair_decoder
),
fun(Values_list) ->
gleam@dynamic@decode:optional_field(
<<"remove_bounds"/utf8>>,
maps:new(),
Bounds_decoder,
fun(Remove_bounds) ->
gleam@dynamic@decode:success(
{Replica_id_str,
Crdt_spec_str,
Key_set_str,
Values_list,
Remove_bounds}
)
end
)
end
)
end
)
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 =:= <<"or_map"/utf8>> of
false ->
{error,
{unable_to_decode,
[{decode_error,
<<"type=or_map"/utf8>>,
Type_tag@1,
[]}]}};
true ->
case Version@1 of
1 ->
case gleam@json:parse(Json_string, State_decoder) of
{error, E@1} ->
{error, E@1};
{ok,
{Replica_id_str@1,
Crdt_spec_str@1,
Key_set_str@1,
Values_list@1,
Remove_bounds@1}} ->
decode_or_map_state(
Replica_id_str@1,
Crdt_spec_str@1,
Key_set_str@1,
Values_list@1,
Remove_bounds@1
)
end;
2 ->
case gleam@json:parse(Json_string, State_decoder) of
{error, E@1} ->
{error, E@1};
{ok,
{Replica_id_str@1,
Crdt_spec_str@1,
Key_set_str@1,
Values_list@1,
Remove_bounds@1}} ->
decode_or_map_state(
Replica_id_str@1,
Crdt_spec_str@1,
Key_set_str@1,
Values_list@1,
Remove_bounds@1
)
end;
_ ->
{error,
{unable_to_decode,
[{decode_error,
<<"v=1 or v=2"/utf8>>,
erlang:integer_to_binary(Version@1),
[<<"v"/utf8>>]}]}}
end
end
end.