From 0c527cf57577565d3c74bf8903a1d5e2e9e08439 Mon Sep 17 00:00:00 2001 From: hywznn Date: Thu, 6 Aug 2026 16:35:02 +0900 Subject: [PATCH 1/3] =?UTF-8?q?feat(airun):=20AI=20=ED=9B=84=EB=B3=B4=20?= =?UTF-8?q?=EA=B2=B0=EC=A0=95=EA=B3=BC=20=EC=97=85=EB=AC=B4=EC=B9=B4?= =?UTF-8?q?=EB=93=9C=20=EC=83=9D=EC=84=B1=EC=9D=84=20=EA=B5=AC=ED=98=84?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../api/AiCandidateDecisionItemRequest.java | 11 + .../api/AiCandidateDecisionItemResponse.java | 14 + .../api/AiCandidateDecisionResponse.java | 25 ++ .../server/airun/api/AiRunController.java | 48 +++ .../api/DecideAiRunCandidatesRequest.java | 17 + .../AiCandidateDecisionCommand.java | 22 ++ .../AiCandidateDecisionResult.java | 29 ++ .../AiCandidateDecisionService.java | 309 +++++++++++++++ .../application/error/AiRunErrorCode.java | 6 +- .../port/AiCandidateDecisionRepository.java | 86 ++++ .../domain/AiCandidateDecisionAction.java | 6 + .../JdbcAiCandidateDecisionRepository.java | 350 +++++++++++++++++ .../persistence/JdbcAiRunRepository.java | 4 +- .../server/audit/domain/AuditAction.java | 1 + .../persistence/JdbcTaskCaseRegistrar.java | 113 +++++- .../AiCandidateTaskCreationService.java | 368 ++++++++++++++++++ .../task/application/error/TaskErrorCode.java | 4 + .../port/AiCandidateTaskCreator.java | 31 ++ .../application/port/TaskCaseRegistrar.java | 6 + ...V23__prepare_ai_candidate_decision_rls.sql | 38 ++ .../V22__create_ai_candidate_decision.sql | 94 +++++ 21 files changed, 1559 insertions(+), 23 deletions(-) create mode 100644 src/main/java/com/fowoco/server/airun/api/AiCandidateDecisionItemRequest.java create mode 100644 src/main/java/com/fowoco/server/airun/api/AiCandidateDecisionItemResponse.java create mode 100644 src/main/java/com/fowoco/server/airun/api/AiCandidateDecisionResponse.java create mode 100644 src/main/java/com/fowoco/server/airun/api/DecideAiRunCandidatesRequest.java create mode 100644 src/main/java/com/fowoco/server/airun/application/AiCandidateDecisionCommand.java create mode 100644 src/main/java/com/fowoco/server/airun/application/AiCandidateDecisionResult.java create mode 100644 src/main/java/com/fowoco/server/airun/application/AiCandidateDecisionService.java create mode 100644 src/main/java/com/fowoco/server/airun/application/port/AiCandidateDecisionRepository.java create mode 100644 src/main/java/com/fowoco/server/airun/domain/AiCandidateDecisionAction.java create mode 100644 src/main/java/com/fowoco/server/airun/infrastructure/persistence/JdbcAiCandidateDecisionRepository.java create mode 100644 src/main/java/com/fowoco/server/task/application/AiCandidateTaskCreationService.java create mode 100644 src/main/java/com/fowoco/server/task/application/port/AiCandidateTaskCreator.java create mode 100644 src/main/resources/db/migration-postgresql/V23__prepare_ai_candidate_decision_rls.sql create mode 100644 src/main/resources/db/migration/V22__create_ai_candidate_decision.sql diff --git a/src/main/java/com/fowoco/server/airun/api/AiCandidateDecisionItemRequest.java b/src/main/java/com/fowoco/server/airun/api/AiCandidateDecisionItemRequest.java new file mode 100644 index 00000000..383e50a2 --- /dev/null +++ b/src/main/java/com/fowoco/server/airun/api/AiCandidateDecisionItemRequest.java @@ -0,0 +1,11 @@ +package com.fowoco.server.airun.api; + +import com.fowoco.server.airun.domain.AiCandidateDecisionAction; +import jakarta.validation.constraints.NotNull; +import java.util.UUID; + +public record AiCandidateDecisionItemRequest( + @NotNull UUID candidateId, + @NotNull AiCandidateDecisionAction action +) { +} diff --git a/src/main/java/com/fowoco/server/airun/api/AiCandidateDecisionItemResponse.java b/src/main/java/com/fowoco/server/airun/api/AiCandidateDecisionItemResponse.java new file mode 100644 index 00000000..daf680d6 --- /dev/null +++ b/src/main/java/com/fowoco/server/airun/api/AiCandidateDecisionItemResponse.java @@ -0,0 +1,14 @@ +package com.fowoco.server.airun.api; + +import com.fowoco.server.airun.application.AiCandidateDecisionResult; +import com.fowoco.server.airun.domain.AiCandidateDecisionAction; +import java.util.UUID; + +public record AiCandidateDecisionItemResponse( + UUID candidateId, + AiCandidateDecisionAction action +) { + static AiCandidateDecisionItemResponse from(AiCandidateDecisionResult.Decision result) { + return new AiCandidateDecisionItemResponse(result.candidateId(), result.action()); + } +} diff --git a/src/main/java/com/fowoco/server/airun/api/AiCandidateDecisionResponse.java b/src/main/java/com/fowoco/server/airun/api/AiCandidateDecisionResponse.java new file mode 100644 index 00000000..8a86c7d3 --- /dev/null +++ b/src/main/java/com/fowoco/server/airun/api/AiCandidateDecisionResponse.java @@ -0,0 +1,25 @@ +package com.fowoco.server.airun.api; + +import com.fowoco.server.airun.application.AiCandidateDecisionResult; +import java.util.List; +import java.util.UUID; + +public record AiCandidateDecisionResponse( + UUID decisionBatchId, + UUID aiRunId, + UUID caseId, + List taskIds, + List decisions, + long runVersion +) { + static AiCandidateDecisionResponse from(AiCandidateDecisionResult result) { + return new AiCandidateDecisionResponse( + result.decisionBatchId(), + result.aiRunId(), + result.caseId(), + result.taskIds(), + result.decisions().stream().map(AiCandidateDecisionItemResponse::from).toList(), + result.runVersion() + ); + } +} diff --git a/src/main/java/com/fowoco/server/airun/api/AiRunController.java b/src/main/java/com/fowoco/server/airun/api/AiRunController.java index 130b61dc..59615d0b 100644 --- a/src/main/java/com/fowoco/server/airun/api/AiRunController.java +++ b/src/main/java/com/fowoco/server/airun/api/AiRunController.java @@ -1,5 +1,7 @@ package com.fowoco.server.airun.api; +import com.fowoco.server.airun.application.AiCandidateDecisionCommand; +import com.fowoco.server.airun.application.AiCandidateDecisionService; import com.fowoco.server.airun.application.AiRunService; import com.fowoco.server.auth.application.ActorContext; import com.fowoco.server.auth.application.port.ActorContextProvider; @@ -33,13 +35,16 @@ public class AiRunController { private final AiRunService aiRunService; + private final AiCandidateDecisionService candidateDecisionService; private final ActorContextProvider actorContextProvider; public AiRunController( AiRunService aiRunService, + AiCandidateDecisionService candidateDecisionService, ActorContextProvider actorContextProvider ) { this.aiRunService = aiRunService; + this.candidateDecisionService = candidateDecisionService; this.actorContextProvider = actorContextProvider; } @@ -116,6 +121,49 @@ public ResponseEntity answer( ))); } + @Operation( + operationId = "decideAiRunCandidates", + summary = "AI 업무 후보 채택·폐기", + description = "HR이 채택한 후보만 Case와 업무카드로 생성합니다. 승인과 발송은 별도입니다." + ) + @ApiResponses({ + @ApiResponse(responseCode = "200", description = "후보 결정 완료"), + @ApiResponse(responseCode = "400", ref = "#/components/responses/BadRequest"), + @ApiResponse(responseCode = "404", ref = "#/components/responses/NotFound"), + @ApiResponse(responseCode = "409", ref = "#/components/responses/Conflict"), + @ApiResponse(responseCode = "422", ref = "#/components/responses/UnprocessableEntity") + }) + @PreAuthorize("hasAnyRole('ADMIN', 'HR')") + @PostMapping( + path = "/{aiRunId}/candidate-decisions", + consumes = MediaType.APPLICATION_JSON_VALUE, + produces = MediaType.APPLICATION_JSON_VALUE + ) + public AiCandidateDecisionResponse decideCandidates( + @PathVariable UUID aiRunId, + @Parameter(description = "같은 후보 결정을 중복 생성하지 않기 위한 키", required = true) + @RequestHeader("Idempotency-Key") String idempotencyKey, + @Valid @RequestBody DecideAiRunCandidatesRequest request, + HttpServletRequest servletRequest + ) { + AiCandidateDecisionCommand command = new AiCandidateDecisionCommand( + request.expectedRunVersion(), + request.decisions().stream() + .map(decision -> new AiCandidateDecisionCommand.Decision( + decision.candidateId(), + decision.action() + )) + .toList() + ); + return AiCandidateDecisionResponse.from(candidateDecisionService.decide( + aiRunId, + idempotencyKey, + command, + actor(), + RequestMetadata.from(servletRequest) + )); + } + private ActorContext actor() { return actorContextProvider.requireCurrentActor(); } diff --git a/src/main/java/com/fowoco/server/airun/api/DecideAiRunCandidatesRequest.java b/src/main/java/com/fowoco/server/airun/api/DecideAiRunCandidatesRequest.java new file mode 100644 index 00000000..36877c97 --- /dev/null +++ b/src/main/java/com/fowoco/server/airun/api/DecideAiRunCandidatesRequest.java @@ -0,0 +1,17 @@ +package com.fowoco.server.airun.api; + +import jakarta.validation.Valid; +import jakarta.validation.constraints.Min; +import jakarta.validation.constraints.NotEmpty; +import jakarta.validation.constraints.NotNull; +import jakarta.validation.constraints.Size; +import java.util.List; + +public record DecideAiRunCandidatesRequest( + @Min(0) long expectedRunVersion, + @NotNull @NotEmpty @Size(max = 20) List<@Valid AiCandidateDecisionItemRequest> decisions +) { + public DecideAiRunCandidatesRequest { + decisions = decisions == null ? null : List.copyOf(decisions); + } +} diff --git a/src/main/java/com/fowoco/server/airun/application/AiCandidateDecisionCommand.java b/src/main/java/com/fowoco/server/airun/application/AiCandidateDecisionCommand.java new file mode 100644 index 00000000..2a100968 --- /dev/null +++ b/src/main/java/com/fowoco/server/airun/application/AiCandidateDecisionCommand.java @@ -0,0 +1,22 @@ +package com.fowoco.server.airun.application; + +import com.fowoco.server.airun.domain.AiCandidateDecisionAction; +import java.util.List; +import java.util.Objects; +import java.util.UUID; + +public record AiCandidateDecisionCommand( + long expectedRunVersion, + List decisions +) { + public AiCandidateDecisionCommand { + decisions = List.copyOf(decisions); + } + + public record Decision(UUID candidateId, AiCandidateDecisionAction action) { + public Decision { + Objects.requireNonNull(candidateId); + Objects.requireNonNull(action); + } + } +} diff --git a/src/main/java/com/fowoco/server/airun/application/AiCandidateDecisionResult.java b/src/main/java/com/fowoco/server/airun/application/AiCandidateDecisionResult.java new file mode 100644 index 00000000..259dfb4e --- /dev/null +++ b/src/main/java/com/fowoco/server/airun/application/AiCandidateDecisionResult.java @@ -0,0 +1,29 @@ +package com.fowoco.server.airun.application; + +import com.fowoco.server.airun.domain.AiCandidateDecisionAction; +import java.util.List; +import java.util.Objects; +import java.util.UUID; + +public record AiCandidateDecisionResult( + UUID decisionBatchId, + UUID aiRunId, + UUID caseId, + List taskIds, + List decisions, + long runVersion +) { + public AiCandidateDecisionResult { + Objects.requireNonNull(decisionBatchId); + Objects.requireNonNull(aiRunId); + taskIds = List.copyOf(taskIds); + decisions = List.copyOf(decisions); + } + + public record Decision(UUID candidateId, AiCandidateDecisionAction action) { + public Decision { + Objects.requireNonNull(candidateId); + Objects.requireNonNull(action); + } + } +} diff --git a/src/main/java/com/fowoco/server/airun/application/AiCandidateDecisionService.java b/src/main/java/com/fowoco/server/airun/application/AiCandidateDecisionService.java new file mode 100644 index 00000000..24649e99 --- /dev/null +++ b/src/main/java/com/fowoco/server/airun/application/AiCandidateDecisionService.java @@ -0,0 +1,309 @@ +package com.fowoco.server.airun.application; + +import com.fowoco.server.aiintegration.application.model.AiAnalysisOutcome; +import com.fowoco.server.airun.application.error.AiRunErrorCode; +import com.fowoco.server.airun.application.port.AiCandidateDecisionRepository; +import com.fowoco.server.airun.application.port.AiCandidateDecisionRepository.DecisionContext; +import com.fowoco.server.airun.application.port.AiCandidateDecisionRepository.NewBatch; +import com.fowoco.server.airun.application.port.AiCandidateDecisionRepository.NewDecision; +import com.fowoco.server.airun.domain.AiCandidateDecisionAction; +import com.fowoco.server.airun.domain.AiRunStatus; +import com.fowoco.server.auth.application.ActorAuthorizer; +import com.fowoco.server.auth.application.ActorContext; +import com.fowoco.server.auth.domain.UserRole; +import com.fowoco.server.audit.application.port.AuditEventRepository; +import com.fowoco.server.audit.domain.ActorType; +import com.fowoco.server.audit.domain.AuditAction; +import com.fowoco.server.audit.domain.AuditEvent; +import com.fowoco.server.audit.domain.AuditTargetType; +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.port.AiCandidateTaskCreator; +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.time.Clock; +import java.time.Instant; +import java.util.Comparator; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.function.Function; +import java.util.stream.Collectors; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +@Service +public class AiCandidateDecisionService { + + private static final String AUDIT_EVENT_VERSION = "1"; + + private final ActorAuthorizer actorAuthorizer; + private final TenantDatabaseContext tenantDatabaseContext; + private final AiCandidateDecisionRepository repository; + private final AiCandidateTaskCreator taskCreator; + private final AuditEventRepository auditEventRepository; + private final UuidGenerator uuidGenerator; + private final Clock clock; + + public AiCandidateDecisionService( + ActorAuthorizer actorAuthorizer, + TenantDatabaseContext tenantDatabaseContext, + AiCandidateDecisionRepository repository, + AiCandidateTaskCreator taskCreator, + AuditEventRepository auditEventRepository, + UuidGenerator uuidGenerator, + Clock clock + ) { + this.actorAuthorizer = actorAuthorizer; + this.tenantDatabaseContext = tenantDatabaseContext; + this.repository = repository; + this.taskCreator = taskCreator; + this.auditEventRepository = auditEventRepository; + this.uuidGenerator = uuidGenerator; + this.clock = clock; + } + + @Transactional + public AiCandidateDecisionResult decide( + UUID aiRunId, + String idempotencyKey, + AiCandidateDecisionCommand command, + ActorContext actor, + RequestMetadata metadata + ) { + tenantDatabaseContext.setCompanyIdForCurrentTransaction(actor.companyId()); + actorAuthorizer.requireHrWrite(actor); + validateCommand(command); + + String keyHash = sha256(normalizeIdempotencyKey(idempotencyKey)); + String payloadHash = payloadHash(command); + DecisionContext context = repository.lockRun(aiRunId, actor.companyId()); + var existing = repository.findBatch(aiRunId, actor.companyId(), keyHash); + if (existing.isPresent()) { + if (!existing.get().payloadHash().equals(payloadHash)) { + throw new ApiException(AiRunErrorCode.AI_RUN_IDEMPOTENCY_CONFLICT); + } + return existing.get().result(); + } + + validateRun(context, command.expectedRunVersion()); + Map candidates = context.candidates().stream() + .collect(Collectors.toMap(AiRunCandidateResult::candidateId, Function.identity())); + List resolved = command.decisions().stream() + .map(decision -> resolve(decision, candidates, actor.companyId())) + .toList(); + List accepted = resolved.stream() + .filter(decision -> decision.action() == AiCandidateDecisionAction.ACCEPT) + .toList(); + if (accepted.size() > 1) { + throw new ApiException(AiRunErrorCode.AI_RUN_INVALID_DECISION); + } + + Instant now = clock.instant(); + UUID batchId = uuidGenerator.generate(); + repository.insertBatch(new NewBatch( + batchId, + aiRunId, + actor.companyId(), + actor.actorId(), + keyHash, + payloadHash, + now + )); + List persisted = resolved.stream() + .map(decision -> persist(batchId, aiRunId, actor.companyId(), decision, now)) + .toList(); + + UUID caseId = null; + List taskIds = List.of(); + if (!accepted.isEmpty()) { + ResolvedDecision acceptedDecision = accepted.get(0); + if (!acceptedDecision.candidate().missingSlots().isEmpty()) { + throw new ApiException(AiRunErrorCode.AI_RUN_CANDIDATE_NOT_READY); + } + AiCandidateTaskCreator.CreationResult created = taskCreator.create( + new AiCandidateTaskCreator.CreationCommand( + aiRunId, + acceptedDecision.candidate().candidateId(), + acceptedDecision.candidate().workerId(), + context.detectedIntent(), + acceptedDecision.candidate().workflowId(), + acceptedDecision.candidate().extractedSlots() + ), + actor, + metadata + ); + caseId = created.caseId(); + taskIds = created.taskIds(); + UUID acceptedDecisionId = persisted.stream() + .filter(decision -> decision.action() == AiCandidateDecisionAction.ACCEPT) + .findFirst() + .orElseThrow() + .decisionId(); + repository.attachTasks(acceptedDecisionId, actor.companyId(), taskIds, now); + } + + long runVersion = repository.completeBatch( + batchId, + aiRunId, + actor.companyId(), + caseId, + command.expectedRunVersion(), + now + ); + appendAudit(aiRunId, actor, metadata, now, accepted.isEmpty(), taskIds.size()); + return new AiCandidateDecisionResult( + batchId, + aiRunId, + caseId, + taskIds, + persisted.stream() + .map(decision -> new AiCandidateDecisionResult.Decision( + decision.candidateId(), + decision.action() + )) + .toList(), + runVersion + ); + } + + private void validateCommand(AiCandidateDecisionCommand command) { + if (command.decisions().isEmpty() || command.decisions().size() > 20) { + throw new ApiException(AiRunErrorCode.AI_RUN_INVALID_DECISION); + } + HashSet candidateIds = new HashSet<>(); + if (command.decisions().stream().anyMatch(decision -> !candidateIds.add(decision.candidateId()))) { + throw new ApiException(AiRunErrorCode.AI_RUN_INVALID_DECISION); + } + } + + private void validateRun(DecisionContext context, long expectedVersion) { + if (context.version() != expectedVersion) { + throw new ApiException(AiRunErrorCode.AI_RUN_VERSION_CONFLICT); + } + if (context.status() != AiRunStatus.SUCCEEDED + || context.outcome() != AiAnalysisOutcome.REVIEW_REQUIRED + || context.detectedIntent() == null + || context.candidates().isEmpty()) { + throw new ApiException(AiRunErrorCode.AI_RUN_DECISION_NOT_ALLOWED); + } + } + + private ResolvedDecision resolve( + AiCandidateDecisionCommand.Decision decision, + Map candidates, + UUID companyId + ) { + AiRunCandidateResult candidate = candidates.get(decision.candidateId()); + if (candidate == null) { + throw new ApiException(AiRunErrorCode.AI_RUN_INVALID_DECISION); + } + if (repository.candidateAlreadyDecided(candidate.candidateId(), companyId)) { + throw new ApiException(AiRunErrorCode.AI_RUN_CANDIDATE_ALREADY_DECIDED); + } + return new ResolvedDecision(candidate, decision.action()); + } + + private PersistedDecision persist( + UUID batchId, + UUID aiRunId, + UUID companyId, + ResolvedDecision decision, + Instant now + ) { + UUID decisionId = uuidGenerator.generate(); + repository.insertDecision(new NewDecision( + decisionId, + batchId, + aiRunId, + decision.candidate().candidateId(), + companyId, + decision.action(), + now + )); + return new PersistedDecision( + decisionId, + decision.candidate().candidateId(), + decision.action() + ); + } + + private String normalizeIdempotencyKey(String key) { + if (key == null || key.isBlank() || key.length() > 100) { + throw new ApiException(AiRunErrorCode.AI_RUN_INVALID_IDEMPOTENCY_KEY); + } + return key.trim(); + } + + private String payloadHash(AiCandidateDecisionCommand command) { + String decisions = command.decisions().stream() + .sorted(Comparator.comparing(decision -> decision.candidateId().toString())) + .map(decision -> decision.candidateId() + ":" + decision.action()) + .collect(Collectors.joining("|")); + return sha256(command.expectedRunVersion() + "|" + decisions); + } + + private String sha256(String value) { + try { + byte[] digest = MessageDigest.getInstance("SHA-256") + .digest(value.getBytes(StandardCharsets.UTF_8)); + return java.util.HexFormat.of().formatHex(digest); + } catch (NoSuchAlgorithmException exception) { + throw new IllegalStateException("SHA-256 must be available", exception); + } + } + + private void appendAudit( + UUID aiRunId, + ActorContext actor, + RequestMetadata metadata, + Instant now, + boolean discardedOnly, + int taskCount + ) { + auditEventRepository.append(new AuditEvent( + uuidGenerator.generate(), + actor.companyId(), + ActorType.HR_USER, + actor.actorId(), + effectiveRole(actor), + AuditAction.AI_RUN_CANDIDATES_DECIDED, + AuditTargetType.AI_RUN, + aiRunId, + metadata.requestId(), + metadata.traceId(), + AUDIT_EVENT_VERSION, + discardedOnly + ? "AI 업무 후보를 폐기함" + : "AI 업무 후보를 채택하고 업무카드 " + taskCount + "건을 생성함", + now + )); + } + + private UserRole effectiveRole(ActorContext actor) { + return actor.roles().stream() + .min(Comparator.comparingInt(role -> switch (role) { + case ADMIN -> 0; + case HR -> 1; + case VIEWER -> 2; + })) + .orElseThrow(); + } + + private record ResolvedDecision( + AiRunCandidateResult candidate, + AiCandidateDecisionAction action + ) { + } + + private record PersistedDecision( + UUID decisionId, + UUID candidateId, + AiCandidateDecisionAction action + ) { + } +} diff --git a/src/main/java/com/fowoco/server/airun/application/error/AiRunErrorCode.java b/src/main/java/com/fowoco/server/airun/application/error/AiRunErrorCode.java index 4fcc97e9..bdcbf18a 100644 --- a/src/main/java/com/fowoco/server/airun/application/error/AiRunErrorCode.java +++ b/src/main/java/com/fowoco/server/airun/application/error/AiRunErrorCode.java @@ -8,9 +8,13 @@ public enum AiRunErrorCode implements ApiErrorCode { AI_RUN_IDEMPOTENCY_CONFLICT(HttpStatus.CONFLICT, "같은 Idempotency-Key가 다른 요청에 이미 사용되었습니다."), AI_RUN_VERSION_CONFLICT(HttpStatus.CONFLICT, "다른 요청에서 먼저 변경했습니다. 최신 상태를 다시 확인해 주세요."), AI_RUN_ANSWERS_NOT_ALLOWED(HttpStatus.UNPROCESSABLE_ENTITY, "현재 상태에서는 추가 답변을 제출할 수 없습니다."), + AI_RUN_DECISION_NOT_ALLOWED(HttpStatus.UNPROCESSABLE_ENTITY, "현재 상태에서는 AI 업무 후보를 결정할 수 없습니다."), + AI_RUN_CANDIDATE_NOT_READY(HttpStatus.UNPROCESSABLE_ENTITY, "누락정보가 있는 AI 업무 후보는 채택할 수 없습니다."), + AI_RUN_CANDIDATE_ALREADY_DECIDED(HttpStatus.CONFLICT, "이미 결정된 AI 업무 후보입니다."), AI_RUN_INVALID_INSTRUCTION(HttpStatus.BAD_REQUEST, "업무 요청 문장을 확인해 주세요."), AI_RUN_INVALID_IDEMPOTENCY_KEY(HttpStatus.BAD_REQUEST, "Idempotency-Key를 확인해 주세요."), - AI_RUN_INVALID_ANSWER(HttpStatus.BAD_REQUEST, "추가 답변의 항목과 값을 확인해 주세요."); + AI_RUN_INVALID_ANSWER(HttpStatus.BAD_REQUEST, "추가 답변의 항목과 값을 확인해 주세요."), + AI_RUN_INVALID_DECISION(HttpStatus.BAD_REQUEST, "AI 업무 후보 결정값을 확인해 주세요."); private final HttpStatus status; private final String defaultMessage; diff --git a/src/main/java/com/fowoco/server/airun/application/port/AiCandidateDecisionRepository.java b/src/main/java/com/fowoco/server/airun/application/port/AiCandidateDecisionRepository.java new file mode 100644 index 00000000..2bcb3f30 --- /dev/null +++ b/src/main/java/com/fowoco/server/airun/application/port/AiCandidateDecisionRepository.java @@ -0,0 +1,86 @@ +package com.fowoco.server.airun.application.port; + +import com.fowoco.server.aiintegration.application.model.AiAnalysisOutcome; +import com.fowoco.server.airun.application.AiCandidateDecisionResult; +import com.fowoco.server.airun.application.AiRunCandidateResult; +import com.fowoco.server.airun.domain.AiCandidateDecisionAction; +import com.fowoco.server.airun.domain.AiRunStatus; +import java.time.Instant; +import java.util.List; +import java.util.Optional; +import java.util.UUID; + +public interface AiCandidateDecisionRepository { + + DecisionContext lockRun(UUID aiRunId, UUID companyId); + + Optional findBatch( + UUID aiRunId, + UUID companyId, + String idempotencyKeyHash + ); + + boolean candidateAlreadyDecided(UUID candidateId, UUID companyId); + + void insertBatch(NewBatch batch); + + void insertDecision(NewDecision decision); + + void attachTasks( + UUID decisionId, + UUID companyId, + List taskIds, + Instant createdAt + ); + + long completeBatch( + UUID decisionBatchId, + UUID aiRunId, + UUID companyId, + UUID caseId, + long expectedRunVersion, + Instant completedAt + ); + + record DecisionContext( + UUID aiRunId, + UUID companyId, + AiRunStatus status, + AiAnalysisOutcome outcome, + String detectedIntent, + long version, + List candidates + ) { + public DecisionContext { + candidates = List.copyOf(candidates); + } + } + + record StoredBatch( + String payloadHash, + AiCandidateDecisionResult result + ) { + } + + record NewBatch( + UUID decisionBatchId, + UUID aiRunId, + UUID companyId, + UUID decidedBy, + String idempotencyKeyHash, + String payloadHash, + Instant createdAt + ) { + } + + record NewDecision( + UUID decisionId, + UUID decisionBatchId, + UUID aiRunId, + UUID candidateId, + UUID companyId, + AiCandidateDecisionAction action, + Instant createdAt + ) { + } +} diff --git a/src/main/java/com/fowoco/server/airun/domain/AiCandidateDecisionAction.java b/src/main/java/com/fowoco/server/airun/domain/AiCandidateDecisionAction.java new file mode 100644 index 00000000..1803be5f --- /dev/null +++ b/src/main/java/com/fowoco/server/airun/domain/AiCandidateDecisionAction.java @@ -0,0 +1,6 @@ +package com.fowoco.server.airun.domain; + +public enum AiCandidateDecisionAction { + ACCEPT, + DISCARD +} diff --git a/src/main/java/com/fowoco/server/airun/infrastructure/persistence/JdbcAiCandidateDecisionRepository.java b/src/main/java/com/fowoco/server/airun/infrastructure/persistence/JdbcAiCandidateDecisionRepository.java new file mode 100644 index 00000000..2f490e5a --- /dev/null +++ b/src/main/java/com/fowoco/server/airun/infrastructure/persistence/JdbcAiCandidateDecisionRepository.java @@ -0,0 +1,350 @@ +package com.fowoco.server.airun.infrastructure.persistence; + +import com.fowoco.server.aiintegration.application.model.AiAnalysisOutcome; +import com.fowoco.server.airun.application.AiCandidateDecisionResult; +import com.fowoco.server.airun.application.AiRunCandidateResult; +import com.fowoco.server.airun.application.error.AiRunErrorCode; +import com.fowoco.server.airun.application.port.AiCandidateDecisionRepository; +import com.fowoco.server.airun.domain.AiCandidateDecisionAction; +import com.fowoco.server.airun.domain.AiRunStatus; +import com.fowoco.server.common.error.ApiException; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.sql.Timestamp; +import java.time.Instant; +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.UUID; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Repository; +import tools.jackson.core.JacksonException; +import tools.jackson.databind.ObjectMapper; + +@Repository +public class JdbcAiCandidateDecisionRepository implements AiCandidateDecisionRepository { + + private final JdbcTemplate jdbcTemplate; + private final ObjectMapper objectMapper; + + public JdbcAiCandidateDecisionRepository(JdbcTemplate jdbcTemplate, ObjectMapper objectMapper) { + this.jdbcTemplate = jdbcTemplate; + this.objectMapper = objectMapper; + } + + @Override + public DecisionContext lockRun(UUID aiRunId, UUID companyId) { + RunRow run = jdbcTemplate.query( + """ + SELECT status, analysis_outcome, detected_intent, version + FROM ai_run + WHERE ai_run_id = ? AND company_id = ? + FOR UPDATE + """, + (resultSet, rowNum) -> new RunRow( + AiRunStatus.valueOf(resultSet.getString("status")), + nullableOutcome(resultSet.getString("analysis_outcome")), + resultSet.getString("detected_intent"), + resultSet.getLong("version") + ), + aiRunId, + companyId + ).stream().findFirst().orElseThrow(() -> new ApiException(AiRunErrorCode.AI_RUN_NOT_FOUND)); + + List candidates = jdbcTemplate.query( + """ + SELECT candidate.ai_candidate_id, candidate.candidate_ref, + candidate.worker_id, candidate.workflow_id, + candidate.extracted_slots_json, candidate.missing_slots_json, + candidate.confidence + FROM ai_candidate candidate + JOIN ai_attempt attempt + ON attempt.ai_attempt_id = candidate.ai_attempt_id + AND attempt.company_id = candidate.company_id + WHERE candidate.ai_run_id = ? + AND candidate.company_id = ? + AND attempt.sequence_no = ( + SELECT MAX(latest.sequence_no) + FROM ai_attempt latest + WHERE latest.ai_run_id = ? AND latest.company_id = ? + ) + ORDER BY candidate.created_at, candidate.candidate_ref + """, + (resultSet, rowNum) -> new AiRunCandidateResult( + uuid(resultSet, "ai_candidate_id"), + resultSet.getString("candidate_ref"), + uuid(resultSet, "worker_id"), + resultSet.getString("workflow_id"), + decodeStringMap(resultSet.getString("extracted_slots_json")), + decodeStringList(resultSet.getString("missing_slots_json")), + resultSet.getBigDecimal("confidence") + ), + aiRunId, + companyId, + aiRunId, + companyId + ); + return new DecisionContext( + aiRunId, + companyId, + run.status(), + run.outcome(), + run.detectedIntent(), + run.version(), + candidates + ); + } + + @Override + public Optional findBatch( + UUID aiRunId, + UUID companyId, + String idempotencyKeyHash + ) { + Optional batch = jdbcTemplate.query( + """ + SELECT decision_batch_id, payload_hash, case_id, resulting_run_version + FROM ai_candidate_decision_batch + WHERE ai_run_id = ? AND company_id = ? AND idempotency_key_hash = ? + """, + (resultSet, rowNum) -> new BatchRow( + uuid(resultSet, "decision_batch_id"), + resultSet.getString("payload_hash"), + nullableUuid(resultSet, "case_id"), + nullableLong(resultSet, "resulting_run_version") + ), + aiRunId, + companyId, + idempotencyKeyHash + ).stream().findFirst(); + return batch.map(row -> storedBatch(row, aiRunId, companyId)); + } + + @Override + public boolean candidateAlreadyDecided(UUID candidateId, UUID companyId) { + Integer count = jdbcTemplate.queryForObject( + """ + SELECT COUNT(*) + FROM ai_candidate_decision + WHERE ai_candidate_id = ? AND company_id = ? + """, + Integer.class, + candidateId, + companyId + ); + return count != null && count > 0; + } + + @Override + public void insertBatch(NewBatch batch) { + jdbcTemplate.update( + """ + INSERT INTO ai_candidate_decision_batch ( + decision_batch_id, ai_run_id, company_id, decided_by, + idempotency_key_hash, payload_hash, created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?) + """, + batch.decisionBatchId(), + batch.aiRunId(), + batch.companyId(), + batch.decidedBy(), + batch.idempotencyKeyHash(), + batch.payloadHash(), + timestamp(batch.createdAt()) + ); + } + + @Override + public void insertDecision(NewDecision decision) { + jdbcTemplate.update( + """ + INSERT INTO ai_candidate_decision ( + decision_id, decision_batch_id, ai_run_id, ai_candidate_id, + company_id, action, created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?) + """, + decision.decisionId(), + decision.decisionBatchId(), + decision.aiRunId(), + decision.candidateId(), + decision.companyId(), + decision.action().name(), + timestamp(decision.createdAt()) + ); + } + + @Override + public void attachTasks( + UUID decisionId, + UUID companyId, + List taskIds, + Instant createdAt + ) { + for (int index = 0; index < taskIds.size(); index++) { + jdbcTemplate.update( + """ + INSERT INTO ai_candidate_decision_task ( + decision_id, task_id, company_id, sequence_no, created_at + ) VALUES (?, ?, ?, ?, ?) + """, + decisionId, + taskIds.get(index), + companyId, + index + 1, + timestamp(createdAt) + ); + } + } + + @Override + public long completeBatch( + UUID decisionBatchId, + UUID aiRunId, + UUID companyId, + UUID caseId, + long expectedRunVersion, + Instant completedAt + ) { + int updated = jdbcTemplate.update( + """ + UPDATE ai_run + SET updated_at = ?, version = version + 1 + WHERE ai_run_id = ? AND company_id = ? AND version = ? + """, + timestamp(completedAt), + aiRunId, + companyId, + expectedRunVersion + ); + if (updated != 1) { + throw new ApiException(AiRunErrorCode.AI_RUN_VERSION_CONFLICT); + } + long resultingVersion = expectedRunVersion + 1; + jdbcTemplate.update( + """ + UPDATE ai_candidate_decision_batch + SET case_id = ?, resulting_run_version = ?, completed_at = ? + WHERE decision_batch_id = ? AND ai_run_id = ? AND company_id = ? + """, + caseId, + resultingVersion, + timestamp(completedAt), + decisionBatchId, + aiRunId, + companyId + ); + return resultingVersion; + } + + private StoredBatch storedBatch(BatchRow batch, UUID aiRunId, UUID companyId) { + UUID batchId = batch.batchId(); + Long resultingVersion = batch.resultingRunVersion(); + if (resultingVersion == null) { + throw new IllegalStateException("candidate decision batch is incomplete"); + } + List decisions = jdbcTemplate.query( + """ + SELECT ai_candidate_id, action + FROM ai_candidate_decision + WHERE decision_batch_id = ? AND company_id = ? + ORDER BY created_at, ai_candidate_id + """, + (decisionSet, rowNum) -> new AiCandidateDecisionResult.Decision( + uuid(decisionSet, "ai_candidate_id"), + AiCandidateDecisionAction.valueOf(decisionSet.getString("action")) + ), + batchId, + companyId + ); + List taskIds = jdbcTemplate.query( + """ + SELECT decision_task.task_id + FROM ai_candidate_decision_task decision_task + JOIN ai_candidate_decision decision + ON decision.decision_id = decision_task.decision_id + AND decision.company_id = decision_task.company_id + WHERE decision.decision_batch_id = ? AND decision.company_id = ? + ORDER BY decision_task.sequence_no + """, + (taskSet, rowNum) -> uuid(taskSet, "task_id"), + batchId, + companyId + ); + return new StoredBatch( + batch.payloadHash(), + new AiCandidateDecisionResult( + batchId, + aiRunId, + batch.caseId(), + taskIds, + decisions, + resultingVersion + ) + ); + } + + @SuppressWarnings("unchecked") + private Map decodeStringMap(String json) { + try { + Map raw = objectMapper.readValue(json, Map.class); + Map result = new LinkedHashMap<>(); + raw.forEach((key, value) -> result.put(key, value == null ? null : value.toString())); + return result; + } catch (JacksonException exception) { + throw new IllegalStateException("stored candidate slots cannot be decoded", exception); + } + } + + @SuppressWarnings("unchecked") + private List decodeStringList(String json) { + try { + List raw = objectMapper.readValue(json, List.class); + List result = new ArrayList<>(); + raw.forEach(value -> result.add(value.toString())); + return result; + } catch (JacksonException exception) { + throw new IllegalStateException("stored candidate missing slots cannot be decoded", exception); + } + } + + private AiAnalysisOutcome nullableOutcome(String value) { + return value == null ? null : AiAnalysisOutcome.valueOf(value); + } + + private UUID nullableUuid(ResultSet resultSet, String column) throws SQLException { + Object value = resultSet.getObject(column); + return value == null ? null : value instanceof UUID uuid ? uuid : UUID.fromString(value.toString()); + } + + private UUID uuid(ResultSet resultSet, String column) throws SQLException { + return Optional.ofNullable(nullableUuid(resultSet, column)) + .orElseThrow(() -> new SQLException(column + " must not be null")); + } + + private Long nullableLong(ResultSet resultSet, String column) throws SQLException { + Object value = resultSet.getObject(column); + return value == null ? null : ((Number) value).longValue(); + } + + private Timestamp timestamp(Instant instant) { + return Timestamp.from(instant); + } + + private record RunRow( + AiRunStatus status, + AiAnalysisOutcome outcome, + String detectedIntent, + long version + ) { + } + + private record BatchRow( + UUID batchId, + String payloadHash, + UUID caseId, + Long resultingRunVersion + ) { + } +} diff --git a/src/main/java/com/fowoco/server/airun/infrastructure/persistence/JdbcAiRunRepository.java b/src/main/java/com/fowoco/server/airun/infrastructure/persistence/JdbcAiRunRepository.java index ffdb0f33..2dbf2614 100644 --- a/src/main/java/com/fowoco/server/airun/infrastructure/persistence/JdbcAiRunRepository.java +++ b/src/main/java/com/fowoco/server/airun/infrastructure/persistence/JdbcAiRunRepository.java @@ -575,7 +575,9 @@ private String detectedIntent(AiAnalysisResponse response) { if (response.contextRequirement() != null) { return response.contextRequirement().detectedIntent(); } - return response.candidates().isEmpty() ? null : response.candidates().get(0).workflowId(); + // ANALYZE candidates carry canonical Workflow IDs, not Intent codes. + // Returning null preserves the Intent detected during PLAN through the COALESCE update. + return null; } private String encode(Object value) { diff --git a/src/main/java/com/fowoco/server/audit/domain/AuditAction.java b/src/main/java/com/fowoco/server/audit/domain/AuditAction.java index ce496922..e957928d 100644 --- a/src/main/java/com/fowoco/server/audit/domain/AuditAction.java +++ b/src/main/java/com/fowoco/server/audit/domain/AuditAction.java @@ -17,6 +17,7 @@ public enum AuditAction { DOCUMENT_REQUEST_DRAFT_SAVED, AI_RUN_CREATED, AI_RUN_ANSWERS_SUBMITTED, + AI_RUN_CANDIDATES_DECIDED, WORKER_LINK_RESPONSE_SUBMITTED, WORKER_LINK_ACCESSED } diff --git a/src/main/java/com/fowoco/server/casework/infrastructure/persistence/JdbcTaskCaseRegistrar.java b/src/main/java/com/fowoco/server/casework/infrastructure/persistence/JdbcTaskCaseRegistrar.java index 931b7aa2..70dc7a0b 100644 --- a/src/main/java/com/fowoco/server/casework/infrastructure/persistence/JdbcTaskCaseRegistrar.java +++ b/src/main/java/com/fowoco/server/casework/infrastructure/persistence/JdbcTaskCaseRegistrar.java @@ -7,6 +7,7 @@ import com.fowoco.server.workflow.domain.WorkflowDefinition; import java.time.LocalDate; import java.util.ArrayList; +import java.util.Comparator; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; @@ -29,14 +30,37 @@ public JdbcTaskCaseRegistrar(JdbcTemplate jdbcTemplate, ObjectMapper objectMappe @Override public void register(Task task, WorkflowDefinition workflow, LocalDate today) { + registerComposite(List.of(new CaseTask(task, workflow)), today); + } + + @Override + public void registerComposite(List caseTasks, LocalDate today) { + if (caseTasks == null || caseTasks.isEmpty()) { + throw new IllegalArgumentException("caseTasks must not be empty"); + } + List orderedTasks = caseTasks.stream() + .sorted(Comparator + .comparingInt((CaseTask caseTask) -> candidateOrder(caseTask.task())) + .thenComparing(caseTask -> caseTask.task().taskId())) + .toList(); + Task first = orderedTasks.get(0).task(); + boolean mixedScope = orderedTasks.stream().anyMatch(caseTask -> { + Task task = caseTask.task(); + return !first.caseId().equals(task.caseId()) + || !first.companyId().equals(task.companyId()) + || !first.workerId().equals(task.workerId()); + }); + if (mixedScope) { + throw new ApiException(TaskErrorCode.CASE_WORKER_MISMATCH); + } List existingWorkerIds = jdbcTemplate.query( "SELECT worker_id FROM workflow_case WHERE case_id = ? AND company_id = ?", (resultSet, rowNumber) -> resultSet.getObject("worker_id", UUID.class), - task.caseId(), - task.companyId() + first.caseId(), + first.companyId() ); if (!existingWorkerIds.isEmpty()) { - if (!existingWorkerIds.get(0).equals(task.workerId())) { + if (!existingWorkerIds.get(0).equals(first.workerId())) { throw new ApiException(TaskErrorCode.CASE_WORKER_MISMATCH); } return; @@ -50,48 +74,95 @@ INSERT INTO workflow_case ( created_at, updated_at, version ) VALUES (?, ?, ?, ?, 'ACTIVE', ?, ?, ?, ?, ?, ?, 0) """, - task.caseId(), - task.companyId(), - task.workerId(), - task.title(), - priority(task.dueDate(), today), - task.workflowCatalogVersion(), - snapshot(task, workflow), - task.createdBy(), - task.createdAt(), - task.updatedAt() + first.caseId(), + first.companyId(), + first.workerId(), + orderedTasks.size() == 1 ? first.title() : "3년 만료 연장 준비", + priority(orderedTasks, today), + first.workflowCatalogVersion(), + snapshot(orderedTasks), + first.createdBy(), + first.createdAt(), + first.updatedAt() ); } - private String snapshot(Task task, WorkflowDefinition workflow) { + private String snapshot(List caseTasks) { + List> steps = java.util.stream.IntStream + .range(0, caseTasks.size()) + .mapToObj(index -> snapshotStep(caseTasks.get(index), index + 1)) + .toList(); + Task first = caseTasks.get(0).task(); + Map snapshot = new LinkedHashMap<>(); + snapshot.put("workflow_catalog_version", first.workflowCatalogVersion()); + snapshot.put("steps", steps); + try { + return objectMapper.writeValueAsString(snapshot); + } catch (JacksonException exception) { + throw new IllegalStateException("task workflow snapshot cannot be encoded", exception); + } + } + + private Map snapshotStep(CaseTask caseTask, int fallbackOrder) { + Task task = caseTask.task(); + WorkflowDefinition workflow = caseTask.workflow(); Map conditions = new LinkedHashMap<>(); conditions.put("required_slots", sorted(workflow.requiredSlots())); conditions.put("completion_evidence", List.copyOf(workflow.completionEvidence())); + Map businessData = businessData(task); + List.of( + "approval_required", + "depends_on_task_id", + "dependency_reason", + "missing_information", + "submission_due_offset_days" + ).forEach(key -> { + if (businessData.containsKey(key)) { + conditions.put(key, businessData.get(key)); + } + }); Map step = new LinkedHashMap<>(); - step.put("order", 1); + step.put("order", candidateOrder(task, fallbackOrder)); step.put("task_id", task.taskId().toString()); step.put("workflow_id", task.workflowId()); step.put("task_type", task.taskType().name()); step.put("required_conditions", conditions); + return Map.copyOf(step); + } - Map snapshot = new LinkedHashMap<>(); - snapshot.put("workflow_catalog_version", task.workflowCatalogVersion()); - snapshot.put("steps", List.of(step)); + @SuppressWarnings("unchecked") + private Map businessData(Task task) { try { - return objectMapper.writeValueAsString(snapshot); + return objectMapper.readValue(task.businessDataJson(), Map.class); } catch (JacksonException exception) { - throw new IllegalStateException("manual task workflow snapshot cannot be encoded", exception); + throw new IllegalStateException("task business data cannot be decoded", exception); } } + private int candidateOrder(Task task) { + return candidateOrder(task, Integer.MAX_VALUE); + } + + private int candidateOrder(Task task, int fallback) { + Object value = businessData(task).get("candidate_order"); + return value instanceof Number number && number.intValue() > 0 + ? number.intValue() + : fallback; + } + private List sorted(Iterable values) { List result = new ArrayList<>(); values.forEach(result::add); return result.stream().sorted().toList(); } - private String priority(LocalDate dueDate, LocalDate today) { + private String priority(List caseTasks, LocalDate today) { + LocalDate dueDate = caseTasks.stream() + .map(caseTask -> caseTask.task().dueDate()) + .filter(java.util.Objects::nonNull) + .min(LocalDate::compareTo) + .orElse(null); if (dueDate == null) { return "NORMAL"; } diff --git a/src/main/java/com/fowoco/server/task/application/AiCandidateTaskCreationService.java b/src/main/java/com/fowoco/server/task/application/AiCandidateTaskCreationService.java new file mode 100644 index 00000000..4ae79c2e --- /dev/null +++ b/src/main/java/com/fowoco/server/task/application/AiCandidateTaskCreationService.java @@ -0,0 +1,368 @@ +package com.fowoco.server.task.application; + +import com.fowoco.server.auth.application.ActorAuthorizer; +import com.fowoco.server.auth.application.ActorContext; +import com.fowoco.server.auth.domain.UserRole; +import com.fowoco.server.audit.application.port.AuditEventRepository; +import com.fowoco.server.audit.domain.ActorType; +import com.fowoco.server.audit.domain.AuditAction; +import com.fowoco.server.audit.domain.AuditEvent; +import com.fowoco.server.audit.domain.AuditTargetType; +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.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.AiCandidateTaskCreator; +import com.fowoco.server.task.application.port.TaskCaseRegistrar; +import com.fowoco.server.task.application.port.TaskChecklistRepository; +import com.fowoco.server.task.application.port.TaskRepository; +import com.fowoco.server.task.domain.Task; +import com.fowoco.server.task.domain.TaskChecklistItem; +import com.fowoco.server.task.domain.TaskSource; +import com.fowoco.server.task.domain.TaskStatus; +import com.fowoco.server.task.domain.TaskType; +import com.fowoco.server.worker.application.WorkerTaskContext; +import com.fowoco.server.worker.application.port.WorkerTaskContextReader; +import com.fowoco.server.workflow.application.WorkflowCatalogService; +import com.fowoco.server.workflow.domain.WorkflowCatalog; +import com.fowoco.server.workflow.domain.WorkflowDefinition; +import java.time.Clock; +import java.time.Instant; +import java.time.LocalDate; +import java.time.format.DateTimeParseException; +import java.util.ArrayList; +import java.util.Comparator; +import java.util.EnumMap; +import java.util.EnumSet; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +@Service +public class AiCandidateTaskCreationService implements AiCandidateTaskCreator { + + private static final String AUDIT_EVENT_VERSION = "1"; + private static final EnumSet EXPIRY_RENEWAL_TASK_TYPES = EnumSet.of( + TaskType.RECONTRACT, + TaskType.STAY_PERIOD_EXTENSION, + TaskType.EMPLOYMENT_PERIOD_EXTENSION + ); + + private final ActorAuthorizer actorAuthorizer; + private final TenantDatabaseContext tenantDatabaseContext; + private final TaskRepository taskRepository; + private final TaskChecklistRepository checklistRepository; + private final TaskCaseRegistrar taskCaseRegistrar; + private final WorkerTaskContextReader workerReader; + private final WorkflowCatalogService catalogService; + private final AuditEventRepository auditRepository; + private final DomainEventPublisher eventPublisher; + private final TaskContentCodec contentCodec; + private final UuidGenerator uuidGenerator; + private final Clock clock; + + public AiCandidateTaskCreationService( + ActorAuthorizer actorAuthorizer, + TenantDatabaseContext tenantDatabaseContext, + TaskRepository taskRepository, + TaskChecklistRepository checklistRepository, + TaskCaseRegistrar taskCaseRegistrar, + WorkerTaskContextReader workerReader, + WorkflowCatalogService catalogService, + AuditEventRepository auditRepository, + DomainEventPublisher eventPublisher, + TaskContentCodec contentCodec, + UuidGenerator uuidGenerator, + Clock clock + ) { + this.actorAuthorizer = actorAuthorizer; + this.tenantDatabaseContext = tenantDatabaseContext; + this.taskRepository = taskRepository; + this.checklistRepository = checklistRepository; + this.taskCaseRegistrar = taskCaseRegistrar; + this.workerReader = workerReader; + this.catalogService = catalogService; + this.auditRepository = auditRepository; + this.eventPublisher = eventPublisher; + this.contentCodec = contentCodec; + this.uuidGenerator = uuidGenerator; + this.clock = clock; + } + + @Override + @Transactional + public CreationResult create( + CreationCommand command, + ActorContext actor, + RequestMetadata metadata + ) { + tenantDatabaseContext.setCompanyIdForCurrentTransaction(actor.companyId()); + actorAuthorizer.requireHrWrite(actor); + WorkerTaskContext worker = workerReader + .findByIdAndCompanyId(command.workerId(), actor.companyId()) + .orElseThrow(() -> new ApiException(TaskErrorCode.WORKER_NOT_FOUND)); + if (!worker.canReceiveNewTask()) { + throw new ApiException(TaskErrorCode.WORKER_NOT_ELIGIBLE); + } + + WorkflowCatalog catalog = catalogService.getActiveCatalog(); + List workflows = catalog.findByIntent(command.detectedIntent()); + if (workflows.isEmpty() + || workflows.stream().noneMatch(workflow -> workflow.workflowId() + .equals(command.candidateWorkflowId()))) { + throw new ApiException(TaskErrorCode.WORKFLOW_TASK_TYPE_MISMATCH); + } + List plans = plans(workflows); + EnumSet plannedTaskTypes = plans.stream() + .map(TaskPlan::taskType) + .collect(java.util.stream.Collectors.toCollection(() -> EnumSet.noneOf(TaskType.class))); + if (!plannedTaskTypes.equals(EXPIRY_RENEWAL_TASK_TYPES)) { + throw new ApiException(TaskErrorCode.WORKFLOW_TASK_TYPE_MISMATCH); + } + + LocalDate dueDate = dueDate(command.extractedSlots(), worker); + Instant now = clock.instant(); + UUID caseId = uuidGenerator.generate(); + Map taskIds = new EnumMap<>(TaskType.class); + plans.forEach(plan -> taskIds.put(plan.taskType(), uuidGenerator.generate())); + + List caseTasks = plans.stream() + .map(plan -> createTask( + plan, + command, + actor, + catalog.bundleVersion(), + caseId, + taskIds, + worker, + dueDate, + now + )) + .toList(); + taskCaseRegistrar.registerComposite(caseTasks, LocalDate.now(clock)); + + List createdTaskIds = new ArrayList<>(); + for (TaskCaseRegistrar.CaseTask caseTask : caseTasks) { + Task saved = taskRepository.save(caseTask.task()); + checklistRepository.saveAll(caseTask.workflow().checklistItems().stream() + .map(template -> TaskChecklistItem.create( + uuidGenerator.generate(), + saved.taskId(), + saved.companyId(), + template.itemCode(), + template.label(), + template.required(), + now + )) + .toList()); + appendAudit(saved, actor, metadata, now); + eventPublisher.publish(TaskDomainEvents.taskCreated( + uuidGenerator.generate(), + saved, + actor, + metadata, + now + )); + createdTaskIds.add(saved.taskId()); + } + return new CreationResult(caseId, createdTaskIds); + } + + private List plans(List workflows) { + Map workflowByType = new EnumMap<>(TaskType.class); + workflows.forEach(workflow -> workflow.supportedTaskTypes().forEach(taskType -> { + if (workflowByType.putIfAbsent(taskType, workflow) != null) { + throw new ApiException(TaskErrorCode.WORKFLOW_TASK_TYPE_MISMATCH); + } + })); + return workflowByType.entrySet().stream() + .map(entry -> new TaskPlan(entry.getKey(), entry.getValue())) + .sorted(Comparator.comparingInt(plan -> order(plan.taskType()))) + .toList(); + } + + private TaskCaseRegistrar.CaseTask createTask( + TaskPlan plan, + CreationCommand command, + ActorContext actor, + String catalogVersion, + UUID caseId, + Map taskIds, + WorkerTaskContext worker, + LocalDate dueDate, + Instant now + ) { + TaskType taskType = plan.taskType(); + Map businessData = businessData(command, taskType, taskIds); + List missingSlots = missingRequiredSlots( + plan.workflow(), + worker, + dueDate, + businessData + ); + if (!missingSlots.isEmpty()) { + throw new ApiException(TaskErrorCode.INVALID_AI_CANDIDATE_TASK_DATA); + } + String title = title(taskType); + String description = description(taskType); + EncodedTaskContent content = contentCodec.encode( + command.workerId(), + plan.workflow().workflowId(), + taskType.name(), + title, + description, + dueDate, + businessData + ); + Task task = Task.create( + taskIds.get(taskType), + actor.companyId(), + command.workerId(), + caseId, + taskType, + plan.workflow().workflowId(), + catalogVersion, + title, + description, + content.businessDataJson(), + content.criticalFingerprint(), + TaskSource.AI_CANDIDATE, + TaskStatus.DRAFT, + dueDate, + actor.actorId(), + now + ); + return new TaskCaseRegistrar.CaseTask(task, plan.workflow()); + } + + private Map businessData( + CreationCommand command, + TaskType taskType, + Map taskIds + ) { + Map data = new LinkedHashMap<>(); + data.put("ai_run_id", command.aiRunId().toString()); + data.put("ai_candidate_id", command.candidateId().toString()); + data.put("source_intent", command.detectedIntent()); + data.put("candidate_order", order(taskType)); + switch (taskType) { + case RECONTRACT -> data.put("approval_required", true); + case STAY_PERIOD_EXTENSION -> data.put("submission_due_offset_days", 7); + case EMPLOYMENT_PERIOD_EXTENSION -> { + data.put("depends_on_task_id", taskIds.get(TaskType.RECONTRACT).toString()); + data.put("dependency_reason", "SIGNED_CONTRACT_REQUIRED"); + } + } + return Map.copyOf(data); + } + + private List missingRequiredSlots( + WorkflowDefinition workflow, + WorkerTaskContext worker, + LocalDate dueDate, + Map businessData + ) { + List missing = new ArrayList<>(); + workflow.requiredSlots().stream().sorted().forEach(slot -> { + boolean present = switch (slot) { + case "worker_id" -> worker.workerId() != null; + case "due_at", "due_date" -> dueDate != null; + case "contract_start_date" -> worker.contractStartDate() != null; + case "contract_end_date" -> worker.contractEndDate() != null; + case "stay_expiry_date" -> worker.stayExpiryDate() != null; + default -> businessData.get(slot) != null; + }; + if (!present) { + missing.add(slot); + } + }); + return List.copyOf(missing); + } + + private LocalDate dueDate( + Map extractedSlots, + WorkerTaskContext worker + ) { + String value = extractedSlots.get("due_at"); + if (value == null || value.isBlank()) { + value = extractedSlots.get("due_date"); + } + try { + if (value != null && !value.isBlank()) { + return LocalDate.parse(value); + } + if (worker.stayExpiryDate() != null) { + return worker.stayExpiryDate(); + } + return worker.contractEndDate(); + } catch (DateTimeParseException exception) { + throw new ApiException(TaskErrorCode.INVALID_AI_CANDIDATE_TASK_DATA); + } + } + + private int order(TaskType taskType) { + return switch (taskType) { + case RECONTRACT -> 1; + case STAY_PERIOD_EXTENSION -> 2; + case EMPLOYMENT_PERIOD_EXTENSION -> 3; + }; + } + + private String title(TaskType taskType) { + return switch (taskType) { + case RECONTRACT -> "재계약 조건 확인"; + case STAY_PERIOD_EXTENSION -> "체류기간 연장 준비"; + case EMPLOYMENT_PERIOD_EXTENSION -> "취업활동기간 연장 준비"; + }; + } + + private String description(TaskType taskType) { + return switch (taskType) { + case RECONTRACT -> "재계약 조건과 계속 고용 의사를 검토합니다."; + case STAY_PERIOD_EXTENSION -> "체류기간 연장에 필요한 정보와 서류를 확인합니다."; + case EMPLOYMENT_PERIOD_EXTENSION -> "재계약 결과를 바탕으로 취업활동기간 연장을 준비합니다."; + }; + } + + private void appendAudit( + Task task, + ActorContext actor, + RequestMetadata metadata, + Instant now + ) { + auditRepository.append(new AuditEvent( + uuidGenerator.generate(), + task.companyId(), + ActorType.HR_USER, + actor.actorId(), + effectiveRole(actor), + AuditAction.TASK_CREATED, + AuditTargetType.TASK, + task.taskId(), + metadata.requestId(), + metadata.traceId(), + AUDIT_EVENT_VERSION, + "AI 후보를 채택하여 업무카드를 생성함", + now + )); + } + + private UserRole effectiveRole(ActorContext actor) { + return actor.roles().stream() + .min(Comparator.comparingInt(role -> switch (role) { + case ADMIN -> 0; + case HR -> 1; + case VIEWER -> 2; + })) + .orElseThrow(); + } + + private record TaskPlan(TaskType taskType, WorkflowDefinition workflow) { + } +} diff --git a/src/main/java/com/fowoco/server/task/application/error/TaskErrorCode.java b/src/main/java/com/fowoco/server/task/application/error/TaskErrorCode.java index 96f7a7c3..f6915feb 100644 --- a/src/main/java/com/fowoco/server/task/application/error/TaskErrorCode.java +++ b/src/main/java/com/fowoco/server/task/application/error/TaskErrorCode.java @@ -20,6 +20,10 @@ public enum TaskErrorCode implements ApiErrorCode { HttpStatus.UNPROCESSABLE_CONTENT, "업무카드에 저장할 수 없는 개인정보 또는 Secret이 포함되어 있습니다." ), + INVALID_AI_CANDIDATE_TASK_DATA( + HttpStatus.UNPROCESSABLE_CONTENT, + "AI 업무 후보의 Workflow 또는 필수정보를 확인해 주세요." + ), CHECKLIST_ITEM_NOT_FOUND(HttpStatus.NOT_FOUND, "체크리스트 항목을 찾을 수 없습니다."), CASE_WORKER_MISMATCH(HttpStatus.CONFLICT, "Case와 업무카드의 근로자가 일치하지 않습니다."), CONCURRENT_MODIFICATION(HttpStatus.CONFLICT, "업무카드가 다른 요청에서 변경되었습니다."), diff --git a/src/main/java/com/fowoco/server/task/application/port/AiCandidateTaskCreator.java b/src/main/java/com/fowoco/server/task/application/port/AiCandidateTaskCreator.java new file mode 100644 index 00000000..57e7fdde --- /dev/null +++ b/src/main/java/com/fowoco/server/task/application/port/AiCandidateTaskCreator.java @@ -0,0 +1,31 @@ +package com.fowoco.server.task.application.port; + +import com.fowoco.server.auth.application.ActorContext; +import com.fowoco.server.common.web.RequestMetadata; +import java.util.List; +import java.util.Map; +import java.util.UUID; + +public interface AiCandidateTaskCreator { + + CreationResult create(CreationCommand command, ActorContext actor, RequestMetadata metadata); + + record CreationCommand( + UUID aiRunId, + UUID candidateId, + UUID workerId, + String detectedIntent, + String candidateWorkflowId, + Map extractedSlots + ) { + public CreationCommand { + extractedSlots = Map.copyOf(extractedSlots); + } + } + + record CreationResult(UUID caseId, List taskIds) { + public CreationResult { + taskIds = List.copyOf(taskIds); + } + } +} diff --git a/src/main/java/com/fowoco/server/task/application/port/TaskCaseRegistrar.java b/src/main/java/com/fowoco/server/task/application/port/TaskCaseRegistrar.java index ede9f5e5..3df077e3 100644 --- a/src/main/java/com/fowoco/server/task/application/port/TaskCaseRegistrar.java +++ b/src/main/java/com/fowoco/server/task/application/port/TaskCaseRegistrar.java @@ -3,8 +3,14 @@ import com.fowoco.server.task.domain.Task; import com.fowoco.server.workflow.domain.WorkflowDefinition; import java.time.LocalDate; +import java.util.List; public interface TaskCaseRegistrar { void register(Task task, WorkflowDefinition workflow, LocalDate today); + + void registerComposite(List caseTasks, LocalDate today); + + record CaseTask(Task task, WorkflowDefinition workflow) { + } } diff --git a/src/main/resources/db/migration-postgresql/V23__prepare_ai_candidate_decision_rls.sql b/src/main/resources/db/migration-postgresql/V23__prepare_ai_candidate_decision_rls.sql new file mode 100644 index 00000000..a938b42c --- /dev/null +++ b/src/main/resources/db/migration-postgresql/V23__prepare_ai_candidate_decision_rls.sql @@ -0,0 +1,38 @@ +CREATE POLICY pl_ai_candidate_decision_batch_tenant_isolation + ON public.ai_candidate_decision_batch + FOR ALL + TO PUBLIC + USING ( + company_id = + NULLIF(pg_catalog.current_setting('app.company_id', true), '')::UUID + ) + WITH CHECK ( + company_id = + NULLIF(pg_catalog.current_setting('app.company_id', true), '')::UUID + ); + +CREATE POLICY pl_ai_candidate_decision_tenant_isolation + ON public.ai_candidate_decision + FOR ALL + TO PUBLIC + USING ( + company_id = + NULLIF(pg_catalog.current_setting('app.company_id', true), '')::UUID + ) + WITH CHECK ( + company_id = + NULLIF(pg_catalog.current_setting('app.company_id', true), '')::UUID + ); + +CREATE POLICY pl_ai_candidate_decision_task_tenant_isolation + ON public.ai_candidate_decision_task + FOR ALL + TO PUBLIC + USING ( + company_id = + NULLIF(pg_catalog.current_setting('app.company_id', true), '')::UUID + ) + WITH CHECK ( + company_id = + NULLIF(pg_catalog.current_setting('app.company_id', true), '')::UUID + ); diff --git a/src/main/resources/db/migration/V22__create_ai_candidate_decision.sql b/src/main/resources/db/migration/V22__create_ai_candidate_decision.sql new file mode 100644 index 00000000..31336615 --- /dev/null +++ b/src/main/resources/db/migration/V22__create_ai_candidate_decision.sql @@ -0,0 +1,94 @@ +ALTER TABLE ai_candidate + ADD CONSTRAINT uq_ai_candidate_id_company + UNIQUE (ai_candidate_id, company_id); + +CREATE TABLE ai_candidate_decision_batch ( + decision_batch_id UUID NOT NULL, + ai_run_id UUID NOT NULL, + company_id UUID NOT NULL, + decided_by UUID NOT NULL, + idempotency_key_hash VARCHAR(64) NOT NULL, + payload_hash VARCHAR(64) NOT NULL, + case_id UUID, + resulting_run_version BIGINT, + created_at TIMESTAMP(6) WITH TIME ZONE NOT NULL, + completed_at TIMESTAMP(6) WITH TIME ZONE, + CONSTRAINT pk_ai_candidate_decision_batch PRIMARY KEY (decision_batch_id), + CONSTRAINT uq_ai_candidate_decision_batch_id_company + UNIQUE (decision_batch_id, company_id), + CONSTRAINT uq_ai_candidate_decision_batch_idempotency + UNIQUE (company_id, ai_run_id, idempotency_key_hash), + CONSTRAINT fk_ai_candidate_decision_batch_run_company + FOREIGN KEY (ai_run_id, company_id) + REFERENCES ai_run (ai_run_id, company_id) ON DELETE CASCADE, + CONSTRAINT fk_ai_candidate_decision_batch_actor_company + FOREIGN KEY (decided_by, company_id) + REFERENCES user_account (user_id, company_id) ON DELETE RESTRICT, + CONSTRAINT fk_ai_candidate_decision_batch_case_company + FOREIGN KEY (case_id, company_id) + REFERENCES workflow_case (case_id, company_id) ON DELETE RESTRICT, + CONSTRAINT ck_ai_candidate_decision_batch_idempotency_hash + CHECK (CHAR_LENGTH(idempotency_key_hash) = 64), + CONSTRAINT ck_ai_candidate_decision_batch_payload_hash + CHECK (CHAR_LENGTH(payload_hash) = 64), + CONSTRAINT ck_ai_candidate_decision_batch_version + CHECK (resulting_run_version IS NULL OR resulting_run_version >= 0), + CONSTRAINT ck_ai_candidate_decision_batch_completion CHECK ( + (completed_at IS NULL AND resulting_run_version IS NULL) + OR (completed_at IS NOT NULL AND resulting_run_version IS NOT NULL) + ) +); + +CREATE TABLE ai_candidate_decision ( + decision_id UUID NOT NULL, + decision_batch_id UUID NOT NULL, + ai_run_id UUID NOT NULL, + ai_candidate_id UUID NOT NULL, + company_id UUID NOT NULL, + action VARCHAR(20) NOT NULL, + created_at TIMESTAMP(6) WITH TIME ZONE NOT NULL, + CONSTRAINT pk_ai_candidate_decision PRIMARY KEY (decision_id), + CONSTRAINT uq_ai_candidate_decision_id_company + UNIQUE (decision_id, company_id), + CONSTRAINT uq_ai_candidate_decision_batch_candidate + UNIQUE (decision_batch_id, ai_candidate_id), + CONSTRAINT uq_ai_candidate_decision_candidate + UNIQUE (company_id, ai_candidate_id), + CONSTRAINT fk_ai_candidate_decision_batch_company + FOREIGN KEY (decision_batch_id, company_id) + REFERENCES ai_candidate_decision_batch (decision_batch_id, company_id) + ON DELETE CASCADE, + CONSTRAINT fk_ai_candidate_decision_run_company + FOREIGN KEY (ai_run_id, company_id) + REFERENCES ai_run (ai_run_id, company_id) ON DELETE CASCADE, + CONSTRAINT fk_ai_candidate_decision_candidate_company + FOREIGN KEY (ai_candidate_id, company_id) + REFERENCES ai_candidate (ai_candidate_id, company_id) ON DELETE RESTRICT, + CONSTRAINT ck_ai_candidate_decision_action + CHECK (action IN ('ACCEPT', 'DISCARD')) +); + +CREATE TABLE ai_candidate_decision_task ( + decision_id UUID NOT NULL, + task_id UUID NOT NULL, + company_id UUID NOT NULL, + sequence_no INTEGER NOT NULL, + created_at TIMESTAMP(6) WITH TIME ZONE NOT NULL, + CONSTRAINT pk_ai_candidate_decision_task + PRIMARY KEY (decision_id, task_id), + CONSTRAINT fk_ai_candidate_decision_task_decision_company + FOREIGN KEY (decision_id, company_id) + REFERENCES ai_candidate_decision (decision_id, company_id) ON DELETE CASCADE, + CONSTRAINT fk_ai_candidate_decision_task_task_company + FOREIGN KEY (task_id, company_id) + REFERENCES task (task_id, company_id) ON DELETE RESTRICT, + CONSTRAINT ck_ai_candidate_decision_task_sequence + CHECK (sequence_no > 0) +); + +CREATE INDEX idx_ai_candidate_decision_batch_run + ON ai_candidate_decision_batch (company_id, ai_run_id, created_at); +CREATE INDEX idx_ai_candidate_decision_run + ON ai_candidate_decision (company_id, ai_run_id, created_at); +CREATE INDEX idx_ai_candidate_decision_task_task + ON ai_candidate_decision_task (company_id, task_id); From e89710a90b91afd9c47f56e55ad94604ee0f5e61 Mon Sep 17 00:00:00 2001 From: hywznn Date: Thu, 6 Aug 2026 16:35:10 +0900 Subject: [PATCH 2/3] =?UTF-8?q?test(airun):=20=ED=9B=84=EB=B3=B4=20?= =?UTF-8?q?=EC=B1=84=ED=83=9D=EA=B3=BC=20=EC=82=AC=EC=97=85=EC=9E=A5=20?= =?UTF-8?q?=EA=B2=A9=EB=A6=AC=20=EC=8B=9C=EB=82=98=EB=A6=AC=EC=98=A4?= =?UTF-8?q?=EB=A5=BC=20=EA=B2=80=EC=A6=9D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../server/PostgreSqlMigrationTests.java | 35 +++ .../server/airun/AiRunApiIntegrationTest.java | 205 ++++++++++++++++++ 2 files changed, 240 insertions(+) diff --git a/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java b/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java index 707563e9..881fccca 100644 --- a/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java +++ b/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java @@ -86,6 +86,9 @@ private void assertSchemaContract(Connection connection) throws SQLException { "ai_attempt", "ai_question", "ai_candidate", + "ai_candidate_decision_batch", + "ai_candidate_decision", + "ai_candidate_decision_task", "workflow_case", "worker_link", "worker_response", @@ -209,6 +212,22 @@ private void assertSchemaContract(Connection connection) throws SQLException { .containsEntry("ai_attempt_id", new ColumnSpec("uuid", false)) .containsEntry("worker_id", new ColumnSpec("uuid", false)) .containsEntry("confidence", new ColumnSpec("numeric", false)); + assertThat(columnSpecs(connection, "ai_candidate_decision_batch")) + .containsEntry("decision_batch_id", new ColumnSpec("uuid", false)) + .containsEntry("ai_run_id", new ColumnSpec("uuid", false)) + .containsEntry("company_id", new ColumnSpec("uuid", false)) + .containsEntry("case_id", new ColumnSpec("uuid", true)) + .containsEntry("resulting_run_version", new ColumnSpec("int8", true)); + assertThat(columnSpecs(connection, "ai_candidate_decision")) + .containsEntry("decision_id", new ColumnSpec("uuid", false)) + .containsEntry("ai_candidate_id", new ColumnSpec("uuid", false)) + .containsEntry("company_id", new ColumnSpec("uuid", false)) + .containsEntry("action", new ColumnSpec("varchar", false)); + assertThat(columnSpecs(connection, "ai_candidate_decision_task")) + .containsEntry("decision_id", new ColumnSpec("uuid", false)) + .containsEntry("task_id", new ColumnSpec("uuid", false)) + .containsEntry("company_id", new ColumnSpec("uuid", false)) + .containsEntry("sequence_no", new ColumnSpec("int4", false)); assertThat(columnSpecs(connection, "worker_link")) .containsEntry("worker_link_id", new ColumnSpec("uuid", false)) .containsEntry("task_id", new ColumnSpec("uuid", false)) @@ -270,7 +289,17 @@ private void assertSchemaContract(Connection connection) throws SQLException { "pk_ai_question", "fk_ai_question_attempt_company", "pk_ai_candidate", + "uq_ai_candidate_id_company", "fk_ai_candidate_worker_company", + "pk_ai_candidate_decision_batch", + "uq_ai_candidate_decision_batch_idempotency", + "fk_ai_candidate_decision_batch_run_company", + "fk_ai_candidate_decision_batch_case_company", + "pk_ai_candidate_decision", + "uq_ai_candidate_decision_candidate", + "fk_ai_candidate_decision_candidate_company", + "pk_ai_candidate_decision_task", + "fk_ai_candidate_decision_task_task_company", "pk_workflow_case", "uq_workflow_case_id_company", "fk_workflow_case_worker_company", @@ -308,6 +337,9 @@ private void assertSchemaContract(Connection connection) throws SQLException { "idx_ai_attempt_run", "idx_ai_question_run", "idx_ai_candidate_run", + "idx_ai_candidate_decision_batch_run", + "idx_ai_candidate_decision_run", + "idx_ai_candidate_decision_task_task", "idx_workflow_case_company_updated", "idx_workflow_case_company_worker", "idx_worker_response_upload_company", @@ -337,6 +369,9 @@ private void assertSchemaContract(Connection connection) throws SQLException { "pl_ai_attempt_tenant_isolation", "pl_ai_question_tenant_isolation", "pl_ai_candidate_tenant_isolation", + "pl_ai_candidate_decision_batch_tenant_isolation", + "pl_ai_candidate_decision_tenant_isolation", + "pl_ai_candidate_decision_task_tenant_isolation", "pl_workflow_case_tenant_isolation", "pl_worker_link_tenant_isolation", "pl_worker_response_tenant_isolation", diff --git a/src/test/java/com/fowoco/server/airun/AiRunApiIntegrationTest.java b/src/test/java/com/fowoco/server/airun/AiRunApiIntegrationTest.java index 43d44a55..eb12ecb7 100644 --- a/src/test/java/com/fowoco/server/airun/AiRunApiIntegrationTest.java +++ b/src/test/java/com/fowoco/server/airun/AiRunApiIntegrationTest.java @@ -72,6 +72,9 @@ void resetAndSeed() { return scriptedResponse(request, runtimeCalls.incrementAndGet()); }); + jdbcTemplate.update("DELETE FROM ai_candidate_decision_task"); + jdbcTemplate.update("DELETE FROM ai_candidate_decision"); + jdbcTemplate.update("DELETE FROM ai_candidate_decision_batch"); jdbcTemplate.update("DELETE FROM ai_candidate"); jdbcTemplate.update("DELETE FROM ai_question"); jdbcTemplate.update("DELETE FROM ai_attempt"); @@ -85,6 +88,7 @@ void resetAndSeed() { jdbcTemplate.update("DELETE FROM task_transition_history"); jdbcTemplate.update("DELETE FROM task_checklist_item"); jdbcTemplate.update("DELETE FROM task"); + jdbcTemplate.update("DELETE FROM workflow_case"); jdbcTemplate.update("DELETE FROM worker_document"); jdbcTemplate.update("DELETE FROM worker"); jdbcTemplate.update("DELETE FROM refresh_token"); @@ -194,6 +198,164 @@ void idempotencyAndCompanyIsolationAreEnforced() throws Exception { assertThat(conflict.statusCode()).isEqualTo(409); } + @Test + void acceptedCandidateCreatesOneCaseAndThreeTasksIdempotently() throws Exception { + reset(runtimeClient); + runtimeCalls.set(0); + when(runtimeClient.analyze(any(), any())).thenAnswer(invocation -> { + AiAnalysisRequest request = invocation.getArgument(0); + return directReviewResponse(request, runtimeCalls.incrementAndGet()); + }); + String tokenA = login(HR_A_EMAIL); + String tokenB = login(HR_B_EMAIL); + HttpResponse reviewed = post( + "/api/v1/ai-runs", + """ + {"instruction":"응웬반A 체류연장 준비해줘"} + """, + tokenA, + "airun-decision-run" + ); + assertThat(reviewed.statusCode()).isEqualTo(202); + assertThat(JsonPath.read(reviewed.body(), "$.analysis_outcome")) + .isEqualTo("REVIEW_REQUIRED"); + assertThat(runtimeCalls).hasValue(2); + UUID aiRunId = UUID.fromString(JsonPath.read(reviewed.body(), "$.ai_run_id")); + UUID candidateId = UUID.fromString(JsonPath.read( + reviewed.body(), + "$.candidates[0].candidate_id" + )); + long expectedVersion = JsonPath.read(reviewed.body(), "$.version").longValue(); + assertThat(JsonPath.read(reviewed.body(), "$.detected_intent")) + .isEqualTo("EXPIRY_RENEWAL"); + + String decisionBody = """ + { + "expected_run_version":%d, + "decisions":[{"candidate_id":"%s","action":"ACCEPT"}] + } + """.formatted(expectedVersion, candidateId); + HttpResponse first = post( + "/api/v1/ai-runs/" + aiRunId + "/candidate-decisions", + decisionBody, + tokenA, + "candidate-decision-001" + ); + HttpResponse repeated = post( + "/api/v1/ai-runs/" + aiRunId + "/candidate-decisions", + decisionBody, + tokenA, + "candidate-decision-001" + ); + + assertThat(first.statusCode()).isEqualTo(200); + assertThat(repeated.statusCode()).isEqualTo(200); + assertThat(JsonPath.read(repeated.body(), "$.decision_batch_id")) + .isEqualTo(JsonPath.read(first.body(), "$.decision_batch_id")); + UUID caseId = UUID.fromString(JsonPath.read(first.body(), "$.case_id")); + assertThat(JsonPath.>read(first.body(), "$.task_ids")).hasSize(3); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM workflow_case WHERE case_id = ? AND company_id = ?", + Integer.class, + caseId, + COMPANY_A + )).isEqualTo(1); + assertThat(jdbcTemplate.queryForList( + """ + SELECT task_type + FROM task + WHERE case_id = ? AND company_id = ? + ORDER BY CASE task_type + WHEN 'RECONTRACT' THEN 1 + WHEN 'STAY_PERIOD_EXTENSION' THEN 2 + ELSE 3 + END + """, + String.class, + caseId, + COMPANY_A + )).containsExactly( + "RECONTRACT", + "STAY_PERIOD_EXTENSION", + "EMPLOYMENT_PERIOD_EXTENSION" + ); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM task WHERE case_id = ? AND source = 'AI_CANDIDATE'", + Integer.class, + caseId + )).isEqualTo(3); + assertThat(jdbcTemplate.queryForList( + "SELECT due_date FROM task WHERE case_id = ? ORDER BY task_id", + LocalDate.class, + caseId + )).containsOnly(LocalDate.of(2026, 9, 30)); + String snapshot = jdbcTemplate.queryForObject( + "SELECT workflow_snapshot_json FROM workflow_case WHERE case_id = ?", + String.class, + caseId + ); + assertThat(JsonPath.>read(snapshot, "$.steps[*].order")) + .containsExactly(1, 2, 3); + assertThat(JsonPath.>read(snapshot, "$.steps[*].task_type")) + .containsExactly( + "RECONTRACT", + "STAY_PERIOD_EXTENSION", + "EMPLOYMENT_PERIOD_EXTENSION" + ); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM ai_candidate_decision_batch WHERE ai_run_id = ?", + Integer.class, + aiRunId + )).isEqualTo(1); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM ai_candidate_decision_task", + Integer.class + )).isEqualTo(3); + assertThat(jdbcTemplate.queryForList( + "SELECT action FROM audit_event WHERE target_id = ? ORDER BY created_at", + String.class, + aiRunId + )).containsExactly( + "AI_RUN_CREATED", + "AI_RUN_CANDIDATES_DECIDED" + ); + + HttpResponse reusedKeyWithDifferentPayload = post( + "/api/v1/ai-runs/" + aiRunId + "/candidate-decisions", + """ + { + "expected_run_version":%d, + "decisions":[{"candidate_id":"%s","action":"DISCARD"}] + } + """.formatted(expectedVersion, candidateId), + tokenA, + "candidate-decision-001" + ); + assertThat(reusedKeyWithDifferentPayload.statusCode()).isEqualTo(409); + + long decidedVersion = JsonPath.read(first.body(), "$.run_version").longValue(); + HttpResponse alreadyDecided = post( + "/api/v1/ai-runs/" + aiRunId + "/candidate-decisions", + """ + { + "expected_run_version":%d, + "decisions":[{"candidate_id":"%s","action":"ACCEPT"}] + } + """.formatted(decidedVersion, candidateId), + tokenA, + "candidate-decision-002" + ); + assertThat(alreadyDecided.statusCode()).isEqualTo(409); + + HttpResponse otherCompany = post( + "/api/v1/ai-runs/" + aiRunId + "/candidate-decisions", + decisionBody, + tokenB, + "candidate-decision-other-company" + ); + assertThat(otherCompany.statusCode()).isEqualTo(404); + } + private AiAnalysisResponse scriptedResponse(AiAnalysisRequest request, int call) { if (call == 1) { return new AiAnalysisResponse( @@ -247,6 +409,49 @@ private AiAnalysisResponse scriptedResponse(AiAnalysisRequest request, int call) ); } + private AiAnalysisResponse directReviewResponse(AiAnalysisRequest request, int call) { + if (call == 1) { + return new AiAnalysisResponse( + request.requestId(), + AiAnalysisOutcome.CONTEXT_REQUIRED, + new AiContextRequirement( + "EXPIRY_RENEWAL", + new BigDecimal("0.96"), + "응웬반A", + Map.of(), + List.of("worker_id", "stay_expiry_date") + ), + List.of(), + List.of(), + List.of(), + versions(), + 1, + 30 + ); + } + return new AiAnalysisResponse( + request.requestId(), + AiAnalysisOutcome.REVIEW_REQUIRED, + null, + List.of(), + List.of(new AiCandidate( + "candidate-1", + WORKER_A, + "WF-STY-001", + Map.of( + "worker_id", WORKER_A.toString(), + "stay_expiry_date", "2026-09-30" + ), + List.of(), + new BigDecimal("0.93") + )), + List.of(), + versions(), + 1, + 25 + ); + } + private AiRuntimeVersions versions() { return new AiRuntimeVersions( "agent-demo-1", From 07bfac7bbc486be9ac4d80cdbc4595c460298e51 Mon Sep 17 00:00:00 2001 From: hywznn Date: Thu, 6 Aug 2026 16:35:16 +0900 Subject: [PATCH 3/3] =?UTF-8?q?docs(ai):=20Intent=EC=99=80=20Workflow=20ID?= =?UTF-8?q?=20=EC=98=88=EC=8B=9C=EB=A5=BC=20=EA=B5=90=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/ai-runtime-contract.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/ai-runtime-contract.md b/docs/ai-runtime-contract.md index ff79639a..8b02daa7 100644 --- a/docs/ai-runtime-contract.md +++ b/docs/ai-runtime-contract.md @@ -149,7 +149,7 @@ Agent가 문서 작성에 요구한 값은 `***`, `OOO`로 바꾸지 않고 원 { "candidateRef": "candidate-1", "workerRef": "30000000-0000-0000-0000-000000000001", - "workflowId": "EXPIRY_RENEWAL", + "workflowId": "WF-STY-001", "extractedSlots": { "stay_expiry_date": "2026-12-31" },