Skip to content
Merged
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
106 changes: 89 additions & 17 deletions app/lib/linear_cli/linear/paginate.ex
Original file line number Diff line number Diff line change
Expand Up @@ -14,22 +14,41 @@ defmodule LinearCli.Linear.Paginate do
`max` records are collected or the API reports no more pages, decoding each raw
node through `decode_fun`.

`field_name` is the top-level response key (e.g. `"teams"`) holding
`edges`/`pageInfo`. `variables_fun` receives the current `after` cursor
(`nil` on the first page) and returns the GraphQL variables map.
`field_name` is the response key holding `edges`/`pageInfo`. Pass a list of
keys for a nested connection, such as `["team", "members"]`.
`variables_fun` receives the current `after` cursor (`nil` on the first page)
and returns the GraphQL variables map.
"""
def all(document, field_name, variables_fun, decode_fun, max \\ 100) do
do_all(document, field_name, variables_fun, decode_fun, nil, max, [])
do_all(document, field_name, variables_fun, decode_fun, %{
after_cursor: nil,
max: max,
acc: [],
seen_cursors: :bounded
})
end

@doc """
Fetches every page of a GraphQL connection and decodes each record.

This variant has no record limit. Use it only for lookup candidate sets that
must be complete before matching, rather than for bounded issue operations.
It rejects a null continuation cursor and every continuation cursor that it
has already used.
"""
def all_pages(document, field_name, variables_fun, decode_fun) do
do_all(document, field_name, variables_fun, decode_fun, nil, :unbounded, [])
do_all(
document,
field_name,
variables_fun,
decode_fun,
%{
after_cursor: nil,
max: :unbounded,
acc: [],
seen_cursors: MapSet.new([nil])
}
)
end

@doc """
Expand All @@ -51,29 +70,71 @@ defmodule LinearCli.Linear.Paginate do
end
end

defp do_all(document, field_name, variables_fun, decode_fun, after_cursor, max, acc) do
with {:ok, data} <- Api.call(document, variables_fun.(after_cursor)),
defp do_all(document, field_name, variables_fun, decode_fun, state) do
with {:ok, data} <- Api.call(document, variables_fun.(state.after_cursor)),
{:ok, %{"edges" => edges, "pageInfo" => page_info}} <-
fetch_connection(data, field_name) do
acc = acc ++ Enum.map(edges, &decode_fun.(&1["node"]))
state = %{state | acc: state.acc ++ Enum.map(edges, &decode_fun.(&1["node"]))}

if reached_limit?(acc, max) or !page_info["hasNextPage"] do
{:ok, take_max(acc, max)}
if reached_limit?(state.acc, state.max) or !page_info["hasNextPage"] do
{:ok, take_max(state.acc, state.max)}
else
next_cursor = page_info["endCursor"]

if next_cursor == after_cursor do
{:error, {:non_advancing_cursor, next_cursor}}
else
do_all(document, field_name, variables_fun, decode_fun, next_cursor, max, acc)
end
continue(document, field_name, variables_fun, decode_fun, state, next_cursor)
end
else
{:error, {:http_error, status, _body}} -> {:error, {:http_error, status}}
{:error, reason} -> {:error, reason}
end
end

defp continue(
document,
field_name,
variables_fun,
decode_fun,
%{after_cursor: after_cursor, seen_cursors: :bounded} = state,
next_cursor
) do
if next_cursor == after_cursor do
{:error, {:non_advancing_cursor, next_cursor}}
else
do_all(document, field_name, variables_fun, decode_fun, %{state | after_cursor: next_cursor})
end
end

defp continue(
_document,
_field_name,
_variables_fun,
_decode_fun,
%{seen_cursors: _},
nil
) do
{:error, {:non_advancing_cursor, nil}}
end

defp continue(
document,
field_name,
variables_fun,
decode_fun,
%{seen_cursors: seen_cursors} = state,
next_cursor
) do
if MapSet.member?(seen_cursors, next_cursor) do
{:error, {:non_advancing_cursor, next_cursor}}
else
state = %{
state
| after_cursor: next_cursor,
seen_cursors: MapSet.put(seen_cursors, next_cursor)
}

do_all(document, field_name, variables_fun, decode_fun, state)
end
end

defp reached_limit?(_acc, :unbounded), do: false
defp reached_limit?(acc, max), do: length(acc) >= max

Expand All @@ -83,7 +144,11 @@ defmodule LinearCli.Linear.Paginate do
# Safely extracts the named connection from the response data. Returns
# {:error, {:unexpected_response, ...}} instead of crashing with KeyError
# when the field is absent or not the expected connection shape.
defp fetch_connection(data, field_name) do
defp fetch_connection(data, field_name) when is_binary(field_name) do
fetch_connection(data, [field_name])
end

defp fetch_connection(data, [field_name]) do
case data do
%{^field_name => %{"edges" => _, "pageInfo" => _} = connection} ->
{:ok, connection}
Expand All @@ -95,4 +160,11 @@ defmodule LinearCli.Linear.Paginate do
{:error, {:unexpected_response, data}}
end
end

defp fetch_connection(data, [field_name | rest]) do
case data do
%{^field_name => nested} -> fetch_connection(nested, rest)
_ -> {:error, {:unexpected_response, data}}
end
end
end
41 changes: 32 additions & 9 deletions app/lib/linear_cli/linear/user.ex
Original file line number Diff line number Diff line change
Expand Up @@ -82,18 +82,11 @@ defmodule LinearCli.Linear.User.Read.ByTeam do
@moduledoc false
use Ash.Resource.ManualRead

def read(query, ecto_query, opts, context) do
LinearCli.Linear.User.Read.ByTeamForLookup.read(query, ecto_query, opts, context)
end
end

defmodule LinearCli.Linear.User.Read.ByTeamForLookup do
@moduledoc false
use Ash.Resource.ManualRead

alias LinearCli.Api
alias LinearCli.Linear.User

# Keep the team-scoped action separate from the strict workspace lookup path.
# EXT-75 owns this action's response behavior.
@document """
query($id: String!, $after: String) {
team(id: $id) {
Expand Down Expand Up @@ -146,3 +139,33 @@ defmodule LinearCli.Linear.User.Read.ByTeamForLookup do
end
end
end

defmodule LinearCli.Linear.User.Read.ByTeamForLookup do
@moduledoc false
use Ash.Resource.ManualRead

alias LinearCli.Linear.Paginate
alias LinearCli.Linear.User

@document """
query($id: String!, $after: String) {
team(id: $id) {
members(first: 50, after: $after) {
edges { node { #{User.base_fields()} } cursor }
pageInfo { hasNextPage endCursor }
}
}
}
"""

def read(query, _ecto_query, _opts, _context) do
team_id = query.arguments.team_id

Paginate.all_pages(
@document,
["team", "members"],
fn after_cursor -> %{"id" => team_id, "after" => after_cursor} end,
&User.from_map/1
)
end
end
116 changes: 116 additions & 0 deletions app/test/linear_cli/linear/paginate_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,19 @@ defmodule LinearCli.Linear.PaginateTest do

defp variables_fun(after_cursor), do: %{"after" => after_cursor}

defp nested_response(ids, has_next_page, end_cursor) do
%{
"data" => %{
"team" => %{
"members" => %{
"edges" => Enum.map(ids, &%{"node" => %{"id" => &1}, "cursor" => "row-#{&1}"}),
"pageInfo" => %{"hasNextPage" => has_next_page, "endCursor" => end_cursor}
}
}
}
}
end

test "uses the default 100-record limit" do
test_pid = self()

Expand Down Expand Up @@ -82,6 +95,109 @@ defmodule LinearCli.Linear.PaginateTest do
Paginate.all_pages("query", "issues", &variables_fun/1, & &1["id"])
end

test "all_pages rejects a cursor that appeared on an earlier page" do
test_pid = self()

Req.Test.stub(LinearCli.Api, fn conn ->
{:ok, body, conn} = Plug.Conn.read_body(conn)
cursor = Jason.decode!(body)["variables"]["after"]
send(test_pid, {:cursor, cursor})

response =
case cursor do
nil -> response([1], true, "c1")
"c1" -> response([2], true, "c2")
"c2" -> response([3], true, "c1")
end

Req.Test.json(conn, response)
end)

assert {:error, {:non_advancing_cursor, "c1"}} =
Paginate.all_pages("query", "issues", &variables_fun/1, & &1["id"])

assert_receive {:cursor, nil}
assert_receive {:cursor, "c1"}
assert_receive {:cursor, "c2"}
refute_receive {:cursor, "c1"}
end

test "all_pages rejects a null continuation cursor" do
test_pid = self()

Req.Test.stub(LinearCli.Api, fn conn ->
{:ok, body, conn} = Plug.Conn.read_body(conn)
cursor = Jason.decode!(body)["variables"]["after"]
send(test_pid, {:cursor, cursor})

response =
case cursor do
nil -> response([1], true, "c1")
"c1" -> response([2], true, nil)
end

Req.Test.json(conn, response)
end)

assert {:error, {:non_advancing_cursor, nil}} =
Paginate.all_pages("query", "issues", &variables_fun/1, & &1["id"])

assert_receive {:cursor, nil}
assert_receive {:cursor, "c1"}
refute_receive {:cursor, nil}
end

test "all_pages reads a nested connection path" do
Req.Test.stub(LinearCli.Api, fn conn ->
{:ok, body, conn} = Plug.Conn.read_body(conn)
cursor = Jason.decode!(body)["variables"]["after"]

response =
case cursor do
nil -> nested_response([1], true, "c1")
"c1" -> nested_response([2], false, "c2")
end

Req.Test.json(conn, response)
end)

assert {:ok, [1, 2]} =
Paginate.all_pages(
"query",
["team", "members"],
&variables_fun/1,
& &1["id"]
)
end

test "bounded all preserves its limit when cursors cycle" do
test_pid = self()

Req.Test.stub(LinearCli.Api, fn conn ->
{:ok, body, conn} = Plug.Conn.read_body(conn)
cursor = Jason.decode!(body)["variables"]["after"]
send(test_pid, {:cursor, cursor})

response =
case cursor do
nil -> response(1..20, true, "c1")
"c1" -> response(21..40, true, "c2")
"c2" -> response(41..60, true, "c1")
end

Req.Test.json(conn, response)
end)

assert {:ok, values} = Paginate.all("query", "issues", &variables_fun/1, & &1["id"], 100)
assert length(values) == 100
assert_receive {:cursor, nil}
assert_receive {:cursor, "c1"}
assert_receive {:cursor, "c2"}
assert_receive {:cursor, "c1"}
assert_receive {:cursor, "c2"}
refute_receive {:cursor, "c1"}
end

test "returns the first page and reports more records without following the cursor" do
test_pid = self()

Expand Down
Loading
Loading