Skip to content

Commit dfe3b9f

Browse files
Merge pull request #25 from CoreyLeath-code/feature/aws-firehose-s3
feat: add Amazon Data Firehose + S3 telemetry mirror
2 parents 9129500 + 93215f3 commit dfe3b9f

15 files changed

Lines changed: 853 additions & 156 deletions

File tree

‎.env.example‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,17 @@ SNOWFLAKE_DATABASE=
2727
SNOWFLAKE_SCHEMA=
2828
SNOWFLAKE_WAREHOUSE=
2929

30+
# ---------------------------------------------------------------------------
31+
# AWS telemetry mirror (optional — disabled for local development)
32+
# ---------------------------------------------------------------------------
33+
# When enabled, accepted inference logs are queued and mirrored to Amazon Data
34+
# Firehose. Use the AWS SDK default credential chain; prefer workload roles/IRSA
35+
# in AWS instead of committing long-lived access keys.
36+
FIREHOSE_ENABLED=false
37+
FIREHOSE_DELIVERY_STREAM=sentinelai-telemetry
38+
FIREHOSE_QUEUE_SIZE=1000
39+
AWS_REGION=us-east-1
40+
3041
# ---------------------------------------------------------------------------
3142
# LLM / Ollama (optional — llm-guard falls back to stub if unreachable)
3243
# ---------------------------------------------------------------------------
Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,75 @@
1+
name: AWS Telemetry Validation
2+
3+
on:
4+
pull_request:
5+
branches: [main]
6+
paths:
7+
- "ingestion-service/**"
8+
- "terraform/**"
9+
- ".github/workflows/aws-telemetry.yml"
10+
push:
11+
branches: [main]
12+
paths:
13+
- "ingestion-service/**"
14+
- "terraform/**"
15+
- ".github/workflows/aws-telemetry.yml"
16+
17+
permissions:
18+
contents: read
19+
20+
jobs:
21+
firehose-producer:
22+
name: Firehose producer unit tests
23+
runs-on: ubuntu-latest
24+
timeout-minutes: 10
25+
26+
defaults:
27+
run:
28+
working-directory: ingestion-service
29+
30+
steps:
31+
- name: Checkout repository
32+
uses: actions/checkout@v4
33+
34+
- name: Set up Go
35+
uses: actions/setup-go@v5
36+
with:
37+
go-version: "1.24.x"
38+
cache-dependency-path: ingestion-service/go.sum
39+
40+
- name: Verify module files are tidy
41+
run: |
42+
go mod tidy
43+
git diff --exit-code -- go.mod go.sum
44+
45+
- name: Run ingestion tests
46+
run: go test ./...
47+
48+
terraform:
49+
name: Terraform fmt and validate
50+
runs-on: ubuntu-latest
51+
timeout-minutes: 10
52+
53+
defaults:
54+
run:
55+
working-directory: terraform
56+
57+
steps:
58+
- name: Checkout repository
59+
uses: actions/checkout@v4
60+
61+
- name: Set up Terraform
62+
uses: hashicorp/setup-terraform@v3
63+
with:
64+
terraform_version: "1.13.3"
65+
66+
- name: Verify Terraform formatting
67+
run: |
68+
terraform fmt -recursive
69+
git diff --exit-code -- .
70+
71+
- name: Initialize providers without backend
72+
run: terraform init -backend=false -input=false
73+
74+
- name: Validate configuration
75+
run: terraform validate

‎README.md‎

Lines changed: 104 additions & 138 deletions
Large diffs are not rendered by default.

‎docker-compose.yml‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,10 @@ services:
2727
WAREHOUSE_MODE: ${WAREHOUSE_MODE:-postgres}
2828
DATABASE_URL: ${DATABASE_URL:-postgres://sentinel:sentinel@postgres:5432/sentinel?sslmode=disable}
2929
PORT: "8080"
30+
FIREHOSE_ENABLED: ${FIREHOSE_ENABLED:-false}
31+
FIREHOSE_DELIVERY_STREAM: ${FIREHOSE_DELIVERY_STREAM:-sentinelai-telemetry}
32+
FIREHOSE_QUEUE_SIZE: ${FIREHOSE_QUEUE_SIZE:-1000}
33+
AWS_REGION: ${AWS_REGION:-us-east-1}
3034
# Disabled by default; CI enables this to verify multiple replicas route through NGINX.
3135
EXPOSE_INSTANCE_ID: ${EXPOSE_INSTANCE_ID:-false}
3236
expose:

‎docs/aws-firehose-s3.md‎

Lines changed: 97 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,97 @@
1+
# Amazon Data Firehose + S3 telemetry path
2+
3+
SentinelAI can optionally mirror accepted inference telemetry from the Go ingestion service to Amazon Data Firehose. Firehose buffers the records and delivers GZIP-compressed newline-delimited JSON (NDJSON) objects to a private S3 telemetry bucket.
4+
5+
This AWS path is **optional**. Local Docker Compose keeps `FIREHOSE_ENABLED=false`, so PostgreSQL remains the default local persistence path and no AWS account is required for development.
6+
7+
## Architecture
8+
9+
```mermaid
10+
flowchart LR
11+
Client[Model / application] --> Gateway[NGINX ingestion gateway]
12+
Gateway --> Go[Go ingestion replicas]
13+
Go --> DB[(PostgreSQL / Snowflake path)]
14+
Go -. bounded fail-open mirror .-> Queue[In-memory Firehose queue]
15+
Queue --> Firehose[Amazon Data Firehose]
16+
Firehose --> S3[(Amazon S3 telemetry lake)]
17+
Firehose --> CW[CloudWatch delivery logs]
18+
```
19+
20+
The Firehose mirror intentionally does not participate in `/ready`. A temporary AWS failure increments `ingestion_firehose_records_total{status="error"}` and is logged, while the primary ingestion response continues to reflect the primary warehouse write. If the bounded queue fills, records are dropped from the mirror and counted with `status="dropped"` rather than allowing telemetry backpressure to take down ingestion.
21+
22+
## Provision the AWS resources
23+
24+
Prerequisites:
25+
26+
- Terraform 1.6+
27+
- AWS credentials available through the standard AWS credential chain
28+
- Permission to create S3, Firehose, CloudWatch Logs, IAM policy/role, and ECR resources
29+
30+
```bash
31+
cd terraform
32+
terraform init
33+
terraform fmt -check
34+
terraform validate
35+
terraform plan
36+
terraform apply
37+
```
38+
39+
The default configuration creates:
40+
41+
- a private, versioned S3 bucket with SSE-S3 encryption;
42+
- a 30-day telemetry lifecycle policy;
43+
- an Amazon Data Firehose delivery stream named `sentinelai-telemetry`;
44+
- 60-second / 5-MiB buffering with GZIP compression;
45+
- time-partitioned S3 keys under `inference/year=.../month=.../day=.../hour=.../`;
46+
- a CloudWatch log group for Firehose delivery errors;
47+
- a Firehose service role with only the S3 and CloudWatch permissions it needs;
48+
- a separate `sentinelai-firehose-writer` IAM policy granting `PutRecord` and `PutRecordBatch` to the SentinelAI workload;
49+
- the existing SentinelAI ECR repository, now defined in a valid Terraform `.tf` file.
50+
51+
Use `terraform output` after apply to retrieve the generated S3 bucket name, stream ARN, stream name, and writer-policy ARN.
52+
53+
## Give the ingestion workload permission
54+
55+
Do not put long-lived AWS access keys in the repository. Attach the Terraform output `firehose_writer_policy_arn` to the workload identity used by SentinelAI. On EKS, the intended production pattern is an IAM role associated with the ingestion service account (IRSA / EKS workload identity).
56+
57+
For local development against a real AWS account, the AWS SDK for Go v2 uses its normal credential provider chain. Keep credentials outside the repo.
58+
59+
## Enable the mirror
60+
61+
```bash
62+
export FIREHOSE_ENABLED=true
63+
export FIREHOSE_DELIVERY_STREAM=sentinelai-telemetry
64+
export FIREHOSE_QUEUE_SIZE=1000
65+
export AWS_REGION=us-east-1
66+
```
67+
68+
Then start SentinelAI and send a normal inference log:
69+
70+
```bash
71+
curl -X POST http://localhost:8080/log \
72+
-H "Content-Type: application/json" \
73+
-d '{"model_id":"demo","model_version":"v1","latency_ms":120,"tokens_in":32,"tokens_out":64,"status":"ok"}'
74+
```
75+
76+
The producer appends a newline to every JSON record before `PutRecord`. That keeps individual events parseable after Firehose concatenates buffered records into S3 objects.
77+
78+
## Observability
79+
80+
The ingestion service exports:
81+
82+
```text
83+
ingestion_firehose_records_total{status="queued"}
84+
ingestion_firehose_records_total{status="delivered"}
85+
ingestion_firehose_records_total{status="error"}
86+
ingestion_firehose_records_total{status="dropped"}
87+
```
88+
89+
CloudWatch delivery logs cover the managed Firehose-to-S3 leg. Application metrics cover the producer-side queue and `PutRecord` result.
90+
91+
## Failure semantics
92+
93+
This feature is a telemetry mirror, not a transactional dual-write guarantee. The in-memory queue is intentionally bounded and is not persisted across process termination. AWS SDK retries may also produce duplicate Firehose records in some failure scenarios. Consumers should therefore treat the S3 telemetry dataset as at-least-once/best-effort observability data and use stable event identifiers if strict de-duplication is later required.
94+
95+
## Cost control
96+
97+
Firehose, S3, CloudWatch Logs, and related data transfer can incur AWS charges. The default 30-day S3 expiration and 14-day CloudWatch log retention are intended to keep a portfolio/dev deployment bounded. Review the Terraform plan and AWS pricing before leaving the stack running.

‎ingestion-service/Dockerfile‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
# ingestion-service — multi-stage Go build
2-
FROM golang:1.21-alpine AS builder
2+
FROM golang:1.24-alpine AS builder
33
WORKDIR /app
44
COPY go.mod go.sum* ./
55
RUN go mod download

‎ingestion-service/firehose.go‎

Lines changed: 144 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,144 @@
1+
package main
2+
3+
import (
4+
"context"
5+
"encoding/json"
6+
"errors"
7+
"fmt"
8+
"log"
9+
"os"
10+
"strconv"
11+
"strings"
12+
"time"
13+
14+
"github.com/aws/aws-sdk-go-v2/config"
15+
"github.com/aws/aws-sdk-go-v2/service/firehose"
16+
"github.com/aws/aws-sdk-go-v2/service/firehose/types"
17+
"github.com/prometheus/client_golang/prometheus"
18+
)
19+
20+
const (
21+
defaultFirehoseQueueSize = 1000
22+
firehosePublishTimeout = 5 * time.Second
23+
)
24+
25+
var firehoseRecordsTotal = prometheus.NewCounterVec(
26+
prometheus.CounterOpts{
27+
Name: "ingestion_firehose_records_total",
28+
Help: "Inference telemetry records handled by the optional Amazon Data Firehose mirror.",
29+
},
30+
[]string{"status"},
31+
)
32+
33+
func init() {
34+
prometheus.MustRegister(firehoseRecordsTotal)
35+
}
36+
37+
type firehoseAPI interface {
38+
PutRecord(context.Context, *firehose.PutRecordInput, ...func(*firehose.Options)) (*firehose.PutRecordOutput, error)
39+
}
40+
41+
type firehosePublisher struct {
42+
client firehoseAPI
43+
streamName string
44+
}
45+
46+
type firehoseDispatcher struct {
47+
publisher *firehosePublisher
48+
queue chan InferenceLog
49+
}
50+
51+
func newFirehoseDispatcher() (*firehoseDispatcher, error) {
52+
if !strings.EqualFold(strings.TrimSpace(os.Getenv("FIREHOSE_ENABLED")), "true") {
53+
return nil, nil
54+
}
55+
56+
streamName := strings.TrimSpace(os.Getenv("FIREHOSE_DELIVERY_STREAM"))
57+
if streamName == "" {
58+
return nil, errors.New("FIREHOSE_DELIVERY_STREAM is required when FIREHOSE_ENABLED=true")
59+
}
60+
61+
region := strings.TrimSpace(os.Getenv("AWS_REGION"))
62+
if region == "" {
63+
region = "us-east-1"
64+
}
65+
66+
queueSize := defaultFirehoseQueueSize
67+
if raw := strings.TrimSpace(os.Getenv("FIREHOSE_QUEUE_SIZE")); raw != "" {
68+
parsed, err := strconv.Atoi(raw)
69+
if err != nil || parsed < 1 {
70+
return nil, fmt.Errorf("FIREHOSE_QUEUE_SIZE must be a positive integer: %q", raw)
71+
}
72+
queueSize = parsed
73+
}
74+
75+
cfg, err := config.LoadDefaultConfig(context.Background(), config.WithRegion(region))
76+
if err != nil {
77+
return nil, fmt.Errorf("load AWS configuration: %w", err)
78+
}
79+
80+
dispatcher := &firehoseDispatcher{
81+
publisher: &firehosePublisher{
82+
client: firehose.NewFromConfig(cfg),
83+
streamName: streamName,
84+
},
85+
queue: make(chan InferenceLog, queueSize),
86+
}
87+
go dispatcher.run()
88+
89+
return dispatcher, nil
90+
}
91+
92+
func (d *firehoseDispatcher) enqueue(entry InferenceLog) {
93+
select {
94+
case d.queue <- entry:
95+
firehoseRecordsTotal.WithLabelValues("queued").Inc()
96+
default:
97+
firehoseRecordsTotal.WithLabelValues("dropped").Inc()
98+
log.Printf("firehose mirror queue full; dropping telemetry record for model=%s", entry.ModelID)
99+
}
100+
}
101+
102+
func (d *firehoseDispatcher) run() {
103+
for entry := range d.queue {
104+
ctx, cancel := context.WithTimeout(context.Background(), firehosePublishTimeout)
105+
err := d.publisher.publish(ctx, entry)
106+
cancel()
107+
108+
if err != nil {
109+
firehoseRecordsTotal.WithLabelValues("error").Inc()
110+
log.Printf("firehose PutRecord failed for model=%s: %v", entry.ModelID, err)
111+
continue
112+
}
113+
firehoseRecordsTotal.WithLabelValues("delivered").Inc()
114+
}
115+
}
116+
117+
func (p *firehosePublisher) publish(ctx context.Context, entry InferenceLog) error {
118+
payload, err := encodeFirehoseRecord(entry)
119+
if err != nil {
120+
return err
121+
}
122+
123+
_, err = p.client.PutRecord(ctx, &firehose.PutRecordInput{
124+
DeliveryStreamName: &p.streamName,
125+
Record: &types.Record{
126+
Data: payload,
127+
},
128+
})
129+
if err != nil {
130+
return fmt.Errorf("put record: %w", err)
131+
}
132+
return nil
133+
}
134+
135+
func encodeFirehoseRecord(entry InferenceLog) ([]byte, error) {
136+
payload, err := json.Marshal(entry)
137+
if err != nil {
138+
return nil, fmt.Errorf("marshal inference log: %w", err)
139+
}
140+
141+
// Firehose concatenates records inside delivered objects. NDJSON keeps each
142+
// inference event independently parseable after buffering and compression.
143+
return append(payload, '\n'), nil
144+
}

0 commit comments

Comments
 (0)