diff --git a/lib/plausible/clickhouse_event_v2.ex b/lib/plausible/clickhouse_event_v2.ex index 7e2d8861a22f..ee90a04b1c10 100644 --- a/lib/plausible/clickhouse_event_v2.ex +++ b/lib/plausible/clickhouse_event_v2.ex @@ -3,6 +3,7 @@ defmodule Plausible.ClickhouseEventV2 do Event schema for when NumericIDs migration is complete """ use Ecto.Schema + use Plausible import Ecto.Changeset @primary_key false @@ -48,6 +49,11 @@ defmodule Plausible.ClickhouseEventV2 do field :acquisition_channel, Ch, type: "LowCardinality(String)", writable: :never + on_ee do + # Field used during event replay + field :replay_session_id, Ch, type: "UInt64" + end + # Virtual field used during event processing field :interactive?, :boolean, default: true, virtual: true, writable: :never end @@ -72,7 +78,7 @@ defmodule Plausible.ClickhouseEventV2 do :revenue_reporting_amount, :revenue_reporting_currency, :interactive? - ] + ] ++ on_ee(do: [:replay_session_id], else: []) ) |> validate_required([:name, :site_id, :hostname, :pathname, :user_id, :timestamp]) end diff --git a/lib/plausible/clickhouse_session_v2.ex b/lib/plausible/clickhouse_session_v2.ex index e9614df5b9bb..379247dce362 100644 --- a/lib/plausible/clickhouse_session_v2.ex +++ b/lib/plausible/clickhouse_session_v2.ex @@ -3,6 +3,7 @@ defmodule Plausible.ClickhouseSessionV2 do Session schema for when NumericIDs migration is complete """ use Ecto.Schema + use Plausible defmodule BoolUInt8 do @moduledoc """ @@ -75,6 +76,11 @@ defmodule Plausible.ClickhouseSessionV2 do field :transferred_from, :string field :acquisition_channel, Ch, type: "LowCardinality(String)", writable: :never + + on_ee do + # Field used during event replay + field :replay_session_id, Ch, type: "UInt64" + end end def random_uint64() do diff --git a/lib/plausible/ingestion/event.ex b/lib/plausible/ingestion/event.ex index ee4d99e1f8c8..ff140b3f93b1 100644 --- a/lib/plausible/ingestion/event.ex +++ b/lib/plausible/ingestion/event.ex @@ -54,13 +54,20 @@ defmodule Plausible.Ingestion.Event do @spec build_and_buffer(Request.t(), Keyword.t()) :: {:ok, %{buffered: [t()], dropped: [t()]}} def build_and_buffer(%Request{domains: domains} = request, context \\ []) do + skip_rate_limit? = + on_ee do + not is_nil(request.replay_session_id) + else + false + end + processed_events = if spam_referrer?(request) do for domain <- domains, do: drop(new(domain, request), :spam_referrer) else Enum.reduce(domains, [], fn domain, acc -> # credo:disable-for-next-line Credo.Check.Refactor.Nesting - case GateKeeper.check(domain) do + case GateKeeper.check(domain, skip_rate_limit?: skip_rate_limit?) do {:allow, site} -> processed = domain @@ -276,7 +283,7 @@ defmodule Plausible.Ingestion.Event do end defp put_basic_info(%__MODULE__{} = event, _context) do - update_event_attrs(event, %{ + attrs = %{ domain: event.domain, site_id: event.site.id, timestamp: event.request.timestamp, @@ -286,7 +293,13 @@ defmodule Plausible.Ingestion.Event do scroll_depth: event.request.scroll_depth, engagement_time: event.request.engagement_time, interactive?: event.request.interactive? - }) + } + + on_ee do + attrs = Map.put(attrs, :replay_session_id, event.request.replay_session_id) + end + + update_event_attrs(event, attrs) end defp put_source_info(%__MODULE__{} = event, _context) do @@ -562,7 +575,17 @@ defmodule Plausible.Ingestion.Event do user_agent = request.user_agent || "" root_domain = get_root_domain(hostname) - SipHash.hash!(salt, user_agent <> request.remote_ip <> domain <> root_domain) + replay_session_id = + on_ee do + to_string(request.replay_session_id) + else + "" + end + + SipHash.hash!( + salt, + user_agent <> request.remote_ip <> domain <> root_domain <> replay_session_id + ) end end diff --git a/lib/plausible/ingestion/request.ex b/lib/plausible/ingestion/request.ex index e8b15a97f152..e0a758fe5ab5 100644 --- a/lib/plausible/ingestion/request.ex +++ b/lib/plausible/ingestion/request.ex @@ -59,6 +59,9 @@ defmodule Plausible.Ingestion.Request do on_ee do field :revenue_source, :map + + # field for replayed events + field :replay_session_id, :integer end field :query_params, :map @@ -91,6 +94,7 @@ defmodule Plausible.Ingestion.Request do |> put_uri(request_body) |> put_hostname() |> put_user_agent(conn) + |> put_replay_data(conn) |> put_request_params(request_body) |> put_referrer(request_body) |> put_pathname() @@ -131,6 +135,44 @@ defmodule Plausible.Ingestion.Request do defp put_revenue_source(changeset, _request_body), do: changeset end + on_ee do + @replay_session_id_header "x-replay-session-id" + @replay_time_header "x-replay-time" + + defp put_replay_data(changeset, conn) do + now = NaiveDateTime.utc_now(:second) + + replay_session_id = + conn + |> Plug.Conn.get_req_header(@replay_session_id_header) + |> List.first() + |> to_integer() + + if replay_session_id do + time = + conn + |> Plug.Conn.get_req_header(@replay_time_header) + |> List.first() + |> NaiveDateTime.from_iso8601!() + + if NaiveDateTime.compare(time, now) in [:lt, :eq] do + changeset + |> Changeset.put_change(:replay_session_id, replay_session_id) + |> Changeset.put_change(:timestamp, time) + else + changeset + end + else + changeset + end + end + + defp to_integer(s) when is_binary(s), do: String.to_integer(s) + defp to_integer(_), do: nil + else + defp put_replay_data(changeset, _conn), do: changeset + end + defp put_remote_ip(changeset, conn) do Changeset.put_change(changeset, :remote_ip, PlausibleWeb.RemoteIP.get(conn)) end diff --git a/lib/plausible/session/cache_store.ex b/lib/plausible/session/cache_store.ex index 05b3998e06ce..a40814d4e7a0 100644 --- a/lib/plausible/session/cache_store.ex +++ b/lib/plausible/session/cache_store.ex @@ -3,6 +3,8 @@ defmodule Plausible.Session.CacheStore do Session management on the basis of incoming events. """ + use Plausible + alias Plausible.Session.WriteBuffer @lock_timeout 1000 @@ -70,7 +72,18 @@ defmodule Plausible.Session.CacheStore do defp find_session(_domain, nil), do: nil defp find_session(event, user_id) do - from_cache = Plausible.Cache.Adapter.get(:sessions, {event.site_id, user_id}) + key = + on_ee do + if event.replay_session_id do + {event.site_id, user_id, event.replay_session_id} + else + {event.site_id, user_id} + end + else + {event.site_id, user_id} + end + + from_cache = Plausible.Cache.Adapter.get(:sessions, key) case from_cache do nil -> @@ -84,7 +97,17 @@ defmodule Plausible.Session.CacheStore do end defp update_session_cache(session) do - key = {session.site_id, session.user_id} + key = + on_ee do + if session.replay_session_id do + {session.site_id, session.user_id, session.replay_session_id} + else + {session.site_id, session.user_id} + end + else + {session.site_id, session.user_id} + end + Plausible.Cache.Adapter.put(:sessions, key, session, dirty?: true) session end @@ -126,7 +149,7 @@ defmodule Plausible.Session.CacheStore do end defp new_session_from_event(event, session_attributes) do - %Plausible.ClickhouseSessionV2{ + new_session = %Plausible.ClickhouseSessionV2{ sign: 1, session_id: Plausible.ClickhouseSessionV2.random_uint64(), hostname: if(event.name == "pageview", do: event.hostname, else: ""), @@ -161,5 +184,11 @@ defmodule Plausible.Session.CacheStore do "entry_meta.key": Map.get(event, :"meta.key"), "entry_meta.value": Map.get(event, :"meta.value") } + + on_ee do + %{new_session | replay_session_id: event.replay_session_id} + else + new_session + end end end diff --git a/lib/plausible/site/gate_keeper.ex b/lib/plausible/site/gate_keeper.ex index 67f79d2790bc..0cac4c1b2706 100644 --- a/lib/plausible/site/gate_keeper.ex +++ b/lib/plausible/site/gate_keeper.ex @@ -44,11 +44,16 @@ defmodule Plausible.Site.GateKeeper do with %Site{team: %{accept_traffic_until: accept_traffic_until}} = site <- Cache.get(domain, Keyword.get(opts, :cache_opts, [])), true <- Plausible.Sites.regular?(site) do - if not is_nil(accept_traffic_until) and - Date.after?(Date.utc_today(), accept_traffic_until) do - :payment_required - else - check_rate_limit(site, opts) + cond do + not is_nil(accept_traffic_until) and + Date.after?(Date.utc_today(), accept_traffic_until) -> + :payment_required + + opts[:skip_rate_limit?] -> + {:allow, site} + + true -> + check_rate_limit(site, opts) end else _ -> diff --git a/test/plausible/ingestion/event_test.exs b/test/plausible/ingestion/event_test.exs index 43a0a06e2d53..93a70d8fc413 100644 --- a/test/plausible/ingestion/event_test.exs +++ b/test/plausible/ingestion/event_test.exs @@ -475,6 +475,108 @@ defmodule Plausible.Ingestion.EventTest do assert Decimal.eq?(event.clickhouse_event.revenue_source_amount, Decimal.new("10.2")) end + @tag :ee_only + test "saves replay session id when passed in headers" do + site = new_site() + + payload = %{ + name: "pageview", + url: "http://#{site.domain}" + } + + conn = + build_conn(:post, "/api/events", payload) + |> Plug.Conn.put_req_header("x-replay-event-id", "123") + |> Plug.Conn.put_req_header("x-replay-session-id", "456") + |> Plug.Conn.put_req_header("x-replay-time", "2026-06-02 12:43:00") + + assert {:ok, request, _conn} = Request.build(conn) + + assert {:ok, %{buffered: [event], dropped: []}} = Event.build_and_buffer(request) + assert event.clickhouse_event.replay_session_id == 456 + assert event.clickhouse_event.timestamp == ~N[2026-06-02 12:43:00] + end + + @tag :ee_only + test "leaves replay session id empty when not passed in headers" do + site = new_site() + + payload = %{ + name: "pageview", + url: "http://#{site.domain}" + } + + conn = + build_conn(:post, "/api/events", payload) + + assert {:ok, request, _conn} = Request.build(conn) + + assert {:ok, %{buffered: [event], dropped: []}} = Event.build_and_buffer(request) + assert is_nil(event.clickhouse_event.replay_session_id) + assert %NaiveDateTime{} = event.clickhouse_event.timestamp + end + + @tag :ee_only + test "replayed events have a distinct user id due to reply session-dependent salt" do + site = new_site() + + payload = %{ + name: "pageview", + url: "http://#{site.domain}" + } + + conn1 = + build_conn(:post, "/api/events", payload) + |> Plug.Conn.put_req_header("user-agent", "Mozilla") + |> Plug.Conn.put_req_header("x-plausible-ip", "1.2.3.4") + + assert {:ok, request1, _conn} = Request.build(conn1) + assert {:ok, %{buffered: [event1], dropped: []}} = Event.build_and_buffer(request1) + + conn2 = + build_conn(:post, "/api/events", payload) + |> Plug.Conn.put_req_header("user-agent", "Mozilla") + |> Plug.Conn.put_req_header("x-plausible-ip", "1.2.3.4") + |> Plug.Conn.put_req_header("x-replay-event-id", "123") + |> Plug.Conn.put_req_header("x-replay-session-id", "456") + |> Plug.Conn.put_req_header("x-replay-time", "2026-06-02 12:43:00") + + assert {:ok, request2, _conn} = Request.build(conn2) + assert {:ok, %{buffered: [event2], dropped: []}} = Event.build_and_buffer(request2) + + conn3 = + build_conn(:post, "/api/events", payload) + |> Plug.Conn.put_req_header("user-agent", "Mozilla") + |> Plug.Conn.put_req_header("x-plausible-ip", "1.2.3.4") + |> Plug.Conn.put_req_header("x-replay-event-id", "124") + |> Plug.Conn.put_req_header("x-replay-session-id", "456") + |> Plug.Conn.put_req_header("x-replay-time", "2026-06-02 12:44:01") + + assert {:ok, request3, _conn} = Request.build(conn3) + assert {:ok, %{buffered: [event3], dropped: []}} = Event.build_and_buffer(request3) + + conn4 = + build_conn(:post, "/api/events", payload) + |> Plug.Conn.put_req_header("user-agent", "Mozilla") + |> Plug.Conn.put_req_header("x-plausible-ip", "1.2.3.4") + |> Plug.Conn.put_req_header("x-replay-event-id", "125") + |> Plug.Conn.put_req_header("x-replay-session-id", "789") + |> Plug.Conn.put_req_header("x-replay-time", "2026-06-02 12:44:01") + + assert {:ok, request4, _conn} = Request.build(conn4) + assert {:ok, %{buffered: [event4], dropped: []}} = Event.build_and_buffer(request4) + + # non-replayed event user id different from replayed event user despite + # identical fingerprint + assert event1.clickhouse_event.user_id != event2.clickhouse_event.user_id + # replayed events with the same session for the same fingerprint have the + # same user id + assert event2.clickhouse_event.user_id == event3.clickhouse_event.user_id + # replayed event from different sessions for the same fingerprint + # have different users + assert event3.clickhouse_event.user_id != event4.clickhouse_event.user_id + end + test "does not save revenue amount when there is no revenue goal" do site = new_site() diff --git a/test/plausible/ingestion/request_test.exs b/test/plausible/ingestion/request_test.exs index 7426926a2546..5b7d58f103cd 100644 --- a/test/plausible/ingestion/request_test.exs +++ b/test/plausible/ingestion/request_test.exs @@ -173,6 +173,62 @@ defmodule Plausible.Ingestion.RequestTest do assert request.tracker_script_version == 137 end + @tag :ee_only + test "parses replay headers if present" do + payload = %{ + name: "pageview", + domain: "dummy.site", + url: "http://dummy.site/index.html" + } + + conn = + build_conn(:post, "/api/events", payload) + |> put_req_header("x-replay-event-id", "123") + |> put_req_header("x-replay-session-id", "456") + |> put_req_header("x-replay-time", "2026-06-02 12:43:00") + + assert {:ok, request, _conn} = Request.build(conn) + assert request.replay_session_id == 456 + assert NaiveDateTime.compare(request.timestamp, ~N[2026-06-02 12:43:00]) == :eq + end + + @tag :ee_only + test "leaves replay fields empty and timestamp intact if headers not present" do + payload = %{ + name: "pageview", + domain: "dummy.site", + url: "http://dummy.site/index.html" + } + + conn = build_conn(:post, "/api/events", payload) + + assert {:ok, request, _conn} = Request.build(conn) + assert is_nil(request.replay_session_id) + assert %NaiveDateTime{} = request.timestamp + end + + @tag :ee_only + test "leaves replay fields empty and timestamp intact if replay time in the future" do + payload = %{ + name: "pageview", + domain: "dummy.site", + url: "http://dummy.site/index.html" + } + + now = NaiveDateTime.utc_now(:second) + time = NaiveDateTime.add(now, 10, :second) + + conn = + build_conn(:post, "/api/events", payload) + |> put_req_header("x-replay-event-id", "123") + |> put_req_header("x-replay-session-id", "456") + |> put_req_header("x-replay-time", NaiveDateTime.to_iso8601(time)) + + assert {:ok, request, _conn} = Request.build(conn) + assert is_nil(request.replay_session_id) + assert NaiveDateTime.compare(request.timestamp, time) == :lt + end + @tag :ee_only test "parses revenue source field from a json string" do payload = %{ @@ -586,25 +642,31 @@ defmodule Plausible.Ingestion.RequestTest do request = request |> Jason.encode!() |> Jason.decode!() - assert Map.drop(request, ["timestamp"]) == %{ - "domains" => ["dummy.site"], - "event_name" => "pageview", - "hash_mode" => 1, - "hostname" => "dummy.site", - "pathname" => "/pictures/index.html", - "props" => %{"abc" => "qwerty", "hello" => "world"}, - "query_params" => %{"baz" => "bam", "foo" => "bar"}, - "referrer" => "https://example.com", - "remote_ip" => "127.0.0.1", - "revenue_source" => %{"amount" => "12.3", "currency" => "USD"}, - "uri" => "https://dummy.site/pictures/index.html?foo=bar&baz=bam", - "user_agent" => "Mozilla", - "ip_classification" => nil, - "scroll_depth" => nil, - "engagement_time" => nil, - "tracker_script_version" => 0, - "interactive?" => true - } + expected = %{ + "domains" => ["dummy.site"], + "event_name" => "pageview", + "hash_mode" => 1, + "hostname" => "dummy.site", + "pathname" => "/pictures/index.html", + "props" => %{"abc" => "qwerty", "hello" => "world"}, + "query_params" => %{"baz" => "bam", "foo" => "bar"}, + "referrer" => "https://example.com", + "remote_ip" => "127.0.0.1", + "revenue_source" => %{"amount" => "12.3", "currency" => "USD"}, + "uri" => "https://dummy.site/pictures/index.html?foo=bar&baz=bam", + "user_agent" => "Mozilla", + "ip_classification" => nil, + "scroll_depth" => nil, + "engagement_time" => nil, + "tracker_script_version" => 0, + "interactive?" => true + } + + on_ee do + expected = Map.put(expected, "replay_session_id", nil) + end + + assert Map.drop(request, ["timestamp"]) == expected assert %NaiveDateTime{} = NaiveDateTime.from_iso8601!(request["timestamp"]) end diff --git a/test/plausible/session/cache_store_test.exs b/test/plausible/session/cache_store_test.exs index 8efb8c90cb28..e07fb28d7a2a 100644 --- a/test/plausible/session/cache_store_test.exs +++ b/test/plausible/session/cache_store_test.exs @@ -246,6 +246,76 @@ defmodule Plausible.Session.CacheStoreTest do # assert Map.get(session, :"entry.meta.value") == ["true", "false"] end + @tag :ee_only + test "creates a session from a replayed event", %{buffer: buffer} do + event = + build(:event, + name: "pageview", + replay_session_id: 456, + "meta.key": ["logged_in", "darkmode"], + "meta.value": ["true", "false"] + ) + + CacheStore.on_event(event, @session_params, nil, buffer_insert: buffer) + + assert_receive({:buffer, :insert, [sessions]}) + assert [session] = sessions + assert session.hostname == event.hostname + + assert session.site_id == event.site_id + + assert session.replay_session_id == 456 + + assert session.user_id == event.user_id + assert session.entry_page == event.pathname + assert session.exit_page == event.pathname + assert session.is_bounce == true + assert session.duration == 0 + assert session.pageviews == 1 + assert session.events == 1 + assert session.referrer == Map.get(@session_params, :referrer) + assert session.referrer_source == Map.get(@session_params, :referrer_source) + assert session.utm_medium == Map.get(@session_params, :utm_medium) + assert session.utm_source == Map.get(@session_params, :utm_source) + assert session.utm_campaign == Map.get(@session_params, :utm_campaign) + assert session.utm_content == Map.get(@session_params, :utm_content) + assert session.utm_term == Map.get(@session_params, :utm_term) + assert session.country_code == Map.get(@session_params, :country_code) + assert session.screen_size == Map.get(@session_params, :screen_size) + assert session.operating_system == Map.get(@session_params, :operating_system) + assert session.operating_system_version == Map.get(@session_params, :operating_system_version) + assert session.browser == Map.get(@session_params, :browser) + assert session.browser_version == Map.get(@session_params, :browser_version) + assert session.timestamp == event.timestamp + assert session.start === event.timestamp + end + + @tag :ee_only + test "updates session counters on replayed event", %{buffer: buffer} do + timestamp = DateTime.utc_now() + + event1 = + build(:event, + name: "pageview", + replay_session_id: 456, + timestamp: timestamp |> NaiveDateTime.shift(second: -10) + ) + + event2 = %{ + event1 + | timestamp: timestamp + } + + CacheStore.on_event(event1, %{}, nil, buffer_insert: buffer) + CacheStore.on_event(event2, %{}, nil, buffer_insert: buffer) + assert_receive({:buffer, :insert, [[_negative_record, session]]}) + assert session.is_bounce == false + assert session.duration == 10 + assert session.pageviews == 2 + assert session.events == 2 + assert session.replay_session_id == 456 + end + test "updates session counters", %{buffer: buffer} do timestamp = DateTime.utc_now() @@ -266,6 +336,33 @@ defmodule Plausible.Session.CacheStoreTest do assert session.events == 2 end + @tag :ee_only + test "treats replayed and normal events as belonging to distinct sessions even if user_id and site_id match", + %{buffer: buffer} do + timestamp = DateTime.utc_now() + + event1 = + build(:event, + name: "pageview", + replay_session_id: 456, + timestamp: timestamp |> NaiveDateTime.shift(second: -10) + ) + + event2 = %{ + event1 + | replay_session_id: nil, + timestamp: timestamp + } + + CacheStore.on_event(event1, %{}, nil, buffer_insert: buffer) + assert_receive({:buffer, :insert, [[session1]]}) + CacheStore.on_event(event2, %{}, nil, buffer_insert: buffer) + assert_receive({:buffer, :insert, [[session2]]}) + assert session1.session_id != session2.session_id + assert session1.replay_session_id == 456 + assert is_nil(session2.replay_session_id) + end + test "does not update session counters on engagement event", %{buffer: buffer} do now = NaiveDateTime.utc_now(:second) pageview = build(:pageview, timestamp: NaiveDateTime.shift(now, second: -10))