Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 12 additions & 1 deletion .dockerignore
Original file line number Diff line number Diff line change
@@ -1,6 +1,17 @@
.git
.github
.venv
.env
.terraform
*.tfstate
*.tfstate.*
.pytest_cache
.ruff_cache
__pycache__
data/raw
data/raw-events
artifacts/models

output
tmp
dist
.lambda-build
32 changes: 29 additions & 3 deletions .env.example
Original file line number Diff line number Diff line change
@@ -1,4 +1,30 @@
DATABASE_URL=postgresql+psycopg://riskqueue:riskqueue@postgres:5432/riskqueue
MODEL_ARTIFACT=artifacts/models/riskqueue.joblib
RISKQUEUE_DEMO_MODE=true
# Local PostgreSQL (copy to .env and replace the placeholder)
POSTGRES_DB=riskqueue
POSTGRES_USER=riskqueue
POSTGRES_PASSWORD=
DATABASE_URL=

# Worker and decision policy
MODEL_VERSION=demo-policy-v1
MODEL_ARTIFACT_DIR=artifacts/models
REVIEW_CAPACITY=750
MANUAL_REVIEW_COST=4
AUTO_CREATE_SCHEMA=false
LOG_LEVEL=INFO

# AWS: use an IAM role in deployed environments; never place access keys here
AWS_REGION=us-east-1
PROCESSING_QUEUE_URL=
OBJECT_STORE_BACKEND=local
LOCAL_OBJECT_STORE_ROOT=data/raw-events
S3_SERVER_SIDE_ENCRYPTION=AES256

# Snowflake analytical sink; remains disabled for normal local development
SNOWFLAKE_ENABLED=false
SNOWFLAKE_ACCOUNT=
SNOWFLAKE_USER=
SNOWFLAKE_PASSWORD=
SNOWFLAKE_WAREHOUSE=
SNOWFLAKE_DATABASE=
SNOWFLAKE_SCHEMA=RISKQUEUE_ANALYTICS
SNOWFLAKE_ROLE=
12 changes: 11 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
@@ -1,13 +1,23 @@
.venv/
__pycache__/
.pytest_cache/
.coverage
htmlcov/
.ruff_cache/
.env
.aws-sam/
.terraform/
*.tfstate
*.tfstate.*
crash.log
lambda-package.zip
dist/
.lambda-build/
data/raw-events/
data/raw/*.csv
data/interim/*
data/processed/*
!data/interim/.gitkeep
!data/processed/.gitkeep
artifacts/models/*
!artifacts/models/.gitkeep

12 changes: 12 additions & 0 deletions Dockerfile.worker
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
FROM python:3.12-slim

WORKDIR /app
COPY pyproject.toml README.md ./
COPY riskqueue ./riskqueue
RUN pip install --no-cache-dir .
COPY sql ./sql

RUN useradd --create-home --uid 10001 riskqueue
USER riskqueue

CMD ["python", "-m", "riskqueue.worker"]
412 changes: 315 additions & 97 deletions README.md

Large diffs are not rendered by default.

40 changes: 34 additions & 6 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -2,11 +2,11 @@ services:
postgres:
image: postgres:16-alpine
environment:
POSTGRES_DB: riskqueue
POSTGRES_USER: riskqueue
POSTGRES_PASSWORD: riskqueue
POSTGRES_DB: ${POSTGRES_DB:-riskqueue}
POSTGRES_USER: ${POSTGRES_USER:-riskqueue}
POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:?Set POSTGRES_PASSWORD in .env}
healthcheck:
test: ["CMD-SHELL", "pg_isready -U riskqueue"]
test: ["CMD-SHELL", "pg_isready -U ${POSTGRES_USER:-riskqueue}"]
interval: 5s
timeout: 3s
retries: 10
Expand All @@ -17,14 +17,42 @@ services:
build: .
ports: ["8000:8000"]
environment:
DATABASE_URL: postgresql+psycopg://riskqueue:riskqueue@postgres:5432/riskqueue
DATABASE_URL: postgresql+psycopg://${POSTGRES_USER:-riskqueue}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB:-riskqueue}
depends_on:
postgres:
condition: service_healthy
dashboard:
build: .
command: ["streamlit", "run", "dashboard/app.py", "--server.address=0.0.0.0"]
ports: ["8501:8501"]
worker:
profiles: ["cloud"]
build:
context: .
dockerfile: Dockerfile.worker
environment:
DATABASE_URL: postgresql+psycopg://${POSTGRES_USER:-riskqueue}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB:-riskqueue}
AWS_REGION: ${AWS_REGION:-us-east-1}
PROCESSING_QUEUE_URL: ${PROCESSING_QUEUE_URL}
OBJECT_STORE_BACKEND: ${OBJECT_STORE_BACKEND:-s3}
LOCAL_OBJECT_STORE_ROOT: /app/data/raw-events
MODEL_VERSION: ${MODEL_VERSION:-demo-policy-v1}
MODEL_ARTIFACT_DIR: ${MODEL_ARTIFACT_DIR:-artifacts/models}
REVIEW_CAPACITY: ${REVIEW_CAPACITY:-750}
MANUAL_REVIEW_COST: ${MANUAL_REVIEW_COST:-4}
SNOWFLAKE_ENABLED: ${SNOWFLAKE_ENABLED:-false}
SNOWFLAKE_ACCOUNT: ${SNOWFLAKE_ACCOUNT:-}
SNOWFLAKE_USER: ${SNOWFLAKE_USER:-}
SNOWFLAKE_PASSWORD: ${SNOWFLAKE_PASSWORD:-}
SNOWFLAKE_WAREHOUSE: ${SNOWFLAKE_WAREHOUSE:-}
SNOWFLAKE_DATABASE: ${SNOWFLAKE_DATABASE:-}
SNOWFLAKE_SCHEMA: ${SNOWFLAKE_SCHEMA:-RISKQUEUE_ANALYTICS}
SNOWFLAKE_ROLE: ${SNOWFLAKE_ROLE:-}
volumes:
- ./data/raw-events:/app/data/raw-events
- ./artifacts/models:/app/artifacts/models:ro
depends_on:
postgres:
condition: service_healthy
volumes:
riskqueue-postgres:

41 changes: 30 additions & 11 deletions docs/ARCHITECTURE.md
Original file line number Diff line number Diff line change
@@ -1,18 +1,37 @@
# Architecture

RiskQueue separates event transport, operational application state, and analytical history.

```mermaid
flowchart LR
A[PaySim CSV] --> B[Contract validation]
B --> C[History-only features]
C --> D[Chronological split]
D --> E[Logistic + boosted models]
E --> F[Calibration and metrics]
F --> G[Cost + capacity policy]
G --> H[FastAPI]
H --> I[(PostgreSQL audit log)]
I --> J[Streamlit operations view]
C --> K[Offline drift checks]
A[Incoming transaction batch] --> B[(S3 raw landing)]
B --> C[Lambda metadata validation]
C --> D[SQS processing queue]
D --> E[Container worker]
D -. retry exhaustion .-> DLQ[Dead-letter queue]
E --> F[Validation and history-only features]
F --> G[Fraud scoring and expected-loss policy]
G --> H[(PostgreSQL operational state)]
G --> I[(Snowflake analytical history)]
H --> J[FastAPI and Streamlit]
H --> K[Prefect scheduled flows]
I --> K
```

The batch feature builder is intentionally explicit rather than streaming. A real deployment would maintain aggregates in online feature infrastructure; this project precomputes them to keep the methodology reproducible and inspectable.
## Event boundary

Lambda performs metadata validation and message creation only. CPU-intensive validation, feature generation, model inference, decision optimization, and persistence run in the container worker.

SQS delivery is at least once. PostgreSQL `processed_events` and deterministic idempotency keys prevent duplicate state. Snowflake `MERGE` prevents duplicate analytical facts. Messages are deleted only after successful processing; Terraform routes repeated failures to an encrypted dead-letter queue.

## Storage responsibilities

**S3** is the versioned raw landing layer. **PostgreSQL** is the OLTP database for processed events, current transaction scores, audit records, analyst queues, and application state. **Snowflake** is the OLAP system for historical scores and daily monitoring aggregates. FastAPI and Streamlit do not require Snowflake for request handling.

## Batch orchestration

Prefect coordinates historical synchronization, monitoring-dataset preparation, and aggregate refreshes. It does not replace SQS for event-level processing.

## Local development

The API, dashboard, deterministic demo, PostgreSQL, and local object-store adapter run without AWS or Snowflake. Cloud clients are dependency-injected in tests, so CI never requires live cloud credentials.
163 changes: 163 additions & 0 deletions infra/terraform/main.tf
Original file line number Diff line number Diff line change
@@ -0,0 +1,163 @@
data "aws_caller_identity" "current" {}

locals {
name = "${var.project_name}-${var.environment}"
}

resource "aws_s3_bucket" "raw_events" {
bucket_prefix = "${local.name}-raw-events-"
}

resource "aws_s3_bucket_versioning" "raw_events" {
bucket = aws_s3_bucket.raw_events.id
versioning_configuration {
status = "Enabled"
}
}

resource "aws_s3_bucket_server_side_encryption_configuration" "raw_events" {
bucket = aws_s3_bucket.raw_events.id
rule {
apply_server_side_encryption_by_default {
sse_algorithm = "AES256"
}
}
}

resource "aws_s3_bucket_public_access_block" "raw_events" {
bucket = aws_s3_bucket.raw_events.id
block_public_acls = true
block_public_policy = true
ignore_public_acls = true
restrict_public_buckets = true
}

resource "aws_sqs_queue" "dead_letter" {
name = "${local.name}-processing-dlq"
message_retention_seconds = var.message_retention_seconds
sqs_managed_sse_enabled = true
}

resource "aws_sqs_queue" "processing" {
name = "${local.name}-processing"
visibility_timeout_seconds = 180
message_retention_seconds = var.message_retention_seconds
receive_wait_time_seconds = 20
sqs_managed_sse_enabled = true
redrive_policy = jsonencode({
deadLetterTargetArn = aws_sqs_queue.dead_letter.arn
maxReceiveCount = var.max_receive_count
})
}

resource "aws_sqs_queue_redrive_allow_policy" "dead_letter" {
queue_url = aws_sqs_queue.dead_letter.id
redrive_allow_policy = jsonencode({
redrivePermission = "byQueue"
sourceQueueArns = [aws_sqs_queue.processing.arn]
})
}

resource "aws_iam_role" "ingestion_lambda" {
name = "${local.name}-ingestion-lambda"
assume_role_policy = jsonencode({
Version = "2012-10-17"
Statement = [{
Effect = "Allow"
Principal = {
Service = "lambda.amazonaws.com"
}
Action = "sts:AssumeRole"
}]
})
}

resource "aws_iam_role_policy" "ingestion_lambda" {
name = "${local.name}-ingestion-lambda"
role = aws_iam_role.ingestion_lambda.id
policy = jsonencode({
Version = "2012-10-17"
Statement = [
{
Sid = "SendProcessingMessages"
Effect = "Allow"
Action = ["sqs:SendMessage"]
Resource = aws_sqs_queue.processing.arn
},
{
Sid = "WriteFunctionLogs"
Effect = "Allow"
Action = [
"logs:CreateLogGroup",
"logs:CreateLogStream",
"logs:PutLogEvents"
]
Resource = "arn:aws:logs:${var.aws_region}:${data.aws_caller_identity.current.account_id}:log-group:/aws/lambda/${local.name}-ingestion:*"
}
]
})
}

resource "aws_iam_policy" "worker" {
name = "${local.name}-worker"
description = "Least-privilege access for the RiskQueue processing worker."
policy = jsonencode({
Version = "2012-10-17"
Statement = [
{
Sid = "ConsumeProcessingMessages"
Effect = "Allow"
Action = [
"sqs:ReceiveMessage",
"sqs:DeleteMessage",
"sqs:ChangeMessageVisibility",
"sqs:GetQueueAttributes"
]
Resource = aws_sqs_queue.processing.arn
},
{
Sid = "ReadRawEvents"
Effect = "Allow"
Action = ["s3:GetObject", "s3:GetObjectVersion"]
Resource = "${aws_s3_bucket.raw_events.arn}/incoming/*"
}
]
})
}

resource "aws_lambda_function" "ingestion" {
function_name = "${local.name}-ingestion"
filename = var.lambda_package_path
source_code_hash = filebase64sha256(var.lambda_package_path)
role = aws_iam_role.ingestion_lambda.arn
runtime = "python3.12"
handler = "riskqueue.cloud.lambda_handler.lambda_handler"
timeout = 30
memory_size = 256

environment {
variables = {
PROCESSING_QUEUE_URL = aws_sqs_queue.processing.url
}
}
}

resource "aws_lambda_permission" "allow_s3" {
statement_id = "AllowS3Invoke"
action = "lambda:InvokeFunction"
function_name = aws_lambda_function.ingestion.function_name
principal = "s3.amazonaws.com"
source_arn = aws_s3_bucket.raw_events.arn
}

resource "aws_s3_bucket_notification" "raw_events" {
bucket = aws_s3_bucket.raw_events.id

lambda_function {
lambda_function_arn = aws_lambda_function.ingestion.arn
events = ["s3:ObjectCreated:*"]
filter_prefix = "incoming/"
}

depends_on = [aws_lambda_permission.allow_s3]
}
24 changes: 24 additions & 0 deletions infra/terraform/outputs.tf
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
output "raw_events_bucket" {
description = "Immutable raw-event landing bucket."
value = aws_s3_bucket.raw_events.id
}

output "processing_queue_url" {
description = "Worker processing queue URL."
value = aws_sqs_queue.processing.url
}

output "dead_letter_queue_url" {
description = "Failed processing message queue URL."
value = aws_sqs_queue.dead_letter.url
}

output "ingestion_lambda_name" {
description = "S3 event ingestion Lambda function."
value = aws_lambda_function.ingestion.function_name
}

output "worker_iam_policy_arn" {
description = "IAM policy to attach to the container worker task role."
value = aws_iam_policy.worker.arn
}
Loading
Loading