Packages
macula
0.44.0
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_content_system/macula_content_store.erl
%%%-------------------------------------------------------------------
%%% @doc
%%% Content store module for Macula content-addressed storage.
%%%
%%% Provides local storage and retrieval of content blocks and manifests.
%%% Uses file-based storage with directory sharding (first 2 hex chars
%%% of hash) for efficient organization.
%%%
%%% == Storage Layout ==
%%% ```
%%% {base_dir}/
%%% ├── blocks/
%%% │ ├── 5d/
%%% │ │ └── 5d41402abc4b2a76b9719d911017c592.blk
%%% │ └── ...
%%% ├── manifests/
%%% │ ├── 7f/
%%% │ │ └── 7f83b1657ff1fc53b92dc18148a1d65d.man
%%% │ └── ...
%%% └── index.dets
%%% '''
%%%
%%% == Example Usage ==
%%% ```
%%% %% Store a block
%%% Data = <<"content">>,
%%% MCID = compute_mcid(Data),
%%% ok = macula_content_store:put_block(MCID, Data),
%%%
%%% %% Retrieve with verification
%%% {ok, Data} = macula_content_store:get_block(MCID).
%%% '''
%%% @end
%%%-------------------------------------------------------------------
-module(macula_content_store).
-behaviour(gen_server).
%% API
-export([
start_link/1,
stop/0,
%% Block operations
put_block/2,
get_block/1,
has_block/1,
delete_block/1,
%% Manifest operations
put_manifest/1,
get_manifest/1,
list_manifests/0,
delete_manifest/1,
%% Maintenance
gc/0,
stats/0,
verify_integrity/0,
%% Path helpers (for testing)
block_path/1,
manifest_path/1
]).
%% gen_server callbacks
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
-define(SERVER, ?MODULE).
-define(BLOCK_EXT, ".blk").
-define(MANIFEST_EXT, ".man").
-record(state, {
base_dir :: string(),
blocks_dir :: string(),
manifests_dir :: string(),
block_index :: ets:tid(), %% MCID -> {Size, Timestamp}
manifest_index :: ets:tid() %% MCID -> {Name, Size, ChunkCount, Timestamp}
}).
%%%===================================================================
%%% API Functions
%%%===================================================================
%% @doc Start the content store server.
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link({local, ?SERVER}, ?MODULE, Opts, []).
%% @doc Stop the content store server.
-spec stop() -> ok.
stop() ->
gen_server:stop(?SERVER).
%% @doc Store a block by its MCID.
%% Verifies the data matches the MCID hash before storing.
-spec put_block(binary(), binary()) -> ok | {error, term()}.
put_block(MCID, Data) ->
gen_server:call(?SERVER, {put_block, MCID, Data}).
%% @doc Retrieve a block by its MCID.
%% Verifies the retrieved data matches the MCID hash.
-spec get_block(binary()) -> {ok, binary()} | {error, not_found | hash_mismatch}.
get_block(MCID) ->
gen_server:call(?SERVER, {get_block, MCID}).
%% @doc Check if a block exists in the store.
-spec has_block(binary()) -> boolean().
has_block(MCID) ->
gen_server:call(?SERVER, {has_block, MCID}).
%% @doc Delete a block from the store.
-spec delete_block(binary()) -> ok.
delete_block(MCID) ->
gen_server:call(?SERVER, {delete_block, MCID}).
%% @doc Store a manifest.
-spec put_manifest(map()) -> ok | {error, term()}.
put_manifest(Manifest) ->
gen_server:call(?SERVER, {put_manifest, Manifest}).
%% @doc Retrieve a manifest by its MCID.
-spec get_manifest(binary()) -> {ok, map()} | {error, not_found}.
get_manifest(MCID) ->
gen_server:call(?SERVER, {get_manifest, MCID}).
%% @doc List all manifest MCIDs.
-spec list_manifests() -> [binary()].
list_manifests() ->
gen_server:call(?SERVER, list_manifests).
%% @doc Delete a manifest from the store.
-spec delete_manifest(binary()) -> ok.
delete_manifest(MCID) ->
gen_server:call(?SERVER, {delete_manifest, MCID}).
%% @doc Garbage collect orphaned blocks not referenced by any manifest.
-spec gc() -> {ok, #{removed := non_neg_integer()}}.
gc() ->
gen_server:call(?SERVER, gc, 60000).
%% @doc Get storage statistics.
-spec stats() -> map().
stats() ->
gen_server:call(?SERVER, stats).
%% @doc Verify integrity of all stored blocks.
-spec verify_integrity() -> {ok, non_neg_integer()} | {error, {corrupted, [binary()]}}.
verify_integrity() ->
gen_server:call(?SERVER, verify_integrity, 60000).
%% @doc Get the file path for a block MCID.
%% Uses first 2 hex chars of hash for directory sharding.
-spec block_path(binary()) -> string().
block_path(<<_Version:8, _Codec:8, Hash:32/binary>>) ->
BaseDir = get_default_base_dir(),
BlocksDir = filename:join(BaseDir, "blocks"),
HexHash = macula_content_hasher:hex_encode(Hash),
<<Shard:2/binary, _Rest/binary>> = HexHash,
filename:join([BlocksDir, binary_to_list(Shard),
binary_to_list(HexHash) ++ ?BLOCK_EXT]).
%% @doc Get the file path for a manifest MCID.
-spec manifest_path(binary()) -> string().
manifest_path(<<_Version:8, _Codec:8, Hash:32/binary>>) ->
BaseDir = get_default_base_dir(),
ManifestsDir = filename:join(BaseDir, "manifests"),
HexHash = macula_content_hasher:hex_encode(Hash),
<<Shard:2/binary, _Rest/binary>> = HexHash,
filename:join([ManifestsDir, binary_to_list(Shard),
binary_to_list(HexHash) ++ ?MANIFEST_EXT]).
%%%===================================================================
%%% gen_server Callbacks
%%%===================================================================
init(Opts) ->
%% Support both store_path and base_dir keys for flexibility
BaseDir = case maps:get(store_path, Opts, undefined) of
undefined -> maps:get(base_dir, Opts, get_default_base_dir());
Path -> Path
end,
BlocksDir = filename:join(BaseDir, "blocks"),
ManifestsDir = filename:join(BaseDir, "manifests"),
%% Ensure directories exist (create if missing)
ensure_dir_exists(BlocksDir),
ensure_dir_exists(ManifestsDir),
%% Create ETS tables for indexing
BlockIndex = ets:new(content_block_index, [set, private]),
ManifestIndex = ets:new(content_manifest_index, [set, private]),
%% Load existing files into index
State = #state{
base_dir = BaseDir,
blocks_dir = BlocksDir,
manifests_dir = ManifestsDir,
block_index = BlockIndex,
manifest_index = ManifestIndex
},
rebuild_index(State),
{ok, State}.
handle_call({put_block, MCID, Data}, _From, State) ->
Result = do_put_block(MCID, Data, State),
{reply, Result, State};
handle_call({get_block, MCID}, _From, State) ->
Result = do_get_block(MCID, State),
{reply, Result, State};
handle_call({has_block, MCID}, _From, State) ->
<<_:8, _:8, Hash:32/binary>> = MCID,
Result = ets:member(State#state.block_index, Hash),
{reply, Result, State};
handle_call({delete_block, MCID}, _From, State) ->
Result = do_delete_block(MCID, State),
{reply, Result, State};
handle_call({put_manifest, Manifest}, _From, State) ->
Result = do_put_manifest(Manifest, State),
{reply, Result, State};
handle_call({get_manifest, MCID}, _From, State) ->
Result = do_get_manifest(MCID, State),
{reply, Result, State};
handle_call(list_manifests, _From, State) ->
MCIDs = [<<1, 16#56, Hash/binary>> || {Hash, _} <- ets:tab2list(State#state.manifest_index)],
{reply, MCIDs, State};
handle_call({delete_manifest, MCID}, _From, State) ->
Result = do_delete_manifest(MCID, State),
{reply, Result, State};
handle_call(gc, _From, State) ->
Result = do_gc(State),
{reply, Result, State};
handle_call(stats, _From, State) ->
Stats = do_stats(State),
{reply, Stats, State};
handle_call(verify_integrity, _From, State) ->
Result = do_verify_integrity(State),
{reply, Result, State}.
handle_cast(_Msg, State) ->
{noreply, State}.
handle_info(_Info, State) ->
{noreply, State}.
terminate(_Reason, _State) ->
ok.
%%%===================================================================
%%% Internal Functions - Directory Management
%%%===================================================================
%% @private Ensure directory exists, creating it if necessary.
ensure_dir_exists(Dir) ->
case filelib:is_dir(Dir) of
true ->
ok;
false ->
%% ensure_dir needs a trailing slash or file to create parent dirs
%% This creates all parent directories
case filelib:ensure_dir(filename:join(Dir, "dummy")) of
ok ->
%% Now create the final directory
case file:make_dir(Dir) of
ok -> ok;
{error, eexist} -> ok; %% Race condition, dir created by another process
{error, Reason} -> {error, Reason}
end;
{error, Reason} ->
{error, Reason}
end
end.
%%%===================================================================
%%% Internal Functions - Block Operations
%%%===================================================================
do_put_block(<<_Version:8, _Codec:8, Hash:32/binary>> = _MCID, Data, State) ->
%% Verify data matches hash
ActualHash = macula_content_hasher:hash(blake3, Data),
case ActualHash =:= Hash of
true ->
Path = compute_block_path(Hash, State),
ok = filelib:ensure_dir(Path),
ok = file:write_file(Path, Data),
ets:insert(State#state.block_index, {Hash, {byte_size(Data), erlang:system_time(second)}}),
ok;
false ->
%% Try SHA256 as fallback
Sha256Hash = macula_content_hasher:hash(sha256, Data),
case Sha256Hash =:= Hash of
true ->
Path = compute_block_path(Hash, State),
ok = filelib:ensure_dir(Path),
ok = file:write_file(Path, Data),
ets:insert(State#state.block_index, {Hash, {byte_size(Data), erlang:system_time(second)}}),
ok;
false ->
{error, hash_mismatch}
end
end.
do_get_block(<<_Version:8, _Codec:8, Hash:32/binary>> = _MCID, State) ->
case ets:lookup(State#state.block_index, Hash) of
[] ->
{error, not_found};
[{Hash, _}] ->
Path = compute_block_path(Hash, State),
case file:read_file(Path) of
{ok, Data} ->
%% Verify hash on retrieval
ActualBlake3 = macula_content_hasher:hash(blake3, Data),
ActualSha256 = macula_content_hasher:hash(sha256, Data),
case ActualBlake3 =:= Hash orelse ActualSha256 =:= Hash of
true -> {ok, Data};
false -> {error, hash_mismatch}
end;
{error, _} ->
{error, not_found}
end
end.
do_delete_block(<<_Version:8, _Codec:8, Hash:32/binary>>, State) ->
Path = compute_block_path(Hash, State),
file:delete(Path),
ets:delete(State#state.block_index, Hash),
ok.
%%%===================================================================
%%% Internal Functions - Manifest Operations
%%%===================================================================
do_put_manifest(Manifest, State) ->
MCID = maps:get(mcid, Manifest),
<<_Version:8, _Codec:8, Hash:32/binary>> = MCID,
Path = compute_manifest_path(Hash, State),
ok = filelib:ensure_dir(Path),
{ok, Encoded} = macula_content_manifest:encode(Manifest),
ok = file:write_file(Path, Encoded),
Info = {
maps:get(name, Manifest),
maps:get(size, Manifest),
maps:get(chunk_count, Manifest),
erlang:system_time(second)
},
ets:insert(State#state.manifest_index, {Hash, Info}),
ok.
do_get_manifest(<<_Version:8, _Codec:8, Hash:32/binary>>, State) ->
case ets:lookup(State#state.manifest_index, Hash) of
[] ->
{error, not_found};
[{Hash, _}] ->
Path = compute_manifest_path(Hash, State),
case file:read_file(Path) of
{ok, Encoded} ->
macula_content_manifest:decode(Encoded);
{error, _} ->
{error, not_found}
end
end.
do_delete_manifest(<<_Version:8, _Codec:8, Hash:32/binary>>, State) ->
Path = compute_manifest_path(Hash, State),
file:delete(Path),
ets:delete(State#state.manifest_index, Hash),
ok.
%%%===================================================================
%%% Internal Functions - Maintenance
%%%===================================================================
do_gc(State) ->
%% Collect all chunk hashes referenced by manifests
ReferencedHashes = collect_referenced_hashes(State),
%% Find orphaned blocks
AllBlocks = [Hash || {Hash, _} <- ets:tab2list(State#state.block_index)],
Orphans = lists:filter(fun(Hash) ->
not sets:is_element(Hash, ReferencedHashes)
end, AllBlocks),
%% Delete orphans
lists:foreach(fun(Hash) ->
MCID = <<1, 16#55, Hash/binary>>,
do_delete_block(MCID, State)
end, Orphans),
{ok, #{removed => length(Orphans)}}.
collect_referenced_hashes(State) ->
Manifests = [Hash || {Hash, _} <- ets:tab2list(State#state.manifest_index)],
lists:foldl(fun(ManifestHash, Acc) ->
MCID = <<1, 16#56, ManifestHash/binary>>,
case do_get_manifest(MCID, State) of
{ok, Manifest} ->
Chunks = maps:get(chunks, Manifest, []),
lists:foldl(fun(ChunkInfo, InnerAcc) ->
ChunkHash = maps:get(hash, ChunkInfo),
sets:add_element(ChunkHash, InnerAcc)
end, Acc, Chunks);
{error, _} ->
Acc
end
end, sets:new(), Manifests).
do_stats(State) ->
BlockCount = ets:info(State#state.block_index, size),
ManifestCount = ets:info(State#state.manifest_index, size),
TotalSize = ets:foldl(fun({_Hash, {Size, _}}, Acc) ->
Acc + Size
end, 0, State#state.block_index),
#{
block_count => BlockCount,
manifest_count => ManifestCount,
total_size => TotalSize
}.
do_verify_integrity(State) ->
Corrupted = ets:foldl(fun({Hash, _}, Acc) ->
MCID = <<1, 16#55, Hash/binary>>,
case do_get_block(MCID, State) of
{ok, _} -> Acc;
{error, hash_mismatch} -> [MCID | Acc];
{error, not_found} -> [MCID | Acc]
end
end, [], State#state.block_index),
case Corrupted of
[] ->
{ok, ets:info(State#state.block_index, size)};
_ ->
{error, {corrupted, Corrupted}}
end.
%%%===================================================================
%%% Internal Functions - Paths
%%%===================================================================
compute_block_path(Hash, State) ->
HexHash = macula_content_hasher:hex_encode(Hash),
<<Shard:2/binary, _Rest/binary>> = HexHash,
filename:join([State#state.blocks_dir, binary_to_list(Shard),
binary_to_list(HexHash) ++ ?BLOCK_EXT]).
compute_manifest_path(Hash, State) ->
HexHash = macula_content_hasher:hex_encode(Hash),
<<Shard:2/binary, _Rest/binary>> = HexHash,
filename:join([State#state.manifests_dir, binary_to_list(Shard),
binary_to_list(HexHash) ++ ?MANIFEST_EXT]).
%%%===================================================================
%%% Internal Functions - Index Rebuilding
%%%===================================================================
rebuild_index(State) ->
%% Scan blocks directory
BlockPattern = filename:join([State#state.blocks_dir, "*", "*" ++ ?BLOCK_EXT]),
BlockFiles = filelib:wildcard(BlockPattern),
lists:foreach(fun(Path) ->
case file:read_file_info(Path) of
{ok, Info} ->
Basename = filename:basename(Path, ?BLOCK_EXT),
case macula_content_hasher:hex_decode(list_to_binary(Basename)) of
{ok, Hash} ->
Size = element(2, Info),
ets:insert(State#state.block_index, {Hash, {Size, 0}});
{error, _} ->
ok
end;
{error, _} ->
ok
end
end, BlockFiles),
%% Scan manifests directory
ManifestPattern = filename:join([State#state.manifests_dir, "*", "*" ++ ?MANIFEST_EXT]),
ManifestFiles = filelib:wildcard(ManifestPattern),
lists:foreach(fun(Path) ->
Basename = filename:basename(Path, ?MANIFEST_EXT),
case macula_content_hasher:hex_decode(list_to_binary(Basename)) of
{ok, Hash} ->
case file:read_file(Path) of
{ok, Encoded} ->
case macula_content_manifest:decode(Encoded) of
{ok, Manifest} ->
Info = {
maps:get(name, Manifest),
maps:get(size, Manifest),
maps:get(chunk_count, Manifest),
0
},
ets:insert(State#state.manifest_index, {Hash, Info});
{error, _} ->
ok
end;
{error, _} ->
ok
end;
{error, _} ->
ok
end
end, ManifestFiles),
ok.
get_default_base_dir() ->
case application:get_env(macula, content_store_dir) of
{ok, Dir} -> Dir;
undefined -> "/var/lib/macula/content"
end.