defmodule S2.S2S.AppendSession do @moduledoc """ Bidirectional streaming append session. Opens a streaming POST to `/v1/streams/{stream}/records` with `Content-Type: s2s/proto`. Supports multiple sequential `append/2` calls on the same session, each sending a framed `AppendInput` and receiving a framed `AppendAck`. ## Process affinity Sessions are NOT safe to share across processes. The underlying Mint connection delivers TCP messages to the owning process's mailbox. Creating a session in one process and calling `append/2` from another will not work — the receiving process won't see the TCP data. """ require Logger alias S2.S2S.Shared @typedoc "An open append session." @type t :: %__MODULE__{} defstruct [ :conn, :request_ref, :basin, :stream, :owner_pid, recv_timeout: 5_000, compression: :none, closed: false, data: <<>> ] @doc """ Open a new streaming append session. ## Options * `:token` — Bearer token for authentication. * `:recv_timeout` — Timeout in milliseconds for receiving responses (default: 5000). * `:compression` — Compression for S2S frames: `:none`, `:gzip`, or `:zstd` (default: `:none`). Returns `{:ok, session}` on success or `{:error, reason}` on failure. On `Mint.HTTP2.request/5` failure, returns `{:error, reason, conn}` so the caller can still manage the connection. """ @spec open(Mint.HTTP2.t(), String.t(), String.t(), keyword()) :: {:ok, t()} | {:error, term()} | {:error, term(), Mint.HTTP2.t()} def open(conn, basin, stream, opts \\ []) do Logger.debug("S2S.AppendSession.open basin=#{basin} stream=#{stream}") path = Shared.records_path(stream) token = Keyword.get(opts, :token) recv_timeout = Keyword.get(opts, :recv_timeout, Shared.default_timeout()) compression = Keyword.get(opts, :compression, :none) headers = Shared.build_headers(basin, token) case Mint.HTTP2.request(conn, "POST", path, headers, :stream) do {:ok, conn, request_ref} -> session = %__MODULE__{ conn: conn, request_ref: request_ref, basin: basin, stream: stream, owner_pid: self(), recv_timeout: recv_timeout, compression: compression } wait_for_headers(session) {:error, conn, reason} -> {:error, reason, conn} end end @doc """ Append a batch of records to the session. Returns `{:ok, ack, session}` on success or `{:error, reason, session}` on failure. The session is marked as closed on error and cannot be reused. """ @spec append(t(), S2.V1.AppendInput.t()) :: {:ok, S2.V1.AppendAck.t(), t()} | {:error, term(), t()} def append(%__MODULE__{closed: true} = session, _input) do {:error, :session_closed, session} end def append(%__MODULE__{} = session, %S2.V1.AppendInput{} = input) do check_owner!(session) case Shared.encode_framed(input, compression: session.compression) do {:error, reason} -> {:error, reason, Shared.close_session(session)} {:ok, body} -> case Mint.HTTP2.stream_request_body(session.conn, session.request_ref, body) do {:ok, conn} -> session = %{session | conn: conn} receive_ack(session) {:error, conn, reason} -> {:error, reason, Shared.close_session(session, conn)} end end end @doc """ Close the append session gracefully by sending EOF on the request body. Returns `{:ok, session}` or `{:error, reason, session}`. """ @spec close(t()) :: {:ok, t()} | {:error, term(), t()} def close(%__MODULE__{closed: true} = session), do: {:ok, session} def close(%__MODULE__{} = session) do case Mint.HTTP2.stream_request_body(session.conn, session.request_ref, :eof) do {:ok, conn} -> session = %{session | conn: conn, closed: true} drain_final_response(session) {:error, conn, reason} -> {:error, reason, Shared.close_session(session, conn)} end end defp wait_for_headers(session) do case Shared.wait_for_headers(session.conn, session.request_ref, session.recv_timeout) do {:ok, 200, _data, conn} -> {:ok, %{session | conn: conn}} {:ok, status, _data, conn} -> {:error, %S2.Error{status: status, message: "unexpected status #{status}"}, conn} {:error, reason, conn} -> {:error, reason, conn} end end defp receive_ack(session) do all_data = session.data case Shared.decode_frame(all_data, S2.V1.AppendAck) do {:ok, ack, rest} -> {:ok, ack, %{session | data: rest}} {:error, reason} -> {:error, reason, Shared.close_session(session)} :incomplete -> dl = Shared.deadline(session.recv_timeout) do_receive_ack(%{session | data: <<>>}, all_data, dl) end end @doc false def handle_ack_response({:ok, conn, responses}, session, acc) do session = %{session | conn: conn} new_data = Shared.extract_data(responses, session.request_ref) all_data = acc <> new_data case Shared.check_buffer_size(all_data) do {:error, :buffer_overflow} -> {:error, :buffer_overflow, Shared.close_session(%{session | data: all_data})} :ok -> done? = Shared.done?(responses, session.request_ref) case Shared.decode_frame(all_data, S2.V1.AppendAck) do {:ok, ack, rest} -> {:ok, ack, %{session | data: rest}} {:error, reason} -> {:error, reason, Shared.close_session(session)} :incomplete when done? -> {:error, :incomplete_frame, Shared.close_session(session)} :incomplete -> {:continue, session, all_data} end end end def handle_ack_response({:error, conn, _error, _responses}, session, _acc) do {:error, :stream_error, Shared.close_session(session, conn)} end def handle_ack_response(:unknown, session, acc) do {:continue, session, acc} end defp do_receive_ack(session, acc, dl) do receive do message -> case handle_ack_response(Mint.HTTP2.stream(session.conn, message), session, acc) do {:continue, session, acc} -> do_receive_ack(session, acc, dl) result -> result end after Shared.remaining(dl) -> {:error, :timeout, Shared.close_session(session)} end end # After sending EOF, drain the server's final response (status/headers/done) # so the HTTP/2 stream is fully closed. Best-effort: if it times out or errors, # we still return the closed session. defp drain_final_response(session) do drain_final_response(session, Shared.deadline(session.recv_timeout)) end @doc false def handle_drain_response({:ok, conn, responses}, session) do session = %{session | conn: conn} if Shared.done?(responses, session.request_ref) do {:ok, session} else {:continue, session} end end def handle_drain_response({:error, conn, _error, _responses}, session) do {:ok, %{session | conn: conn}} end def handle_drain_response(:unknown, session) do {:continue, session} end defp drain_final_response(session, dl) do receive do message -> case handle_drain_response(Mint.HTTP2.stream(session.conn, message), session) do {:continue, session} -> drain_final_response(session, dl) result -> result end after Shared.remaining(dl) -> {:ok, session} end end defp check_owner!(%__MODULE__{owner_pid: pid}) do Shared.assert_owner!(pid, "AppendSession") end end