Skip to content

Commit 3ebcf7d

Browse files
committed
test: pace LMI execution history reads
Collect each case's own executions during teardown and leave full-run history/log collection to the final workflow step. Pace history reads and bound throttle-only retries, including reads through the public cloud runner, without retrying invocation or callback writes. Cover pagination, throttling exhaustion, permission errors, and scoped versus full-run collection with focused tests.
1 parent 9150c87 commit 3ebcf7d

4 files changed

Lines changed: 158 additions & 11 deletions

File tree

‎lmi_tests/README.md‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -185,6 +185,12 @@ later case's short observation budget. Checkpoint holds and returns are correlat
185185
to the same request, and stale-attempt effects are checked before waiting for a
186186
retry that might itself be unable to start on a pinned worker.
187187

188+
Per-case teardown collects only that case's executions. Complete-run histories
189+
and CloudWatch logs are collected once by the final workflow step. Both runner
190+
and artifact history reads are paced at one request per second, with at most five
191+
throttling attempts within a 20-second retry budget (in addition to the finite
192+
transport timeout). Invocation and callback writes are not retried by this layer.
193+
188194
## Independent budgets and results
189195

190196
| Budget | Default |

‎lmi_tests/cloud.py‎

Lines changed: 61 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -14,13 +14,27 @@
1414
from lmi_tests.evidence import CollectionError, check_controls, select
1515

1616

17+
class RunnerClient:
18+
"""Share paced history reads with the public runner; other APIs are unchanged."""
19+
20+
def __init__(self, cloud):
21+
self.cloud = cloud
22+
23+
def __getattr__(self, name):
24+
return getattr(self.cloud.lam, name)
25+
26+
def get_durable_execution_history(self, **kwargs):
27+
return self.cloud.history_page(**kwargs)
28+
29+
1730
class Cloud:
1831
def __init__(self, manifest):
1932
self.manifest = manifest
2033
self.lam, self.s3, self.logs = client("lambda"), client("s3"), client("logs")
2134
self.events = {}
2235
self.invocations = []
2336
self.gates = set()
37+
self._next_history_read = 0.0
2438

2539
def verify(self):
2640
for key, arn in self.manifest["functions"].items():
@@ -53,7 +67,7 @@ def start(self, scenario, fixture="normal", gate=None):
5367
runner = DurableFunctionCloudTestRunner(
5468
self.manifest["functions"][fixture], region=self.manifest["region"]
5569
)
56-
runner.lambda_client = self.lam # bounded transport, including async start
70+
runner.lambda_client = RunnerClient(self)
5771
try:
5872
arn = runner.run_async(input=payload)
5973
if not arn:
@@ -106,7 +120,7 @@ def release(self, gate):
106120
Bucket=self.manifest["bucket"], Key="control/" + gate, Body=b"release"
107121
)
108122

109-
def refresh(self, markers=None):
123+
def refresh(self, markers=None, *, validate_controls=True):
110124
# A new pytest fixture must not scan every earlier case before it can
111125
# observe its first event. Explicit [] requests the complete run for collection.
112126
if markers is None:
@@ -149,7 +163,7 @@ def read(key):
149163
)
150164
save(ARTIFACTS / "events.json", events)
151165
selected = [e for e in events if not markers or e["marker"] in markers]
152-
if markers:
166+
if markers and validate_controls:
153167
check_controls(selected)
154168
return selected
155169

@@ -196,13 +210,39 @@ def finish(self, item, status="SUCCEEDED"):
196210
self.phase(item, "WRAPPER_RETURN")
197211
return self.history(item)
198212

213+
def history_page(self, **request):
214+
"""Bound read-only throttling retries without replaying invocation writes."""
215+
deadline = time.monotonic() + 20
216+
last_error = None
217+
for attempt in range(5):
218+
delay = max(0, self._next_history_read - time.monotonic())
219+
if time.monotonic() + delay >= deadline:
220+
break
221+
if delay:
222+
time.sleep(delay)
223+
try:
224+
return self.lam.get_durable_execution_history(**request)
225+
except ClientError as error:
226+
if error.response["Error"]["Code"] not in {
227+
"TooManyRequestsException",
228+
"ThrottlingException",
229+
}:
230+
raise
231+
last_error = error
232+
finally:
233+
self._next_history_read = time.monotonic() + 1
234+
self._next_history_read = time.monotonic() + min(2**attempt, 4)
235+
raise CollectionError(
236+
"Execution-history API throttle budget exhausted"
237+
) from last_error
238+
199239
def history(self, item):
200240
events, marker = [], None
201241
while True:
202242
request = {"DurableExecutionArn": item["arn"], "IncludeExecutionData": True}
203243
if marker:
204244
request["Marker"] = marker
205-
result = self.lam.get_durable_execution_history(**request)
245+
result = self.history_page(**request)
206246
events.extend(result.get("Events", []))
207247
marker = result.get("NextMarker")
208248
if not marker:
@@ -229,25 +269,36 @@ def observed():
229269
message="No service invocation timeout evidence for " + request,
230270
)
231271

232-
def collect(self):
272+
def collect(self, *, full_run=True):
233273
errors = []
274+
if not full_run and not self.invocations:
275+
return
234276
try:
235-
events = self.refresh(markers=[])
277+
events = self.refresh(
278+
markers=[] if full_run else [i["marker"] for i in self.invocations],
279+
validate_controls=False,
280+
)
236281
# Independent collect steps can recover execution ARNs without the
237282
# pytest process or successful invoke response.
238283
items = {
239284
e["execution"]: {"arn": e["execution"], "marker": e["marker"]}
240285
for e in events
241286
}
242-
for item in items.values():
287+
except Exception as error:
288+
errors.append(str(error))
289+
items = {}
290+
if not full_run:
291+
items.update({i["arn"]: i for i in self.invocations})
292+
for item in items.values():
293+
try:
243294
self.history(item)
244295
save(
245296
ARTIFACTS / f"executions/{item['marker']}.json",
246297
self.lam.get_durable_execution(DurableExecutionArn=item["arn"]),
247298
)
248-
except Exception as error:
249-
errors.append(str(error))
250-
for key, arn in self.manifest["functions"].items():
299+
except Exception as error:
300+
errors.append(str(error))
301+
for key, arn in self.manifest["functions"].items() if full_run else []:
251302
try:
252303
config = self.lam.get_function_configuration(FunctionName=arn)
253304
logs = []

‎lmi_tests/tests/e2e/conftest.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,4 +50,4 @@ def cloud(deployment):
5050
try:
5151
driver.release_all()
5252
finally:
53-
driver.collect()
53+
driver.collect(full_run=False)

‎lmi_tests/tests/unit/test_cloud.py‎

Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44

55
import boto3
66
from botocore.stub import ANY, Stubber
7+
from botocore.exceptions import ClientError
78
import pytest
89

910
from lmi_tests import cloud as module
@@ -191,3 +192,92 @@ def test_polling_reads_only_current_case_and_surfaces_control_errors(driver):
191192
driver.s3.get_paginator.return_value.paginate.assert_called_with(
192193
Bucket="unit-bucket", Prefix="events/"
193194
)
195+
196+
197+
class Clock:
198+
def __init__(self):
199+
self.now = 0.0
200+
201+
def monotonic(self):
202+
return self.now
203+
204+
def time(self):
205+
return self.now
206+
207+
def sleep(self, seconds):
208+
self.now += seconds
209+
210+
211+
def throttle_error():
212+
return ClientError(
213+
{"Error": {"Code": "TooManyRequestsException", "Message": "Rate exceeded"}},
214+
"GetDurableExecutionHistory",
215+
)
216+
217+
218+
def test_history_retries_preserve_pages_and_pace_requests(driver, monkeypatch):
219+
clock = Clock()
220+
monkeypatch.setattr(module, "time", clock)
221+
driver.lam.get_durable_execution_history = Mock(
222+
side_effect=[
223+
throttle_error(),
224+
{"Events": [{"EventId": 1}], "NextMarker": "page2"},
225+
{"Events": [{"EventId": 2}]},
226+
]
227+
)
228+
assert driver.history({"arn": "execution", "marker": "case"}) == [
229+
{"EventId": 1},
230+
{"EventId": 2},
231+
]
232+
calls = driver.lam.get_durable_execution_history.call_args_list
233+
assert calls[0] == calls[1]
234+
assert calls[2].kwargs["Marker"] == "page2"
235+
assert clock.now >= 2
236+
237+
238+
def test_history_throttling_has_a_finite_budget(driver, monkeypatch):
239+
clock = Clock()
240+
monkeypatch.setattr(module, "time", clock)
241+
driver.lam.get_durable_execution_history = Mock(side_effect=throttle_error())
242+
with pytest.raises(CollectionError, match="throttle budget"):
243+
driver.history_page(DurableExecutionArn="execution")
244+
assert driver.lam.get_durable_execution_history.call_count == 5
245+
assert clock.now < 20
246+
247+
248+
def test_history_does_not_retry_permission_errors_or_invocations(driver, monkeypatch):
249+
monkeypatch.setattr(module, "time", Clock())
250+
driver.lam.invoke = Mock(
251+
return_value={"StatusCode": 202, "DurableExecutionArn": "execution"}
252+
)
253+
item = driver.start("success")
254+
error = ClientError(
255+
{"Error": {"Code": "AccessDeniedException"}}, "GetDurableExecutionHistory"
256+
)
257+
driver.lam.get_durable_execution_history = Mock(side_effect=error)
258+
with pytest.raises(ClientError):
259+
item["runner"].lambda_client.get_durable_execution_history(
260+
DurableExecutionArn="execution"
261+
)
262+
driver.lam.invoke.assert_called_once()
263+
driver.lam.get_durable_execution_history.assert_called_once()
264+
265+
266+
def test_case_collection_is_scoped_but_final_collection_is_complete(driver):
267+
driver.invocations = [{"arn": "execution", "marker": "case"}]
268+
driver.refresh = Mock(return_value=[{"execution": "execution", "marker": "case"}])
269+
driver.history = Mock(return_value=[])
270+
driver.lam.get_durable_execution = Mock(return_value={"Status": "SUCCEEDED"})
271+
driver.logs = Mock()
272+
driver.collect(full_run=False)
273+
driver.refresh.assert_called_once_with(markers=["case"], validate_controls=False)
274+
driver.history.assert_called_once()
275+
driver.logs.get_paginator.assert_not_called()
276+
driver.manifest["created"] = 0
277+
driver.lam.get_function_configuration = Mock(
278+
return_value={"LoggingConfig": {"LogGroup": "group"}}
279+
)
280+
driver.logs.get_paginator.return_value.paginate.return_value = []
281+
driver.collect()
282+
driver.refresh.assert_called_with(markers=[], validate_controls=False)
283+
driver.logs.get_paginator.assert_called_once_with("filter_log_events")

0 commit comments

Comments
 (0)