riskqueue

RiskQueue

Event-driven fraud decisioning under constrained analyst capacity

RiskQueue is a cloud-ready transaction risk and fraud-operations system. It preserves raw events, validates and enriches transactions, scores fraud risk, ranks cases by expected financial loss, records operational state, and publishes analytical history without making synchronous application requests depend on a warehouse.

The application combines AWS S3, Lambda, SQS, a containerized Python worker, PostgreSQL, Snowflake, Prefect, FastAPI, Streamlit, Terraform, scikit-learn, and XGBoost-ready modeling.

RiskQueue generated dashboard overview

[!IMPORTANT] The checked-in metrics are reproducible results from a deterministic synthetic demo. They are not PaySim results and do not establish performance for a real financial institution.

Problem

Fraud operations face two simultaneous constraints: transactions have unequal financial exposure, and analysts can investigate only a limited number of alerts. A useful system must do more than classify transactions. It must preserve source evidence, process events reliably, prevent duplicate financial records, expose current analyst state quickly, and retain analytical history for monitoring and policy evaluation.

RiskQueue separates these concerns:

Verified demo results

Evaluation Result
Held-out synthetic transactions 1,749
Fraud rate 3.49%
Expected-loss value capture at 750 reviews 97.2%
Probability-only value capture at 750 reviews 92.3%
Expected-loss improvement 4.9 percentage points
Amount drift PSI in shifted demo 0.19 — watch
Local in-process batch API p95 at 1,000 10.90 ms

Expected-loss ranking also captures 69.8% of fraud value in 100 reviews and 94.7% in 500 reviews. These values come from the checked-in generated artifacts and are not cloud throughput claims.

Architecture

flowchart LR
    A[Incoming transaction batch] --> B[(Amazon S3 raw events)]
    B --> C[AWS Lambda ingestion]
    C --> D[Amazon SQS processing queue]
    D --> E[Containerized Python worker]
    D -. repeated failure .-> DLQ[Dead-letter queue]
    E --> F[Contract validation]
    F --> G[History-only features]
    G --> H[Fraud scoring]
    H --> I[Expected-loss policy]
    I --> J[(PostgreSQL OLTP)]
    H --> K[(Snowflake OLAP)]
    J --> L[Analyst review queue]
    L --> M[FastAPI and Streamlit]
    K --> N[Prefect batch flows]
    J --> N

Event flow

1. Raw landing in S3

Transaction batches land under the S3 incoming/ prefix. Bucket versioning, server-side encryption, and public-access blocking are provisioned by Terraform. The storage adapter attaches ingestion timestamp and source metadata and uses conditional writes for immutable ingestion.

LocalObjectStore provides the same byte-oriented interface for local development and unit tests. Local writes use exclusive file creation and persist a metadata sidecar.

2. Lightweight Lambda ingestion

S3 object creation invokes a small Lambda handler. The function:

  1. validates the S3 event or supported request envelope;
  2. URL-decodes the object reference;
  3. creates a versioned processing message;
  4. derives a deterministic event identifier when one is not supplied;
  5. sends the message to SQS.

The Lambda function does not load models, create features, or perform database work.

3. SQS delivery and retry isolation

Messages contain:

The worker validates every message with Pydantic. A message is deleted only after successful processing. Failed messages remain visible for retry and move to the dead-letter queue after the Terraform-configured receive count.

4. Containerized processing worker

The worker:

  1. reads and validates the SQS message;
  2. checks PostgreSQL for an existing event ID;
  3. retrieves the raw object;
  4. validates the inference transaction contract;
  5. builds the existing chronological behavioral features;
  6. uses a configured model artifact or deterministic development scorer;
  7. calculates expected loss and queue priority;
  8. writes operational transaction and review state;
  9. incrementally merges analytical scoring history into Snowflake when enabled;
  10. commits the processed-event marker and deletes the SQS message.

The worker uses the existing feature, scoring, and review-queue modules rather than maintaining separate cloud-only model logic.

Reliability

Idempotency

SQS provides at-least-once delivery, so duplicate messages are expected. RiskQueue protects processing at two levels:

A repeated event returns a duplicate result without producing another operational transaction, analytical fact, or analyst queue row. Database changes are committed only after all configured processing stages succeed.

Retries and dead letters

Data architecture

PostgreSQL: operational OLTP state

PostgreSQL remains the application database. It stores:

FastAPI, Streamlit, and analyst workflows should use PostgreSQL for low-latency operational reads and writes.

Snowflake: analytical OLAP history

Snowflake stores append-oriented historical and aggregate data:

The worker uses an idempotent MERGE. Snowflake can be disabled for local development, and FastAPI request latency does not depend on it. The schema is defined in sql/snowflake_schema.sql.

ML decision system

RiskQueue preserves the existing supervised binary-classification workflow:

Leakage-prone post-transaction balance fields and isFlaggedFraud are excluded from model features. Equal-hour transactions see state frozen at the beginning of the hour.

The decision layer makes the operating objective explicit:

expected_loss = fraud_probability × transaction_amount × loss_fraction
expected_review_value = expected_loss − manual_review_cost

Analyst queues support probability, expected-loss, and expected-review-value ranking under a configurable capacity.

Prefect batch workflows

Prefect is used for scheduled analytics, not individual SQS events.

Flow Purpose
historical_warehouse_sync Incrementally extracts recent PostgreSQL operational scores and merges them into Snowflake
monitoring_dataset_flow Creates a parameterized PSI monitoring artifact from reference and current datasets
daily_analytics_flow Refreshes daily Snowflake model-monitoring aggregates

Run flows locally from Python:

python -c "from riskqueue.orchestration.flows import historical_warehouse_sync; historical_warehouse_sync(lookback_hours=24)"

Prefect receives credentials through the process environment or a deployment secret mechanism. No credentials are stored in flow source.

API and dashboard

FastAPI provides:

The six-view Streamlit dashboard covers executive results, model performance, analyst queues, decision policy, explainability, and monitoring.

Testing

The suite contains 73 tests and requires no live AWS or Snowflake account. Cloud behavior is verified through local adapters and injected clients.

Coverage includes:

Run the quality gates:

uv sync --extra dev
uv run ruff check .
uv run pytest --cov=riskqueue --cov-report=term-missing

Local setup

Copy the safe configuration template and set a local database secret:

cp .env.example .env

Generate the existing deterministic demo:

uv sync --extra dev
uv run python -m scripts.run_demo
uv run streamlit run dashboard/app.py

Run the API:

uv run uvicorn riskqueue.api.main:app --reload

Run PostgreSQL, the API, and the dashboard without an AWS or Snowflake account:

docker compose up --build postgres api dashboard

The local object store defaults to data/raw-events/. Snowflake is disabled unless SNOWFLAKE_ENABLED=true.

The SQS worker is an optional Compose profile because it requires an actual queue endpoint:

docker compose --profile cloud up --build worker

Cloud setup

Prerequisites

Build the lightweight Lambda package:

bash scripts/build_lambda_package.sh

Provision AWS resources:

cd infra/terraform
terraform init
terraform fmt -check
terraform validate
terraform plan
terraform apply

Terraform does not create, destroy, or mutate resources automatically from application startup. Review every plan before applying it.

Initialize Snowflake separately with an appropriately privileged deployment identity:

snowsql -f sql/snowflake_schema.sql

Attach Terraform’s worker_iam_policy_arn output to the IAM role used by the worker container. Deploy the worker image from Dockerfile.worker in the container platform of your choice.

Terraform resources

infra/terraform/ provisions:

Terraform state, crash logs, plan artifacts, credentials, and Lambda build output are ignored by Git.

Configuration and security

Never commit .env, AWS keys, Snowflake passwords, database passwords, Terraform state, or generated deployment packages.

Important environment variables are documented in .env.example. In deployed environments:

Repository structure

Path Purpose
riskqueue/cloud/ S3/local storage, SQS transport, message contracts, Lambda ingestion
riskqueue/worker.py Idempotent asynchronous processing worker
riskqueue/warehouse/ Optional Snowflake analytical sink
riskqueue/orchestration/ Prefect batch and monitoring flows
riskqueue/data/ Training and inference validation
riskqueue/features/ Leakage-safe static and behavioral features
riskqueue/modeling/ Training, calibration support, metrics, model artifacts
riskqueue/decisions/ Expected loss, thresholds, review queues
riskqueue/api/ FastAPI scoring and queue endpoints
riskqueue/db/ SQLAlchemy operational models and sessions
dashboard/ Six-view Streamlit application
infra/terraform/ AWS S3, Lambda, SQS, DLQ, and IAM infrastructure
sql/schema.sql PostgreSQL operational schema
sql/snowflake_schema.sql Snowflake analytical schema and aggregates
tests/ Local deterministic unit and integration tests

Limitations

The checked-in results use synthetic data. Historical features are batch-computed rather than maintained in an online feature store. The development API scorer remains deterministic when a trained model artifact is not configured. The project does not provision PostgreSQL, Snowflake, or the worker compute platform through the included AWS module. Costs, review capacity, and loss fraction are illustrative and must be validated for a real deployment.

Data attribution

The intended external modeling dataset is PaySim by E. A. Lopez-Rojas, A. Elmir, and S. Axelsson (2016). Setup and repository policy are documented in data/README.md; the dataset is not redistributed here.