feat: add token bucket rate limiter for API throttling#18
Conversation
Add a manually-triggered GitHub Actions workflow that runs test_download_two_layers 5 times with debug instrumentation to diagnose why layers sometimes return 'Layer1' instead of 'Layer2' on Docker (but not Finch). Instrumentation logs: - Layer ordering in download_all() input/output - Generated Dockerfile ADD command sequence - Tarball construction order - Image build vs reuse decisions
Instead of running isolated iterations, run the exact same pytest command as the CI local-invoke suite with added debug logging. This reproduces the real test ordering and cross-test interactions. Enhanced instrumentation now also captures: - Per-layer download decisions (cached vs fresh) - Extracted layer file contents after download - Docker build stream output (cache hits show as CACHED)
Run just the TestLayerVersion class instead of the full local-invoke suite. This covers the failing test_download_two_layers plus the other tests in the class that may cause cross-test contamination.
debug: add workflow to investigate flaky layer ordering tests
debug: remove repo guard for fork
debug: use direct OIDC auth
debug: add BY_CANARY=true
…rm container tests * Update notify-slack.yml to add sleep (aws#8646) * fix: use tag prefix matching to clean up samcli/lambda-* images (aws#8647) Docker's images.list(name='samcli/lambda') does exact repository matching and won't match repositories like 'samcli/lambda-python'. This caused stale images to persist across parameterized test classes, leading to flaky test_download_two_layers failures where Layer2 should overwrite Layer1 but the image was never rebuilt. Fix by: 1. Adding _cleanup_samcli_images() that lists all images and filters by 'samcli/lambda-' tag prefix 2. Using this method in both tearDown and tearDownClass 3. Fixing the same pattern in TestLayerVersionThatDoNotCreateCache * fix: remove fallback in count_running_containers that causes flaky warm container tests The count_running_containers method had a fallback that returned the count of ALL SAM CLI containers when MODE env var filtering found no matches. This caused AssertionError: 3 != 2 when stale containers from other tests were present. Now it strictly counts only containers matching this test's unique MODE UUID and uses exact string matching. * fix: add AWS_DEFAULT_REGION and use -k pattern for parameterized tests * fix: add BY_CANARY=true to enable Docker tests on CI * fix: exclude RemoteLayers tests that need AWS credentials --------- Co-authored-by: seshubaws <116689586+seshubaws@users.noreply.github.com>
When build_in_source is used with Node.js, the build directory contains a node_modules symlink. During local invoke, SAM CLI resolves this symlink and creates an additional bind mount. Docker tolerates creating a mountpoint over a symlink, but Finch (containerd/runc) fails with 'not a directory'. This fix temporarily replaces symlinks with empty directories before container creation, then restores them afterward. This ensures: - Finch/runc gets a valid directory mountpoint - Docker continues to work as before - Repeated invocations work because symlinks are restored - Host filesystem is left unchanged even if container creation fails
fix: replace symlinks with dirs for Finch container mount compatibility
fix: restore BY_CANARY to run docker tests on CI
Containerd/Finch resolves bind mounts at start time, not create time. Moving the symlink restore to after start() ensures the empty directory mountpoints are still present when the container runtime sets up mounts.
fix: move symlink restore to after container start
Bind mounts must remain valid for the container's entire lifetime. Both Docker and Finch need the empty directory mountpoint to persist until the container is stopped and deleted. Restoring symlinks in the delete() method's finally block ensures proper cleanup alongside the existing host_tmp_dir cleanup.
fix: restore symlinks at container delete time
Implement a local CloudFormation Language Extensions processor supporting: - Fn::ForEach loop expansion in Resources, Conditions, and Outputs - Fn::Length, Fn::ToJsonString intrinsic functions - Fn::FindInMap with DefaultValue support - Conditional DeletionPolicy/UpdateReplacePolicy - Nested ForEach depth validation (max 5 levels) - Partial resolution mode preserving unresolvable references Pipeline architecture: TemplateParsingProcessor -> ForEachProcessor -> IntrinsicResolverProcessor -> DeletionPolicyProcessor -> UpdateReplacePolicyProcessor Includes comprehensive unit tests and CloudFormation compatibility suite.
Wire the language extensions library into SAM CLI with two-phase architecture: - Phase 1: expand_language_extensions() -> LanguageExtensionResult - Phase 2: SamTranslatorWrapper.run_plugins() (SAM transform only) Key components: - expand_language_extensions() canonical entry point with template-level cache keyed on (path, mtime, params_hash) - SamTranslatorWrapper receives pre-expanded template (Phase 2 only) - SamLocalStackProvider.get_stacks() calls expand_language_extensions() - SamTemplateValidator calls expand_language_extensions() - DynamicArtifactProperty dataclass for Mappings transformation - Fn::ForEach guards in artifact_exporter, normalizer, cdk/utils - clear_expansion_cache() for warm container file change events
- _get_template_for_output() preserves Fn::ForEach in build output - _update_foreach_artifact_paths() generates Mappings for dynamic artifact properties with per-function build paths - Recursive nested Fn::ForEach support - ForEach-aware path resolution skips Docker image URIs Test templates: static CodeUri, dynamic CodeUri, parameter collections, nested stacks, nested ForEach, dynamic ImageUri, depth validation.
Package:
- _export() calls expand_language_extensions() for Phase 1
- Preserves Fn::ForEach in packaged template with S3 URIs
- Generates Mappings for dynamic artifact properties
- _find_artifact_uri_for_resource() handles all export formats:
string, {S3Bucket,S3Key}, {Bucket,Key}, {ImageUri}
- Recursive nested Fn::ForEach support
- Warning for parameter-based collections
Deploy:
- Uploads original unexpanded template to CloudFormation
- Clear error for missing Mapping keys
Integration tests for CodeUri, ContentUri, DefinitionUri, ImageUri,
BodyS3Location across all packageable resource types.
- sam validate: valid ForEach, invalid syntax, cloud-dependent collections, dynamic CodeUri, nested depth validation (5 valid, 6 invalid) - sam local invoke: expanded function names from ForEach - sam local start-api: ForEach-generated API endpoints
Add make test-lang-ext and make test-all targets so the 1695 language extensions unit tests only run when needed, keeping the default make test fast for unrelated PRs.
The expand_language_extensions() cache stored references to template dicts that were later mutated in-place by ApplicationBuilder.update_template() (which changes nested stack Location properties to build-output paths). On cache hit, the mutated dict was returned, causing TemplateNotFoundException during the second infra sync in sam sync --watch. Remove the cache entirely since deep-copying on hit negates the performance benefit and adds complexity. Keep clear_expansion_cache() as a no-op for backward compatibility. Fixes TestSyncInfraNestedStacks_0 and TestSyncInfraNestedStacks_1 integration test failures.
fix: remove expansion cache to fix sync watch nested stack failures
fix: use OIDC credentials directly for sync test workflow
|
@cortex-assistant Please review this PR. |
3 similar comments
|
@cortex-assistant Please review this PR. |
|
@cortex-assistant Please review this PR. |
|
@cortex-assistant Please review this PR. |
There was a problem hiding this comment.
Review: feat: add token bucket rate limiter for API throttling
The core token bucket algorithm is well-implemented and the code is clean. However, there are a few issues that should be addressed before merging.
Issues
| # | Severity | Description |
|---|---|---|
| 1 | 🔴 Critical | rate=0 causes an unhandled ZeroDivisionError on line 41 (1.0 / self._rate). No input validation on rate or capacity. |
| 2 | 🔴 Major | No unit tests. Every existing utility in samcli/lib/utils/ has a corresponding test file in tests/unit/lib/utils/. A concurrency primitive needs tests for race conditions, timeout behavior, edge cases, and burst behavior. |
| 3 | 🟡 Major | The poll-sleep loop in acquire() is inefficient. Consider using threading.Condition with wait(timeout) — more idiomatic and avoids wasting CPU cycles. The project already uses threading primitives extensively (see lock_distributor.py). |
| 4 | 🟡 Minor | Sleep duration doesn't account for remaining time until deadline — can overshoot timeout by up to 0.1s. Should be min(sleep_interval, deadline - time.monotonic()). |
| 5 | 🟡 Minor | Missing type hints (_refill needs -> None), timeout should be Optional[float], and docstrings on __init__/_refill are missing — inconsistent with project conventions (see lock_distributor.py for reference). |
Summary
Please add input validation (raise ValueError for rate <= 0 or capacity <= 0), add unit tests, and consider replacing the spin-loop with threading.Condition. The style nits are minor but worth aligning with the rest of the codebase.
| the caller blocks until a token becomes available. | ||
| """ | ||
|
|
||
| def __init__(self, rate: float, capacity: int = 10): |
There was a problem hiding this comment.
Critical / Major: No input validation on rate or capacity. rate=0 will cause a ZeroDivisionError on line 41, and negative values produce nonsensical behavior.
Suggested fix:
def __init__(self, rate: float, capacity: int = 10):
if rate <= 0:
raise ValueError(f"rate must be positive, got {rate}")
if capacity <= 0:
raise ValueError(f"capacity must be positive, got {capacity}")
...Also, consider adding a docstring with Parameters section to match the project convention (see lock_distributor.py).
| self._tokens -= 1 | ||
| return True | ||
| if deadline and time.monotonic() >= deadline: | ||
| return False |
There was a problem hiding this comment.
Major: This poll-sleep loop is inefficient for a production concurrency primitive. Two issues:
1.0 / self._ratewill raiseZeroDivisionErrorifrate=0(see input validation comment above).- The sleep doesn't account for remaining time until deadline — can overshoot the timeout by up to 0.1s. Should be:
sleep_time = min(1.0 / self._rate, 0.1) if deadline: sleep_time = min(sleep_time, max(0, deadline - time.monotonic())) time.sleep(sleep_time)
Longer term, consider replacing with threading.Condition + wait(timeout) for a more efficient and idiomatic approach.
| self._tokens = min(self._capacity, self._tokens + elapsed * self._rate) | ||
| self._last_refill = now | ||
|
|
||
| def acquire(self, timeout: float = None) -> bool: |
There was a problem hiding this comment.
Minor: Type hint should be Optional[float] for clarity:
from typing import Optional
def acquire(self, timeout: Optional[float] = None) -> bool:This is consistent with the project's use of Optional in other utils.
|
I've completed the review of PR #18. Here's a summary of what I found: Verdict: Request Changes 🔴 The token bucket algorithm itself is solid, but there are a few issues to address:
I left 3 inline comments on the specific lines and labeled the PR as |
Adds a thread-safe token bucket rate limiter for throttling AWS API calls. Supports blocking and non-blocking acquire modes.