Repository navigation
Production hardening, LocalStack E2E, and Step Functions contract fixes
Added
- LocalStack Community deployment (
deploy/localstack/) — Build script (auto-detects host arch arm64/amd64), Python deploy script using boto3 to match production Terraform resource shape, and Makefile targets for full local E2E testing. Creates DynamoDB tables with streams, EventBridge custom bus + rules, SQS alert queue, 6 Lambda functions, Step Functions state machine, IAM dummy roles, and event source mappings. Enables running real interlock Lambdas locally against LocalStack for integration testing without AWS costs. SKIP_SCHEDULERenv var —sla-monitorLambda no-ops EventBridge Scheduler calls (CreateSchedule,DeleteSchedule) whenSKIP_SCHEDULER=true. Allows deployment to LocalStack Community (Scheduler is Pro-only) while preserving production behavior. Guards applied ininternal/lambda/sla_monitor.go,internal/lambda/sla/, andinternal/lambda/watchdog/.- Dead-letter queue subsystem (
internal/dlq/) — Typed error classification (transient vs permanent), SQS routing with slog fallback when SQS is unreachable, ULID-based record IDs, and per-record metrics counter interface. Includes no-op router for testing and dry-run modes. - Stream batch handler (
internal/handler/) — Implements AWSReportBatchItemFailuresfor partial batch processing. UsesSequenceNumber(notEventID) per AWS contract. Enforces accounting invariant: processed + dlq_routed + batch_failures == total. - Lambda context middleware (
internal/aws/lambda/) — Derivescontext.WithTimeoutfrom Lambda's remaining execution time minus a configurable safety buffer (default 500ms). Floors at 50ms to prevent zero/negative timeouts. - OpenTelemetry initialization (
internal/telemetry/) — OTLP gRPC trace and metric exporters with graceful no-op fallback whenOTEL_EXPORTER_OTLP_ENDPOINTis unset. Defines 6 application metrics: records processed, stage duration, rules evaluated, DLQ routed, worker pool active, circuit breaker state. - Structured logging with correlation IDs (
internal/telemetry/) — slog handler wrapper that injectscorrelation_idfrom context into every log record. JSON output with source location. - Circuit breaker for HTTP evaluators (
internal/client/) — Wraps external HTTP calls withsony/gobreaker. Configurable trip thresholds, nil-safe defaults, and state introspection. - Exponential backoff retry (
internal/resilience/) — Context-aware retry with jitter clamping to [0, 1], propertime.NewTimercleanup (no timer leaks), and configurable max retries/delay. - Bounded worker pool (
internal/concurrency/) — Thin wrapper arounderrgroup+semaphore.Weightedfor bounded concurrent processing within Lambda executions. - CI quality gates — Makefile
audittarget usinggolangci-lint(reads.golangci.yml) andgo test -race. GitHub Actions workflow updated to usemake auditas a blocking gate. - DLQ audit tracker (
internal/audit/) — Record lifecycle tracking with RWMutex-protected state map, valid transition enforcement (PENDING→ACKED/REJECTED), duplicate detection, and reconciliation reporting for data loss detection. - Hardening config (
internal/config/) — Centralized env-var-based configuration for timeouts, worker pools, DLQ, and circuit breaker thresholds with validation at startup. - Pipeline stage decorators (
internal/pipeline/) — ComposableWithTimeoutandComposedecorators for pipeline stage wrapping. Context pre-cancellation check avoids unnecessary goroutine allocation. - Serverless health checks (
internal/handler/) — EventBridge__ping__handler with pluggableHealthCheckerinterface returning provider connectivity status. - CPU profiler (
internal/handler/) — Captures pprof CPU profiles on__profile__payloads and uploads to S3 with collision-resistant timestamped keys. - Integration and fault injection tests (
tests/integration/) — Mixed-batch stream processing, DLQ router failures, circuit breaker state transitions, retry exhaustion, and context cancellation under fault injection.
Changed
- HTTP trigger retry —
ExecuteHTTPretries transient failures (5xx, network errors) with exponential backoff. Request body resets between attempts. Permanent errors (4xx) skip retry. - Alert dispatcher circuit breaker — Slack HTTP client wrapped with gobreaker to prevent cascade during Slack outages.
- Stream router correlation IDs — Per-record correlation IDs injected into context for structured log tracing across services.
- Telemetry flush per invocation — OTel providers flush (not shutdown) per Lambda invocation to survive environment reuse across warm starts.
Fixed
- SFN SLA cancel states no longer fail with
States.Runtime—CancelSLASchedulesandCancelSLAOnCompleteTriggerFailurereferenced$.config.sla.deadline,$.config.sla.expectedDuration,$.config.sla.maxDurationand$.sensorArrivalAt, buttypes.SLAConfigmarks every fieldomitempty, so an absolute-only SLA never emittedmaxDuration, a relative-only SLA never emitteddeadline, andsensorArrivalAtwas usually absent. The resultingStates.Runtimeis not retriable and is not caught byCatch: ["States.ALL"], so every SLA-configured execution failed afterCompleteTrigger. The SFN input now uses dedicatedSFNInput/SFNConfig/SFNSLAtypes that always emit those keys, and both cancel states now forwardtimezonesohandleSLACancelrecomputes deadlines in the configured zone. - Trigger run IDs are extracted for every trigger type —
ExtractRunIDonly matchedrunId,jobRunId,glue_job_run_id,executionArn,stepIdanddagRunId, while the step-function, EMR, EMR Serverless, Airflow and Databricks executors emitsfn_execution_arn,emr_step_id,emr_sl_job_run_id,airflow_dag_run_idanddatabricks_run_id. The empty run ID was then omitted from the Lambda result and theCheckJobstate raisedStates.Runtimeafter the external job had already been launched.OrchestratorOutput.RunIDis now always marshaled and an empty metadata map takes the sync-sentinel path. - Absolute SLA deadlines are anchored to the execution date —
CalculateAbsoluteDeadlinerolled an explicit execution date forward by 24h (or 1h for:MMdeadlines) whenever the deadline had already passed.sla-monitorcanceltherefore publishedSLA_METfor runs that finished late, thereconcilebreach branch was unreachable, and the watchdog scheduled breach alerts a day late. Roll-forward now applies only when no execution date is supplied, and usesAddDateso the wall-clock time survives DST transitions. - Orchestrator evaluate/trigger results satisfy the state machine contract —
handleEvaluatereturned a result withoutstatuson storage failures, so theIsReadyChoice dereferenced a missing$.evaluateResult.status; it now always emits a status.handleTriggerreturned a partial result with a nil error on configuration failures, soHasTriggerResultsawIsPresent=trueandCheckJobdereferenced a missingrunId, killing the execution with theTRIGGER#lock stuck inRUNNINGuntil TTL; those paths now return a Lambda error soTrigger's Retry/Catch routes toTriggerRetryExhausted, which releases the lock. - Time-dependent tests —
TestSLAMonitor_Calculate_ReturnsRFC3339(failing since 2026-06-16),TestSLAMonitor_Reconcile_ReturnsDeadlines,TestSLAMonitor_Cancel_RecalculatesWhenTimesNotProvidedand three watchdog proactive-SLA tests now injectDeps.NowFuncinstead of relying on the wall clock. - SLA cancel verdict uses the T+1 date for sensor-triggered daily pipelines —
handleSLACanceland the dry-run SLA projection recomputed the absolute deadline frominput.Date/date(the data date D), but sensor-triggered daily pipelines run T+1 — data for date D arrives on D+1, and the watchdog's proactive SLA scheduling already anchors the deadline to D+1. Now that an explicit date no longer rolls forward when its deadline has passed,cancelpublished a falseSLA_BREACHfor pipelines that finished on D+1 before their deadline. The newinternal/lambda.ResolveSLADatecentralizes the watchdog's T+1 rule (cron pipelines and hourly:MMdeadlines are unaffected) and is now shared by the watchdog,sla.handleSLACancel, andstream.publishDryRunSLAProjection.orchestrator.handleEvaluatestorage failures (GetConfigerror, config not found,GetAllSensorserror) are now also logged at error level instead of surfacing only in the returnedstatus: "error"payload.
Release notes: verdicts for pipelines with a non-UTC sla.timezone are now computed in that zone; absolute SLA deadlines no longer roll forward to the next day when past; in-flight executions started before deploy are unaffected by the state-machine change; the sla-monitor reconcile mode has no production invoker today and is not a safety net for deadlines missed while the watchdog was down.
Dependencies
go.opentelemetry.io/otelv1.43.0 (traces + metrics)go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc(OTLP gRPC trace export)go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc(OTLP gRPC metric export)github.com/sony/gobreaker(circuit breaker)github.com/oklog/ulid/v2(DLQ record IDs)golang.org/x/sync(errgroup + semaphore for worker pool)