refactor: support multi-task Flowable approval semantics

This commit is contained in:
selfrelease
2026-07-18 08:43:45 +08:00
parent c5d80aba43
commit e46aa00e80
10 changed files with 333 additions and 20 deletions
+1
View File
@@ -28,6 +28,7 @@ export GRADLE_USER_HOME=/tmp/aioa-gradle-home
- UUIDv7、写入幂等、租户隔离和乐观锁冲突处理
- 请假提交、撤回、状态时间线与事务内审计
- Flowable 7.2 部门主管审批 BPMN、待办、批准和驳回
- 复杂流程安全边界:中间串行/并行任务不提前结束业务申请
Flowable 开发环境自动维护 `flowable` Schema。生产环境必须设置 `FLOWABLE_SCHEMA_UPDATE=false`,并通过受控数据库变更流程管理 Flowable 表结构。
@@ -69,7 +69,6 @@ class ApprovalTaskService(
}
val operation = if (approved) "LEAVE_APPROVE" else "LEAVE_REJECT"
val action = if (approved) "LEAVE_REQUEST_APPROVED" else "LEAVE_REQUEST_REJECTED"
val targetStatus = if (approved) LeaveStatus.APPROVED else LeaveStatus.REJECTED
val fingerprint = fingerprint(operation, taskId, expectedVersion, comment)
val claim = leaveRequestRepository.claimOperation(
@@ -86,18 +85,37 @@ class ApprovalTaskService(
throw ApiException(HttpStatus.CONFLICT, "APPROVAL_TASK_COMPLETED", "审批任务已处理")
}
val transitioned = leaveRequestRepository.transitionForApprover(
tenantId = actor.tenantId,
actorId = actor.id,
id = request.id,
expectedVersion = expectedVersion,
toStatus = targetStatus,
eventId = UuidV7.generate(),
eventType = action,
traceId = MDC.get(TraceIdFilter.MDC_TRACE_ID) ?: "unknown",
) ?: transitionFailure(actor, request.id, expectedVersion)
val completion = workflowGateway.completeTask(taskId, approved, comment)
val traceId = MDC.get(TraceIdFilter.MDC_TRACE_ID) ?: "unknown"
val action = when {
completion.processEnded && approved -> "LEAVE_REQUEST_APPROVED"
completion.processEnded -> "LEAVE_REQUEST_REJECTED"
approved -> "LEAVE_APPROVAL_TASK_APPROVED"
else -> "LEAVE_APPROVAL_TASK_REJECTED"
}
val result = if (completion.processEnded) {
leaveRequestRepository.transitionForApprover(
tenantId = actor.tenantId,
actorId = actor.id,
id = request.id,
expectedVersion = expectedVersion,
toStatus = targetStatus,
eventId = UuidV7.generate(),
eventType = action,
traceId = traceId,
) ?: transitionFailure(actor, request.id, expectedVersion)
} else {
leaveRequestRepository.appendWorkflowEvent(
eventId = UuidV7.generate(),
tenantId = actor.tenantId,
leaveRequestId = request.id,
actorId = actor.id,
eventType = action,
traceId = traceId,
)
request
}
workflowGateway.completeTask(taskId, approved, comment)
auditService.recordSuccess(
actor = actor,
action = action,
@@ -106,12 +124,14 @@ class ApprovalTaskService(
idempotencyKey = idempotencyKey,
details = mapOf(
"taskId" to taskId,
"toStatus" to targetStatus.name,
"version" to transitioned.version,
"decision" to if (approved) "APPROVED" else "REJECTED",
"processEnded" to completion.processEnded,
"toStatus" to result.status.name,
"version" to result.version,
"comment" to comment?.trim(),
),
)
return transitioned
return result
}
private fun transitionFailure(actor: CurrentUser, id: java.util.UUID, expectedVersion: Long): Nothing {
@@ -105,6 +105,15 @@ interface LeaveRequestRepository {
eventType: String,
traceId: String,
): LeaveRequest?
fun appendWorkflowEvent(
eventId: UUID,
tenantId: UUID,
leaveRequestId: UUID,
actorId: UUID,
eventType: String,
traceId: String,
)
}
data class IdempotentDraft(
@@ -291,6 +291,30 @@ class JooqLeaveRequestRepository(
return transitioned
}
override fun appendWorkflowEvent(
eventId: UUID,
tenantId: UUID,
leaveRequestId: UUID,
actorId: UUID,
eventType: String,
traceId: String,
) {
dsl.execute(
"""
INSERT INTO business.leave_request_event (
id, tenant_id, leave_request_id, actor_id, event_type,
from_status, to_status, trace_id
) VALUES (?, ?, ?, ?, ?, 'PENDING', 'PENDING', ?)
""".trimIndent(),
eventId,
tenantId,
leaveRequestId,
actorId,
eventType,
traceId,
)
}
private fun map(record: Record): LeaveRequest = LeaveRequest(
id = record.get("id", UUID::class.java)!!,
tenantId = record.get("tenant_id", UUID::class.java)!!,
@@ -15,7 +15,7 @@ interface LeaveWorkflowGateway {
fun resolveTask(taskId: String): WorkflowTask?
fun completeTask(taskId: String, approved: Boolean, comment: String?)
fun completeTask(taskId: String, approved: Boolean, comment: String?): TaskCompletion
fun cancelProcess(processInstanceId: String, reason: String)
}
@@ -34,3 +34,8 @@ data class WorkflowTask(
val createdAt: Instant,
val completed: Boolean,
)
data class TaskCompletion(
val processInstanceId: String,
val processEnded: Boolean,
)
@@ -3,6 +3,7 @@ package com.all8ai.aioa.workflow.infrastructure
import com.all8ai.aioa.workflow.domain.LeaveWorkflowGateway
import com.all8ai.aioa.workflow.domain.StartedProcess
import com.all8ai.aioa.workflow.domain.WorkflowTask
import com.all8ai.aioa.workflow.domain.TaskCompletion
import org.flowable.engine.HistoryService
import org.flowable.engine.RuntimeService
import org.flowable.engine.TaskService
@@ -70,13 +71,17 @@ class FlowableLeaveWorkflowGateway(
)
}
override fun completeTask(taskId: String, approved: Boolean, comment: String?) {
override fun completeTask(taskId: String, approved: Boolean, comment: String?): TaskCompletion {
val task = taskService.createTaskQuery().taskId(taskId).singleResult()
?: error("Active workflow task not found: $taskId")
if (!comment.isNullOrBlank()) {
val task = taskService.createTaskQuery().taskId(taskId).singleResult()
?: error("Active workflow task not found: $taskId")
taskService.addComment(taskId, task.processInstanceId, comment.trim())
}
taskService.complete(taskId, mapOf("approved" to approved))
val processEnded = runtimeService.createProcessInstanceQuery()
.processInstanceId(task.processInstanceId)
.singleResult() == null
return TaskCompletion(task.processInstanceId, processEnded)
}
override fun cancelProcess(processInstanceId: String, reason: String) {
@@ -0,0 +1,199 @@
package com.all8ai.aioa.approval.application
import com.all8ai.aioa.approval.domain.DraftContent
import com.all8ai.aioa.approval.domain.IdempotentDraft
import com.all8ai.aioa.approval.domain.LeaveRequest
import com.all8ai.aioa.approval.domain.LeaveRequestEvent
import com.all8ai.aioa.approval.domain.LeaveRequestRepository
import com.all8ai.aioa.approval.domain.LeaveStatus
import com.all8ai.aioa.approval.domain.LeaveType
import com.all8ai.aioa.approval.domain.OperationClaim
import com.all8ai.aioa.audit.application.AuditService
import com.all8ai.aioa.audit.domain.AuditEvent
import com.all8ai.aioa.audit.domain.AuditEventRepository
import com.all8ai.aioa.identity.domain.CurrentUser
import com.all8ai.aioa.workflow.domain.LeaveWorkflowGateway
import com.all8ai.aioa.workflow.domain.StartedProcess
import com.all8ai.aioa.workflow.domain.TaskCompletion
import com.all8ai.aioa.workflow.domain.WorkflowTask
import org.assertj.core.api.Assertions.assertThat
import org.junit.jupiter.api.Test
import java.time.Instant
import java.util.UUID
class ApprovalTaskServiceTest {
private val tenantId = UUID.randomUUID()
private val applicantId = UUID.randomUUID()
private val managerId = UUID.randomUUID()
private val leaveId = UUID.randomUUID()
private val actor = CurrentUser(
managerId, tenantId, "manager", "主管", null, null, null,
setOf("employee", "department_manager"),
)
private val pending = LeaveRequest(
id = leaveId,
tenantId = tenantId,
applicantId = applicantId,
type = LeaveType.ANNUAL,
startsAt = Instant.parse("2026-07-23T01:00:00Z"),
endsAt = Instant.parse("2026-07-23T09:00:00Z"),
reason = "测试",
status = LeaveStatus.PENDING,
createdAt = Instant.parse("2026-07-18T00:00:00Z"),
updatedAt = Instant.parse("2026-07-18T00:00:00Z"),
version = 1,
processInstanceId = "process-1",
processDefinitionId = "serialApproval:1:test",
)
private val workflowTask = WorkflowTask(
id = "task-1",
name = "一级审批",
processInstanceId = "process-1",
leaveRequestId = leaveId,
assigneeId = managerId,
createdAt = Instant.parse("2026-07-18T00:00:00Z"),
completed = false,
)
@Test
fun `keeps business pending when a serial or parallel task is not the final task`() {
val fixture = fixture(processEnded = false)
val result = fixture.service.approve(actor, "task-1", "approval-task-key-0001", 1, "一级同意")
assertThat(result.status).isEqualTo(LeaveStatus.PENDING)
assertThat(result.version).isEqualTo(1)
assertThat(fixture.repository.intermediateEventType).isEqualTo("LEAVE_APPROVAL_TASK_APPROVED")
assertThat(fixture.repository.finalTransitionCalled).isFalse()
assertThat(fixture.audits.single().action).isEqualTo("LEAVE_APPROVAL_TASK_APPROVED")
}
@Test
fun `sets final status only when the workflow process has ended`() {
val fixture = fixture(processEnded = true)
val result = fixture.service.approve(actor, "task-1", "approval-task-key-0002", 1, "最终同意")
assertThat(result.status).isEqualTo(LeaveStatus.APPROVED)
assertThat(result.version).isEqualTo(2)
assertThat(fixture.repository.finalTransitionCalled).isTrue()
assertThat(fixture.repository.intermediateEventType).isNull()
assertThat(fixture.audits.single().action).isEqualTo("LEAVE_REQUEST_APPROVED")
}
private fun fixture(processEnded: Boolean): Fixture {
val repository = FakeRepository(pending)
val audits = mutableListOf<AuditEvent>()
val workflow = FakeWorkflowGateway(workflowTask, processEnded)
return Fixture(
ApprovalTaskService(workflow, repository, AuditService(AuditEventRepository { audits += it })),
repository,
audits,
)
}
private data class Fixture(
val service: ApprovalTaskService,
val repository: FakeRepository,
val audits: List<AuditEvent>,
)
private class FakeWorkflowGateway(
private val task: WorkflowTask,
private val processEnded: Boolean,
) : LeaveWorkflowGateway {
override fun startLeaveApproval(
tenantId: UUID,
leaveRequestId: UUID,
applicantId: UUID,
approverId: UUID,
): StartedProcess = error("Not used")
override fun listAssignedTasks(assigneeId: UUID): List<WorkflowTask> = listOf(task)
override fun resolveTask(taskId: String): WorkflowTask? = task.takeIf { it.id == taskId }
override fun completeTask(taskId: String, approved: Boolean, comment: String?) =
TaskCompletion(task.processInstanceId, processEnded)
override fun cancelProcess(processInstanceId: String, reason: String) = Unit
}
private class FakeRepository(private var request: LeaveRequest) : LeaveRequestRepository {
var intermediateEventType: String? = null
var finalTransitionCalled = false
override fun findById(tenantId: UUID, id: UUID): LeaveRequest? =
request.takeIf { it.tenantId == tenantId && it.id == id }
override fun claimOperation(
id: UUID,
tenantId: UUID,
actorId: UUID,
operation: String,
idempotencyKey: String,
fingerprint: String,
resourceId: UUID,
) = OperationClaim(true, fingerprint, resourceId)
override fun transitionForApprover(
tenantId: UUID,
actorId: UUID,
id: UUID,
expectedVersion: Long,
toStatus: LeaveStatus,
eventId: UUID,
eventType: String,
traceId: String,
): LeaveRequest? {
finalTransitionCalled = true
if (request.version != expectedVersion || request.status != LeaveStatus.PENDING) return null
request = request.copy(status = toStatus, version = request.version + 1)
return request
}
override fun appendWorkflowEvent(
eventId: UUID,
tenantId: UUID,
leaveRequestId: UUID,
actorId: UUID,
eventType: String,
traceId: String,
) {
intermediateEventType = eventType
}
override fun createDraft(
id: UUID,
tenantId: UUID,
applicantId: UUID,
idempotencyKey: String,
fingerprint: String,
content: DraftContent,
): IdempotentDraft = error("Not used")
override fun updateDraft(
tenantId: UUID,
applicantId: UUID,
id: UUID,
expectedVersion: Long,
content: DraftContent,
): LeaveRequest? = error("Not used")
override fun findOwn(tenantId: UUID, applicantId: UUID, id: UUID): LeaveRequest? = null
override fun listOwn(tenantId: UUID, applicantId: UUID, limit: Int): List<LeaveRequest> = emptyList()
override fun transitionOwn(
tenantId: UUID,
actorId: UUID,
id: UUID,
expectedVersion: Long,
fromStatus: LeaveStatus,
toStatus: LeaveStatus,
eventId: UUID,
eventType: String,
traceId: String,
): LeaveRequest? = error("Not used")
override fun listTimeline(tenantId: UUID, leaveRequestId: UUID): List<LeaveRequestEvent> = emptyList()
override fun attachWorkflow(
tenantId: UUID,
leaveRequestId: UUID,
processInstanceId: String,
processDefinitionId: String,
) = Unit
}
}
@@ -15,6 +15,7 @@ import com.all8ai.aioa.approval.domain.ApprovalRoutingRepository
import com.all8ai.aioa.workflow.domain.LeaveWorkflowGateway
import com.all8ai.aioa.workflow.domain.StartedProcess
import com.all8ai.aioa.workflow.domain.WorkflowTask
import com.all8ai.aioa.workflow.domain.TaskCompletion
import com.all8ai.aioa.identity.domain.CurrentUser
import com.all8ai.aioa.shared.web.ApiException
import org.assertj.core.api.Assertions.assertThat
@@ -144,7 +145,8 @@ class LeaveRequestServiceTest {
override fun listAssignedTasks(assigneeId: UUID): List<WorkflowTask> = emptyList()
override fun resolveTask(taskId: String): WorkflowTask? = null
override fun completeTask(taskId: String, approved: Boolean, comment: String?) = Unit
override fun completeTask(taskId: String, approved: Boolean, comment: String?) =
TaskCompletion("process-test", processEnded = true)
override fun cancelProcess(processInstanceId: String, reason: String) = Unit
}
@@ -266,5 +268,14 @@ class LeaveRequestServiceTest {
eventType: String,
traceId: String,
): LeaveRequest? = null
override fun appendWorkflowEvent(
eventId: UUID,
tenantId: UUID,
leaveRequestId: UUID,
actorId: UUID,
eventType: String,
traceId: String,
) = Unit
}
}
+1
View File
@@ -27,6 +27,7 @@
- [x] 提交、撤回、状态时间线和同步审计
- [x] Flowable 部门主管审批流程发布与执行
- [x] 发起、主管待办、批准、驳回、撤回和时间线
- [x] 多任务流程的中间审批事件与最终流程结束判定
- 附件、通知、弱网恢复和幂等处理
## M3AI 最小闭环
+38
View File
@@ -0,0 +1,38 @@
# Flowable 复杂流程能力与业务边界
## 支持能力
Flowable 可以承载:
- 串行审批:多个用户任务依次执行。
- 并行审批:Parallel Gateway 同时创建多个任务,全部完成后汇聚。
- 条件并行:Inclusive Gateway 根据条件创建一组分支。
- 条件路由:Exclusive Gateway 根据金额、类型、风险等选择唯一路径。
- 会签和或签:并行或串行 Multi-instance User Task,使用完成条件控制通过比例。
- 子流程和复用流程:Embedded Subprocess 与 Call Activity。
- 超时、提醒和升级:Timer Boundary Event、定时作业和升级路径。
- 决策表:DMN 管理金额、岗位、地区和风险等审批规则。
- 撤回、终止和人工干预:运行时实例管理与完整历史记录。
- 流程版本:新实例使用新定义,运行中的实例继续引用原版本。
## AIOA 状态同步规则
Flowable 负责任务和流程路径,`business.leave_request` 仍是请假业务的权威数据。
- 用户完成中间审批任务时,申请保持 `PENDING`
- 中间任务只写 `LEAVE_APPROVAL_TASK_APPROVED``LEAVE_APPROVAL_TASK_REJECTED` 时间线事件。
- 只有整个流程实例结束后,业务状态才变为 `APPROVED``REJECTED`
- 申请人撤回时,业务状态和 Flowable 运行实例在同一事务中结束。
- 每个任务动作均重新校验实际 assignee、租户、版本和幂等键。
- 申请人不得审批自己的申请。
## 建模约束
- BPMN 只保存 `tenantId``businessId``applicantId`、审批人标识和路由所需的少量变量。
- 表单内容、附件和权威状态不能长期放入流程变量。
- 并行网关必须成对建模,避免产生无法汇聚的执行路径。
- 会签必须明确完成条件,例如全部通过、超过半数或任一通过。
- 驳回路径必须明确是结束流程、退回上一步还是返回申请人修改。
- 流程发布前必须覆盖通过、驳回、超时、撤回和无审批人等路径测试。
当前仓库部署的是单主管审批定义,但后端状态同步已经按多任务流程设计,不会在第一个串行或并行任务完成时提前结束申请。