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
Original file line number Diff line number Diff line change
Expand Up @@ -574,16 +574,21 @@ def default_pagination(
query_parameters: dict[str, Any] | None = None,
responses: Callable[[], list[dict[str, Any]] | None] = lambda: [],
) -> tuple[Any, dict[str, Any] | None]:
"""
Resolve the url and query parameters of the page following ``response``.

The ``$skip`` offset is derived from ``query_parameters`` rather than from ``responses``:
callers accumulate whatever their own callbacks produced, so the entries are not guaranteed
to be the raw pages this offset would have to be counted from.
"""
if isinstance(response, dict):
odata_count = response.get("@odata.count")
if odata_count and query_parameters:
top = query_parameters.get("$top")

if top and odata_count:
if len(response.get("value", [])) == top:
results = responses()
skip = sum([len(result["value"]) for result in results]) + top if results else top # type: ignore
query_parameters["$skip"] = skip
query_parameters["$skip"] = (query_parameters.get("$skip") or 0) + top
return url, query_parameters
return response.get("@odata.nextLink"), query_parameters
return None, query_parameters
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -187,14 +187,22 @@ def execute_complete(
self,
context: Context,
event: dict[Any, Any] | None = None,
query_parameters: dict[str, Any] | None = None,
) -> Any:
"""
Execute callback when MSGraphTrigger finishes execution.

This method gets executed automatically when MSGraphTrigger completes its execution.

:param query_parameters: The query parameters the completed page was requested with, passed
back by :meth:`trigger_next_link`. The operator is rebuilt from the serialized Dag on
every page, so without them ``self.query_parameters`` would still describe the first page.
"""
self.log.debug("context: %s", context)

if query_parameters is not None:
self.query_parameters = query_parameters

if event:
self.log.debug("%s completed with %s: %s", self.task_id, event.get("status"), event)

Expand Down Expand Up @@ -315,7 +323,6 @@ def paginate(
response=response,
url=operator.url,
query_parameters=operator.query_parameters,
responses=lambda: operator.pull_xcom(context),
)

def trigger_next_link(self, response, method_name: str, context: Context) -> None:
Expand Down Expand Up @@ -355,4 +362,5 @@ def trigger_next_link(self, response, method_name: str, context: Context) -> Non
pagination_link=True,
),
method_name=method_name,
kwargs={"query_parameters": query_parameters},
)
Original file line number Diff line number Diff line change
Expand Up @@ -442,6 +442,44 @@ async def test_paginated_run(self):
assert isinstance(actual, list)
assert actual == [users, next_users]

@pytest.mark.asyncio
@pytest.mark.parametrize(
("query_parameters", "expected_skips"),
[
pytest.param(
{"$top": 12, "$count": True},
["", "&%24skip=12", "&%24skip=24"],
id="from_the_start",
),
pytest.param(
{"$top": 12, "$count": True, "$skip": 100},
["&%24skip=100", "&%24skip=112", "&%24skip=124"],
id="from_a_user_supplied_offset",
),
],
)
async def test_paginated_run_advances_the_skip_offset_by_a_single_page(
self, query_parameters, expected_skips
):
messages = load_json_from_resources(dirname(__file__), "..", "resources", "messages.json")
second_messages = load_json_from_resources(
dirname(__file__), "..", "resources", "second_messages.json"
)
third_messages = load_json_from_resources(dirname(__file__), "..", "resources", "third_messages.json")
response = mock_json_response(200, messages, second_messages, third_messages)

with patch_hook_and_request_adapter(response) as mocks:
mock_get_http_response = mocks[-1]
hook = KiotaRequestAdapterHook(conn_id="msgraph_api")

# paginated_run mutates the query parameters it is given, so hand it a copy rather than
# the dict pytest built once at collection time.
await hook.paginated_run(url="users/messages", query_parameters=dict(query_parameters))

urls = [call.args[0].url for call in mock_get_http_response.call_args_list]

assert urls == [f"users/messages?%24top=12&%24count=true{skip}" for skip in expected_skips]

@pytest.mark.asyncio
async def test_paginated_run_refuses_cross_host_next_link(self):
first_page = {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -356,10 +356,44 @@ def test_trigger_next_link_forwards_the_path_parameters(self):
assert trigger.url == "users/{user_id}/mailFolders/{mailFolder_id}/messages"
assert trigger.path_parameters == path_parameters

def test_execute_complete_advances_the_skip_offset_it_was_resumed_with(self):
# The operator is rebuilt from the serialized Dag on every page, so its own query parameters
# describe the first request; only the resume kwargs know where the completed page came from.
operator = MSGraphAsyncOperator(
task_id="messages",
conn_id="msgraph_api",
url="users/messages",
query_parameters={"$top": 12, "$count": True},
)
context = mock_context(task=operator)
second_messages = load_json_from_resources(
dirname(__file__), "..", "resources", "second_messages.json"
)
event = {"status": "success", "type": "builtins.dict", "response": json.dumps(second_messages)}

with mock.patch.object(operator, "defer") as mock_defer:
operator.execute_complete(
context=context,
event=event,
query_parameters={"$top": 12, "$count": True, "$skip": 12},
)

assert mock_defer.call_args.kwargs["trigger"].query_parameters == {
"$top": 12,
"$count": True,
"$skip": 24,
}
assert mock_defer.call_args.kwargs["kwargs"] == {
"query_parameters": {"$top": 12, "$count": True, "$skip": 24}
}

def test_skip_pagination_expands_the_url_template_on_every_page(self):
messages = load_json_from_resources(dirname(__file__), "..", "resources", "messages.json")
next_messages = load_json_from_resources(dirname(__file__), "..", "resources", "next_messages.json")
response = mock_json_response(200, messages, next_messages)
second_messages = load_json_from_resources(
dirname(__file__), "..", "resources", "second_messages.json"
)
third_messages = load_json_from_resources(dirname(__file__), "..", "resources", "third_messages.json")
response = mock_json_response(200, messages, second_messages, third_messages)

with patch_hook_and_request_adapter(response) as (*_, mock_get_http_response):
operator = MSGraphAsyncOperator(
Expand All @@ -378,6 +412,7 @@ def test_skip_pagination_expands_the_url_template_on_every_page(self):
assert urls == [
"users/48d31887-5fad-4d73-a9f5-3c356e68a038/mailFolders/inbox/messages?%24top=12&%24count=true",
"users/48d31887-5fad-4d73-a9f5-3c356e68a038/mailFolders/inbox/messages?%24top=12&%24count=true&%24skip=12",
"users/48d31887-5fad-4d73-a9f5-3c356e68a038/mailFolders/inbox/messages?%24top=12&%24count=true&%24skip=24",
]

def test_pagination_issues_every_page_with_the_configured_request(self):
Expand Down
Original file line number Diff line number Diff line change
@@ -1 +1 @@
{"@odata.context": "https://graph.microsoft.com/v1.0/$metadata#users('48d31887-5fad-4d73-a9f5-3c356e68a038')/mailFolders('inbox')/messages", "@odata.count": 18, "@odata.nextLink": "https://graph.microsoft.com/v1.0/users('48d31887-5fad-4d73-a9f5-3c356e68a038')/mailFolders('inbox')/messages?%24top=12&%24count=true&%24skip=12", "value": [{"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwR4Hg\"", "id": "AAMkAGUAAAwTW09AAA=", "subject": "Weekly status report", "receivedDateTime": "2026-08-14T08:02:11Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwR7Km\"", "id": "AAMkAGUAAAwTW1BAAA=", "subject": "Re: Budget approval", "receivedDateTime": "2026-08-14T09:47:03Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSBd2\"", "id": "AAMkAGUAAAwTW2CAAA=", "subject": "Lunch tomorrow?", "receivedDateTime": "2026-08-14T11:15:52Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSFqP\"", "id": "AAMkAGUAAAwTW3DAAA=", "subject": "Deployment window moved", "receivedDateTime": "2026-08-14T13:20:07Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSJt4\"", "id": "AAMkAGUAAAwTW4EAAA=", "subject": "Re: Q3 roadmap", "receivedDateTime": "2026-08-14T16:38:44Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSNc1\"", "id": "AAMkAGUAAAwTW5FAAA=", "subject": "Invoice 2026-0814", "receivedDateTime": "2026-08-15T07:05:19Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSRkD\"", "id": "AAMkAGUAAAwTW6GAAA=", "subject": "Welcome to the team", "receivedDateTime": "2026-08-15T09:12:36Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSVpq\"", "id": "AAMkAGUAAAwTW7HAAA=", "subject": "Security review scheduled", "receivedDateTime": "2026-08-15T10:44:58Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSZ2X\"", "id": "AAMkAGUAAAwTW8IAAA=", "subject": "Re: Vacation request", "receivedDateTime": "2026-08-15T14:29:13Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSdG9\"", "id": "AAMkAGUAAAwTW9JAAA=", "subject": "Quarterly all-hands", "receivedDateTime": "2026-08-16T08:51:40Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwShNe\"", "id": "AAMkAGUAAAwTXAKAAA=", "subject": "Please review the draft", "receivedDateTime": "2026-08-16T12:06:25Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSlz0\"", "id": "AAMkAGUAAAwTXBLAAA=", "subject": "Reminder: timesheet due", "receivedDateTime": "2026-08-16T15:33:02Z", "isRead": false}]}
{"@odata.context": "https://graph.microsoft.com/v1.0/$metadata#users('48d31887-5fad-4d73-a9f5-3c356e68a038')/mailFolders('inbox')/messages", "@odata.count": 30, "@odata.nextLink": "https://graph.microsoft.com/v1.0/users('48d31887-5fad-4d73-a9f5-3c356e68a038')/mailFolders('inbox')/messages?%24top=12&%24count=true&%24skip=12", "value": [{"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwR4Hg\"", "id": "AAMkAGUAAAwTW09AAA=", "subject": "Weekly status report", "receivedDateTime": "2026-08-14T08:02:11Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwR7Km\"", "id": "AAMkAGUAAAwTW1BAAA=", "subject": "Re: Budget approval", "receivedDateTime": "2026-08-14T09:47:03Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSBd2\"", "id": "AAMkAGUAAAwTW2CAAA=", "subject": "Lunch tomorrow?", "receivedDateTime": "2026-08-14T11:15:52Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSFqP\"", "id": "AAMkAGUAAAwTW3DAAA=", "subject": "Deployment window moved", "receivedDateTime": "2026-08-14T13:20:07Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSJt4\"", "id": "AAMkAGUAAAwTW4EAAA=", "subject": "Re: Q3 roadmap", "receivedDateTime": "2026-08-14T16:38:44Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSNc1\"", "id": "AAMkAGUAAAwTW5FAAA=", "subject": "Invoice 2026-0814", "receivedDateTime": "2026-08-15T07:05:19Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSRkD\"", "id": "AAMkAGUAAAwTW6GAAA=", "subject": "Welcome to the team", "receivedDateTime": "2026-08-15T09:12:36Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSVpq\"", "id": "AAMkAGUAAAwTW7HAAA=", "subject": "Security review scheduled", "receivedDateTime": "2026-08-15T10:44:58Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSZ2X\"", "id": "AAMkAGUAAAwTW8IAAA=", "subject": "Re: Vacation request", "receivedDateTime": "2026-08-15T14:29:13Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSdG9\"", "id": "AAMkAGUAAAwTW9JAAA=", "subject": "Quarterly all-hands", "receivedDateTime": "2026-08-16T08:51:40Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwShNe\"", "id": "AAMkAGUAAAwTXAKAAA=", "subject": "Please review the draft", "receivedDateTime": "2026-08-16T12:06:25Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSlz0\"", "id": "AAMkAGUAAAwTXBLAAA=", "subject": "Reminder: timesheet due", "receivedDateTime": "2026-08-16T15:33:02Z", "isRead": false}]}
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
{"@odata.context": "https://graph.microsoft.com/v1.0/$metadata#users('48d31887-5fad-4d73-a9f5-3c356e68a038')/mailFolders('inbox')/messages", "@odata.count": 30, "@odata.nextLink": "https://graph.microsoft.com/v1.0/users('48d31887-5fad-4d73-a9f5-3c356e68a038')/mailFolders('inbox')/messages?%24top=12&%24count=true&%24skip=24", "value": [{"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwTA01\"", "id": "AAMkAGUAAAwUW09AAA=", "subject": "Re: Reminder: timesheet due", "receivedDateTime": "2026-08-16T16:47:18Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwTEjR\"", "id": "AAMkAGUAAAwUW1BAAA=", "subject": "Build 2026.8.16 succeeded", "receivedDateTime": "2026-08-16T17:22:05Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwTIvz\"", "id": "AAMkAGUAAAwUW2CAAA=", "subject": "Re: Quarterly all-hands", "receivedDateTime": "2026-08-16T18:03:41Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwTM8h\"", "id": "AAMkAGUAAAwUW3DAAA=", "subject": "Nightly backup report", "receivedDateTime": "2026-08-16T19:15:56Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwTQlN\"", "id": "AAMkAGUAAAwUW4EAAA=", "subject": "Re: Security review scheduled", "receivedDateTime": "2026-08-16T20:38:12Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwTUyF\"", "id": "AAMkAGUAAAwUW5FAAA=", "subject": "On-call handover notes", "receivedDateTime": "2026-08-16T21:04:33Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwTY3B\"", "id": "AAMkAGUAAAwUW6GAAA=", "subject": "Expense report approved", "receivedDateTime": "2026-08-16T22:19:27Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwTcQd\"", "id": "AAMkAGUAAAwUW7HAAA=", "subject": "Re: Lunch tomorrow?", "receivedDateTime": "2026-08-17T05:41:08Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwTgpX\"", "id": "AAMkAGUAAAwUW8IAAA=", "subject": "Sprint planning agenda", "receivedDateTime": "2026-08-17T06:12:49Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwTk7M\"", "id": "AAMkAGUAAAwUW9JAAA=", "subject": "Password expires in 7 days", "receivedDateTime": "2026-08-17T07:03:15Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwToLW\"", "id": "AAMkAGUAAAwUXAKAAA=", "subject": "Re: Welcome to the team", "receivedDateTime": "2026-08-17T07:58:30Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwTsdG\"", "id": "AAMkAGUAAAwUXBLAAA=", "subject": "Conference room booking confirmed", "receivedDateTime": "2026-08-17T08:39:52Z", "isRead": false}]}
Original file line number Diff line number Diff line change
@@ -1 +1 @@
{"@odata.context": "https://graph.microsoft.com/v1.0/$metadata#users('48d31887-5fad-4d73-a9f5-3c356e68a038')/mailFolders('inbox')/messages", "@odata.count": 18, "value": [{"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSp4B\"", "id": "AAMkAGUAAAwTXCMAAA=", "subject": "Re: Deployment window moved", "receivedDateTime": "2026-08-17T09:18:47Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwStQm\"", "id": "AAMkAGUAAAwTXDNAAA=", "subject": "Offsite logistics", "receivedDateTime": "2026-08-17T11:52:30Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSxfL\"", "id": "AAMkAGUAAAwTXEOAAA=", "subject": "Re: Please review the draft", "receivedDateTime": "2026-08-18T08:24:15Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwS1uT\"", "id": "AAMkAGUAAAwTXFPAAA=", "subject": "New badge photo needed", "receivedDateTime": "2026-08-18T13:41:09Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwS5Ci\"", "id": "AAMkAGUAAAwTXGQAAA=", "subject": "Re: Invoice 2026-0814", "receivedDateTime": "2026-08-19T10:07:53Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwS9YR\"", "id": "AAMkAGUAAAwTXHRAAA=", "subject": "Parking permit renewal", "receivedDateTime": "2026-08-19T16:55:21Z", "isRead": false}]}
{"@odata.context": "https://graph.microsoft.com/v1.0/$metadata#users('48d31887-5fad-4d73-a9f5-3c356e68a038')/mailFolders('inbox')/messages", "@odata.count": 30, "value": [{"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSp4B\"", "id": "AAMkAGUAAAwTXCMAAA=", "subject": "Re: Deployment window moved", "receivedDateTime": "2026-08-17T09:18:47Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwStQm\"", "id": "AAMkAGUAAAwTXDNAAA=", "subject": "Offsite logistics", "receivedDateTime": "2026-08-17T11:52:30Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwSxfL\"", "id": "AAMkAGUAAAwTXEOAAA=", "subject": "Re: Please review the draft", "receivedDateTime": "2026-08-18T08:24:15Z", "isRead": true}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwS1uT\"", "id": "AAMkAGUAAAwTXFPAAA=", "subject": "New badge photo needed", "receivedDateTime": "2026-08-18T13:41:09Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwS5Ci\"", "id": "AAMkAGUAAAwTXGQAAA=", "subject": "Re: Invoice 2026-0814", "receivedDateTime": "2026-08-19T10:07:53Z", "isRead": false}, {"@odata.etag": "W/\"CQAAABYAAADHcgC8Hl9tRZ/hc1wEUs1TAAAwS9YR\"", "id": "AAMkAGUAAAwTXHRAAA=", "subject": "Parking permit renewal", "receivedDateTime": "2026-08-19T16:55:21Z", "isRead": false}]}