diff --git a/assets/js/source_lv_hooks.js b/assets/js/source_lv_hooks.js
index 55c59d6e23..49a4baf459 100644
--- a/assets/js/source_lv_hooks.js
+++ b/assets/js/source_lv_hooks.js
@@ -26,7 +26,7 @@ hooks.SourceLogsSearchList = {
const hook = this
activateDelegatedTooltips(this.el, '[data-toggle="tooltip"]')
- window.scrollTo(0, document.body.scrollHeight)
+ // window.scrollTo(0, document.body.scrollHeight)
const observer =
new IntersectionObserver((entries, observer) => {
diff --git a/lib/logflare/logs/log_event.ex b/lib/logflare/logs/log_event.ex
index e0bd3c796a..e6a001ef66 100644
--- a/lib/logflare/logs/log_event.ex
+++ b/lib/logflare/logs/log_event.ex
@@ -25,7 +25,6 @@ defmodule Logflare.LogEvent do
field :body, :map, default: %{}
field :valid, :boolean
field :drop, :boolean, default: false
- field :is_from_stale_query, :boolean
field :timestamp_inferred, :boolean, default: false
field :ingested_at, :utc_datetime_usec
field :source_uuid, Ecto.UUID.Atom
diff --git a/lib/logflare/logs/logs_search.ex b/lib/logflare/logs/logs_search.ex
index f47d205fd4..4cfba12ccb 100644
--- a/lib/logflare/logs/logs_search.ex
+++ b/lib/logflare/logs/logs_search.ex
@@ -44,6 +44,7 @@ defmodule Logflare.Logs.Search do
%{error: nil} = so <- apply_local_timestamp_correction(so),
%{error: nil} = so <- apply_timestamp_filter_rules(so),
%{error: nil} = so <- apply_filters(so),
+ %{error: nil} = so <- apply_cursor(so),
%{error: nil} = so <- apply_select_rules(so),
%{error: nil} = so <- do_query(so),
%{error: nil} = so <- apply_warning_conditions(so),
diff --git a/lib/logflare/logs/search_operation.ex b/lib/logflare/logs/search_operation.ex
index 26acb1225c..f830f26d1d 100644
--- a/lib/logflare/logs/search_operation.ex
+++ b/lib/logflare/logs/search_operation.ex
@@ -35,8 +35,30 @@ defmodule Logflare.Logs.SearchOperation do
field :chart_data_shape_id, atom(), default: nil, enforce: true
field :type, :events | :aggregates
field :status, {atom(), String.t() | [String.t()]}
+ field :event_page_request, event_page_request()
+ field :event_page_result, event_page_result()
+ field :has_more_events?, boolean(), default: false
end
+ @type event_cursor :: %{timestamp: integer(), id: String.t()}
+ @type event_page_intent :: :within_range | :extend_previous | :extend_next
+ @type event_page_request :: %{
+ intent: event_page_intent(),
+ boundary: integer() | DateTime.t() | NaiveDateTime.t() | nil,
+ cursor: event_cursor() | nil
+ }
+ @type event_page_result :: %{
+ request: event_page_request(),
+ has_more?: boolean(),
+ cursor: event_cursor() | nil
+ }
+
+ @spec event_page_direction(event_page_intent()) :: :previous | :next
+ def event_page_direction(:extend_next), do: :next
+
+ def event_page_direction(intent) when intent in [:within_range, :extend_previous],
+ do: :previous
+
def new(params) do
so = struct(__MODULE__, params)
diff --git a/lib/logflare/logs/search_operations.ex b/lib/logflare/logs/search_operations.ex
index 6f4c54387d..dacddfab9e 100644
--- a/lib/logflare/logs/search_operations.ex
+++ b/lib/logflare/logs/search_operations.ex
@@ -47,6 +47,32 @@ defmodule Logflare.Logs.SearchOperations do
@spec max_chart_ticks :: integer()
def max_chart_ticks, do: @default_max_n_chart_ticks
+ @spec default_limit :: pos_integer()
+ def default_limit, do: @default_limit
+
+ @spec fetch_limit :: pos_integer()
+ def fetch_limit, do: @default_limit + 1
+
+ @spec event_page_params(SO.t(), SO.event_page_intent(), SO.event_cursor() | nil) ::
+ {:ok, %{event_page_request: SO.event_page_request(), tailing?: false}} | :error
+ def event_page_params(%SO{tailing?: false} = so, intent, cursor)
+ when (intent == :within_range and is_map(cursor)) or
+ (intent in [:extend_previous, :extend_next] and (is_map(cursor) or is_nil(cursor))) do
+ with {:ok, boundary} <- event_page_boundary(so, intent, cursor) do
+ {:ok,
+ %{
+ event_page_request: %{
+ intent: intent,
+ boundary: boundary,
+ cursor: cursor
+ },
+ tailing?: false
+ }}
+ end
+ end
+
+ def event_page_params(%SO{}, _intent, _cursor), do: :error
+
@spec do_query(SO.t()) :: SO.t()
def do_query(%SO{} = so) do
with {:ok, response} <- execute_backend_query(so) do
@@ -122,15 +148,109 @@ defmodule Logflare.Logs.SearchOperations do
@spec apply_query_defaults(SO.t()) :: SO.t()
def apply_query_defaults(%SO{} = so) do
+ intent = if so.event_page_request, do: so.event_page_request.intent, else: :within_range
+ direction = SO.event_page_direction(intent)
+
query =
from(table_name(so))
|> select(%{})
- |> order_by([t], desc: t.timestamp)
- |> limit(@default_limit)
+ |> order_events(direction)
+ |> limit(^fetch_limit())
%{so | query: query}
end
+ defp event_page_boundary(%SO{} = so, intent, cursor)
+ when intent in [:extend_previous, :extend_next] and
+ (is_map(cursor) or intent == :extend_next) do
+ case event_range_bounds(so) do
+ {:ok, bounds} -> {:ok, Map.fetch!(bounds, event_page_boundary_key(intent))}
+ :error -> {:ok, nil}
+ end
+ end
+
+ defp event_page_boundary(%SO{} = so, intent, _cursor) do
+ with {:ok, bounds} <- event_range_bounds(so) do
+ boundary =
+ case intent do
+ :within_range -> nil
+ intent -> Map.fetch!(bounds, event_page_boundary_key(intent))
+ end
+
+ {:ok, boundary}
+ end
+ end
+
+ defp event_page_boundary_key(:extend_previous), do: :min
+ defp event_page_boundary_key(:extend_next), do: :max
+
+ defp event_range_bounds(%SO{} = so) do
+ if bounded_timestamp_filters?(so.lql_ts_filters) do
+ bounds =
+ SearchOperationHelpers.get_min_max_filter_timestamps(
+ so.lql_ts_filters,
+ chart_period(so)
+ )
+
+ {:ok, Map.take(bounds, [:min, :max])}
+ else
+ :error
+ end
+ end
+
+ defp bounded_timestamp_filters?(filters) do
+ Enum.any?(filters, &(&1.operator == :range and length(&1.values || []) == 2)) or
+ (Enum.any?(filters, &(&1.operator in [:>, :>=])) and
+ Enum.any?(filters, &(&1.operator in [:<, :<=])))
+ end
+
+ defp order_events(query, :previous),
+ do: order_by(query, [t], desc: t.timestamp, desc: t.id)
+
+ defp order_events(query, :next), do: order_by(query, [t], asc: t.timestamp, asc: t.id)
+
+ @spec apply_cursor(SO.t()) :: SO.t()
+ def apply_cursor(%SO{event_page_request: %{cursor: nil}} = so), do: so
+
+ def apply_cursor(%SO{event_page_request: %{intent: intent, cursor: cursor}} = so) do
+ %{timestamp: timestamp, id: id} = cursor
+ timestamp = normalize_event_timestamp(so.backend_type, timestamp)
+ direction = SO.event_page_direction(intent)
+
+ %{so | query: where(so.query, ^cursor_condition(direction, timestamp, id))}
+ end
+
+ def apply_cursor(%SO{} = so), do: so
+
+ defp normalize_event_timestamp(:postgres, timestamp) when is_integer(timestamp),
+ do: DateTime.from_unix!(timestamp, :microsecond)
+
+ defp normalize_event_timestamp(_backend_type, timestamp), do: timestamp
+
+ defp cursor_condition(:previous, timestamp, id) when is_integer(timestamp) do
+ dynamic(
+ [t],
+ t.timestamp < fragment("TIMESTAMP_MICROS(?)", ^timestamp) or
+ (t.timestamp == fragment("TIMESTAMP_MICROS(?)", ^timestamp) and t.id < ^id)
+ )
+ end
+
+ defp cursor_condition(:next, timestamp, id) when is_integer(timestamp) do
+ dynamic(
+ [t],
+ t.timestamp > fragment("TIMESTAMP_MICROS(?)", ^timestamp) or
+ (t.timestamp == fragment("TIMESTAMP_MICROS(?)", ^timestamp) and t.id > ^id)
+ )
+ end
+
+ defp cursor_condition(:previous, timestamp, id) do
+ dynamic([t], t.timestamp < ^timestamp or (t.timestamp == ^timestamp and t.id < ^id))
+ end
+
+ defp cursor_condition(:next, timestamp, id) do
+ dynamic([t], t.timestamp > ^timestamp or (t.timestamp == ^timestamp and t.id > ^id))
+ end
+
@spec apply_halt_conditions(SO.t()) :: SO.t()
def apply_halt_conditions(%SO{} = so) do
chart_period = chart_period(so)
@@ -239,6 +359,11 @@ defmodule Logflare.Logs.SearchOperations do
%{so | rows: rows}
end
+ def apply_timestamp_filter_rules(%SO{type: :events, event_page_request: %{intent: intent}} = so)
+ when intent in [:extend_previous, :extend_next] do
+ apply_event_page_boundary(so)
+ end
+
def apply_timestamp_filter_rules(%SO{backend_type: :postgres, type: :events} = so) do
%{so | query: apply_postgres_event_timestamp_filter_rules(so)}
end
@@ -342,6 +467,82 @@ defmodule Logflare.Logs.SearchOperations do
%{so | query: q}
end
+ defp apply_event_page_boundary(
+ %SO{
+ event_page_request: %{intent: intent, boundary: nil, cursor: %{timestamp: timestamp}}
+ } = so
+ ) do
+ direction = SO.event_page_direction(intent)
+ timestamp = normalize_event_timestamp(so.backend_type, timestamp)
+
+ %{so | query: apply_event_page_partition_filter(so.query, so, direction, timestamp)}
+ end
+
+ defp apply_event_page_boundary(%SO{event_page_request: %{boundary: nil}} = so), do: so
+
+ defp apply_event_page_boundary(
+ %SO{
+ event_page_request: %{intent: intent, boundary: boundary, cursor: %{}},
+ query: query
+ } = so
+ ) do
+ boundary = normalize_event_timestamp(so.backend_type, boundary)
+ direction = SO.event_page_direction(intent)
+ %{so | query: apply_event_page_partition_filter(query, so, direction, boundary)}
+ end
+
+ defp apply_event_page_boundary(
+ %SO{event_page_request: %{intent: intent, boundary: boundary}} = so
+ ) do
+ boundary = normalize_event_timestamp(so.backend_type, boundary)
+ direction = SO.event_page_direction(intent)
+ query = where(so.query, ^boundary_condition(direction, boundary))
+
+ %{so | query: apply_event_page_partition_filter(query, so, direction, boundary)}
+ end
+
+ defp boundary_condition(:previous, boundary) when is_integer(boundary),
+ do: dynamic([t], t.timestamp < fragment("TIMESTAMP_MICROS(?)", ^boundary))
+
+ defp boundary_condition(:next, boundary) when is_integer(boundary),
+ do: dynamic([t], t.timestamp > fragment("TIMESTAMP_MICROS(?)", ^boundary))
+
+ defp boundary_condition(:previous, boundary), do: dynamic([t], t.timestamp < ^boundary)
+ defp boundary_condition(:next, boundary), do: dynamic([t], t.timestamp > ^boundary)
+
+ defp apply_event_page_partition_filter(query, %SO{backend_type: :postgres}, _, _), do: query
+
+ defp apply_event_page_partition_filter(
+ query,
+ %SO{partition_by: :timestamp},
+ direction,
+ boundary
+ ) do
+ date = boundary |> event_boundary_datetime() |> DateTime.to_date()
+
+ case direction do
+ :previous -> where(query, [t], fragment("EXTRACT(DATE FROM ?)", t.timestamp) <= ^date)
+ :next -> where(query, [t], fragment("EXTRACT(DATE FROM ?)", t.timestamp) >= ^date)
+ end
+ end
+
+ defp apply_event_page_partition_filter(query, %SO{partition_by: :pseudo}, direction, boundary) do
+ date = boundary |> event_boundary_datetime() |> DateTime.to_date()
+
+ case direction do
+ :previous -> where(query, partition_date() <= ^date or in_streaming_buffer())
+ :next -> where(query, partition_date() >= ^date or in_streaming_buffer())
+ end
+ end
+
+ defp event_boundary_datetime(timestamp) when is_integer(timestamp),
+ do: DateTime.from_unix!(timestamp, :microsecond)
+
+ defp event_boundary_datetime(%DateTime{} = timestamp), do: timestamp
+
+ defp event_boundary_datetime(%NaiveDateTime{} = timestamp),
+ do: DateTime.from_naive!(timestamp, "Etc/UTC")
+
defp apply_bq_aggregate_timestamp_filters(query, so, filters, chart_period) do
period = to_bq_interval_token(chart_period)
tick_count = SearchOperationHelpers.default_period_tick_count(chart_period)
diff --git a/lib/logflare/logs/search_query_executor.ex b/lib/logflare/logs/search_query_executor.ex
index 2c7ef8251c..b0d914c235 100644
--- a/lib/logflare/logs/search_query_executor.ex
+++ b/lib/logflare/logs/search_query_executor.ex
@@ -11,6 +11,7 @@ defmodule Logflare.Logs.SearchQueryExecutor do
alias Logflare.LogEvent
alias Logflare.Logs.Search
alias Logflare.Logs.SearchOperation, as: SO
+ alias Logflare.Logs.SearchOperations
alias Logflare.Utils.Tasks
@query_timeout 30_000
@@ -37,11 +38,19 @@ defmodule Logflare.Logs.SearchQueryExecutor do
end
def query(pid, params) do
- GenServer.call(pid, {:query, params}, @query_timeout)
+ GenServer.call(pid, {:query, query_params(params)}, @query_timeout)
+ end
+
+ def query_page(pid, params, intent, cursor) do
+ GenServer.call(
+ pid,
+ {:query_page, query_params(params), intent, cursor},
+ @query_timeout
+ )
end
def query_agg(pid, params) do
- GenServer.call(pid, {:query_agg, params}, @query_timeout)
+ GenServer.call(pid, {:query_agg, query_params(params)}, @query_timeout)
end
def cancel_agg(pid) do
@@ -75,6 +84,25 @@ defmodule Logflare.Logs.SearchQueryExecutor do
{:reply, :ok, %{state | event_task: {new_ref, new_params}}}
end
+ def handle_call({:query_page, params, intent, cursor}, {lv_pid, _ref}, state) do
+ search_op = SO.new(params)
+
+ case SearchOperations.event_page_params(search_op, intent, cursor) do
+ {:ok, page_params} ->
+ {ref, _params} = state.event_task
+
+ if ref, do: Task.shutdown(ref, :brutal_kill)
+
+ search_op = struct(search_op, page_params)
+ new_ref = start_search_task(lv_pid, search_op)
+
+ {:reply, :ok, %{state | event_task: {new_ref, page_params}}}
+
+ :error ->
+ {:reply, :error, state}
+ end
+ end
+
def handle_call({:query_agg, new_params}, {lv_pid, _ref}, state) do
{ref, _params} = state.agg_task
@@ -118,61 +146,40 @@ defmodule Logflare.Logs.SearchQueryExecutor do
end
@impl true
- def handle_info({_ref, {:search_result, lv_pid, %{events: events_so}}}, state) do
- Logger.debug(
- "SearchQueryExecutor: Getting search events for #{pid_to_string(lv_pid)} / #{state.source_id} source..."
- )
-
- {_ref, params} = state.event_task
-
- rows = Enum.map(events_so.rows, &LogEvent.make_from_db(&1, %{source: params.source}))
-
- old_rows = if params.search_op_log_events, do: params.search_op_log_events.rows, else: []
-
- # prevents removal of log events loaded
- # during initial tailing query
- log_events =
- old_rows
- |> Enum.reject(& &1.is_from_stale_query)
- |> Enum.concat(rows)
- |> Enum.uniq_by(&{&1.body, &1.id})
- |> Enum.sort_by(& &1.body["timestamp"], &>=/2)
- |> Enum.take(100)
-
- send(
- state.caller,
- {:search_result,
- %{
- events: %{events_so | rows: log_events}
- }}
- )
-
- {:noreply, %{state | event_task: {nil, nil}}}
+ def handle_info({ref, {:search_result, lv_pid, %{events: events_so}}}, state) do
+ if active_task?(state.event_task, ref) do
+ handle_event_result(lv_pid, events_so, state)
+ else
+ {:noreply, state}
+ end
end
@impl true
- def handle_info({_ref, {:search_result, lv_pid, %{aggregates: aggregates_so}}}, state) do
- Logger.debug(
- "SearchQueryExecutor: Getting search aggregates for #{pid_to_string(lv_pid)} / #{state.source_id} source..."
- )
-
- {_ref, _params} = state.agg_task
-
- send(
- state.caller,
- {:search_result,
- %{
- aggregates: aggregates_so
- }}
- )
-
- {:noreply, %{state | agg_task: {nil, nil}}}
+ def handle_info({ref, {:search_result, lv_pid, %{aggregates: aggregates_so}}}, state) do
+ if active_task?(state.agg_task, ref) do
+ handle_aggregate_result(lv_pid, aggregates_so, state)
+ else
+ {:noreply, state}
+ end
end
@impl true
- def handle_info({_ref, {:search_error, _lv_pid, %SO{} = search_op}}, state) do
- send(state.caller, {:search_error, search_op})
- {:noreply, state}
+ def handle_info({ref, {:search_error, _lv_pid, %SO{type: :aggregates} = search_op}}, state) do
+ if active_task?(state.agg_task, ref) do
+ send(state.caller, {:search_error, search_op})
+ {:noreply, %{state | agg_task: {nil, nil}}}
+ else
+ {:noreply, state}
+ end
+ end
+
+ def handle_info({ref, {:search_error, _lv_pid, %SO{} = search_op}}, state) do
+ if active_task?(state.event_task, ref) do
+ send(state.caller, {:search_error, search_op})
+ {:noreply, %{state | event_task: {nil, nil}}}
+ else
+ {:noreply, state}
+ end
end
# handles task shutdown messages
@@ -188,9 +195,53 @@ defmodule Logflare.Logs.SearchQueryExecutor do
{:noreply, state}
end
+ defp handle_event_result(lv_pid, events_so, state) do
+ Logger.debug(
+ "SearchQueryExecutor: Getting search events for #{pid_to_string(lv_pid)} / #{state.source_id} source..."
+ )
+
+ page_size = SearchOperations.default_limit()
+ raw_rows = events_so.rows
+ has_sentinel_row? = has_sentinel_row?(raw_rows)
+
+ page_rows =
+ raw_rows
+ |> Enum.take(page_size)
+ |> Enum.map(&LogEvent.make_from_db(&1, %{source: events_so.source}))
+ |> uniq_sort_log_events()
+
+ event_page_result =
+ event_page_result(events_so.event_page_request, page_rows, has_sentinel_row?)
+
+ events_so = %{
+ events_so
+ | rows: page_rows,
+ has_more_events?: has_sentinel_row?,
+ event_page_result: event_page_result
+ }
+
+ send(state.caller, {:search_result, %{events: events_so}})
+
+ {:noreply, %{state | event_task: {nil, nil}}}
+ end
+
+ defp handle_aggregate_result(lv_pid, aggregates_so, state) do
+ Logger.debug(
+ "SearchQueryExecutor: Getting search aggregates for #{pid_to_string(lv_pid)} / #{state.source_id} source..."
+ )
+
+ send(state.caller, {:search_result, %{aggregates: aggregates_so}})
+
+ {:noreply, %{state | agg_task: {nil, nil}}}
+ end
+
+ def start_search_task(lv_pid, %SO{} = so), do: do_start_search_task(lv_pid, so)
+
def start_search_task(lv_pid, params) do
- so = SO.new(params)
+ do_start_search_task(lv_pid, SO.new(params))
+ end
+ defp do_start_search_task(lv_pid, so) do
Tasks.async(fn ->
so
|> Search.search()
@@ -219,4 +270,46 @@ defmodule Logflare.Logs.SearchQueryExecutor do
end
end)
end
+
+ defp uniq_sort_log_events(log_events) do
+ log_events
+ |> Enum.uniq_by(&{&1.body["timestamp"], event_id(&1)})
+ |> Enum.sort_by(&{&1.body["timestamp"], event_id(&1)}, :desc)
+ end
+
+ defp has_sentinel_row?(rows), do: length(rows) >= SearchOperations.fetch_limit()
+
+ defp log_event_cursor(%LogEvent{} = event) do
+ %{timestamp: event.body["timestamp"], id: event_id(event)}
+ end
+
+ defp event_page_result(nil, _rows, _has_more?), do: nil
+
+ defp event_page_result(request, rows, has_more?) do
+ %{
+ request: request,
+ has_more?: has_more?,
+ cursor: page_edge_cursor(rows, SO.event_page_direction(request.intent))
+ }
+ end
+
+ defp page_edge_cursor([], _direction), do: nil
+ defp page_edge_cursor(rows, :previous), do: rows |> List.last() |> log_event_cursor()
+ defp page_edge_cursor(rows, :next), do: rows |> List.first() |> log_event_cursor()
+
+ defp active_task?({%Task{ref: ref}, _params}, ref), do: true
+ defp active_task?({ref, _params}, ref) when is_reference(ref), do: true
+ defp active_task?(_task, _ref), do: false
+
+ defp event_id(%LogEvent{id: id, body: body}), do: id || body["id"]
+
+ defp query_params(params) do
+ Map.drop(params, [
+ :range_extension_patch,
+ :search_op,
+ :search_op_log_events,
+ :search_op_log_aggregates,
+ :streams
+ ])
+ end
end
diff --git a/lib/logflare/lql/rules.ex b/lib/logflare/lql/rules.ex
index 4d65938e9d..cf29d64fbf 100644
--- a/lib/logflare/lql/rules.ex
+++ b/lib/logflare/lql/rules.ex
@@ -16,6 +16,8 @@ defmodule Logflare.Lql.Rules do
@type lql_rule :: ChartRule.t() | FilterRule.t() | FromRule.t() | SelectRule.t()
@type lql_rules :: [lql_rule()]
+ @type timestamp_extension_direction :: :previous | :next
+ @type timestamp_value :: integer() | Date.t() | DateTime.t() | NaiveDateTime.t()
# =============================================================================
# Rule Type Extraction
@@ -276,6 +278,50 @@ defmodule Logflare.Lql.Rules do
|> Enum.concat(new_timestamp_rules)
end
+ @doc """
+ Replaces timestamp filters with an absolute range extended to an event timestamp.
+
+ The previous direction replaces the lower edge, while the next direction
+ replaces the upper edge. All other rules are preserved.
+ """
+ @spec extend_timestamp_range(
+ lql_rules(),
+ timestamp_extension_direction(),
+ timestamp_value(),
+ Calendar.time_zone()
+ ) :: lql_rules()
+ def extend_timestamp_range(lql_rules, direction, event_timestamp, timezone \\ "Etc/UTC")
+ when is_list(lql_rules) and direction in [:previous, :next] and is_binary(timezone) do
+ timestamp_values =
+ lql_rules
+ |> get_timestamp_filters()
+ |> Enum.flat_map(×tamp_values/1)
+ |> Enum.map(&normalize_timestamp/1)
+
+ case timestamp_values do
+ [] ->
+ lql_rules
+
+ timestamp_values ->
+ event_timestamp = normalize_timestamp(event_timestamp, timezone)
+
+ values =
+ case direction do
+ :previous -> [event_timestamp, Enum.max(timestamp_values)]
+ :next -> [Enum.min(timestamp_values), event_timestamp]
+ end
+
+ timestamp_rule =
+ FilterRule.build(
+ path: "timestamp",
+ operator: :range,
+ values: values
+ )
+
+ update_timestamp_rules(lql_rules, [timestamp_rule])
+ end
+ end
+
@doc """
Creates new timestamp filters by jumping forward or backward in time.
@@ -299,6 +345,40 @@ defmodule Logflare.Lql.Rules do
FilterRule.shorthand_timestamp?(filter_rule)
end
+ defp timestamp_values(%FilterRule{value: value, values: values}) do
+ [value | List.wrap(values)]
+ |> Enum.reject(&is_nil/1)
+ end
+
+ defp normalize_timestamp(timestamp) when is_integer(timestamp) do
+ timestamp
+ |> DateTime.from_unix!(:microsecond)
+ |> DateTime.to_naive()
+ end
+
+ defp normalize_timestamp(%DateTime{} = timestamp) do
+ timestamp
+ |> DateTime.shift_zone!("Etc/UTC")
+ |> DateTime.to_naive()
+ end
+
+ defp normalize_timestamp(%NaiveDateTime{} = timestamp), do: timestamp
+ defp normalize_timestamp(%Date{} = timestamp), do: NaiveDateTime.new!(timestamp, ~T[00:00:00])
+
+ defp normalize_timestamp(timestamp, timezone) when is_integer(timestamp) do
+ timestamp
+ |> DateTime.from_unix!(:microsecond)
+ |> normalize_timestamp(timezone)
+ end
+
+ defp normalize_timestamp(%DateTime{} = timestamp, timezone) do
+ timestamp
+ |> DateTime.shift_zone!(timezone)
+ |> DateTime.to_naive()
+ end
+
+ defp normalize_timestamp(timestamp, _timezone), do: normalize_timestamp(timestamp)
+
# =============================================================================
# LQL Parser Warnings
# =============================================================================
diff --git a/lib/logflare_web/live/log_event_live.ex b/lib/logflare_web/live/log_event_live.ex
index 3c2d22688f..e8fc8198a5 100644
--- a/lib/logflare_web/live/log_event_live.ex
+++ b/lib/logflare_web/live/log_event_live.ex
@@ -24,6 +24,8 @@ defmodule LogflareWeb.LogEventLive do
lql = params["lql"] || ""
is_tailing = params["tailing?"] == "true"
+ search_timezone = params["tz"] || preferred_timezone(socket.assigns)
+
opts =
[
source: source,
@@ -49,7 +51,7 @@ defmodule LogflareWeb.LogEventLive do
|> assign(:log_event_id, params["uuid"])
|> assign(:lql, lql)
|> assign(:tailing?, is_tailing)
- |> assign(:tz, params["tz"])
+ |> assign(:search_timezone, search_timezone)
|> assign(:timestamp, timestamp)
{:ok, socket}
@@ -68,4 +70,15 @@ defmodule LogflareWeb.LogEventLive do
defp maybe_put_timestamp(opts, timestamp),
do: Keyword.put(opts, :timestamp, DateTime.truncate(timestamp, :second))
+
+ @spec preferred_timezone(map()) :: String.t()
+ defp preferred_timezone(%{team_user: %{preferences: %{timezone: timezone}}})
+ when is_binary(timezone),
+ do: timezone
+
+ defp preferred_timezone(%{user: %{preferences: %{timezone: timezone}}})
+ when is_binary(timezone),
+ do: timezone
+
+ defp preferred_timezone(_assigns), do: "Etc/UTC"
end
diff --git a/lib/logflare_web/live/log_event_live/search_log_event_viewer_component.ex b/lib/logflare_web/live/log_event_live/search_log_event_viewer_component.ex
index 33d8a0c3ae..0ef4c9285d 100644
--- a/lib/logflare_web/live/log_event_live/search_log_event_viewer_component.ex
+++ b/lib/logflare_web/live/log_event_live/search_log_event_viewer_component.ex
@@ -37,7 +37,7 @@ defmodule LogflareWeb.Search.LogEventViewerComponent do
params =
event_params(assigns)
|> Map.merge(%{log_event_id: id, timestamp: d})
- |> Map.put(:lql, assigns.params["lql"] || "")
+ |> Map.put(:lql, assigns.lql)
socket =
socket
@@ -94,11 +94,7 @@ defmodule LogflareWeb.Search.LogEventViewerComponent do
@impl true
def render(%{source: source, log_event: %LE{body: body} = le} = assigns) do
- tz =
- if assigns.team_user,
- do: Map.get(assigns.team_user.preferences || %{}, :timezone, "Etc/UTC"),
- else: Map.get(assigns.user.preferences || %{}, :timezone, "Etc/UTC")
-
+ tz = assigns.search_timezone
timestamp = Timex.from_unix(body["timestamp"], :microsecond)
local_timestamp =
@@ -110,7 +106,7 @@ defmodule LogflareWeb.Search.LogEventViewerComponent do
LogView.render("log_event_body.html",
source: source,
source_schema_flat_map: assigns.source_schema_flat_map,
- search_params: assigns.search_params,
+ search_params: %{"tz" => tz},
team: assigns.team,
body: body,
fmt_body: BqSchema.encode_metadata(body),
@@ -119,8 +115,7 @@ defmodule LogflareWeb.Search.LogEventViewerComponent do
lql: assigns.lql,
lql_schema: get_lql_schema(source),
timestamp: timestamp,
- local_timezone: tz,
- search_timezone: assigns.search_params["tz"] || tz,
+ local_timezone: assigns.search_timezone,
local_timestamp: local_timestamp
)
end
@@ -136,13 +131,12 @@ defmodule LogflareWeb.Search.LogEventViewerComponent do
team = socket.assigns[:team] || assigns[:team]
source = socket.assigns[:source] || assigns[:source]
timestamp = socket.assigns[:timestamp] || assigns[:timestamp]
- lql = socket.assigns[:lql] || assigns[:lql] || assigns.params["lql"] || ""
+ lql = assigns[:lql] || socket.assigns[:lql] || ""
source_schema_flat_map =
socket.assigns[:source_schema_flat_map] || assigns[:source_schema_flat_map]
- search_params =
- socket.assigns[:search_params] || extract_search_params(assigns)
+ search_timezone = assigns[:search_timezone] || socket.assigns[:search_timezone] || "Etc/UTC"
socket
|> assign(:user, user)
@@ -152,7 +146,7 @@ defmodule LogflareWeb.Search.LogEventViewerComponent do
|> assign(:timestamp, timestamp)
|> assign(:lql, lql)
|> assign(:source_schema_flat_map, source_schema_flat_map)
- |> assign(:search_params, search_params)
+ |> assign(:search_timezone, search_timezone)
|> assign(:error, nil)
end
@@ -172,9 +166,4 @@ defmodule LogflareWeb.Search.LogEventViewerComponent do
_ -> SchemaBuilder.initial_table_schema()
end
end
-
- defp extract_search_params(%{params: params}) when is_map(params),
- do: Map.take(params, ["tz"])
-
- defp extract_search_params(_assigns), do: %{}
end
diff --git a/lib/logflare_web/live/search_live/event_context_component.ex b/lib/logflare_web/live/search_live/event_context_component.ex
index f678c0e327..0c1a041180 100644
--- a/lib/logflare_web/live/search_live/event_context_component.ex
+++ b/lib/logflare_web/live/search_live/event_context_component.ex
@@ -14,18 +14,19 @@ defmodule LogflareWeb.SearchLive.EventContextComponent do
%{
params: %{
"log-event-timestamp" => log_timestamp,
- "log-event-id" => log_event_id,
- "querystring" => query_string,
- "source-id" => source_id,
- "timezone" => timezone
- }
+ "log-event-id" => log_event_id
+ },
+ querystring: query_string,
+ search_timezone: timezone,
+ source: source
} = assigns
+ source_id = source.id
+
event_timestamp = log_timestamp |> String.to_integer() |> Timex.from_unix(:microsecond)
lql_rules =
- Sources.get_source_for_lv_param(source_id)
- |> prepare_lql_rules(query_string, event_timestamp)
+ prepare_lql_rules(source, query_string, event_timestamp)
{:ok,
socket
@@ -33,7 +34,7 @@ defmodule LogflareWeb.SearchLive.EventContextComponent do
|> assign(target_event_id: log_event_id, timezone: timezone)
|> assign(is_truncated_before: false)
|> assign(is_truncated_after: false)
- |> assign(source: Sources.get_source_for_lv_param(source_id))
+ |> assign(source: source)
|> assign(:logs, AsyncResult.loading())
|> start_async(:logs, fn ->
search_logs(log_event_id, event_timestamp, source_id, lql_rules)
diff --git a/lib/logflare_web/live/search_live/form_components.ex b/lib/logflare_web/live/search_live/form_components.ex
index dfeeb35893..d9d3cc5f23 100644
--- a/lib/logflare_web/live/search_live/form_components.ex
+++ b/lib/logflare_web/live/search_live/form_components.ex
@@ -135,7 +135,6 @@ defmodule LogflareWeb.SearchLive.FormComponents do
attr :uri_params, :map, required: true
attr :lql_rules, :list, required: true
attr :user, Logflare.User, required: true
- attr :search_op_log_events, :any, default: nil
attr :search_op_log_aggregates, :any, default: nil
attr :has_results?, :boolean
attr :source, Logflare.Sources.Source, required: true
diff --git a/lib/logflare_web/live/search_live/log_event_components.ex b/lib/logflare_web/live/search_live/log_event_components.ex
index b6da8d20e3..21d5f8341b 100644
--- a/lib/logflare_web/live/search_live/log_event_components.ex
+++ b/lib/logflare_web/live/search_live/log_event_components.ex
@@ -21,22 +21,47 @@ defmodule LogflareWeb.SearchLive.LogEventComponents do
@default_empty_event_message "(empty event message)"
attr :search_op_log_events, :map, default: nil
+ attr :search_op_log_aggregates, :map, default: nil
+ attr :log_events, :any, default: []
attr :last_query_completed_at, :any, default: nil
attr :loading, :boolean, required: true
+ attr :pagination_available?, :boolean, default: false
+ attr :unbounded_pagination?, :boolean, default: false
+ attr :event_page_loading, :atom, default: nil
+ attr :next_events_exhausted?, :boolean, default: false
attr :search_timezone, :string, required: true
- attr :tailing?, :boolean, required: true
- attr :querystring, :string, required: true
attr :empty_event_message_placeholder, :string, default: @default_empty_event_message
attr :source_schema_flat_map, :map, default: %{}
attr :search_op, Logflare.Logs.SearchOperation
def results_list(assigns) do
- assigns = assign(assigns, :select_fields, build_select_fields(assigns.search_op))
+ assigns =
+ assigns
+ |> assign(:select_fields, build_select_fields(assigns.search_op))
+ |> assign(:event_page_loading?, not is_nil(assigns.event_page_loading))
+ |> assign(
+ :top_pagination_available?,
+ assigns.pagination_available? or
+ (assigns.unbounded_pagination? and
+ match?(%{has_more_events?: true}, assigns.search_op_log_events))
+ )
~H"""
-
-
- <.log_event :for={log <- @search_op_log_events.rows} timezone={@search_timezone} log_event={log} select_fields={build_select_fields(@search_op)} source_schema_flat_map={@source_schema_flat_map}>
+
+ <.load_more_button
+ :if={@top_pagination_available?}
+ id="load-more-events-top"
+ intent={if(@pagination_available? and @search_op_log_events.has_more_events?, do: "within_range", else: "extend_previous")}
+ load_enabled={not @loading and is_nil(@event_page_loading)}
+ click={
+ JS.dispatch("logflare:before-log-prepend", to: "#source-logs-search-list")
+ |> JS.push("load_events")
+ }
+ class="tw-my-2"
+ />
+
+ <.empty_result_list :if={not @loading} search_op_log_events={@search_op_log_events} search_op_log_aggregates={@search_op_log_aggregates} />
+ <.log_event :for={{dom_id, log} <- @log_events} id={dom_id} data-event-id={event_id(log)} data-event-timestamp={log.body["timestamp"]} timezone={@search_timezone} log_event={log} select_fields={build_select_fields(@search_op)} source_schema_flat_map={@source_schema_flat_map}>
{log.body["event_message"]}
<:actions phx-no-format>
@@ -47,24 +72,18 @@ defmodule LogflareWeb.SearchLive.LogEventComponents do
title="Log Event"
phx-value-log-event-id={log.id}
phx-value-log-event-timestamp={log.body["timestamp"]}
- phx-value-lql={@querystring}
- phx-value-tailing?={@tailing?}
- phx-value-tz={@search_timezone}
>
view
<.modal_link
component={LogflareWeb.SearchLive.EventContextComponent}
- click={JS.push("soft_pause")}
- close={if(@tailing?, do: JS.push("soft_play", target: "#source-logs-search-control") |> JS.push("close"), else: nil)}
+ click={JS.push("open_event_context")}
+ close={JS.push("close_event_context", target: "#source-logs-search-control") |> JS.push("close")}
class="tw-text-[0.65rem]"
modal_id={:log_event_context_viewer}
title="View Event Context"
phx-value-log-event-id={log.id}
- phx-value-source-id={@search_op.source.id}
phx-value-log-event-timestamp={log.body["timestamp"]}
- phx-value-timezone={@search_timezone}
- phx-value-querystring={@querystring}
>
context
@@ -87,6 +106,24 @@ defmodule LogflareWeb.SearchLive.LogEventComponents do
+ <.load_more_button :if={@pagination_available? or @unbounded_pagination?} id="load-more-events-bottom" intent="extend_next" load_enabled={not @loading and is_nil(@event_page_loading) and not @next_events_exhausted?} disabled={@loading or @event_page_loading? or @next_events_exhausted?} />
+
+ """
+ end
+
+ attr :id, :string, required: true
+ attr :intent, :string, values: ~w(within_range extend_previous extend_next), required: true
+ attr :class, :any, default: nil
+ attr :load_enabled, :boolean, required: true
+ attr :disabled, :boolean, default: false
+ attr :click, :any, default: "load_events"
+
+ def load_more_button(assigns) do
+ ~H"""
+
+
"""
end
@@ -283,7 +320,7 @@ defmodule LogflareWeb.SearchLive.LogEventComponents do
)
~H"""
-
+
No events matching your query
@@ -302,10 +339,7 @@ defmodule LogflareWeb.SearchLive.LogEventComponents do
"""
end
- defp show_empty_results?(%{rows: rows})
- when is_list(rows), do: Enum.empty?(rows)
-
- defp show_empty_results?(_), do: false
+ defp event_id(%Logflare.LogEvent{id: id, body: body}), do: id || body["id"]
def extended_search_lql(datetime) do
new_rule =
diff --git a/lib/logflare_web/live/search_live/logs_search_lv.ex b/lib/logflare_web/live/search_live/logs_search_lv.ex
index 611222faaa..b84cdc0b4f 100644
--- a/lib/logflare_web/live/search_live/logs_search_lv.ex
+++ b/lib/logflare_web/live/search_live/logs_search_lv.ex
@@ -11,6 +11,7 @@ defmodule LogflareWeb.Source.SearchLV do
alias Logflare.Backends.QueryError
alias Logflare.Billing
+ alias Logflare.Logs.SearchOperation
alias Logflare.Logs.SearchQueryExecutor
alias Logflare.Logs.SearchOperations
alias Logflare.Logs.SearchUtils
@@ -35,6 +36,7 @@ defmodule LogflareWeb.Source.SearchLV do
require Logger
+ @log_event_stream_limit 5_000
@tail_search_interval 1000
@user_idle_interval :timer.minutes(2)
@@ -88,10 +90,15 @@ defmodule LogflareWeb.Source.SearchLV do
# loading states
loading: true,
chart_loading: true,
+ event_page_loading: nil,
+ event_page_cursors: %{previous: nil, next: nil},
+ next_events_exhausted?: false,
+ range_extension_patch: nil,
# tailing states
tailing_initial?: true,
tailing_timer: nil,
tailing?: tailing?,
+ resume_tailing_after_context?: false,
# search states
search_op: nil,
search_op_error: nil,
@@ -106,6 +113,8 @@ defmodule LogflareWeb.Source.SearchLV do
saved_searches: saved_searches(source),
force_query: Map.get(params, "force", "false") == "true"
)
+ |> stream_configure(:log_events, dom_id: &log_event_dom_id/1)
+ |> stream(:log_events, [])
|> maybe_assign_user_timezone(team_user, user)
end
@@ -151,6 +160,37 @@ defmodule LogflareWeb.Source.SearchLV do
{:noreply, push_patch(socket, to: path, replace: true)}
end
+ def handle_params(
+ %{"querystring" => qs} = params,
+ uri,
+ %{
+ assigns: %{
+ range_extension_patch: %{
+ querystring: expected_querystring,
+ rows: rows,
+ direction: direction
+ }
+ }
+ } = socket
+ )
+ when is_binary(expected_querystring) do
+ if qs == expected_querystring do
+ {:noreply,
+ socket
+ |> put_event_page(rows, direction)
+ |> assign(:event_page_loading, nil)
+ |> assign(:range_extension_patch, nil)
+ |> assign(uri: URI.parse(uri), uri_params: params, querystring: qs)}
+ else
+ socket =
+ socket
+ |> assign(:event_page_loading, nil)
+ |> assign(:range_extension_patch, nil)
+
+ handle_params(params, uri, socket)
+ end
+ end
+
def handle_params(%{"querystring" => qs} = params, uri, socket) do
source = socket.assigns.source
@@ -180,23 +220,14 @@ defmodule LogflareWeb.Source.SearchLV do
{:ok, socket} <- check_suggested_keys(lql_rules, source, socket) do
qs = Lql.encode!(lql_rules)
- search_op_log_events =
- if socket.assigns.search_op_log_events do
- rows = socket.assigns.search_op_log_events.rows
- events = Enum.map(rows, &Map.put(&1, :is_from_stale_query, true))
- Map.put(socket.assigns.search_op_log_events, :rows, events)
- else
- socket.assigns.search_op_log_events
- end
-
socket =
socket
|> assign(:loading, true)
|> assign(:chart_loading, true)
+ |> reset_event_pagination()
|> assign(:tailing_initial?, true)
|> assign(:lql_rules, lql_rules)
|> assign(:querystring, qs)
- |> assign(:search_op_log_events, search_op_log_events)
if connected?(socket) do
kickoff_queries(source.token, socket.assigns)
@@ -245,6 +276,9 @@ defmodule LogflareWeb.Source.SearchLV do
search_op_error: @search_op_error,
team_user: @team_user,
team: @team,
+ lql: @querystring,
+ querystring: @querystring,
+ search_timezone: @search_timezone,
close: @modal.body[:close],
return_to: @modal.body.return_to
)}
@@ -257,15 +291,18 @@ defmodule LogflareWeb.Source.SearchLV do
-
@@ -376,6 +413,31 @@ defmodule LogflareWeb.Source.SearchLV do
{:noreply, socket}
end
+ def handle_event(
+ "load_events",
+ %{"intent" => intent},
+ %{assigns: %{loading: false, event_page_loading: nil, tailing?: false}} = socket
+ ) do
+ with {:ok, intent} <- event_page_intent(intent),
+ {:ok, cursor} <- event_page_cursor(socket.assigns, intent),
+ :ok <- event_page_available(socket.assigns, intent),
+ :ok <-
+ SearchQueryExecutor.query_page(
+ socket.assigns.executor_pid,
+ socket.assigns,
+ intent,
+ cursor
+ ) do
+ maybe_cancel_tailing_timer(socket)
+
+ {:noreply, assign(socket, :event_page_loading, intent)}
+ else
+ _ -> {:noreply, socket}
+ end
+ end
+
+ def handle_event("load_events", _params, socket), do: {:noreply, socket}
+
def handle_event(direction, _, socket) when direction in ["backwards", "forwards"] do
rules = socket.assigns.lql_rules
@@ -417,6 +479,28 @@ defmodule LogflareWeb.Source.SearchLV do
soft_pause(ev, socket)
end
+ def handle_event("open_event_context", _, socket) do
+ resume_tailing? = socket.assigns.tailing?
+
+ socket =
+ socket
+ |> assign(:resume_tailing_after_context?, resume_tailing?)
+ |> pause_tailing()
+
+ {:noreply, socket}
+ end
+
+ def handle_event("close_event_context", _, socket) do
+ socket =
+ if socket.assigns.resume_tailing_after_context? do
+ resume_tailing(socket)
+ else
+ socket
+ end
+
+ {:noreply, assign(socket, :resume_tailing_after_context?, false)}
+ end
+
def handle_event("hard_play" = ev, _, socket) do
hard_play(ev, socket)
end
@@ -597,6 +681,7 @@ defmodule LogflareWeb.Source.SearchLV do
|> assign(:lql_rules, lql_rules)
|> assign(:loading, true)
|> assign(:chart_loading, true)
+ |> reset_event_pagination()
|> clear_flash()
|> push_patch_with_params(%{querystring: qs, tailing?: socket.assigns.tailing?})
else
@@ -604,6 +689,203 @@ defmodule LogflareWeb.Source.SearchLV do
end
end
+ defp reset_event_pagination(socket) do
+ socket
+ |> assign(:event_page_loading, nil)
+ |> assign(:event_page_cursors, %{previous: nil, next: nil})
+ |> assign(:range_extension_patch, nil)
+ |> assign(:next_events_exhausted?, false)
+ end
+
+ defp event_page_intent("within_range"), do: {:ok, :within_range}
+ defp event_page_intent("extend_previous"), do: {:ok, :extend_previous}
+ defp event_page_intent("extend_next"), do: {:ok, :extend_next}
+ defp event_page_intent(_intent), do: :error
+
+ @spec event_page_cursor(map(), SearchOperation.event_page_intent()) ::
+ {:ok, SearchOperation.event_cursor() | nil}
+ defp event_page_cursor(%{event_page_cursors: cursors}, :extend_next),
+ do: {:ok, cursors.next}
+
+ defp event_page_cursor(%{event_page_cursors: cursors}, intent)
+ when intent in [:within_range, :extend_previous],
+ do: {:ok, cursors.previous}
+
+ defp event_page_available(
+ %{search_op_log_events: %{has_more_events?: true}},
+ :within_range
+ ),
+ do: :ok
+
+ defp event_page_available(
+ %{search_op_log_events: %{has_more_events?: false}},
+ :extend_previous
+ ),
+ do: :ok
+
+ defp event_page_available(
+ %{
+ tailing?: tailing?,
+ lql_rules: rules,
+ search_op_log_events: %{has_more_events?: true}
+ },
+ :extend_previous
+ ) do
+ if unbounded_event_pagination?(tailing?, rules), do: :ok, else: :error
+ end
+
+ defp event_page_available(%{next_events_exhausted?: false}, :extend_next), do: :ok
+ defp event_page_available(_assigns, _intent), do: :error
+
+ defp bounded_event_pagination?(true, _rules), do: false
+
+ defp bounded_event_pagination?(false, rules) do
+ rules
+ |> Rules.get_timestamp_filters()
+ |> bounded_timestamp_filters?()
+ end
+
+ defp unbounded_event_pagination?(true, _rules), do: false
+
+ defp unbounded_event_pagination?(false, rules),
+ do: not bounded_event_pagination?(false, rules)
+
+ defp bounded_timestamp_filters?(filters) do
+ Enum.any?(filters, &(&1.operator == :range and length(&1.values || []) == 2)) or
+ (Enum.any?(filters, &(&1.operator in [:>, :>=])) and
+ Enum.any?(filters, &(&1.operator in [:<, :<=])))
+ end
+
+ defp put_event_page(socket, rows, :previous) do
+ socket
+ |> stream(:log_events, rows, at: -1)
+ |> put_event_page_cursor(:previous, List.last(rows))
+ end
+
+ defp put_event_page(socket, rows, :next) do
+ socket
+ |> stream(:log_events, Enum.reverse(rows), at: 0)
+ |> put_event_page_cursor(:next, List.first(rows))
+ end
+
+ @spec put_event_page_cursor(
+ Phoenix.LiveView.Socket.t(),
+ SearchOperation.event_page_direction(),
+ Logflare.LogEvent.t() | nil
+ ) :: Phoenix.LiveView.Socket.t()
+ defp put_event_page_cursor(socket, _direction, nil), do: socket
+
+ defp put_event_page_cursor(socket, direction, event) do
+ cursors = Map.put(socket.assigns.event_page_cursors, direction, event_cursor(event))
+ assign(socket, :event_page_cursors, cursors)
+ end
+
+ @spec event_cursor(Logflare.LogEvent.t() | nil) :: SearchOperation.event_cursor() | nil
+ defp event_cursor(nil), do: nil
+
+ defp event_cursor(%{id: id, body: body}) do
+ %{id: id || body["id"], timestamp: body["timestamp"]}
+ end
+
+ defp put_search_events(socket, [])
+ when socket.assigns.tailing? and not socket.assigns.tailing_initial?,
+ do: socket
+
+ defp put_search_events(socket, rows)
+ when socket.assigns.tailing? and not socket.assigns.tailing_initial? do
+ rows
+ |> Enum.with_index()
+ |> Enum.reduce(socket, fn {row, index}, socket ->
+ stream_insert(socket, :log_events, row, at: index, limit: @log_event_stream_limit)
+ end)
+ end
+
+ defp put_search_events(socket, rows) do
+ socket
+ |> stream(:log_events, rows, reset: true)
+ |> assign(:event_page_cursors, %{
+ previous: rows |> List.last() |> event_cursor(),
+ next: rows |> List.first() |> event_cursor()
+ })
+ end
+
+ defp event_search_metadata(events_op, has_more_events? \\ nil) do
+ has_more_events? =
+ if is_nil(has_more_events?), do: events_op.has_more_events?, else: has_more_events?
+
+ %{events_op | rows: [], has_more_events?: has_more_events?}
+ end
+
+ defp apply_event_page_result(socket, events_op, direction) do
+ has_more_events? =
+ case direction do
+ :previous -> events_op.has_more_events?
+ :next -> socket.assigns.search_op_log_events.has_more_events?
+ end
+
+ events_metadata = event_search_metadata(events_op, has_more_events?)
+
+ socket
+ |> put_event_page(events_op.rows, direction)
+ |> assign(:search_op, events_metadata)
+ |> assign(:search_op_error, nil)
+ |> assign(:search_op_log_events, events_metadata)
+ |> assign(:event_page_loading, nil)
+ |> assign(:last_query_completed_at, DateTime.utc_now())
+ end
+
+ defp apply_range_extension_result(socket, events_op) do
+ %{request: request, cursor: cursor, has_more?: has_more?} = events_op.event_page_result
+ direction = SearchOperation.event_page_direction(request.intent)
+ has_more_events? = socket.assigns.search_op_log_events.has_more_events?
+ events_metadata = event_search_metadata(events_op, has_more_events?)
+
+ socket =
+ socket
+ |> assign(:search_op, events_metadata)
+ |> assign(:search_op_error, nil)
+ |> assign(:search_op_log_events, events_metadata)
+ |> assign(:last_query_completed_at, DateTime.utc_now())
+ |> put_extension_state(direction, not has_more?)
+
+ if is_nil(cursor) do
+ socket
+ |> put_event_page(events_op.rows, direction)
+ |> assign(:event_page_loading, nil)
+ else
+ lql_rules =
+ socket.assigns.lql_rules
+ |> Rules.extend_timestamp_range(
+ direction,
+ cursor.timestamp,
+ socket.assigns.search_timezone
+ )
+ |> maybe_adjust_chart_period()
+
+ querystring = Lql.encode!(lql_rules)
+
+ socket =
+ socket
+ |> assign(:lql_rules, lql_rules)
+ |> assign(:querystring, querystring)
+ |> assign(:chart_loading, true)
+ |> assign(:range_extension_patch, %{
+ querystring: querystring,
+ rows: events_op.rows,
+ direction: direction
+ })
+
+ SearchQueryExecutor.query_agg(socket.assigns.executor_pid, socket.assigns)
+
+ push_patch_with_params(socket, %{querystring: querystring, tailing?: false})
+ end
+ end
+
+ defp put_extension_state(socket, :previous, _exhausted?), do: socket
+
+ defp put_extension_state(socket, :next, exhausted?),
+ do: assign(socket, :next_events_exhausted?, exhausted?)
+
def handle_info(:soft_pause = ev, socket) do
soft_pause(ev, socket)
end
@@ -645,38 +927,87 @@ defmodule LogflareWeb.Source.SearchLV do
{:noreply, socket}
end
- def handle_info({:search_result, %{events: events_op} = search_result}, socket) do
+ def handle_info(
+ {:search_result,
+ %{events: %{event_page_result: %{request: %{intent: :within_range}}} = events_op}},
+ %{assigns: %{event_page_loading: :within_range}} = socket
+ ) do
+ {:noreply, apply_event_page_result(socket, events_op, :previous)}
+ end
+
+ def handle_info(
+ {:search_result,
+ %{
+ events:
+ %{
+ event_page_result: %{
+ request: %{intent: intent, boundary: nil}
+ }
+ } = events_op
+ }},
+ %{assigns: %{event_page_loading: intent}} = socket
+ )
+ when intent in [:extend_previous, :extend_next] do
+ direction = SearchOperation.event_page_direction(intent)
+ {:noreply, apply_event_page_result(socket, events_op, direction)}
+ end
+
+ def handle_info(
+ {:search_result,
+ %{events: %{event_page_result: %{request: %{intent: intent}}} = events_op}},
+ %{assigns: %{event_page_loading: intent}} = socket
+ )
+ when intent in [:extend_previous, :extend_next] do
+ {:noreply, apply_range_extension_result(socket, events_op)}
+ end
+
+ def handle_info(
+ {:search_result, %{events: %{event_page_result: %{}}}},
+ socket
+ ),
+ do: {:noreply, socket}
+
+ def handle_info(
+ {:search_result, %{events: %{event_page_result: nil} = events_op} = search_result},
+ socket
+ ) do
tailing_timer =
if socket.assigns.tailing? do
Process.send_after(self(), :schedule_tail_search, @tail_search_interval)
end
+ events_metadata = event_search_metadata(events_op)
+
socket =
socket
- |> assign(:search_op, events_op)
+ |> reset_event_pagination()
+ |> put_search_events(events_op.rows)
+ |> assign(:search_op, events_metadata)
|> assign(:search_op_error, nil)
- |> assign(:search_op_log_events, search_result.events)
+ |> assign(:search_op_log_events, events_metadata)
|> assign(:tailing_timer, tailing_timer)
|> assign(:loading, false)
|> assign(:tailing_initial?, false)
|> assign(:last_query_completed_at, DateTime.utc_now())
socket =
- cond do
- match?({:warning, _}, search_result.events.status) ->
- {:warning, message} = search_result.events.status
- put_flash(socket, :info, message)
-
- msg = warning_message(socket.assigns, search_result) ->
- put_flash(socket, :warning, msg)
-
- true ->
- socket
+ if match?({:warning, _}, search_result.events.status) do
+ {:warning, message} = search_result.events.status
+ put_flash(socket, :info, message)
+ else
+ socket
end
{:noreply, socket}
end
+ def handle_info(
+ {:search_error, %{event_page_request: %{intent: intent}}},
+ %{assigns: %{event_page_loading: active_intent}} = socket
+ )
+ when intent != active_intent,
+ do: {:noreply, socket}
+
def handle_info({:search_error, search_op}, socket) do
socket =
case search_op.error do
@@ -686,6 +1017,7 @@ defmodule LogflareWeb.Source.SearchLV do
socket
|> assign(loading: false)
|> assign(chart_loading: false)
+ |> assign(event_page_loading: nil)
|> put_halt_flash_message(search_op)
err ->
@@ -694,6 +1026,7 @@ defmodule LogflareWeb.Source.SearchLV do
socket
|> assign(loading: false)
|> assign(chart_loading: false)
+ |> assign(event_page_loading: nil)
|> put_flash_query_error(err)
end
@@ -739,6 +1072,7 @@ defmodule LogflareWeb.Source.SearchLV do
socket
|> assign(:loading, true)
+ |> reset_event_pagination()
|> assign(:tailing_initial?, true)
|> clear_flash()
|> assign(:lql_rules, lql_rules)
@@ -827,26 +1161,6 @@ defmodule LogflareWeb.Source.SearchLV do
push_patch(socket, to: path, replace: false)
end
- defp warning_message(assigns, search_op) do
- tailing? = assigns.tailing?
- querystring = assigns.querystring
- log_events_empty? = Enum.empty?(search_op.events.rows)
-
- cond do
- log_events_empty? and not tailing? ->
- "No log events matching your search query."
-
- log_events_empty? and tailing? ->
- "No log events matching your search query."
-
- querystring == "" and log_events_empty? and tailing? ->
- "No log events ingested during last 24 hours. Try searching over a longer time period, and clicking the bar chart to drill down."
-
- true ->
- nil
- end
- end
-
defp adjust_timestamp_rules(timestamp_rules, search_timezone) do
case Timex.Timezone.get(search_timezone) do
{:error, _} -> timestamp_rules
@@ -979,6 +1293,7 @@ defmodule LogflareWeb.Source.SearchLV do
|> assign(:tailing?, false)
|> assign(:loading, false)
|> assign(:chart_loading, false)
+ |> assign(:event_page_loading, nil)
|> put_flash(:error, error)
end
@@ -993,17 +1308,7 @@ defmodule LogflareWeb.Source.SearchLV do
{:noreply, error_socket(socket, "Tailing is disabled for this source")}
end
- defp soft_play(_ev, %{assigns: prev_assigns} = socket) do
- %{source: %{token: stoken} = _source} = prev_assigns
-
- kickoff_queries(stoken, socket.assigns)
-
- socket =
- socket
- |> assign(:tailing?, true)
-
- {:noreply, socket}
- end
+ defp soft_play(_ev, socket), do: {:noreply, resume_tailing(socket)}
defp soft_pause(
_ev,
@@ -1012,15 +1317,25 @@ defmodule LogflareWeb.Source.SearchLV do
{:noreply, socket}
end
- defp soft_pause(_ev, %{assigns: %{source: _source, executor_pid: executor_pid}} = socket) do
+ defp soft_pause(_ev, socket), do: {:noreply, pause_tailing(socket)}
+
+ defp pause_tailing(%{assigns: %{tailing?: false}} = socket), do: socket
+
+ defp pause_tailing(%{assigns: %{executor_pid: executor_pid}} = socket) do
maybe_cancel_tailing_timer(socket)
SearchQueryExecutor.cancel_query(executor_pid)
- socket =
- socket
- |> assign(:tailing?, false)
+ socket
+ |> assign(:tailing?, false)
+ |> reset_event_pagination()
+ end
- {:noreply, socket}
+ defp resume_tailing(socket) do
+ kickoff_queries(socket.assigns.source.token, socket.assigns)
+
+ socket
+ |> assign(:tailing?, true)
+ |> reset_event_pagination()
end
defp hard_play(
@@ -1038,6 +1353,7 @@ defmodule LogflareWeb.Source.SearchLV do
socket =
socket
|> assign(:tailing?, true)
+ |> reset_event_pagination()
|> push_patch_with_params(%{
querystring: prev_assigns.querystring,
tailing?: true
@@ -1134,6 +1450,11 @@ defmodule LogflareWeb.Source.SearchLV do
)
end
+ @spec log_event_dom_id(Logflare.LogEvent.t()) :: String.t()
+ defp log_event_dom_id(%Logflare.LogEvent{id: id, body: %{"timestamp" => timestamp}}) do
+ "log-events-#{id}-#{timestamp}"
+ end
+
@spec querystring_or_default(String.t(), Logflare.Sources.Source.t()) :: String.t()
defp querystring_or_default("", source), do: source.default_search_lql || ""
defp querystring_or_default(qs, _source), do: qs
diff --git a/lib/logflare_web/templates/log/log_event.html.heex b/lib/logflare_web/templates/log/log_event.html.heex
index cc0d089876..127ccb8676 100644
--- a/lib/logflare_web/templates/log/log_event.html.heex
+++ b/lib/logflare_web/templates/log/log_event.html.heex
@@ -15,18 +15,17 @@
module={LogflareWeb.Search.LogEventViewerComponent}
id={:log_event_viewer}
{%{
- user: @user,
- source_schema_flat_map: @source_schema_flat_map,
- source: @source,
- timestamp: @timestamp,
- log_event: @log_event,
- params: %{
- "log-event-id" => @log_event_id,
- "log-event-timestamp" => @timestamp,
- "lql" => assigns[:lql],
- "tailing?" => assigns[:tailing?],
- "tz" => assigns[:tz]
- }
- }}
+ user: @user,
+ source_schema_flat_map: @source_schema_flat_map,
+ source: @source,
+ timestamp: @timestamp,
+ log_event: @log_event,
+ lql: @lql,
+ search_timezone: @search_timezone,
+ params: %{
+ "log-event-id" => @log_event_id,
+ "log-event-timestamp" => @timestamp
+ }
+ }}
/>
diff --git a/test/logflare/log_event_test.exs b/test/logflare/log_event_test.exs
index 0ded31871e..8cbf86591b 100644
--- a/test/logflare/log_event_test.exs
+++ b/test/logflare/log_event_test.exs
@@ -22,7 +22,6 @@ defmodule Logflare.LogEventTest do
drop: false,
id: id,
ingested_at: _,
- is_from_stale_query: nil,
source_id: source_id,
valid: true,
pipeline_error: nil,
@@ -655,7 +654,6 @@ defmodule Logflare.LogEventTest do
drop: false,
id: id,
ingested_at: _,
- is_from_stale_query: nil,
valid: true,
pipeline_error: nil,
via_rule_id: nil
diff --git a/test/logflare/lql/rules_test.exs b/test/logflare/lql/rules_test.exs
index f98bd01cbf..e224a184e7 100644
--- a/test/logflare/lql/rules_test.exs
+++ b/test/logflare/lql/rules_test.exs
@@ -565,6 +565,112 @@ defmodule Logflare.Lql.RulesTest do
end
end
+ describe "extend_timestamp_range/3" do
+ test "replaces a concrete timestamp range and changes only the previous edge" do
+ timestamp_filter = %FilterRule{
+ path: "timestamp",
+ operator: :range,
+ values: [~N[2026-07-13 10:00:00], ~N[2026-07-13 11:00:00]]
+ }
+
+ extra_timestamp_filter = %FilterRule{
+ path: "timestamp",
+ operator: :<,
+ value: ~N[2026-07-13 12:00:00]
+ }
+
+ message_filter = %FilterRule{path: "event_message", operator: :=, value: "error"}
+ chart_rule = %ChartRule{aggregate: :count, path: "timestamp", period: :minute}
+ select_rule = %SelectRule{path: "metadata.request_id", wildcard: false}
+
+ lql_rules = [
+ timestamp_filter,
+ message_filter,
+ chart_rule,
+ extra_timestamp_filter,
+ select_rule
+ ]
+
+ result =
+ Rules.extend_timestamp_range(
+ lql_rules,
+ :previous,
+ DateTime.to_unix(~U[2026-07-13 09:30:00Z], :microsecond)
+ )
+
+ assert Rules.get_timestamp_filters(result) == [
+ %FilterRule{
+ path: "timestamp",
+ operator: :range,
+ values: [~N[2026-07-13 09:30:00.000000], ~N[2026-07-13 12:00:00]],
+ modifiers: %{}
+ }
+ ]
+
+ assert Rules.get_metadata_and_message_filters(result) == [message_filter]
+ assert Rules.get_chart_rule(result) == chart_rule
+ assert Rules.get_select_rules(result) == [select_rule]
+ end
+
+ test "materializes shorthand and changes only the next edge" do
+ shorthand_filter = %FilterRule{
+ path: "timestamp",
+ operator: :range,
+ values: [~U[2026-07-13 10:00:00Z], ~U[2026-07-13 11:00:00Z]],
+ shorthand: "this@day"
+ }
+
+ metadata_filter = %FilterRule{path: "metadata.level", operator: :=, value: "info"}
+ chart_rule = %ChartRule{aggregate: :count, path: "timestamp", period: :hour}
+
+ result =
+ Rules.extend_timestamp_range(
+ [shorthand_filter, metadata_filter, chart_rule],
+ :next,
+ ~U[2026-07-13 11:30:00Z]
+ )
+
+ assert Rules.get_timestamp_filters(result) == [
+ %FilterRule{
+ path: "timestamp",
+ operator: :range,
+ values: [~N[2026-07-13 10:00:00], ~N[2026-07-13 11:30:00]],
+ modifiers: %{}
+ }
+ ]
+
+ assert Rules.get_metadata_and_message_filters(result) == [metadata_filter]
+ assert Rules.get_chart_rule(result) == chart_rule
+ end
+
+ test "converts an event cursor into the search display timezone" do
+ timestamp_filter = %FilterRule{
+ path: "timestamp",
+ operator: :range,
+ values: [~N[2026-07-13 09:31:50], ~N[2026-07-13 09:32:00]]
+ }
+
+ event_timestamp = DateTime.to_unix(~U[2026-07-13 01:32:07Z], :microsecond)
+
+ result =
+ Rules.extend_timestamp_range(
+ [timestamp_filter],
+ :next,
+ event_timestamp,
+ "Asia/Singapore"
+ )
+
+ assert Rules.get_timestamp_filters(result) == [
+ %FilterRule{
+ path: "timestamp",
+ operator: :range,
+ values: [~N[2026-07-13 09:31:50], ~N[2026-07-13 09:32:07.000000]],
+ modifiers: %{}
+ }
+ ]
+ end
+ end
+
describe "jump_timestamp/2" do
test "creates new timestamp range by jumping forwards" do
timestamp_filter = %FilterRule{
diff --git a/test/logflare_web/live/log_event_live_test.exs b/test/logflare_web/live/log_event_live_test.exs
index 67426d6a4e..d7ed57d5f6 100644
--- a/test/logflare_web/live/log_event_live_test.exs
+++ b/test/logflare_web/live/log_event_live_test.exs
@@ -8,7 +8,7 @@ defmodule LogflareWeb.LogEventLiveTest do
setup %{conn: conn} do
insert(:plan)
- user = insert(:user)
+ user = insert(:user, preferences: build(:user_preferences, timezone: "Singapore"))
source = insert(:source, user: user)
insert(:source_schema, source: source)
conn = login_user(conn, user)
@@ -33,12 +33,14 @@ defmodule LogflareWeb.LogEventLiveTest do
TestUtils.gen_bq_response([%{"id" => le.id, "event_message" => le.body["event_message"]}])}
end)
- {:ok, _view, _html} =
+ {:ok, view, _html} =
live(
conn,
~p"/sources/#{source.id}/event?#{%{timestamp: "2024-01-10T20:13:03Z", uuid: le.id}}"
)
+ assert view |> element("a", "inspect") |> render() =~ "tz=Singapore"
+
TestUtils.retry_assert(fn ->
assert_receive {:query, ^ref, body}
assert Enum.any?(body.queryParameters, &(&1.parameterValue.value =~ "2024-01-"))
diff --git a/test/logflare_web/live/search_live/log_event_components_test.exs b/test/logflare_web/live/search_live/log_event_components_test.exs
index 9cdf5bb41c..1186994727 100644
--- a/test/logflare_web/live/search_live/log_event_components_test.exs
+++ b/test/logflare_web/live/search_live/log_event_components_test.exs
@@ -13,6 +13,7 @@ defmodule LogflareWeb.SearchLive.LogEventComponentsTest do
@default_attrs %{
search_op_log_events: nil,
+ log_events: [],
last_query_completed_at: nil,
loading: false,
search_timezone: "Etc/UTC",
@@ -23,39 +24,6 @@ defmodule LogflareWeb.SearchLive.LogEventComponentsTest do
search_op: nil
}
- defmodule TestLive do
- use LogflareWeb, :live_view
-
- def render(assigns) do
- ~H"""
-
-
-
- """
- end
-
- def mount(_params, session, socket) do
- {:ok,
- assign(socket,
- search_op_log_events: session["search_op_log_events"],
- search_op: session["search_op"],
- last_query_completed_at: session["last_query_completed_at"],
- loading: session["loading"],
- search_timezone: session["search_timezone"],
- tailing?: session["tailing?"],
- querystring: session["querystring"]
- )}
- end
- end
-
describe "results_list/1" do
setup do
user = insert(:user)
@@ -95,7 +63,8 @@ defmodule LogflareWeb.SearchLive.LogEventComponentsTest do
render_component(&LogEventComponents.results_list/1, %{
@default_attrs
| search_op: %{source: source, lql_rules: lql_rules, search_timezone: "Etc/UTC"},
- search_op_log_events: search_op_log_events
+ search_op_log_events: search_op_log_events,
+ log_events: stream_entries(search_op_log_events.rows)
})
assert html =~ "Log message 1"
@@ -110,7 +79,7 @@ defmodule LogflareWeb.SearchLive.LogEventComponentsTest do
search_op_log_events: %{rows: []}
})
- assert html =~ ~r|id="logs-list" class="(.*)blurred"|
+ assert html =~ ~r|id="logs-list".*class="(.*)blurred"|
end
test "renders empty state when no log events", %{source: source, lql_rules: lql_rules} do
@@ -147,7 +116,8 @@ defmodule LogflareWeb.SearchLive.LogEventComponentsTest do
render_component(&LogEventComponents.results_list/1, %{
@default_attrs
| search_op: %{source: source, lql_rules: lql_rules, search_timezone: "Etc/UTC"},
- search_op_log_events: search_op_log_events
+ search_op_log_events: search_op_log_events,
+ log_events: stream_entries(search_op_log_events.rows)
})
assert html =~ "(empty event message)"
@@ -183,7 +153,8 @@ defmodule LogflareWeb.SearchLive.LogEventComponentsTest do
render_component(&LogEventComponents.results_list/1, %{
@default_attrs
| search_op: %{source: source, lql_rules: lql_rules, search_timezone: "Etc/UTC"},
- search_op_log_events: search_op_log_events
+ search_op_log_events: search_op_log_events,
+ log_events: stream_entries(search_op_log_events.rows)
})
assert html =~ "Normal log message"
@@ -339,4 +310,10 @@ defmodule LogflareWeb.SearchLive.LogEventComponentsTest do
"""
end
end
+
+ defp stream_entries(log_events) do
+ log_events
+ |> Enum.with_index()
+ |> Enum.map(fn {log_event, index} -> {"log-events-#{index}", log_event} end)
+ end
end
diff --git a/test/logflare_web/live/search_live/logs_search_lv_test.exs b/test/logflare_web/live/search_live/logs_search_lv_test.exs
index 3c8ec9cffe..0f5a19cdb1 100644
--- a/test/logflare_web/live/search_live/logs_search_lv_test.exs
+++ b/test/logflare_web/live/search_live/logs_search_lv_test.exs
@@ -1052,7 +1052,6 @@ defmodule LogflareWeb.Source.SearchLVTest do
|> render()
assert link =~ ~r/phx-value-log-event-timestamp="\d+/
- assert link =~ ~r/phx-value-lql="\w+/
end
@tag source_schema:
@@ -1563,6 +1562,43 @@ defmodule LogflareWeb.Source.SearchLVTest do
assert get_view_assigns(view).tailing?
end
+ test "closing context does not resume a paused search", %{
+ conn: conn,
+ source: source
+ } do
+ {:ok, view, _html} = live_with_redirect(conn, ~p"/sources/#{source.id}/search")
+
+ view |> TestUtils.wait_for_render("#logs-list li:first-of-type a")
+ assert get_view_assigns(view).tailing?
+
+ render_click(view, "soft_pause", %{})
+ refute get_view_assigns(view).tailing?
+
+ render_click(view, "open_event_context", %{})
+
+ render_click(view, "close_event_context", %{})
+
+ refute get_view_assigns(view).tailing?
+ end
+
+ test "closing context resumes a search that was live", %{
+ conn: conn,
+ source: source
+ } do
+ {:ok, view, _html} = live_with_redirect(conn, ~p"/sources/#{source.id}/search")
+
+ view |> TestUtils.wait_for_render("#logs-list li:first-of-type a")
+ assert get_view_assigns(view).tailing?
+
+ render_click(view, "open_event_context", %{})
+
+ refute get_view_assigns(view).tailing?
+
+ render_click(view, "close_event_context", %{})
+
+ assert get_view_assigns(view).tailing?
+ end
+
test "datetime_update", %{conn: conn, source: source} do
{:ok, view, _html} =
live_with_redirect(conn, Routes.live_path(conn, SearchLV, source, querystring: "error"))
@@ -1764,7 +1800,7 @@ defmodule LogflareWeb.Source.SearchLVTest do
end
end
- describe "single tenant searching with postgres backend" do
+ describe "search tasks with postgres backend" do
TestUtils.setup_single_tenant(seed_user: true, backend_type: :postgres)
setup do
@@ -1823,6 +1859,113 @@ defmodule LogflareWeb.Source.SearchLVTest do
assert view |> element("#logs-list-container") |> render() =~ matching_message
refute view |> element("#logs-list-container") |> render() =~ non_matching_message
end
+
+ test "tailing retains rows on an empty result", %{conn: conn, source: source} do
+ {:ok, view, _html} = live_with_redirect(conn, Routes.live_path(conn, SearchLV, source.id))
+
+ view
+ |> TestUtils.wait_for_render("#logs-list-container li")
+
+ events_op = get_view_assigns(view).search_op_log_events
+ initial_log_event_ids = log_event_ids(view)
+
+ send(view.pid, {:search_result, %{events: %{events_op | rows: []}}})
+
+ assert log_event_ids(view) == initial_log_event_ids
+ end
+
+ test "tailing inserts a late event at its timestamp position", %{
+ conn: conn,
+ source: source
+ } do
+ timestamp_a = DateTime.utc_now() |> DateTime.to_unix(:microsecond)
+ timestamp_b = timestamp_a + 1
+ timestamp_c = timestamp_a + 2
+ message_prefix = "postgres-tail-order-#{System.unique_integer([:positive])}"
+
+ event_a =
+ build(:log_event,
+ source: source,
+ id: "ordered-event-a",
+ timestamp: timestamp_a,
+ message: "#{message_prefix} A"
+ )
+
+ event_b =
+ build(:log_event,
+ source: source,
+ id: "ordered-event-b",
+ timestamp: timestamp_b,
+ message: "#{message_prefix} B"
+ )
+
+ event_c =
+ build(:log_event,
+ source: source,
+ id: "ordered-event-c",
+ timestamp: timestamp_c,
+ message: "#{message_prefix} C"
+ )
+
+ assert {:ok, 2} = Backends.ingest_logs([event_c, event_a], source)
+
+ {:ok, view, _html} = live_with_redirect(conn, Routes.live_path(conn, SearchLV, source.id))
+
+ %{executor_pid: search_executor_pid} = get_view_assigns(view)
+ allow_sandbox(search_executor_pid)
+
+ render_change(view, :start_search, %{"querystring" => message_prefix})
+
+ view
+ |> TestUtils.wait_for_render("#log-events-ordered-event-a-#{timestamp_a}")
+
+ assert log_event_ids(view) ==
+ [
+ "log-events-ordered-event-c-#{timestamp_c}",
+ "log-events-ordered-event-a-#{timestamp_a}"
+ ]
+
+ assert {:ok, 1} = Backends.ingest_logs([event_b], source)
+
+ send(view.pid, :schedule_tail_search)
+
+ view
+ |> TestUtils.wait_for_render("#log-events-ordered-event-b-#{timestamp_b}")
+
+ assert log_event_ids(view) ==
+ [
+ "log-events-ordered-event-c-#{timestamp_c}",
+ "log-events-ordered-event-b-#{timestamp_b}",
+ "log-events-ordered-event-a-#{timestamp_a}"
+ ]
+ end
+
+ test "changing timezone immediately clears displayed search results", %{
+ conn: conn,
+ source: source,
+ matching_message: matching_message
+ } do
+ {:ok, view, _html} = live_with_redirect(conn, Routes.live_path(conn, SearchLV, source.id))
+
+ %{executor_pid: search_executor_pid} = get_view_assigns(view)
+ allow_sandbox(search_executor_pid)
+
+ render_change(view, :start_search, %{"querystring" => matching_message})
+
+ assert view
+ |> TestUtils.wait_for_render("#logs-list > li")
+ |> has_element?("#logs-list > li", matching_message)
+
+ html =
+ view
+ |> element("#results-actions")
+ |> render_change(%{search_timezone: "Singapore"})
+
+ assert html
+ |> Floki.parse_document!()
+ |> Floki.find("#logs-list > li")
+ |> Enum.empty?()
+ end
end
describe "single tenant postgres shorthand timestamp search" do
@@ -2452,6 +2595,15 @@ defmodule LogflareWeb.Source.SearchLVTest do
:sys.get_state(view.pid).socket.assigns
end
+ defp log_event_ids(view) do
+ view
+ |> element("#logs-list")
+ |> render()
+ |> Floki.parse_document!()
+ |> Floki.find("#logs-list > li")
+ |> Floki.attribute("id")
+ end
+
defp find_search_form_value(html, selector) do
{:ok, document} = Floki.parse_document(html)