Packages

Conflict-free replicated data types (CRDTs) for Gleam — umbrella package that includes all lattice sub-packages

Current section

Files

Jump to
lattice_crdt src lattice@or_map.erl
Raw

src/lattice@or_map.erl

-module(lattice@or_map).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/lattice/or_map.gleam").
-export([new/2, update/3, get/2, remove/2, keys/1, values/1, merge/2, to_json/1, 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/crdt\n"
" import lattice/g_counter\n"
" import lattice/or_map\n"
"\n"
" let map = or_map.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"
).
-type o_r_map() :: {o_r_map,
binary(),
lattice@crdt:crdt_spec(),
lattice@or_set:o_r_set(binary()),
gleam@dict:dict(binary(), lattice@crdt:crdt())}.
-file("src/lattice/or_map.gleam", 49).
-spec spec_to_string(lattice@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/or_map.gleam", 61).
-spec string_to_spec(binary()) -> {ok, lattice@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/or_map.gleam", 78).
?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(binary(), lattice@crdt:crdt_spec()) -> o_r_map().
new(Replica_id, Crdt_spec) ->
{o_r_map, Replica_id, Crdt_spec, lattice@or_set:new(Replica_id), maps:new()}.
-file("src/lattice/or_map.gleam", 92).
?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@crdt:crdt()) -> lattice@crdt:crdt())
) -> o_r_map().
update(Map, Key, F) ->
Current = case gleam_stdlib:map_get(erlang:element(5, Map), Key) of
{ok, Crdt_val} ->
Crdt_val;
{error, _} ->
lattice@crdt:default_crdt(
erlang:element(3, Map),
erlang:element(2, Map)
)
end,
Updated = F(Current),
{o_r_map,
erlang:element(2, Map),
erlang:element(3, Map),
lattice@or_set:add(erlang:element(4, Map), Key),
gleam@dict:insert(erlang:element(5, Map), Key, Updated)}.
-file("src/lattice/or_map.gleam", 111).
?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@crdt:crdt()} | {error, nil}.
get(Map, Key) ->
case lattice@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/or_map.gleam", 127).
?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 in the values dict so it can participate in\n"
" per-key merge if the key is concurrently re-added on another replica.\n"
).
-spec remove(o_r_map(), binary()) -> o_r_map().
remove(Map, Key) ->
{o_r_map,
erlang:element(2, Map),
erlang:element(3, Map),
lattice@or_set:remove(erlang:element(4, Map), Key),
erlang:element(5, Map)}.
-file("src/lattice/or_map.gleam", 134).
?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@or_set:value(erlang:element(4, Map))).
-file("src/lattice/or_map.gleam", 141).
?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@crdt:crdt()).
values(Map) ->
Active_keys = lattice@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/or_map.gleam", 159).
?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"
" Merge is commutative, associative, and idempotent (a valid CRDT join).\n"
).
-spec merge(o_r_map(), o_r_map()) -> o_r_map().
merge(A, B) ->
Merged_key_set = lattice@or_set:merge(
erlang:element(4, A),
erlang:element(4, B)
),
All_value_keys = gleam@list:unique(
lists:append(
maps:keys(erlang:element(5, A)),
maps:keys(erlang:element(5, B))
)
),
Merged_values = gleam@list:fold(
All_value_keys,
maps:new(),
fun(Acc, Key) ->
Merged_crdt = case {gleam_stdlib:map_get(erlang:element(5, A), Key),
gleam_stdlib:map_get(erlang:element(5, B), Key)} of
{{ok, Ca}, {ok, Cb}} ->
lattice@crdt:merge(Ca, Cb);
{{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/or_map"/utf8>>,
function => <<"merge"/utf8>>,
line => 170})
end,
gleam@dict:insert(Acc, Key, Merged_crdt)
end
),
{o_r_map,
erlang:element(2, A),
erlang:element(3, A),
Merged_key_set,
Merged_values}.
-file("src/lattice/or_map.gleam", 190).
?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\": 1, \"state\": {\"replica_id\": \"...\", \"crdt_spec\": \"...\", \"key_set\": \"...\", \"values\": [...]}}`\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, Replica_id, Crdt_spec, Key_set, Values} = 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@crdt:to_json(Crdt_val))
)}]
)
end
),
gleam@json:object(
[{<<"type"/utf8>>, gleam@json:string(<<"or_map"/utf8>>)},
{<<"v"/utf8>>, gleam@json:int(1)},
{<<"state"/utf8>>,
gleam@json:object(
[{<<"replica_id"/utf8>>, gleam@json:string(Replica_id)},
{<<"crdt_spec"/utf8>>,
gleam@json:string(spec_to_string(Crdt_spec))},
{<<"key_set"/utf8>>,
gleam@json:string(
gleam@json:to_string(
lattice@or_set:to_json(Key_set)
)
)},
{<<"values"/utf8>>, Values_json}]
)}]
).
-file("src/lattice/or_map.gleam", 219).
?DOC(
" Decode an `ORMap` from a JSON string produced by `to_json`.\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,
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) ->
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:success(
{Replica_id,
Crdt_spec_str,
Key_set_str,
Values_list}
)
end
)
end
)
end
)
end
)
end,
fun(State) -> gleam@dynamic@decode:success(State) end
)
end,
case gleam@json:parse(Json_string, Decoder) of
{error, E} ->
{error, E};
{ok, {Replica_id@1, Crdt_spec_str@1, Key_set_str@1, Values_list@1}} ->
case string_to_spec(Crdt_spec_str@1) of
{error, _} ->
{error,
{unable_to_decode,
[{decode_error,
<<"known CrdtSpec"/utf8>>,
Crdt_spec_str@1,
[<<"state"/utf8>>, <<"crdt_spec"/utf8>>]}]}};
{ok, Crdt_spec} ->
case lattice@or_set:from_json(Key_set_str@1) of
{error, E@1} ->
{error, E@1};
{ok, Key_set} ->
Values_result = gleam@list:try_map(
Values_list@1,
fun(Pair) ->
{Key@1, Crdt_str@1} = Pair,
case lattice@crdt:from_json(Crdt_str@1) of
{ok, C} ->
{ok, {Key@1, C}};
{error, E@2} ->
{error, E@2}
end
end
),
case Values_result of
{error, E@3} ->
{error, E@3};
{ok, Pairs} ->
{ok,
{o_r_map,
Replica_id@1,
Crdt_spec,
Key_set,
maps:from_list(Pairs)}}
end
end
end
end.