From b6a26f8cc27c41f163b6ef12994c600a99b048a2 Mon Sep 17 00:00:00 2001 From: hanwenli Date: Mon, 28 Sep 2026 13:36:11 -0700 Subject: [PATCH 1/2] Revert "[integ-tests] Add Next token polling and start time to reduce the API response size" This reverts commit 95da8700d5368ef80e0699105ed514f9b4168202. --- .../tests/createami/test_createami.py | 32 ++----------------- 1 file changed, 2 insertions(+), 30 deletions(-) diff --git a/tests/integration-tests/tests/createami/test_createami.py b/tests/integration-tests/tests/createami/test_createami.py index 8f759d980e..69e17e892f 100644 --- a/tests/integration-tests/tests/createami/test_createami.py +++ b/tests/integration-tests/tests/createami/test_createami.py @@ -350,41 +350,13 @@ def _test_list_image_log_streams(image): assert_that(stream_names).contains(expected_log_stream) -def _read_all_log_events(image, log_stream_name, **args): - """Read image build log events, following CloudWatch's pagination token to collect the full list. - - ``GetLogEvents`` can return a partially full or even empty page while more events remain reachable - via the returned token; an empty page does not mean the stream is empty (see - https://docs.aws.amazon.com/AmazonCloudWatchLogs/latest/APIReference/API_GetLogEvents.html). - This is common for a tail read (``start_from_head`` unset), whose first page can come back empty. - Follow the token until it stops advancing, so the whole requested window is read and a spurious - empty page does not surface as an empty result. - """ - # Forward reads advance via nextToken, backward (tail) reads via prevToken. Either way each page - # moves away from where the read started, so the events stay ordered from that starting point. - token_key = "nextToken" if args.get("start_from_head") else "prevToken" - events = [] - token = args.get("next_token") - while True: - response = image.get_log_events(log_stream_name, **{**args, "next_token": token}) - events.extend(response["events"]) - next_token = response.get(token_key) - if not next_token or next_token == token: - break - token = next_token - return events - - def _test_get_image_log_events(image): """Test pcluster get-image-log-events functionality.""" logging.info("Testing that pcluster get-image-log-events is working as expected") log_stream_name = f"{get_installed_parallelcluster_base_version()}/1" - # Get the first event to establish time boundary for testing. + # Get the first event to establish time boundary for testing initial_events = image.get_log_events(log_stream_name, limit=1, start_from_head=True) - assert_that(initial_events["events"]).described_as( - f"no log events found in stream {log_stream_name}; cannot establish first event" - ).is_not_empty() first_event = initial_events["events"][0] first_event_time_str = first_event["timestamp"] first_event_time = date_parse(first_event_time_str) @@ -404,7 +376,7 @@ def _test_get_image_log_events(image): ] for args, expect_first, expect_count in test_cases: - events = _read_all_log_events(image, log_stream_name, **args) + events = image.get_log_events(log_stream_name, **args)["events"] if expect_count is not None: assert_that(events).is_length(expect_count) From d99384ef07e49f5be061c69db4497b15d981bd21 Mon Sep 17 00:00:00 2001 From: hanwenli Date: Mon, 28 Sep 2026 06:47:52 -0700 Subject: [PATCH 2/2] [integ-test] Make get-cluster/image-log-events integ tests robust to empty pages GetLogEvents may return partially full or empty pages even when more events are available; `limit` is only an upper bound. The previous checks assumed `limit=1` always returned exactly one event and indexed `events[0]`, making them flaky. Replace the filter matrix with a full pagination of the log stream (until the returned token equals the one passed in) and assert the stream is not empty. Arguments like `limit` are already covered by unit tests, and don't have to be in integration tests. --- .../tests/cli_commands/test_cli_commands.py | 45 +++++-------------- .../tests/createami/test_createami.py | 45 +++++-------------- 2 files changed, 24 insertions(+), 66 deletions(-) diff --git a/tests/integration-tests/tests/cli_commands/test_cli_commands.py b/tests/integration-tests/tests/cli_commands/test_cli_commands.py index 9a94645226..75f4a5fe04 100644 --- a/tests/integration-tests/tests/cli_commands/test_cli_commands.py +++ b/tests/integration-tests/tests/cli_commands/test_cli_commands.py @@ -9,7 +9,6 @@ # or in the "LICENSE.txt" file accompanying this file. # This file is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, express or implied. # See the License for the specific language governing permissions and limitations under the License. -import datetime import json import logging import os as os_lib @@ -21,7 +20,6 @@ import botocore import pytest from assertpy import assert_that -from dateutil.parser import parse as date_parse from framework.credential_providers import run_pcluster_command from remote_command_executor import RemoteCommandExecutor from utils import ( @@ -405,37 +403,18 @@ def _test_pcluster_get_cluster_log_events(cluster): cluster_info = cluster.describe_cluster() cfn_init_log_stream = instance_stream_name(cluster_info["headNode"], "cfn-init") - # Get the first event to establish time boundary for testing - initial_events = cluster.get_log_events(cfn_init_log_stream, limit=1, start_from_head=True) - first_event = initial_events["events"][0] - first_event_time_str = first_event["timestamp"] - first_event_time = date_parse(first_event_time_str) - before_first = (first_event_time - datetime.timedelta(seconds=1)).isoformat() - after_first = (first_event_time + datetime.timedelta(seconds=1)).isoformat() - - # args, expect_first, expect_count - test_cases = [ - ({}, None, None), - ({"limit": 1}, False, 1), - ({"limit": 2, "start_from_head": True}, True, 2), - ({"limit": 1, "start_time": before_first, "end_time": after_first, "start_from_head": True}, True, 1), - ({"limit": 1, "end_time": before_first}, None, 0), - ({"limit": 1, "start_time": after_first, "start_from_head": True}, False, 1), - ({"limit": 1, "next_token": initial_events["nextToken"]}, False, 1), - ({"limit": 1, "next_token": initial_events["nextToken"], "start_from_head": True}, False, 1), - ] - - for args, expect_first, expect_count in test_cases: - events = cluster.get_log_events(cfn_init_log_stream, **args)["events"] - - if expect_count is not None: - assert_that(events).is_length(expect_count) - - if expect_first is True: - assert_that(events[0]["message"]).is_equal_to(first_event["message"]) - - if expect_first is False: - assert_that(events[0]["message"]).is_not_equal_to(first_event["message"]) + # GetLogEvents can return empty pages while more events are available, + # so page forward until the token stops changing, which marks the end of the stream. + events = [] + next_token = None + while True: + response = cluster.get_log_events(cfn_init_log_stream, start_from_head=True, next_token=next_token) + events.extend(response["events"]) + if not response.get("nextToken") or response["nextToken"] == next_token: + break + next_token = response["nextToken"] + + assert_that(events).is_not_empty() def _test_pcluster_get_cluster_stack_events(cluster): diff --git a/tests/integration-tests/tests/createami/test_createami.py b/tests/integration-tests/tests/createami/test_createami.py index 69e17e892f..41245fcf44 100644 --- a/tests/integration-tests/tests/createami/test_createami.py +++ b/tests/integration-tests/tests/createami/test_createami.py @@ -9,7 +9,6 @@ # or in the "LICENSE.txt" file accompanying this file. # This file is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, express or implied. # See the License for the specific language governing permissions and limitations under the License. -import datetime import json import logging import os @@ -23,7 +22,6 @@ from assertpy import assert_that, soft_assertions from botocore.exceptions import ClientError from cfn_stacks_factory import CfnStack -from dateutil.parser import parse as date_parse from remote_command_executor import RemoteCommandExecutor from retrying import retry from time_utils import minutes, seconds @@ -355,37 +353,18 @@ def _test_get_image_log_events(image): logging.info("Testing that pcluster get-image-log-events is working as expected") log_stream_name = f"{get_installed_parallelcluster_base_version()}/1" - # Get the first event to establish time boundary for testing - initial_events = image.get_log_events(log_stream_name, limit=1, start_from_head=True) - first_event = initial_events["events"][0] - first_event_time_str = first_event["timestamp"] - first_event_time = date_parse(first_event_time_str) - before_first = (first_event_time - datetime.timedelta(seconds=1)).isoformat() - after_first = (first_event_time + datetime.timedelta(seconds=1)).isoformat() - - # args, expect_first, expect_count - test_cases = [ - ({}, None, None), - ({"limit": 1}, False, 1), - ({"limit": 2, "start_from_head": True}, True, 2), - ({"limit": 1, "start_time": before_first, "end_time": after_first, "start_from_head": True}, True, 1), - ({"limit": 1, "end_time": before_first}, None, 0), - ({"limit": 1, "start_time": after_first, "start_from_head": True}, False, 1), - ({"limit": 1, "next_token": initial_events["nextToken"]}, False, 1), - ({"limit": 1, "next_token": initial_events["nextToken"], "start_from_head": True}, False, 1), - ] - - for args, expect_first, expect_count in test_cases: - events = image.get_log_events(log_stream_name, **args)["events"] - - if expect_count is not None: - assert_that(events).is_length(expect_count) - - if expect_first is True: - assert_that(events[0]["message"]).contains(first_event["message"]) - - if expect_first is False: - assert_that(events[0]["message"]).does_not_contain(first_event["message"]) + # GetLogEvents can return empty pages while more events are available, + # so page forward until the token stops changing, which marks the end of the stream. + events = [] + next_token = None + while True: + response = image.get_log_events(log_stream_name, start_from_head=True, next_token=next_token) + events.extend(response["events"]) + if not response.get("nextToken") or response["nextToken"] == next_token: + break + next_token = response["nextToken"] + + assert_that(events).is_not_empty() def _set_s3_bucket_policy(bucket_name, partition, region):