mirror of
https://github.com/Sea-Haven-Industries/procurement-ingest.git
synced 2026-09-30 06:03:14 +00:00
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.
This commit is contained in:
parent
2f2fc83a82
commit
7569bd8250
10 changed files with 304 additions and 22 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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`
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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": []}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
):
|
||||
|
|
|
|||
|
|
@ -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": []}
|
||||
|
|
|
|||
|
|
@ -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 --------------------------------------------------------------
|
||||
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue