diff --git a/.env.example b/.env.example index e4662ff..8a30db3 100644 --- a/.env.example +++ b/.env.example @@ -12,6 +12,16 @@ DB_MIGRATION_PASSWORD= # React 개발 서버 또는 배포 Client 주소를 쉼표로 구분합니다. CORS_ALLOWED_ORIGINS=http://localhost:3000,http://localhost:5173 +# Transactional Outbox worker 설정입니다. +# 일반 실행에서는 켜 두며, 운영 장애 조사 중 자동 처리를 멈춰야 할 때만 false로 둡니다. +OUTBOX_ENABLED=true +OUTBOX_POLL_INTERVAL=1s +OUTBOX_BATCH_SIZE=20 +OUTBOX_LEASE_DURATION=30s +OUTBOX_MAX_ATTEMPTS=8 +OUTBOX_INITIAL_BACKOFF=1s +OUTBOX_MAX_BACKOFF=5m + # Workflow Catalog는 fowoco/knowledge release로부터 만든 Server용 read-only projection입니다. # local/test는 저장소의 개발용 DRAFT projection을 사용합니다. # prod에서는 RELEASED projection 파일 위치를 반드시 지정해야 합니다. diff --git a/README.md b/README.md index 70e315a..671139e 100644 --- a/README.md +++ b/README.md @@ -15,10 +15,10 @@ FOWOCO는 단순 번역 서비스가 아닙니다. 체류·계약·서류·신 | 기술 | Java 17, Spring Boot 4.1.0, Gradle | | 구현 API | Health, Auth 5개, Task Workflow 7개, Approval·Audit 8개. 전체 계약은 실행 중인 Swagger에서 확인 | | 계획 API | Wiki API 카탈로그와 관련 Issue에서 설계·추적 | -| 로컬 DB | H2 + Flyway Auth·Company·Worker core·Task core·Approval·Audit schema | +| 로컬 DB | H2 + Flyway Auth·Company·Worker·Task·Approval·Audit·Outbox schema | | 개발·배포 DB | PostgreSQL + Flyway | | 보안 | JWT Access Token, `ADMIN`·`HR`·`VIEWER` 역할, `company_id` 기반 ActorContext | -| 개발 기반 | Swagger UI, 공통 오류, `request_id`, CI 구성 완료 | +| 개발 기반 | Swagger UI, 공통 오류, `request_id`, CI와 Transactional Outbox 구성 완료 | | AI·Workflow | Knowledge Catalog projection, Task·Checklist·승인·감사와 AI Runtime 계약·방어 검증 구현. Remote 연동·AiRun은 후속 Issue | 계획 문서는 현재 동작하는 API가 아닙니다. 구현의 원본은 코드·테스트와 실행 시 생성되는 OpenAPI이고, 장기 아키텍처 결정은 [ADR](docs/adr/README.md), 계획 범위와 예시는 [API 카탈로그](https://github.com/fowoco/server/wiki/09-API-Specification)와 Issue에서 확인합니다. @@ -140,6 +140,33 @@ POST /api/v1/tasks/{taskId}/approval-requests - 상태 변경, 승인 기록, 감사 이벤트는 같은 DB transaction에 기록되므로 중간 하나가 실패하면 함께 되돌아갑니다. - `GET /api/v1/tasks/{taskId}/activities`는 화면용 안전 타임라인이고, `GET /api/v1/audit-events`는 ADMIN용 필터·cursor 조회입니다. 내부 snapshot 원문은 두 API에 노출하지 않습니다. +### 이벤트 유실 방지와 재처리 + +Task 생성·취소처럼 후속 처리가 필요한 변경은 업무 데이터와 `event_publication`을 +하나의 DB transaction에 저장합니다. 따라서 서버가 commit 직후 종료되어도 +이벤트가 사라지지 않습니다. + +```text +업무 transaction +→ 업무 데이터 + event_publication 함께 commit +→ Outbox worker가 lease 획득 +→ handler 실행 + event_consumption 기록 +→ COMPLETED +``` + +- 일시적 실패는 지수 backoff 후 `RETRY_WAIT`에서 다시 처리합니다. +- 처리 중 서버가 종료되면 lease 만료 후 다른 서버가 이어받습니다. +- `(event_id, handler_name)` 완료 기록으로 이미 성공한 handler를 다시 실행하지 않습니다. +- 재시도 한도 초과, 잘못된 payload 같은 영구 실패는 버리지 않고 + `REVIEW_REQUIRED`로 남깁니다. +- Event payload는 기능별 allow-list를 통과한 작은 업무값만 허용합니다. 이름·이메일, + 전화번호, 여권번호, 토큰, 비밀번호, 전체 Prompt는 저장할 수 없습니다. + +현재 Task 모듈은 `TaskCreated`, `TaskCancelled`를 발행합니다. 새 handler를 추가하는 +방법, 장애 확인 SQL과 설정값은 +[Transactional Outbox 운영 가이드](docs/reliability/transactional-outbox.md)를 +확인합니다. + ### AI Runtime 계약 기반 Server는 AI Runtime에 보낼 수 있는 field를 typed DTO로 제한하고, 전송 전과 응답 후에 @@ -204,6 +231,7 @@ export DEMO_SEED_ADMIN_PASSWORD='로컬 또는 배포 Secret의 12자 이상 값 | Flyway | H2는 공통 migration만, PostgreSQL은 공통 및 PostgreSQL 전용 migration을 순서대로 실행합니다. | `db/migration`, `db/migration-postgresql` | | Workflow Catalog | Knowledge release의 Server용 read-only projection을 시작 시 검증합니다. | `workflow/` | | AI Runtime 계약 | 외부 AI 요청·응답을 allow-list와 version으로 다시 검증합니다. | `aiintegration/` | +| Transactional Outbox | 업무 변경과 후속 이벤트를 함께 저장하고 lease·재시도·멱등 기록으로 복구합니다. | `reliability/` | | Security | JWT에서 ActorContext와 역할을 만들고 VIEWER의 쓰기 요청을 기본 차단합니다. | `SecurityConfig` | | Swagger | Controller의 API 설명을 브라우저 문서로 보여줍니다. | `OpenApiConfig` | | 공통 오류 | 모든 실패를 같은 JSON 구조로 반환합니다. | `common/error` | @@ -331,7 +359,8 @@ server/ │ │ ├── V3__create_worker_document.sql # Worker·Document metadata │ │ ├── V4__create_task_workflow_core.sql # Task·Checklist·전이 이력 │ │ ├── V5__create_approval_audit.sql # 승인·제출·증빙·감사 - │ │ └── V6__add_user_display_name.sql # 회원가입 담당자 표시 이름 + │ │ ├── V6__add_user_display_name.sql # 회원가입 담당자 표시 이름 + │ │ └── V7__create_event_outbox.sql # 내구성 이벤트·handler 완료 기록 │ └── migration-postgresql/ # RLS 등 PostgreSQL 전용 migration └── test/ └── java/com/fowoco/server/ @@ -340,7 +369,8 @@ server/ ├── worker/ ├── task/ ├── aiintegration/ - └── airun/ + ├── airun/ + └── reliability/ ``` 기능 코드가 생기면 해당 기능 안에서 다음 방향으로 확장합니다. diff --git a/docs/adr/0003-task-airun-event-and-retry-model.md b/docs/adr/0003-task-airun-event-and-retry-model.md index e39b5cb..c519308 100644 --- a/docs/adr/0003-task-airun-event-and-retry-model.md +++ b/docs/adr/0003-task-airun-event-and-retry-model.md @@ -200,7 +200,11 @@ allow-list payload - 영구 실패를 성공으로 바꾸지 않고 운영 검토 대상으로 남깁니다. - Event에는 JWT, Worker Link 원본 token, 민감 식별정보, 전체 Prompt를 넣지 않습니다. -MVP는 `DomainEventPublisher` Port 뒤의 **PostgreSQL-backed durable publication**을 사용합니다. #25의 작은 spike에서 Spring Modulith Event Publication Registry와 Spring Boot 4.1 호환성을 확인하고, 적합하면 해당 adapter를 사용하며 부적합하면 Transactional Outbox adapter를 구현합니다. Adapter 선택 때문에 Domain 코드를 바꾸지 않으며 Kafka·RabbitMQ를 필수로 도입하지 않습니다. +MVP는 `DomainEventPublisher` Port 뒤의 **PostgreSQL-backed durable publication**을 +사용합니다. #25 구현에서는 현재 모듈 경계와 Spring Boot 4.1 호환성을 유지하면서 +lease, 실패 분류, 지수 backoff, handler별 멱등 완료 기록을 명시적으로 통제할 수 있는 +Transactional Outbox adapter를 선택했습니다. Adapter 선택 때문에 Domain 코드를 +바꾸지 않으며 Kafka·RabbitMQ를 필수로 도입하지 않습니다. 대표 event 이름은 과거형의 versioned business fact로 작성합니다. diff --git a/docs/database/postgresql-rls-rollout.md b/docs/database/postgresql-rls-rollout.md index 91051f4..76dcc24 100644 --- a/docs/database/postgresql-rls-rollout.md +++ b/docs/database/postgresql-rls-rollout.md @@ -19,7 +19,7 @@ RLS는 기존 `ActorContext`, Repository의 `company_id` 조건, tenant-aware DB transaction-local tenant context와 connection pool 비누수 테스트만 준비합니다. 아직 policy를 만들거나 RLS를 활성화하지 않습니다. -현재 `main`의 V1~V6에는 아래 12개 tenant table이 존재합니다. 기반 단계의 제한 +현재 `main`의 V1~V7에는 아래 14개 tenant table이 존재합니다. 기반 단계의 제한 role 테스트는 이 전체 범위에 업무 DML만 허용하고, table owner·DDL·`TRUNCATE`· `REFERENCES` 권한과 RLS 우회 권한이 없음을 확인합니다. @@ -27,6 +27,15 @@ role 테스트는 이 전체 범위에 업무 DML만 허용하고, table owner· - `worker`, `worker_document` - `task`, `task_checklist_item`, `task_transition_history` - `approval_request`, `external_submission`, `task_evidence`, `audit_event` +- `event_publication`, `event_consumption` + +`event_publication`은 여러 tenant의 미완료 row를 찾는 background queue이므로 일반 +요청 table과 같은 policy를 바로 활성화하면 worker가 아무 이벤트도 claim하지 못할 수 +있습니다. RLS 활성화 전 #34에서 “처리 가능한 `event_id + company_id`만 반환하는 최소 +claim 함수” 또는 동등한 제한된 queue bootstrap 계약을 먼저 확정합니다. claim 뒤 +handler·완료·실패 transaction은 event에 저장된 `company_id`를 tenant context로 +설정하고 일반 policy를 따릅니다. Runtime role에 전체 Outbox RLS 우회 권한을 주지는 +않습니다. ## 설정 계약 @@ -70,6 +79,8 @@ DDL, `TRUNCATE`, `REFERENCES` 권한을 갖지 않습니다. 실제 값은 배 - commit, rollback, 예외, timeout 뒤 같은 physical connection을 재사용해도 이전 context가 남지 않습니다. - Login·Refresh와 구현된 Worker Link 정상 흐름이 유지됩니다. +- Outbox claim이 다른 tenant payload를 노출하지 않고, claim된 event의 handler + transaction이 해당 `company_id` context에서만 실행됩니다. - 오류 응답과 일반 로그에 SQL, JWT, token, email, 개인정보가 노출되지 않습니다. 로컬 또는 CI PostgreSQL 기반 검증은 다음 환경변수를 사용합니다. diff --git a/docs/reliability/transactional-outbox.md b/docs/reliability/transactional-outbox.md new file mode 100644 index 0000000..87cdb13 --- /dev/null +++ b/docs/reliability/transactional-outbox.md @@ -0,0 +1,138 @@ +# Transactional Outbox 운영 가이드 + +## 한눈에 보기 + +Transactional Outbox는 “DB 저장은 성공했는데 후속 이벤트가 사라지는 문제”를 +막습니다. 업무 데이터와 발행할 이벤트를 같은 transaction에 저장한 뒤, 별도 worker가 +안전하게 처리합니다. + +```text +Task 생성·취소 +→ Task·Audit·event_publication을 한 번에 commit +→ Outbox worker가 처리할 row를 lease +→ 각 handler 실행 +→ event_consumption으로 handler 완료를 기록 +→ event_publication을 COMPLETED로 종료 +``` + +Kafka나 RabbitMQ가 필요한 구조는 아닙니다. MVP는 기존 PostgreSQL을 내구성 저장소로 +사용하고, 호출 코드는 `DomainEventPublisher` Port에만 의존합니다. + +## 테이블을 쉽게 이해하기 + +| 테이블 | 의미 | +| --- | --- | +| `event_publication` | 처리해야 할 이벤트와 현재 상태, 시도 횟수, 다음 시각, lease를 저장 | +| `event_consumption` | 어느 handler가 어느 이벤트를 이미 성공했는지 저장 | + +`event_publication`의 상태는 다음과 같습니다. + +| 상태 | 의미 | +| --- | --- | +| `PENDING` | 아직 처리하지 않음 | +| `PROCESSING` | 특정 서버가 제한 시간 동안 처리 권한을 가짐 | +| `RETRY_WAIT` | 일시 실패 후 다음 처리 시각을 기다림 | +| `COMPLETED` | 모든 handler 처리가 끝남 | +| `REVIEW_REQUIRED` | 자동 처리를 멈추고 개발자·운영자 확인이 필요함 | + +## 새 이벤트를 발행하는 방법 + +1. 기능 모듈에서 과거형 업무 사실 이름을 정합니다. 예: `TaskCreated`. +2. payload field를 상수 allow-list로 선언합니다. +3. `DomainEventEnvelope`를 만들고 application service의 기존 `@Transactional` + method 안에서 `DomainEventPublisher.publish()`를 호출합니다. +4. 업무 저장, Audit, Event 중 하나라도 실패하면 전체가 rollback되는 통합 테스트를 + 작성합니다. + +transaction 밖에서 `publish()`하면 서버가 즉시 거부합니다. 이벤트를 먼저 commit한 뒤 +업무 저장을 따로 수행하면 원자성이 깨지므로 금지합니다. + +## 새 handler를 추가하는 방법 + +`DomainEventHandler`를 구현한 Spring Bean을 기능 모듈의 infrastructure 또는 +application adapter에 둡니다. + +- `handlerName()`은 배포 뒤 의미를 바꾸지 않는 고유 이름을 사용합니다. +- `supports(eventType)`으로 처리할 이벤트를 명시합니다. +- `handle(event)`는 다른 모듈 Entity를 직접 수정하지 않고 해당 모듈의 application + command 또는 Port를 호출합니다. +- handler의 DB 변경과 `event_consumption` 저장은 같은 transaction입니다. +- 외부 API는 상대 시스템에도 `event_id` 또는 별도 업무 unique key를 idempotency + key로 전달해야 합니다. DB 완료 기록만으로 상대 시스템의 중복 실행까지 되돌릴 수는 + 없습니다. +- 일시 오류는 `RetryableEventHandlingException`, 입력·계약 오류처럼 반복해도 + 성공하지 않는 오류는 `NonRetryableEventHandlingException`으로 분류합니다. + +이미 완료된 `(event_id, handler_name)`은 재전달 시 건너뜁니다. handler 이름을 +변경하면 새 handler로 인식되므로 단순 refactoring 때 이름을 바꾸지 않습니다. + +## 개인정보와 로그 규칙 + +Event payload는 `SafeEventPayload.of(allowedFields, values)`를 통과해야 합니다. + +- 필요한 작은 업무값만 넣습니다. 예: `status`, `workflow_id`, `task_type`. +- Worker 이름·이메일·전화번호·여권번호·외국인등록번호·계좌번호를 넣지 않습니다. +- JWT, Worker Link 원본 token, 비밀번호, API Key, 전체 Prompt를 넣지 않습니다. +- 큰 객체, 중첩 JSON, Entity 전체를 넣지 않습니다. +- 실패 로그에는 payload나 예외 원문 대신 `event_id`, `event_type`, 안전한 + `error_code`, 시도 횟수만 기록합니다. + +## 기본 설정 + +| 환경변수 | 기본값 | 설명 | +| --- | --- | --- | +| `OUTBOX_ENABLED` | `true` | scheduler 실행 여부 | +| `OUTBOX_POLL_INTERVAL` | `1s` | 처리할 이벤트를 확인하는 간격 | +| `OUTBOX_BATCH_SIZE` | `20` | 한 번에 lease할 최대 row 수 | +| `OUTBOX_LEASE_DURATION` | `30s` | 한 서버의 처리 권한 유효시간 | +| `OUTBOX_MAX_ATTEMPTS` | `8` | 자동 시도 한도 | +| `OUTBOX_INITIAL_BACKOFF` | `1s` | 첫 재시도 대기시간 | +| `OUTBOX_MAX_BACKOFF` | `5m` | 재시도 대기시간 상한 | + +lease는 정상 handler 최대 처리시간보다 길어야 합니다. 값을 줄이기 전에 느린 handler와 +외부 API timeout을 확인합니다. + +## 장애 확인 + +Actuator가 노출되는 내부 운영 환경에서는 다음 Micrometer 지표를 확인합니다. + +- `fowoco.outbox.publications.backlog`: 미완료 이벤트 수 +- `fowoco.outbox.publications.oldest.delay.seconds`: 가장 오래된 미완료 이벤트 지연 +- `fowoco.outbox.publications.processed{result=completed|retry|review_required}`: + 처리 결과 누적 수 + +DB에서는 payload를 출력하지 않고 상태와 안전한 오류만 확인합니다. + +```sql +SELECT event_id, company_id, event_type, status, attempt_count, + next_attempt_at, lease_owner, lease_expires_at, last_error_code, updated_at +FROM event_publication +WHERE status <> 'COMPLETED' +ORDER BY occurred_at; +``` + +`PROCESSING` lease가 만료되면 다음 poll에서 자동 복구됩니다. `RETRY_WAIT`도 +`next_attempt_at` 이후 자동 처리됩니다. `REVIEW_REQUIRED`는 원인을 수정했다고 해서 +DB를 임의로 `PENDING`으로 바꾸지 않습니다. 별도 관리 command와 감사로그가 구현되기 +전에는 담당 개발자가 원인을 확인하고 forward migration 또는 후속 Issue로 복구합니다. + +`OUTBOX_ENABLED=false`는 자동 처리를 멈출 뿐 새 이벤트 저장을 막지 않습니다. 장애 중 +이벤트가 계속 누적될 수 있으므로 backlog를 함께 관찰하고, 수정 배포 후 다시 활성화해 +순서대로 처리합니다. + +## 테스트 기준 + +- H2 통합 테스트: 업무 변경·이벤트 원자성, 재시도 rollback, 중복 전달, lease 만료 + 복구, 수동 검토 전환 +- PostgreSQL CI 테스트: V7 migration, 상태 CHECK, tenant-aware FK, handler unique + constraint, index +- 기능 통합 테스트: 실제 command가 올바른 event type과 최소 payload를 발행하는지 + 검증 + +로컬 전체 검증: + +```bash +./gradlew clean test +``` + +CI는 PostgreSQL 17 service에도 모든 Flyway migration을 적용하고 계약을 검증합니다. diff --git a/src/main/java/com/fowoco/server/reliability/application/NonRetryableEventHandlingException.java b/src/main/java/com/fowoco/server/reliability/application/NonRetryableEventHandlingException.java new file mode 100644 index 0000000..4014f14 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/NonRetryableEventHandlingException.java @@ -0,0 +1,22 @@ +package com.fowoco.server.reliability.application; + +public class NonRetryableEventHandlingException extends RuntimeException { + + private final String errorCode; + + public NonRetryableEventHandlingException(String errorCode) { + super("Non-retryable event handler failure."); + this.errorCode = requireErrorCode(errorCode); + } + + public String errorCode() { + return errorCode; + } + + private static String requireErrorCode(String errorCode) { + if (errorCode == null || !errorCode.matches("^[A-Z][A-Z0-9_]{1,79}$")) { + throw new IllegalArgumentException("Invalid safe event error code."); + } + return errorCode; + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxBackoffPolicy.java b/src/main/java/com/fowoco/server/reliability/application/OutboxBackoffPolicy.java new file mode 100644 index 0000000..9aba8bd --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/OutboxBackoffPolicy.java @@ -0,0 +1,35 @@ +package com.fowoco.server.reliability.application; + +import com.fowoco.server.reliability.config.OutboxProperties; +import java.time.Duration; +import org.springframework.stereotype.Component; + +@Component +public class OutboxBackoffPolicy { + + private final Duration initialBackoff; + private final Duration maxBackoff; + + public OutboxBackoffPolicy(OutboxProperties properties) { + initialBackoff = properties.getInitialBackoff(); + maxBackoff = properties.getMaxBackoff(); + if (maxBackoff.compareTo(initialBackoff) < 0) { + throw new IllegalArgumentException( + "Outbox maxBackoff cannot be shorter than initialBackoff." + ); + } + } + + public Duration delayForAttempt(int attemptCount) { + if (attemptCount < 1) { + throw new IllegalArgumentException("attemptCount must be positive"); + } + long multiplier = 1L << Math.min(attemptCount - 1, 30); + try { + Duration calculated = initialBackoff.multipliedBy(multiplier); + return calculated.compareTo(maxBackoff) > 0 ? maxBackoff : calculated; + } catch (ArithmeticException exception) { + return maxBackoff; + } + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxClaimService.java b/src/main/java/com/fowoco/server/reliability/application/OutboxClaimService.java new file mode 100644 index 0000000..f79876a --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/OutboxClaimService.java @@ -0,0 +1,52 @@ +package com.fowoco.server.reliability.application; + +import com.fowoco.server.reliability.application.port.EventPublicationRepository; +import com.fowoco.server.reliability.config.OutboxProperties; +import com.fowoco.server.reliability.domain.EventPublication; +import java.time.Clock; +import java.time.Instant; +import java.util.ArrayList; +import java.util.List; +import java.util.UUID; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +@Service +public class OutboxClaimService { + + private final EventPublicationRepository repository; + private final OutboxProperties properties; + private final OutboxMetrics metrics; + private final Clock clock; + + public OutboxClaimService( + EventPublicationRepository repository, + OutboxProperties properties, + OutboxMetrics metrics, + Clock clock + ) { + this.repository = repository; + this.properties = properties; + this.metrics = metrics; + this.clock = clock; + } + + @Transactional + public List claimBatch(String owner) { + Instant now = clock.instant(); + List candidates = + repository.lockClaimable(now, properties.getBatchSize()); + List claimed = new ArrayList<>(candidates.size()); + for (EventPublication publication : candidates) { + publication.claim(owner, now, properties.getLeaseDuration()); + if (publication.attemptCount() > properties.getMaxAttempts()) { + publication.requireReview(owner, "EVENT_ATTEMPTS_EXHAUSTED", now); + metrics.recordReviewRequired(); + } else { + claimed.add(publication.eventId()); + } + repository.save(publication); + } + return List.copyOf(claimed); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxCompletionTransaction.java b/src/main/java/com/fowoco/server/reliability/application/OutboxCompletionTransaction.java new file mode 100644 index 0000000..f8628dd --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/OutboxCompletionTransaction.java @@ -0,0 +1,38 @@ +package com.fowoco.server.reliability.application; + +import com.fowoco.server.reliability.application.port.EventPublicationRepository; +import com.fowoco.server.reliability.domain.EventPublication; +import java.time.Clock; +import java.time.Instant; +import java.util.UUID; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; + +@Service +public class OutboxCompletionTransaction { + + private final EventPublicationRepository repository; + private final OutboxMetrics metrics; + private final Clock clock; + + public OutboxCompletionTransaction( + EventPublicationRepository repository, + OutboxMetrics metrics, + Clock clock + ) { + this.repository = repository; + this.metrics = metrics; + this.clock = clock; + } + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void complete(UUID eventId, String owner) { + Instant now = clock.instant(); + EventPublication publication = repository.findByIdForUpdate(eventId) + .orElseThrow(() -> new IllegalStateException("Event publication not found.")); + publication.complete(owner, now); + repository.save(publication); + metrics.recordCompleted(); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxFailureClassifier.java b/src/main/java/com/fowoco/server/reliability/application/OutboxFailureClassifier.java new file mode 100644 index 0000000..5fee47b --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/OutboxFailureClassifier.java @@ -0,0 +1,35 @@ +package com.fowoco.server.reliability.application; + +import com.fowoco.server.reliability.infrastructure.serialization.EventPayloadDecodingException; +import org.springframework.stereotype.Component; + +@Component +public class OutboxFailureClassifier { + + public FailureClassification classify(Throwable failure) { + Throwable cause = unwrap(failure); + if (cause instanceof NonRetryableEventHandlingException exception) { + return new FailureClassification(exception.errorCode(), false); + } + if (cause instanceof EventPayloadDecodingException) { + return new FailureClassification("EVENT_PAYLOAD_INVALID", false); + } + if (cause instanceof RetryableEventHandlingException exception) { + return new FailureClassification(exception.errorCode(), true); + } + return new FailureClassification("EVENT_HANDLER_FAILED", true); + } + + private Throwable unwrap(Throwable failure) { + Throwable current = failure; + while (current.getCause() != null + && (current instanceof java.util.concurrent.CompletionException + || current instanceof java.lang.reflect.UndeclaredThrowableException)) { + current = current.getCause(); + } + return current; + } + + public record FailureClassification(String errorCode, boolean retryable) { + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxFailureTransaction.java b/src/main/java/com/fowoco/server/reliability/application/OutboxFailureTransaction.java new file mode 100644 index 0000000..558edb4 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/OutboxFailureTransaction.java @@ -0,0 +1,70 @@ +package com.fowoco.server.reliability.application; + +import com.fowoco.server.reliability.application.OutboxFailureClassifier.FailureClassification; +import com.fowoco.server.reliability.application.port.EventPublicationRepository; +import com.fowoco.server.reliability.config.OutboxProperties; +import com.fowoco.server.reliability.domain.EventPublication; +import java.time.Clock; +import java.time.Instant; +import java.util.UUID; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; + +@Service +public class OutboxFailureTransaction { + + private final EventPublicationRepository repository; + private final OutboxFailureClassifier classifier; + private final OutboxBackoffPolicy backoffPolicy; + private final OutboxProperties properties; + private final OutboxMetrics metrics; + private final Clock clock; + + public OutboxFailureTransaction( + EventPublicationRepository repository, + OutboxFailureClassifier classifier, + OutboxBackoffPolicy backoffPolicy, + OutboxProperties properties, + OutboxMetrics metrics, + Clock clock + ) { + this.repository = repository; + this.classifier = classifier; + this.backoffPolicy = backoffPolicy; + this.properties = properties; + this.metrics = metrics; + this.clock = clock; + } + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public FailureOutcome recordFailure( + UUID eventId, + String owner, + Throwable failure + ) { + Instant now = clock.instant(); + EventPublication publication = repository.findByIdForUpdate(eventId) + .orElseThrow(() -> new IllegalStateException("Event publication not found.")); + FailureClassification classification = classifier.classify(failure); + boolean exhausted = publication.attemptCount() >= properties.getMaxAttempts(); + if (!classification.retryable() || exhausted) { + publication.requireReview(owner, classification.errorCode(), now); + repository.save(publication); + metrics.recordReviewRequired(); + return new FailureOutcome(classification.errorCode(), false); + } + publication.retry( + owner, + classification.errorCode(), + now.plus(backoffPolicy.delayForAttempt(publication.attemptCount())), + now + ); + repository.save(publication); + metrics.recordRetry(); + return new FailureOutcome(classification.errorCode(), true); + } + + public record FailureOutcome(String errorCode, boolean retryScheduled) { + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxHandlerRegistry.java b/src/main/java/com/fowoco/server/reliability/application/OutboxHandlerRegistry.java new file mode 100644 index 0000000..df72566 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/OutboxHandlerRegistry.java @@ -0,0 +1,35 @@ +package com.fowoco.server.reliability.application; + +import com.fowoco.server.reliability.application.port.DomainEventHandler; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import org.springframework.stereotype.Component; + +@Component +public class OutboxHandlerRegistry { + + private final List handlers; + + public OutboxHandlerRegistry(List handlers) { + Set names = new HashSet<>(); + handlers.forEach(handler -> { + String name = handler.handlerName(); + if (name == null || name.isBlank() || name.length() > 120) { + throw new IllegalArgumentException( + "Event handler name must be 1 to 120 characters." + ); + } + if (!names.add(name)) { + throw new IllegalStateException("Duplicate event handler name: " + name); + } + }); + this.handlers = List.copyOf(handlers); + } + + public List handlersFor(String eventType) { + return handlers.stream() + .filter(handler -> handler.supports(eventType)) + .toList(); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxHandlerTransaction.java b/src/main/java/com/fowoco/server/reliability/application/OutboxHandlerTransaction.java new file mode 100644 index 0000000..5996ed5 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/OutboxHandlerTransaction.java @@ -0,0 +1,68 @@ +package com.fowoco.server.reliability.application; + +import com.fowoco.server.common.id.UuidGenerator; +import com.fowoco.server.reliability.application.port.DomainEventHandler; +import com.fowoco.server.reliability.application.port.EventConsumptionRepository; +import com.fowoco.server.reliability.application.port.EventPublicationRepository; +import com.fowoco.server.reliability.domain.DomainEventEnvelope; +import com.fowoco.server.reliability.domain.EventConsumption; +import com.fowoco.server.reliability.domain.EventPublication; +import com.fowoco.server.reliability.infrastructure.serialization.EventPayloadCodec; +import java.time.Clock; +import java.time.Instant; +import java.util.UUID; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; + +@Service +public class OutboxHandlerTransaction { + + private final EventPublicationRepository publicationRepository; + private final EventConsumptionRepository consumptionRepository; + private final EventPayloadCodec payloadCodec; + private final UuidGenerator uuidGenerator; + private final Clock clock; + + public OutboxHandlerTransaction( + EventPublicationRepository publicationRepository, + EventConsumptionRepository consumptionRepository, + EventPayloadCodec payloadCodec, + UuidGenerator uuidGenerator, + Clock clock + ) { + this.publicationRepository = publicationRepository; + this.consumptionRepository = consumptionRepository; + this.payloadCodec = payloadCodec; + this.uuidGenerator = uuidGenerator; + this.clock = clock; + } + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public boolean deliver( + UUID eventId, + String owner, + DomainEventHandler handler + ) { + Instant now = clock.instant(); + EventPublication publication = publicationRepository.findByIdForUpdate(eventId) + .orElseThrow(() -> new IllegalStateException("Event publication not found.")); + publication.requireActiveLease(owner, now); + String handlerName = handler.handlerName(); + if (consumptionRepository.existsByEventIdAndHandlerName(eventId, handlerName)) { + return false; + } + DomainEventEnvelope event = publication.toEnvelope( + payloadCodec.decode(publication.payloadJson()) + ); + handler.handle(event); + consumptionRepository.save(new EventConsumption( + uuidGenerator.generate(), + eventId, + publication.companyId(), + handlerName, + now + )); + return true; + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxMetrics.java b/src/main/java/com/fowoco/server/reliability/application/OutboxMetrics.java new file mode 100644 index 0000000..9e1f8a8 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/OutboxMetrics.java @@ -0,0 +1,66 @@ +package com.fowoco.server.reliability.application; + +import com.fowoco.server.reliability.application.port.EventPublicationRepository; +import io.micrometer.core.instrument.Counter; +import io.micrometer.core.instrument.Gauge; +import io.micrometer.core.instrument.MeterRegistry; +import java.time.Clock; +import java.time.Duration; +import org.springframework.stereotype.Component; + +@Component +public class OutboxMetrics { + + private final Counter completed; + private final Counter retried; + private final Counter reviewRequired; + + public OutboxMetrics( + MeterRegistry meterRegistry, + EventPublicationRepository repository, + Clock clock + ) { + completed = counter(meterRegistry, "completed"); + retried = counter(meterRegistry, "retry"); + reviewRequired = counter(meterRegistry, "review_required"); + Gauge.builder( + "fowoco.outbox.publications.backlog", + repository, + EventPublicationRepository::countOutstanding + ) + .description("Outstanding durable event publications") + .register(meterRegistry); + Gauge.builder( + "fowoco.outbox.publications.oldest.delay.seconds", + repository, + candidate -> candidate.findOldestOutstandingOccurredAt() + .map(occurredAt -> Math.max( + 0.0, + Duration.between(occurredAt, clock.instant()).toMillis() + / 1000.0 + )) + .orElse(0.0) + ) + .description("Age in seconds of the oldest outstanding publication") + .register(meterRegistry); + } + + public void recordCompleted() { + completed.increment(); + } + + public void recordRetry() { + retried.increment(); + } + + public void recordReviewRequired() { + reviewRequired.increment(); + } + + private Counter counter(MeterRegistry registry, String result) { + return Counter.builder("fowoco.outbox.publications.processed") + .description("Durable event publication processing outcomes") + .tag("result", result) + .register(registry); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxProcessor.java b/src/main/java/com/fowoco/server/reliability/application/OutboxProcessor.java new file mode 100644 index 0000000..ad59e48 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/OutboxProcessor.java @@ -0,0 +1,84 @@ +package com.fowoco.server.reliability.application; + +import com.fowoco.server.reliability.application.port.DomainEventHandler; +import com.fowoco.server.reliability.config.OutboxWorkerIdentity; +import com.fowoco.server.reliability.domain.EventPublication; +import java.util.List; +import java.util.UUID; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Service; + +@Service +public class OutboxProcessor { + + private static final Logger log = LoggerFactory.getLogger(OutboxProcessor.class); + + private final OutboxWorkerIdentity workerIdentity; + private final OutboxClaimService claimService; + private final OutboxReadService readService; + private final OutboxHandlerRegistry handlerRegistry; + private final OutboxHandlerTransaction handlerTransaction; + private final OutboxCompletionTransaction completionTransaction; + private final OutboxFailureTransaction failureTransaction; + + public OutboxProcessor( + OutboxWorkerIdentity workerIdentity, + OutboxClaimService claimService, + OutboxReadService readService, + OutboxHandlerRegistry handlerRegistry, + OutboxHandlerTransaction handlerTransaction, + OutboxCompletionTransaction completionTransaction, + OutboxFailureTransaction failureTransaction + ) { + this.workerIdentity = workerIdentity; + this.claimService = claimService; + this.readService = readService; + this.handlerRegistry = handlerRegistry; + this.handlerTransaction = handlerTransaction; + this.completionTransaction = completionTransaction; + this.failureTransaction = failureTransaction; + } + + public int processAvailable() { + List eventIds = claimService.claimBatch(workerIdentity.value()); + eventIds.forEach(this::processOne); + return eventIds.size(); + } + + private void processOne(UUID eventId) { + EventPublication publication = readService.requirePublication(eventId); + try { + List handlers = + handlerRegistry.handlersFor(publication.eventType()); + for (DomainEventHandler handler : handlers) { + handlerTransaction.deliver(eventId, workerIdentity.value(), handler); + } + completionTransaction.complete(eventId, workerIdentity.value()); + } catch (RuntimeException failure) { + try { + OutboxFailureTransaction.FailureOutcome outcome = + failureTransaction.recordFailure( + eventId, + workerIdentity.value(), + failure + ); + log.warn( + "Outbox event processing failed: eventId={}, eventType={}, " + + "attempt={}, errorCode={}, retryScheduled={}", + eventId, + publication.eventType(), + publication.attemptCount(), + outcome.errorCode(), + outcome.retryScheduled() + ); + } catch (RuntimeException recordingFailure) { + log.error( + "Outbox failure state could not be recorded: eventId={}, eventType={}", + eventId, + publication.eventType() + ); + } + } + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxReadService.java b/src/main/java/com/fowoco/server/reliability/application/OutboxReadService.java new file mode 100644 index 0000000..82e1590 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/OutboxReadService.java @@ -0,0 +1,23 @@ +package com.fowoco.server.reliability.application; + +import com.fowoco.server.reliability.application.port.EventPublicationRepository; +import com.fowoco.server.reliability.domain.EventPublication; +import java.util.UUID; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +@Service +public class OutboxReadService { + + private final EventPublicationRepository repository; + + public OutboxReadService(EventPublicationRepository repository) { + this.repository = repository; + } + + @Transactional(readOnly = true) + public EventPublication requirePublication(UUID eventId) { + return repository.findById(eventId) + .orElseThrow(() -> new IllegalStateException("Event publication not found.")); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/RetryableEventHandlingException.java b/src/main/java/com/fowoco/server/reliability/application/RetryableEventHandlingException.java new file mode 100644 index 0000000..03f693c --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/RetryableEventHandlingException.java @@ -0,0 +1,22 @@ +package com.fowoco.server.reliability.application; + +public class RetryableEventHandlingException extends RuntimeException { + + private final String errorCode; + + public RetryableEventHandlingException(String errorCode) { + super("Retryable event handler failure."); + this.errorCode = requireErrorCode(errorCode); + } + + public String errorCode() { + return errorCode; + } + + private static String requireErrorCode(String errorCode) { + if (errorCode == null || !errorCode.matches("^[A-Z][A-Z0-9_]{1,79}$")) { + throw new IllegalArgumentException("Invalid safe event error code."); + } + return errorCode; + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/port/DomainEventHandler.java b/src/main/java/com/fowoco/server/reliability/application/port/DomainEventHandler.java new file mode 100644 index 0000000..55dcbe9 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/port/DomainEventHandler.java @@ -0,0 +1,12 @@ +package com.fowoco.server.reliability.application.port; + +import com.fowoco.server.reliability.domain.DomainEventEnvelope; + +public interface DomainEventHandler { + + String handlerName(); + + boolean supports(String eventType); + + void handle(DomainEventEnvelope event); +} diff --git a/src/main/java/com/fowoco/server/reliability/application/port/DomainEventPublisher.java b/src/main/java/com/fowoco/server/reliability/application/port/DomainEventPublisher.java new file mode 100644 index 0000000..7b236bb --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/port/DomainEventPublisher.java @@ -0,0 +1,9 @@ +package com.fowoco.server.reliability.application.port; + +import com.fowoco.server.reliability.domain.DomainEventEnvelope; + +@FunctionalInterface +public interface DomainEventPublisher { + + void publish(DomainEventEnvelope event); +} diff --git a/src/main/java/com/fowoco/server/reliability/application/port/EventConsumptionRepository.java b/src/main/java/com/fowoco/server/reliability/application/port/EventConsumptionRepository.java new file mode 100644 index 0000000..97bbdbd --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/port/EventConsumptionRepository.java @@ -0,0 +1,11 @@ +package com.fowoco.server.reliability.application.port; + +import com.fowoco.server.reliability.domain.EventConsumption; +import java.util.UUID; + +public interface EventConsumptionRepository { + + boolean existsByEventIdAndHandlerName(UUID eventId, String handlerName); + + EventConsumption save(EventConsumption consumption); +} diff --git a/src/main/java/com/fowoco/server/reliability/application/port/EventPublicationRepository.java b/src/main/java/com/fowoco/server/reliability/application/port/EventPublicationRepository.java new file mode 100644 index 0000000..5f40baf --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/port/EventPublicationRepository.java @@ -0,0 +1,24 @@ +package com.fowoco.server.reliability.application.port; + +import com.fowoco.server.reliability.domain.EventPublication; +import java.time.Instant; +import java.util.List; +import java.util.Optional; +import java.util.UUID; + +public interface EventPublicationRepository { + + EventPublication append(EventPublication publication); + + EventPublication save(EventPublication publication); + + List lockClaimable(Instant now, int limit); + + Optional findById(UUID eventId); + + Optional findByIdForUpdate(UUID eventId); + + long countOutstanding(); + + Optional findOldestOutstandingOccurredAt(); +} diff --git a/src/main/java/com/fowoco/server/reliability/config/OutboxConfiguration.java b/src/main/java/com/fowoco/server/reliability/config/OutboxConfiguration.java new file mode 100644 index 0000000..3a5d29b --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/config/OutboxConfiguration.java @@ -0,0 +1,18 @@ +package com.fowoco.server.reliability.config; + +import java.util.UUID; +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.annotation.EnableScheduling; + +@Configuration(proxyBeanMethods = false) +@EnableScheduling +@EnableConfigurationProperties(OutboxProperties.class) +public class OutboxConfiguration { + + @Bean + public OutboxWorkerIdentity outboxWorkerIdentity() { + return new OutboxWorkerIdentity("server-" + UUID.randomUUID()); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/config/OutboxProperties.java b/src/main/java/com/fowoco/server/reliability/config/OutboxProperties.java new file mode 100644 index 0000000..313fd9a --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/config/OutboxProperties.java @@ -0,0 +1,85 @@ +package com.fowoco.server.reliability.config; + +import java.time.Duration; +import org.springframework.boot.context.properties.ConfigurationProperties; + +@ConfigurationProperties(prefix = "app.reliability.outbox") +public class OutboxProperties { + + private boolean enabled = true; + private Duration pollInterval = Duration.ofSeconds(1); + private int batchSize = 20; + private Duration leaseDuration = Duration.ofSeconds(30); + private int maxAttempts = 8; + private Duration initialBackoff = Duration.ofSeconds(1); + private Duration maxBackoff = Duration.ofMinutes(5); + + public boolean isEnabled() { + return enabled; + } + + public void setEnabled(boolean enabled) { + this.enabled = enabled; + } + + public Duration getPollInterval() { + return pollInterval; + } + + public void setPollInterval(Duration pollInterval) { + this.pollInterval = requirePositive(pollInterval, "pollInterval"); + } + + public int getBatchSize() { + return batchSize; + } + + public void setBatchSize(int batchSize) { + if (batchSize < 1 || batchSize > 500) { + throw new IllegalArgumentException("batchSize must be between 1 and 500"); + } + this.batchSize = batchSize; + } + + public Duration getLeaseDuration() { + return leaseDuration; + } + + public void setLeaseDuration(Duration leaseDuration) { + this.leaseDuration = requirePositive(leaseDuration, "leaseDuration"); + } + + public int getMaxAttempts() { + return maxAttempts; + } + + public void setMaxAttempts(int maxAttempts) { + if (maxAttempts < 1 || maxAttempts > 100) { + throw new IllegalArgumentException("maxAttempts must be between 1 and 100"); + } + this.maxAttempts = maxAttempts; + } + + public Duration getInitialBackoff() { + return initialBackoff; + } + + public void setInitialBackoff(Duration initialBackoff) { + this.initialBackoff = requirePositive(initialBackoff, "initialBackoff"); + } + + public Duration getMaxBackoff() { + return maxBackoff; + } + + public void setMaxBackoff(Duration maxBackoff) { + this.maxBackoff = requirePositive(maxBackoff, "maxBackoff"); + } + + private static Duration requirePositive(Duration duration, String field) { + if (duration == null || duration.isZero() || duration.isNegative()) { + throw new IllegalArgumentException(field + " must be positive"); + } + return duration; + } +} diff --git a/src/main/java/com/fowoco/server/reliability/config/OutboxWorkerIdentity.java b/src/main/java/com/fowoco/server/reliability/config/OutboxWorkerIdentity.java new file mode 100644 index 0000000..14d8ef5 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/config/OutboxWorkerIdentity.java @@ -0,0 +1,11 @@ +package com.fowoco.server.reliability.config; + +public record OutboxWorkerIdentity(String value) { + + public OutboxWorkerIdentity { + if (value == null || value.isBlank() || value.length() > 128) { + throw new IllegalArgumentException("Outbox worker identity must be 1 to 128 characters."); + } + value = value.trim(); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/domain/DomainEventEnvelope.java b/src/main/java/com/fowoco/server/reliability/domain/DomainEventEnvelope.java new file mode 100644 index 0000000..9eda8bf --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/domain/DomainEventEnvelope.java @@ -0,0 +1,56 @@ +package com.fowoco.server.reliability.domain; + +import java.time.Instant; +import java.util.Objects; +import java.util.UUID; +import java.util.regex.Pattern; + +public record DomainEventEnvelope( + UUID eventId, + String eventType, + String payloadVersion, + String aggregateType, + UUID aggregateId, + UUID companyId, + EventActorType actorType, + UUID actorId, + String requestId, + String traceId, + Instant occurredAt, + SafeEventPayload payload +) { + + private static final Pattern TYPE_NAME = Pattern.compile("^[A-Z][A-Za-z0-9]{1,99}$"); + private static final Pattern PAYLOAD_VERSION = Pattern.compile("^[1-9][0-9]{0,9}$"); + private static final Pattern TRACE_ID = Pattern.compile("^[0-9a-f]{32}$"); + + public DomainEventEnvelope { + Objects.requireNonNull(eventId, "eventId must not be null"); + requireType(eventType, "eventType"); + if (payloadVersion == null || !PAYLOAD_VERSION.matcher(payloadVersion).matches()) { + throw new IllegalArgumentException("payloadVersion must be a positive integer string"); + } + requireType(aggregateType, "aggregateType"); + Objects.requireNonNull(aggregateId, "aggregateId must not be null"); + Objects.requireNonNull(companyId, "companyId must not be null"); + Objects.requireNonNull(actorType, "actorType must not be null"); + if (actorType != EventActorType.SYSTEM_RULE && actorId == null) { + throw new IllegalArgumentException("actorId is required for non-system events"); + } + if (requestId == null || requestId.isBlank() || requestId.length() > 128) { + throw new IllegalArgumentException("requestId must be 1 to 128 characters"); + } + requestId = requestId.trim(); + if (traceId != null && !TRACE_ID.matcher(traceId).matches()) { + throw new IllegalArgumentException("traceId must be a 32-character lowercase hex value"); + } + Objects.requireNonNull(occurredAt, "occurredAt must not be null"); + Objects.requireNonNull(payload, "payload must not be null"); + } + + private static void requireType(String value, String fieldName) { + if (value == null || !TYPE_NAME.matcher(value).matches()) { + throw new IllegalArgumentException(fieldName + " must use a stable type name"); + } + } +} diff --git a/src/main/java/com/fowoco/server/reliability/domain/EventActorType.java b/src/main/java/com/fowoco/server/reliability/domain/EventActorType.java new file mode 100644 index 0000000..87c4815 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/domain/EventActorType.java @@ -0,0 +1,8 @@ +package com.fowoco.server.reliability.domain; + +public enum EventActorType { + HR_USER, + WORKER_LINK, + AI_AGENT, + SYSTEM_RULE +} diff --git a/src/main/java/com/fowoco/server/reliability/domain/EventConsumption.java b/src/main/java/com/fowoco/server/reliability/domain/EventConsumption.java new file mode 100644 index 0000000..749072c --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/domain/EventConsumption.java @@ -0,0 +1,25 @@ +package com.fowoco.server.reliability.domain; + +import java.time.Instant; +import java.util.Objects; +import java.util.UUID; + +public record EventConsumption( + UUID consumptionId, + UUID eventId, + UUID companyId, + String handlerName, + Instant completedAt +) { + + public EventConsumption { + Objects.requireNonNull(consumptionId, "consumptionId must not be null"); + Objects.requireNonNull(eventId, "eventId must not be null"); + Objects.requireNonNull(companyId, "companyId must not be null"); + if (handlerName == null || handlerName.isBlank() || handlerName.length() > 120) { + throw new IllegalArgumentException("handlerName must be 1 to 120 characters"); + } + handlerName = handlerName.trim(); + Objects.requireNonNull(completedAt, "completedAt must not be null"); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/domain/EventPublication.java b/src/main/java/com/fowoco/server/reliability/domain/EventPublication.java new file mode 100644 index 0000000..599a0d2 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/domain/EventPublication.java @@ -0,0 +1,378 @@ +package com.fowoco.server.reliability.domain; + +import java.time.Duration; +import java.time.Instant; +import java.util.Objects; +import java.util.UUID; +import java.util.regex.Pattern; + +public final class EventPublication { + + private static final Pattern ERROR_CODE = Pattern.compile("^[A-Z][A-Z0-9_]{1,79}$"); + + private final UUID eventId; + private final UUID companyId; + private final String eventType; + private final String payloadVersion; + private final String aggregateType; + private final UUID aggregateId; + private final EventActorType actorType; + private final UUID actorId; + private final String requestId; + private final String traceId; + private final String payloadJson; + private final Instant occurredAt; + private final Instant createdAt; + private EventPublicationStatus status; + private int attemptCount; + private Instant nextAttemptAt; + private String leaseOwner; + private Instant leaseExpiresAt; + private String lastErrorCode; + private Instant completedAt; + private Instant updatedAt; + private long version; + + private EventPublication( + UUID eventId, + UUID companyId, + String eventType, + String payloadVersion, + String aggregateType, + UUID aggregateId, + EventActorType actorType, + UUID actorId, + String requestId, + String traceId, + String payloadJson, + EventPublicationStatus status, + int attemptCount, + Instant nextAttemptAt, + String leaseOwner, + Instant leaseExpiresAt, + String lastErrorCode, + Instant completedAt, + Instant occurredAt, + Instant createdAt, + Instant updatedAt, + long version + ) { + this.eventId = Objects.requireNonNull(eventId); + this.companyId = Objects.requireNonNull(companyId); + this.eventType = requireText(eventType); + this.payloadVersion = requireText(payloadVersion); + this.aggregateType = requireText(aggregateType); + this.aggregateId = Objects.requireNonNull(aggregateId); + this.actorType = Objects.requireNonNull(actorType); + this.actorId = actorId; + this.requestId = requireText(requestId); + this.traceId = traceId; + this.payloadJson = requireText(payloadJson); + this.status = Objects.requireNonNull(status); + this.attemptCount = attemptCount; + this.nextAttemptAt = nextAttemptAt; + this.leaseOwner = leaseOwner; + this.leaseExpiresAt = leaseExpiresAt; + this.lastErrorCode = lastErrorCode; + this.completedAt = completedAt; + this.occurredAt = Objects.requireNonNull(occurredAt); + this.createdAt = Objects.requireNonNull(createdAt); + this.updatedAt = Objects.requireNonNull(updatedAt); + this.version = version; + } + + public static EventPublication pending( + DomainEventEnvelope event, + String payloadJson, + Instant publishedAt + ) { + Objects.requireNonNull(event); + Objects.requireNonNull(publishedAt); + if (publishedAt.isBefore(event.occurredAt())) { + throw new IllegalArgumentException("publishedAt cannot precede occurredAt"); + } + return new EventPublication( + event.eventId(), + event.companyId(), + event.eventType(), + event.payloadVersion(), + event.aggregateType(), + event.aggregateId(), + event.actorType(), + event.actorId(), + event.requestId(), + event.traceId(), + payloadJson, + EventPublicationStatus.PENDING, + 0, + publishedAt, + null, + null, + null, + null, + event.occurredAt(), + publishedAt, + publishedAt, + 0 + ); + } + + public static EventPublication restore( + UUID eventId, + UUID companyId, + String eventType, + String payloadVersion, + String aggregateType, + UUID aggregateId, + EventActorType actorType, + UUID actorId, + String requestId, + String traceId, + String payloadJson, + EventPublicationStatus status, + int attemptCount, + Instant nextAttemptAt, + String leaseOwner, + Instant leaseExpiresAt, + String lastErrorCode, + Instant completedAt, + Instant occurredAt, + Instant createdAt, + Instant updatedAt, + long version + ) { + return new EventPublication( + eventId, + companyId, + eventType, + payloadVersion, + aggregateType, + aggregateId, + actorType, + actorId, + requestId, + traceId, + payloadJson, + status, + attemptCount, + nextAttemptAt, + leaseOwner, + leaseExpiresAt, + lastErrorCode, + completedAt, + occurredAt, + createdAt, + updatedAt, + version + ); + } + + public boolean isClaimableAt(Instant now) { + Objects.requireNonNull(now); + return switch (status) { + case PENDING, RETRY_WAIT -> + nextAttemptAt != null && !nextAttemptAt.isAfter(now); + case PROCESSING -> + leaseExpiresAt != null && !leaseExpiresAt.isAfter(now); + case COMPLETED, REVIEW_REQUIRED -> false; + }; + } + + public void claim(String owner, Instant now, Duration leaseDuration) { + String normalizedOwner = requireOwner(owner); + Objects.requireNonNull(now); + requirePositive(leaseDuration, "leaseDuration"); + if (!isClaimableAt(now)) { + throw new IllegalStateException("Event publication is not claimable."); + } + status = EventPublicationStatus.PROCESSING; + attemptCount++; + nextAttemptAt = null; + leaseOwner = normalizedOwner; + leaseExpiresAt = now.plus(leaseDuration); + lastErrorCode = null; + completedAt = null; + updatedAt = now; + } + + public void complete(String owner, Instant now) { + requireActiveLease(owner, now); + status = EventPublicationStatus.COMPLETED; + clearLease(); + nextAttemptAt = null; + lastErrorCode = null; + completedAt = now; + updatedAt = now; + } + + public void retry(String owner, String errorCode, Instant nextAttempt, Instant now) { + requireActiveLease(owner, now); + if (nextAttempt == null || !nextAttempt.isAfter(now)) { + throw new IllegalArgumentException("nextAttempt must be after now"); + } + status = EventPublicationStatus.RETRY_WAIT; + clearLease(); + nextAttemptAt = nextAttempt; + lastErrorCode = requireErrorCode(errorCode); + completedAt = null; + updatedAt = now; + } + + public void requireReview(String owner, String errorCode, Instant now) { + requireActiveLease(owner, now); + status = EventPublicationStatus.REVIEW_REQUIRED; + clearLease(); + nextAttemptAt = null; + lastErrorCode = requireErrorCode(errorCode); + completedAt = null; + updatedAt = now; + } + + public void requireActiveLease(String owner, Instant now) { + String normalizedOwner = requireOwner(owner); + Objects.requireNonNull(now); + if (status != EventPublicationStatus.PROCESSING + || !normalizedOwner.equals(leaseOwner) + || leaseExpiresAt == null + || !leaseExpiresAt.isAfter(now)) { + throw new IllegalStateException("Event publication lease is not active."); + } + } + + private void clearLease() { + leaseOwner = null; + leaseExpiresAt = null; + } + + private static String requireErrorCode(String errorCode) { + if (errorCode == null || !ERROR_CODE.matcher(errorCode).matches()) { + throw new IllegalArgumentException("Invalid safe event error code."); + } + return errorCode; + } + + private static String requireOwner(String owner) { + if (owner == null || owner.isBlank() || owner.length() > 128) { + throw new IllegalArgumentException("lease owner must be 1 to 128 characters"); + } + return owner.trim(); + } + + private static String requireText(String value) { + if (value == null || value.isBlank()) { + throw new IllegalArgumentException("value must not be blank"); + } + return value.trim(); + } + + private static void requirePositive(Duration duration, String name) { + if (duration == null || duration.isZero() || duration.isNegative()) { + throw new IllegalArgumentException(name + " must be positive"); + } + } + + public DomainEventEnvelope toEnvelope(SafeEventPayload payload) { + return new DomainEventEnvelope( + eventId, + eventType, + payloadVersion, + aggregateType, + aggregateId, + companyId, + actorType, + actorId, + requestId, + traceId, + occurredAt, + payload + ); + } + + public UUID eventId() { + return eventId; + } + + public UUID companyId() { + return companyId; + } + + public String eventType() { + return eventType; + } + + public String payloadVersion() { + return payloadVersion; + } + + public String aggregateType() { + return aggregateType; + } + + public UUID aggregateId() { + return aggregateId; + } + + public EventActorType actorType() { + return actorType; + } + + public UUID actorId() { + return actorId; + } + + public String requestId() { + return requestId; + } + + public String traceId() { + return traceId; + } + + public String payloadJson() { + return payloadJson; + } + + public EventPublicationStatus status() { + return status; + } + + public int attemptCount() { + return attemptCount; + } + + public Instant nextAttemptAt() { + return nextAttemptAt; + } + + public String leaseOwner() { + return leaseOwner; + } + + public Instant leaseExpiresAt() { + return leaseExpiresAt; + } + + public String lastErrorCode() { + return lastErrorCode; + } + + public Instant completedAt() { + return completedAt; + } + + public Instant occurredAt() { + return occurredAt; + } + + public Instant createdAt() { + return createdAt; + } + + public Instant updatedAt() { + return updatedAt; + } + + public long version() { + return version; + } +} diff --git a/src/main/java/com/fowoco/server/reliability/domain/EventPublicationStatus.java b/src/main/java/com/fowoco/server/reliability/domain/EventPublicationStatus.java new file mode 100644 index 0000000..97512de --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/domain/EventPublicationStatus.java @@ -0,0 +1,9 @@ +package com.fowoco.server.reliability.domain; + +public enum EventPublicationStatus { + PENDING, + PROCESSING, + RETRY_WAIT, + COMPLETED, + REVIEW_REQUIRED +} diff --git a/src/main/java/com/fowoco/server/reliability/domain/SafeEventPayload.java b/src/main/java/com/fowoco/server/reliability/domain/SafeEventPayload.java new file mode 100644 index 0000000..0adf6cb --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/domain/SafeEventPayload.java @@ -0,0 +1,146 @@ +package com.fowoco.server.reliability.domain; + +import java.time.temporal.TemporalAccessor; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Objects; +import java.util.Set; +import java.util.UUID; +import java.util.regex.Pattern; + +/** + * Event payload that accepts only explicitly allowed, non-sensitive fields and scalar values. + */ +public final class SafeEventPayload { + + private static final Pattern FIELD_NAME = Pattern.compile("^[a-z][a-z0-9_]{0,49}$"); + private static final int MAX_COLLECTION_SIZE = 50; + private static final int MAX_STRING_LENGTH = 500; + private static final Pattern REGISTRATION_NUMBER = + Pattern.compile("(? SENSITIVE_NAME_PARTS = Set.of( + "account_number", + "alien_registration", + "bank_account", + "display_name", + "email", + "legal_name", + "passport", + "password", + "phone", + "prompt", + "resident_number", + "token", + "jwt" + ); + + private final Map values; + + private SafeEventPayload(Map values) { + this.values = values; + } + + public static SafeEventPayload empty() { + return new SafeEventPayload(Map.of()); + } + + public static SafeEventPayload of( + Set allowedFields, + Map rawValues + ) { + Objects.requireNonNull(allowedFields, "allowedFields must not be null"); + Objects.requireNonNull(rawValues, "rawValues must not be null"); + Set immutableAllowedFields = Set.copyOf(allowedFields); + immutableAllowedFields.forEach(SafeEventPayload::validateFieldName); + + Map safeValues = new LinkedHashMap<>(); + rawValues.forEach((field, value) -> { + validateFieldName(field); + if (!immutableAllowedFields.contains(field)) { + throw new IllegalArgumentException( + "Event payload field is not allow-listed: " + field + ); + } + if (value != null) { + safeValues.put(field, canonicalize(value)); + } + }); + return new SafeEventPayload(Collections.unmodifiableMap(safeValues)); + } + + public Map values() { + return values; + } + + private static void validateFieldName(String field) { + if (field == null || !FIELD_NAME.matcher(field).matches()) { + throw new IllegalArgumentException("Invalid event payload field name."); + } + String normalized = field.toLowerCase(Locale.ROOT); + if (SENSITIVE_NAME_PARTS.stream().anyMatch(normalized::contains)) { + throw new IllegalArgumentException( + "Sensitive event payload field is not allowed: " + field + ); + } + } + + private static Object canonicalize(Object value) { + if (value instanceof String text) { + if (text.length() > MAX_STRING_LENGTH) { + throw new IllegalArgumentException("Event payload string is too long."); + } + if (containsSensitiveValue(text)) { + throw new IllegalArgumentException( + "Sensitive event payload value is not allowed." + ); + } + return text; + } + if (value instanceof Number || value instanceof Boolean) { + return value; + } + if (value instanceof Enum enumValue) { + return enumValue.name(); + } + if (value instanceof UUID || value instanceof TemporalAccessor) { + return value.toString(); + } + if (value instanceof Collection collection) { + if (collection.size() > MAX_COLLECTION_SIZE) { + throw new IllegalArgumentException("Event payload collection is too large."); + } + List safeItems = new ArrayList<>(collection.size()); + collection.forEach(item -> { + if (item == null || item instanceof Map || item instanceof Collection) { + throw new IllegalArgumentException( + "Event payload collections support non-null scalar items only." + ); + } + safeItems.add(canonicalize(item)); + }); + return List.copyOf(safeItems); + } + throw new IllegalArgumentException( + "Unsupported event payload value type: " + value.getClass().getSimpleName() + ); + } + + private static boolean containsSensitiveValue(String value) { + return REGISTRATION_NUMBER.matcher(value).find() + || PHONE_NUMBER.matcher(value).find() + || BEARER_TOKEN.matcher(value).find() + || SECRET_ASSIGNMENT.matcher(value).find(); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/EventConsumptionJpaEntity.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/EventConsumptionJpaEntity.java new file mode 100644 index 0000000..62f940f --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/EventConsumptionJpaEntity.java @@ -0,0 +1,47 @@ +package com.fowoco.server.reliability.infrastructure.persistence; + +import com.fowoco.server.reliability.domain.EventConsumption; +import jakarta.persistence.Column; +import jakarta.persistence.Entity; +import jakarta.persistence.Id; +import jakarta.persistence.Table; +import java.time.Instant; +import java.util.UUID; + +@Entity +@Table(name = "event_consumption") +class EventConsumptionJpaEntity { + + @Id + @Column(name = "consumption_id", nullable = false, updatable = false) + private UUID consumptionId; + @Column(name = "event_id", nullable = false, updatable = false) + private UUID eventId; + @Column(name = "company_id", nullable = false, updatable = false) + private UUID companyId; + @Column(name = "handler_name", nullable = false, length = 120, updatable = false) + private String handlerName; + @Column(name = "completed_at", nullable = false, updatable = false) + private Instant completedAt; + + protected EventConsumptionJpaEntity() { + } + + EventConsumptionJpaEntity(EventConsumption consumption) { + consumptionId = consumption.consumptionId(); + eventId = consumption.eventId(); + companyId = consumption.companyId(); + handlerName = consumption.handlerName(); + completedAt = consumption.completedAt(); + } + + EventConsumption toDomain() { + return new EventConsumption( + consumptionId, + eventId, + companyId, + handlerName, + completedAt + ); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/EventPublicationJpaEntity.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/EventPublicationJpaEntity.java new file mode 100644 index 0000000..18a2413 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/EventPublicationJpaEntity.java @@ -0,0 +1,127 @@ +package com.fowoco.server.reliability.infrastructure.persistence; + +import com.fowoco.server.reliability.domain.EventActorType; +import com.fowoco.server.reliability.domain.EventPublication; +import com.fowoco.server.reliability.domain.EventPublicationStatus; +import jakarta.persistence.Column; +import jakarta.persistence.Entity; +import jakarta.persistence.EnumType; +import jakarta.persistence.Enumerated; +import jakarta.persistence.Id; +import jakarta.persistence.Table; +import jakarta.persistence.Version; +import java.time.Instant; +import java.util.UUID; + +@Entity +@Table(name = "event_publication") +class EventPublicationJpaEntity { + + @Id + @Column(name = "event_id", nullable = false, updatable = false) + private UUID eventId; + @Column(name = "company_id", nullable = false, updatable = false) + private UUID companyId; + @Column(name = "event_type", nullable = false, length = 100, updatable = false) + private String eventType; + @Column(name = "payload_version", nullable = false, length = 20, updatable = false) + private String payloadVersion; + @Column(name = "aggregate_type", nullable = false, length = 60, updatable = false) + private String aggregateType; + @Column(name = "aggregate_id", nullable = false, updatable = false) + private UUID aggregateId; + @Enumerated(EnumType.STRING) + @Column(name = "actor_type", nullable = false, length = 30, updatable = false) + private EventActorType actorType; + @Column(name = "actor_id", updatable = false) + private UUID actorId; + @Column(name = "request_id", nullable = false, length = 128, updatable = false) + private String requestId; + @Column(name = "trace_id", length = 64, updatable = false) + private String traceId; + @Column(name = "payload_json", nullable = false, columnDefinition = "TEXT", updatable = false) + private String payloadJson; + @Enumerated(EnumType.STRING) + @Column(name = "status", nullable = false, length = 30) + private EventPublicationStatus status; + @Column(name = "attempt_count", nullable = false) + private int attemptCount; + @Column(name = "next_attempt_at") + private Instant nextAttemptAt; + @Column(name = "lease_owner", length = 128) + private String leaseOwner; + @Column(name = "lease_expires_at") + private Instant leaseExpiresAt; + @Column(name = "last_error_code", length = 80) + private String lastErrorCode; + @Column(name = "completed_at") + private Instant completedAt; + @Column(name = "occurred_at", nullable = false, updatable = false) + private Instant occurredAt; + @Column(name = "created_at", nullable = false, updatable = false) + private Instant createdAt; + @Column(name = "updated_at", nullable = false) + private Instant updatedAt; + @Version + @Column(name = "version", nullable = false) + private long version; + + protected EventPublicationJpaEntity() { + } + + EventPublicationJpaEntity(EventPublication publication) { + eventId = publication.eventId(); + companyId = publication.companyId(); + eventType = publication.eventType(); + payloadVersion = publication.payloadVersion(); + aggregateType = publication.aggregateType(); + aggregateId = publication.aggregateId(); + actorType = publication.actorType(); + actorId = publication.actorId(); + requestId = publication.requestId(); + traceId = publication.traceId(); + payloadJson = publication.payloadJson(); + occurredAt = publication.occurredAt(); + createdAt = publication.createdAt(); + version = publication.version(); + apply(publication); + } + + void apply(EventPublication publication) { + status = publication.status(); + attemptCount = publication.attemptCount(); + nextAttemptAt = publication.nextAttemptAt(); + leaseOwner = publication.leaseOwner(); + leaseExpiresAt = publication.leaseExpiresAt(); + lastErrorCode = publication.lastErrorCode(); + completedAt = publication.completedAt(); + updatedAt = publication.updatedAt(); + } + + EventPublication toDomain() { + return EventPublication.restore( + eventId, + companyId, + eventType, + payloadVersion, + aggregateType, + aggregateId, + actorType, + actorId, + requestId, + traceId, + payloadJson, + status, + attemptCount, + nextAttemptAt, + leaseOwner, + leaseExpiresAt, + lastErrorCode, + completedAt, + occurredAt, + createdAt, + updatedAt, + version + ); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaEventConsumptionRepository.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaEventConsumptionRepository.java new file mode 100644 index 0000000..c7aeecc --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaEventConsumptionRepository.java @@ -0,0 +1,28 @@ +package com.fowoco.server.reliability.infrastructure.persistence; + +import com.fowoco.server.reliability.application.port.EventConsumptionRepository; +import com.fowoco.server.reliability.domain.EventConsumption; +import java.util.UUID; +import org.springframework.stereotype.Repository; + +@Repository +public class JpaEventConsumptionRepository implements EventConsumptionRepository { + + private final SpringDataEventConsumptionJpaRepository repository; + + public JpaEventConsumptionRepository( + SpringDataEventConsumptionJpaRepository repository + ) { + this.repository = repository; + } + + @Override + public boolean existsByEventIdAndHandlerName(UUID eventId, String handlerName) { + return repository.existsByEventIdAndHandlerName(eventId, handlerName); + } + + @Override + public EventConsumption save(EventConsumption consumption) { + return repository.saveAndFlush(new EventConsumptionJpaEntity(consumption)).toDomain(); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaEventPublicationRepository.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaEventPublicationRepository.java new file mode 100644 index 0000000..256053f --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaEventPublicationRepository.java @@ -0,0 +1,82 @@ +package com.fowoco.server.reliability.infrastructure.persistence; + +import com.fowoco.server.reliability.application.port.EventPublicationRepository; +import com.fowoco.server.reliability.domain.EventPublication; +import com.fowoco.server.reliability.domain.EventPublicationStatus; +import java.time.Instant; +import java.util.List; +import java.util.Optional; +import java.util.UUID; +import org.springframework.data.domain.PageRequest; +import org.springframework.stereotype.Repository; + +@Repository +public class JpaEventPublicationRepository implements EventPublicationRepository { + + private static final List READY_STATUSES = List.of( + EventPublicationStatus.PENDING, + EventPublicationStatus.RETRY_WAIT + ); + private static final List OUTSTANDING_STATUSES = List.of( + EventPublicationStatus.PENDING, + EventPublicationStatus.PROCESSING, + EventPublicationStatus.RETRY_WAIT, + EventPublicationStatus.REVIEW_REQUIRED + ); + + private final SpringDataEventPublicationJpaRepository repository; + + public JpaEventPublicationRepository( + SpringDataEventPublicationJpaRepository repository + ) { + this.repository = repository; + } + + @Override + public EventPublication append(EventPublication publication) { + return repository.saveAndFlush(new EventPublicationJpaEntity(publication)).toDomain(); + } + + @Override + public EventPublication save(EventPublication publication) { + EventPublicationJpaEntity entity = repository.findById(publication.eventId()) + .orElseThrow(() -> new IllegalStateException("Event publication not found.")); + entity.apply(publication); + return repository.saveAndFlush(entity).toDomain(); + } + + @Override + public List lockClaimable(Instant now, int limit) { + return repository.lockClaimable( + READY_STATUSES, + EventPublicationStatus.PROCESSING, + now, + PageRequest.of(0, limit) + ) + .stream() + .map(EventPublicationJpaEntity::toDomain) + .toList(); + } + + @Override + public Optional findById(UUID eventId) { + return repository.findById(eventId) + .map(EventPublicationJpaEntity::toDomain); + } + + @Override + public Optional findByIdForUpdate(UUID eventId) { + return repository.findByIdForUpdate(eventId) + .map(EventPublicationJpaEntity::toDomain); + } + + @Override + public long countOutstanding() { + return repository.countByStatusIn(OUTSTANDING_STATUSES); + } + + @Override + public Optional findOldestOutstandingOccurredAt() { + return Optional.ofNullable(repository.findOldestOccurredAt(OUTSTANDING_STATUSES)); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/SpringDataEventConsumptionJpaRepository.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/SpringDataEventConsumptionJpaRepository.java new file mode 100644 index 0000000..9ba8cd4 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/SpringDataEventConsumptionJpaRepository.java @@ -0,0 +1,10 @@ +package com.fowoco.server.reliability.infrastructure.persistence; + +import java.util.UUID; +import org.springframework.data.jpa.repository.JpaRepository; + +interface SpringDataEventConsumptionJpaRepository + extends JpaRepository { + + boolean existsByEventIdAndHandlerName(UUID eventId, String handlerName); +} diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/SpringDataEventPublicationJpaRepository.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/SpringDataEventPublicationJpaRepository.java new file mode 100644 index 0000000..a493bf7 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/SpringDataEventPublicationJpaRepository.java @@ -0,0 +1,57 @@ +package com.fowoco.server.reliability.infrastructure.persistence; + +import com.fowoco.server.reliability.domain.EventPublicationStatus; +import jakarta.persistence.LockModeType; +import java.time.Instant; +import java.util.Collection; +import java.util.List; +import java.util.Optional; +import java.util.UUID; +import org.springframework.data.domain.Pageable; +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Lock; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; + +interface SpringDataEventPublicationJpaRepository + extends JpaRepository { + + @Lock(LockModeType.PESSIMISTIC_WRITE) + @Query(""" + SELECT publication + FROM EventPublicationJpaEntity publication + WHERE ( + publication.status IN :readyStatuses + AND publication.nextAttemptAt <= :now + ) OR ( + publication.status = :processingStatus + AND publication.leaseExpiresAt <= :now + ) + ORDER BY publication.occurredAt ASC, publication.eventId ASC + """) + List lockClaimable( + @Param("readyStatuses") Collection readyStatuses, + @Param("processingStatus") EventPublicationStatus processingStatus, + @Param("now") Instant now, + Pageable pageable + ); + + @Lock(LockModeType.PESSIMISTIC_WRITE) + @Query(""" + SELECT publication + FROM EventPublicationJpaEntity publication + WHERE publication.eventId = :eventId + """) + Optional findByIdForUpdate(@Param("eventId") UUID eventId); + + long countByStatusIn(Collection statuses); + + @Query(""" + SELECT MIN(publication.occurredAt) + FROM EventPublicationJpaEntity publication + WHERE publication.status IN :statuses + """) + Instant findOldestOccurredAt( + @Param("statuses") Collection statuses + ); +} diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/publishing/JpaDomainEventPublisher.java b/src/main/java/com/fowoco/server/reliability/infrastructure/publishing/JpaDomainEventPublisher.java new file mode 100644 index 0000000..f467313 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/infrastructure/publishing/JpaDomainEventPublisher.java @@ -0,0 +1,44 @@ +package com.fowoco.server.reliability.infrastructure.publishing; + +import com.fowoco.server.reliability.application.port.DomainEventPublisher; +import com.fowoco.server.reliability.application.port.EventPublicationRepository; +import com.fowoco.server.reliability.domain.DomainEventEnvelope; +import com.fowoco.server.reliability.domain.EventPublication; +import com.fowoco.server.reliability.infrastructure.serialization.EventPayloadCodec; +import java.time.Clock; +import java.time.Instant; +import org.springframework.stereotype.Component; +import org.springframework.transaction.support.TransactionSynchronizationManager; + +@Component +public class JpaDomainEventPublisher implements DomainEventPublisher { + + private final EventPublicationRepository repository; + private final EventPayloadCodec payloadCodec; + private final Clock clock; + + public JpaDomainEventPublisher( + EventPublicationRepository repository, + EventPayloadCodec payloadCodec, + Clock clock + ) { + this.repository = repository; + this.payloadCodec = payloadCodec; + this.clock = clock; + } + + @Override + public void publish(DomainEventEnvelope event) { + if (!TransactionSynchronizationManager.isActualTransactionActive()) { + throw new IllegalStateException( + "Durable domain events require an active transaction." + ); + } + Instant publishedAt = clock.instant(); + repository.append(EventPublication.pending( + event, + payloadCodec.encode(event.payload()), + publishedAt + )); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/serialization/EventPayloadCodec.java b/src/main/java/com/fowoco/server/reliability/infrastructure/serialization/EventPayloadCodec.java new file mode 100644 index 0000000..25abec8 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/infrastructure/serialization/EventPayloadCodec.java @@ -0,0 +1,35 @@ +package com.fowoco.server.reliability.infrastructure.serialization; + +import com.fowoco.server.reliability.domain.SafeEventPayload; +import java.util.Map; +import org.springframework.stereotype.Component; +import tools.jackson.core.JacksonException; +import tools.jackson.databind.ObjectMapper; + +@Component +public class EventPayloadCodec { + + private final ObjectMapper objectMapper; + + public EventPayloadCodec(ObjectMapper objectMapper) { + this.objectMapper = objectMapper; + } + + public String encode(SafeEventPayload payload) { + try { + return objectMapper.writeValueAsString(payload.values()); + } catch (JacksonException exception) { + throw new IllegalArgumentException("Event payload cannot be encoded.", exception); + } + } + + @SuppressWarnings("unchecked") + public SafeEventPayload decode(String payloadJson) { + try { + Map values = objectMapper.readValue(payloadJson, Map.class); + return SafeEventPayload.of(values.keySet(), values); + } catch (JacksonException | IllegalArgumentException exception) { + throw new EventPayloadDecodingException(exception); + } + } +} diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/serialization/EventPayloadDecodingException.java b/src/main/java/com/fowoco/server/reliability/infrastructure/serialization/EventPayloadDecodingException.java new file mode 100644 index 0000000..afb4a9c --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/infrastructure/serialization/EventPayloadDecodingException.java @@ -0,0 +1,8 @@ +package com.fowoco.server.reliability.infrastructure.serialization; + +public class EventPayloadDecodingException extends RuntimeException { + + public EventPayloadDecodingException(Throwable cause) { + super("Stored event payload is invalid.", cause); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/scheduling/OutboxScheduler.java b/src/main/java/com/fowoco/server/reliability/scheduling/OutboxScheduler.java new file mode 100644 index 0000000..472d704 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/scheduling/OutboxScheduler.java @@ -0,0 +1,26 @@ +package com.fowoco.server.reliability.scheduling; + +import com.fowoco.server.reliability.application.OutboxProcessor; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +@Component +@ConditionalOnProperty( + name = "app.reliability.outbox.enabled", + havingValue = "true", + matchIfMissing = true +) +public class OutboxScheduler { + + private final OutboxProcessor processor; + + public OutboxScheduler(OutboxProcessor processor) { + this.processor = processor; + } + + @Scheduled(fixedDelayString = "${app.reliability.outbox.poll-interval:1s}") + public void processAvailableEvents() { + processor.processAvailable(); + } +} diff --git a/src/main/java/com/fowoco/server/task/application/TaskDomainEvents.java b/src/main/java/com/fowoco/server/task/application/TaskDomainEvents.java new file mode 100644 index 0000000..bed3515 --- /dev/null +++ b/src/main/java/com/fowoco/server/task/application/TaskDomainEvents.java @@ -0,0 +1,110 @@ +package com.fowoco.server.task.application; + +import com.fowoco.server.auth.application.ActorContext; +import com.fowoco.server.common.web.RequestMetadata; +import com.fowoco.server.reliability.domain.DomainEventEnvelope; +import com.fowoco.server.reliability.domain.EventActorType; +import com.fowoco.server.reliability.domain.SafeEventPayload; +import com.fowoco.server.task.domain.Task; +import com.fowoco.server.task.domain.TaskStatus; +import java.time.Instant; +import java.util.Map; +import java.util.Set; +import java.util.UUID; + +final class TaskDomainEvents { + + private static final String PAYLOAD_VERSION = "1"; + private static final String AGGREGATE_TYPE = "Task"; + private static final Set TASK_CREATED_FIELDS = Set.of( + "task_type", + "status", + "workflow_id", + "workflow_catalog_version", + "source" + ); + private static final Set TASK_CANCELLED_FIELDS = Set.of( + "previous_status", + "status" + ); + + private TaskDomainEvents() { + } + + static DomainEventEnvelope taskCreated( + UUID eventId, + Task task, + ActorContext actor, + RequestMetadata metadata, + Instant occurredAt + ) { + return envelope( + eventId, + "TaskCreated", + task, + actor, + metadata, + occurredAt, + SafeEventPayload.of( + TASK_CREATED_FIELDS, + Map.of( + "task_type", task.taskType(), + "status", task.status(), + "workflow_id", task.workflowId(), + "workflow_catalog_version", task.workflowCatalogVersion(), + "source", task.source() + ) + ) + ); + } + + static DomainEventEnvelope taskCancelled( + UUID eventId, + Task task, + TaskStatus previousStatus, + ActorContext actor, + RequestMetadata metadata, + Instant occurredAt + ) { + return envelope( + eventId, + "TaskCancelled", + task, + actor, + metadata, + occurredAt, + SafeEventPayload.of( + TASK_CANCELLED_FIELDS, + Map.of( + "previous_status", previousStatus, + "status", task.status() + ) + ) + ); + } + + private static DomainEventEnvelope envelope( + UUID eventId, + String eventType, + Task task, + ActorContext actor, + RequestMetadata metadata, + Instant occurredAt, + SafeEventPayload payload + ) { + return new DomainEventEnvelope( + eventId, + eventType, + PAYLOAD_VERSION, + AGGREGATE_TYPE, + task.taskId(), + task.companyId(), + EventActorType.HR_USER, + actor.actorId(), + metadata.requestId(), + metadata.traceId(), + occurredAt, + payload + ); + } +} diff --git a/src/main/java/com/fowoco/server/task/application/TaskWorkflowService.java b/src/main/java/com/fowoco/server/task/application/TaskWorkflowService.java index ac6a993..00248e2 100644 --- a/src/main/java/com/fowoco/server/task/application/TaskWorkflowService.java +++ b/src/main/java/com/fowoco/server/task/application/TaskWorkflowService.java @@ -12,6 +12,7 @@ import com.fowoco.server.common.error.ApiException; import com.fowoco.server.common.id.UuidGenerator; import com.fowoco.server.common.web.RequestMetadata; +import com.fowoco.server.reliability.application.port.DomainEventPublisher; import com.fowoco.server.task.application.TaskContentCodec.EncodedTaskContent; import com.fowoco.server.task.application.error.TaskErrorCode; import com.fowoco.server.task.application.port.TaskChecklistRepository; @@ -51,6 +52,7 @@ public class TaskWorkflowService { private final WorkflowCatalogService catalogService; private final ApprovalControlPort approvalControl; private final AuditEventRepository auditRepository; + private final DomainEventPublisher eventPublisher; private final TaskContentCodec contentCodec; private final UuidGenerator uuidGenerator; private final Clock clock; @@ -64,6 +66,7 @@ public TaskWorkflowService( WorkflowCatalogService catalogService, ApprovalControlPort approvalControl, AuditEventRepository auditRepository, + DomainEventPublisher eventPublisher, TaskContentCodec contentCodec, UuidGenerator uuidGenerator, Clock clock @@ -76,6 +79,7 @@ public TaskWorkflowService( this.catalogService = catalogService; this.approvalControl = approvalControl; this.auditRepository = auditRepository; + this.eventPublisher = eventPublisher; this.contentCodec = contentCodec; this.uuidGenerator = uuidGenerator; this.clock = clock; @@ -154,6 +158,13 @@ public TaskResult create( metadata, now ); + eventPublisher.publish(TaskDomainEvents.taskCreated( + uuidGenerator.generate(), + savedTask, + actor, + metadata, + now + )); return toResult(savedTask, checklistItems, worker, workflow); } @@ -406,6 +417,14 @@ public TaskResult cancel( metadata, now ); + eventPublisher.publish(TaskDomainEvents.taskCancelled( + uuidGenerator.generate(), + savedTask, + previous, + actor, + metadata, + now + )); return toResult( savedTask, checklistRepository.findAllByTaskIdAndCompanyId(taskId, actor.companyId()), diff --git a/src/main/resources/application.yaml b/src/main/resources/application.yaml index 5be0ad9..fa70e39 100644 --- a/src/main/resources/application.yaml +++ b/src/main/resources/application.yaml @@ -37,6 +37,15 @@ springdoc: path: /swagger-ui.html app: + reliability: + outbox: + enabled: ${OUTBOX_ENABLED:true} + poll-interval: ${OUTBOX_POLL_INTERVAL:1s} + batch-size: ${OUTBOX_BATCH_SIZE:20} + lease-duration: ${OUTBOX_LEASE_DURATION:30s} + max-attempts: ${OUTBOX_MAX_ATTEMPTS:8} + initial-backoff: ${OUTBOX_INITIAL_BACKOFF:1s} + max-backoff: ${OUTBOX_MAX_BACKOFF:5m} workflow: catalog: location: ${WORKFLOW_CATALOG_LOCATION:classpath:workflow/catalog-projection.local.json} @@ -136,6 +145,9 @@ spring: console: enabled: false app: + reliability: + outbox: + enabled: false auth: jwt: secret-base64: ${JWT_SECRET_BASE64:AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA=} diff --git a/src/main/resources/db/migration/V7__create_event_outbox.sql b/src/main/resources/db/migration/V7__create_event_outbox.sql new file mode 100644 index 0000000..a14a6ea --- /dev/null +++ b/src/main/resources/db/migration/V7__create_event_outbox.sql @@ -0,0 +1,105 @@ +CREATE TABLE event_publication ( + event_id UUID NOT NULL, + company_id UUID NOT NULL, + event_type VARCHAR(100) NOT NULL, + payload_version VARCHAR(20) NOT NULL, + aggregate_type VARCHAR(60) NOT NULL, + aggregate_id UUID NOT NULL, + actor_type VARCHAR(30) NOT NULL, + actor_id UUID, + request_id VARCHAR(128) NOT NULL, + trace_id VARCHAR(64), + payload_json TEXT NOT NULL, + status VARCHAR(30) NOT NULL, + attempt_count INTEGER NOT NULL DEFAULT 0, + next_attempt_at TIMESTAMP(6) WITH TIME ZONE, + lease_owner VARCHAR(128), + lease_expires_at TIMESTAMP(6) WITH TIME ZONE, + last_error_code VARCHAR(80), + completed_at TIMESTAMP(6) WITH TIME ZONE, + occurred_at TIMESTAMP(6) WITH TIME ZONE NOT NULL, + created_at TIMESTAMP(6) WITH TIME ZONE NOT NULL, + updated_at TIMESTAMP(6) WITH TIME ZONE NOT NULL, + version BIGINT NOT NULL DEFAULT 0, + CONSTRAINT pk_event_publication PRIMARY KEY (event_id), + CONSTRAINT uq_event_publication_id_company UNIQUE (event_id, company_id), + CONSTRAINT fk_event_publication_company + FOREIGN KEY (company_id) REFERENCES company (company_id) ON DELETE RESTRICT, + CONSTRAINT ck_event_publication_type_not_blank + CHECK (CHAR_LENGTH(TRIM(event_type)) > 0), + CONSTRAINT ck_event_publication_payload_version_not_blank + CHECK (CHAR_LENGTH(TRIM(payload_version)) > 0), + CONSTRAINT ck_event_publication_aggregate_type_not_blank + CHECK (CHAR_LENGTH(TRIM(aggregate_type)) > 0), + CONSTRAINT ck_event_publication_actor_type CHECK ( + actor_type IN ('HR_USER', 'WORKER_LINK', 'AI_AGENT', 'SYSTEM_RULE') + ), + CONSTRAINT ck_event_publication_request_not_blank + CHECK (CHAR_LENGTH(TRIM(request_id)) > 0), + CONSTRAINT ck_event_publication_payload_not_blank + CHECK (CHAR_LENGTH(TRIM(payload_json)) > 0), + CONSTRAINT ck_event_publication_status CHECK ( + status IN ( + 'PENDING', + 'PROCESSING', + 'RETRY_WAIT', + 'COMPLETED', + 'REVIEW_REQUIRED' + ) + ), + CONSTRAINT ck_event_publication_attempt_count CHECK (attempt_count >= 0), + CONSTRAINT ck_event_publication_version CHECK (version >= 0), + CONSTRAINT ck_event_publication_next_attempt CHECK ( + (status IN ('PENDING', 'RETRY_WAIT') AND next_attempt_at IS NOT NULL) + OR (status NOT IN ('PENDING', 'RETRY_WAIT') AND next_attempt_at IS NULL) + ), + CONSTRAINT ck_event_publication_lease CHECK ( + (status = 'PROCESSING' + AND lease_owner IS NOT NULL + AND CHAR_LENGTH(TRIM(lease_owner)) > 0 + AND lease_expires_at IS NOT NULL) + OR (status <> 'PROCESSING' + AND lease_owner IS NULL + AND lease_expires_at IS NULL) + ), + CONSTRAINT ck_event_publication_error CHECK ( + (status IN ('RETRY_WAIT', 'REVIEW_REQUIRED') + AND last_error_code IS NOT NULL + AND CHAR_LENGTH(TRIM(last_error_code)) > 0) + OR (status NOT IN ('RETRY_WAIT', 'REVIEW_REQUIRED') + AND last_error_code IS NULL) + ), + CONSTRAINT ck_event_publication_completed CHECK ( + (status = 'COMPLETED' AND completed_at IS NOT NULL) + OR (status <> 'COMPLETED' AND completed_at IS NULL) + ), + CONSTRAINT ck_event_publication_timestamps CHECK ( + created_at >= occurred_at AND updated_at >= created_at + ) +); + +CREATE TABLE event_consumption ( + consumption_id UUID NOT NULL, + event_id UUID NOT NULL, + company_id UUID NOT NULL, + handler_name VARCHAR(120) NOT NULL, + completed_at TIMESTAMP(6) WITH TIME ZONE NOT NULL, + CONSTRAINT pk_event_consumption PRIMARY KEY (consumption_id), + CONSTRAINT uq_event_consumption_event_handler UNIQUE (event_id, handler_name), + CONSTRAINT fk_event_consumption_publication + FOREIGN KEY (event_id, company_id) + REFERENCES event_publication (event_id, company_id) ON DELETE RESTRICT, + CONSTRAINT ck_event_consumption_handler_not_blank + CHECK (CHAR_LENGTH(TRIM(handler_name)) > 0) +); + +CREATE INDEX idx_event_publication_claim + ON event_publication (status, next_attempt_at, lease_expires_at, occurred_at); +CREATE INDEX idx_event_publication_company_time + ON event_publication (company_id, occurred_at); +CREATE INDEX idx_event_publication_request + ON event_publication (company_id, request_id); +CREATE INDEX idx_event_publication_review + ON event_publication (status, updated_at); +CREATE INDEX idx_event_consumption_company_event + ON event_consumption (company_id, event_id); diff --git a/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java b/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java index 6996592..d7676c8 100644 --- a/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java +++ b/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java @@ -26,6 +26,7 @@ class PostgreSqlMigrationTests { private static final String USER_B = "22000000-0000-0000-0000-000000000002"; private static final String WORKER_A = "12000000-0000-0000-0000-000000000001"; private static final String TASK_A = "13000000-0000-0000-0000-000000000001"; + private static final String EVENT_A = "18000000-0000-0000-0000-000000000001"; private static final String TOKEN_HASH_A = "a".repeat(64); @Test @@ -72,7 +73,9 @@ private void assertSchemaContract(Connection connection) throws SQLException { "approval_request", "external_submission", "task_evidence", - "audit_event" + "audit_event", + "event_publication", + "event_consumption" ); assertThat(columnSpecs(connection, "company")) @@ -133,6 +136,22 @@ private void assertSchemaContract(Connection connection) throws SQLException { .containsEntry("company_id", new ColumnSpec("uuid", false)) .containsEntry("request_id", new ColumnSpec("varchar", false)) .containsEntry("trace_id", new ColumnSpec("varchar", true)); + assertThat(columnSpecs(connection, "event_publication")) + .containsEntry("event_id", new ColumnSpec("uuid", false)) + .containsEntry("company_id", new ColumnSpec("uuid", false)) + .containsEntry("event_type", new ColumnSpec("varchar", false)) + .containsEntry("payload_json", new ColumnSpec("text", false)) + .containsEntry("status", new ColumnSpec("varchar", false)) + .containsEntry("attempt_count", new ColumnSpec("int4", false)) + .containsEntry("next_attempt_at", new ColumnSpec("timestamptz", true)) + .containsEntry("lease_expires_at", new ColumnSpec("timestamptz", true)) + .containsEntry("version", new ColumnSpec("int8", false)); + assertThat(columnSpecs(connection, "event_consumption")) + .containsEntry("consumption_id", new ColumnSpec("uuid", false)) + .containsEntry("event_id", new ColumnSpec("uuid", false)) + .containsEntry("company_id", new ColumnSpec("uuid", false)) + .containsEntry("handler_name", new ColumnSpec("varchar", false)) + .containsEntry("completed_at", new ColumnSpec("timestamptz", false)); assertThat(constraintNames(connection)) .contains( @@ -150,7 +169,13 @@ private void assertSchemaContract(Connection connection) throws SQLException { "fk_task_created_by_company", "fk_approval_request_task_company", "fk_approval_request_requester_company", - "fk_audit_event_company" + "fk_audit_event_company", + "pk_event_publication", + "uq_event_publication_id_company", + "fk_event_publication_company", + "pk_event_consumption", + "uq_event_consumption_event_handler", + "fk_event_consumption_publication" ); assertThat(indexNames(connection)) .contains( @@ -162,7 +187,10 @@ private void assertSchemaContract(Connection connection) throws SQLException { "idx_worker_document_company_status", "idx_task_company_status_due", "idx_approval_request_task_status", - "idx_audit_event_company_time" + "idx_audit_event_company_time", + "idx_event_publication_claim", + "idx_event_publication_company_time", + "idx_event_consumption_company_event" ); } @@ -228,6 +256,27 @@ INSERT INTO approval_request ( CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP ) """.formatted(TASK_A, COMPANY_A, "f".repeat(64), USER_A)); + execute(connection, """ + INSERT INTO event_publication ( + event_id, company_id, event_type, payload_version, + aggregate_type, aggregate_id, actor_type, request_id, + payload_json, status, attempt_count, next_attempt_at, + occurred_at, created_at, updated_at + ) VALUES ( + '%s', '%s', 'TaskCreated', '1', + 'Task', '%s', 'SYSTEM_RULE', 'migration-test-request', + '{}', 'PENDING', 0, CURRENT_TIMESTAMP, + CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP + ) + """.formatted(EVENT_A, COMPANY_A, TASK_A)); + execute(connection, """ + INSERT INTO event_consumption ( + consumption_id, event_id, company_id, handler_name, completed_at + ) VALUES ( + '19000000-0000-0000-0000-000000000001', + '%s', '%s', 'migration-test-handler', CURRENT_TIMESTAMP + ) + """.formatted(EVENT_A, COMPANY_A)); assertSqlState(connection, "23505", """ INSERT INTO user_account ( @@ -356,6 +405,36 @@ INSERT INTO audit_event ( '1', 'approved', CURRENT_TIMESTAMP ) """.formatted(COMPANY_A, USER_A, TASK_A)); + assertSqlState(connection, "23503", """ + INSERT INTO event_consumption ( + consumption_id, event_id, company_id, handler_name, completed_at + ) VALUES ( + '19000000-0000-0000-0000-000000000002', + '%s', '%s', 'wrong-tenant-handler', CURRENT_TIMESTAMP + ) + """.formatted(EVENT_A, COMPANY_B)); + assertSqlState(connection, "23505", """ + INSERT INTO event_consumption ( + consumption_id, event_id, company_id, handler_name, completed_at + ) VALUES ( + '19000000-0000-0000-0000-000000000003', + '%s', '%s', 'migration-test-handler', CURRENT_TIMESTAMP + ) + """.formatted(EVENT_A, COMPANY_A)); + assertSqlState(connection, "23514", """ + INSERT INTO event_publication ( + event_id, company_id, event_type, payload_version, + aggregate_type, aggregate_id, actor_type, request_id, + payload_json, status, attempt_count, + occurred_at, created_at, updated_at + ) VALUES ( + '18000000-0000-0000-0000-000000000002', + '%s', 'TaskCreated', '1', 'Task', '%s', + 'SYSTEM_RULE', 'invalid-state-request', '{}', + 'UNKNOWN', 0, + CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP + ) + """.formatted(COMPANY_A, TASK_A)); } private Set tableNames(Connection connection) throws SQLException { diff --git a/src/test/java/com/fowoco/server/approval/ApprovalAuditIntegrationTest.java b/src/test/java/com/fowoco/server/approval/ApprovalAuditIntegrationTest.java index ad24081..3eb7845 100644 --- a/src/test/java/com/fowoco/server/approval/ApprovalAuditIntegrationTest.java +++ b/src/test/java/com/fowoco/server/approval/ApprovalAuditIntegrationTest.java @@ -67,6 +67,8 @@ class ApprovalAuditIntegrationTest { @BeforeEach void resetAndSeed() { + jdbcTemplate.update("DELETE FROM event_consumption"); + jdbcTemplate.update("DELETE FROM event_publication"); jdbcTemplate.update("DELETE FROM audit_event"); jdbcTemplate.update("DELETE FROM task_evidence"); jdbcTemplate.update("DELETE FROM external_submission"); diff --git a/src/test/java/com/fowoco/server/approval/ApprovalAuditRollbackIntegrationTest.java b/src/test/java/com/fowoco/server/approval/ApprovalAuditRollbackIntegrationTest.java index a51ca35..6f7a460 100644 --- a/src/test/java/com/fowoco/server/approval/ApprovalAuditRollbackIntegrationTest.java +++ b/src/test/java/com/fowoco/server/approval/ApprovalAuditRollbackIntegrationTest.java @@ -48,6 +48,8 @@ class ApprovalAuditRollbackIntegrationTest { @BeforeEach void seedTask() { + jdbcTemplate.update("DELETE FROM event_consumption"); + jdbcTemplate.update("DELETE FROM event_publication"); jdbcTemplate.update("DELETE FROM task_transition_history"); jdbcTemplate.update("DELETE FROM approval_request"); jdbcTemplate.update("DELETE FROM task"); diff --git a/src/test/java/com/fowoco/server/auth/AuthRefreshPostgreSqlConcurrencyTest.java b/src/test/java/com/fowoco/server/auth/AuthRefreshPostgreSqlConcurrencyTest.java index ed69621..93c7792 100644 --- a/src/test/java/com/fowoco/server/auth/AuthRefreshPostgreSqlConcurrencyTest.java +++ b/src/test/java/com/fowoco/server/auth/AuthRefreshPostgreSqlConcurrencyTest.java @@ -320,6 +320,8 @@ SELECT COUNT(*) } private void clearAuthenticationRows() { + jdbcTemplate.update("DELETE FROM event_consumption"); + jdbcTemplate.update("DELETE FROM event_publication"); jdbcTemplate.update("DELETE FROM refresh_token"); jdbcTemplate.update("DELETE FROM user_account"); jdbcTemplate.update("DELETE FROM company"); diff --git a/src/test/java/com/fowoco/server/auth/AuthSecurityIntegrationTest.java b/src/test/java/com/fowoco/server/auth/AuthSecurityIntegrationTest.java index bd98e63..bea525e 100644 --- a/src/test/java/com/fowoco/server/auth/AuthSecurityIntegrationTest.java +++ b/src/test/java/com/fowoco/server/auth/AuthSecurityIntegrationTest.java @@ -68,6 +68,8 @@ class AuthSecurityIntegrationTest { @BeforeAll void seedCompaniesAndUsers() { + jdbcTemplate.update("DELETE FROM event_consumption"); + jdbcTemplate.update("DELETE FROM event_publication"); jdbcTemplate.update("DELETE FROM refresh_token"); jdbcTemplate.update("DELETE FROM user_account"); jdbcTemplate.update("DELETE FROM company"); diff --git a/src/test/java/com/fowoco/server/auth/SignupIntegrationTest.java b/src/test/java/com/fowoco/server/auth/SignupIntegrationTest.java index 30f86fe..188e832 100644 --- a/src/test/java/com/fowoco/server/auth/SignupIntegrationTest.java +++ b/src/test/java/com/fowoco/server/auth/SignupIntegrationTest.java @@ -50,6 +50,8 @@ class SignupIntegrationTest { @BeforeEach void cleanSignupData() { + jdbcTemplate.update("DELETE FROM event_consumption"); + jdbcTemplate.update("DELETE FROM event_publication"); jdbcTemplate.update("DELETE FROM refresh_token"); jdbcTemplate.update("DELETE FROM user_account"); jdbcTemplate.update("DELETE FROM company"); diff --git a/src/test/java/com/fowoco/server/common/security/PostgreSqlTenantDatabaseContextTest.java b/src/test/java/com/fowoco/server/common/security/PostgreSqlTenantDatabaseContextTest.java index b92a64a..f3257b9 100644 --- a/src/test/java/com/fowoco/server/common/security/PostgreSqlTenantDatabaseContextTest.java +++ b/src/test/java/com/fowoco/server/common/security/PostgreSqlTenantDatabaseContextTest.java @@ -53,7 +53,9 @@ class PostgreSqlTenantDatabaseContextTest { "approval_request", "external_submission", "task_evidence", - "audit_event" + "audit_event", + "event_publication", + "event_consumption" }; private static final String TENANT_TABLE_SQL = "public." + String.join(", public.", TENANT_TABLES); diff --git a/src/test/java/com/fowoco/server/reliability/OutboxIntegrationTest.java b/src/test/java/com/fowoco/server/reliability/OutboxIntegrationTest.java new file mode 100644 index 0000000..ac0a153 --- /dev/null +++ b/src/test/java/com/fowoco/server/reliability/OutboxIntegrationTest.java @@ -0,0 +1,383 @@ +package com.fowoco.server.reliability; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import com.fowoco.server.reliability.application.NonRetryableEventHandlingException; +import com.fowoco.server.reliability.application.OutboxClaimService; +import com.fowoco.server.reliability.application.OutboxProcessor; +import com.fowoco.server.reliability.application.RetryableEventHandlingException; +import com.fowoco.server.reliability.application.port.DomainEventHandler; +import com.fowoco.server.reliability.application.port.DomainEventPublisher; +import com.fowoco.server.reliability.domain.DomainEventEnvelope; +import com.fowoco.server.reliability.domain.EventActorType; +import com.fowoco.server.reliability.domain.SafeEventPayload; +import java.time.Clock; +import java.util.Map; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.atomic.AtomicInteger; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.context.TestConfiguration; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Import; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.test.context.ActiveProfiles; +import org.springframework.transaction.support.TransactionTemplate; + +@ActiveProfiles("test") +@SpringBootTest(properties = { + "app.reliability.outbox.enabled=false", + "app.reliability.outbox.lease-duration=30s", + "app.reliability.outbox.max-attempts=3", + "app.reliability.outbox.initial-backoff=1s", + "app.reliability.outbox.max-backoff=4s" +}) +@Import(OutboxIntegrationTest.ReliabilityTestConfiguration.class) +class OutboxIntegrationTest { + + private static final UUID COMPANY_ID = + UUID.fromString("71000000-0000-0000-0000-000000000001"); + + @Autowired + private JdbcTemplate jdbcTemplate; + + @Autowired + private TransactionTemplate transactionTemplate; + + @Autowired + private DomainEventPublisher eventPublisher; + + @Autowired + private OutboxProcessor processor; + + @Autowired + private OutboxClaimService claimService; + + @Autowired + private TestEventHandler handler; + + @Autowired + private Clock clock; + + @BeforeEach + void resetState() { + jdbcTemplate.execute(""" + CREATE TABLE IF NOT EXISTS reliability_test_effect ( + event_id UUID NOT NULL PRIMARY KEY, + handled_at TIMESTAMP(6) WITH TIME ZONE NOT NULL + ) + """); + jdbcTemplate.update("DELETE FROM reliability_test_effect"); + jdbcTemplate.update("DELETE FROM event_consumption"); + jdbcTemplate.update("DELETE FROM event_publication"); + jdbcTemplate.update("DELETE FROM company WHERE company_id = ?", COMPANY_ID); + jdbcTemplate.update( + """ + INSERT INTO company ( + company_id, name, status, created_at, updated_at, version + ) VALUES (?, 'Outbox 원자성 테스트', 'ACTIVE', + CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, 0) + """, + COMPANY_ID + ); + handler.reset(); + } + + @Test + void businessChangeAndEventPublicationCommitOrRollbackTogether() { + DomainEventEnvelope rolledBackEvent = event(); + + assertThatThrownBy(() -> transactionTemplate.executeWithoutResult(status -> { + jdbcTemplate.update( + "UPDATE company SET name = '롤백 대상' WHERE company_id = ?", + COMPANY_ID + ); + eventPublisher.publish(rolledBackEvent); + throw new IllegalStateException("force rollback"); + })).isInstanceOf(IllegalStateException.class); + + assertThat(companyName()).isEqualTo("Outbox 원자성 테스트"); + assertThat(publicationCount(rolledBackEvent.eventId())).isZero(); + + DomainEventEnvelope committedEvent = event(); + transactionTemplate.executeWithoutResult(status -> { + jdbcTemplate.update( + "UPDATE company SET name = '커밋 완료' WHERE company_id = ?", + COMPANY_ID + ); + eventPublisher.publish(committedEvent); + }); + + assertThat(companyName()).isEqualTo("커밋 완료"); + assertThat(publicationStatus(committedEvent.eventId())).isEqualTo("PENDING"); + } + + @Test + void transientFailureRollsBackSideEffectAndRetriesSafely() { + handler.failRetryablyOnce(); + DomainEventEnvelope event = event(); + publish(event); + + assertThat(processor.processAvailable()).isEqualTo(1); + + assertThat(publicationStatus(event.eventId())).isEqualTo("RETRY_WAIT"); + assertThat(publicationAttemptCount(event.eventId())).isEqualTo(1); + assertThat(lastErrorCode(event.eventId())).isEqualTo("TEST_DEPENDENCY_UNAVAILABLE"); + assertThat(effectCount(event.eventId())).isZero(); + assertThat(consumptionCount(event.eventId())).isZero(); + + makeImmediatelyClaimable(event.eventId()); + + assertThat(processor.processAvailable()).isEqualTo(1); + assertThat(publicationStatus(event.eventId())).isEqualTo("COMPLETED"); + assertThat(publicationAttemptCount(event.eventId())).isEqualTo(2); + assertThat(effectCount(event.eventId())).isEqualTo(1); + assertThat(consumptionCount(event.eventId())).isEqualTo(1); + assertThat(handler.invocationCount()).isEqualTo(2); + } + + @Test + void completedHandlerIsNotExecutedAgainAfterDuplicateDelivery() { + DomainEventEnvelope event = event(); + publish(event); + processor.processAvailable(); + + jdbcTemplate.update( + """ + UPDATE event_publication + SET status = 'RETRY_WAIT', + next_attempt_at = CURRENT_TIMESTAMP, + completed_at = NULL, + last_error_code = 'DUPLICATE_DELIVERY_SIMULATION', + updated_at = CURRENT_TIMESTAMP, + version = version + 1 + WHERE event_id = ? + """, + event.eventId() + ); + + assertThat(processor.processAvailable()).isEqualTo(1); + + assertThat(publicationStatus(event.eventId())).isEqualTo("COMPLETED"); + assertThat(publicationAttemptCount(event.eventId())).isEqualTo(2); + assertThat(effectCount(event.eventId())).isEqualTo(1); + assertThat(consumptionCount(event.eventId())).isEqualTo(1); + assertThat(handler.invocationCount()).isEqualTo(1); + } + + @Test + void expiredLeaseIsRecoveredAfterServerRestart() { + DomainEventEnvelope event = event(); + publish(event); + + assertThat(claimService.claimBatch("stopped-server")) + .containsExactly(event.eventId()); + assertThat(publicationStatus(event.eventId())).isEqualTo("PROCESSING"); + + jdbcTemplate.update( + """ + UPDATE event_publication + SET lease_expires_at = DATEADD('SECOND', -1, CURRENT_TIMESTAMP), + updated_at = CURRENT_TIMESTAMP, + version = version + 1 + WHERE event_id = ? + """, + event.eventId() + ); + + assertThat(processor.processAvailable()).isEqualTo(1); + + assertThat(publicationStatus(event.eventId())).isEqualTo("COMPLETED"); + assertThat(publicationAttemptCount(event.eventId())).isEqualTo(2); + assertThat(effectCount(event.eventId())).isEqualTo(1); + } + + @Test + void nonRetryableFailureMovesEventToManualReview() { + handler.failPermanently(); + DomainEventEnvelope event = event(); + publish(event); + + assertThat(processor.processAvailable()).isEqualTo(1); + + assertThat(publicationStatus(event.eventId())).isEqualTo("REVIEW_REQUIRED"); + assertThat(lastErrorCode(event.eventId())).isEqualTo("TEST_PAYLOAD_REJECTED"); + assertThat(effectCount(event.eventId())).isZero(); + assertThat(consumptionCount(event.eventId())).isZero(); + } + + @Test + void publishingOutsideBusinessTransactionIsRejected() { + assertThatThrownBy(() -> eventPublisher.publish(event())) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("active transaction"); + } + + private void publish(DomainEventEnvelope event) { + transactionTemplate.executeWithoutResult(status -> eventPublisher.publish(event)); + } + + private DomainEventEnvelope event() { + return new DomainEventEnvelope( + UUID.randomUUID(), + TestEventHandler.EVENT_TYPE, + "1", + "ReliabilityProbe", + UUID.randomUUID(), + COMPANY_ID, + EventActorType.SYSTEM_RULE, + null, + "reliability-test-" + UUID.randomUUID(), + null, + clock.instant(), + SafeEventPayload.of( + Set.of("result"), + Map.of("result", "READY") + ) + ); + } + + private String companyName() { + return jdbcTemplate.queryForObject( + "SELECT name FROM company WHERE company_id = ?", + String.class, + COMPANY_ID + ); + } + + private int publicationCount(UUID eventId) { + return jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM event_publication WHERE event_id = ?", + Integer.class, + eventId + ); + } + + private String publicationStatus(UUID eventId) { + return jdbcTemplate.queryForObject( + "SELECT status FROM event_publication WHERE event_id = ?", + String.class, + eventId + ); + } + + private int publicationAttemptCount(UUID eventId) { + return jdbcTemplate.queryForObject( + "SELECT attempt_count FROM event_publication WHERE event_id = ?", + Integer.class, + eventId + ); + } + + private String lastErrorCode(UUID eventId) { + return jdbcTemplate.queryForObject( + "SELECT last_error_code FROM event_publication WHERE event_id = ?", + String.class, + eventId + ); + } + + private int effectCount(UUID eventId) { + return jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM reliability_test_effect WHERE event_id = ?", + Integer.class, + eventId + ); + } + + private int consumptionCount(UUID eventId) { + return jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM event_consumption WHERE event_id = ?", + Integer.class, + eventId + ); + } + + private void makeImmediatelyClaimable(UUID eventId) { + jdbcTemplate.update( + """ + UPDATE event_publication + SET next_attempt_at = CURRENT_TIMESTAMP, + updated_at = CURRENT_TIMESTAMP, + version = version + 1 + WHERE event_id = ? + """, + eventId + ); + } + + @TestConfiguration(proxyBeanMethods = false) + static class ReliabilityTestConfiguration { + + @Bean + TestEventHandler testEventHandler(JdbcTemplate jdbcTemplate) { + return new TestEventHandler(jdbcTemplate); + } + } + + static final class TestEventHandler implements DomainEventHandler { + + private static final String EVENT_TYPE = "ReliabilityTestRequested"; + + private final JdbcTemplate jdbcTemplate; + private final AtomicInteger invocationCount = new AtomicInteger(); + private final AtomicInteger retryableFailuresRemaining = new AtomicInteger(); + private volatile boolean permanentFailure; + + private TestEventHandler(JdbcTemplate jdbcTemplate) { + this.jdbcTemplate = jdbcTemplate; + } + + @Override + public String handlerName() { + return "reliability-test-handler-v1"; + } + + @Override + public boolean supports(String eventType) { + return EVENT_TYPE.equals(eventType); + } + + @Override + public void handle(DomainEventEnvelope event) { + invocationCount.incrementAndGet(); + jdbcTemplate.update( + """ + INSERT INTO reliability_test_effect (event_id, handled_at) + VALUES (?, CURRENT_TIMESTAMP) + """, + event.eventId() + ); + if (permanentFailure) { + throw new NonRetryableEventHandlingException("TEST_PAYLOAD_REJECTED"); + } + if (retryableFailuresRemaining.getAndUpdate(value -> Math.max(0, value - 1)) > 0) { + throw new RetryableEventHandlingException( + "TEST_DEPENDENCY_UNAVAILABLE" + ); + } + } + + void failRetryablyOnce() { + retryableFailuresRemaining.set(1); + } + + void failPermanently() { + permanentFailure = true; + } + + int invocationCount() { + return invocationCount.get(); + } + + void reset() { + invocationCount.set(0); + retryableFailuresRemaining.set(0); + permanentFailure = false; + } + } +} diff --git a/src/test/java/com/fowoco/server/reliability/application/OutboxBackoffPolicyTest.java b/src/test/java/com/fowoco/server/reliability/application/OutboxBackoffPolicyTest.java new file mode 100644 index 0000000..589ec1d --- /dev/null +++ b/src/test/java/com/fowoco/server/reliability/application/OutboxBackoffPolicyTest.java @@ -0,0 +1,24 @@ +package com.fowoco.server.reliability.application; + +import static org.assertj.core.api.Assertions.assertThat; + +import com.fowoco.server.reliability.config.OutboxProperties; +import java.time.Duration; +import org.junit.jupiter.api.Test; + +class OutboxBackoffPolicyTest { + + @Test + void doublesDelayPerAttemptAndCapsAtMaximum() { + OutboxProperties properties = new OutboxProperties(); + properties.setInitialBackoff(Duration.ofSeconds(2)); + properties.setMaxBackoff(Duration.ofSeconds(10)); + OutboxBackoffPolicy policy = new OutboxBackoffPolicy(properties); + + assertThat(policy.delayForAttempt(1)).isEqualTo(Duration.ofSeconds(2)); + assertThat(policy.delayForAttempt(2)).isEqualTo(Duration.ofSeconds(4)); + assertThat(policy.delayForAttempt(3)).isEqualTo(Duration.ofSeconds(8)); + assertThat(policy.delayForAttempt(4)).isEqualTo(Duration.ofSeconds(10)); + assertThat(policy.delayForAttempt(100)).isEqualTo(Duration.ofSeconds(10)); + } +} diff --git a/src/test/java/com/fowoco/server/reliability/domain/EventPublicationTest.java b/src/test/java/com/fowoco/server/reliability/domain/EventPublicationTest.java new file mode 100644 index 0000000..5bde7ef --- /dev/null +++ b/src/test/java/com/fowoco/server/reliability/domain/EventPublicationTest.java @@ -0,0 +1,85 @@ +package com.fowoco.server.reliability.domain; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.time.Duration; +import java.time.Instant; +import java.util.UUID; +import org.junit.jupiter.api.Test; + +class EventPublicationTest { + + private static final Instant OCCURRED_AT = Instant.parse("2026-07-24T00:00:00Z"); + + @Test + void anExpiredLeaseCanBeClaimedByAnotherWorker() { + EventPublication publication = pendingPublication(); + publication.claim("server-blue", OCCURRED_AT, Duration.ofSeconds(30)); + + assertThat(publication.isClaimableAt(OCCURRED_AT.plusSeconds(29))).isFalse(); + assertThat(publication.isClaimableAt(OCCURRED_AT.plusSeconds(30))).isTrue(); + + publication.claim( + "server-green", + OCCURRED_AT.plusSeconds(30), + Duration.ofSeconds(30) + ); + + assertThat(publication.status()).isEqualTo(EventPublicationStatus.PROCESSING); + assertThat(publication.leaseOwner()).isEqualTo("server-green"); + assertThat(publication.attemptCount()).isEqualTo(2); + } + + @Test + void aWorkerCannotCompleteAnotherWorkersActiveLease() { + EventPublication publication = pendingPublication(); + publication.claim("server-blue", OCCURRED_AT, Duration.ofSeconds(30)); + + assertThatThrownBy(() -> publication.complete( + "server-green", + OCCURRED_AT.plusSeconds(1) + )) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("lease"); + } + + @Test + void retryClearsLeaseAndSchedulesTheNextAttempt() { + EventPublication publication = pendingPublication(); + publication.claim("server-blue", OCCURRED_AT, Duration.ofSeconds(30)); + Instant nextAttempt = OCCURRED_AT.plusSeconds(10); + + publication.retry( + "server-blue", + "DEPENDENCY_TEMPORARY_FAILURE", + nextAttempt, + OCCURRED_AT.plusSeconds(1) + ); + + assertThat(publication.status()).isEqualTo(EventPublicationStatus.RETRY_WAIT); + assertThat(publication.leaseOwner()).isNull(); + assertThat(publication.nextAttemptAt()).isEqualTo(nextAttempt); + assertThat(publication.lastErrorCode()).isEqualTo("DEPENDENCY_TEMPORARY_FAILURE"); + } + + private EventPublication pendingPublication() { + UUID companyId = UUID.randomUUID(); + UUID eventId = UUID.randomUUID(); + DomainEventEnvelope event = new DomainEventEnvelope( + eventId, + "TaskCreated", + "1", + "Task", + UUID.randomUUID(), + companyId, + EventActorType.SYSTEM_RULE, + null, + "request-1", + "0123456789abcdef0123456789abcdef", + OCCURRED_AT, + SafeEventPayload.empty() + ); + return EventPublication.pending(event, "{}", OCCURRED_AT); + } +} diff --git a/src/test/java/com/fowoco/server/reliability/domain/SafeEventPayloadTest.java b/src/test/java/com/fowoco/server/reliability/domain/SafeEventPayloadTest.java new file mode 100644 index 0000000..ec00a9b --- /dev/null +++ b/src/test/java/com/fowoco/server/reliability/domain/SafeEventPayloadTest.java @@ -0,0 +1,78 @@ +package com.fowoco.server.reliability.domain; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.time.LocalDate; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.UUID; +import org.junit.jupiter.api.Test; + +class SafeEventPayloadTest { + + @Test + void acceptsOnlyAllowListedScalarBusinessFields() { + UUID workflowId = UUID.randomUUID(); + + SafeEventPayload payload = SafeEventPayload.of( + Set.of("status", "due_date", "workflow_ids"), + Map.of( + "status", EventPublicationStatus.PENDING, + "due_date", LocalDate.of(2027, 3, 1), + "workflow_ids", List.of(workflowId) + ) + ); + + assertThat(payload.values()) + .containsEntry("status", "PENDING") + .containsEntry("due_date", "2027-03-01") + .containsEntry("workflow_ids", List.of(workflowId.toString())); + } + + @Test + void rejectsAFieldThatWasNotExplicitlyAllowed() { + assertThatThrownBy(() -> SafeEventPayload.of( + Set.of("status"), + Map.of("status", "READY", "task_title", "체류기간 연장") + )) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("not allow-listed"); + } + + @Test + void rejectsSensitiveFieldNamesEvenWhenAllowListed() { + assertThatThrownBy(() -> SafeEventPayload.of( + Set.of("passport_number"), + Map.of("passport_number", "M12345678") + )) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Sensitive"); + } + + @Test + void rejectsSensitiveValuesHiddenUnderAnOtherwiseSafeField() { + assertThatThrownBy(() -> SafeEventPayload.of( + Set.of("result"), + Map.of("result", "연락처 010-1234-5678") + )) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Sensitive"); + } + + @Test + void rejectsNestedObjectsAndOversizedCollections() { + assertThatThrownBy(() -> SafeEventPayload.of( + Set.of("metadata"), + Map.of("metadata", Map.of("status", "READY")) + )).isInstanceOf(IllegalArgumentException.class); + + assertThatThrownBy(() -> SafeEventPayload.of( + Set.of("workflow_ids"), + Map.of("workflow_ids", java.util.stream.IntStream.range(0, 51).boxed().toList()) + )) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("too large"); + } +} diff --git a/src/test/java/com/fowoco/server/task/TaskWorkflowIntegrationTest.java b/src/test/java/com/fowoco/server/task/TaskWorkflowIntegrationTest.java index 72f8f09..792c9d4 100644 --- a/src/test/java/com/fowoco/server/task/TaskWorkflowIntegrationTest.java +++ b/src/test/java/com/fowoco/server/task/TaskWorkflowIntegrationTest.java @@ -59,6 +59,8 @@ class TaskWorkflowIntegrationTest { @BeforeEach void resetAndSeed() { + jdbcTemplate.update("DELETE FROM event_consumption"); + jdbcTemplate.update("DELETE FROM event_publication"); jdbcTemplate.update("DELETE FROM audit_event"); jdbcTemplate.update("DELETE FROM task_evidence"); jdbcTemplate.update("DELETE FROM external_submission"); @@ -98,6 +100,12 @@ void supportsTheCatalogTaskChecklistAndCancelApiFlow() throws Exception { assertThat(JsonPath.read(created.body(), "$.status")).isEqualTo("DRAFT"); assertThat(JsonPath.read(created.body(), "$.workflow_catalog_version")) .isEqualTo("0.2.0"); + assertThat(jdbcTemplate.queryForList( + "SELECT event_type FROM event_publication " + + "WHERE aggregate_id = ? ORDER BY occurred_at", + String.class, + taskId + )).containsExactly("TaskCreated"); List checklistIds = JsonPath.read( created.body(), "$.checklist_items[*].checklist_item_id" @@ -130,6 +138,12 @@ void supportsTheCatalogTaskChecklistAndCancelApiFlow() throws Exception { ); assertThat(updated.statusCode()).isEqualTo(200); assertThat(JsonPath.read(updated.body(), "$.version").longValue()).isEqualTo(1); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM event_publication " + + "WHERE aggregate_id = ? AND event_type = 'TaskCancelled'", + Integer.class, + taskId + )).isZero(); HttpResponse checklist = patch( "/api/v1/tasks/" + taskId + "/checklist-items/" + checklistIds.get(0), @@ -162,6 +176,21 @@ void supportsTheCatalogTaskChecklistAndCancelApiFlow() throws Exception { Integer.class, taskId )).isEqualTo(4); + assertThat(jdbcTemplate.queryForList( + "SELECT event_type FROM event_publication " + + "WHERE aggregate_id = ?", + String.class, + taskId + )).containsExactlyInAnyOrder("TaskCreated", "TaskCancelled"); + assertThat(jdbcTemplate.queryForObject( + "SELECT payload_json FROM event_publication " + + "WHERE aggregate_id = ? AND event_type = 'TaskCancelled'", + String.class, + taskId + )) + .contains("\"previous_status\":\"NEEDS_INFO\"") + .contains("\"status\":\"CANCELLED\"") + .doesNotContain("display_name", "passport", "phone", "token"); } @Test diff --git a/src/test/java/com/fowoco/server/worker/WorkerDocumentSecurityIntegrationTest.java b/src/test/java/com/fowoco/server/worker/WorkerDocumentSecurityIntegrationTest.java index b80c4fd..299e7c2 100644 --- a/src/test/java/com/fowoco/server/worker/WorkerDocumentSecurityIntegrationTest.java +++ b/src/test/java/com/fowoco/server/worker/WorkerDocumentSecurityIntegrationTest.java @@ -48,6 +48,8 @@ class WorkerDocumentSecurityIntegrationTest { @BeforeAll void seedCompaniesAndUsers() { + jdbcTemplate.update("DELETE FROM event_consumption"); + jdbcTemplate.update("DELETE FROM event_publication"); jdbcTemplate.update("DELETE FROM task_evidence"); jdbcTemplate.update("DELETE FROM external_submission"); jdbcTemplate.update("DELETE FROM approval_request"); diff --git a/src/test/java/com/fowoco/server/worker/WorkerSecurityIntegrationTest.java b/src/test/java/com/fowoco/server/worker/WorkerSecurityIntegrationTest.java index 00a7e2c..abdaee6 100644 --- a/src/test/java/com/fowoco/server/worker/WorkerSecurityIntegrationTest.java +++ b/src/test/java/com/fowoco/server/worker/WorkerSecurityIntegrationTest.java @@ -61,6 +61,8 @@ class WorkerSecurityIntegrationTest { @BeforeAll void seedCompaniesAndUsers() { + jdbcTemplate.update("DELETE FROM event_consumption"); + jdbcTemplate.update("DELETE FROM event_publication"); jdbcTemplate.update("DELETE FROM task_evidence"); jdbcTemplate.update("DELETE FROM external_submission"); jdbcTemplate.update("DELETE FROM approval_request");