Skip to content

Folders and files

NameName
Last commit message
Last commit date

Latest commit

 

History

16 Commits
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

🔨 Anvil

Distributed Job Processing & Workflow Orchestration Platform

Submit a job. Don't block the request. Get results when they're ready.

Status Java 21 Spring Boot React 19 222 Tests Passing Chaos Tested


ArchitectureFeaturesTech StackGetting StartedTestingProject Structure


Overview

Anvil is a horizontally scalable, fault-tolerant job processing platform built for high-throughput asynchronous work. Clients submit heavy computations — report generation, AI content creation, CSV imports, bulk email campaigns — and receive an immediate tracking ID. A distributed worker pool processes jobs in the background with automatic retries, priority queuing, real-time progress updates, and full operational visibility through an admin dashboard.

Built for: Backend engineers, full-stack developers, and platform teams who need reliable async job orchestration.


Architecture

System Pipeline Architecture

Anvil architecture


System Topology

graph TB
    subgraph Client["Client Layer"]
        Browser["React SPA<br/>TypeScript + Tailwind"]
    end

    subgraph API["API Layer"]
        REST["REST Controllers<br/>/api/v1/*"]
        WS["WebSocket (STOMP)<br/>/ws"]
        Auth["JWT Auth<br/>Spring Security"]
    end

    subgraph Core["Core Engine"]
        Service["JobService"]
        SM["JobStateMachine"]
        Outbox["Transactional Outbox<br/>PostgreSQL"]
        Relay["Outbox Relay"]
    end

    subgraph Queue["Queue Layer"]
        Redis["Redis 7<br/>Priority Queues"]
        Scheduler["Cron Scheduler"]
        Retry["Retry Scheduler"]
    end

    subgraph Workers["Worker Pool"]
        WR["Worker Runner"]
        Watchdog["Worker Watchdog<br/>Orphan Recovery"]
        Handler["Job Handlers<br/>Pluggable Interface"]
    end

    subgraph Data["Data Layer"]
        PG["PostgreSQL 16<br/>Flyway Migrations"]
        R["Redis 7<br/>Heartbeats + Claims"]
    end

    subgraph Obs["Observability"]
        Prom["Prometheus Metrics"]
        Logs["Structured JSON Logs<br/>Correlation IDs"]
        Health["Kubernetes Probes<br/>/actuator/health"]
    end

    Browser -->|REST| REST
    Browser -->|WebSocket| WS
    REST --> Auth
    Auth --> Service
    Service --> SM
    SM --> Outbox
    Outbox --> Relay
    Relay --> Redis
    Scheduler -->|cron due| SM
    Retry -->|backoff expires| SM
    Redis -->|claim| WR
    WR --> Handler
    WR -->|heartbeat| R
    Watchdog -->|reclaim orphans| SM
    Service --> PG
    SM --> PG
    WS --> Browser
    REST -.-> Obs
Loading

Data Flow — Job Lifecycle

sequenceDiagram
    participant C as Client
    participant API as REST API
    participant DB as PostgreSQL
    participant Outbox as Outbox Relay
    participant R as Redis
    participant W as Worker
    participant H as Job Handler
    participant WS as WebSocket

    C->>API: POST /api/v1/jobs
    API->>DB: BEGIN (INSERT job + outbox entry)
    API-->>C: 201 { id, status: CREATED }
    Note over Outbox: Relay polls every 1s
    Outbox->>R: Enqueue job ID
    Outbox->>DB: DELETE outbox entry

    Note over W: Worker polls every 2s
    W->>R: BLPOP + SADD claim
    R-->>W: Job ID
    W->>DB: status → RUNNING
    W->>WS: Push progress (0%)

    W->>H: execute(payload)
    loop Progress Updates
        H-->>WS: Push progress (25%, 50%, 75%)
    end
    H-->>W: result
    W->>DB: status → COMPLETED, result saved
    W->>WS: Push status: COMPLETED
    W->>R: Remove from claimed set
Loading

Job State Machine

stateDiagram-v2
    [*] --> CREATED: Job submitted
    CREATED --> QUEUED: Outbox relay enqueues
    QUEUED --> RUNNING: Worker claims job
    RUNNING --> COMPLETED: Handler succeeds
    RUNNING --> FAILED: Handler throws

    FAILED --> RETRYING: Retries left & backoff elapsed
    RETRYING --> QUEUED: Re-enqueued

    FAILED --> FAILED_PERMANENTLY: Max retries exceeded
    FAILED_PERMANENTLY --> [*]: Moved to DLQ

    RUNNING --> CANCELLING: User cancels
    CANCELLING --> CANCELLED: Confirmed

    CREATED --> QUEUED: Cron scheduler fires
    note right of CREATED: Cron jobs stay here\nuntil next_fire_at
Loading

Container Architecture

graph LR
    subgraph Docker["Docker Compose"]
        subgraph FE["Frontend"]
            NGINX["nginx:alpine<br/>:80"]
        end
        subgraph BE["Backend"]
            JAVA["eclipse-temurin:21-jre-alpine<br/>:8080"]
        end
        subgraph DB["Database"]
            PG["postgres:16-alpine<br/>:5432"]
        end
        subgraph Cache["Cache"]
            REDIS["redis:7-alpine<br/>:6379"]
        end
    end

    NGINX -->|proxy /api/*| JAVA
    NGINX -->|WebSocket upgrade| JAVA
    JAVA --> PG
    JAVA --> REDIS
Loading

Features

Core

Feature Description
Asynchronous Processing Submit jobs via REST, get tracking ID immediately, poll or WebSocket for results
Priority Queues HIGH / MEDIUM / LOW priority with aging — high-priority jobs execute first
Real-Time Progress WebSocket (STOMP + SockJS) pushes live progress bars, status changes, and messages
Cron Scheduling Standard 5-field cron expressions with automatic re-firing on completion
One-Shot Scheduling Schedule a job for a specific future datetime
Automatic Retries Configurable max retries with exponential backoff, automatic re-enqueue
Dead Letter Queue Permanent failures isolated with full failure history for debugging
Job Cancellation Cancel queued or running jobs via API
Pluggable Handlers Add new job types by implementing one interface — zero changes to queue/worker/scheduler

Admin Console

Page Description
Overview Dashboard Live stats: jobs by priority, running/completed/failed counts, worker utilization, DLQ size
Worker Management List all worker nodes with status, heartbeat age, current job
Dead Letter Queue Browse failed jobs, inspect failure history, requeue or discard
Audit Log Filterable log of all system actions with actor, target, and metadata

Reliability

Pattern Implementation
Transactional Outbox DB write + queue enqueue are atomic — no ghost jobs, no data loss
Orphan Reclamation Worker watchdog detects crashed workers every 15s, re-enqueues stalled jobs
Crash Recovery On restart, all QUEUED jobs are re-enqueued to Redis (survives Redis restarts)
Graceful Shutdown SIGTERM stops accepting new work, finishes current job, then exits cleanly
Chaos Tested 20 kill/restart cycles, 60 jobs, 0 jobs lost

Observability

Signal Tool
Metrics Micrometer + Prometheus: queue depth, job throughput, execution latency histograms, worker health
Logging Structured JSON (logstash-logback-encoder) with correlation IDs and job IDs in MDC
Health Kubernetes-ready /actuator/health/liveness and /actuator/health/readiness
Tracing Every request gets an X-Correlation-Id header, propagated through all logs

Tech Stack

Layer Technologies
Backend Java 21, Spring Boot 3.5, Spring Security, Spring WebSocket, Maven
Database PostgreSQL 16, Flyway (schema migrations), hypersistence-utils (JSONB)
Cache & Queue Redis 7 (Streams + Sets for priority queue, heartbeats, distributed locks)
Real-Time STOMP over SockJS with @MessageMapping and topic subscriptions
Auth JWT access tokens (15 min) + refresh tokens (7 days, server-side revocable)
Frontend React 19, TypeScript 6, Vite 8, Tailwind CSS 4
Monitoring Micrometer, Prometheus, Logstash Logback Encoder
Testing JUnit 5, Mockito, Testcontainers (Postgres + Redis), Vitest, React Testing Library
Infra Docker, Docker Compose, multi-stage builds

Quick Start

Docker (Recommended)

git clone https://github.com/Abdul-Rafy2005/anvil.git
cd anvil
docker compose up -d
Service URL
Frontend Dashboard http://localhost
REST API http://localhost:8080/api/v1
Swagger Docs http://localhost:8080/swagger-ui/index.html
Prometheus Metrics http://localhost:8080/actuator/prometheus

Local Development

# 1. Start infrastructure
docker compose up -d postgres redis

# 2. Start backend (separate terminal)
cd backend
./mvnw spring-boot:run

# 3. Start frontend (separate terminal)
cd frontend
npm install && npm run dev

Testing

222 tests, 0 failures — unit, integration, and contract tests across the full stack.

Backend — 180 Tests

Category Framework What's Covered
Unit Tests JUnit 5 + Mockito State machine transitions, retry backoff, cron parsing, password encoding, JWT generation
Integration Tests Testcontainers Real Postgres + Redis: queue ops, worker orchestration, outbox relay, WebSocket security, orphan reclamation
Contract Tests MockMvc Every REST endpoint: happy path, validation, 401/403, 404, pagination

Frontend — 42 Tests

Category Framework What's Covered
Component Tests Vitest + RTL Job list, job detail, create form, admin dashboard, worker list, DLQ, audit log
Hook Tests Vitest WebSocket connection, URL scheme, null handling
Route Guard Tests Vitest Admin access control, redirect behavior

Run Tests

# Backend (requires Docker for Testcontainers)
cd backend && mvn test

# Frontend
cd frontend && npm test

Load & Chaos Testing

Test Result
Load Test 100 concurrent connections, 60s → p50: 12ms, p95: 17ms, p99: 34ms, ~3,857 submissions/min
Chaos Test 20 kill/restart cycles, 60 jobs → 0 jobs lost, all reached terminal state

Project Structure

anvil/
├── backend/                          # Java 21 + Spring Boot 3.5
│   └── src/main/java/com/anvil/
│       ├── api/                      # REST controllers, DTOs, exception handlers
│       ├── auth/                     # JWT provider, Spring Security config, rate limiter
│       ├── job/
│       │   ├── domain/               # Job entity, JobStatus enum, JobStateMachine
│       │   ├── handler/              # JobHandler interface + 6 concrete handlers
│       │   ├── service/              # JobService (create/cancel/retry business rules)
│       │   └── repository/           # Spring Data JPA repositories
│       ├── queue/                    # Queue interface + Redis implementation + OutboxRelay
│       ├── worker/                   # WorkerRunner (poll/claim/execute) + WorkerWatchdog
│       ├── scheduler/                # CronScheduler + RetryScheduler
│       ├── notification/             # WebSocket STOMP push + email service
│       ├── admin/                    # AdminStatsService, DeadLetterService
│       ├── audit/                    # AuditLogService (append-only)
│       └── config/                   # Security, WebSocket, Metrics, Graceful Shutdown
│
├── frontend/                         # React 19 + TypeScript + Vite + Tailwind
│   └── src/
│       ├── pages/                    # Login, Register, JobList, JobDetail, CreateJob
│       │   └── admin/                # Overview, Workers, Dlq, AuditLog
│       ├── components/               # Layout, Badges, ProgressBar, Feedback, AdminRoute
│       ├── hooks/                    # useWebSocket (STOMP/SockJS)
│       ├── contexts/                 # AuthContext (JWT + refresh)
│       └── api/                      # API client with auth interceptors
│
├── scripts/                          # Load test & chaos test PowerShell scripts
├── docs/                             # PRD, Implementation Plan, Phase Summaries
└── docker-compose.yml                # 4 services: postgres, redis, backend, frontend

API Reference

Authentication

# Register
POST /api/v1/auth/register
{ "email": "user@example.com", "password": "secret" }

# Login (returns access + refresh tokens)
POST /api/v1/auth/login
{ "email": "user@example.com", "password": "secret" }

# Refresh access token
POST /api/v1/auth/refresh
{ "refreshToken": "..." }

Jobs

# Create job (immediate)
POST /api/v1/jobs
{ "jobType": "REPORT_GENERATION", "payload": "{\"format\":\"PDF\"}", "priority": "HIGH" }

# Create recurring job
POST /api/v1/jobs
{ "jobType": "CSV_IMPORT", "payload": "{}", "cronExpression": "0 */6 * * *" }

# List jobs (paginated, filterable)
GET /api/v1/jobs?page=0&size=20&status=RUNNING&jobType=REPORT_GENERATION

# Get job detail (includes result, progress, retry info)
GET /api/v1/jobs/{id}

# Cancel job
POST /api/v1/jobs/{id}/cancel

Admin

GET /api/v1/admin/stats/overview     # Dashboard metrics
GET /api/v1/admin/workers            # Worker node list
GET /api/v1/admin/dlq                # Dead letter queue
POST /api/v1/admin/dlq/{id}/requeue  # Requeue failed job
DELETE /api/v1/admin/dlq/{id}        # Discard failed job
GET /api/v1/admin/audit              # Audit log (filterable)

Available Job Types

Type Handler Description
CSV_IMPORT CsvImportHandler Process CSV data rows
EMAIL_CAMPAIGN EmailCampaignHandler Send bulk emails
FILE_COMPRESSION FileCompressionHandler Compress files into ZIP
IMAGE_PROCESSING ImageProcessingHandler Convert images to WebP
REPORT_GENERATION ReportGenerationHandler Generate PDF/CSV reports
AI_CONTENT_GENERATION AiContentGenerationHandler Generate text content

Design Decisions

Decision Rationale
Transactional Outbox over dual-write Prevents the classic "written to DB but never enqueued" bug when a process crashes between DB commit and Redis push
Redis over RabbitMQ for v1 Simpler ops, sufficient for v1 throughput; queue interface is abstracted for future swap
STOMP over raw WebSocket Built-in pub/sub, topic-based routing, simpler client code, Spring native support
Spring Cron (5-field) Industry standard, well-supported by cron-utils library, familiar to ops teams
Multi-stage Docker builds Smaller runtime images (21 JRE Alpine vs full JDK), faster container startup
Testcontainers over mocked DB Integration tests run against real Postgres/Redis — no "works in tests, fails in prod" surprises
JobHandler interface over switch/case Adding a job type = adding one class. Zero changes to queue, scheduler, or worker code

Performance Benchmarks

Load Test (100 concurrent, 60s):
  Submission latency:  p50 = 12ms  |  p95 = 17ms  |  p99 = 34ms
  Throughput:          ~3,857 submissions/min (API accept rate)
  Error rate:          0%

Chaos Test (20 iterations):
  Jobs submitted:      60
  Jobs lost:           0
  Recovery method:     Worker watchdog re-enqueue on restart

Environment Configuration

Variable Default Description
SPRING_DATASOURCE_URL jdbc:postgresql://localhost:5432/anvil Database URL
SPRING_REDIS_HOST localhost Redis host
JWT_SECRET (must be set) HMAC secret for JWT signing
WORKER_POLL_INTERVAL_MS 2000 Worker poll frequency
WORKER_HEARTBEAT_TIMEOUT_MS 30000 Worker considered dead after this
JOB_MAX_RETRIES 3 Default max retry count
JOB_RETRY_BASE_DELAY_MS 1000 Base delay for exponential backoff

License

Distributed under the MIT License. See LICENSE for more information.


Built with a focus on clean architecture, distributed systems reliability, and production-ready engineering patterns.

About

A scalable, job-type agnostic background job processing system for asynchronous task execution, scheduling, retries, progress tracking, and result retrieval.

Resources

Stars

Watchers

Forks

Releases

Packages

Contributors

Languages