Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions lib/optimal_engine/api/router.ex
Original file line number Diff line number Diff line change
Expand Up @@ -2266,6 +2266,7 @@ defmodule OptimalEngine.API.Router do
tenant_id: queue.tenant_id,
workspace_id: queue.workspace_id,
count: queue.count,
returned: queue.returned,
review_counts: queue.review_counts,
lifecycle_counts: queue.lifecycle_counts,
claims: Enum.map(queue.claims, &claim_to_map/1)
Expand Down
36 changes: 18 additions & 18 deletions lib/optimal_engine/memory_core/claim_review.ex
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ defmodule OptimalEngine.MemoryCore.ClaimReview do
tenant_id: String.t(),
workspace_id: String.t(),
count: non_neg_integer(),
returned: non_neg_integer(),
review_counts: map(),
lifecycle_counts: map(),
claims: [map()]
Expand All @@ -33,21 +34,27 @@ defmodule OptimalEngine.MemoryCore.ClaimReview do
tenant_id = Keyword.get(opts, :tenant_id, "default")
workspace_id = Keyword.get(opts, :workspace_id, "default")

with {:ok, claims} <-
Store.list_claims(
tenant_id: tenant_id,
workspace_id: workspace_id,
review_status: Keyword.get(opts, :review_status),
lifecycle_state: Keyword.get(opts, :lifecycle_state),
limit: Keyword.get(opts, :limit, 100)
) do
filters = [
tenant_id: tenant_id,
workspace_id: workspace_id,
review_status: Keyword.get(opts, :review_status),
lifecycle_state: Keyword.get(opts, :lifecycle_state)
]

# `count`/`*_counts` come from a grouped count over ALL matching rows;
# `claims` is the page the caller asked for. Before the store honoured
# `:limit` these were the same list — now the totals stay honest while the
# page stays bounded.
with {:ok, claims} <- Store.list_claims(filters ++ [limit: Keyword.get(opts, :limit, 100)]),
{:ok, counts} <- Store.count_claims(filters) do
{:ok,
%{
tenant_id: tenant_id,
workspace_id: workspace_id,
count: length(claims),
review_counts: count_by(claims, :review_status),
lifecycle_counts: count_by(claims, :lifecycle_state),
count: counts.total,
returned: length(claims),
review_counts: counts.review_counts,
lifecycle_counts: counts.lifecycle_counts,
claims: claims
}}
end
Expand Down Expand Up @@ -238,13 +245,6 @@ defmodule OptimalEngine.MemoryCore.ClaimReview do
end
end

defp count_by(claims, field) do
claims
|> Enum.map(&Map.get(&1, field))
|> Enum.reject(&is_nil/1)
|> Enum.frequencies()
end

defp fact_opts(opts, policy) do
[
actor_id: Keyword.get(opts, :actor_id),
Expand Down
76 changes: 65 additions & 11 deletions lib/optimal_engine/memory_core/store.ex
Original file line number Diff line number Diff line change
Expand Up @@ -758,29 +758,83 @@ defmodule OptimalEngine.MemoryCore.Store do
end

def list_claims(workspace_id, opts) when is_binary(workspace_id) and is_list(opts) do
{clauses, params} =
filter_clauses(
[
{"tenant_id =", Keyword.get(opts, :tenant_id)},
{"lifecycle_state =", Keyword.get(opts, :lifecycle_state)},
{"review_status =", Keyword.get(opts, :review_status)},
{"source_package_id =", Keyword.get(opts, :source_package_id)}
],
[workspace_id]
)
{clauses, params} = claim_filter_clauses(workspace_id, opts)
{limit_clause, params} = limit_clause(Keyword.get(opts, :limit), params)

sql = """
SELECT #{Enum.join(@claim_columns, ", ")}
FROM claims
WHERE workspace_id = ?1 #{clauses}
ORDER BY created_at, id
ORDER BY created_at, id#{limit_clause}
"""

with {:ok, rows} <- Store.raw_query(sql, params) do
{:ok, Enum.map(rows, &claim_from_row/1)}
end
end

@doc """
Counts claims matching the same filters as `list_claims/2`, without hydrating
rows. Returns the total plus per-status breakdowns, so a caller that pages
with `:limit` can still report honest totals.
"""
@spec count_claims(keyword()) ::
{:ok,
%{
total: non_neg_integer(),
review_counts: %{optional(String.t()) => non_neg_integer()},
lifecycle_counts: %{optional(String.t()) => non_neg_integer()}
}}
| {:error, term()}
def count_claims(opts \\ []) when is_list(opts) do
workspace_id = Keyword.get(opts, :workspace_id, "default")
{clauses, params} = claim_filter_clauses(workspace_id, opts)

sql = """
SELECT review_status, lifecycle_state, COUNT(*)
FROM claims
WHERE workspace_id = ?1 #{clauses}
GROUP BY review_status, lifecycle_state
"""

with {:ok, rows} <- Store.raw_query(sql, params) do
{:ok,
%{
total: rows |> Enum.map(fn [_, _, n] -> n end) |> Enum.sum(),
review_counts: sum_grouped(rows, fn [review, _, n] -> {review, n} end),
lifecycle_counts: sum_grouped(rows, fn [_, lifecycle, n] -> {lifecycle, n} end)
}}
end
end

defp claim_filter_clauses(workspace_id, opts) do
filter_clauses(
[
{"tenant_id =", Keyword.get(opts, :tenant_id)},
{"lifecycle_state =", Keyword.get(opts, :lifecycle_state)},
{"review_status =", Keyword.get(opts, :review_status)},
{"source_package_id =", Keyword.get(opts, :source_package_id)}
],
[workspace_id]
)
end

defp sum_grouped(rows, mapper) do
Enum.reduce(rows, %{}, fn row, acc ->
{key, n} = mapper.(row)
Map.update(acc, to_string(key), n, &(&1 + n))
end)
end

# A missing or non-positive :limit keeps the query unbounded — existing
# callers that never paged keep their behaviour; a positive limit becomes a
# bound SQL parameter instead of being silently dropped.
defp limit_clause(limit, params) when is_integer(limit) and limit > 0 do
{"\nLIMIT ?#{length(params) + 1}", params ++ [limit]}
end

defp limit_clause(_absent_or_invalid, params), do: {"", params}

@doc """
Conditionally transitions a Claim out of review.

Expand Down
19 changes: 19 additions & 0 deletions test/memory_core/claim_review_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -328,6 +328,25 @@ defmodule OptimalEngine.MemoryCore.ClaimReviewTest do
)
end

test "queue honours :limit as a SQL bound while totals stay honest" do
workspace_id = ws()

for i <- 1..7 do
assert {:ok, _claim} =
extracted_claim(workspace_id, claim_text: "Queue limit case #{i}")
end

assert {:ok, queue} = MemoryCore.claim_review_queue(workspace_id: workspace_id, limit: 3)

# The page is bounded by the requested limit...
assert length(queue.claims) == 3
assert queue.returned == 3
# ...while count and the breakdowns keep describing ALL matching claims.
assert queue.count == 7
assert queue.review_counts["unreviewed"] == 7
assert queue.lifecycle_counts["pending"] == 7
end

defp extracted_claim(workspace_id, opts) do
source =
MemoryCore.source_package_from_text(Keyword.fetch!(opts, :claim_text),
Expand Down