defmodule Electric.Client.Poll do @moduledoc """ Poll-based API for fetching shape changes. This module provides explicit request-response semantics for fetching changes from Electric, as an alternative to the streaming API. ## Usage # Create initial state state = ShapeState.new() # Make a polling request case Poll.request(client, state) do {:ok, messages, new_state} -> # Process messages, use new_state for next poll ... {:must_refetch, messages, new_state} -> # Shape was reset, clear local state and process messages ... {:error, error} -> # Handle error ... end ## Behavior - First request (when `up_to_date?: false`): Makes a non-live request to get initial snapshot - Subsequent requests (when `up_to_date?: true`): Makes a live request that long-polls until changes arrive - Handles synthetic deletes from move-out events - Returns updated state for the next request """ alias Electric.Client alias Electric.Client.ExpiredShapesCache alias Electric.Client.Fetch alias Electric.Client.Message alias Electric.Client.ShapeKey alias Electric.Client.ShapeState alias Electric.Client.TagTracker @max_stale_retries 3 @type poll_result :: {:ok, [Client.message()], ShapeState.t()} | {:must_refetch, [Client.message()], ShapeState.t()} | {:error, Client.Error.t()} @doc """ Make a single polling request to fetch shape changes. ## Arguments * `client` - The Electric client * `shape` - The shape definition (or a client pre-configured for a shape) * `state` - The current polling state (use `ShapeState.new()` for initial request) * `opts` - Options: * `:replica` - `:default` or `:full` (default: `:default`) ## Returns * `{:ok, messages, new_state}` - Success, messages received * `{:must_refetch, messages, new_state}` - Shape was reset (409), state has been cleared * `{:error, error}` - Error occurred ## Examples state = ShapeState.new() {:ok, messages, state} = Poll.request(client, state, replica: :full) # Process messages... # Poll again for more changes {:ok, messages, state} = Poll.request(client, state, replica: :full) """ @spec request(Client.t(), ShapeState.t(), keyword()) :: poll_result() def request(%Client{} = client, %ShapeState{} = state, opts \\ []) do replica = Keyword.get(opts, :replica, :default) shape_key = ShapeKey.canonical(client.endpoint, client.params) request = build_request(client, state, replica, shape_key) case Fetch.request(client, request) do %Fetch.Response{status: status} = resp when status in 200..299 -> validate_headers!(resp, state) handle_success(resp, client, state, shape_key) {:error, %Fetch.Response{status: 409} = resp} -> handle_must_refetch(resp, client, state, shape_key) {:error, %Fetch.Response{body: body} = resp} -> {:error, %Client.Error{message: unwrap_error(body), resp: resp}} {:error, error} -> {:error, %Client.Error{message: "Unable to retrieve data", resp: error}} end end defp build_request(client, state, replica, shape_key) do %{ shape_handle: shape_handle, offset: offset, up_to_date?: up_to_date?, next_cursor: cursor, stale_cache_buster: cache_buster } = state # Build additional params for cache busting cache_busting_params = %{} |> maybe_add_expired_handle(shape_key) |> maybe_add_cache_buster(cache_buster) # Merge cache busting params into client params before building request client_with_cache_params = Client.merge_params(client, cache_busting_params) Client.request(client_with_cache_params, offset: offset, shape_handle: shape_handle, replica: replica, live: up_to_date?, next_cursor: cursor ) end defp maybe_add_expired_handle(params, shape_key) do case ExpiredShapesCache.get_expired_handle(shape_key) do nil -> params expired -> Map.put(params, "expired_handle", expired) end end defp maybe_add_cache_buster(params, nil), do: params defp maybe_add_cache_buster(params, buster), do: Map.put(params, "cache-buster", buster) defp handle_success(resp, client, state, shape_key) do response_handle = resp.shape_handle expired_handle = ExpiredShapesCache.get_expired_handle(shape_key) # Check for stale CDN response — always enter stale-retry to add a cache # buster. Without this, the CDN keeps serving the same stale response and # the client loops infinitely (the URL never changes). cond do response_handle == expired_handle -> handle_stale_response(state, shape_key) # Normal: process response true -> process_success_response(resp, client, state, response_handle) end end defp process_success_response(resp, client, state, shape_handle) do final_offset = last_offset(resp, state.offset) next_cursor = resp.next_cursor state = %{state | shape_handle: shape_handle, next_cursor: next_cursor, offset: final_offset} state = handle_schema(resp, client, state) state = ShapeState.clear_stale_retry(state) %{value_mapper_fun: value_mapper_fun} = state {messages, new_state} = resp.body |> ensure_enum() |> Enum.flat_map(&Message.parse(&1, shape_handle, value_mapper_fun, resp.request_timestamp)) |> process_messages(state) {:ok, messages, new_state} end defp handle_stale_response(state, shape_key) do cond do state.stale_cache_retry_count < @max_stale_retries -> {:stale_retry, ShapeState.enter_stale_retry(state)} state.self_heal_attempted? -> {:error, %Client.Error{ message: "CDN continues serving stale cached responses after #{@max_stale_retries} " <> "retry attempts and one self-heal attempt" }} true -> # Self-heal: clear the expired entry from local cache so the next # request omits the `expired_handle` param. Since the server never # reuses handles (SPEC.md S0), the next response should bypass stale # detection. If it doesn't (broken CDN), we error on the next pass via # the `self_heal_attempted?` branch above. ExpiredShapesCache.clear_handle(shape_key) new_state = ShapeState.enter_stale_retry(%{ state | self_heal_attempted?: true, stale_cache_retry_count: 0 }) {:stale_retry, new_state} end end defp handle_must_refetch(resp, client, state, shape_key) do # Mark the old handle as expired if state.shape_handle do ExpiredShapesCache.mark_expired(shape_key, state.shape_handle) end handle = shape_handle(resp) || "#{String.trim_trailing(state.shape_handle || "", "-next")}-next" new_state = ShapeState.reset(state, handle) new_state = handle_schema(resp, client, new_state) new_state = ShapeState.clear_stale_retry(new_state) # Add a cache-buster on every 409 so that the next request URL cannot match # a URL the CDN has cached. Without this, a CDN that strips the # `expired_handle` query param from its cache key keeps serving the same # cached 409 indefinitely. new_state = %{new_state | stale_cache_buster: ShapeState.generate_cache_buster()} # Always emit a synthetic must-refetch control message rather than # forwarding whatever the server (or a misbehaving proxy) put in the 409 # body. Subscribers must receive the signal to clear local state on every # 409, even when the body is empty or stripped of the control message. # Any data rows present in the body refer to the old, expired handle, so # discarding them is correct. messages = [ %Message.ControlMessage{ control: :must_refetch, handle: handle, request_timestamp: resp.request_timestamp } ] {:must_refetch, messages, new_state} end defp process_messages(messages, state) do {processed_messages, new_state} = Enum.reduce(messages, {[], state}, fn msg, {msgs_acc, state_acc} -> case handle_message(msg, state_acc) do {:message, processed_msg, new_state} -> {[processed_msg | msgs_acc], new_state} {:messages, processed_msgs, new_state} -> {Enum.reverse(processed_msgs) ++ msgs_acc, new_state} {:skip, new_state} -> {msgs_acc, new_state} end end) {Enum.reverse(processed_messages), new_state} end defp handle_message(%Message.ControlMessage{control: :up_to_date} = msg, state) do {:message, msg, %{state | up_to_date?: true}} end defp handle_message(%Message.ControlMessage{control: :snapshot_end}, state) do {:skip, state} end defp handle_message(%Message.ChangeMessage{} = msg, state) do {tag_to_keys, key_data, disjunct_positions} = TagTracker.update_tag_index( state.tag_to_keys, state.key_data, state.disjunct_positions, msg ) {:message, msg, %{ state | tag_to_keys: tag_to_keys, key_data: key_data, disjunct_positions: disjunct_positions }} end defp handle_message( %Message.MoveOutMessage{patterns: patterns, request_timestamp: request_timestamp}, state ) do {synthetic_deletes, tag_to_keys, key_data} = TagTracker.generate_synthetic_deletes( state.tag_to_keys, state.key_data, state.disjunct_positions, patterns, request_timestamp ) {:messages, synthetic_deletes, %{state | tag_to_keys: tag_to_keys, key_data: key_data}} end defp handle_message( %Message.MoveInMessage{patterns: patterns}, state ) do {tag_to_keys, key_data} = TagTracker.handle_move_in( state.tag_to_keys, state.key_data, patterns ) {:skip, %{state | tag_to_keys: tag_to_keys, key_data: key_data}} end defp handle_schema(%Fetch.Response{schema: schema}, client, %{value_mapper_fun: nil} = state) when is_map(schema) do {parser_module, parser_opts} = client.parser value_mapper_fun = parser_module.for_schema(schema, parser_opts) %{state | schema: schema, value_mapper_fun: value_mapper_fun} end defp handle_schema(_resp, _client, state) do state end defp ensure_enum(body) do case Enumerable.impl_for(body) do nil -> List.wrap(body) Enumerable.Map -> List.wrap(body) _impl -> body end end defp validate_headers!(%Fetch.Response{} = resp, %ShapeState{} = state) do # Validate required Electric response fields, matching the TypeScript # client's createFetchWithResponseHeadersCheck middleware. # # We check parsed struct fields (not raw HTTP headers) so this works # for both real HTTP responses and the Mock fetch implementation. # # Rules (based on TypeScript's createFetchWithResponseHeadersCheck): # - All responses: electric-handle, electric-offset # - Non-live responses: electric-schema (server always sends it, but we # only validate when we don't have one yet since it's first-write-wins) # - Live responses: electric-cursor (CDN cache buster) is_live? = state.up_to_date? missing = [] missing = if is_nil(resp.shape_handle), do: ["electric-handle" | missing], else: missing missing = if is_nil(resp.last_offset), do: ["electric-offset" | missing], else: missing missing = if not is_live? and is_nil(state.schema) and is_nil(resp.schema), do: ["electric-schema" | missing], else: missing missing = if is_live? and is_nil(resp.next_cursor), do: ["electric-cursor" | missing], else: missing if missing != [] do raise Client.Error, message: "Response is missing required Electric header(s): #{Enum.join(missing, ", ")}. " <> "This usually indicates a proxy or CDN misconfiguration — " <> "check that your proxy forwards all Electric headers " <> "(electric-handle, electric-offset, electric-schema, electric-cursor) to the client." end end defp shape_handle(%Fetch.Response{shape_handle: shape_handle}) do shape_handle end defp last_offset(%Fetch.Response{last_offset: nil}, offset), do: offset defp last_offset(%Fetch.Response{last_offset: offset}, _offset), do: offset defp unwrap_error([]), do: "Unknown error" defp unwrap_error([msg]), do: msg defp unwrap_error([_ | _] = msgs), do: msgs defp unwrap_error(msg), do: msg end