|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| -module(seb_pnp_bridge).
|
| -behaviour(gen_server).
|
|
|
| -export([start_link/1]).
|
| -export([init/1, handle_call/3, handle_cast/2, handle_info/2,
|
| terminate/2, code_change/3]).
|
|
|
| -define(SERVER, ?MODULE).
|
| -define(POLL_MS, 10000).
|
| -define(ATTACK_THRESHOLD, -1.0).
|
|
|
| -record(state, {
|
| conv_log :: string(),
|
| last_pos :: non_neg_integer(),
|
| kernel :: reference(),
|
| universe_sum :: float()
|
| }).
|
|
|
|
|
|
|
|
|
|
|
| start_link(ConvLogPath) ->
|
| gen_server:start_link({local, ?SERVER}, ?MODULE, [ConvLogPath], []).
|
|
|
|
|
|
|
|
|
|
|
| init([ConvLogPath]) ->
|
| {ok, Handle} = seb_kernel_nif:init_kernel(4, 0),
|
| erlang:send_after(?POLL_MS, self(), poll),
|
| {ok, #state{
|
| conv_log = ConvLogPath,
|
| last_pos = 0,
|
| kernel = Handle,
|
| universe_sum = 0.0
|
| }}.
|
|
|
| handle_info(poll, State) ->
|
| NewState = poll_and_seal(State),
|
| erlang:send_after(?POLL_MS, self(), poll),
|
| {noreply, NewState};
|
|
|
| handle_info(_Info, State) ->
|
| {noreply, State}.
|
|
|
| handle_call(_Req, _From, State) ->
|
| {reply, ok, State}.
|
|
|
| handle_cast(_Msg, State) ->
|
| {noreply, State}.
|
|
|
| terminate(_Reason, _State) -> ok.
|
| code_change(_OldVsn, State, _Extra) -> {ok, State}.
|
|
|
|
|
|
|
|
|
|
|
| poll_and_seal(#state{conv_log = Path, last_pos = Pos,
|
| kernel = Handle, universe_sum = Sum} = State) ->
|
| case file:open(Path, [read, binary]) of
|
| {error, _} ->
|
| State;
|
| {ok, Fd} ->
|
| {ok, _} = file:position(Fd, Pos),
|
| {Lines, NewPos} = read_lines(Fd, Pos),
|
| file:close(Fd),
|
| process_entries(Lines, State#state{last_pos = NewPos})
|
| end.
|
|
|
| read_lines(Fd, Pos) ->
|
| read_lines(Fd, Pos, []).
|
|
|
| read_lines(Fd, Pos, Acc) ->
|
| case file:read_line(Fd) of
|
| eof -> {lists:reverse(Acc), Pos};
|
| {error, _} -> {lists:reverse(Acc), Pos};
|
| {ok, Line} ->
|
| {ok, NewPos} = file:position(Fd, cur),
|
| read_lines(Fd, NewPos, [Line | Acc])
|
| end.
|
|
|
| process_entries([], State) -> State;
|
| process_entries([Line | Rest], State) ->
|
| case catch jiffy:decode(Line, [return_maps]) of
|
| {'EXIT', _} ->
|
| process_entries(Rest, State);
|
| Entry ->
|
| NewState = seal_entry(Entry, State),
|
| process_entries(Rest, NewState)
|
| end.
|
|
|
| seal_entry(Entry, #state{kernel = Handle, universe_sum = Sum} = State) ->
|
| Delta = maps:get(<<"universeSumDelta">>, Entry, 0.0),
|
| ProblemId = maps:get(<<"problemId">>, Entry, <<>>),
|
| Timestamp = maps:get(<<"timestamp">>, Entry, <<>>),
|
| NewSum = Sum + Delta,
|
|
|
|
|
| Payload = build_payload(Entry, Delta),
|
|
|
|
|
| EventType = if Delta < 0 -> 16#0401; true -> 16#0400 end,
|
| Header = build_header(EventType, byte_size(Payload)),
|
|
|
|
|
| {ok, PrevTip} = seb_kernel_nif:get_tip_hash(Handle),
|
|
|
| Footer = <<PrevTip/binary, (binary:copy(<<0>>, 32))/binary, (binary:copy(<<0>>, 64))/binary>>,
|
|
|
| case seb_kernel_nif:append_event(Handle, Header, Payload, Footer) of
|
| {ok, Offset} ->
|
| if Delta < 0 ->
|
| error_logger:warning_msg(
|
| "seb_pnp_bridge: ATTACK EVENT sealed at offset ~p, delta=~p, problem=~s~n",
|
| [Offset, Delta, ProblemId]),
|
| handle_attack(Handle, NewSum);
|
| true ->
|
| ok
|
| end;
|
| {error, Reason} ->
|
| error_logger:error_msg(
|
| "seb_pnp_bridge: failed to seal ~s: ~p~n", [ProblemId, Reason])
|
| end,
|
|
|
| State#state{universe_sum = NewSum}.
|
|
|
| handle_attack(Handle, Sum) ->
|
|
|
| case seb_kernel_nif:verify_chain(Handle) of
|
| {ok, Count} ->
|
| error_logger:warning_msg(
|
| "seb_pnp_bridge: chain intact (~p records), sum=~p~n", [Count, Sum]);
|
| {error, Reason} ->
|
| error_logger:error_msg(
|
| "seb_pnp_bridge: CHAIN INTEGRITY FAILURE: ~p — escalating~n", [Reason]),
|
| exit({chain_integrity_failure, Sum})
|
| end,
|
|
|
| if Sum < ?ATTACK_THRESHOLD ->
|
| error_logger:error_msg(
|
| "seb_pnp_bridge: universe_sum=~p < threshold ~p — HALT~n",
|
| [Sum, ?ATTACK_THRESHOLD]),
|
| exit({attack_threshold_exceeded, Sum});
|
| true -> ok
|
| end.
|
|
|
|
|
| build_header(EventType, PayloadSize) ->
|
| Timestamp = erlang:system_time(nanosecond),
|
| <<EventType:64/little-unsigned,
|
| Timestamp:64/little-unsigned,
|
| 0:64/little-unsigned,
|
| PayloadSize:32/little-unsigned,
|
| 0:32/little-unsigned,
|
| 0:64/little-unsigned,
|
| 0:64/little-unsigned,
|
| 0:64/little-unsigned,
|
| 0:32/little-unsigned,
|
| 0:32/little-unsigned>>.
|
|
|
|
|
| build_payload(Entry, Delta) ->
|
| EventType = if Delta < 0 -> 16#0401; true -> 16#0400 end,
|
| Timestamp = erlang:system_time(nanosecond),
|
| ProblemId = maps:get(<<"problemId">>, Entry, <<>>),
|
| PidHash = crypto:hash(sha256, ProblemId),
|
| DeltaBin = <<Delta:64/little-float>>,
|
| Reserved = binary:copy(<<0>>, 8),
|
| <<EventType:64/little-unsigned,
|
| Timestamp:64/little-unsigned,
|
| PidHash/binary,
|
| DeltaBin/binary,
|
| Reserved/binary>>.
|
|
|