-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 -> <<<>/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.