From 7569bd8250c1a435b55ded0941f31935e8c9e2c2 Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Fri, 24 Jul 2026 15:23:24 -0400 Subject: [PATCH] harden(webhook): resolve /sh-security-review findings (1 confirmed medium + cheap fixes) High-recall detector fan-out (injection/authz/secrets-crypto/iac-iam/logic) + proof-or-kill verifier. Gate PASSES: 1 confirmed medium, 0 confirmed critical/high. Confirmed finding fixed; several unverified-but-cheap hardenings applied since the emitter ships dark and activation is weeks out. - CONFIRMED medium (confused deputy): the rotation Lambda's generated invoke permission for secretsmanager.amazonaws.com carried no SourceAccount/SourceArn, so any account's Secrets Manager could invoke the rotator. Patched the generated CfnPermission in place (a second permission would be additive, not restrictive) to pin account + this secret ARN. - delivery + replay: refuse to follow receiver 3xx redirects (no-redirect opener) so live X-SH-* auth headers can't be forwarded to a receiver-chosen Location and an http:// Location can't slip past the https guard. Fixed the "unfollowed 3xx" comment that was factually wrong. - delivery: classify 401/403 as retryable (invalidate key cache + retry in order) instead of parking -- transient auth failures (rotation outran the TTL cache, clock skew) are availability events, not contract bugs. - envelope: build_event now genuinely total (guarded eventID / ApproximateCreationDateTime subscripts) per its own never-raise contract. - handler: catch-all so an unexpected per-record error (e.g. SQS park failure) reports only that record instead of failing the whole batch (which would re-deliver every earlier success for 24h); per-invocation emit/skip batch summary so a systemic silent drop is queryable/alarmable. - rotator: narrow the AWSCURRENT-read except to ResourceNotFound/JSONDecode (transient SM/KMS errors re-raise so the overlap key isn't silently dropped); kid uniqueness checked against ALL retained kids with a random suffix on collision (never reissue a kid for a different secret). - contract: skeleton-upsert required on ANY unknown work_order_id (not just comment-before-create) + monotonicity guard (ignore older updated_at), so a parked created or an out-of-order replay can't corrupt receiver state. Unverified/refuted findings left as-is with rationale: the two "high" logic claims (whole-batch crash triggers, ordering violation) were refuted on reachability (real stream records carry required fields; persistence writes strings only; full-state idempotent upsert absorbs the ordering gap). Signed kid/version binding (AUTHZ-002) declined: coordinated contract change, not cheap, no exploit with one algorithm/key. --- cdk/wo_stack.py | 38 +++++++++++++++++ docs/shoc-webhook-contract.md | 9 ++-- lambdas/wo/shoc_emitter/delivery.py | 52 +++++++++++++++++++++-- lambdas/wo/shoc_emitter/envelope.py | 11 ++++- lambdas/wo/shoc_emitter/handler.py | 56 ++++++++++++++++++++++--- lambdas/wo/shoc_hmac_rotator/handler.py | 16 +++++-- scripts/replay_shoc_webhooks.py | 18 +++++++- tests/test_shoc_emitter_delivery.py | 39 ++++++++++++++++- tests/test_shoc_emitter_handler.py | 34 +++++++++++++++ tests/test_shoc_hmac_rotator.py | 53 +++++++++++++++++++++-- 10 files changed, 304 insertions(+), 22 deletions(-) diff --git a/cdk/wo_stack.py b/cdk/wo_stack.py index cb02eb2..7d6bbc0 100644 --- a/cdk/wo_stack.py +++ b/cdk/wo_stack.py @@ -434,6 +434,32 @@ class WorkorderIngestStack(Stack): _add_shoc_webhook_emitter(self, work_orders_table, comments_table, alarm_topic) +def _scope_rotation_invoke_permission(stack, rotator_fn, secret): + """Add SourceAccount/SourceArn to the generated rotation invoke permission. + + ``add_rotation_schedule`` emits an ``AWS::Lambda::Permission`` for the + ``secretsmanager.amazonaws.com`` service principal with no source + conditions. Rather than add a second (additive) permission, find that + generated ``CfnPermission`` and pin it to this account + secret so only + this secret's Secrets Manager can invoke the rotator. + """ + patched = False + for child in stack.node.find_all(): + if ( + isinstance(child, lambda_.CfnPermission) + and child.principal == "secretsmanager.amazonaws.com" + and stack.resolve(child.function_name) + == stack.resolve(rotator_fn.function_arn) + ): + child.source_account = stack.account + child.source_arn = secret.secret_arn + patched = True + if not patched: # fail loud if CDK changes the generated shape on upgrade + raise RuntimeError( + "rotation invoke CfnPermission not found; cannot scope source conditions" + ) + + def _add_shoc_webhook_emitter(stack, work_orders_table, comments_table, alarm_topic): """SHOC webhook emitter (docs/shoc-webhook-plan.md Phases 1-4). @@ -609,6 +635,18 @@ def _add_shoc_webhook_emitter(stack, work_orders_table, comments_table, alarm_to rotation_lambda=shoc_hmac_rotator, automatically_after=Duration.days(30), ) + # Scope the Secrets-Manager-service invoke permission to THIS secret + # (cross-review FIX / confused-deputy): add_rotation_schedule emits an + # AWS::Lambda::Permission for secretsmanager.amazonaws.com with no + # SourceAccount/SourceArn, so any account's Secrets Manager could invoke + # the rotator by pointing a foreign secret's RotationLambdaARN at it. + # Lambda permissions are additive (OR), so a second scoped permission + # would NOT revoke the unscoped one -- patch the generated permission in + # place. source_arn pins the invoker to this secret; source_account is the + # belt-and-braces account bound. (Blast radius was already contained by + # the rotator role being resource-scoped to this secret, but this closes + # the unauthenticated invoke primitive per AWS rotation guidance.) + _scope_rotation_invoke_permission(stack, shoc_hmac_rotator, shoc_hmac_secret) # --- Standard per-Lambda alarms: workorder-shoc-hmac-rotator --- # errors + throttles + p99 duration. No DLQ alarm: rotation is invoked diff --git a/docs/shoc-webhook-contract.md b/docs/shoc-webhook-contract.md index d83af83..0eba8d3 100644 --- a/docs/shoc-webhook-contract.md +++ b/docs/shoc-webhook-contract.md @@ -90,8 +90,10 @@ Emitted from `WorkOrderComments` inserts — one per source email (comments, but ## 5. Ordering and delivery semantics - **At-least-once.** Duplicates are possible on retry; dedupe on `delivery_id`. -- **Per-work-order, per-event-family ordering is guaranteed:** all `work_order.*` state events for a given `work_order_id` arrive in commit order (a `cancelled` never precedes its `created`). Comment events are likewise ordered among themselves per work order. -- **Cross-family ordering is NOT guaranteed:** a `comment_added` may occasionally arrive before the `created` for its work order (separate streams). **Requirement:** on a comment for an unknown `work_order_id`, upsert a skeleton work order and let the state event backfill it — the same semantics the pipeline itself uses for out-of-order source emails. +- **Per-work-order, per-event-family ordering is guaranteed on the live feed:** all `work_order.*` state events for a given `work_order_id` arrive in commit order (a `cancelled` never precedes its `created`). Comment events are likewise ordered among themselves per work order. **Caveat:** this holds for delivered events. A rare non-retryable `4xx` on one event (a contract bug) parks *that* event and lets later events proceed, so the receiver can in principle see a `work_order.updated`/`cancelled` for a WO whose `created` was parked. Because every `work_order.*` body is **full current state** (not a diff), the two requirements below make this safe regardless. +- **Requirement — skeleton-upsert on ANY unknown `work_order_id`:** on any `work_order.*` event (state or comment) referencing a `work_order_id` you have not seen, upsert the record from the event's own `data` (or a skeleton for a bare comment) rather than dropping it. Full-state bodies mean a later event fully reconstructs a WO whose `created` never arrived — the same out-of-order semantics the pipeline itself uses for source emails. +- **Requirement — monotonicity (do not regress state):** ignore any `work_order.*` state event whose `data.updated_at` is **older** than the `updated_at` already stored for that WO. This makes an operator replay (see §8) or a rare out-of-order delivery idempotent and non-regressing: a stale snapshot can never overwrite newer state. Comments are append-only (dedupe on `comment_id`), so this applies only to WO state events. +- **Cross-family ordering is NOT guaranteed:** a `comment_added` may occasionally arrive before the `created` for its work order (separate streams) — covered by the skeleton-upsert requirement above. - Expected volume ≈ 22,900 events/month (~760/day); >90% are comment events. Bursts of a few events/second are possible. ## 6. Authentication — HMAC signature @@ -145,7 +147,8 @@ If SHOC is down longer than the retry window or deliveries are parked, an operat - [ ] Raw-body HMAC verification, constant-time compare, ±300s timestamp window, fail-closed - [ ] Cross-account secret fetch with ≤5-min cache + refresh-on-unknown-kid - [ ] Dedupe store on `delivery_id` (and/or `comment_id`) -- [ ] Skeleton-upsert on comment-before-create +- [ ] Skeleton-upsert on ANY unknown `work_order_id` (state events too, not just comment-before-create) +- [ ] Monotonicity guard: ignore WO state events whose `data.updated_at` is older than stored (§5) - [ ] Fast `2xx` ack (<10s), internal queue if processing is slow - [ ] Timestamp parsing accepts `+00:00` offsets - [ ] Reject unknown `schema_version` diff --git a/lambdas/wo/shoc_emitter/delivery.py b/lambdas/wo/shoc_emitter/delivery.py index d94643c..bce3414 100644 --- a/lambdas/wo/shoc_emitter/delivery.py +++ b/lambdas/wo/shoc_emitter/delivery.py @@ -44,6 +44,29 @@ USER_AGENT = "workorder-shoc-emitter/1" # 2xx = delivered; 429 and 5xx retry; everything else parks (contract sec. 7). _HTTP_SUCCESS_RANGE = range(HTTPStatus.OK, HTTPStatus.MULTIPLE_CHOICES) _HTTP_SERVER_ERROR_MIN = HTTPStatus.INTERNAL_SERVER_ERROR +# Transient auth failures: a stale cached key mid-rotation, a receiver +# secret-fetch blip, or clock skew past the +/-300s window. These are +# availability events, not payload contract bugs, so they retry (in order) +# after the key cache is dropped -- never park, which would strand the +# delivery out of the ordered feed until a manual replay. +_HTTP_AUTH_FAILURES = (HTTPStatus.UNAUTHORIZED, HTTPStatus.FORBIDDEN) + + +class _NoRedirectHandler(urllib.request.HTTPRedirectHandler): + """Refuse to follow receiver redirects. + + The default opener follows 3xx transparently, which would (a) forward the + live X-SH-* auth headers to a receiver-chosen Location and (b) let an + http:// Location slip past the https-only guard on the configured URL. We + only ever POST to the one configured endpoint; a redirect is a receiver + misconfiguration, so surface the 3xx as an HTTPError and let it park. + """ + + def redirect_request(self, *args, **kwargs): + return None + + +_opener = urllib.request.build_opener(_NoRedirectHandler) # Refresh the cached key material this often so a rotated secret propagates # without waiting for the execution environment to recycle (contract requires @@ -81,6 +104,19 @@ def _get_hmac_keys() -> list[dict]: and now - _hmac_keys_cached_at < _HMAC_KEYS_CACHE_TTL_SECONDS ): return _hmac_keys_cache + return _refresh_hmac_keys() + + +def _invalidate_hmac_keys() -> None: + """Drop the cached key material so the next sign re-fetches (rotation).""" + global _hmac_keys_cache + _hmac_keys_cache = None + + +def _refresh_hmac_keys() -> list[dict]: + """Force a fresh Secrets Manager fetch, bypassing the TTL cache.""" + global _hmac_keys_cache, _hmac_keys_cached_at + now = time.monotonic() secrets = boto3.client("secretsmanager") try: secret = secrets.get_secret_value(SecretId=HMAC_SECRET_ARN) @@ -138,7 +174,10 @@ def deliver(envelope: dict) -> tuple[str, int]: method="POST", ) try: - with urllib.request.urlopen(request, timeout=POST_TIMEOUT_SECONDS) as response: + # _opener refuses redirects, so a 3xx surfaces here as an HTTPError + # rather than silently re-issuing the request (with its auth headers) + # to a receiver-chosen Location. + with _opener.open(request, timeout=POST_TIMEOUT_SECONDS) as response: status_code = response.status except urllib.error.HTTPError as exc: status_code = exc.code @@ -149,11 +188,18 @@ def deliver(envelope: dict) -> tuple[str, int]: raise RetryableDeliveryError( f"receiver returned {status_code}", status_code=status_code ) from exc + if status_code in _HTTP_AUTH_FAILURES: + # Drop the cached key so an in-order retry re-signs with the + # current secret (handles a rotation that outran the TTL cache). + _invalidate_hmac_keys() + raise RetryableDeliveryError( + f"receiver auth failure {status_code}", status_code=status_code + ) from exc return ("rejected", status_code) except (TimeoutError, urllib.error.URLError, OSError) as exc: raise RetryableDeliveryError(f"connection error: {exc}") from exc if status_code in _HTTP_SUCCESS_RANGE: return ("delivered", status_code) - # A non-2xx that urlopen did not raise for (e.g. an unfollowed 3xx): - # contract bug on someone's side -- park it, never block the shard. + # A non-2xx that the opener did not raise for -- contract bug on + # someone's side. Park it, never block the shard. return ("rejected", status_code) diff --git a/lambdas/wo/shoc_emitter/envelope.py b/lambdas/wo/shoc_emitter/envelope.py index 2682486..5b9997e 100644 --- a/lambdas/wo/shoc_emitter/envelope.py +++ b/lambdas/wo/shoc_emitter/envelope.py @@ -141,13 +141,20 @@ def build_event(record: dict) -> dict | None: fields = COMMENT_DATA_FIELDS else: fields = WO_DATA_FIELDS + # eventID and ApproximateCreationDateTime are present on every real + # DynamoDB stream record; guarding them keeps this function total (the + # "never raise" contract above) rather than trusting a hard subscript. + event_id = record.get("eventID") + approx_creation = stream.get("ApproximateCreationDateTime") + if event_id is None or approx_creation is None: + return None # ApproximateCreationDateTime arrives as epoch seconds (float/Decimal). occurred_at = datetime.fromtimestamp( - float(stream["ApproximateCreationDateTime"]), tz=timezone.utc + float(approx_creation), tz=timezone.utc ).isoformat() return { "schema_version": SCHEMA_VERSION, - "delivery_id": record["eventID"], + "delivery_id": event_id, "event_type": event_type, "occurred_at": occurred_at, "source": SOURCE, diff --git a/lambdas/wo/shoc_emitter/handler.py b/lambdas/wo/shoc_emitter/handler.py index 5225508..1f3a920 100644 --- a/lambdas/wo/shoc_emitter/handler.py +++ b/lambdas/wo/shoc_emitter/handler.py @@ -88,20 +88,49 @@ def _park_rejected(event: dict, status_code: int): ) +def _batch_summary(records: int, emitted: int, skipped: int) -> str: + """Per-invocation emit/skip tally. + + Skips (REMOVE, echo guard, comment MODIFYs, unmappable records) return no + delivery and would otherwise be invisible: a systemic classification break + -- e.g. a table rename desyncing the eventSourceARN parse -- would drop + 100%% of events while every batch still reports success. Logging the ratio + makes that queryable in Logs Insights and alarmable. + """ + return json.dumps( + { + "event": "shoc_batch_summary", + "records": records, + "emitted": emitted, + "skipped": skipped, + } + ) + + def handler(event, context): """Lambda entry point. Triggered by the two WO-table stream ESMs.""" - for record in event.get("Records", []): + records = event.get("Records", []) + emitted = 0 + skipped = 0 + for record in records: webhook_event = envelope.build_event(record) if webhook_event is None: + skipped += 1 continue + emitted += 1 start = time.monotonic() try: outcome, status_code = delivery.deliver(webhook_event) + latency_ms = int((time.monotonic() - start) * 1000) + logger.info(_delivery_log(webhook_event, status_code, latency_ms, outcome)) + if outcome == "rejected": + _park_rejected(webhook_event, status_code) except delivery.RetryableDeliveryError as exc: latency_ms = int((time.monotonic() - start) * 1000) logger.warning( _delivery_log(webhook_event, exc.status_code, latency_ms, "retryable") ) + logger.info(_batch_summary(len(records), emitted, skipped)) # Stop here: reporting this record's sequence number makes the ESM # retry from it in order; earlier successes are not re-delivered. return { @@ -109,8 +138,25 @@ def handler(event, context): {"itemIdentifier": record["dynamodb"]["SequenceNumber"]} ] } - latency_ms = int((time.monotonic() - start) * 1000) - logger.info(_delivery_log(webhook_event, status_code, latency_ms, outcome)) - if outcome == "rejected": - _park_rejected(webhook_event, status_code) + except Exception: + # Catch-all so an unexpected error (e.g. an SQS park failure) + # cannot escape and fail the WHOLE invocation -- that would make + # the ESM re-deliver every earlier success in the batch for up to + # 24h. Report only THIS record so the ESM retries from it in order. + logger.exception( + json.dumps( + { + "event": "shoc_delivery_error", + "delivery_id": webhook_event["delivery_id"], + "event_type": webhook_event["event_type"], + } + ) + ) + logger.info(_batch_summary(len(records), emitted, skipped)) + return { + "batchItemFailures": [ + {"itemIdentifier": record["dynamodb"]["SequenceNumber"]} + ] + } + logger.info(_batch_summary(len(records), emitted, skipped)) return {"batchItemFailures": []} diff --git a/lambdas/wo/shoc_hmac_rotator/handler.py b/lambdas/wo/shoc_hmac_rotator/handler.py index 7c05d56..8bcaad0 100644 --- a/lambdas/wo/shoc_hmac_rotator/handler.py +++ b/lambdas/wo/shoc_hmac_rotator/handler.py @@ -61,7 +61,12 @@ def _current_keys(client, secret_id: str) -> list: try: raw = client.get_secret_value(SecretId=secret_id, VersionStage="AWSCURRENT") parsed = json.loads(raw["SecretString"]) - except Exception: + except (client.exceptions.ResourceNotFoundException, json.JSONDecodeError): + # Genuinely-absent or corrupt current value -> start a fresh list. + # NB: this does NOT catch transient failures (throttling, KMS/IAM + # blips): those re-raise so Secrets Manager marks the rotation failed + # and retries with the prior AWSCURRENT intact, rather than silently + # dropping the overlap key and stranding in-flight deliveries. logger.warning( json.dumps( {"event": "hmac_rotation_current_unreadable", "fallback": "empty"} @@ -110,9 +115,14 @@ def _create_secret(client, secret_id: str, token: str) -> None: pass current_keys = _current_keys(client, secret_id) now = datetime.now(timezone.utc) + existing_kids = {key["kid"] for key in current_keys} new_kid = now.strftime(KID_FORMAT) - if current_keys and current_keys[0]["kid"] == new_kid: - new_kid = now.strftime(KID_FORMAT_INTRA_HOUR) + if new_kid in existing_kids: + # Collision against ANY retained kid (not just keys[0]): repeated + # same-hour forced rotations must never reissue a kid for a different + # secret, or a receiver's cached kid->secret map desyncs. Append a + # random suffix guaranteed distinct from the retained set. + new_kid = f"{now.strftime(KID_FORMAT_INTRA_HOUR)}{secrets.token_hex(2)}" new_key = {"kid": new_kid, "secret": secrets.token_hex(SECRET_NUM_BYTES)} client.put_secret_value( SecretId=secret_id, diff --git a/scripts/replay_shoc_webhooks.py b/scripts/replay_shoc_webhooks.py index 6c8a0a1..b41a8a1 100644 --- a/scripts/replay_shoc_webhooks.py +++ b/scripts/replay_shoc_webhooks.py @@ -79,6 +79,22 @@ HTTP_2XX_MAX = 300 HTTP_TOO_MANY_REQUESTS = 429 HTTP_5XX_MIN = 500 + +class _NoRedirectHandler(urllib.request.HTTPRedirectHandler): + """Refuse to follow receiver redirects (matches the emitter's opener). + + Following a 3xx would forward the live X-SH-* auth headers to a + receiver-chosen Location and could downgrade an http:// target past the + https-only --url check; a redirect is a receiver misconfiguration, so let + it surface as an HTTPError status. + """ + + def redirect_request(self, *args, **kwargs): + return None + + +_opener = urllib.request.build_opener(_NoRedirectHandler) + # data payload field lists -- pinned to the emitter's envelope module. Absent # attributes are emitted as null. source_email_s3_key is deliberately # EXCLUDED (internal S3 pointer, never leaves the account). @@ -175,7 +191,7 @@ def post_event(url, raw_body_bytes, kid, secret_hex): method="POST", ) try: - with urllib.request.urlopen(request, timeout=POST_TIMEOUT_SECONDS) as resp: + with _opener.open(request, timeout=POST_TIMEOUT_SECONDS) as resp: return resp.status except urllib.error.HTTPError as exc: return exc.code diff --git a/tests/test_shoc_emitter_delivery.py b/tests/test_shoc_emitter_delivery.py index efdbe14..16f0cc3 100644 --- a/tests/test_shoc_emitter_delivery.py +++ b/tests/test_shoc_emitter_delivery.py @@ -125,7 +125,9 @@ def signing_keys(monkeypatch, delivery): def _patch_urlopen(monkeypatch, delivery, fn): - monkeypatch.setattr(delivery.urllib.request, "urlopen", fn) + # deliver() posts through the no-redirect opener, not the module-level + # urlopen, so patch the opener's open method. + monkeypatch.setattr(delivery._opener, "open", fn) @pytest.mark.parametrize("status", [200, 204]) @@ -177,6 +179,41 @@ def test_other_4xx_is_rejected_not_raised(monkeypatch, delivery, signing_keys, s assert delivery.deliver(ENVELOPE) == ("rejected", status) +@pytest.mark.parametrize("status", [401, 403]) +def test_auth_failures_are_retryable_and_invalidate_cache( + monkeypatch, delivery, signing_keys, status +): + # A transient auth failure (stale cached key mid-rotation, receiver + # secret-fetch blip, clock skew) must retry in order -- NOT park -- and + # drop the key cache so the retry re-signs with the current secret. + invalidated = {"called": False} + + def _mark(*_a, **_k): + invalidated["called"] = True + + monkeypatch.setattr(delivery, "_invalidate_hmac_keys", _mark) + + def _raise(request, timeout): + raise urllib.error.HTTPError( + delivery.SHOC_WEBHOOK_URL, status, "unauthorized", None, io.BytesIO(b"") + ) + + _patch_urlopen(monkeypatch, delivery, _raise) + with pytest.raises(delivery.RetryableDeliveryError) as exc: + delivery.deliver(ENVELOPE) + assert exc.value.status_code == status + assert invalidated["called"] is True + + +def test_redirects_are_not_followed(): + # The opener must refuse 3xx so auth headers are never forwarded to a + # receiver-chosen Location. redirect_request returning None makes urllib + # raise instead of following. + delivery = load_lambda_module("wo", "shoc_emitter/delivery") + handler = delivery._NoRedirectHandler() + assert handler.redirect_request(None, None, 302, "Found", {}, "http://evil") is None + + def test_request_signed_with_keys0_and_contract_headers( monkeypatch, delivery, signing_keys ): diff --git a/tests/test_shoc_emitter_handler.py b/tests/test_shoc_emitter_handler.py index 3a62328..1a4a713 100644 --- a/tests/test_shoc_emitter_handler.py +++ b/tests/test_shoc_emitter_handler.py @@ -171,6 +171,40 @@ def test_skip_records_produce_no_delivery_calls(monkeypatch, handler_mod): assert fake_sqs.sent == [] +def test_unexpected_error_reports_record_not_whole_batch(monkeypatch, handler_mod): + # An unexpected exception (here: SQS park failure on a 4xx rejection) must + # NOT escape the loop -- that would fail the invocation and make the ESM + # re-deliver every earlier success for 24h. The offending record is + # reported so the ESM retries from it in order; evt-1 is not re-delivered. + event = { + "Records": [ + _wo_record("evt-1", "401"), + _wo_record("evt-2", "402"), + _wo_record("evt-3", "403"), + ] + } + attempted, fake_sqs = _wire( + monkeypatch, + handler_mod, + { + "evt-1": ("delivered", 200), + "evt-2": ("rejected", 400), + "evt-3": ("delivered", 200), + }, + ) + + def _boom(QueueUrl, MessageBody): # noqa: N803 (boto3 kwargs) + raise RuntimeError("sqs unavailable") + + monkeypatch.setattr(fake_sqs, "send_message", _boom) + + result = handler_mod.handler(event, None) + + assert result == {"batchItemFailures": [{"itemIdentifier": "402"}]} + # Stopped at the failing record; evt-3 not attempted. + assert attempted == ["evt-1", "evt-2"] + + def test_empty_batch_returns_no_failures(monkeypatch, handler_mod): attempted, _ = _wire(monkeypatch, handler_mod, {}) assert handler_mod.handler({"Records": []}, None) == {"batchItemFailures": []} diff --git a/tests/test_shoc_hmac_rotator.py b/tests/test_shoc_hmac_rotator.py index 1fd7b94..919b37c 100644 --- a/tests/test_shoc_hmac_rotator.py +++ b/tests/test_shoc_hmac_rotator.py @@ -191,19 +191,41 @@ def test_create_is_noop_when_token_already_current(monkeypatch, rotator, fixed_n assert fake.put_calls == [] -def test_create_kid_collision_same_hour_gets_minute_suffix( +def test_create_kid_collision_same_hour_gets_distinct_suffix( monkeypatch, rotator, fixed_now ): - # Forced re-rotation within the same hour: keys[0].kid already equals the - # hour-format kid, so the new kid extends to minutes for uniqueness. + # Forced re-rotation within the same hour: the hour-format kid already + # exists, so the new kid extends to minutes + a random suffix, guaranteed + # distinct from every retained kid (never reissue a kid for a new secret). fake = _bootstrap_fake({"keys": [{"kid": "2026-07-24T15", "secret": "c" * 64}]}) _run(monkeypatch, rotator, fake, "createSecret") staged = json.loads(fake.put_calls[0]["SecretString"]) - assert staged["keys"][0]["kid"] == "2026-07-24T1507" + new_kid = staged["keys"][0]["kid"] + assert new_kid.startswith("2026-07-24T1507") + assert new_kid != "2026-07-24T15" assert staged["keys"][1]["kid"] == "2026-07-24T15" +def test_create_kid_collision_against_second_key_also_avoided( + monkeypatch, rotator, fixed_now +): + # Uniqueness is checked against ALL retained kids, not just keys[0]: if the + # hour-format kid matches keys[1] (a third same-hour rotation), it still + # gets a distinct suffix rather than being reissued. + fake = _bootstrap_fake( + { + "keys": [ + {"kid": "2026-07-24T1500ab", "secret": "a" * 64}, + {"kid": "2026-07-24T15", "secret": "b" * 64}, + ] + } + ) + _run(monkeypatch, rotator, fake, "createSecret") + new_kid = json.loads(fake.put_calls[0]["SecretString"])["keys"][0]["kid"] + assert new_kid not in {"2026-07-24T15", "2026-07-24T1500ab"} + + def test_create_malformed_current_value_starts_fresh_list( monkeypatch, rotator, fixed_now ): @@ -216,6 +238,29 @@ def test_create_malformed_current_value_starts_fresh_list( assert HEX_64_RE.fullmatch(staged["keys"][0]["secret"]) +def test_create_transient_current_read_error_propagates_not_swallowed( + monkeypatch, rotator, fixed_now +): + # A transient AWSCURRENT read failure (throttle/KMS blip) must NOT be + # swallowed into an empty key list -- that would drop the overlap key and + # strand in-flight deliveries. It re-raises so Secrets Manager fails and + # retries the rotation with the prior AWSCURRENT intact. + class _ThrottlingFake(FakeSecretsManager): + def get_secret_value(self, SecretId, VersionId=None, VersionStage=None): # noqa: N803 + if VersionStage == "AWSCURRENT": + raise RuntimeError("ThrottlingException") + return super().get_secret_value( + SecretId, VersionId=VersionId, VersionStage=VersionStage + ) + + fake = _ThrottlingFake( + {"cur-1": {"stages": {"AWSCURRENT"}, "value": json.dumps({"keys": []})}} + ) + with pytest.raises(RuntimeError, match="ThrottlingException"): + _run(monkeypatch, rotator, fake, "createSecret") + assert fake.put_calls == [] + + # --- testSecret --------------------------------------------------------------