diff --git a/.github/workflows/database-docs.yml b/.github/workflows/database-docs.yml
index 95ba84d..67cf4d4 100644
--- a/.github/workflows/database-docs.yml
+++ b/.github/workflows/database-docs.yml
@@ -8,6 +8,7 @@ on:
- 'src/main/java/**'
- 'src/main/resources/application.yaml'
- 'src/main/resources/db/migration/**'
+ - 'src/main/resources/db/migration-postgresql/**'
- 'scripts/api-docs/**'
- 'scripts/db-docs/**'
- '.github/workflows/database-docs.yml'
@@ -18,6 +19,7 @@ on:
- 'src/main/java/**'
- 'src/main/resources/application.yaml'
- 'src/main/resources/db/migration/**'
+ - 'src/main/resources/db/migration-postgresql/**'
- 'scripts/api-docs/**'
- 'scripts/db-docs/**'
- '.github/workflows/database-docs.yml'
diff --git a/docs/database-documentation.md b/docs/database-documentation.md
index dc95e72..6fe99f4 100644
--- a/docs/database-documentation.md
+++ b/docs/database-documentation.md
@@ -5,7 +5,8 @@ PostgreSQL에 처음부터 적용한 결과로 생성합니다.
- 팀 공유 사이트:
- API Swagger 문서:
-- 변경의 원본: `src/main/resources/db/migration`
+- 공통 변경의 원본: `src/main/resources/db/migration`
+- PostgreSQL 전용 변경의 원본: `src/main/resources/db/migration-postgresql`
- 구조 결정의 원본: `docs/adr`
문서는 구조를 쉽게 찾기 위한 보조 수단입니다. 문서 화면에서 DB를 변경할 수
@@ -29,6 +30,7 @@ PostgreSQL에 처음부터 적용한 결과로 생성합니다.
```text
src/main/resources/db/migration/**
+src/main/resources/db/migration-postgresql/**
scripts/db-docs/**
.github/workflows/database-docs.yml
```
@@ -37,7 +39,7 @@ Workflow는 다음 순서로 동작합니다.
```text
빈 PostgreSQL 시작
-→ Flyway migrate
+→ 공통·PostgreSQL 전용 Flyway migrate
→ Flyway validate
→ SchemaSpy HTML 생성
→ Migration 이력 페이지 생성
diff --git a/docs/database/postgresql-rls-rollout.md b/docs/database/postgresql-rls-rollout.md
index 76dcc24..2c25ba2 100644
--- a/docs/database/postgresql-rls-rollout.md
+++ b/docs/database/postgresql-rls-rollout.md
@@ -16,18 +16,33 @@ RLS는 기존 `ActorContext`, Repository의 `company_id` 조건, tenant-aware DB
placeholder를 만들지 않습니다.
현재 기반 단계에서는 runtime/Flyway 설정 경계, PostgreSQL 전용 Flyway location,
-transaction-local tenant context와 connection pool 비누수 테스트만 준비합니다.
-아직 policy를 만들거나 RLS를 활성화하지 않습니다.
-
-현재 `main`의 V1~V7에는 아래 14개 tenant table이 존재합니다. 기반 단계의 제한
-role 테스트는 이 전체 범위에 업무 DML만 허용하고, table owner·DDL·`TRUNCATE`·
+transaction-local tenant context와 connection pool 비누수 테스트를 준비했습니다.
+JWT로 인증된 Worker·Task·Approval·Audit 업무 transaction은 요청 값이 아니라
+`ActorContext.companyId`를 transaction-local context의 신뢰 원본으로 사용합니다.
+H2는 PostgreSQL custom setting을 흉내 내지 않고 transaction 경계만 검증합니다.
+`V10`에서 bootstrap 함수와 tenant 테이블 RLS policy를 생성했으며, RLS는 아직 활성화하지 않았습니다.
+
+로그인·Refresh Token·Logout은 tenant context가 생기기 전 최소 bootstrap 조회가
+필요합니다. Issue #34 작성 뒤 추가된 사업장 회원가입도 새 tenant 행을 처음 만드는
+별도 bootstrap 흐름으로 함께 검토해야 합니다. Worker Link는 해당 기능이 구현된 뒤
+같은 기준으로 확장합니다.
+
+현재 `main`의 V1~V9에는 `company_id`를 직접 보유한 아래 16개 tenant table과,
+부모 초안의 tenant를 따르는 `document_request_draft_type`이 존재합니다. 기반 단계의
+제한 role 테스트는 이 전체 범위에 업무 DML만 허용하고, table owner·DDL·`TRUNCATE`·
`REFERENCES` 권한과 RLS 우회 권한이 없음을 확인합니다.
- `company`, `user_account`, `refresh_token`
-- `worker`, `worker_document`
+- `worker`, `worker_document`, `stored_file`
- `task`, `task_checklist_item`, `task_transition_history`
- `approval_request`, `external_submission`, `task_evidence`, `audit_event`
- `event_publication`, `event_consumption`
+- `document_request_draft`, `document_request_draft_type`
+
+`document_request_draft_type`에는 `company_id`가 없으므로 부모
+`document_request_draft`의 `draft_id`와 현재 tenant context를 확인하는 `EXISTS`
+policy를 사용합니다. 이 예외는 #57의 스키마와 JPA collection-table 계약을 유지하면서
+자식 테이블 직접 접근도 격리하기 위한 것입니다.
`event_publication`은 여러 tenant의 미완료 row를 찾는 background queue이므로 일반
요청 table과 같은 policy를 바로 활성화하면 worker가 아무 이벤트도 claim하지 못할 수
@@ -60,13 +75,16 @@ DDL, `TRUNCATE`, `REFERENCES` 권한을 갖지 않습니다. 실제 값은 배
## Staging 적용 순서
1. 대상 table과 tenant-aware FK·UNIQUE 제약이 `main`에 병합됐는지 확인합니다.
-2. 준비 migration에서 bootstrap 함수와 policy를 만들되 RLS는 켜지 않습니다.
-3. tenant context와 bootstrap 호환 코드를 배포합니다.
-4. #9에서 분리된 runtime role, 최소 GRANT와 Secret을 적용합니다.
-5. RLS 비활성 상태에서 Login·Refresh·tenant A/B·connection pool 회귀 테스트를
+2. 인증된 업무 transaction이 `ActorContext.companyId`를 context로 설정하는지
+ 검증합니다.
+3. 준비 migration에서 Login·Refresh·Outbox bootstrap 함수와 tenant 테이블 policy를 생성하되,
+ RLS는 활성화하지 않습니다.
+4. bootstrap 호환 코드를 배포합니다.
+5. #9에서 분리된 runtime role, 최소 GRANT와 Secret을 적용합니다.
+6. RLS 비활성 상태에서 Signup·Login·Refresh·tenant A/B·connection pool 회귀 테스트를
실행합니다.
-6. 별도 forward migration으로 `ENABLE ROW LEVEL SECURITY`를 적용합니다.
-7. 제한된 runtime role로 Smoke Test를 실행합니다.
+7. 별도 forward migration으로 `ENABLE ROW LEVEL SECURITY`를 적용합니다.
+8. 제한된 runtime role로 Smoke Test를 실행합니다.
## Smoke Test
diff --git a/scripts/db-docs/generate-site.test.mjs b/scripts/db-docs/generate-site.test.mjs
index 47b65f2..f0bb7eb 100644
--- a/scripts/db-docs/generate-site.test.mjs
+++ b/scripts/db-docs/generate-site.test.mjs
@@ -15,7 +15,7 @@ test('Flyway JSON을 안전한 DB 문서 사이트로 변환한다', async () =>
await mkdir(path.join(output, 'schema'), { recursive: true })
await writeFile(path.join(output, 'schema', 'index.html'), 'SchemaSpy')
await writeFile(infoFile, JSON.stringify({
- schemaVersion: '5',
+ schemaVersion: '10',
schemaName: 'public',
flywayVersion: '12.4.0',
migrations: [
@@ -37,6 +37,15 @@ test('Flyway JSON을 안전한 DB 문서 사이트로 변환한다', async () =>
executionTime: 0,
filepath: '/private/path/V6__next.sql',
},
+ {
+ version: '10',
+ description: 'prepare postgresql rls',
+ type: 'SQL',
+ state: 'Success',
+ installedOnUTC: '2026-07-29T00:00:00Z',
+ executionTime: 21,
+ filepath: '/flyway/sql/postgresql/V10__prepare_postgresql_rls.sql',
+ },
],
}))
@@ -55,12 +64,14 @@ test('Flyway JSON을 안전한 DB 문서 사이트로 변환한다', async () =>
const metadata = JSON.parse(await readFile(path.join(output, 'metadata.json'), 'utf8'))
assert.match(index, /현재 Schema Version/)
- assert.match(index, /성공 1개 · 대기 1개/)
+ assert.match(index, /성공 2개 · 대기 1개/)
assert.match(migrations, /baseline <safe>/)
+ assert.match(migrations, /prepare postgresql rls/)
assert.doesNotMatch(migrations, /private\/path/)
- assert.equal(metadata.schema_version, '5')
+ assert.doesNotMatch(migrations, /flyway\/sql\/postgresql/)
+ assert.equal(metadata.schema_version, '10')
assert.deepEqual(metadata.migration_counts, {
- success: 1,
+ success: 2,
pending: 1,
attention_required: 0,
})
diff --git a/scripts/db-docs/generate.sh b/scripts/db-docs/generate.sh
index c0e49b4..3b4c6da 100755
--- a/scripts/db-docs/generate.sh
+++ b/scripts/db-docs/generate.sh
@@ -4,7 +4,8 @@ set -euo pipefail
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
REPOSITORY_ROOT="$(cd "${SCRIPT_DIR}/../.." && pwd)"
OUTPUT_ROOT="${REPOSITORY_ROOT}/build/db-docs"
-MIGRATION_DIRECTORY="${REPOSITORY_ROOT}/src/main/resources/db/migration"
+COMMON_MIGRATION_DIRECTORY="${REPOSITORY_ROOT}/src/main/resources/db/migration"
+POSTGRESQL_MIGRATION_DIRECTORY="${REPOSITORY_ROOT}/src/main/resources/db/migration-postgresql"
FLYWAY_IMAGE="${DB_DOCS_FLYWAY_IMAGE:-flyway/flyway:12.4.0}"
SCHEMASPY_IMAGE="${DB_DOCS_SCHEMASPY_IMAGE:-schemaspy/schemaspy:7.0.2}"
@@ -31,6 +32,10 @@ if ! command -v node >/dev/null 2>&1; then
echo "[db-docs] Node.js를 찾지 못했습니다. Node.js 24 이상을 설치해 주세요." >&2
exit 1
fi
+if [[ ! -d "${COMMON_MIGRATION_DIRECTORY}" || ! -d "${POSTGRESQL_MIGRATION_DIRECTORY}" ]]; then
+ echo "[db-docs] 공통·PostgreSQL 전용 Migration 경로가 모두 필요합니다." >&2
+ exit 1
+fi
if ! docker info >/dev/null 2>&1; then
echo "[db-docs] Docker가 실행 중이 아닙니다." >&2
exit 1
@@ -79,11 +84,15 @@ mkdir -p "${OUTPUT_ROOT}/site/schema"
chmod 0777 "${OUTPUT_ROOT}/site/schema"
JDBC_URL="jdbc:postgresql://${DATABASE_HOST}:${DATABASE_PORT}/${DATABASE_NAME}"
+FLYWAY_MOUNT_ARGUMENTS=(
+ -v "${COMMON_MIGRATION_DIRECTORY}:/flyway/sql/common:ro"
+ -v "${POSTGRESQL_MIGRATION_DIRECTORY}:/flyway/sql/postgresql:ro"
+)
FLYWAY_ARGUMENTS=(
"-url=${JDBC_URL}"
"-user=${DATABASE_USER}"
"-password=${DB_DOCS_PASSWORD}"
- "-locations=filesystem:/flyway/sql"
+ "-locations=filesystem:/flyway/sql/common,filesystem:/flyway/sql/postgresql"
"-defaultSchema=public"
"-schemas=public"
"-connectRetries=20"
@@ -92,7 +101,7 @@ FLYWAY_ARGUMENTS=(
echo "[db-docs] 빈 PostgreSQL에 Flyway Migration을 적용합니다."
docker run --rm \
"${NETWORK_ARGUMENTS[@]}" \
- -v "${MIGRATION_DIRECTORY}:/flyway/sql:ro" \
+ "${FLYWAY_MOUNT_ARGUMENTS[@]}" \
"${FLYWAY_IMAGE}" \
"${FLYWAY_ARGUMENTS[@]}" \
migrate
@@ -100,14 +109,14 @@ docker run --rm \
echo "[db-docs] 적용된 Migration과 저장소 checksum을 검증합니다."
docker run --rm \
"${NETWORK_ARGUMENTS[@]}" \
- -v "${MIGRATION_DIRECTORY}:/flyway/sql:ro" \
+ "${FLYWAY_MOUNT_ARGUMENTS[@]}" \
"${FLYWAY_IMAGE}" \
"${FLYWAY_ARGUMENTS[@]}" \
validate
docker run --rm \
"${NETWORK_ARGUMENTS[@]}" \
- -v "${MIGRATION_DIRECTORY}:/flyway/sql:ro" \
+ "${FLYWAY_MOUNT_ARGUMENTS[@]}" \
"${FLYWAY_IMAGE}" \
"${FLYWAY_ARGUMENTS[@]}" \
-outputType=json \
diff --git a/src/main/java/com/fowoco/server/approval/application/ApprovalService.java b/src/main/java/com/fowoco/server/approval/application/ApprovalService.java
index 02d73bf..df6d13d 100644
--- a/src/main/java/com/fowoco/server/approval/application/ApprovalService.java
+++ b/src/main/java/com/fowoco/server/approval/application/ApprovalService.java
@@ -17,6 +17,7 @@
import com.fowoco.server.auth.domain.UserRole;
import com.fowoco.server.common.error.ApiException;
import com.fowoco.server.common.id.UuidGenerator;
+import com.fowoco.server.common.security.TenantDatabaseContext;
import com.fowoco.server.common.web.RequestMetadata;
import com.fowoco.server.task.application.error.TaskErrorCode;
import com.fowoco.server.task.application.TaskReadinessChecker;
@@ -40,6 +41,7 @@ public class ApprovalService implements ApprovalControlPort {
private static final String AUDIT_EVENT_VERSION = "1";
private final ActorAuthorizer actorAuthorizer;
+ private final TenantDatabaseContext tenantDatabaseContext;
private final TaskRepository taskRepository;
private final TaskTransitionRecorder transitionRecorder;
private final TaskReadinessChecker taskReadinessChecker;
@@ -53,6 +55,7 @@ public class ApprovalService implements ApprovalControlPort {
public ApprovalService(
ActorAuthorizer actorAuthorizer,
+ TenantDatabaseContext tenantDatabaseContext,
TaskRepository taskRepository,
TaskTransitionRecorder transitionRecorder,
TaskReadinessChecker taskReadinessChecker,
@@ -65,6 +68,7 @@ public ApprovalService(
Clock clock
) {
this.actorAuthorizer = actorAuthorizer;
+ this.tenantDatabaseContext = tenantDatabaseContext;
this.taskRepository = taskRepository;
this.transitionRecorder = transitionRecorder;
this.taskReadinessChecker = taskReadinessChecker;
@@ -84,6 +88,7 @@ public ApprovalResult requestApproval(
ActorContext actor,
RequestMetadata metadata
) {
+ bindTenant(actor);
actorAuthorizer.requireHrWrite(actor);
approvalRepository.findPendingByTaskIdAndCompanyId(taskId, actor.companyId())
.ifPresent(ignored -> {
@@ -141,6 +146,7 @@ public ApprovalResult approve(
ActorContext actor,
RequestMetadata metadata
) {
+ bindTenant(actor);
actorAuthorizer.requireHrWrite(actor);
Task task = requireTask(taskId, actor.companyId());
requireTaskVersion(task, command.expectedVersion());
@@ -178,6 +184,7 @@ public ApprovalResult reject(
ActorContext actor,
RequestMetadata metadata
) {
+ bindTenant(actor);
actorAuthorizer.requireHrWrite(actor);
Task task = requireTask(taskId, actor.companyId());
requireTaskVersion(task, command.expectedVersion());
@@ -207,6 +214,7 @@ public TaskActionResult recordExternalSubmission(
ActorContext actor,
RequestMetadata metadata
) {
+ bindTenant(actor);
actorAuthorizer.requireHrWrite(actor);
Task task = requireTask(taskId, actor.companyId());
requireValidApproval(task);
@@ -257,6 +265,7 @@ public TaskActionResult recordEvidence(
ActorContext actor,
RequestMetadata metadata
) {
+ bindTenant(actor);
actorAuthorizer.requireHrWrite(actor);
Task task = requireTask(taskId, actor.companyId());
if (task.status() != TaskStatus.APPROVED
@@ -303,6 +312,7 @@ public TaskActionResult complete(
ActorContext actor,
RequestMetadata metadata
) {
+ bindTenant(actor);
actorAuthorizer.requireHrWrite(actor);
Task task = requireTask(taskId, actor.companyId());
boolean approved = hasValidApproval(
@@ -341,6 +351,7 @@ public boolean hasValidApproval(
long contentRevision,
String criticalFingerprint
) {
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(companyId);
return approvalRepository.findLatestApprovedByTaskIdAndCompanyId(taskId, companyId)
.filter(approval -> approval.isValidFor(contentRevision, criticalFingerprint))
.isPresent();
@@ -355,6 +366,7 @@ public void invalidateForCriticalChange(
Instant occurredAt,
RequestMetadata metadata
) {
+ bindTenant(actor);
actorAuthorizer.requireHrWrite(actor);
Task task = requireTask(taskId, actor.companyId());
List active = invalidateActiveApprovals(
@@ -384,6 +396,7 @@ public Task replaceReviewAfterCriticalChange(
Instant occurredAt,
RequestMetadata metadata
) {
+ bindTenant(actor);
actorAuthorizer.requireHrWrite(actor);
Task task = requireTask(taskId, actor.companyId());
List invalidated = invalidateActiveApprovals(
@@ -557,6 +570,10 @@ private UserRole effectiveRole(ActorContext actor) {
.orElseThrow();
}
+ private void bindTenant(ActorContext actor) {
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(actor.companyId());
+ }
+
private int rolePriority(UserRole role) {
return switch (role) {
case ADMIN -> 0;
diff --git a/src/main/java/com/fowoco/server/audit/application/AuditQueryService.java b/src/main/java/com/fowoco/server/audit/application/AuditQueryService.java
index 741b617..b180632 100644
--- a/src/main/java/com/fowoco/server/audit/application/AuditQueryService.java
+++ b/src/main/java/com/fowoco/server/audit/application/AuditQueryService.java
@@ -11,6 +11,7 @@
import com.fowoco.server.auth.domain.UserRole;
import com.fowoco.server.common.error.ApiException;
import com.fowoco.server.common.error.ErrorCode;
+import com.fowoco.server.common.security.TenantDatabaseContext;
import com.fowoco.server.task.application.error.TaskErrorCode;
import com.fowoco.server.task.application.port.TaskRepository;
import java.time.Instant;
@@ -25,17 +26,20 @@ public class AuditQueryService {
private static final int MAX_PAGE_SIZE = 100;
private final ActorAuthorizer actorAuthorizer;
+ private final TenantDatabaseContext tenantDatabaseContext;
private final TaskRepository taskRepository;
private final AuditEventRepository auditRepository;
private final AuditCursorCodec cursorCodec;
public AuditQueryService(
ActorAuthorizer actorAuthorizer,
+ TenantDatabaseContext tenantDatabaseContext,
TaskRepository taskRepository,
AuditEventRepository auditRepository,
AuditCursorCodec cursorCodec
) {
this.actorAuthorizer = actorAuthorizer;
+ this.tenantDatabaseContext = tenantDatabaseContext;
this.taskRepository = taskRepository;
this.auditRepository = auditRepository;
this.cursorCodec = cursorCodec;
@@ -43,6 +47,7 @@ public AuditQueryService(
@Transactional(readOnly = true)
public List getTaskActivities(UUID taskId, ActorContext actor) {
+ bindTenant(actor);
actorAuthorizer.requireAnyRole(actor, UserRole.ADMIN, UserRole.HR, UserRole.VIEWER);
taskRepository.findByIdAndCompanyId(taskId, actor.companyId())
.orElseThrow(() -> new ApiException(TaskErrorCode.TASK_NOT_FOUND));
@@ -64,6 +69,7 @@ public AuditPageResult search(
int requestedLimit,
ActorContext actor
) {
+ bindTenant(actor);
actorAuthorizer.requireAnyRole(actor, UserRole.ADMIN);
if (createdFrom != null && createdTo != null && createdFrom.isAfter(createdTo)) {
throw new ApiException(ErrorCode.INVALID_REQUEST);
@@ -95,4 +101,8 @@ public AuditPageResult search(
private String normalizeTraceId(String traceId) {
return traceId == null || traceId.isBlank() ? null : traceId.trim();
}
+
+ private void bindTenant(ActorContext actor) {
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(actor.companyId());
+ }
}
diff --git a/src/main/java/com/fowoco/server/auth/application/AuthService.java b/src/main/java/com/fowoco/server/auth/application/AuthService.java
index afa373f..427090b 100644
--- a/src/main/java/com/fowoco/server/auth/application/AuthService.java
+++ b/src/main/java/com/fowoco/server/auth/application/AuthService.java
@@ -4,6 +4,7 @@
import com.fowoco.server.auth.application.error.InvalidRefreshTokenException;
import com.fowoco.server.auth.application.port.AccessTokenIssuer;
import com.fowoco.server.auth.application.port.AuthAuditPort;
+import com.fowoco.server.auth.application.port.AuthTenantBootstrap;
import com.fowoco.server.auth.application.port.PasswordVerifier;
import com.fowoco.server.auth.application.port.RefreshTokenGenerator;
import com.fowoco.server.auth.application.port.RefreshTokenHashPort;
@@ -13,11 +14,13 @@
import com.fowoco.server.auth.domain.UserAccount;
import com.fowoco.server.common.error.ApiException;
import com.fowoco.server.common.id.UuidGenerator;
+import com.fowoco.server.common.security.TenantDatabaseContext;
import com.fowoco.server.company.application.CompanyAuthenticationReader;
import com.fowoco.server.company.application.CompanyAuthenticationSnapshot;
import java.time.Clock;
import java.time.Instant;
import java.util.Optional;
+import java.util.UUID;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@@ -25,6 +28,8 @@
public class AuthService {
private final UserAccountRepository userAccountRepository;
+ private final AuthTenantBootstrap authTenantBootstrap;
+ private final TenantDatabaseContext tenantDatabaseContext;
private final CompanyAuthenticationReader companyAuthenticationReader;
private final PasswordVerifier passwordVerifier;
private final AccessTokenIssuer accessTokenIssuer;
@@ -39,6 +44,8 @@ public class AuthService {
public AuthService(
UserAccountRepository userAccountRepository,
+ AuthTenantBootstrap authTenantBootstrap,
+ TenantDatabaseContext tenantDatabaseContext,
CompanyAuthenticationReader companyAuthenticationReader,
PasswordVerifier passwordVerifier,
AccessTokenIssuer accessTokenIssuer,
@@ -52,6 +59,8 @@ public AuthService(
Clock clock
) {
this.userAccountRepository = userAccountRepository;
+ this.authTenantBootstrap = authTenantBootstrap;
+ this.tenantDatabaseContext = tenantDatabaseContext;
this.companyAuthenticationReader = companyAuthenticationReader;
this.passwordVerifier = passwordVerifier;
this.accessTokenIssuer = accessTokenIssuer;
@@ -68,7 +77,18 @@ public AuthService(
@Transactional
public LoginResult login(LoginCommand command) {
String normalizedEmail = UserAccount.normalizeEmail(command.email());
- Optional userAccountCandidate = userAccountRepository.findByNormalizedEmail(normalizedEmail);
+ Optional companyIdCandidate =
+ authTenantBootstrap.findCompanyIdByNormalizedEmail(normalizedEmail);
+ if (companyIdCandidate.isEmpty()) {
+ passwordVerifier.performDummyCheck(command.password());
+ throw invalidCredentialsWithAudit();
+ }
+
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(
+ companyIdCandidate.orElseThrow()
+ );
+ Optional userAccountCandidate =
+ userAccountRepository.findByNormalizedEmail(normalizedEmail);
if (userAccountCandidate.isEmpty()) {
passwordVerifier.performDummyCheck(command.password());
diff --git a/src/main/java/com/fowoco/server/auth/application/RefreshTokenLogoutTransaction.java b/src/main/java/com/fowoco/server/auth/application/RefreshTokenLogoutTransaction.java
index cc50ef3..7663a14 100644
--- a/src/main/java/com/fowoco/server/auth/application/RefreshTokenLogoutTransaction.java
+++ b/src/main/java/com/fowoco/server/auth/application/RefreshTokenLogoutTransaction.java
@@ -1,11 +1,14 @@
package com.fowoco.server.auth.application;
import com.fowoco.server.auth.application.port.AuthAuditPort;
+import com.fowoco.server.auth.application.port.AuthTenantBootstrap;
import com.fowoco.server.auth.application.port.RefreshTokenRepository;
import com.fowoco.server.auth.domain.RefreshToken;
+import com.fowoco.server.common.security.TenantDatabaseContext;
import java.time.Clock;
import java.time.Instant;
import java.util.Optional;
+import java.util.UUID;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@@ -13,21 +16,40 @@
public class RefreshTokenLogoutTransaction {
private final RefreshTokenRepository refreshTokenRepository;
+ private final AuthTenantBootstrap authTenantBootstrap;
+ private final TenantDatabaseContext tenantDatabaseContext;
private final AuthAuditPort authAuditPort;
private final Clock clock;
public RefreshTokenLogoutTransaction(
RefreshTokenRepository refreshTokenRepository,
+ AuthTenantBootstrap authTenantBootstrap,
+ TenantDatabaseContext tenantDatabaseContext,
AuthAuditPort authAuditPort,
Clock clock
) {
this.refreshTokenRepository = refreshTokenRepository;
+ this.authTenantBootstrap = authTenantBootstrap;
+ this.tenantDatabaseContext = tenantDatabaseContext;
this.authAuditPort = authAuditPort;
this.clock = clock;
}
@Transactional
public void revokeIfKnown(String tokenHash) {
+ Optional companyIdCandidate =
+ authTenantBootstrap.findCompanyIdByRefreshTokenHash(tokenHash);
+ if (companyIdCandidate.isEmpty()) {
+ authAuditPort.record(AuthAuditEvent.anonymous(
+ AuthAuditEvent.Action.LOGOUT_COMPLETED,
+ clock.instant()
+ ));
+ return;
+ }
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(
+ companyIdCandidate.orElseThrow()
+ );
+
Optional refreshTokenCandidate =
refreshTokenRepository.findByTokenHashWithFamilyLock(tokenHash);
Instant now = clock.instant();
diff --git a/src/main/java/com/fowoco/server/auth/application/RefreshTokenRotationTransaction.java b/src/main/java/com/fowoco/server/auth/application/RefreshTokenRotationTransaction.java
index 228f283..140556a 100644
--- a/src/main/java/com/fowoco/server/auth/application/RefreshTokenRotationTransaction.java
+++ b/src/main/java/com/fowoco/server/auth/application/RefreshTokenRotationTransaction.java
@@ -2,16 +2,19 @@
import com.fowoco.server.auth.application.port.AccessTokenIssuer;
import com.fowoco.server.auth.application.port.AuthAuditPort;
+import com.fowoco.server.auth.application.port.AuthTenantBootstrap;
import com.fowoco.server.auth.application.port.RefreshTokenGenerator;
import com.fowoco.server.auth.application.port.RefreshTokenRepository;
import com.fowoco.server.auth.application.port.UserAccountRepository;
import com.fowoco.server.auth.domain.RefreshToken;
import com.fowoco.server.auth.domain.UserAccount;
import com.fowoco.server.common.id.UuidGenerator;
+import com.fowoco.server.common.security.TenantDatabaseContext;
import com.fowoco.server.company.application.CompanyAuthenticationReader;
import java.time.Clock;
import java.time.Instant;
import java.util.Optional;
+import java.util.UUID;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@@ -19,6 +22,8 @@
public class RefreshTokenRotationTransaction {
private final RefreshTokenRepository refreshTokenRepository;
+ private final AuthTenantBootstrap authTenantBootstrap;
+ private final TenantDatabaseContext tenantDatabaseContext;
private final UserAccountRepository userAccountRepository;
private final CompanyAuthenticationReader companyAuthenticationReader;
private final AccessTokenIssuer accessTokenIssuer;
@@ -29,6 +34,8 @@ public class RefreshTokenRotationTransaction {
public RefreshTokenRotationTransaction(
RefreshTokenRepository refreshTokenRepository,
+ AuthTenantBootstrap authTenantBootstrap,
+ TenantDatabaseContext tenantDatabaseContext,
UserAccountRepository userAccountRepository,
CompanyAuthenticationReader companyAuthenticationReader,
AccessTokenIssuer accessTokenIssuer,
@@ -38,6 +45,8 @@ public RefreshTokenRotationTransaction(
Clock clock
) {
this.refreshTokenRepository = refreshTokenRepository;
+ this.authTenantBootstrap = authTenantBootstrap;
+ this.tenantDatabaseContext = tenantDatabaseContext;
this.userAccountRepository = userAccountRepository;
this.companyAuthenticationReader = companyAuthenticationReader;
this.accessTokenIssuer = accessTokenIssuer;
@@ -49,6 +58,19 @@ public RefreshTokenRotationTransaction(
@Transactional
public RefreshOutcome rotate(String tokenHash) {
+ Optional companyIdCandidate =
+ authTenantBootstrap.findCompanyIdByRefreshTokenHash(tokenHash);
+ if (companyIdCandidate.isEmpty()) {
+ authAuditPort.record(AuthAuditEvent.anonymous(
+ AuthAuditEvent.Action.REFRESH_REJECTED,
+ clock.instant()
+ ));
+ return RefreshOutcome.rejected(RefreshOutcome.Status.INVALID);
+ }
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(
+ companyIdCandidate.orElseThrow()
+ );
+
Optional presentedTokenCandidate =
refreshTokenRepository.findByTokenHashWithFamilyLock(tokenHash);
Instant now = clock.instant();
diff --git a/src/main/java/com/fowoco/server/auth/application/SignupService.java b/src/main/java/com/fowoco/server/auth/application/SignupService.java
index 5e459fc..281637d 100644
--- a/src/main/java/com/fowoco/server/auth/application/SignupService.java
+++ b/src/main/java/com/fowoco/server/auth/application/SignupService.java
@@ -2,12 +2,14 @@
import com.fowoco.server.auth.application.error.AuthErrorCode;
import com.fowoco.server.auth.application.port.AuthAuditPort;
+import com.fowoco.server.auth.application.port.AuthTenantBootstrap;
import com.fowoco.server.auth.application.port.PasswordHasher;
import com.fowoco.server.auth.application.port.UserAccountRepository;
import com.fowoco.server.auth.domain.UserAccount;
import com.fowoco.server.auth.domain.UserRole;
import com.fowoco.server.common.error.ApiException;
import com.fowoco.server.common.id.UuidGenerator;
+import com.fowoco.server.common.security.TenantDatabaseContext;
import com.fowoco.server.company.application.port.CompanyRepository;
import com.fowoco.server.company.domain.Company;
import java.time.Clock;
@@ -21,6 +23,8 @@ public class SignupService {
private final CompanyRepository companyRepository;
private final UserAccountRepository userAccountRepository;
+ private final AuthTenantBootstrap authTenantBootstrap;
+ private final TenantDatabaseContext tenantDatabaseContext;
private final PasswordHasher passwordHasher;
private final AuthAuditPort authAuditPort;
private final UuidGenerator uuidGenerator;
@@ -29,6 +33,8 @@ public class SignupService {
public SignupService(
CompanyRepository companyRepository,
UserAccountRepository userAccountRepository,
+ AuthTenantBootstrap authTenantBootstrap,
+ TenantDatabaseContext tenantDatabaseContext,
PasswordHasher passwordHasher,
AuthAuditPort authAuditPort,
UuidGenerator uuidGenerator,
@@ -36,6 +42,8 @@ public SignupService(
) {
this.companyRepository = companyRepository;
this.userAccountRepository = userAccountRepository;
+ this.authTenantBootstrap = authTenantBootstrap;
+ this.tenantDatabaseContext = tenantDatabaseContext;
this.passwordHasher = passwordHasher;
this.authAuditPort = authAuditPort;
this.uuidGenerator = uuidGenerator;
@@ -45,7 +53,7 @@ public SignupService(
@Transactional
public SignupResult signup(SignupCommand command) {
String normalizedEmail = UserAccount.normalizeEmail(command.email());
- if (userAccountRepository.existsByNormalizedEmail(normalizedEmail)) {
+ if (authTenantBootstrap.findCompanyIdByNormalizedEmail(normalizedEmail).isPresent()) {
throw duplicateEmail();
}
@@ -55,6 +63,7 @@ public SignupResult signup(SignupCommand command) {
command.companyName(),
now
);
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(company.companyId());
UserAccount initialAdmin = UserAccount.create(
uuidGenerator.generate(),
company.companyId(),
diff --git a/src/main/java/com/fowoco/server/auth/application/port/AuthTenantBootstrap.java b/src/main/java/com/fowoco/server/auth/application/port/AuthTenantBootstrap.java
new file mode 100644
index 0000000..03774bf
--- /dev/null
+++ b/src/main/java/com/fowoco/server/auth/application/port/AuthTenantBootstrap.java
@@ -0,0 +1,14 @@
+package com.fowoco.server.auth.application.port;
+
+import java.util.Optional;
+import java.util.UUID;
+
+/**
+ * Resolves only the tenant identifier needed to enter an authentication transaction.
+ */
+public interface AuthTenantBootstrap {
+
+ Optional findCompanyIdByNormalizedEmail(String normalizedEmail);
+
+ Optional findCompanyIdByRefreshTokenHash(String tokenHash);
+}
diff --git a/src/main/java/com/fowoco/server/auth/infrastructure/persistence/JpaAuthTenantBootstrap.java b/src/main/java/com/fowoco/server/auth/infrastructure/persistence/JpaAuthTenantBootstrap.java
new file mode 100644
index 0000000..7d113bd
--- /dev/null
+++ b/src/main/java/com/fowoco/server/auth/infrastructure/persistence/JpaAuthTenantBootstrap.java
@@ -0,0 +1,61 @@
+package com.fowoco.server.auth.infrastructure.persistence;
+
+import com.fowoco.server.auth.application.port.AuthTenantBootstrap;
+import jakarta.persistence.EntityManager;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.UUID;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.stereotype.Repository;
+
+/**
+ * H2/local bootstrap adapter used where PostgreSQL SECURITY DEFINER functions are unavailable.
+ */
+@Repository
+@ConditionalOnProperty(
+ name = "app.database.tenant-context-mode",
+ havingValue = "transaction-only",
+ matchIfMissing = true
+)
+public class JpaAuthTenantBootstrap implements AuthTenantBootstrap {
+
+ private final EntityManager entityManager;
+
+ public JpaAuthTenantBootstrap(EntityManager entityManager) {
+ this.entityManager = entityManager;
+ }
+
+ @Override
+ public Optional findCompanyIdByNormalizedEmail(String normalizedEmail) {
+ Objects.requireNonNull(normalizedEmail, "normalizedEmail must not be null");
+ return entityManager.createQuery(
+ """
+ select account.companyId
+ from UserAccountJpaEntity account
+ where account.normalizedEmail = :normalizedEmail
+ """,
+ UUID.class
+ )
+ .setParameter("normalizedEmail", normalizedEmail)
+ .setMaxResults(1)
+ .getResultStream()
+ .findFirst();
+ }
+
+ @Override
+ public Optional findCompanyIdByRefreshTokenHash(String tokenHash) {
+ Objects.requireNonNull(tokenHash, "tokenHash must not be null");
+ return entityManager.createQuery(
+ """
+ select token.companyId
+ from RefreshTokenJpaEntity token
+ where token.tokenHash = :tokenHash
+ """,
+ UUID.class
+ )
+ .setParameter("tokenHash", tokenHash)
+ .setMaxResults(1)
+ .getResultStream()
+ .findFirst();
+ }
+}
diff --git a/src/main/java/com/fowoco/server/auth/infrastructure/persistence/PostgreSqlAuthTenantBootstrap.java b/src/main/java/com/fowoco/server/auth/infrastructure/persistence/PostgreSqlAuthTenantBootstrap.java
new file mode 100644
index 0000000..4c74854
--- /dev/null
+++ b/src/main/java/com/fowoco/server/auth/infrastructure/persistence/PostgreSqlAuthTenantBootstrap.java
@@ -0,0 +1,58 @@
+package com.fowoco.server.auth.infrastructure.persistence;
+
+import com.fowoco.server.auth.application.port.AuthTenantBootstrap;
+import jakarta.persistence.EntityManager;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.UUID;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.stereotype.Repository;
+
+/**
+ * PostgreSQL bootstrap adapter restricted to company-id-only SECURITY DEFINER functions.
+ */
+@Repository
+@ConditionalOnProperty(
+ name = "app.database.tenant-context-mode",
+ havingValue = "postgresql"
+)
+public class PostgreSqlAuthTenantBootstrap implements AuthTenantBootstrap {
+
+ private static final String EMAIL_BOOTSTRAP_SQL = """
+ SELECT public.bootstrap_company_id_by_normalized_email(?1)
+ """;
+ private static final String REFRESH_TOKEN_BOOTSTRAP_SQL = """
+ SELECT public.bootstrap_company_id_by_refresh_token_hash(?1)
+ """;
+
+ private final EntityManager entityManager;
+
+ public PostgreSqlAuthTenantBootstrap(EntityManager entityManager) {
+ this.entityManager = entityManager;
+ }
+
+ @Override
+ public Optional findCompanyIdByNormalizedEmail(String normalizedEmail) {
+ Objects.requireNonNull(normalizedEmail, "normalizedEmail must not be null");
+ return queryCompanyId(EMAIL_BOOTSTRAP_SQL, normalizedEmail);
+ }
+
+ @Override
+ public Optional findCompanyIdByRefreshTokenHash(String tokenHash) {
+ Objects.requireNonNull(tokenHash, "tokenHash must not be null");
+ return queryCompanyId(REFRESH_TOKEN_BOOTSTRAP_SQL, tokenHash);
+ }
+
+ private Optional queryCompanyId(String sql, String lookupValue) {
+ Object result = entityManager.createNativeQuery(sql)
+ .setParameter(1, lookupValue)
+ .getSingleResult();
+ if (result == null) {
+ return Optional.empty();
+ }
+ if (result instanceof UUID companyId) {
+ return Optional.of(companyId);
+ }
+ return Optional.of(UUID.fromString(result.toString()));
+ }
+}
diff --git a/src/main/java/com/fowoco/server/common/security/PostgreSqlTenantDatabaseContext.java b/src/main/java/com/fowoco/server/common/security/PostgreSqlTenantDatabaseContext.java
index 6821539..8fabb19 100644
--- a/src/main/java/com/fowoco/server/common/security/PostgreSqlTenantDatabaseContext.java
+++ b/src/main/java/com/fowoco/server/common/security/PostgreSqlTenantDatabaseContext.java
@@ -3,6 +3,7 @@
import jakarta.persistence.EntityManager;
import java.util.Objects;
import java.util.UUID;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Component;
import org.springframework.transaction.support.TransactionSynchronizationManager;
@@ -10,6 +11,10 @@
* PostgreSQL tenant context backed by a transaction-local custom setting.
*/
@Component
+@ConditionalOnProperty(
+ name = "app.database.tenant-context-mode",
+ havingValue = "postgresql"
+)
public final class PostgreSqlTenantDatabaseContext implements TenantDatabaseContext {
private static final String READ_COMPANY_ID_SQL = """
diff --git a/src/main/java/com/fowoco/server/common/security/TransactionOnlyTenantDatabaseContext.java b/src/main/java/com/fowoco/server/common/security/TransactionOnlyTenantDatabaseContext.java
new file mode 100644
index 0000000..6e70865
--- /dev/null
+++ b/src/main/java/com/fowoco/server/common/security/TransactionOnlyTenantDatabaseContext.java
@@ -0,0 +1,29 @@
+package com.fowoco.server.common.security;
+
+import java.util.Objects;
+import java.util.UUID;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.stereotype.Component;
+import org.springframework.transaction.support.TransactionSynchronizationManager;
+
+/**
+ * Validates tenant transaction boundaries on databases that do not support PostgreSQL settings.
+ */
+@Component
+@ConditionalOnProperty(
+ name = "app.database.tenant-context-mode",
+ havingValue = "transaction-only",
+ matchIfMissing = true
+)
+public final class TransactionOnlyTenantDatabaseContext implements TenantDatabaseContext {
+
+ @Override
+ public void setCompanyIdForCurrentTransaction(UUID companyId) {
+ Objects.requireNonNull(companyId, "companyId must not be null");
+ if (!TransactionSynchronizationManager.isActualTransactionActive()) {
+ throw new IllegalStateException(
+ "Tenant database context requires an active transaction."
+ );
+ }
+ }
+}
diff --git a/src/main/java/com/fowoco/server/document/api/DocumentController.java b/src/main/java/com/fowoco/server/document/api/DocumentController.java
index fd308ee..0bd8be4 100644
--- a/src/main/java/com/fowoco/server/document/api/DocumentController.java
+++ b/src/main/java/com/fowoco/server/document/api/DocumentController.java
@@ -1,5 +1,6 @@
package com.fowoco.server.document.api;
+import com.fowoco.server.auth.application.ActorContext;
import com.fowoco.server.auth.application.port.ActorContextProvider;
import com.fowoco.server.document.application.DocumentPageResult;
import com.fowoco.server.document.application.DocumentService;
@@ -75,7 +76,7 @@ public DocumentPageResponse list(
@Parameter(description = "페이지당 항목 수 (1~100)")
@RequestParam(defaultValue = "20") @Min(1) @Max(100) int size
) {
- UUID companyId = actorContextProvider.requireCurrentActor().companyId();
+ ActorContext actor = actorContextProvider.requireCurrentActor();
WorkerDocumentSearchQuery query = new WorkerDocumentSearchQuery(
workerId,
documentType,
@@ -84,7 +85,7 @@ public DocumentPageResponse list(
page,
size
);
- DocumentPageResult result = documentService.findPage(companyId, query);
+ DocumentPageResult result = documentService.findPage(actor, query);
List items = result.items().stream()
.map(document -> DocumentItemResponse.from(
document,
diff --git a/src/main/java/com/fowoco/server/document/api/DocumentReadinessController.java b/src/main/java/com/fowoco/server/document/api/DocumentReadinessController.java
index 718b125..dd8630e 100644
--- a/src/main/java/com/fowoco/server/document/api/DocumentReadinessController.java
+++ b/src/main/java/com/fowoco/server/document/api/DocumentReadinessController.java
@@ -1,5 +1,6 @@
package com.fowoco.server.document.api;
+import com.fowoco.server.auth.application.ActorContext;
import com.fowoco.server.auth.application.port.ActorContextProvider;
import com.fowoco.server.document.application.DocumentReadinessResult;
import com.fowoco.server.document.application.DocumentReadinessService;
@@ -60,8 +61,8 @@ public DocumentReadinessController(
public DocumentReadinessResponse get(
@Parameter(description = "업무 ID") @PathVariable UUID taskId
) {
- UUID companyId = actorContextProvider.requireCurrentActor().companyId();
- DocumentReadinessResult result = documentReadinessService.calculate(taskId, companyId);
+ ActorContext actor = actorContextProvider.requireCurrentActor();
+ DocumentReadinessResult result = documentReadinessService.calculate(taskId, actor);
return new DocumentReadinessResponse(
result.required(),
result.available(),
diff --git a/src/main/java/com/fowoco/server/document/application/DocumentReadinessService.java b/src/main/java/com/fowoco/server/document/application/DocumentReadinessService.java
index 84c7830..88bcb4d 100644
--- a/src/main/java/com/fowoco/server/document/application/DocumentReadinessService.java
+++ b/src/main/java/com/fowoco/server/document/application/DocumentReadinessService.java
@@ -1,6 +1,8 @@
package com.fowoco.server.document.application;
+import com.fowoco.server.auth.application.ActorContext;
import com.fowoco.server.common.error.ApiException;
+import com.fowoco.server.common.security.TenantDatabaseContext;
import com.fowoco.server.document.domain.ChecklistItemDocumentMapper;
import com.fowoco.server.task.application.error.TaskErrorCode;
import com.fowoco.server.task.application.port.TaskChecklistRepository;
@@ -26,22 +28,27 @@ public class DocumentReadinessService {
private final TaskRepository taskRepository;
private final TaskChecklistRepository taskChecklistRepository;
private final WorkerDocumentRepository workerDocumentRepository;
+ private final TenantDatabaseContext tenantDatabaseContext;
private final Clock clock;
public DocumentReadinessService(
TaskRepository taskRepository,
TaskChecklistRepository taskChecklistRepository,
WorkerDocumentRepository workerDocumentRepository,
+ TenantDatabaseContext tenantDatabaseContext,
Clock clock
) {
this.taskRepository = taskRepository;
this.taskChecklistRepository = taskChecklistRepository;
this.workerDocumentRepository = workerDocumentRepository;
+ this.tenantDatabaseContext = tenantDatabaseContext;
this.clock = clock;
}
@Transactional(readOnly = true)
- public DocumentReadinessResult calculate(UUID taskId, UUID companyId) {
+ public DocumentReadinessResult calculate(UUID taskId, ActorContext actor) {
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(actor.companyId());
+ UUID companyId = actor.companyId();
Task task = taskRepository.findByIdAndCompanyId(taskId, companyId)
.orElseThrow(() -> new ApiException(TaskErrorCode.TASK_NOT_FOUND));
diff --git a/src/main/java/com/fowoco/server/document/application/DocumentRequestDraftService.java b/src/main/java/com/fowoco/server/document/application/DocumentRequestDraftService.java
index dcc6fc5..4fa3804 100644
--- a/src/main/java/com/fowoco/server/document/application/DocumentRequestDraftService.java
+++ b/src/main/java/com/fowoco/server/document/application/DocumentRequestDraftService.java
@@ -9,6 +9,7 @@
import com.fowoco.server.auth.domain.UserRole;
import com.fowoco.server.common.error.ApiException;
import com.fowoco.server.common.id.UuidGenerator;
+import com.fowoco.server.common.security.TenantDatabaseContext;
import com.fowoco.server.common.web.RequestMetadata;
import com.fowoco.server.document.application.error.DocumentErrorCode;
import com.fowoco.server.document.application.port.DocumentRequestDraftRepository;
@@ -30,6 +31,7 @@ public class DocumentRequestDraftService {
private final TaskRepository taskRepository;
private final DocumentRequestDraftRepository documentRequestDraftRepository;
private final AuditEventRepository auditRepository;
+ private final TenantDatabaseContext tenantDatabaseContext;
private final UuidGenerator uuidGenerator;
private final Clock clock;
@@ -37,12 +39,14 @@ public DocumentRequestDraftService(
TaskRepository taskRepository,
DocumentRequestDraftRepository documentRequestDraftRepository,
AuditEventRepository auditRepository,
+ TenantDatabaseContext tenantDatabaseContext,
UuidGenerator uuidGenerator,
Clock clock
) {
this.taskRepository = taskRepository;
this.documentRequestDraftRepository = documentRequestDraftRepository;
this.auditRepository = auditRepository;
+ this.tenantDatabaseContext = tenantDatabaseContext;
this.uuidGenerator = uuidGenerator;
this.clock = clock;
}
@@ -53,11 +57,12 @@ public DocumentRequestDraft upsert(
ActorContext actor,
RequestMetadata metadata
) {
- taskRepository.findByIdAndCompanyId(command.taskId(), command.companyId())
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(actor.companyId());
+ taskRepository.findByIdAndCompanyId(command.taskId(), actor.companyId())
.orElseThrow(() -> new ApiException(TaskErrorCode.TASK_NOT_FOUND));
Optional existing = documentRequestDraftRepository
- .findByTaskIdAndCompanyId(command.taskId(), command.companyId());
+ .findByTaskIdAndCompanyId(command.taskId(), actor.companyId());
Instant now = clock.instant();
DocumentRequestDraft saved;
@@ -66,7 +71,7 @@ public DocumentRequestDraft upsert(
DocumentRequestDraft draft = DocumentRequestDraft.create(
uuidGenerator.generate(),
command.taskId(),
- command.companyId(),
+ actor.companyId(),
command.language(),
command.documentTypes(),
command.message(),
diff --git a/src/main/java/com/fowoco/server/document/application/DocumentService.java b/src/main/java/com/fowoco/server/document/application/DocumentService.java
index 9d592c9..93c13b6 100644
--- a/src/main/java/com/fowoco/server/document/application/DocumentService.java
+++ b/src/main/java/com/fowoco/server/document/application/DocumentService.java
@@ -1,5 +1,7 @@
package com.fowoco.server.document.application;
+import com.fowoco.server.auth.application.ActorContext;
+import com.fowoco.server.common.security.TenantDatabaseContext;
import com.fowoco.server.worker.application.WorkerDocumentSearchQuery;
import com.fowoco.server.worker.application.port.WorkerDocumentRepository;
import com.fowoco.server.worker.application.port.WorkerRepository;
@@ -19,17 +21,22 @@ public class DocumentService {
private final WorkerDocumentRepository workerDocumentRepository;
private final WorkerRepository workerRepository;
+ private final TenantDatabaseContext tenantDatabaseContext;
public DocumentService(
WorkerDocumentRepository workerDocumentRepository,
- WorkerRepository workerRepository
+ WorkerRepository workerRepository,
+ TenantDatabaseContext tenantDatabaseContext
) {
this.workerDocumentRepository = workerDocumentRepository;
this.workerRepository = workerRepository;
+ this.tenantDatabaseContext = tenantDatabaseContext;
}
@Transactional(readOnly = true)
- public DocumentPageResult findPage(UUID companyId, WorkerDocumentSearchQuery query) {
+ public DocumentPageResult findPage(ActorContext actor, WorkerDocumentSearchQuery query) {
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(actor.companyId());
+ UUID companyId = actor.companyId();
List items = workerDocumentRepository.findPage(companyId, query);
long totalElements = workerDocumentRepository.countPage(companyId, query);
diff --git a/src/main/java/com/fowoco/server/file/application/FileService.java b/src/main/java/com/fowoco/server/file/application/FileService.java
index 61fdfd4..1a9b0b1 100644
--- a/src/main/java/com/fowoco/server/file/application/FileService.java
+++ b/src/main/java/com/fowoco/server/file/application/FileService.java
@@ -9,6 +9,7 @@
import com.fowoco.server.auth.domain.UserRole;
import com.fowoco.server.common.error.ApiException;
import com.fowoco.server.common.id.UuidGenerator;
+import com.fowoco.server.common.security.TenantDatabaseContext;
import com.fowoco.server.common.web.RequestMetadata;
import com.fowoco.server.file.application.error.FileErrorCode;
import com.fowoco.server.file.application.port.FileStorage;
@@ -48,6 +49,7 @@ public class FileService {
private final TaskRepository taskRepository;
private final WorkerRepository workerRepository;
private final AuditEventRepository auditRepository;
+ private final TenantDatabaseContext tenantDatabaseContext;
private final UuidGenerator uuidGenerator;
private final Clock clock;
@@ -57,6 +59,7 @@ public FileService(
TaskRepository taskRepository,
WorkerRepository workerRepository,
AuditEventRepository auditRepository,
+ TenantDatabaseContext tenantDatabaseContext,
UuidGenerator uuidGenerator,
Clock clock
) {
@@ -65,12 +68,15 @@ public FileService(
this.taskRepository = taskRepository;
this.workerRepository = workerRepository;
this.auditRepository = auditRepository;
+ this.tenantDatabaseContext = tenantDatabaseContext;
this.uuidGenerator = uuidGenerator;
this.clock = clock;
}
@Transactional
public StoredFile upload(FileCreateCommand command, ActorContext actor, RequestMetadata metadata) {
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(actor.companyId());
+ UUID companyId = actor.companyId();
if (command.size() > MAX_FILE_SIZE_BYTES) {
throw new ApiException(FileErrorCode.FILE_TOO_LARGE);
}
@@ -78,11 +84,11 @@ public StoredFile upload(FileCreateCommand command, ActorContext actor, RequestM
throw new ApiException(FileErrorCode.UNSUPPORTED_FILE_TYPE);
}
if (command.taskId() != null) {
- taskRepository.findByIdAndCompanyId(command.taskId(), command.companyId())
+ taskRepository.findByIdAndCompanyId(command.taskId(), companyId)
.orElseThrow(() -> new ApiException(TaskErrorCode.TASK_NOT_FOUND));
}
if (command.workerId() != null) {
- workerRepository.findByWorkerIdAndCompanyId(command.workerId(), command.companyId())
+ workerRepository.findByWorkerIdAndCompanyId(command.workerId(), companyId)
.orElseThrow(() -> new ApiException(WorkerErrorCode.WORKER_NOT_FOUND));
}
@@ -92,7 +98,7 @@ public StoredFile upload(FileCreateCommand command, ActorContext actor, RequestM
StoredFile storedFile = StoredFile.create(
storedFileId,
- command.companyId(),
+ companyId,
command.name(),
command.mimeType(),
command.size(),
diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxClaimService.java b/src/main/java/com/fowoco/server/reliability/application/OutboxClaimService.java
index f79876a..1c953fb 100644
--- a/src/main/java/com/fowoco/server/reliability/application/OutboxClaimService.java
+++ b/src/main/java/com/fowoco/server/reliability/application/OutboxClaimService.java
@@ -1,10 +1,8 @@
package com.fowoco.server.reliability.application;
-import com.fowoco.server.reliability.application.port.EventPublicationRepository;
+import com.fowoco.server.reliability.application.port.OutboxClaimBootstrap;
+import com.fowoco.server.reliability.application.port.OutboxClaimBootstrap.ClaimResult;
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;
@@ -14,39 +12,45 @@
@Service
public class OutboxClaimService {
- private final EventPublicationRepository repository;
+ private final OutboxClaimBootstrap claimBootstrap;
private final OutboxProperties properties;
private final OutboxMetrics metrics;
- private final Clock clock;
public OutboxClaimService(
- EventPublicationRepository repository,
+ OutboxClaimBootstrap claimBootstrap,
OutboxProperties properties,
- OutboxMetrics metrics,
- Clock clock
+ OutboxMetrics metrics
) {
- this.repository = repository;
+ this.claimBootstrap = claimBootstrap;
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);
+ public List claimBatch(String owner) {
+ List results = claimBootstrap.claim(
+ owner,
+ properties.getLeaseDuration(),
+ properties.getBatchSize(),
+ properties.getMaxAttempts()
+ );
+ List claimed = new ArrayList<>(results.size());
+ for (ClaimResult result : results) {
+ if (result.reviewRequired()) {
metrics.recordReviewRequired();
} else {
- claimed.add(publication.eventId());
+ claimed.add(new ClaimedEvent(result.eventId(), result.companyId()));
}
- repository.save(publication);
}
return List.copyOf(claimed);
}
+
+ public record ClaimedEvent(UUID eventId, UUID companyId) {
+
+ public ClaimedEvent {
+ if (eventId == null || companyId == null) {
+ throw new IllegalArgumentException("claimed event identifiers must not be null");
+ }
+ }
+ }
}
diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxCompletionTransaction.java b/src/main/java/com/fowoco/server/reliability/application/OutboxCompletionTransaction.java
index f8628dd..3568033 100644
--- a/src/main/java/com/fowoco/server/reliability/application/OutboxCompletionTransaction.java
+++ b/src/main/java/com/fowoco/server/reliability/application/OutboxCompletionTransaction.java
@@ -1,8 +1,9 @@
package com.fowoco.server.reliability.application;
+import com.fowoco.server.common.security.TenantDatabaseContext;
import com.fowoco.server.reliability.application.port.EventPublicationRepository;
+import com.fowoco.server.reliability.application.port.OutboxTimeSource;
import com.fowoco.server.reliability.domain.EventPublication;
-import java.time.Clock;
import java.time.Instant;
import java.util.UUID;
import org.springframework.stereotype.Service;
@@ -13,24 +14,29 @@
public class OutboxCompletionTransaction {
private final EventPublicationRepository repository;
+ private final TenantDatabaseContext tenantDatabaseContext;
private final OutboxMetrics metrics;
- private final Clock clock;
+ private final OutboxTimeSource timeSource;
public OutboxCompletionTransaction(
EventPublicationRepository repository,
+ TenantDatabaseContext tenantDatabaseContext,
OutboxMetrics metrics,
- Clock clock
+ OutboxTimeSource timeSource
) {
this.repository = repository;
+ this.tenantDatabaseContext = tenantDatabaseContext;
this.metrics = metrics;
- this.clock = clock;
+ this.timeSource = timeSource;
}
@Transactional(propagation = Propagation.REQUIRES_NEW)
- public void complete(UUID eventId, String owner) {
- Instant now = clock.instant();
- EventPublication publication = repository.findByIdForUpdate(eventId)
+ public void complete(UUID eventId, UUID companyId, String owner) {
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(companyId);
+ EventPublication publication = repository
+ .findByIdAndCompanyIdForUpdate(eventId, companyId)
.orElseThrow(() -> new IllegalStateException("Event publication not found."));
+ Instant now = timeSource.now();
publication.complete(owner, now);
repository.save(publication);
metrics.recordCompleted();
diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxFailureTransaction.java b/src/main/java/com/fowoco/server/reliability/application/OutboxFailureTransaction.java
index 558edb4..0113a96 100644
--- a/src/main/java/com/fowoco/server/reliability/application/OutboxFailureTransaction.java
+++ b/src/main/java/com/fowoco/server/reliability/application/OutboxFailureTransaction.java
@@ -1,10 +1,11 @@
package com.fowoco.server.reliability.application;
+import com.fowoco.server.common.security.TenantDatabaseContext;
import com.fowoco.server.reliability.application.OutboxFailureClassifier.FailureClassification;
import com.fowoco.server.reliability.application.port.EventPublicationRepository;
+import com.fowoco.server.reliability.application.port.OutboxTimeSource;
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;
@@ -15,37 +16,43 @@
public class OutboxFailureTransaction {
private final EventPublicationRepository repository;
+ private final TenantDatabaseContext tenantDatabaseContext;
private final OutboxFailureClassifier classifier;
private final OutboxBackoffPolicy backoffPolicy;
private final OutboxProperties properties;
private final OutboxMetrics metrics;
- private final Clock clock;
+ private final OutboxTimeSource timeSource;
public OutboxFailureTransaction(
EventPublicationRepository repository,
+ TenantDatabaseContext tenantDatabaseContext,
OutboxFailureClassifier classifier,
OutboxBackoffPolicy backoffPolicy,
OutboxProperties properties,
OutboxMetrics metrics,
- Clock clock
+ OutboxTimeSource timeSource
) {
this.repository = repository;
+ this.tenantDatabaseContext = tenantDatabaseContext;
this.classifier = classifier;
this.backoffPolicy = backoffPolicy;
this.properties = properties;
this.metrics = metrics;
- this.clock = clock;
+ this.timeSource = timeSource;
}
@Transactional(propagation = Propagation.REQUIRES_NEW)
public FailureOutcome recordFailure(
UUID eventId,
+ UUID companyId,
String owner,
Throwable failure
) {
- Instant now = clock.instant();
- EventPublication publication = repository.findByIdForUpdate(eventId)
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(companyId);
+ EventPublication publication = repository
+ .findByIdAndCompanyIdForUpdate(eventId, companyId)
.orElseThrow(() -> new IllegalStateException("Event publication not found."));
+ Instant now = timeSource.now();
FailureClassification classification = classifier.classify(failure);
boolean exhausted = publication.attemptCount() >= properties.getMaxAttempts();
if (!classification.retryable() || exhausted) {
diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxHandlerTransaction.java b/src/main/java/com/fowoco/server/reliability/application/OutboxHandlerTransaction.java
index 5996ed5..5ba776a 100644
--- a/src/main/java/com/fowoco/server/reliability/application/OutboxHandlerTransaction.java
+++ b/src/main/java/com/fowoco/server/reliability/application/OutboxHandlerTransaction.java
@@ -1,14 +1,15 @@
package com.fowoco.server.reliability.application;
import com.fowoco.server.common.id.UuidGenerator;
+import com.fowoco.server.common.security.TenantDatabaseContext;
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.application.port.OutboxTimeSource;
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;
@@ -20,33 +21,39 @@ public class OutboxHandlerTransaction {
private final EventPublicationRepository publicationRepository;
private final EventConsumptionRepository consumptionRepository;
+ private final TenantDatabaseContext tenantDatabaseContext;
private final EventPayloadCodec payloadCodec;
private final UuidGenerator uuidGenerator;
- private final Clock clock;
+ private final OutboxTimeSource timeSource;
public OutboxHandlerTransaction(
EventPublicationRepository publicationRepository,
EventConsumptionRepository consumptionRepository,
+ TenantDatabaseContext tenantDatabaseContext,
EventPayloadCodec payloadCodec,
UuidGenerator uuidGenerator,
- Clock clock
+ OutboxTimeSource timeSource
) {
this.publicationRepository = publicationRepository;
this.consumptionRepository = consumptionRepository;
+ this.tenantDatabaseContext = tenantDatabaseContext;
this.payloadCodec = payloadCodec;
this.uuidGenerator = uuidGenerator;
- this.clock = clock;
+ this.timeSource = timeSource;
}
@Transactional(propagation = Propagation.REQUIRES_NEW)
public boolean deliver(
UUID eventId,
+ UUID companyId,
String owner,
DomainEventHandler handler
) {
- Instant now = clock.instant();
- EventPublication publication = publicationRepository.findByIdForUpdate(eventId)
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(companyId);
+ EventPublication publication = publicationRepository
+ .findByIdAndCompanyIdForUpdate(eventId, companyId)
.orElseThrow(() -> new IllegalStateException("Event publication not found."));
+ Instant now = timeSource.now();
publication.requireActiveLease(owner, now);
String handlerName = handler.handlerName();
if (consumptionRepository.existsByEventIdAndHandlerName(eventId, handlerName)) {
diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxMetrics.java b/src/main/java/com/fowoco/server/reliability/application/OutboxMetrics.java
index 9e1f8a8..06409a6 100644
--- a/src/main/java/com/fowoco/server/reliability/application/OutboxMetrics.java
+++ b/src/main/java/com/fowoco/server/reliability/application/OutboxMetrics.java
@@ -1,6 +1,6 @@
package com.fowoco.server.reliability.application;
-import com.fowoco.server.reliability.application.port.EventPublicationRepository;
+import com.fowoco.server.reliability.application.port.OutboxBacklogReader;
import io.micrometer.core.instrument.Counter;
import io.micrometer.core.instrument.Gauge;
import io.micrometer.core.instrument.MeterRegistry;
@@ -17,7 +17,7 @@ public class OutboxMetrics {
public OutboxMetrics(
MeterRegistry meterRegistry,
- EventPublicationRepository repository,
+ OutboxBacklogReader backlogReader,
Clock clock
) {
completed = counter(meterRegistry, "completed");
@@ -25,14 +25,14 @@ public OutboxMetrics(
reviewRequired = counter(meterRegistry, "review_required");
Gauge.builder(
"fowoco.outbox.publications.backlog",
- repository,
- EventPublicationRepository::countOutstanding
+ backlogReader,
+ OutboxBacklogReader::countOutstanding
)
.description("Outstanding durable event publications")
.register(meterRegistry);
Gauge.builder(
"fowoco.outbox.publications.oldest.delay.seconds",
- repository,
+ backlogReader,
candidate -> candidate.findOldestOutstandingOccurredAt()
.map(occurredAt -> Math.max(
0.0,
diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxProcessor.java b/src/main/java/com/fowoco/server/reliability/application/OutboxProcessor.java
index ad59e48..0eb60ee 100644
--- a/src/main/java/com/fowoco/server/reliability/application/OutboxProcessor.java
+++ b/src/main/java/com/fowoco/server/reliability/application/OutboxProcessor.java
@@ -4,7 +4,6 @@
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;
@@ -41,32 +40,46 @@ public OutboxProcessor(
}
public int processAvailable() {
- List eventIds = claimService.claimBatch(workerIdentity.value());
- eventIds.forEach(this::processOne);
- return eventIds.size();
+ List claimedEvents =
+ claimService.claimBatch(workerIdentity.value());
+ claimedEvents.forEach(this::processOne);
+ return claimedEvents.size();
}
- private void processOne(UUID eventId) {
- EventPublication publication = readService.requirePublication(eventId);
+ private void processOne(OutboxClaimService.ClaimedEvent claimedEvent) {
+ EventPublication publication = readService.requirePublication(
+ claimedEvent.eventId(),
+ claimedEvent.companyId()
+ );
try {
List handlers =
handlerRegistry.handlersFor(publication.eventType());
for (DomainEventHandler handler : handlers) {
- handlerTransaction.deliver(eventId, workerIdentity.value(), handler);
+ handlerTransaction.deliver(
+ claimedEvent.eventId(),
+ claimedEvent.companyId(),
+ workerIdentity.value(),
+ handler
+ );
}
- completionTransaction.complete(eventId, workerIdentity.value());
+ completionTransaction.complete(
+ claimedEvent.eventId(),
+ claimedEvent.companyId(),
+ workerIdentity.value()
+ );
} catch (RuntimeException failure) {
try {
OutboxFailureTransaction.FailureOutcome outcome =
failureTransaction.recordFailure(
- eventId,
+ claimedEvent.eventId(),
+ claimedEvent.companyId(),
workerIdentity.value(),
failure
);
log.warn(
"Outbox event processing failed: eventId={}, eventType={}, "
+ "attempt={}, errorCode={}, retryScheduled={}",
- eventId,
+ claimedEvent.eventId(),
publication.eventType(),
publication.attemptCount(),
outcome.errorCode(),
@@ -75,7 +88,7 @@ private void processOne(UUID eventId) {
} catch (RuntimeException recordingFailure) {
log.error(
"Outbox failure state could not be recorded: eventId={}, eventType={}",
- eventId,
+ claimedEvent.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
index 82e1590..23d3efa 100644
--- a/src/main/java/com/fowoco/server/reliability/application/OutboxReadService.java
+++ b/src/main/java/com/fowoco/server/reliability/application/OutboxReadService.java
@@ -1,5 +1,6 @@
package com.fowoco.server.reliability.application;
+import com.fowoco.server.common.security.TenantDatabaseContext;
import com.fowoco.server.reliability.application.port.EventPublicationRepository;
import com.fowoco.server.reliability.domain.EventPublication;
import java.util.UUID;
@@ -10,14 +11,20 @@
public class OutboxReadService {
private final EventPublicationRepository repository;
+ private final TenantDatabaseContext tenantDatabaseContext;
- public OutboxReadService(EventPublicationRepository repository) {
+ public OutboxReadService(
+ EventPublicationRepository repository,
+ TenantDatabaseContext tenantDatabaseContext
+ ) {
this.repository = repository;
+ this.tenantDatabaseContext = tenantDatabaseContext;
}
@Transactional(readOnly = true)
- public EventPublication requirePublication(UUID eventId) {
- return repository.findById(eventId)
+ public EventPublication requirePublication(UUID eventId, UUID companyId) {
+ tenantDatabaseContext.setCompanyIdForCurrentTransaction(companyId);
+ return repository.findByIdAndCompanyId(eventId, companyId)
.orElseThrow(() -> new IllegalStateException("Event publication not found."));
}
}
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
index 5f40baf..9f1cb1f 100644
--- a/src/main/java/com/fowoco/server/reliability/application/port/EventPublicationRepository.java
+++ b/src/main/java/com/fowoco/server/reliability/application/port/EventPublicationRepository.java
@@ -14,9 +14,9 @@ public interface EventPublicationRepository {
List lockClaimable(Instant now, int limit);
- Optional findById(UUID eventId);
+ Optional findByIdAndCompanyId(UUID eventId, UUID companyId);
- Optional findByIdForUpdate(UUID eventId);
+ Optional findByIdAndCompanyIdForUpdate(UUID eventId, UUID companyId);
long countOutstanding();
diff --git a/src/main/java/com/fowoco/server/reliability/application/port/OutboxBacklogReader.java b/src/main/java/com/fowoco/server/reliability/application/port/OutboxBacklogReader.java
new file mode 100644
index 0000000..853984a
--- /dev/null
+++ b/src/main/java/com/fowoco/server/reliability/application/port/OutboxBacklogReader.java
@@ -0,0 +1,14 @@
+package com.fowoco.server.reliability.application.port;
+
+import java.time.Instant;
+import java.util.Optional;
+
+/**
+ * Reads payload-free, cross-tenant backlog aggregates for operational metrics.
+ */
+public interface OutboxBacklogReader {
+
+ long countOutstanding();
+
+ Optional findOldestOutstandingOccurredAt();
+}
diff --git a/src/main/java/com/fowoco/server/reliability/application/port/OutboxClaimBootstrap.java b/src/main/java/com/fowoco/server/reliability/application/port/OutboxClaimBootstrap.java
new file mode 100644
index 0000000..c78ffcf
--- /dev/null
+++ b/src/main/java/com/fowoco/server/reliability/application/port/OutboxClaimBootstrap.java
@@ -0,0 +1,27 @@
+package com.fowoco.server.reliability.application.port;
+
+import java.time.Duration;
+import java.util.List;
+import java.util.UUID;
+
+/**
+ * Claims cross-tenant outbox work without exposing event payloads.
+ */
+public interface OutboxClaimBootstrap {
+
+ List claim(
+ String owner,
+ Duration leaseDuration,
+ int batchSize,
+ int maxAttempts
+ );
+
+ record ClaimResult(UUID eventId, UUID companyId, boolean reviewRequired) {
+
+ public ClaimResult {
+ if (eventId == null || companyId == null) {
+ throw new IllegalArgumentException("claimed event identifiers must not be null");
+ }
+ }
+ }
+}
diff --git a/src/main/java/com/fowoco/server/reliability/application/port/OutboxTimeSource.java b/src/main/java/com/fowoco/server/reliability/application/port/OutboxTimeSource.java
new file mode 100644
index 0000000..6076f37
--- /dev/null
+++ b/src/main/java/com/fowoco/server/reliability/application/port/OutboxTimeSource.java
@@ -0,0 +1,8 @@
+package com.fowoco.server.reliability.application.port;
+
+import java.time.Instant;
+
+public interface OutboxTimeSource {
+
+ Instant now();
+}
diff --git a/src/main/java/com/fowoco/server/reliability/config/OutboxProperties.java b/src/main/java/com/fowoco/server/reliability/config/OutboxProperties.java
index 313fd9a..9151558 100644
--- a/src/main/java/com/fowoco/server/reliability/config/OutboxProperties.java
+++ b/src/main/java/com/fowoco/server/reliability/config/OutboxProperties.java
@@ -6,6 +6,9 @@
@ConfigurationProperties(prefix = "app.reliability.outbox")
public class OutboxProperties {
+ private static final Duration MIN_LEASE_DURATION = Duration.ofMillis(1);
+ private static final Duration MAX_LEASE_DURATION = Duration.ofDays(1);
+
private boolean enabled = true;
private Duration pollInterval = Duration.ofSeconds(1);
private int batchSize = 20;
@@ -46,7 +49,16 @@ public Duration getLeaseDuration() {
}
public void setLeaseDuration(Duration leaseDuration) {
- this.leaseDuration = requirePositive(leaseDuration, "leaseDuration");
+ if (leaseDuration == null
+ || leaseDuration.compareTo(MIN_LEASE_DURATION) < 0
+ || leaseDuration.compareTo(MAX_LEASE_DURATION) > 0
+ || leaseDuration.getNano() % 1_000_000 != 0) {
+ throw new IllegalArgumentException(
+ "leaseDuration must be between 1 millisecond and 1 day "
+ + "and aligned to whole milliseconds"
+ );
+ }
+ this.leaseDuration = leaseDuration;
}
public int getMaxAttempts() {
diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/ClockOutboxTimeSource.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/ClockOutboxTimeSource.java
new file mode 100644
index 0000000..bf2e107
--- /dev/null
+++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/ClockOutboxTimeSource.java
@@ -0,0 +1,27 @@
+package com.fowoco.server.reliability.infrastructure.persistence;
+
+import com.fowoco.server.reliability.application.port.OutboxTimeSource;
+import java.time.Clock;
+import java.time.Instant;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.stereotype.Component;
+
+@Component
+@ConditionalOnProperty(
+ name = "app.database.tenant-context-mode",
+ havingValue = "transaction-only",
+ matchIfMissing = true
+)
+public class ClockOutboxTimeSource implements OutboxTimeSource {
+
+ private final Clock clock;
+
+ public ClockOutboxTimeSource(Clock clock) {
+ this.clock = clock;
+ }
+
+ @Override
+ public Instant now() {
+ return clock.instant();
+ }
+}
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
index 256053f..9688e00 100644
--- a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaEventPublicationRepository.java
+++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaEventPublicationRepository.java
@@ -39,7 +39,10 @@ public EventPublication append(EventPublication publication) {
@Override
public EventPublication save(EventPublication publication) {
- EventPublicationJpaEntity entity = repository.findById(publication.eventId())
+ EventPublicationJpaEntity entity = repository.findByEventIdAndCompanyId(
+ publication.eventId(),
+ publication.companyId()
+ )
.orElseThrow(() -> new IllegalStateException("Event publication not found."));
entity.apply(publication);
return repository.saveAndFlush(entity).toDomain();
@@ -59,14 +62,17 @@ public List lockClaimable(Instant now, int limit) {
}
@Override
- public Optional findById(UUID eventId) {
- return repository.findById(eventId)
+ public Optional findByIdAndCompanyId(UUID eventId, UUID companyId) {
+ return repository.findByEventIdAndCompanyId(eventId, companyId)
.map(EventPublicationJpaEntity::toDomain);
}
@Override
- public Optional findByIdForUpdate(UUID eventId) {
- return repository.findByIdForUpdate(eventId)
+ public Optional findByIdAndCompanyIdForUpdate(
+ UUID eventId,
+ UUID companyId
+ ) {
+ return repository.findByIdAndCompanyIdForUpdate(eventId, companyId)
.map(EventPublicationJpaEntity::toDomain);
}
diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaOutboxBacklogReader.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaOutboxBacklogReader.java
new file mode 100644
index 0000000..392dbb7
--- /dev/null
+++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaOutboxBacklogReader.java
@@ -0,0 +1,33 @@
+package com.fowoco.server.reliability.infrastructure.persistence;
+
+import com.fowoco.server.reliability.application.port.EventPublicationRepository;
+import com.fowoco.server.reliability.application.port.OutboxBacklogReader;
+import java.time.Instant;
+import java.util.Optional;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.stereotype.Repository;
+
+@Repository
+@ConditionalOnProperty(
+ name = "app.database.tenant-context-mode",
+ havingValue = "transaction-only",
+ matchIfMissing = true
+)
+public class JpaOutboxBacklogReader implements OutboxBacklogReader {
+
+ private final EventPublicationRepository repository;
+
+ public JpaOutboxBacklogReader(EventPublicationRepository repository) {
+ this.repository = repository;
+ }
+
+ @Override
+ public long countOutstanding() {
+ return repository.countOutstanding();
+ }
+
+ @Override
+ public Optional findOldestOutstandingOccurredAt() {
+ return repository.findOldestOutstandingOccurredAt();
+ }
+}
diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaOutboxClaimBootstrap.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaOutboxClaimBootstrap.java
new file mode 100644
index 0000000..3038244
--- /dev/null
+++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaOutboxClaimBootstrap.java
@@ -0,0 +1,63 @@
+package com.fowoco.server.reliability.infrastructure.persistence;
+
+import com.fowoco.server.reliability.application.port.EventPublicationRepository;
+import com.fowoco.server.reliability.application.port.OutboxClaimBootstrap;
+import com.fowoco.server.reliability.application.port.OutboxTimeSource;
+import com.fowoco.server.reliability.domain.EventPublication;
+import java.time.Duration;
+import java.time.Instant;
+import java.util.ArrayList;
+import java.util.List;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.stereotype.Repository;
+
+/**
+ * H2/local claim adapter preserving the domain claim behavior without PostgreSQL functions.
+ */
+@Repository
+@ConditionalOnProperty(
+ name = "app.database.tenant-context-mode",
+ havingValue = "transaction-only",
+ matchIfMissing = true
+)
+public class JpaOutboxClaimBootstrap implements OutboxClaimBootstrap {
+
+ private static final String ATTEMPTS_EXHAUSTED = "EVENT_ATTEMPTS_EXHAUSTED";
+
+ private final EventPublicationRepository repository;
+ private final OutboxTimeSource timeSource;
+
+ public JpaOutboxClaimBootstrap(
+ EventPublicationRepository repository,
+ OutboxTimeSource timeSource
+ ) {
+ this.repository = repository;
+ this.timeSource = timeSource;
+ }
+
+ @Override
+ public List claim(
+ String owner,
+ Duration leaseDuration,
+ int batchSize,
+ int maxAttempts
+ ) {
+ Instant now = timeSource.now();
+ List candidates = repository.lockClaimable(now, batchSize);
+ List results = new ArrayList<>(candidates.size());
+ for (EventPublication publication : candidates) {
+ publication.claim(owner, now, leaseDuration);
+ boolean reviewRequired = publication.attemptCount() > maxAttempts;
+ if (reviewRequired) {
+ publication.requireReview(owner, ATTEMPTS_EXHAUSTED, now);
+ }
+ repository.save(publication);
+ results.add(new ClaimResult(
+ publication.eventId(),
+ publication.companyId(),
+ reviewRequired
+ ));
+ }
+ return List.copyOf(results);
+ }
+}
diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/PostgreSqlOutboxBacklogReader.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/PostgreSqlOutboxBacklogReader.java
new file mode 100644
index 0000000..bf71325
--- /dev/null
+++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/PostgreSqlOutboxBacklogReader.java
@@ -0,0 +1,47 @@
+package com.fowoco.server.reliability.infrastructure.persistence;
+
+import com.fowoco.server.reliability.application.port.OutboxBacklogReader;
+import jakarta.persistence.EntityManager;
+import java.time.Instant;
+import java.time.OffsetDateTime;
+import java.util.Optional;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.stereotype.Repository;
+
+@Repository
+@ConditionalOnProperty(
+ name = "app.database.tenant-context-mode",
+ havingValue = "postgresql"
+)
+public class PostgreSqlOutboxBacklogReader implements OutboxBacklogReader {
+
+ private final EntityManager entityManager;
+
+ public PostgreSqlOutboxBacklogReader(EntityManager entityManager) {
+ this.entityManager = entityManager;
+ }
+
+ @Override
+ public long countOutstanding() {
+ return ((Number) entityManager.createNativeQuery(
+ "SELECT public.bootstrap_count_outstanding_event_publications()"
+ ).getSingleResult()).longValue();
+ }
+
+ @Override
+ public Optional findOldestOutstandingOccurredAt() {
+ Object result = entityManager.createNativeQuery(
+ "SELECT public.bootstrap_oldest_outstanding_event_occurred_at()"
+ ).getSingleResult();
+ if (result == null) {
+ return Optional.empty();
+ }
+ if (result instanceof Instant instant) {
+ return Optional.of(instant);
+ }
+ if (result instanceof OffsetDateTime offsetDateTime) {
+ return Optional.of(offsetDateTime.toInstant());
+ }
+ throw new IllegalStateException("Unexpected PostgreSQL timestamp mapping.");
+ }
+}
diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/PostgreSqlOutboxClaimBootstrap.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/PostgreSqlOutboxClaimBootstrap.java
new file mode 100644
index 0000000..a87b4b9
--- /dev/null
+++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/PostgreSqlOutboxClaimBootstrap.java
@@ -0,0 +1,61 @@
+package com.fowoco.server.reliability.infrastructure.persistence;
+
+import com.fowoco.server.reliability.application.port.OutboxClaimBootstrap;
+import jakarta.persistence.EntityManager;
+import java.time.Duration;
+import java.util.List;
+import java.util.UUID;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.stereotype.Repository;
+
+/**
+ * PostgreSQL claim adapter returning only identifiers from the restricted bootstrap function.
+ */
+@Repository
+@ConditionalOnProperty(
+ name = "app.database.tenant-context-mode",
+ havingValue = "postgresql"
+)
+public class PostgreSqlOutboxClaimBootstrap implements OutboxClaimBootstrap {
+
+ private static final String CLAIM_SQL = """
+ SELECT event_id, company_id, review_required
+ FROM public.bootstrap_claim_event_publications(?1, ?2, ?3, ?4)
+ """;
+
+ private final EntityManager entityManager;
+
+ public PostgreSqlOutboxClaimBootstrap(EntityManager entityManager) {
+ this.entityManager = entityManager;
+ }
+
+ @Override
+ public List claim(
+ String owner,
+ Duration leaseDuration,
+ int batchSize,
+ int maxAttempts
+ ) {
+ @SuppressWarnings("unchecked")
+ List