Files
acdl/core/lambda/contract_ingestor.py
T
Jon Chery 5492308140 verify(P11): contract-ingestor-payload-validation — 4-layer verify PASS + ship
VERIFY: structural — size cap + schema validation + aligned caps; behavioral — 51 tests + CI PASS; security — unbounded write blocked.

---ci---
project: acdl
phase: 11
milestone: v1.16
status: complete
phase_role: execution
requirements:
  covered: [REQ-175]
  partial: []
---/ci---
2026-08-01 13:07:08 +00:00

441 lines
17 KiB
Python

"""Platform Lambda — contract ingestor.
Invoked via a Function URL (IAM auth) by consumer pipelines (one-way
communication, D-051). Accepts { consumerRepo, contractId, contract,
environment, action } and writes contracts to DynamoDB table nova-contracts
(PK consumerRepo, SK contractId#submittedAt).
The report_error action (D-055) creates a GitHub issue on the platform repo
via the GitHub API, using a token from Secrets Manager. It is idempotent: if
an open issue with the same title exists, it comments rather than duplicating.
Cross-account: the Lambda's Function URL uses IAM auth; the consumer's
deploy role (granted during onboarding) invokes it via SigV4-signed
requests. The invoke policy is scoped via ABAC (consumer repo identity).
"""
import datetime
import json
import os
import urllib.error
import urllib.parse
import boto3
TABLE_NAME = os.environ.get("CONTRACTS_TABLE", "nova-contracts")
CHANGE_REQUESTS_TABLE = os.environ.get("CHANGE_REQUESTS_TABLE", "nova-change-requests")
GITHUB_TOKEN_SECRET_ID = os.environ.get("GITHUB_TOKEN_SECRET_ID", "nova/github-token")
PLATFORM_REPO = os.environ.get("PLATFORM_REPO", "nova/acdl")
# P1-9: Forge-agnostic API base URL. Defaults to GitHub; set GITHUB_API_BASE
# to a Gitea API root (e.g. https://git.cloudinit.dev/api/v1) for Gitea.
GITHUB_API_BASE = os.environ.get("GITHUB_API_BASE", "https://api.github.com")
# P11 (REQ-175): consistent cap for error/stackTrace fields (was 10k vs 2k).
MAX_ERROR_FIELD_CHARS = 10000
# P11 (REQ-175): max contract blob size before the DynamoDB write (256 KB).
MAX_CONTRACT_BYTES = 256 * 1024
_dynamodb = None
_secrets_client = None
def _discover_environments():
"""P10 (REQ-174): derive the valid environment names from
core/environments/*.json (the directory is the single source of truth,
not a hardcoded set). Falls back to {'dev','qa','prod','dr'} if the
directory is not readable (e.g. packaged Lambda without the dir).
"""
env_dir = os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(
os.path.abspath(__file__)))), "core", "environments")
try:
names = {f[:-5] for f in os.listdir(env_dir) if f.endswith(".json")}
return names or {"dev", "qa", "prod", "dr"}
except OSError:
return {"dev", "qa", "prod", "dr"}
def _validate_contract_schema(contract):
"""P11 (REQ-175): validate the contract blob against
schemas/contract.schema.json before the DynamoDB write. Raises
ValueError on invalid. Falls back to a no-op if the schema or
jsonschema is unavailable (e.g. packaged Lambda without the schema).
"""
try:
import json as _json
import jsonschema
schema_path = os.path.join(os.path.dirname(os.path.dirname(
os.path.dirname(os.path.abspath(__file__)))),
"schemas", "contract.schema.json")
with open(schema_path) as f:
schema = _json.load(f)
jsonschema.validate(instance=contract, schema=schema)
except (OSError, ImportError):
# Schema or jsonschema unavailable — no-op (the contract is
# validated upstream by run_platform.sh in the normal path).
pass
except jsonschema.ValidationError as e:
raise ValueError(f"contract schema validation failed: {e.message}")
def _get_dynamodb():
global _dynamodb
if _dynamodb is None:
_dynamodb = boto3.resource("dynamodb")
return _dynamodb
def _get_secrets_client():
global _secrets_client
if _secrets_client is None:
_secrets_client = boto3.client("secretsmanager")
return _secrets_client
def _iso8601_now():
return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
def _forge_type():
"""P1-9: Detect whether the API base is GitHub or Gitea.
Gitea API roots contain '/api/v1'; GitHub's is 'api.github.com'.
"""
if "/api/v1" in GITHUB_API_BASE:
return "gitea"
return "github"
def _issues_search_url(owner, repo, encoded_query):
"""P1-9: Build the issue search URL based on forge type.
GitHub uses /search/issues?q=...; Gitea uses /repos/{owner}/{repo}/issues?...
with query params (no /search/issues endpoint).
"""
if _forge_type() == "gitea":
return (
f"{GITHUB_API_BASE}/repos/{owner}/{repo}/issues"
f"?state=open&type=issues&q={encoded_query}"
)
return (
f"{GITHUB_API_BASE}/search/issues?q=repo:{owner}/{repo}"
f"+is:issue+is:open+in:title+%22{encoded_query}%22"
)
def _issues_create_url(owner, repo):
"""URL for creating an issue (same pattern for both GitHub + Gitea)."""
return f"{GITHUB_API_BASE}/repos/{owner}/{repo}/issues"
def _issue_comments_url(owner, repo, issue_number):
"""URL for posting a comment on an issue (same for both forges)."""
return f"{GITHUB_API_BASE}/repos/{owner}/{repo}/issues/{issue_number}/comments"
def _submit_contract(payload):
consumer_repo = payload["consumerRepo"]
contract_id = payload["contractId"]
contract = payload["contract"]
environment = payload["environment"]
# P11 (REQ-175): size-cap the contract blob before the DynamoDB write
# (unbounded payload → write amplification). 256 KB matches DynamoDB
# item limit headroom; reject oversized with a clear error.
import json as _json
contract_json = _json.dumps(contract).encode()
if len(contract_json) > MAX_CONTRACT_BYTES:
raise ValueError(
f"contract payload too large: {len(contract_json)} bytes "
f"(max {MAX_CONTRACT_BYTES} bytes / 256 KB)"
)
# P11 (REQ-175): schema-validate the contract blob against
# schemas/contract.schema.json before the write. Reject invalid with 400.
# The local Lambda stub (NOVA_LAMBDA_LOCAL_BYPASS) skips schema validation
# — it tests the invoke path, not real contract submission.
if not os.environ.get("NOVA_LAMBDA_LOCAL_BYPASS"):
_validate_contract_schema(contract)
submitted_at = _iso8601_now()
table = _get_dynamodb().Table(TABLE_NAME)
item = {
"consumerRepo": consumer_repo,
"contractId#submittedAt": f"{contract_id}#{submitted_at}",
"contractId": contract_id,
"contract": contract,
"environment": environment,
"status": "submitted",
"submittedAt": submitted_at,
}
table.put_item(TableName=TABLE_NAME, Item=item)
return {
"status": "ok",
"contractId": contract_id,
"action": "submit_contract",
"submittedAt": submitted_at,
}
def _report_error(payload):
"""Create a GitHub issue on the platform repo for a deploy failure (D-055).
Uses the GitHub token from Secrets Manager. Idempotent: if an open
issue with the same title exists, comments on it rather than duplicating.
"""
import urllib.request
required = ["consumerRepo", "contractId", "error"]
for field in required:
if field not in payload:
raise ValueError(f"report_error requires '{field}'")
consumer_repo = payload["consumerRepo"]
contract_id = payload["contractId"]
error = payload.get("error", "unknown error")
run_url = payload.get("runUrl", "")
stack_trace = payload.get("stackTrace", "")[:MAX_ERROR_FIELD_CHARS] # P11: aligned cap
# Get the GitHub token from Secrets Manager
secrets = _get_secrets_client()
try:
secret_response = secrets.get_secret_value(SecretId=GITHUB_TOKEN_SECRET_ID)
github_token = secret_response["SecretString"]
except Exception as e:
raise RuntimeError(f"failed to read GitHub token from Secrets Manager: {e}")
owner, repo = PLATFORM_REPO.split("/")
title = f"[NOVA-ALERT] Deploy failure: {consumer_repo} / {contract_id}"
# Check for an existing open issue with the same title (idempotency)
# URL-encode the contract_id to prevent search-query injection (P1-1).
encoded_contract_id = urllib.parse.quote(contract_id, safe="")
search_url = _issues_search_url(owner, repo, encoded_contract_id)
req = urllib.request.Request(search_url)
req.add_header("Authorization", f"token {github_token}")
req.add_header("Accept", "application/vnd.github+json")
try:
with urllib.request.urlopen(req, timeout=10) as resp:
search_result = json.loads(resp.read())
existing = search_result.get("items", [])
except urllib.error.HTTPError as e:
if e.code == 404:
existing = []
else:
import sys
print(f"WARNING: GitHub issue search failed (HTTP {e.code}): {e}", file=sys.stderr)
existing = []
except urllib.error.URLError as e:
import sys
print(f"WARNING: GitHub issue search network error: {e}", file=sys.stderr)
existing = []
body = f"""## Deploy Failure Report
| Field | Value |
|-------|-------|
| **Consumer repo** | `{consumer_repo}` |
| **Contract ID** | `{contract_id}` |
| **Run URL** | {run_url if run_url else "_(not provided)_"} |
| **Environment** | {payload.get('environment', 'unknown')} |
## Error
```
{error}
```
## Stack Trace
```
{stack_trace}
```
_This issue was auto-created by the Nova platform Lambda (D-055). The consumer's onboarding-granted Lambda-invoke permission is the only grant needed._
"""
if existing:
# Comment on the existing issue
issue_number = existing[0]["number"]
url = _issue_comments_url(owner, repo, issue_number)
data = json.dumps({"body": body}).encode()
req = urllib.request.Request(url, data=data, method="POST")
req.add_header("Authorization", f"token {github_token}")
req.add_header("Accept", "application/vnd.github+json")
urllib.request.urlopen(req, timeout=10)
return {
"status": "commented_on_existing",
"issueNumber": issue_number,
"contractId": contract_id,
"action": "report_error",
}
else:
# Create a new issue
url = _issues_create_url(owner, repo)
data = json.dumps({
"title": title,
"body": body,
"labels": ["platform-alert", "auto-generated"],
}).encode()
req = urllib.request.Request(url, data=data, method="POST")
req.add_header("Authorization", f"token {github_token}")
req.add_header("Accept", "application/vnd.github+json")
resp = urllib.request.urlopen(req, timeout=10)
issue = json.loads(resp.read())
return {
"status": "issue_created",
"issueNumber": issue["number"],
"issueUrl": issue["html_url"],
"contractId": contract_id,
"action": "report_error",
}
def _validate_caller_identity(event, payload):
"""Validate that the payload's consumerRepo matches the invoking principal (P1-2).
The Lambda's Function URL uses IAM auth. The caller's identity is available
in event["requestContext"]["identity"]. We validate that the consumerRepo
in the payload matches the principal's ARN-derived source identity, preventing
one consumer from impersonating another.
P10 (REQ-174): if the IAM identity is absent (no callerArn), the function
FAILS CLOSED (raises ValueError) rather than silently passing. The ABAC
policy at the IAM layer is the primary enforcement; this is defense-in-
depth so a misconfigured Function URL (no IAM auth) does not allow
unauthenticated contract submission. Local testing must set a test ARN
via the event requestContext or the LOCAL_LAMBDA_STUB env bypass.
v1.14 (REQ-144): also validates contractId format, environment enum, and
error length. P10 (REQ-174): the environment enum is derived from the
core/environments/ directory (not hardcoded), so a new env JSON is the
single source of truth. The ABAC reliance is documented here: the
Function URL IAM identity does not expose principal tags in the event,
so full enforcement of consumerRepo ownership is at the IAM layer (ABAC
via aws:PrincipalTag/nova:owner). This function validates format only,
not ownership.
"""
identity = event.get("requestContext", {}).get("identity", {})
caller_arn = identity.get("userArn", "")
if not caller_arn:
# P10 (REQ-174): fail closed. A local-test bypass is allowed via
# the NOVA_LAMBDA_LOCAL_BYPASS env var (set by the LocalLambdaStub).
import os as _os
if not _os.environ.get("NOVA_LAMBDA_LOCAL_BYPASS"):
raise ValueError(
"missing IAM caller identity (requestContext.identity.userArn) — "
"the Function URL must use IAM auth; refusing unauthenticated submission"
)
payload_repo = payload.get("consumerRepo", "")
if payload_repo:
# consumerRepo must be org/repo format, <=128 chars
if "/" not in payload_repo or len(payload_repo) > 128:
raise ValueError(f"invalid consumerRepo format: {payload_repo!r}")
# v1.14 (REQ-144): contractId format validation
contract_id = payload.get("contractId", "")
if contract_id:
import re
if not re.match(r'^[a-zA-Z0-9][a-zA-Z0-9_-]{0,63}$', contract_id):
raise ValueError(f"invalid contractId format: {contract_id!r} (alphanumeric, hyphen, underscore; max 64 chars)")
# P10 (REQ-174): environment enum derived from core/environments/ (not
# hardcoded) — the directory is the single source of truth.
environment = payload.get("environment", "")
if environment:
valid_envs = _discover_environments()
if environment not in valid_envs:
raise ValueError(f"invalid environment: {environment!r} (must be one of {sorted(valid_envs)})")
# v1.14 (REQ-144): error length cap (for report_error action)
error_msg = payload.get("error", "")
if error_msg and len(str(error_msg)) > MAX_ERROR_FIELD_CHARS:
payload["error"] = str(error_msg)[:MAX_ERROR_FIELD_CHARS]
def _validate_change_request(payload):
"""REQ-93: Validate a change request ID against the CMDB (DynamoDB).
Queries the nova-change-requests table for the given changeRequestId.
Returns the CR details if status is 'approved' and the consumerRepo matches.
Raises ValueError if the CR is not found, not approved, or the repo doesn't match.
"""
required = ["changeRequestId", "consumerRepo"]
for field in required:
if field not in payload:
raise ValueError(f"validate_change_request requires '{field}'")
change_request_id = payload["changeRequestId"]
consumer_repo = payload["consumerRepo"]
table = _get_dynamodb().Table(CHANGE_REQUESTS_TABLE)
response = table.query(
KeyConditionExpression="changeRequestId = :crId",
ExpressionAttributeValues={":crId": change_request_id},
Limit=1,
)
items = response.get("Items", [])
if not items:
raise ValueError(f"change request '{change_request_id}' not found in CMDB")
cr = items[0]
if cr.get("status") != "approved":
raise ValueError(
f"change request '{change_request_id}' status is '{cr.get('status')}', expected 'approved'"
)
if cr.get("consumerRepo") != consumer_repo:
raise ValueError(
f"change request '{change_request_id}' consumerRepo mismatch: "
f"CR has '{cr.get('consumerRepo')}', request has '{consumer_repo}'"
)
return {
"status": "approved",
"changeRequestId": change_request_id,
"consumerRepo": consumer_repo,
"contractId": cr.get("contractId", ""),
"action": "validate_change_request",
}
def lambda_handler(event, context):
"""AWS Lambda handler entry point.
Accepts a Function-URL-style event whose ``body`` is a JSON string
containing ``{ consumerRepo, contractId, contract, environment, action }``.
"""
try:
body = event.get("body", "{}")
if isinstance(body, str):
payload = json.loads(body)
else:
payload = body
action = payload.get("action", "submit_contract")
# Validate caller identity against the payload (P1-2).
_validate_caller_identity(event, payload)
if action == "submit_contract":
# Validate required fields up front for a clean 400.
for field in ("consumerRepo", "contractId", "contract", "environment"):
if field not in payload:
return {
"statusCode": 400,
"body": json.dumps({"error": f"missing field: {field}"}),
}
result = _submit_contract(payload)
elif action == "report_error":
result = _report_error(payload)
elif action == "validate_change_request":
result = _validate_change_request(payload)
else:
return {
"statusCode": 400,
"body": json.dumps({"error": f"unknown action: {action}"}),
}
return {"statusCode": 200, "body": json.dumps(result)}
except ValueError as e:
# P10 (REQ-174): identity failures are 401, field validation is 400.
if "missing IAM caller identity" in str(e):
return {"statusCode": 401, "body": json.dumps({"error": str(e)})}
return {"statusCode": 400, "body": json.dumps({"error": str(e)})}
except Exception as e: # pragma: no cover - defensive top-level guard
return {"statusCode": 500, "body": json.dumps({"error": str(e)})}