Current section
Files
Jump to
Current section
Files
src/goose.erl
-module(goose).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/goose.gleam").
-export([default_config/0, build_url/1, connect/3, start_consumer/2, parse_event/1]).
-export_type([jetstream_event/0, commit_data/0, identity_data/0, account_data/0, jetstream_config/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.
-type jetstream_event() :: {commit_event, binary(), integer(), commit_data()} |
{identity_event, binary(), integer(), identity_data()} |
{account_event, binary(), integer(), account_data()} |
{unknown_event, binary()}.
-type commit_data() :: {commit_data,
binary(),
binary(),
binary(),
binary(),
gleam@option:option(gleam@dynamic:dynamic_()),
gleam@option:option(binary())}.
-type identity_data() :: {identity_data,
binary(),
binary(),
integer(),
binary()}.
-type account_data() :: {account_data, boolean(), binary(), integer(), binary()}.
-type jetstream_config() :: {jetstream_config,
binary(),
list(binary()),
list(binary()),
gleam@option:option(integer()),
gleam@option:option(integer()),
boolean(),
boolean()}.
-file("src/goose.gleam", 51).
?DOC(" Create a default configuration for US East endpoint\n").
-spec default_config() -> jetstream_config().
default_config() ->
{jetstream_config,
<<"wss://jetstream2.us-east.bsky.network/subscribe"/utf8>>,
[],
[],
none,
none,
false,
false}.
-file("src/goose.gleam", 64).
?DOC(" Build the WebSocket URL with query parameters\n").
-spec build_url(jetstream_config()) -> binary().
build_url(Config) ->
Base = erlang:element(2, Config),
Mut_params = [],
Mut_params@1 = case erlang:element(3, Config) of
[] ->
Mut_params;
Collections ->
Collection_params = gleam@list:map(
Collections,
fun(Col) -> <<"wantedCollections="/utf8, Col/binary>> end
),
lists:append(Collection_params, Mut_params)
end,
Mut_params@2 = case erlang:element(4, Config) of
[] ->
Mut_params@1;
Dids ->
Did_params = gleam@list:map(
Dids,
fun(Did) -> <<"wantedDids="/utf8, Did/binary>> end
),
lists:append(Did_params, Mut_params@1)
end,
Mut_params@3 = case erlang:element(5, Config) of
none ->
Mut_params@2;
{some, Cursor_val} ->
lists:append(
[<<"cursor="/utf8, (gleam@string:inspect(Cursor_val))/binary>>],
Mut_params@2
)
end,
Mut_params@4 = case erlang:element(6, Config) of
none ->
Mut_params@3;
{some, Size_val} ->
lists:append(
[<<"maxMessageSizeBytes="/utf8,
(gleam@string:inspect(Size_val))/binary>>],
Mut_params@3
)
end,
Mut_params@5 = case erlang:element(7, Config) of
false ->
lists:append([<<"compress=false"/utf8>>], Mut_params@4);
true ->
lists:append([<<"compress=true"/utf8>>], Mut_params@4)
end,
Mut_params@6 = case erlang:element(8, Config) of
false ->
lists:append([<<"requireHello=false"/utf8>>], Mut_params@5);
true ->
lists:append([<<"requireHello=true"/utf8>>], Mut_params@5)
end,
case Mut_params@6 of
[] ->
Base;
Params ->
<<<<Base/binary, "?"/utf8>>/binary,
(gleam@string:join(lists:reverse(Params), <<"&"/utf8>>))/binary>>
end.
-file("src/goose.gleam", 124).
?DOC(" Connect to Jetstream WebSocket using Erlang gun library\n").
-spec connect(binary(), gleam@erlang@process:pid_(), boolean()) -> {ok,
gleam@erlang@process:pid_()} |
{error, gleam@dynamic:dynamic_()}.
connect(Url, Handler_pid, Compress) ->
goose_ws_ffi:connect(Url, Handler_pid, Compress).
-file("src/goose.gleam", 151).
?DOC(" Receive loop for WebSocket messages\n").
-spec receive_loop(fun((binary()) -> nil)) -> nil.
receive_loop(On_event) ->
case goose_ffi:receive_ws_message() of
{ok, Text} ->
On_event(Text),
receive_loop(On_event);
{error, _} ->
receive_loop(On_event)
end.
-file("src/goose.gleam", 131).
?DOC(" Start consuming the Jetstream feed\n").
-spec start_consumer(jetstream_config(), fun((binary()) -> nil)) -> nil.
start_consumer(Config, On_event) ->
Url = build_url(Config),
Self = erlang:self(),
Result = goose_ws_ffi:connect(Url, Self, erlang:element(7, Config)),
case Result of
{ok, _} ->
receive_loop(On_event);
{error, Err} ->
gleam_stdlib:println(<<"Failed to connect to Jetstream"/utf8>>),
gleam_stdlib:println_error(gleam@string:inspect(Err))
end.
-file("src/goose.gleam", 208).
?DOC(" Decoder for commit with record (create/update operations)\n").
-spec commit_with_record_decoder() -> gleam@dynamic@decode:decoder(commit_data()).
commit_with_record_decoder() ->
gleam@dynamic@decode:field(
<<"rev"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Rev) ->
gleam@dynamic@decode:field(
<<"operation"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Operation) ->
gleam@dynamic@decode:field(
<<"collection"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Collection) ->
gleam@dynamic@decode:field(
<<"rkey"/utf8>>,
{decoder,
fun gleam@dynamic@decode:decode_string/1},
fun(Rkey) ->
gleam@dynamic@decode:field(
<<"record"/utf8>>,
{decoder,
fun gleam@dynamic@decode:decode_dynamic/1},
fun(Record) ->
gleam@dynamic@decode:field(
<<"cid"/utf8>>,
{decoder,
fun gleam@dynamic@decode:decode_string/1},
fun(Cid) ->
gleam@dynamic@decode:success(
{commit_data,
Rev,
Operation,
Collection,
Rkey,
{some, Record},
{some, Cid}}
)
end
)
end
)
end
)
end
)
end
)
end
).
-file("src/goose.gleam", 226).
?DOC(" Decoder for commit without record (delete operations)\n").
-spec commit_without_record_decoder() -> gleam@dynamic@decode:decoder(commit_data()).
commit_without_record_decoder() ->
gleam@dynamic@decode:field(
<<"rev"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Rev) ->
gleam@dynamic@decode:field(
<<"operation"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Operation) ->
gleam@dynamic@decode:field(
<<"collection"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Collection) ->
gleam@dynamic@decode:field(
<<"rkey"/utf8>>,
{decoder,
fun gleam@dynamic@decode:decode_string/1},
fun(Rkey) ->
gleam@dynamic@decode:success(
{commit_data,
Rev,
Operation,
Collection,
Rkey,
none,
none}
)
end
)
end
)
end
)
end
).
-file("src/goose.gleam", 199).
?DOC(" Decoder for commit data - handles both create/update (with record) and delete (without)\n").
-spec commit_data_decoder() -> gleam@dynamic@decode:decoder(commit_data()).
commit_data_decoder() ->
gleam@dynamic@decode:one_of(
commit_with_record_decoder(),
[commit_without_record_decoder()]
).
-file("src/goose.gleam", 191).
?DOC(" Decoder for commit events\n").
-spec commit_event_decoder() -> gleam@dynamic@decode:decoder(jetstream_event()).
commit_event_decoder() ->
gleam@dynamic@decode:field(
<<"did"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Did) ->
gleam@dynamic@decode:field(
<<"time_us"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_int/1},
fun(Time_us) ->
gleam@dynamic@decode:field(
<<"commit"/utf8>>,
commit_data_decoder(),
fun(Commit) ->
gleam@dynamic@decode:success(
{commit_event, Did, Time_us, Commit}
)
end
)
end
)
end
).
-file("src/goose.gleam", 250).
?DOC(" Decoder for identity data\n").
-spec identity_data_decoder() -> gleam@dynamic@decode:decoder(identity_data()).
identity_data_decoder() ->
gleam@dynamic@decode:field(
<<"did"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Did) ->
gleam@dynamic@decode:field(
<<"handle"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Handle) ->
gleam@dynamic@decode:field(
<<"seq"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_int/1},
fun(Seq) ->
gleam@dynamic@decode:field(
<<"time"/utf8>>,
{decoder,
fun gleam@dynamic@decode:decode_string/1},
fun(Time) ->
gleam@dynamic@decode:success(
{identity_data, Did, Handle, Seq, Time}
)
end
)
end
)
end
)
end
).
-file("src/goose.gleam", 242).
?DOC(" Decoder for identity events\n").
-spec identity_event_decoder() -> gleam@dynamic@decode:decoder(jetstream_event()).
identity_event_decoder() ->
gleam@dynamic@decode:field(
<<"did"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Did) ->
gleam@dynamic@decode:field(
<<"time_us"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_int/1},
fun(Time_us) ->
gleam@dynamic@decode:field(
<<"identity"/utf8>>,
identity_data_decoder(),
fun(Identity) ->
gleam@dynamic@decode:success(
{identity_event, Did, Time_us, Identity}
)
end
)
end
)
end
).
-file("src/goose.gleam", 267).
?DOC(" Decoder for account data\n").
-spec account_data_decoder() -> gleam@dynamic@decode:decoder(account_data()).
account_data_decoder() ->
gleam@dynamic@decode:field(
<<"active"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_bool/1},
fun(Active) ->
gleam@dynamic@decode:field(
<<"did"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Did) ->
gleam@dynamic@decode:field(
<<"seq"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_int/1},
fun(Seq) ->
gleam@dynamic@decode:field(
<<"time"/utf8>>,
{decoder,
fun gleam@dynamic@decode:decode_string/1},
fun(Time) ->
gleam@dynamic@decode:success(
{account_data, Active, Did, Seq, Time}
)
end
)
end
)
end
)
end
).
-file("src/goose.gleam", 259).
?DOC(" Decoder for account events\n").
-spec account_event_decoder() -> gleam@dynamic@decode:decoder(jetstream_event()).
account_event_decoder() ->
gleam@dynamic@decode:field(
<<"did"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_string/1},
fun(Did) ->
gleam@dynamic@decode:field(
<<"time_us"/utf8>>,
{decoder, fun gleam@dynamic@decode:decode_int/1},
fun(Time_us) ->
gleam@dynamic@decode:field(
<<"account"/utf8>>,
account_data_decoder(),
fun(Account) ->
gleam@dynamic@decode:success(
{account_event, Did, Time_us, Account}
)
end
)
end
)
end
).
-file("src/goose.gleam", 170).
?DOC(" Parse a JSON event string into a JetstreamEvent\n").
-spec parse_event(binary()) -> jetstream_event().
parse_event(Json_string) ->
case gleam@json:parse(Json_string, commit_event_decoder()) of
{ok, Event} ->
Event;
{error, _} ->
case gleam@json:parse(Json_string, identity_event_decoder()) of
{ok, Event@1} ->
Event@1;
{error, _} ->
case gleam@json:parse(Json_string, account_event_decoder()) of
{ok, Event@2} ->
Event@2;
{error, _} ->
{unknown_event, Json_string}
end
end
end.