From 6a4a115260b55d11d6d7b6e3c1602fd3ca001a33 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=BB=84=E8=87=AA=E8=BE=BE?= <980324341@qq.com> Date: Sun, 30 Aug 2026 20:39:31 +0800 Subject: [PATCH] =?UTF-8?q?task-75:=20=E8=99=9A=E6=8B=9F=E7=BA=BF=E7=A8=8B?= =?UTF-8?q?=E4=BB=BB=E5=8A=A1=E7=AD=89=E5=BE=85=E9=98=9F=E5=88=97=E4=B8=8A?= =?UTF-8?q?=E9=99=90=E4=B8=8E=E6=8B=92=E7=BB=9D/=E5=BB=B6=E8=BF=9F?= =?UTF-8?q?=E6=8C=87=E6=A0=87?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../aiimage/config/CozeTaskQueueGate.java | 114 +++++++++ .../aiimage/config/TaskFileJobConfig.java | 9 +- .../aiimage/config/CozeTaskQueueGateTest.java | 218 ++++++++++++++++++ .../config/ThreadPoolIsolationConfigTest.java | 2 +- 4 files changed, 340 insertions(+), 3 deletions(-) create mode 100644 backend-java/src/main/java/com/nanri/aiimage/config/CozeTaskQueueGate.java create mode 100644 backend-java/src/test/java/com/nanri/aiimage/config/CozeTaskQueueGateTest.java diff --git a/backend-java/src/main/java/com/nanri/aiimage/config/CozeTaskQueueGate.java b/backend-java/src/main/java/com/nanri/aiimage/config/CozeTaskQueueGate.java new file mode 100644 index 00000000..82bb7373 --- /dev/null +++ b/backend-java/src/main/java/com/nanri/aiimage/config/CozeTaskQueueGate.java @@ -0,0 +1,114 @@ +package com.nanri.aiimage.config; + +import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.Timer; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.ObjectProvider; +import org.springframework.core.task.TaskExecutor; +import org.springframework.core.task.TaskRejectedException; + +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +/** + * Task 75:虚拟线程任务排队闸门。Coze 执行池的信号量只限制"正在执行"的 + * 并发度,提交侧仍会在虚拟线程里无限排队。此闸门在提交时统计"已受理未启动" + * 的等待数,达到上限立即拒绝并记录指标,防止等待队列无界堆积: + * + */ +@Slf4j +public class CozeTaskQueueGate implements TaskExecutor { + + private final TaskExecutor delegate; + private final int maxWaiting; + private final ObjectProvider meterRegistryProvider; + private final AtomicInteger waiting = new AtomicInteger(); + + public CozeTaskQueueGate(TaskExecutor delegate, int maxWaiting, + ObjectProvider meterRegistryProvider) { + this.delegate = delegate; + this.maxWaiting = Math.max(1, maxWaiting); + this.meterRegistryProvider = meterRegistryProvider; + } + + public int waiting() { + return waiting.get(); + } + + @Override + public void execute(Runnable command) { + if (command == null) { + recordRejected("invalid-input"); + throw new IllegalArgumentException("coze 任务不能为 null"); + } + if (waiting.get() >= maxWaiting) { + recordRejected("queue-full"); + log.warn("[coze-task][gate] waiting queue full, reject submit waiting={} limit={}", + waiting.get(), maxWaiting); + throw new TaskRejectedException("coze 等待队列已满,limit=" + maxWaiting + + ", waiting=" + waiting.get()); + } + waiting.incrementAndGet(); + long queuedAt = System.nanoTime(); + try { + delegate.execute(() -> { + long waitNanos = System.nanoTime() - queuedAt; + try { + run(command); + } finally { + waiting.decrementAndGet(); + recordQueueWait(waitNanos); + } + }); + } catch (RuntimeException ex) { + waiting.decrementAndGet(); + recordQueueWait(System.nanoTime() - queuedAt); + recordRejected("delegate-rejected"); + log.warn("[coze-task][gate] delegate rejected submit waiting={} limit={} msg={}", + waiting.get(), maxWaiting, ex.getMessage(), ex); + throw ex; + } + } + + private void run(Runnable command) { + long startedAt = System.nanoTime(); + try { + command.run(); + } finally { + recordExecution(System.nanoTime() - startedAt); + } + } + + private void recordRejected(String reason) { + MeterRegistry registry = meterRegistry(); + if (registry != null) { + registry.counter("aiimage.coze-task.submit.rejected.total", "reason", reason).increment(); + } + } + + private void recordQueueWait(long waitNanos) { + MeterRegistry registry = meterRegistry(); + if (registry != null && waitNanos >= 0L) { + Timer.builder("aiimage.coze-task.queue.wait.duration") + .register(registry) + .record(waitNanos, TimeUnit.NANOSECONDS); + } + } + + private void recordExecution(long durationNanos) { + MeterRegistry registry = meterRegistry(); + if (registry != null && durationNanos >= 0L) { + Timer.builder("aiimage.coze-task.execution.duration") + .register(registry) + .record(durationNanos, TimeUnit.NANOSECONDS); + } + } + + private MeterRegistry meterRegistry() { + return meterRegistryProvider == null ? null : meterRegistryProvider.getIfAvailable(); + } +} diff --git a/backend-java/src/main/java/com/nanri/aiimage/config/TaskFileJobConfig.java b/backend-java/src/main/java/com/nanri/aiimage/config/TaskFileJobConfig.java index f8736d7f..017e7189 100644 --- a/backend-java/src/main/java/com/nanri/aiimage/config/TaskFileJobConfig.java +++ b/backend-java/src/main/java/com/nanri/aiimage/config/TaskFileJobConfig.java @@ -1,5 +1,7 @@ package com.nanri.aiimage.config; +import io.micrometer.core.instrument.MeterRegistry; +import org.springframework.beans.factory.ObjectProvider; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -40,9 +42,11 @@ public class TaskFileJobConfig { @Bean("cozeTaskExecutor") public TaskExecutor cozeTaskExecutor( ExecutorService cozeVirtualThreadExecutor, - @Value("${aiimage.coze-task.max-concurrent:12}") int maxConcurrent) { + @Value("${aiimage.coze-task.max-concurrent:12}") int maxConcurrent, + @Value("${aiimage.coze-task.max-waiting:1000}") int maxWaiting, + ObjectProvider meterRegistryProvider) { Semaphore semaphore = new Semaphore(Math.max(1, maxConcurrent)); - return new ConcurrentTaskExecutor(command -> { + TaskExecutor semaphoreLimited = new ConcurrentTaskExecutor(command -> { if (command == null) { throw new IllegalArgumentException("coze 任务不能为 null"); } @@ -61,5 +65,6 @@ public class TaskFileJobConfig { } }); }); + return new CozeTaskQueueGate(semaphoreLimited, maxWaiting, meterRegistryProvider); } } diff --git a/backend-java/src/test/java/com/nanri/aiimage/config/CozeTaskQueueGateTest.java b/backend-java/src/test/java/com/nanri/aiimage/config/CozeTaskQueueGateTest.java new file mode 100644 index 00000000..4995a55d --- /dev/null +++ b/backend-java/src/test/java/com/nanri/aiimage/config/CozeTaskQueueGateTest.java @@ -0,0 +1,218 @@ +package com.nanri.aiimage.config; + +import io.micrometer.core.instrument.Counter; +import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.Timer; +import io.micrometer.core.instrument.simple.SimpleMeterRegistry; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.ObjectProvider; +import org.springframework.core.task.TaskExecutor; +import org.springframework.core.task.TaskRejectedException; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +/** + * Task 75:为虚拟线程任务增加等待队列上限与拒绝/延迟指标。 + * CozeTaskQueueGate 在信号量限流(并发上限)之外增加排队闸门: + * 已提交未启动(含等待信号量)的任务数量达到 max-waiting 上限时 + * 立即拒绝新提交并记录拒绝指标;每次提交/执行记录等待耗时与执行耗时 + * 指标;任务执行完成、失败或执行器拒绝后排队名额必须释放。 + */ +class CozeTaskQueueGateTest { + + private SimpleMeterRegistry registry; + private ObjectProvider meterRegistryProvider; + private List captured; + private TaskExecutor capturingDelegate; + + @BeforeEach + @SuppressWarnings("unchecked") + void setUp() { + registry = new SimpleMeterRegistry(); + meterRegistryProvider = mock(ObjectProvider.class); + when(meterRegistryProvider.getIfAvailable()).thenReturn(registry); + captured = new ArrayList<>(); + capturingDelegate = captured::add; + } + + private CozeTaskQueueGate gate(int maxWaiting) { + return new CozeTaskQueueGate(capturingDelegate, maxWaiting, meterRegistryProvider); + } + + private long rejectedCount(String reason) { + Counter counter = registry.find("aiimage.coze-task.submit.rejected.total") + .tag("reason", reason).counter(); + return counter == null ? 0 : (long) counter.count(); + } + + private long executionCount() { + Timer timer = registry.find("aiimage.coze-task.execution.duration").timer(); + return timer == null ? 0 : timer.count(); + } + + private long waitCount() { + Timer timer = registry.find("aiimage.coze-task.queue.wait.duration").timer(); + return timer == null ? 0 : timer.count(); + } + + @Test + void test_task_075_metrics_normal_default_path() { + // 默认路径:任务正常受理并执行,等待/执行指标各记一次,排队名额释放。 + CozeTaskQueueGate gate = gate(4); + AtomicInteger ran = new AtomicInteger(); + gate.execute(ran::incrementAndGet); + + assertEquals(1, captured.size(), "任务入队"); + assertEquals(1, gate.waiting(), "未启动任务计数为 1"); + captured.get(0).run(); + assertEquals(1, ran.get(), "任务已执行"); + assertEquals(0, gate.waiting(), "执行后名额释放"); + assertEquals(1, executionCount(), "执行耗时指标记录一次"); + assertEquals(1, waitCount(), "等待耗时指标记录一次"); + assertEquals(0, rejectedCount("queue-full")); + } + + @Test + void test_task_075_metrics_normal_multiple_items() { + // 批量场景:多个任务依次受理,执行顺序稳定不丢失,指标逐条记录。 + CozeTaskQueueGate gate = gate(8); + List order = new ArrayList<>(); + for (int i = 1; i <= 3; i++) { + int id = i; + gate.execute(() -> order.add(id)); + } + assertEquals(3, captured.size()); + for (Runnable runnable : captured) { + runnable.run(); + } + assertEquals(List.of(1, 2, 3), order, "执行顺序与提交顺序一致"); + assertEquals(0, gate.waiting()); + assertEquals(3, executionCount()); + assertEquals(3, waitCount()); + } + + @Test + void test_task_075_metrics_normal_repeated_operation_is_idempotent() { + // 幂等:同一任务对象重复提交各自独立执行,不合并、不丢失。 + CozeTaskQueueGate gate = gate(4); + AtomicInteger ran = new AtomicInteger(); + Runnable task = ran::incrementAndGet; + gate.execute(task); + gate.execute(task); + assertEquals(2, captured.size(), "同一任务重复提交各入队一次"); + captured.forEach(Runnable::run); + assertEquals(2, ran.get()); + assertEquals(2, executionCount()); + } + + @Test + void test_task_075_metrics_boundary_empty_input() { + // 空输入:null 任务直接拒绝,不计数、不产生指标。 + CozeTaskQueueGate gate = gate(4); + IllegalArgumentException ex = assertThrows(IllegalArgumentException.class, + () -> gate.execute(null)); + assertTrue(ex.getMessage().contains("不能为 null"), "可识别错误消息"); + assertTrue(captured.isEmpty()); + assertEquals(0, gate.waiting()); + assertEquals(0, rejectedCount("queue-full")); + assertEquals(0, executionCount()); + } + + @Test + void test_task_075_metrics_boundary_single_item() { + // 单元素:单任务直接受理执行,不依赖批量路径。 + CozeTaskQueueGate gate = gate(1); + AtomicInteger ran = new AtomicInteger(); + gate.execute(ran::incrementAndGet); + captured.get(0).run(); + assertEquals(1, ran.get()); + assertEquals(0, gate.waiting()); + assertEquals(1, executionCount()); + } + + @Test + void test_task_075_metrics_boundary_limit_and_overflow() { + // 上限/超限:等待队列达到上限后新提交被拒绝(记录拒绝指标), + // 排队名额释放后可再次受理,等待数不发生无界增长。 + CozeTaskQueueGate gate = gate(1); + AtomicInteger ran = new AtomicInteger(); + + gate.execute(ran::incrementAndGet); + assertEquals(1, gate.waiting()); + + TaskRejectedException firstReject = assertThrows(TaskRejectedException.class, + () -> gate.execute(ran::incrementAndGet), "队列满拒绝新提交"); + assertTrue(firstReject.getMessage().contains("limit=1"), "错误消息含上限值"); + TaskRejectedException secondReject = assertThrows(TaskRejectedException.class, + () -> gate.execute(ran::incrementAndGet)); + assertEquals(2, rejectedCount("queue-full"), "两次拒绝各记一次指标"); + assertEquals(1, captured.size(), "被拒绝的任务不入队"); + + captured.get(0).run(); + assertEquals(0, gate.waiting(), "执行后名额释放"); + gate.execute(ran::incrementAndGet); + assertEquals(2, captured.size(), "超限后仍可继续受理"); + assertEquals(1, gate.waiting(), "等待数不超过上限"); + captured.get(1).run(); + assertEquals(2, ran.get(), "再次受理的任务正常执行"); + assertEquals(0, gate.waiting(), "执行后名额再次释放"); + } + + @Test + void test_task_075_metrics_invalid_input_rejected() { + // 非法参数:等待上限 0/负值统一钳制到 1(首个受理、第二个拒绝), + // null 任务拒绝且不产生指标。 + CozeTaskQueueGate zeroLimit = gate(0); + zeroLimit.execute(() -> { }); + assertThrows(TaskRejectedException.class, () -> zeroLimit.execute(() -> { })); + assertEquals(1, captured.size(), "上限 0 回退到 1"); + + captured.clear(); + CozeTaskQueueGate negativeLimit = gate(-5); + negativeLimit.execute(() -> { }); + assertThrows(TaskRejectedException.class, () -> negativeLimit.execute(() -> { })); + assertEquals(1, captured.size(), "负值回退到 1"); + + assertThrows(IllegalArgumentException.class, () -> negativeLimit.execute(null)); + assertEquals(0, executionCount(), "非法参数不产生执行指标"); + } + + @Test + void test_task_075_metrics_dependency_failure_releases_resources() { + // 依赖失败:执行器拒绝时名额释放、记录拒绝指标,恢复后再次提交成功; + // 任务执行抛异常时名额同样释放,后续任务不受影响。 + CozeTaskQueueGate gate = new CozeTaskQueueGate(command -> { + throw new TaskRejectedException("executor full"); + }, 4, meterRegistryProvider); + assertThrows(TaskRejectedException.class, () -> gate.execute(() -> { })); + assertEquals(0, gate.waiting(), "拒绝后名额释放"); + assertEquals(1, rejectedCount("delegate-rejected"), "执行器拒绝单独计数"); + + CozeTaskQueueGate recovered = gate(4); + AtomicInteger ran = new AtomicInteger(); + recovered.execute(ran::incrementAndGet); + captured.get(0).run(); + assertEquals(1, ran.get(), "恢复后提交成功"); + + captured.clear(); + CozeTaskQueueGate failingTaskGate = gate(4); + failingTaskGate.execute(() -> { + throw new IllegalStateException("task boom"); + }); + assertThrows(IllegalStateException.class, () -> captured.get(0).run(), + "任务异常向上传播(生产环境由虚拟线程吞掉)"); + assertEquals(0, failingTaskGate.waiting(), "任务抛异常后名额释放"); + failingTaskGate.execute(ran::incrementAndGet); + captured.get(1).run(); + assertEquals(2, ran.get(), "异常后新任务可受理执行"); + } +} diff --git a/backend-java/src/test/java/com/nanri/aiimage/config/ThreadPoolIsolationConfigTest.java b/backend-java/src/test/java/com/nanri/aiimage/config/ThreadPoolIsolationConfigTest.java index 73ce8636..cec4cf79 100644 --- a/backend-java/src/test/java/com/nanri/aiimage/config/ThreadPoolIsolationConfigTest.java +++ b/backend-java/src/test/java/com/nanri/aiimage/config/ThreadPoolIsolationConfigTest.java @@ -62,7 +62,7 @@ class ThreadPoolIsolationConfigTest { } private TaskExecutor newCoze(ExecutorService virtualExecutor, int maxConcurrent) { - return new TaskFileJobConfig().cozeTaskExecutor(virtualExecutor, maxConcurrent); + return new TaskFileJobConfig().cozeTaskExecutor(virtualExecutor, maxConcurrent, 4, null); } @Test