From 26150b68e890263934907d2c73c009ff474f825a Mon Sep 17 00:00:00 2001 From: matevip Date: Fri, 4 Sep 2026 04:18:41 -0400 Subject: [PATCH] feat(team): add controlled worker intervention --- .../java/vip/mate/agent/AgentService.java | 28 +- .../vip/mate/agent/ChatResultCollector.java | 39 +++ .../agent/graph/NodeStreamingChatHelper.java | 10 +- .../mate/agent/graph/node/ReasoningNode.java | 26 ++ .../approval/ApprovalWorkflowService.java | 73 ++++- .../mate/channel/web/ChatStreamTracker.java | 144 +++++++++ .../cron/service/CronJobLifecycleService.java | 35 ++ .../vip/mate/cron/service/CronJobRunner.java | 5 + .../llm/service/ModelDiscoveryService.java | 8 +- .../mate/team/controller/TeamController.java | 82 +++++ .../mate/team/model/TeamTaskEventEntity.java | 1 + .../vip/mate/team/model/TeamTaskStatus.java | 2 + .../team/service/TeamDispatchService.java | 28 +- .../team/service/TeamRunStateMachine.java | 2 + .../mate/team/service/TeamRunViewFactory.java | 14 +- .../mate/team/service/TeamTaskService.java | 303 +++++++++++++++++- .../TeamWorkerInterventionService.java | 288 +++++++++++++++++ .../TeamWorkerReplayPersistenceService.java | 34 ++ .../mate/agent/ChatResultCollectorTest.java | 40 +++ .../agent/graph/ErrorClassificationTest.java | 21 ++ .../graph/node/ReasoningNodeOutputTest.java | 20 ++ .../ApprovalWorkflowServiceResolveTest.java | 44 +++ .../ChatStreamTrackerContentBatchTest.java | 151 +++++++++ .../ChatStreamTrackerOrphanPolicyTest.java | 6 +- .../service/CronJobLifecycleFenceTest.java | 20 ++ .../service/CronJobOriginPropagationTest.java | 33 ++ .../ModelDiscoveryServiceTestPromptTest.java | 9 + .../team/controller/TeamControllerTest.java | 52 ++- .../service/TeamDispatchServiceEventTest.java | 4 +- .../team/service/TeamDispatchServiceTest.java | 28 +- .../team/service/TeamRunStateMachineTest.java | 2 + .../team/service/TeamRunViewFactoryTest.java | 33 ++ .../team/service/TeamTaskServiceTest.java | 151 +++++++++ .../TeamWorkerInterventionServiceTest.java | 284 ++++++++++++++++ ...eamWorkerReplayPersistenceServiceTest.java | 30 ++ mateclaw-ui/src/api/index.ts | 6 + .../components/team-run/TeamRunAttention.vue | 25 +- .../src/components/team-run/TeamRunDetail.vue | 43 +++ .../src/components/team-run/TeamRunDrawer.vue | 6 + .../components/team-run/TeamRunTaskList.vue | 4 +- .../__tests__/teamRunComponents.test.ts | 22 ++ .../teamRunProjectionPrimitives.test.ts | 28 ++ .../team-run/teamRunAttentionHandlers.ts | 15 +- mateclaw-ui/src/i18n/locales/en-US.ts | 9 + mateclaw-ui/src/i18n/locales/zh-CN.ts | 9 + mateclaw-ui/src/stores/useTeamStore.ts | 2 +- mateclaw-ui/src/views/Teams.vue | 79 ++++- 47 files changed, 2244 insertions(+), 54 deletions(-) create mode 100644 mateclaw-server/src/main/java/vip/mate/agent/ChatResultCollector.java create mode 100644 mateclaw-server/src/main/java/vip/mate/team/service/TeamWorkerInterventionService.java create mode 100644 mateclaw-server/src/main/java/vip/mate/team/service/TeamWorkerReplayPersistenceService.java create mode 100644 mateclaw-server/src/test/java/vip/mate/agent/ChatResultCollectorTest.java create mode 100644 mateclaw-server/src/test/java/vip/mate/channel/web/ChatStreamTrackerContentBatchTest.java create mode 100644 mateclaw-server/src/test/java/vip/mate/team/service/TeamWorkerInterventionServiceTest.java create mode 100644 mateclaw-server/src/test/java/vip/mate/team/service/TeamWorkerReplayPersistenceServiceTest.java diff --git a/mateclaw-server/src/main/java/vip/mate/agent/AgentService.java b/mateclaw-server/src/main/java/vip/mate/agent/AgentService.java index 0d83c1a6..54044fae 100644 --- a/mateclaw-server/src/main/java/vip/mate/agent/AgentService.java +++ b/mateclaw-server/src/main/java/vip/mate/agent/AgentService.java @@ -28,7 +28,6 @@ import vip.mate.workspace.conversation.repository.ConversationMapper; import java.util.List; import java.util.Locale; -import java.time.Duration; import java.util.Map; import java.nio.file.Path; import java.util.concurrent.ConcurrentHashMap; @@ -524,23 +523,7 @@ public class AgentService { * {@code _usage_final} event for token and model attribution. */ private ChatResult collectChatResult(Flux stream) { - StringBuilder content = new StringBuilder(); - final int[] usage = {0, 0}; - final String[] modelInfo = {null, null}; - stream.doOnNext(delta -> { - if (delta.isEvent() && "_usage_final".equals(delta.eventType())) { - Map data = delta.eventData(); - usage[0] = ((Number) data.getOrDefault("promptTokens", 0)).intValue(); - usage[1] = ((Number) data.getOrDefault("completionTokens", 0)).intValue(); - Object model = data.get("runtimeModelName"); - Object provider = data.get("runtimeProviderId"); - if (model != null) modelInfo[0] = model.toString(); - if (provider != null) modelInfo[1] = provider.toString(); - } else if (delta.content() != null) { - content.append(delta.content()); - } - }).blockLast(Duration.ofMinutes(10)); - return new ChatResult(content.toString(), usage[0], usage[1], modelInfo[0], modelInfo[1]); + return ChatResultCollector.collect(stream); } /** @@ -970,10 +953,15 @@ public class AgentService { * post-approval replays). */ public record ChatResult(String content, int promptTokens, int completionTokens, - String runtimeModel, String runtimeProvider) { + String runtimeModel, String runtimeProvider, String finishReason) { + + public ChatResult(String content, int promptTokens, int completionTokens, + String runtimeModel, String runtimeProvider) { + this(content, promptTokens, completionTokens, runtimeModel, runtimeProvider, null); + } public static ChatResult contentOnly(String content) { - return new ChatResult(content != null ? content : "", 0, 0, null, null); + return new ChatResult(content != null ? content : "", 0, 0, null, null, null); } } } diff --git a/mateclaw-server/src/main/java/vip/mate/agent/ChatResultCollector.java b/mateclaw-server/src/main/java/vip/mate/agent/ChatResultCollector.java new file mode 100644 index 00000000..e7bb0ab1 --- /dev/null +++ b/mateclaw-server/src/main/java/vip/mate/agent/ChatResultCollector.java @@ -0,0 +1,39 @@ +package vip.mate.agent; + +import reactor.core.publisher.Flux; + +import java.time.Duration; +import java.util.Map; + +/** Collapses a structured agent stream without discarding terminal metadata. */ +final class ChatResultCollector { + + private ChatResultCollector() { + } + + static AgentService.ChatResult collect(Flux stream) { + StringBuilder content = new StringBuilder(); + final int[] usage = {0, 0}; + final String[] modelInfo = {null, null}; + final String[] finishReason = {null}; + stream.doOnNext(delta -> { + if (delta.isEvent() && "_usage_final".equals(delta.eventType())) { + Map data = delta.eventData() != null ? delta.eventData() : Map.of(); + usage[0] = ((Number) data.getOrDefault("promptTokens", 0)).intValue(); + usage[1] = ((Number) data.getOrDefault("completionTokens", 0)).intValue(); + Object model = data.get("runtimeModelName"); + Object provider = data.get("runtimeProviderId"); + if (model != null) modelInfo[0] = model.toString(); + if (provider != null) modelInfo[1] = provider.toString(); + } else if (delta.isEvent() && "finish_reason".equals(delta.eventType())) { + Map data = delta.eventData(); + Object reason = data != null ? data.get("reason") : null; + if (reason != null) finishReason[0] = reason.toString(); + } else if (delta.content() != null) { + content.append(delta.content()); + } + }).blockLast(Duration.ofMinutes(10)); + return new AgentService.ChatResult(content.toString(), usage[0], usage[1], + modelInfo[0], modelInfo[1], finishReason[0]); + } +} diff --git a/mateclaw-server/src/main/java/vip/mate/agent/graph/NodeStreamingChatHelper.java b/mateclaw-server/src/main/java/vip/mate/agent/graph/NodeStreamingChatHelper.java index 24255cd7..50e643e0 100644 --- a/mateclaw-server/src/main/java/vip/mate/agent/graph/NodeStreamingChatHelper.java +++ b/mateclaw-server/src/main/java/vip/mate/agent/graph/NodeStreamingChatHelper.java @@ -26,6 +26,7 @@ import java.time.format.DateTimeFormatter; import java.time.format.DateTimeParseException; import java.util.ArrayList; import java.util.List; +import java.util.Locale; import java.util.Map; import java.util.concurrent.CancellationException; import java.util.concurrent.CountDownLatch; @@ -554,6 +555,7 @@ public class NodeStreamingChatHelper { // ("credit balance is too low") use these phrases in 402-class responses. // Chinese provider patterns (Zhipu 1113, DashScope, general) — same hard // failure semantics: retrying the same provider won't refill the balance. + String lowerMsg = msg.toLowerCase(Locale.ROOT); if (msg.contains("402") || msg.contains("insufficient_quota") || msg.contains("credit balance is too low") || msg.contains("billing_error") || msg.contains("billing_hard_limit_reached") @@ -562,7 +564,13 @@ public class NodeStreamingChatHelper { || msg.contains("余额不足") || msg.contains("请充值") || msg.contains("\"code\":\"1113\"") || msg.contains("\"code\":1113") || msg.contains("AccountBalanceNotEnough") - || msg.contains("balance not enough")) { + || msg.contains("balance not enough") + || lowerMsg.contains("invalidsubscription") + || lowerMsg.contains("subscription has expired") + || lowerMsg.contains("arrearage") + || lowerMsg.contains("account is in good standing") + || lowerMsg.contains("insufficient_balance") + || lowerMsg.contains("insufficient balance")) { return ErrorType.BILLING; } // RFC-009 P3.2: MODEL_NOT_FOUND — provider rejects the requested model id. diff --git a/mateclaw-server/src/main/java/vip/mate/agent/graph/node/ReasoningNode.java b/mateclaw-server/src/main/java/vip/mate/agent/graph/node/ReasoningNode.java index 316b3d2d..d546aef0 100644 --- a/mateclaw-server/src/main/java/vip/mate/agent/graph/node/ReasoningNode.java +++ b/mateclaw-server/src/main/java/vip/mate/agent/graph/node/ReasoningNode.java @@ -1244,6 +1244,28 @@ public class ReasoningNode implements NodeAction { .build(); } + // Compatibility safety net for providers/adapters that return the + // runtime's reserved error placeholder as an HTTP-successful content + // response. Without this guard the long-form completion gate treats + // the placeholder as a short draft and can repeat it until the graph's + // iteration cap. Cron and other synchronous callers consume the + // resulting structured ERROR_FALLBACK; they do not need to infer from + // user-facing text. + if (isRuntimeErrorPlaceholder(result.text())) { + String errorText = result.text(); + log.error("[ReasoningNode] Runtime error placeholder returned as normal content; failing turn"); + return reasonOutput() + .needsToolCall(false) + .shouldSummarize(false) + .finalAnswer(errorText) + .llmCallCount(nextLlmCallCount) + .finishReason(FinishReason.ERROR_FALLBACK) + .contentStreamed(true) + .thinkingStreamed(result.thinking() != null && !result.thinking().isEmpty()) + .mergeUsage(state, result) + .build(); + } + if (result.partial()) { int partialChars = result.text() != null ? result.text().length() : 0; log.warn("[ReasoningNode] Partial LLM result ({} chars), treating as final answer", partialChars); @@ -1429,6 +1451,10 @@ public class ReasoningNode implements NodeAction { } } + static boolean isRuntimeErrorPlaceholder(String text) { + return text != null && text.stripLeading().startsWith("[错误]"); + } + private static String evidenceWarning(List unsupportedReferences) { return "\n\n[证据不足] 以下引用未出现在本轮已读取/搜索到的工具证据中,或缺少有效来源标注:" + String.join(", ", unsupportedReferences) diff --git a/mateclaw-server/src/main/java/vip/mate/approval/ApprovalWorkflowService.java b/mateclaw-server/src/main/java/vip/mate/approval/ApprovalWorkflowService.java index 238e17fd..e9f475bb 100644 --- a/mateclaw-server/src/main/java/vip/mate/approval/ApprovalWorkflowService.java +++ b/mateclaw-server/src/main/java/vip/mate/approval/ApprovalWorkflowService.java @@ -412,6 +412,24 @@ public class ApprovalWorkflowService implements ApplicationRunner { "consumed", /* removeFromMap */ true); } + /** Claim one exact approval before a team worker executes its guarded tool. */ + @Transactional + public ResolveOutcome claimForReplay(String pendingId, String userId) { + return performResolve(pendingId, userId, "APPROVED", MetadataDecision.APPROVED, + "approved", /* removeFromMap */ false); + } + + /** Consume an approval previously claimed by {@link #claimForReplay}. */ + @Transactional + public ResolveOutcome consumeReplayClaim(String pendingId, String userId) { + PendingApproval target = getReplayClaim(pendingId).orElse(null); + if (target == null) { + return ResolveOutcome.alreadyResolved(pendingId); + } + return performResolveOnSnapshot(target, userId, "APPROVED", "CONSUMED", + MetadataDecision.APPROVED, "consumed", /* removeFromMap */ true); + } + /** * Consume the earliest already-{@code approved} record for the conversation + * tool — used when an out-of-band approval (e.g. /approve text command flow that @@ -423,7 +441,7 @@ public class ApprovalWorkflowService implements ApplicationRunner { if (target == null) { return ResolveOutcome.alreadyResolved(null); } - return performResolveOnSnapshot(target, null, "CONSUMED", MetadataDecision.APPROVED, + return performResolveOnSnapshot(target, null, "APPROVED", "CONSUMED", MetadataDecision.APPROVED, "consumed", /* removeFromMap */ true); } @@ -447,7 +465,7 @@ public class ApprovalWorkflowService implements ApplicationRunner { List outcomes = new java.util.ArrayList<>(targets.size()); for (PendingApproval target : targets) { try { - ResolveOutcome outcome = performResolveOnSnapshot(target, userId, "DENIED", + ResolveOutcome outcome = performResolveOnSnapshot(target, userId, "PENDING", "DENIED", MetadataDecision.DENIED, "denied", /* removeFromMap */ true); if (outcome.dbSynced()) outcomes.add(outcome); } catch (Exception e) { @@ -475,7 +493,7 @@ public class ApprovalWorkflowService implements ApplicationRunner { if (targets.isEmpty()) return List.of(); List outcomes = new java.util.ArrayList<>(targets.size()); for (PendingApproval target : targets) { - ResolveOutcome outcome = performResolveOnSnapshot(target, null, "SUPERSEDED", + ResolveOutcome outcome = performResolveOnSnapshot(target, null, "PENDING", "SUPERSEDED", MetadataDecision.DENIED, "superseded", /* removeFromMap */ true); if (outcome.dbSynced()) outcomes.add(outcome); } @@ -644,12 +662,13 @@ public class ApprovalWorkflowService implements ApplicationRunner { pendingId, snapshot != null, snapshot != null ? snapshot.getStatus() : "n/a"); return ResolveOutcome.alreadyResolved(pendingId); } - return performResolveOnSnapshot(snapshot, userId, dbStatus, metaDecision, + return performResolveOnSnapshot(snapshot, userId, "PENDING", dbStatus, metaDecision, snapshotStatus, removeFromMap); } private ResolveOutcome performResolveOnSnapshot(PendingApproval snapshot, String userId, - String dbStatus, MetadataDecision metaDecision, + String expectedDbStatus, String dbStatus, + MetadataDecision metaDecision, String snapshotStatus, boolean removeFromMap) { // Phase 1 — DB UPDATE (conditional). The eq("PENDING") guard makes the call // idempotent: if another path already won, we get rows=0 and bail without @@ -658,7 +677,7 @@ public class ApprovalWorkflowService implements ApplicationRunner { try { LambdaUpdateWrapper wrapper = new LambdaUpdateWrapper() .eq(ToolApprovalEntity::getPendingId, snapshot.getPendingId()) - .eq(ToolApprovalEntity::getStatus, "PENDING") + .eq(ToolApprovalEntity::getStatus, expectedDbStatus) .set(ToolApprovalEntity::getStatus, dbStatus) .set(ToolApprovalEntity::getResolvedAt, LocalDateTime.now()); if (userId != null) { @@ -672,8 +691,8 @@ public class ApprovalWorkflowService implements ApplicationRunner { throw e; } if (rows == 0) { - log.info("[ApprovalWorkflow] resolve no-op for {}: DB row not in PENDING (concurrent resolve)", - snapshot.getPendingId()); + log.info("[ApprovalWorkflow] resolve no-op for {}: DB row not in {} (concurrent resolve)", + snapshot.getPendingId(), expectedDbStatus); return ResolveOutcome.alreadyResolved(snapshot.getPendingId()); } @@ -809,6 +828,44 @@ public class ApprovalWorkflowService implements ApplicationRunner { return approvalService.getPending(pendingId); } + /** + * Recover an exact APPROVED replay claim from memory or DB. APPROVED claims are + * intentionally durable so a worker replay can be finalized after a restart + * without reopening the approval to denial. + */ + public java.util.Optional getReplayClaim(String pendingId) { + PendingApproval inMemory = approvalService.getPending(pendingId) + .filter(pending -> "approved".equals(pending.getStatus())) + .orElse(null); + if (inMemory != null) { + return java.util.Optional.of(inMemory); + } + ToolApprovalEntity entity = approvalMapper.selectOne( + new LambdaQueryWrapper() + .eq(ToolApprovalEntity::getPendingId, pendingId) + .eq(ToolApprovalEntity::getStatus, "APPROVED")); + if (entity == null) { + return java.util.Optional.empty(); + } + Instant createdAt = entity.getCreatedAt() == null + ? Instant.now() + : entity.getCreatedAt().atZone(ZoneId.systemDefault()).toInstant(); + PendingApproval snapshot = new PendingApproval(entity.getPendingId(), + entity.getConversationId(), entity.getUserId(), entity.getToolName(), + entity.getToolArguments(), entity.getSummary(), createdAt, "approved"); + snapshot.setToolCallPayload(entity.getToolCallPayload()); + snapshot.setSiblingToolCalls(entity.getSiblingToolCalls()); + snapshot.setAgentId(entity.getAgentId()); + snapshot.setChannelType(entity.getChannelType()); + snapshot.setRequesterName(entity.getRequesterName()); + snapshot.setReplyTarget(entity.getReplyTarget()); + snapshot.setFindingsJson(entity.getFindingsJson()); + snapshot.setMaxSeverity(entity.getMaxSeverity()); + snapshot.setSummary(entity.getSummary()); + snapshot.setChatOrigin(entity.getChatOrigin()); + return java.util.Optional.of(snapshot); + } + public PendingApproval findPendingByConversation(String conversationId) { return approvalService.findPendingByConversation(conversationId); } diff --git a/mateclaw-server/src/main/java/vip/mate/channel/web/ChatStreamTracker.java b/mateclaw-server/src/main/java/vip/mate/channel/web/ChatStreamTracker.java index 122f30e0..bd1fdcf7 100644 --- a/mateclaw-server/src/main/java/vip/mate/channel/web/ChatStreamTracker.java +++ b/mateclaw-server/src/main/java/vip/mate/channel/web/ChatStreamTracker.java @@ -1,5 +1,6 @@ package vip.mate.channel.web; +import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import jakarta.annotation.PreDestroy; import lombok.extern.slf4j.Slf4j; @@ -88,6 +89,18 @@ public class ChatStreamTracker { @Value("${mateclaw.stream.iteration-events:true}") private boolean iterationEventsEnabled = true; + /** + * Coalesce the tiny token fragments produced by streaming model clients + * before assigning an SSE id and touching the replay buffer. This keeps + * rendering responsive while avoiding thousands of emitter writes for a + * single long answer. + */ + @Value("${mateclaw.stream.content-batch-ms:25}") + private long contentBatchMs = 25L; + + @Value("${mateclaw.stream.content-batch-chars:256}") + private int contentBatchChars = 256; + /** * Heartbeat cadence (seconds) before the first model token arrives. Short * because pre-token gaps strand the UI on a blank "正在生成中" placeholder @@ -127,6 +140,11 @@ public class ChatStreamTracker { this.iterationEventsEnabled = enabled; } + void setContentBatchingForTesting(long flushMs, int maxChars) { + this.contentBatchMs = Math.max(1L, flushMs); + this.contentBatchChars = Math.max(1, maxChars); + } + public boolean isIterationEventsEnabled() { return iterationEventsEnabled; } @@ -219,6 +237,11 @@ public class ChatStreamTracker { /** 已广播的 pending approval ID 集合(用于幂等去重) */ final java.util.Set broadcastedApprovalIds = java.util.concurrent.ConcurrentHashMap.newKeySet(); + /** Pending visible answer text waiting for the SSE coalescing window. Guarded by lock. */ + String pendingContentField; + final StringBuilder pendingContent = new StringBuilder(); + ScheduledFuture pendingContentFlush; + /** 创建时间(用于 stale 检测和清理) */ final long createdAt = System.currentTimeMillis(); @@ -748,6 +771,99 @@ public class ChatStreamTracker { return true; } + private record ContentDelta(String field, String text) {} + + /** + * Buffer only the two established visible-content wire shapes: + * {@code {"delta":"..."}} (workspace chat) and + * {@code {"text":"..."}} (embedded webchat). Payloads with extra + * metadata stay on the ordinary path so batching never discards fields. + */ + private boolean tryBufferContentDelta(RunState state, String eventName, + String jsonData, boolean skipBuffer) { + if (!"content_delta".equals(eventName) || skipBuffer || state == null) { + return false; + } + ContentDelta delta = parseContentDelta(jsonData); + if (delta == null) { + return false; + } + + boolean flushNow = false; + synchronized (state.lock) { + if (!isCurrent(state) || state.done) { + return true; + } + // A conversation uses one wire field for a run. If a caller does + // switch shapes, flush the old batch and deliver the new payload + // unchanged rather than mixing contracts. + if (state.pendingContentField != null + && !state.pendingContentField.equals(delta.field())) { + return false; + } + state.lastEventAt = System.currentTimeMillis(); + state.pendingContentField = delta.field(); + state.pendingContent.append(delta.text()); + if (state.pendingContent.length() >= Math.max(1, contentBatchChars)) { + flushNow = true; + } else if (state.pendingContentFlush == null + || state.pendingContentFlush.isDone()) { + state.pendingContentFlush = heartbeatScheduler.schedule( + () -> flushPendingContent(state), + Math.max(1L, contentBatchMs), TimeUnit.MILLISECONDS); + } + } + if (flushNow) { + flushPendingContent(state); + } + return true; + } + + private ContentDelta parseContentDelta(String jsonData) { + if (jsonData == null || jsonData.isEmpty()) return null; + try { + JsonNode node = objectMapper.readTree(jsonData); + if (node == null || !node.isObject() || node.size() != 1) return null; + String field = node.has("delta") ? "delta" : node.has("text") ? "text" : null; + if (field == null || !node.path(field).isTextual()) return null; + String text = node.path(field).textValue(); + return text == null || text.isEmpty() ? null : new ContentDelta(field, text); + } catch (Exception ignored) { + return null; + } + } + + /** Snapshot under the run lock, then emit through the fenced raw path. */ + private void flushPendingContent(RunState state) { + String field; + String text; + synchronized (state.lock) { + if (state.pendingContent.length() == 0) { + if (state.pendingContentFlush != null) { + state.pendingContentFlush.cancel(false); + state.pendingContentFlush = null; + } + state.pendingContentField = null; + return; + } + field = state.pendingContentField; + text = state.pendingContent.toString(); + state.pendingContent.setLength(0); + state.pendingContentField = null; + if (state.pendingContentFlush != null) { + state.pendingContentFlush.cancel(false); + state.pendingContentFlush = null; + } + } + try { + String json = objectMapper.writeValueAsString(Map.of(field, text)); + broadcastNow(new RunHandle(state), "content_delta", json, false); + } catch (Exception e) { + log.warn("Failed to flush content batch for {}: {}", + state.conversationId, e.getMessage()); + } + } + /** * 广播事件到所有订阅者并缓存到 buffer. *

@@ -777,6 +893,17 @@ public class ChatStreamTracker { public void broadcast(RunHandle handle, String eventName, String jsonData, boolean skipBuffer) { if (handle == null) return; + RunState state = handle.state; + if (tryBufferContentDelta(state, eventName, jsonData, skipBuffer)) { + return; + } + if (!"heartbeat".equals(eventName)) { + flushPendingContent(state); + } + broadcastNow(handle, eventName, jsonData, skipBuffer); + } + + private void broadcastNow(RunHandle handle, String eventName, String jsonData, boolean skipBuffer) { RunState state = handle.state; boolean isDone = "done".equals(eventName); boolean isPostTurnEvent = "goal_continuation".equals(eventName) @@ -854,6 +981,18 @@ public class ChatStreamTracker { */ public void broadcast(String conversationId, String eventName, String jsonData, boolean skipBuffer) { RunState state = runs.get(conversationId); + if (state == null) return; + if (tryBufferContentDelta(state, eventName, jsonData, skipBuffer)) { + return; + } + if (!"heartbeat".equals(eventName)) { + flushPendingContent(state); + } + broadcastNow(conversationId, eventName, jsonData, skipBuffer); + } + + private void broadcastNow(String conversationId, String eventName, String jsonData, boolean skipBuffer) { + RunState state = runs.get(conversationId); boolean isDone = "done".equals(eventName); boolean isPostTurnEvent = "goal_continuation".equals(eventName) @@ -1312,6 +1451,10 @@ public class ChatStreamTracker { private boolean complete(RunState state) { String conversationId = state.conversationId; + // Some terminal paths do not publish a done envelope. Flush visible + // text while the run is still live so the scheduled batch cannot be + // rejected after state.done flips below. + flushPendingContent(state); ScheduledFuture oldHeartbeat; synchronized (state.lock) { if (!isCurrent(state)) { @@ -1354,6 +1497,7 @@ public class ChatStreamTracker { if (state == null) { return new CompletionResult(true); } + flushPendingContent(state); ScheduledFuture oldHeartbeat; synchronized (state.lock) { if (!isCurrent(state)) { diff --git a/mateclaw-server/src/main/java/vip/mate/cron/service/CronJobLifecycleService.java b/mateclaw-server/src/main/java/vip/mate/cron/service/CronJobLifecycleService.java index b90f6393..33998212 100644 --- a/mateclaw-server/src/main/java/vip/mate/cron/service/CronJobLifecycleService.java +++ b/mateclaw-server/src/main/java/vip/mate/cron/service/CronJobLifecycleService.java @@ -138,6 +138,41 @@ public class CronJobLifecycleService { .set(CronJobRunEntity::getErrorMessage, StrUtil.maxLength(message, 1000))); } + /** + * T-fail — terminal graph failure that arrived as structured stream metadata + * rather than a thrown exception. Persist the diagnostic assistant message + * for conversation coherence, but never publish success, memory, or delivery + * events for an {@code error_fallback} result. + */ + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void finishRunFailed(CronJobRunEntity run, AssistantMessage result, + String conversationId, AgentService.ChatResult chatResult) { + String convId = conversationId != null ? conversationId : run.getConversationId(); + String text = result != null && result.getText() != null ? result.getText() : ""; + int totalTokens = chatResult != null + ? chatResult.promptTokens() + chatResult.completionTokens() : 0; + int updated = runMapper.update(null, new LambdaUpdateWrapper() + .eq(CronJobRunEntity::getId, run.getId()) + .eq(CronJobRunEntity::getStatus, "running") + .set(CronJobRunEntity::getStatus, "failed") + .set(CronJobRunEntity::getFinishedAt, LocalDateTime.now()) + .set(CronJobRunEntity::getErrorMessage, StrUtil.maxLength(text, 1000)) + .set(totalTokens > 0, CronJobRunEntity::getTokenUsage, totalTokens)); + if (updated == 0) { + log.warn("[CronLifecycle] Run {} lost its running fence before graph failure; dropping late result", + run.getId()); + return; + } + + if (chatResult != null) { + conversationService.saveMessage(convId, "assistant", text, null, "error", + chatResult.promptTokens(), chatResult.completionTokens(), + chatResult.runtimeModel(), chatResult.runtimeProvider()); + } else { + conversationService.saveMessage(convId, "assistant", text, null, "error"); + } + } + /** * Insert a {@code running} run row for a task type that does not produce * a conversation (e.g. {@code wiki_process}). No header / user message is diff --git a/mateclaw-server/src/main/java/vip/mate/cron/service/CronJobRunner.java b/mateclaw-server/src/main/java/vip/mate/cron/service/CronJobRunner.java index 11004e11..5ce5e76e 100644 --- a/mateclaw-server/src/main/java/vip/mate/cron/service/CronJobRunner.java +++ b/mateclaw-server/src/main/java/vip/mate/cron/service/CronJobRunner.java @@ -8,6 +8,7 @@ import org.springframework.ai.chat.messages.AssistantMessage; import org.springframework.stereotype.Component; import vip.mate.agent.AgentService; import vip.mate.agent.context.ChatOrigin; +import vip.mate.agent.graph.state.FinishReason; import vip.mate.cron.CronChatOriginFactory; import vip.mate.cron.model.CronJobEntity; import vip.mate.dashboard.model.CronJobRunEntity; @@ -168,6 +169,10 @@ public class CronJobRunner { // T2 — short tx try { + if (FinishReason.ERROR_FALLBACK.getValue().equals(chatResult.finishReason())) { + lifecycle.finishRunFailed(run, result, conversationId, chatResult); + return; + } lifecycle.finishRunAndPublish(job, run, userMessage, result, conversationId, silent, chatResult); } catch (Exception e) { log.error("[CronRunner] T2 finishRunAndPublish failed for job {}: {}", job.getId(), e.getMessage(), e); diff --git a/mateclaw-server/src/main/java/vip/mate/llm/service/ModelDiscoveryService.java b/mateclaw-server/src/main/java/vip/mate/llm/service/ModelDiscoveryService.java index e274869b..8fcb2e61 100644 --- a/mateclaw-server/src/main/java/vip/mate/llm/service/ModelDiscoveryService.java +++ b/mateclaw-server/src/main/java/vip/mate/llm/service/ModelDiscoveryService.java @@ -695,7 +695,13 @@ public class ModelDiscoveryService { requestBody.put("model", modelId); requestBody.put("messages", List.of(Map.of("role", "user", "content", "请回复:连接正常"))); requestBody.put("max_tokens", 10); - requestBody.put("temperature", 0); + Object probeTemperature; + if (ModelFamily.detect(modelId).fixedTemperatureOne()) { + probeTemperature = 1.0d; + } else { + probeTemperature = 0; + } + requestBody.put("temperature", probeTemperature); return requestBody; } diff --git a/mateclaw-server/src/main/java/vip/mate/team/controller/TeamController.java b/mateclaw-server/src/main/java/vip/mate/team/controller/TeamController.java index abad3cd7..835b3462 100644 --- a/mateclaw-server/src/main/java/vip/mate/team/controller/TeamController.java +++ b/mateclaw-server/src/main/java/vip/mate/team/controller/TeamController.java @@ -24,6 +24,7 @@ import vip.mate.team.service.TeamEventChannel; import vip.mate.team.service.TeamManualTaskService; import vip.mate.team.service.TeamService; import vip.mate.team.service.TeamTaskService; +import vip.mate.team.service.TeamWorkerInterventionService; import vip.mate.workspace.core.annotation.RequireWorkspaceRole; import java.security.Principal; @@ -56,6 +57,7 @@ public class TeamController { private final TeamDispatchService dispatchService; private final TeamAnnounceService announceService; private final TeamEventChannel eventChannel; + private final TeamWorkerInterventionService workerInterventionService; private final AgentMapper agentMapper; // ==================== team CRUD ==================== @@ -213,6 +215,60 @@ public class TeamController { }); } + @Operation(summary = "批准 worker 工具调用并在原会话恢复执行") + @PostMapping("/{id}/tasks/{taskId}/worker/approve") + @RequireWorkspaceRole("admin") + public R approveWorkerTool(@PathVariable Long id, @PathVariable Long taskId, + @RequestBody WorkerApprovalRequest req, + Principal principal) { + return workerGuarded(() -> { + requireTeam(id); + requireTask(id, taskId); + if (req == null || req.getPendingId() == null || req.getPendingId().isBlank()) { + throw new IllegalArgumentException("pending approval id is required"); + } + TeamTaskEntity task = workerInterventionService.approve(id, taskId, + req.getPendingId().strip(), principalName(principal)); + return R.ok(toTaskVO(task)); + }); + } + + @Operation(summary = "拒绝 worker 工具调用") + @PostMapping("/{id}/tasks/{taskId}/worker/deny") + @RequireWorkspaceRole("admin") + public R denyWorkerTool(@PathVariable Long id, @PathVariable Long taskId, + @RequestBody WorkerApprovalRequest req, + Principal principal) { + return workerGuarded(() -> { + requireTeam(id); + requireTask(id, taskId); + if (req == null || req.getPendingId() == null || req.getPendingId().isBlank()) { + throw new IllegalArgumentException("pending approval id is required"); + } + TeamTaskEntity task = workerInterventionService.deny(id, taskId, + req.getPendingId().strip(), principalName(principal)); + return R.ok(toTaskVO(task)); + }); + } + + @Operation(summary = "向 worker 原会话发送任务级补充指令") + @PostMapping("/{id}/tasks/{taskId}/worker/feedback") + @RequireWorkspaceRole("admin") + public R feedbackWorker(@PathVariable Long id, @PathVariable Long taskId, + @RequestBody WorkerFeedbackRequest req, + Principal principal) { + return workerGuarded(() -> { + requireTeam(id); + requireTask(id, taskId); + if (req == null || req.getMessage() == null || req.getMessage().isBlank()) { + throw new IllegalArgumentException("feedback is required"); + } + TeamTaskEntity task = workerInterventionService.feedback(id, taskId, + req.getMessage(), principalName(principal)); + return R.ok(toTaskVO(task)); + }); + } + @Operation(summary = "驳回 in_review 任务") @PostMapping("/{id}/tasks/{taskId}/reject") @RequireWorkspaceRole("admin") @@ -310,6 +366,10 @@ public class TeamController { eventChannel.publishTaskEvent(taskService.getTask(taskId), event, Map.of()); } + private String principalName(Principal principal) { + return principal != null && principal.getName() != null ? principal.getName() : "admin"; + } + @Operation(summary = "添加评论") @PostMapping("/{id}/tasks/{taskId}/comments") @RequireWorkspaceRole("admin") @@ -352,6 +412,18 @@ public class TeamController { } } + /** Intervention endpoints expose recoverable client states instead of generic 500s. */ + private R workerGuarded(Supplier> action) { + try { + return action.get(); + } catch (IllegalArgumentException error) { + int code = error.getMessage() != null && error.getMessage().contains("not found") ? 404 : 400; + return R.fail(code, error.getMessage()); + } catch (IllegalStateException error) { + return R.fail(409, error.getMessage()); + } + } + private TeamTaskEntity requireTask(Long teamId, Long taskId) { TeamTaskEntity task = taskService.getTask(taskId); if (task == null || !task.getTeamId().equals(teamId)) { @@ -495,4 +567,14 @@ public class TeamController { public static class CommentRequest { private String content; } + + @Data + public static class WorkerApprovalRequest { + private String pendingId; + } + + @Data + public static class WorkerFeedbackRequest { + private String message; + } } diff --git a/mateclaw-server/src/main/java/vip/mate/team/model/TeamTaskEventEntity.java b/mateclaw-server/src/main/java/vip/mate/team/model/TeamTaskEventEntity.java index abb7a23e..b69bb3b1 100644 --- a/mateclaw-server/src/main/java/vip/mate/team/model/TeamTaskEventEntity.java +++ b/mateclaw-server/src/main/java/vip/mate/team/model/TeamTaskEventEntity.java @@ -27,6 +27,7 @@ public class TeamTaskEventEntity { public static final String DELIVERABLE = "deliverable"; public static final String COMPLETED = "completed"; public static final String IN_REVIEW = "in_review"; + public static final String AWAITING_APPROVAL = "awaiting_approval"; public static final String FAILED = "failed"; public static final String CANCELLED = "cancelled"; public static final String APPROVED = "approved"; diff --git a/mateclaw-server/src/main/java/vip/mate/team/model/TeamTaskStatus.java b/mateclaw-server/src/main/java/vip/mate/team/model/TeamTaskStatus.java index 97dc5c0e..daeeb68a 100644 --- a/mateclaw-server/src/main/java/vip/mate/team/model/TeamTaskStatus.java +++ b/mateclaw-server/src/main/java/vip/mate/team/model/TeamTaskStatus.java @@ -9,6 +9,7 @@ import java.util.Set; * pending ──claim/assign──▶ in_progress ──complete──▶ completed * │ │ (require_approval) ▶ in_review ──approve──▶ completed * │ │ └──reject───▶ cancelled + * │ ├──guarded tool──▶ awaiting_approval ──approve──▶ in_progress * │ ├──blocker/error──▶ failed ──retry──▶ pending * │ └──lease expired──▶ stale ──retry──▶ pending * ├──blocked_by set──▶ blocked ──all blockers released──▶ pending @@ -21,6 +22,7 @@ public final class TeamTaskStatus { public static final String PENDING = "pending"; public static final String IN_PROGRESS = "in_progress"; + public static final String AWAITING_APPROVAL = "awaiting_approval"; public static final String IN_REVIEW = "in_review"; public static final String COMPLETED = "completed"; public static final String FAILED = "failed"; diff --git a/mateclaw-server/src/main/java/vip/mate/team/service/TeamDispatchService.java b/mateclaw-server/src/main/java/vip/mate/team/service/TeamDispatchService.java index 92da9640..da7dd714 100644 --- a/mateclaw-server/src/main/java/vip/mate/team/service/TeamDispatchService.java +++ b/mateclaw-server/src/main/java/vip/mate/team/service/TeamDispatchService.java @@ -10,6 +10,8 @@ import org.springframework.transaction.event.TransactionPhase; import org.springframework.transaction.event.TransactionalEventListener; import vip.mate.team.event.TeamTasksDelegatedEvent; import vip.mate.agent.AgentService; +import vip.mate.approval.ApprovalWorkflowService; +import vip.mate.approval.PendingApproval; import vip.mate.channel.web.ChatStreamTracker; import vip.mate.team.model.AgentTeamEntity; import vip.mate.team.model.TeamTaskEntity; @@ -107,6 +109,7 @@ public class TeamDispatchService { private final ChatStreamTracker streamTracker; private final TeamAnnounceService announceService; private final TeamEventChannel eventChannel; + private final ApprovalWorkflowService approvalService; /** Members with a run currently in flight in this JVM (belt-and-braces on top of hasActiveTask). */ private final Set runningMembers = ConcurrentHashMap.newKeySet(); @@ -215,9 +218,7 @@ public class TeamDispatchService { streamTracker.incrementFlux(childConvId); // Renew the execution lease while the member works; the conditional // UPDATE inside renewLock makes this a no-op once the task settles. - heartbeat = HEARTBEAT_SCHEDULER.scheduleAtFixedRate( - () -> taskService.renewLock(task.getId()), - HEARTBEAT_MINUTES, HEARTBEAT_MINUTES, TimeUnit.MINUTES); + heartbeat = startLeaseHeartbeat(task.getId()); broadcast(task, "team_task_dispatched", Map.of()); log.info("Team {} task #{} dispatched to agent {} (conv {})", teamId, task.getTaskNumber(), memberId, childConvId); @@ -235,6 +236,20 @@ public class TeamDispatchService { conversationService.saveMessage(childConvId, "assistant", reply); } + PendingApproval pending = approvalService.findPendingByConversation(childConvId); + if (pending != null) { + String summary = pending.getSummary() == null || pending.getSummary().isBlank() + ? pending.getReason() : pending.getSummary(); + if (taskService.parkForToolApproval(task.getId(), pending.getPendingId(), summary)) { + TeamTaskEntity parked = taskService.getTask(task.getId()); + broadcast(parked != null ? parked : task, "team_task_awaiting_approval", + Map.of("pendingId", pending.getPendingId(), + "toolName", pending.getToolName() == null ? "" : pending.getToolName(), + "summary", summary == null ? "Tool approval required" : summary)); + } + return; + } + settleOutcome(task, reply); } catch (Exception e) { log.warn("Team {} task #{} member run ended exceptionally: {}", teamId, @@ -259,6 +274,13 @@ public class TeamDispatchService { } } + /** Share the same DB-backed lease heartbeat with controlled worker replays. */ + ScheduledFuture startLeaseHeartbeat(Long taskId) { + return HEARTBEAT_SCHEDULER.scheduleAtFixedRate( + () -> taskService.renewLock(taskId), + HEARTBEAT_MINUTES, HEARTBEAT_MINUTES, TimeUnit.MINUTES); + } + /** * Ask the member conversation executing this task to stop at the next graph * node boundary (cancel path). No-op when the task never dispatched or the diff --git a/mateclaw-server/src/main/java/vip/mate/team/service/TeamRunStateMachine.java b/mateclaw-server/src/main/java/vip/mate/team/service/TeamRunStateMachine.java index 23f9d920..ef4e0ad6 100644 --- a/mateclaw-server/src/main/java/vip/mate/team/service/TeamRunStateMachine.java +++ b/mateclaw-server/src/main/java/vip/mate/team/service/TeamRunStateMachine.java @@ -16,12 +16,14 @@ public final class TeamRunStateMachine { TeamTaskStatus.PENDING, TeamTaskStatus.BLOCKED, TeamTaskStatus.IN_PROGRESS, + TeamTaskStatus.AWAITING_APPROVAL, TeamTaskStatus.STALE ); private static final Set KNOWN_TASK_STATUSES = Set.of( TeamTaskStatus.PENDING, TeamTaskStatus.BLOCKED, TeamTaskStatus.IN_PROGRESS, + TeamTaskStatus.AWAITING_APPROVAL, TeamTaskStatus.IN_REVIEW, TeamTaskStatus.COMPLETED, TeamTaskStatus.FAILED, diff --git a/mateclaw-server/src/main/java/vip/mate/team/service/TeamRunViewFactory.java b/mateclaw-server/src/main/java/vip/mate/team/service/TeamRunViewFactory.java index 720fedb9..63735020 100644 --- a/mateclaw-server/src/main/java/vip/mate/team/service/TeamRunViewFactory.java +++ b/mateclaw-server/src/main/java/vip/mate/team/service/TeamRunViewFactory.java @@ -177,6 +177,8 @@ final class TeamRunViewFactory { List items = new ArrayList<>(); for (TeamTaskEntity task : tasks) { String type = switch (task.getStatus()) { + case TeamTaskStatus.AWAITING_APPROVAL -> replayOutcomeUncertain(task) + ? "replay_uncertain" : "approval"; case TeamTaskStatus.IN_REVIEW -> "review"; case TeamTaskStatus.FAILED -> "failure"; case TeamTaskStatus.BLOCKED -> "blocked"; @@ -185,7 +187,8 @@ final class TeamRunViewFactory { }; if (type != null) { String message = text(task.getReason()); - int priority = TeamTaskStatus.IN_REVIEW.equals(task.getStatus()) ? 0 : 20; + int priority = TeamTaskStatus.IN_REVIEW.equals(task.getStatus()) + || TeamTaskStatus.AWAITING_APPROVAL.equals(task.getStatus()) ? 0 : 20; items.add(new TeamRunView.AttentionItem("task:" + task.getId() + ":" + type, type, priority == 0 ? "action" : "error", priority, task.getId(), message == null ? task.getSubject() : message, task.getUpdateTime())); @@ -205,6 +208,15 @@ final class TeamRunViewFactory { return List.copyOf(items); } + private static boolean replayOutcomeUncertain(TeamTaskEntity task) { + try { + JSONObject approval = JSONUtil.parseObj(task.getMetadata()).getJSONObject("toolApproval"); + return approval != null && approval.getBool("replayOutcomeUncertain", false); + } catch (RuntimeException invalidMetadata) { + return false; + } + } + private static TeamRunView.Liveness liveness(String status, LocalDateTime lastActivity, List tasks) { if (TeamRunStatus.isTerminal(status)) { diff --git a/mateclaw-server/src/main/java/vip/mate/team/service/TeamTaskService.java b/mateclaw-server/src/main/java/vip/mate/team/service/TeamTaskService.java index e18846c9..7cfcbf2e 100644 --- a/mateclaw-server/src/main/java/vip/mate/team/service/TeamTaskService.java +++ b/mateclaw-server/src/main/java/vip/mate/team/service/TeamTaskService.java @@ -45,6 +45,8 @@ import java.util.regex.Pattern; @RequiredArgsConstructor public class TeamTaskService { + private static final int MAX_STAGED_REPLAY_RESULT_CHARS = 8000; + private static final Pattern CHECKPOINT_RANGE = Pattern.compile( "(?i)R(\\d{3})\\s*[-–—]\\s*R(\\d{3})"); @@ -375,6 +377,275 @@ public class TeamTaskService { return updated; } + /** Park a running worker task until its guarded tool call receives a human decision. */ + public boolean parkForToolApproval(Long taskId, String pendingId, String summary) { + if (pendingId == null || pendingId.isBlank()) { + throw new IllegalArgumentException("pending approval id is required"); + } + TeamTaskEntity task = taskMapper.selectById(taskId); + if (task == null) { + return false; + } + JSONObject metadata; + try { + metadata = task.getMetadata() == null || task.getMetadata().isBlank() + ? new JSONObject() : JSONUtil.parseObj(task.getMetadata()); + } catch (RuntimeException invalid) { + metadata = new JSONObject(); + } + String detail = summary == null || summary.isBlank() + ? "Tool approval required" : summary.strip(); + metadata.set("toolApproval", new JSONObject() + .set("pendingId", pendingId) + .set("summary", detail)); + boolean parked = taskMapper.update(null, Wrappers.lambdaUpdate() + .eq(TeamTaskEntity::getId, taskId) + .in(TeamTaskEntity::getStatus, TeamTaskStatus.IN_PROGRESS, + TeamTaskStatus.AWAITING_APPROVAL) + .set(TeamTaskEntity::getStatus, TeamTaskStatus.AWAITING_APPROVAL) + .set(TeamTaskEntity::getReason, detail) + .set(TeamTaskEntity::getMetadata, metadata.toString()) + .set(TeamTaskEntity::getLockExpiresAt, null)) == 1; + if (parked) { + recordEvent(task.getTeamId(), taskId, TeamTaskEventEntity.AWAITING_APPROVAL, + AUTHOR_SYSTEM, null, pendingId + " — " + detail); + projectTask(taskId); + } + return parked; + } + + /** Resume the exact guarded tool request currently recorded on a parked task. */ + public boolean resumeAfterToolApproval(Long taskId, String pendingId) { + if (pendingId == null || pendingId.isBlank()) { + throw new IllegalArgumentException("pending approval id is required"); + } + TeamTaskEntity task = requireTask(taskId); + if (!TeamTaskStatus.AWAITING_APPROVAL.equals(task.getStatus())) { + throw new IllegalStateException("task #" + task.getTaskNumber() + + " is not awaiting tool approval"); + } + JSONObject metadata = parseMetadata(task.getMetadata()); + JSONObject approval = metadata.getJSONObject("toolApproval"); + String currentPendingId = approval == null ? null : approval.getStr("pendingId"); + if (!pendingId.equals(currentPendingId)) { + throw new IllegalStateException("tool approval is no longer current for task #" + + task.getTaskNumber()); + } + approval.set("replayInProgress", true); + metadata.set("toolApproval", approval); + boolean resumed = taskMapper.update(null, Wrappers.lambdaUpdate() + .eq(TeamTaskEntity::getId, taskId) + .eq(TeamTaskEntity::getStatus, TeamTaskStatus.AWAITING_APPROVAL) + .set(TeamTaskEntity::getStatus, TeamTaskStatus.IN_PROGRESS) + .set(TeamTaskEntity::getOwnerAgentId, task.getAssigneeAgentId()) + .set(TeamTaskEntity::getReason, null) + .set(TeamTaskEntity::getMetadata, metadata.toString()) + .set(TeamTaskEntity::getLockExpiresAt, newLease())) == 1; + if (resumed) { + projectTask(taskId); + } + return resumed; + } + + /** Settle a parked guarded tool request as denied without executing it. */ + public boolean denyToolApproval(Long taskId, String pendingId, String requester) { + TeamTaskEntity task = requireTask(taskId); + if (!TeamTaskStatus.AWAITING_APPROVAL.equals(task.getStatus())) { + throw new IllegalStateException("task #" + task.getTaskNumber() + + " is not awaiting tool approval"); + } + JSONObject metadata = parseMetadata(task.getMetadata()); + JSONObject approval = metadata.getJSONObject("toolApproval"); + if (approval == null || !Objects.equals(pendingId, approval.getStr("pendingId"))) { + throw new IllegalStateException("tool approval is no longer current for task #" + + task.getTaskNumber()); + } + metadata.remove("toolApproval"); + String reason = "Tool request denied by " + + (requester == null || requester.isBlank() ? "user" : requester); + boolean denied = taskMapper.update(null, Wrappers.lambdaUpdate() + .eq(TeamTaskEntity::getId, taskId) + .eq(TeamTaskEntity::getStatus, TeamTaskStatus.AWAITING_APPROVAL) + .set(TeamTaskEntity::getStatus, TeamTaskStatus.FAILED) + .set(TeamTaskEntity::getReason, reason) + .set(TeamTaskEntity::getMetadata, metadata.toString()) + .set(TeamTaskEntity::getLockExpiresAt, null)) == 1; + if (denied) { + recordEvent(task.getTeamId(), taskId, TeamTaskEventEntity.FAILED, + AUTHOR_USER, requester, reason); + projectTask(taskId); + } + return denied; + } + + /** Durably stage a successful replay before consuming its approval. */ + public boolean stageToolReplayResult(Long taskId, String pendingId, String reply) { + TeamTaskEntity task = requireTask(taskId); + JSONObject metadata = parseMetadata(task.getMetadata()); + JSONObject currentApproval = metadata.getJSONObject("toolApproval"); + if (!TeamTaskStatus.IN_PROGRESS.equals(task.getStatus()) + || currentApproval == null + || !Objects.equals(pendingId, currentApproval.getStr("pendingId"))) { + throw new IllegalStateException("tool approval is no longer current for task #" + + task.getTaskNumber()); + } + JSONObject approval = new JSONObject() + .set("pendingId", pendingId) + .set("summary", "Approved tool completed; finalizing result") + .set("replayResult", truncate(reply == null ? "" : reply, + MAX_STAGED_REPLAY_RESULT_CHARS)); + metadata.set("toolApproval", approval); + boolean staged = taskMapper.update(null, Wrappers.lambdaUpdate() + .eq(TeamTaskEntity::getId, taskId) + .eq(TeamTaskEntity::getStatus, TeamTaskStatus.IN_PROGRESS) + .set(TeamTaskEntity::getStatus, TeamTaskStatus.AWAITING_APPROVAL) + .set(TeamTaskEntity::getReason, "Approved tool completed; finalizing result") + .set(TeamTaskEntity::getMetadata, metadata.toString()) + .set(TeamTaskEntity::getLockExpiresAt, null)) == 1; + if (staged) { + projectTask(taskId); + } + return staged; + } + + /** Park a failed replay without allowing an automatic second execution. */ + public boolean parkToolReplayUncertain(Long taskId, String pendingId, String detail) { + TeamTaskEntity task = requireTask(taskId); + JSONObject metadata = parseMetadata(task.getMetadata()); + JSONObject approval = metadata.getJSONObject("toolApproval"); + if (approval == null || !Objects.equals(pendingId, approval.getStr("pendingId"))) { + return false; + } + approval.set("replayInProgress", false); + approval.set("replayOutcomeUncertain", true); + approval.set("summary", detail); + metadata.set("toolApproval", approval); + boolean parked = taskMapper.update(null, Wrappers.lambdaUpdate() + .eq(TeamTaskEntity::getId, taskId) + .eq(TeamTaskEntity::getStatus, TeamTaskStatus.IN_PROGRESS) + .set(TeamTaskEntity::getStatus, TeamTaskStatus.AWAITING_APPROVAL) + .set(TeamTaskEntity::getReason, detail) + .set(TeamTaskEntity::getMetadata, metadata.toString()) + .set(TeamTaskEntity::getLockExpiresAt, null)) == 1; + if (parked) { + recordEvent(task.getTeamId(), taskId, TeamTaskEventEntity.AWAITING_APPROVAL, + AUTHOR_SYSTEM, null, detail); + projectTask(taskId); + } + return parked; + } + + public String stagedToolReplayResult(TeamTaskEntity task) { + JSONObject approval = task == null ? null : parseMetadata(task.getMetadata()) + .getJSONObject("toolApproval"); + return approval != null && approval.containsKey("replayResult") + ? approval.getStr("replayResult", "") : null; + } + + public boolean isToolReplayMessagePersisted(TeamTaskEntity task) { + JSONObject approval = task == null ? null : parseMetadata(task.getMetadata()) + .getJSONObject("toolApproval"); + return approval != null && approval.getBool("messagePersisted", false); + } + + public boolean isToolReplayOutcomeUncertain(TeamTaskEntity task) { + JSONObject approval = task == null ? null : parseMetadata(task.getMetadata()) + .getJSONObject("toolApproval"); + return approval != null && approval.getBool("replayOutcomeUncertain", false); + } + + public boolean markToolReplayMessagePersisted(Long taskId, String pendingId) { + TeamTaskEntity task = requireTask(taskId); + JSONObject metadata = parseMetadata(task.getMetadata()); + JSONObject approval = metadata.getJSONObject("toolApproval"); + if (approval == null || !Objects.equals(pendingId, approval.getStr("pendingId"))) { + return false; + } + approval.set("messagePersisted", true); + metadata.set("toolApproval", approval); + return taskMapper.update(null, Wrappers.lambdaUpdate() + .eq(TeamTaskEntity::getId, taskId) + .eq(TeamTaskEntity::getStatus, TeamTaskStatus.AWAITING_APPROVAL) + .set(TeamTaskEntity::getMetadata, metadata.toString())) == 1; + } + + /** Stop an already-claimed replay after a failed execution attempt. */ + public boolean abortClaimedToolReplay(Long taskId, String pendingId, String requester) { + TeamTaskEntity task = requireTask(taskId); + if (!TeamTaskStatus.AWAITING_APPROVAL.equals(task.getStatus())) { + throw new IllegalStateException("task #" + task.getTaskNumber() + + " is not awaiting replay recovery"); + } + JSONObject metadata = parseMetadata(task.getMetadata()); + JSONObject approval = metadata.getJSONObject("toolApproval"); + if (approval == null || !Objects.equals(pendingId, approval.getStr("pendingId")) + || approval.containsKey("replayResult")) { + throw new IllegalStateException("tool replay is no longer abortable for task #" + + task.getTaskNumber()); + } + metadata.remove("toolApproval"); + String actor = requester == null || requester.isBlank() ? "user" : requester; + String reason = "Approved tool replay aborted by " + actor + + "; the previous execution outcome may be uncertain"; + boolean aborted = taskMapper.update(null, Wrappers.lambdaUpdate() + .eq(TeamTaskEntity::getId, taskId) + .eq(TeamTaskEntity::getStatus, TeamTaskStatus.AWAITING_APPROVAL) + .set(TeamTaskEntity::getStatus, TeamTaskStatus.FAILED) + .set(TeamTaskEntity::getReason, reason) + .set(TeamTaskEntity::getMetadata, metadata.toString()) + .set(TeamTaskEntity::getLockExpiresAt, null)) == 1; + if (aborted) { + recordEvent(task.getTeamId(), taskId, TeamTaskEventEntity.FAILED, + AUTHOR_USER, requester, reason); + projectTask(taskId); + } + return aborted; + } + + private static String truncate(String value, int maxChars) { + return value.length() <= maxChars ? value : value.substring(0, maxChars); + } + + /** Reopen a settled worker task for one deliberate, task-scoped follow-up turn. */ + public boolean resumeForWorkerFeedback(Long taskId) { + TeamTaskEntity task = requireTask(taskId); + if (TeamTaskStatus.IN_PROGRESS.equals(task.getStatus())) { + throw new IllegalStateException("worker task is already running"); + } + if (TeamTaskStatus.AWAITING_APPROVAL.equals(task.getStatus())) { + throw new IllegalStateException("resolve the pending tool approval before sending feedback"); + } + if (TeamTaskStatus.CANCELLED.equals(task.getStatus()) + || TeamTaskStatus.PENDING.equals(task.getStatus()) + || TeamTaskStatus.BLOCKED.equals(task.getStatus())) { + throw new IllegalStateException("task #" + task.getTaskNumber() + + " cannot accept worker feedback while " + task.getStatus()); + } + boolean resumed = taskMapper.update(null, Wrappers.lambdaUpdate() + .eq(TeamTaskEntity::getId, taskId) + .in(TeamTaskEntity::getStatus, TeamTaskStatus.COMPLETED, TeamTaskStatus.FAILED, + TeamTaskStatus.STALE, TeamTaskStatus.IN_REVIEW) + .set(TeamTaskEntity::getStatus, TeamTaskStatus.IN_PROGRESS) + .set(TeamTaskEntity::getOwnerAgentId, task.getAssigneeAgentId()) + .set(TeamTaskEntity::getReason, null) + .set(TeamTaskEntity::getLockExpiresAt, newLease())) == 1; + if (resumed) { + projectTask(taskId); + } + return resumed; + } + + private static JSONObject parseMetadata(String raw) { + if (raw == null || raw.isBlank()) { + return new JSONObject(); + } + try { + return JSONUtil.parseObj(raw); + } catch (RuntimeException invalid) { + return new JSONObject(); + } + } + /** Extend the execution lease (runner heartbeat). */ public void renewLock(Long taskId) { taskMapper.update(null, Wrappers.lambdaUpdate() @@ -684,20 +955,48 @@ public class TeamTaskService { .eq(TeamTaskEntity::getStatus, TeamTaskStatus.IN_PROGRESS) .isNotNull(TeamTaskEntity::getLockExpiresAt) .lt(TeamTaskEntity::getLockExpiresAt, LocalDateTime.now())); + int staleCount = 0; + int uncertainReplayCount = 0; for (TeamTaskEntity task : expired) { + JSONObject metadata = parseMetadata(task.getMetadata()); + JSONObject approval = metadata.getJSONObject("toolApproval"); + if (approval != null && approval.getBool("replayInProgress", false)) { + approval.set("replayInProgress", false); + approval.set("replayOutcomeUncertain", true); + approval.set("summary", "Approved tool replay was interrupted; outcome is uncertain"); + metadata.set("toolApproval", approval); + String reason = "Approved tool replay lease expired; stop the replay or verify its outcome manually"; + int rows = taskMapper.update(null, Wrappers.lambdaUpdate() + .eq(TeamTaskEntity::getId, task.getId()) + .eq(TeamTaskEntity::getStatus, TeamTaskStatus.IN_PROGRESS) + .set(TeamTaskEntity::getStatus, TeamTaskStatus.AWAITING_APPROVAL) + .set(TeamTaskEntity::getReason, reason) + .set(TeamTaskEntity::getMetadata, metadata.toString()) + .set(TeamTaskEntity::getLockExpiresAt, null)); + if (rows == 1) { + uncertainReplayCount++; + recordEvent(task.getTeamId(), task.getId(), + TeamTaskEventEntity.AWAITING_APPROVAL, + AUTHOR_SYSTEM, null, reason); + projectTask(task); + } + continue; + } int rows = taskMapper.update(null, Wrappers.lambdaUpdate() .eq(TeamTaskEntity::getId, task.getId()) .eq(TeamTaskEntity::getStatus, TeamTaskStatus.IN_PROGRESS) .set(TeamTaskEntity::getStatus, TeamTaskStatus.STALE) .set(TeamTaskEntity::getReason, "execution lease expired")); if (rows == 1) { + staleCount++; recordEvent(task.getTeamId(), task.getId(), TeamTaskEventEntity.STALE, AUTHOR_SYSTEM, null, "execution lease expired"); projectTask(task); } } - if (!expired.isEmpty()) { - log.warn("Marked {} team task(s) stale after lease expiry", expired.size()); + if (staleCount > 0 || uncertainReplayCount > 0) { + log.warn("Recovered expired team task leases: stale={}, replayOutcomeUncertain={}", + staleCount, uncertainReplayCount); } return expired; } diff --git a/mateclaw-server/src/main/java/vip/mate/team/service/TeamWorkerInterventionService.java b/mateclaw-server/src/main/java/vip/mate/team/service/TeamWorkerInterventionService.java new file mode 100644 index 00000000..47021d39 --- /dev/null +++ b/mateclaw-server/src/main/java/vip/mate/team/service/TeamWorkerInterventionService.java @@ -0,0 +1,288 @@ +package vip.mate.team.service; + +import cn.hutool.json.JSONUtil; +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; +import vip.mate.agent.AgentService; +import vip.mate.agent.context.ChatOrigin; +import vip.mate.agent.runtime.ConversationTurnGate; +import vip.mate.approval.ApprovalWorkflowService; +import vip.mate.approval.PendingApproval; +import vip.mate.approval.ResolveOutcome; +import vip.mate.channel.web.ChatStreamTracker; +import vip.mate.team.model.TeamTaskEntity; +import vip.mate.workspace.conversation.ConversationService; + +import java.util.Map; +import java.util.Objects; +import java.util.concurrent.ScheduledFuture; + +/** Controlled write path for a delegated worker conversation. */ +@Service +@RequiredArgsConstructor +public class TeamWorkerInterventionService { + + static final String REPLAY_PROMPT = "继续执行已批准的工具调用。"; + + private final TeamTaskService taskService; + private final TeamWorkerConversationGovernanceService governanceService; + private final ApprovalWorkflowService approvalService; + private final AgentService agentService; + private final ConversationService conversationService; + private final ConversationTurnGate turnGate; + private final ChatStreamTracker streamTracker; + private final TeamDispatchService dispatchService; + private final TeamAnnounceService announceService; + private final TeamEventChannel eventChannel; + private final TeamWorkerReplayPersistenceService replayPersistenceService; + + public TeamTaskEntity approve(Long teamId, Long taskId, String pendingId, String requester) { + Intervention intervention = requireIntervention(teamId, taskId); + if (!vip.mate.team.model.TeamTaskStatus.AWAITING_APPROVAL.equals( + intervention.task().getStatus())) { + return intervention.task(); + } + if (taskService.isToolReplayOutcomeUncertain(intervention.task())) { + throw new IllegalStateException( + "tool replay outcome is uncertain; stop it or verify the side effect manually"); + } + PendingApproval pending = requireReplayApproval(intervention, pendingId); + ScheduledFuture heartbeat = null; + try (ConversationTurnGate.Permit permit = reserve(intervention.conversationId())) { + requireMemberIdle(intervention); + String reply = taskService.stagedToolReplayResult(intervention.task()); + if (reply == null) { + PendingApproval claimedPending = claimReplay( + intervention, pendingId, requester, pending); + if (!taskService.resumeAfterToolApproval(taskId, pendingId)) { + throw new IllegalStateException( + "worker task changed while replay was being claimed"); + } + heartbeat = dispatchService.startLeaseHeartbeat(taskId); + conversationService.removeApprovalPlaceholders(intervention.conversationId()); + ChatOrigin origin = approvalService.restoreChatOrigin(claimedPending.getChatOrigin()); + AgentService.ChatResult result; + try { + result = turnGate.withPermit(permit, () -> agentService.chatWithReplayWithUsage( + intervention.agentId(), REPLAY_PROMPT, intervention.conversationId(), + claimedPending.getToolCallPayload(), origin)); + } catch (RuntimeException error) { + taskService.parkToolReplayUncertain(taskId, pendingId, + "Approved tool replay failed and its outcome is uncertain: " + + safeMessage(error)); + throw error; + } + reply = result == null ? "" : result.content(); + if (!taskService.stageToolReplayResult(taskId, pendingId, reply)) { + throw new IllegalStateException("tool replay completed but its result could not be staged"); + } + replayPersistenceService.persist(taskId, pendingId, + intervention.conversationId(), reply, result); + } else if (!taskService.isToolReplayMessagePersisted(intervention.task()) + && !reply.isBlank()) { + replayPersistenceService.persist(taskId, pendingId, + intervention.conversationId(), reply, null); + } + ResolveOutcome consumed = approvalService.consumeReplayClaim(pendingId, requester); + if (!consumed.isConsumed()) { + throw new IllegalStateException("approved tool replay could not be finalized"); + } + if (!taskService.resumeAfterToolApproval(taskId, pendingId)) { + throw new IllegalStateException("worker task changed while replay was being finalized"); + } + settleOrPark(intervention, reply); + return taskService.getTask(taskId); + } finally { + if (heartbeat != null) { + heartbeat.cancel(false); + } + } + } + + @Transactional + public TeamTaskEntity deny(Long teamId, Long taskId, String pendingId, String requester) { + Intervention intervention = requireIntervention(teamId, taskId); + if (!vip.mate.team.model.TeamTaskStatus.AWAITING_APPROVAL.equals( + intervention.task().getStatus())) { + return intervention.task(); + } + if (taskService.stagedToolReplayResult(intervention.task()) != null) { + throw new IllegalStateException("approved tool already executed; finalize its result instead"); + } + PendingApproval approval = requireReplayApproval(intervention, pendingId); + try (ConversationTurnGate.Permit ignored = reserve(intervention.conversationId())) { + conversationService.removeApprovalPlaceholders(intervention.conversationId()); + String event; + if ("approved".equals(approval.getStatus())) { + if (!taskService.abortClaimedToolReplay(taskId, pendingId, requester)) { + throw new IllegalStateException("worker task changed while replay was being stopped"); + } + ResolveOutcome consumed = approvalService.consumeReplayClaim(pendingId, requester); + if (!consumed.isConsumed()) { + throw new IllegalStateException("claimed tool replay could not be stopped"); + } + event = "team_task_tool_replay_aborted"; + } else { + ResolveOutcome outcome = approvalService.resolve(pendingId, requester, "denied"); + if (outcome.isAlreadyResolved()) { + throw new IllegalStateException("tool approval is no longer pending"); + } + if (!taskService.denyToolApproval(taskId, pendingId, requester)) { + throw new IllegalStateException("worker task changed while approval was being denied"); + } + event = "team_task_tool_denied"; + } + TeamTaskEntity settled = taskService.getTask(taskId); + eventChannel.publishTaskEvent(settled, event, Map.of("pendingId", pendingId)); + announceService.announceTaskSettled(settled); + dispatchService.requestDispatch(teamId); + return settled; + } + } + + public TeamTaskEntity feedback(Long teamId, Long taskId, String message, String requester) { + String feedback = message == null ? "" : message.strip(); + if (feedback.isEmpty()) { + throw new IllegalArgumentException("feedback is required"); + } + if (feedback.length() > 4000) { + throw new IllegalArgumentException("feedback must be at most 4000 characters"); + } + Intervention intervention = requireIntervention(teamId, taskId); + if (approvalService.findPendingByConversation(intervention.conversationId()) != null) { + throw new IllegalStateException("resolve the pending tool approval before sending feedback"); + } + try (ConversationTurnGate.Permit permit = reserve(intervention.conversationId())) { + requireMemberIdle(intervention); + if (!taskService.resumeForWorkerFeedback(taskId)) { + throw new IllegalStateException("worker task changed before feedback could start"); + } + conversationService.saveMessage(intervention.conversationId(), "user", feedback); + var agent = agentService.getAgent(intervention.agentId()); + Long workspaceId = agent == null ? null : agent.getWorkspaceId(); + ChatOrigin origin = ChatOrigin.web( + intervention.conversationId(), requester, workspaceId, null); + AgentService.ChatResult result; + try { + result = turnGate.withPermit(permit, () -> agentService.chatWithUsage( + intervention.agentId(), feedback, intervention.conversationId(), origin)); + } catch (RuntimeException error) { + taskService.failTask(taskId, "worker feedback failed: " + safeMessage(error)); + throw error; + } + String reply = persistAssistant(intervention.conversationId(), result); + settleOrPark(intervention, reply); + return taskService.getTask(taskId); + } + } + + private Intervention requireIntervention(Long teamId, Long taskId) { + TeamTaskEntity task = taskService.getTask(taskId); + if (task == null || !teamId.equals(task.getTeamId()) || task.getRunId() == null + || task.getConversationId() == null || task.getConversationId().isBlank()) { + throw new IllegalArgumentException("worker conversation not found for this task"); + } + TeamWorkerConversationContext context = governanceService.resolve( + task.getConversationId(), task.getRunId(), taskId) + .filter(candidate -> teamId.equals(candidate.teamId()) + && Objects.equals(task.getAssigneeAgentId(), candidate.agentId())) + .orElseThrow(() -> new IllegalArgumentException( + "worker conversation not found for this task")); + return new Intervention(task, context.conversationId(), context.agentId()); + } + + private PendingApproval requireReplayApproval(Intervention intervention, String pendingId) { + requireCurrentPendingId(intervention, pendingId); + return approvalService.getPending(pendingId) + .filter(pending -> intervention.conversationId().equals(pending.getConversationId())) + .filter(pending -> "pending".equals(pending.getStatus()) + || "approved".equals(pending.getStatus())) + .or(() -> approvalService.getReplayClaim(pendingId) + .filter(pending -> intervention.conversationId() + .equals(pending.getConversationId()))) + .orElseThrow(() -> new IllegalStateException( + "tool approval is no longer pending or claimed")); + } + + private void requireCurrentPendingId(Intervention intervention, String pendingId) { + if (pendingId == null || pendingId.isBlank()) { + throw new IllegalArgumentException("pending approval id is required"); + } + String currentPendingId = null; + try { + var metadata = JSONUtil.parseObj(intervention.task().getMetadata()); + var approval = metadata.getJSONObject("toolApproval"); + currentPendingId = approval == null ? null : approval.getStr("pendingId"); + } catch (RuntimeException ignored) { + // Missing or malformed task metadata means the client cannot prove + // that this approval is the one the task is parked on. + } + if (!pendingId.equals(currentPendingId)) { + throw new IllegalStateException("tool approval is no longer current for this task"); + } + } + + private PendingApproval claimReplay(Intervention intervention, String pendingId, + String requester, PendingApproval pending) { + if ("pending".equals(pending.getStatus())) { + ResolveOutcome claimed = approvalService.claimForReplay(pendingId, requester); + if (claimed.isAlreadyResolved()) { + throw new IllegalStateException("tool approval was resolved concurrently"); + } + } + return approvalService.getReplayClaim(pendingId) + .filter(candidate -> intervention.conversationId() + .equals(candidate.getConversationId())) + .orElseThrow(() -> new IllegalStateException( + "approved tool replay claim could not be recovered")); + } + + private void requireMemberIdle(Intervention intervention) { + if (taskService.hasActiveTask(intervention.task().getTeamId(), intervention.agentId())) { + throw new IllegalStateException("worker agent is already executing another team task"); + } + } + + private ConversationTurnGate.Permit reserve(String conversationId) { + ConversationTurnGate.Permit permit = turnGate.tryAcquire(conversationId); + if (permit == null || streamTracker.isRunning(conversationId)) { + if (permit != null) { + permit.close(); + } + throw new IllegalStateException("worker conversation is already running"); + } + return permit; + } + + private String persistAssistant(String conversationId, AgentService.ChatResult result) { + String reply = result == null ? "" : result.content(); + if (reply != null && !reply.isBlank()) { + conversationService.saveMessage(conversationId, "assistant", reply, null, "completed", + result.promptTokens(), result.completionTokens(), + result.runtimeModel(), result.runtimeProvider()); + } + return reply; + } + + private void settleOrPark(Intervention intervention, String reply) { + PendingApproval next = approvalService.findPendingByConversation(intervention.conversationId()); + if (next != null) { + String summary = next.getSummary() == null || next.getSummary().isBlank() + ? next.getReason() : next.getSummary(); + taskService.parkForToolApproval(intervention.task().getId(), next.getPendingId(), summary); + eventChannel.publishTaskEvent(taskService.getTask(intervention.task().getId()), + "team_task_awaiting_approval", Map.of("pendingId", next.getPendingId())); + return; + } + dispatchService.settleOutcome(intervention.task(), reply); + dispatchService.requestDispatch(intervention.task().getTeamId()); + } + + private static String safeMessage(RuntimeException error) { + return error.getMessage() == null ? error.getClass().getSimpleName() : error.getMessage(); + } + + private record Intervention(TeamTaskEntity task, String conversationId, Long agentId) { + } +} diff --git a/mateclaw-server/src/main/java/vip/mate/team/service/TeamWorkerReplayPersistenceService.java b/mateclaw-server/src/main/java/vip/mate/team/service/TeamWorkerReplayPersistenceService.java new file mode 100644 index 00000000..b33cbf6c --- /dev/null +++ b/mateclaw-server/src/main/java/vip/mate/team/service/TeamWorkerReplayPersistenceService.java @@ -0,0 +1,34 @@ +package vip.mate.team.service; + +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; +import vip.mate.agent.AgentService; +import vip.mate.workspace.conversation.ConversationService; + +/** Atomically records a replay reply and its task-level idempotency marker. */ +@Service +@RequiredArgsConstructor +public class TeamWorkerReplayPersistenceService { + + private final ConversationService conversationService; + private final TeamTaskService taskService; + + @Transactional + public void persist(Long taskId, String pendingId, String conversationId, + String reply, AgentService.ChatResult result) { + if (reply == null || reply.isBlank()) { + return; + } + if (result == null) { + conversationService.saveMessage(conversationId, "assistant", reply); + } else { + conversationService.saveMessage(conversationId, "assistant", reply, + null, "completed", result.promptTokens(), result.completionTokens(), + result.runtimeModel(), result.runtimeProvider()); + } + if (!taskService.markToolReplayMessagePersisted(taskId, pendingId)) { + throw new IllegalStateException("tool replay message marker could not be persisted"); + } + } +} diff --git a/mateclaw-server/src/test/java/vip/mate/agent/ChatResultCollectorTest.java b/mateclaw-server/src/test/java/vip/mate/agent/ChatResultCollectorTest.java new file mode 100644 index 00000000..0a732f06 --- /dev/null +++ b/mateclaw-server/src/test/java/vip/mate/agent/ChatResultCollectorTest.java @@ -0,0 +1,40 @@ +package vip.mate.agent; + +import org.junit.jupiter.api.Test; +import reactor.core.publisher.Flux; + +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +class ChatResultCollectorTest { + + @Test + void preservesStructuredFinishReasonAlongsideContentAndUsage() { + AgentService.ChatResult result = ChatResultCollector.collect(Flux.just( + new AgentService.StreamDelta("partial failure", null), + AgentService.StreamDelta.event("finish_reason", + Map.of("reason", "error_fallback")), + AgentService.StreamDelta.event("_usage_final", Map.of( + "promptTokens", 12, + "completionTokens", 3, + "runtimeModelName", "model-a", + "runtimeProviderId", "provider-a")))); + + assertEquals("partial failure", result.content()); + assertEquals(12, result.promptTokens()); + assertEquals(3, result.completionTokens()); + assertEquals("model-a", result.runtimeModel()); + assertEquals("provider-a", result.runtimeProvider()); + assertEquals("error_fallback", result.finishReason()); + } + + @Test + void lastFinishReasonWinsForReplayCompatibleStreams() { + AgentService.ChatResult result = ChatResultCollector.collect(Flux.just( + AgentService.StreamDelta.event("finish_reason", Map.of("reason", "incomplete")), + AgentService.StreamDelta.event("finish_reason", Map.of("reason", "normal")))); + + assertEquals("normal", result.finishReason()); + } +} diff --git a/mateclaw-server/src/test/java/vip/mate/agent/graph/ErrorClassificationTest.java b/mateclaw-server/src/test/java/vip/mate/agent/graph/ErrorClassificationTest.java index 3fe7b558..27a0a520 100644 --- a/mateclaw-server/src/test/java/vip/mate/agent/graph/ErrorClassificationTest.java +++ b/mateclaw-server/src/test/java/vip/mate/agent/graph/ErrorClassificationTest.java @@ -246,6 +246,27 @@ class ErrorClassificationTest { classify(new RuntimeException("AccountBalanceNotEnough: balance not enough"))); } + @Test + @DisplayName("Volcengine InvalidSubscription → BILLING instead of CLIENT_ERROR") + void volcengineExpiredCodingPlanIsBilling() throws Exception { + assertEquals(NodeStreamingChatHelper.ErrorType.BILLING, + classify(new RuntimeException("400 InvalidSubscription: CodingPlan subscription has expired"))); + } + + @Test + @DisplayName("DashScope Arrearage / good-standing error → BILLING") + void dashscopeArrearageIsBilling() throws Exception { + assertEquals(NodeStreamingChatHelper.ErrorType.BILLING, + classify(new RuntimeException("400 Arrearage: Access denied, make sure your account is in good standing"))); + } + + @Test + @DisplayName("MiniMax insufficient balance body → BILLING") + void minimaxInsufficientBalanceIsBilling() throws Exception { + assertEquals(NodeStreamingChatHelper.ErrorType.BILLING, + classify(new RuntimeException("insufficient_balance_error: insufficient balance (1008)"))); + } + // ===== Infrastructure-fatal errors → AUTH_ERROR (HARD, no same-model retry) ===== // // DNS / TLS-trust failures do not self-heal on retry. They are routed through diff --git a/mateclaw-server/src/test/java/vip/mate/agent/graph/node/ReasoningNodeOutputTest.java b/mateclaw-server/src/test/java/vip/mate/agent/graph/node/ReasoningNodeOutputTest.java index b4f58a2b..6a1e96ce 100644 --- a/mateclaw-server/src/test/java/vip/mate/agent/graph/node/ReasoningNodeOutputTest.java +++ b/mateclaw-server/src/test/java/vip/mate/agent/graph/node/ReasoningNodeOutputTest.java @@ -190,6 +190,26 @@ class ReasoningNodeOutputTest { "Continuation prompt should ask the model to keep writing instead of ending the run"); } + @Test + @DisplayName("runtime error placeholder cannot satisfy or continue a long-form request") + void longFormTextRequest_runtimeErrorPlaceholderFailsImmediately() throws Exception { + String internalError = "[错误] Bad request: account subscription expired"; + NodeStreamingChatHelper.StreamResult result = new NodeStreamingChatHelper.StreamResult( + internalError, "", new AssistantMessage(internalError), + List.of(), false, 100, 0); + when(streamingHelper.streamCall(any(), any(), anyString(), anyString())).thenReturn(result); + + Map state = baseStateMap(); + state.put(USER_MESSAGE, "请输出不少于 8000 字的技术报告"); + state.put(MAX_ITERATIONS, 150); + Map output = createNode().apply(new OverAllState(state)); + + assertEquals(false, output.get(CONTINUE_REASONING)); + assertEquals("error_fallback", output.get(FINISH_REASON)); + assertEquals(internalError, output.get(FINAL_ANSWER)); + assertNull(output.get("long_form_draft")); + } + @Test @DisplayName("long-form continuation persists all chunks as one final answer") void longFormTextRequest_combinesContinuationChunksInFinalAnswer() throws Exception { diff --git a/mateclaw-server/src/test/java/vip/mate/approval/ApprovalWorkflowServiceResolveTest.java b/mateclaw-server/src/test/java/vip/mate/approval/ApprovalWorkflowServiceResolveTest.java index b19a6469..01e7e70f 100644 --- a/mateclaw-server/src/test/java/vip/mate/approval/ApprovalWorkflowServiceResolveTest.java +++ b/mateclaw-server/src/test/java/vip/mate/approval/ApprovalWorkflowServiceResolveTest.java @@ -208,6 +208,50 @@ class ApprovalWorkflowServiceResolveTest { verifyNoInteractions(conversationService); } + @Test + @DisplayName("team replay claims PENDING before execution and consumes exact APPROVED claim") + void teamReplayClaimIsDurableAndSingleShot() { + PendingApproval pending = seedPending("pid-team", "conv-team", "shell"); + pending.setToolCallPayload("{\"name\":\"shell\"}"); + when(approvalMapper.update(isNull(), any(Wrapper.class))).thenReturn(1); + when(conversationService.markPendingApprovalsResolved( + eq("conv-team"), eq(Set.of("pid-team")), eq(MetadataDecision.APPROVED))) + .thenReturn(1); + + ResolveOutcome claimed = workflow.claimForReplay("pid-team", "alice"); + ResolveOutcome consumed = workflow.consumeReplayClaim("pid-team", "alice"); + + assertThat(claimed.decision()).isEqualTo("approved"); + assertThat(pending.getStatus()).isEqualTo("consumed"); + assertThat(consumed.isConsumed()).isTrue(); + assertThat(consumed.consumedSnapshot().getToolCallPayload()) + .isEqualTo("{\"name\":\"shell\"}"); + assertThat(approvalService.getPending("pid-team")).isEmpty(); + verify(approvalMapper, times(2)).update(isNull(), any(Wrapper.class)); + } + + @Test + @DisplayName("an APPROVED replay claim is recoverable from DB after restart") + void replayClaimRecoversFromDatabase() { + ToolApprovalEntity entity = new ToolApprovalEntity(); + entity.setPendingId("pid-restart"); + entity.setConversationId("conv-restart"); + entity.setUserId("alice"); + entity.setAgentId("201"); + entity.setToolName("shell"); + entity.setToolArguments("{}"); + entity.setToolCallPayload("{\"name\":\"shell\"}"); + entity.setSummary("approved replay"); + entity.setStatus("APPROVED"); + when(approvalMapper.selectOne(any())).thenReturn(entity); + + PendingApproval recovered = workflow.getReplayClaim("pid-restart").orElseThrow(); + + assertThat(recovered.getStatus()).isEqualTo("approved"); + assertThat(recovered.getConversationId()).isEqualTo("conv-restart"); + assertThat(recovered.getToolCallPayload()).isEqualTo("{\"name\":\"shell\"}"); + } + @Test @DisplayName("cancelStalePending issues a SUPERSEDED outcome per pending in the conversation") void cancelStalePendingMultipleEntries() { diff --git a/mateclaw-server/src/test/java/vip/mate/channel/web/ChatStreamTrackerContentBatchTest.java b/mateclaw-server/src/test/java/vip/mate/channel/web/ChatStreamTrackerContentBatchTest.java new file mode 100644 index 00000000..ca844110 --- /dev/null +++ b/mateclaw-server/src/test/java/vip/mate/channel/web/ChatStreamTrackerContentBatchTest.java @@ -0,0 +1,151 @@ +package vip.mate.channel.web; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import org.junit.jupiter.api.Test; +import org.springframework.web.servlet.mvc.method.annotation.ResponseBodyEmitter; +import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; + +import java.io.IOException; +import java.util.List; +import java.util.Set; +import java.util.concurrent.CopyOnWriteArrayList; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class ChatStreamTrackerContentBatchTest { + + private static final ObjectMapper MAPPER = new ObjectMapper(); + + @Test + void adjacentContentDeltasFlushAsOneTimedBatch() throws Exception { + ChatStreamTracker tracker = tracker(25, 256); + CapturingEmitter emitter = attach(tracker, "timed"); + + tracker.broadcast("timed", "content_delta", "{\"delta\":\"你\"}"); + tracker.broadcast("timed", "content_delta", "{\"delta\":\"好\"}"); + + awaitEventCount(emitter, 1); + assertEquals(1, emitter.events.size()); + assertEquals("content_delta", emitter.events.getFirst().name()); + assertEquals("你好", text(emitter.events.getFirst().data(), "delta")); + } + + @Test + void characterLimitFlushesWithoutWaitingForTimer() throws Exception { + ChatStreamTracker tracker = tracker(60_000, 4); + CapturingEmitter emitter = attach(tracker, "bounded"); + + tracker.broadcast("bounded", "content_delta", "{\"delta\":\"ab\"}"); + tracker.broadcast("bounded", "content_delta", "{\"delta\":\"cd\"}"); + + assertEquals(1, emitter.events.size()); + assertEquals("abcd", text(emitter.events.getFirst().data(), "delta")); + } + + @Test + void lifecycleEventFlushesContentFirstAndDoneIsReplayable() throws Exception { + ChatStreamTracker tracker = tracker(60_000, 256); + CapturingEmitter live = attach(tracker, "ordered"); + + tracker.broadcast("ordered", "content_delta", "{\"delta\":\"answer\"}"); + tracker.broadcast("ordered", "phase", "{\"phase\":\"complete\"}"); + tracker.broadcast("ordered", "done", "{\"status\":\"completed\"}"); + + assertEquals(List.of("content_delta", "phase", "done"), live.names()); + + CapturingEmitter replay = new CapturingEmitter(); + assertTrue(tracker.attach("ordered", replay)); + assertEquals(List.of("content_delta", "phase", "done"), replay.names()); + assertEquals("answer", text(replay.events.getFirst().data(), "delta")); + } + + @Test + void webchatTextPayloadKeepsItsWireField() throws Exception { + ChatStreamTracker tracker = tracker(60_000, 4); + ChatStreamTracker.RunHandle handle = tracker.register("webchat"); + CapturingEmitter emitter = new CapturingEmitter(); + tracker.attach(handle, emitter); + + tracker.broadcast(handle, "content_delta", "{\"text\":\"ab\"}"); + tracker.broadcast(handle, "content_delta", "{\"text\":\"cd\"}"); + + assertEquals(1, emitter.events.size()); + assertEquals("abcd", text(emitter.events.getFirst().data(), "text")); + } + + @Test + void lifecycleCompletionFlushesPendingContentEvenWithoutDoneEnvelope() throws Exception { + ChatStreamTracker tracker = tracker(60_000, 256); + ChatStreamTracker.RunHandle handle = tracker.register("complete"); + CapturingEmitter emitter = new CapturingEmitter(); + tracker.attach(handle, emitter); + + tracker.broadcast(handle, "content_delta", "{\"delta\":\"partial\"}"); + tracker.complete(handle); + + assertEquals(1, emitter.events.size()); + assertEquals("partial", text(emitter.events.getFirst().data(), "delta")); + } + + private static ChatStreamTracker tracker(long flushMs, int maxChars) { + ChatStreamTracker tracker = new ChatStreamTracker(MAPPER); + tracker.setContentBatchingForTesting(flushMs, maxChars); + return tracker; + } + + private static CapturingEmitter attach(ChatStreamTracker tracker, String conversationId) { + tracker.register(conversationId); + CapturingEmitter emitter = new CapturingEmitter(); + tracker.attach(conversationId, emitter); + return emitter; + } + + private static void awaitEventCount(CapturingEmitter emitter, int expected) throws InterruptedException { + long deadline = System.currentTimeMillis() + 1_000; + while (emitter.events.size() < expected && System.currentTimeMillis() < deadline) { + Thread.sleep(5); + } + } + + private static String text(String json, String field) throws Exception { + JsonNode node = MAPPER.readTree(json); + return node.path(field).asText(); + } + + private record Event(String name, String data) {} + + private static final class CapturingEmitter extends SseEmitter { + private final List events = new CopyOnWriteArrayList<>(); + + CapturingEmitter() { + super(60_000L); + } + + List names() { + return events.stream().map(Event::name).toList(); + } + + @Override + public void send(SseEventBuilder builder) throws IOException { + Set entries = builder.build(); + String name = ""; + String payload = ""; + boolean expectPayload = false; + for (ResponseBodyEmitter.DataWithMediaType entry : entries) { + if (!(entry.getData() instanceof String text)) continue; + if (text.contains("event:") && text.contains("data:")) { + int start = text.indexOf("event:") + 6; + int end = text.indexOf('\n', start); + name = text.substring(start, end < 0 ? text.length() : end).trim(); + expectPayload = true; + } else if (expectPayload && !"\n\n".equals(text)) { + payload = text; + expectPayload = false; + } + } + events.add(new Event(name, payload)); + } + } +} diff --git a/mateclaw-server/src/test/java/vip/mate/channel/web/ChatStreamTrackerOrphanPolicyTest.java b/mateclaw-server/src/test/java/vip/mate/channel/web/ChatStreamTrackerOrphanPolicyTest.java index d52bf1c6..e0771ee1 100644 --- a/mateclaw-server/src/test/java/vip/mate/channel/web/ChatStreamTrackerOrphanPolicyTest.java +++ b/mateclaw-server/src/test/java/vip/mate/channel/web/ChatStreamTrackerOrphanPolicyTest.java @@ -430,9 +430,13 @@ class ChatStreamTrackerOrphanPolicyTest { Map runs = (Map) runsField.get(tracker); ChatStreamTracker.RunState state = runs.get(cid); + long deadline = System.currentTimeMillis() + 1_000L; + while (state.subscribersZeroSince == null && System.currentTimeMillis() < deadline) { + Thread.sleep(5L); + } synchronized (state.lock) { assertNotNull(state.subscribersZeroSince, - "removing the final dead subscriber must arm the orphan clock"); + "removing the final dead subscriber must arm the orphan clock within the batch window"); state.subscribersZeroSince = System.currentTimeMillis() - 3_000L; } diff --git a/mateclaw-server/src/test/java/vip/mate/cron/service/CronJobLifecycleFenceTest.java b/mateclaw-server/src/test/java/vip/mate/cron/service/CronJobLifecycleFenceTest.java index 397d9503..c81255a8 100644 --- a/mateclaw-server/src/test/java/vip/mate/cron/service/CronJobLifecycleFenceTest.java +++ b/mateclaw-server/src/test/java/vip/mate/cron/service/CronJobLifecycleFenceTest.java @@ -9,6 +9,7 @@ import org.junit.jupiter.api.Test; import org.mockito.ArgumentCaptor; import org.springframework.ai.chat.messages.AssistantMessage; import org.springframework.context.ApplicationEventPublisher; +import vip.mate.agent.AgentService; import vip.mate.cron.model.CronJobEntity; import vip.mate.dashboard.model.CronJobRunEntity; import vip.mate.dashboard.repository.CronJobRunMapper; @@ -53,6 +54,25 @@ class CronJobLifecycleFenceTest { verifyNoInteractions(fixture.conversations, fixture.completionPublisher, fixture.events); } + @Test + void graphErrorPersistsAnErrorMessageWithoutPublishingSuccessEvents() { + Fixture fixture = new Fixture(); + when(fixture.mapper.update(isNull(), any(Wrapper.class))).thenReturn(1); + AgentService.ChatResult failed = new AgentService.ChatResult( + "[错误] account expired", 12, 3, "model-a", "provider-a", "error_fallback"); + + fixture.service.finishRunFailed(run(), new AssistantMessage(failed.content()), + "cron-1", failed); + + @SuppressWarnings("rawtypes") + ArgumentCaptor captor = ArgumentCaptor.forClass(Wrapper.class); + verify(fixture.mapper).update(isNull(), captor.capture()); + assertTrue(captor.getValue().getSqlSegment().contains("status")); + verify(fixture.conversations).saveMessage("cron-1", "assistant", failed.content(), + null, "error", 12, 3, "model-a", "provider-a"); + verifyNoInteractions(fixture.completionPublisher, fixture.events); + } + private static CronJobRunEntity run() { CronJobRunEntity run = new CronJobRunEntity(); run.setId(42L); diff --git a/mateclaw-server/src/test/java/vip/mate/cron/service/CronJobOriginPropagationTest.java b/mateclaw-server/src/test/java/vip/mate/cron/service/CronJobOriginPropagationTest.java index f5cdab56..985f14a5 100644 --- a/mateclaw-server/src/test/java/vip/mate/cron/service/CronJobOriginPropagationTest.java +++ b/mateclaw-server/src/test/java/vip/mate/cron/service/CronJobOriginPropagationTest.java @@ -117,6 +117,39 @@ class CronJobOriginPropagationTest { verify(lifecycle).markRunFailed(eq(run), any(IllegalStateException.class)); } + @Test + void runnerTreatsStructuredGraphErrorAsFailedTerminalState() { + CronJobLifecycleService lifecycle = mock(CronJobLifecycleService.class); + CronRunHeartbeatService heartbeat = mock(CronRunHeartbeatService.class); + CronRunHeartbeatService.Lease lease = mock(CronRunHeartbeatService.Lease.class); + AgentService agentService = mock(AgentService.class); + CronChatOriginFactory originFactory = mock(CronChatOriginFactory.class); + CronConversationResolver resolver = mock(CronConversationResolver.class); + CronJobEntity job = job(); + CronJobRunEntity run = new CronJobRunEntity(); + run.setId(55L); + ChatOrigin origin = ChatOrigin.cron(CONVERSATION_ID, WORKSPACE_ID, null, null, null); + AgentService.ChatResult failed = new AgentService.ChatResult( + "[错误] account expired", 12, 3, "model-a", "provider-a", "error_fallback"); + when(resolver.resolve(job)).thenReturn(CONVERSATION_ID); + when(lifecycle.startRun(job, "do work", "scheduled", CONVERSATION_ID)) + .thenReturn(new CronJobLifecycleService.StartResult(run, MESSAGE_ID)); + when(originFactory.from(job, CONVERSATION_ID, MESSAGE_ID)).thenReturn(origin); + when(heartbeat.begin(55L)).thenReturn(lease); + when(agentService.chatWithUsage(eq(AGENT_ID), anyString(), eq(CONVERSATION_ID), eq(origin))) + .thenReturn(failed); + CronJobRunner runner = new CronJobRunner(lifecycle, heartbeat, agentService, originFactory, resolver, + mock(WikiProcessingService.class), new ObjectMapper()); + + runner.executeJob(job); + + verify(lifecycle).finishRunFailed(eq(run), any(org.springframework.ai.chat.messages.AssistantMessage.class), + eq(CONVERSATION_ID), eq(failed)); + verify(lifecycle, never()).finishRunAndPublish(eq(job), eq(run), anyString(), + any(org.springframework.ai.chat.messages.AssistantMessage.class), eq(CONVERSATION_ID), + any(Boolean.class), eq(failed)); + } + private static CronJobEntity job() { CronJobEntity job = new CronJobEntity(); job.setId(JOB_ID); diff --git a/mateclaw-server/src/test/java/vip/mate/llm/service/ModelDiscoveryServiceTestPromptTest.java b/mateclaw-server/src/test/java/vip/mate/llm/service/ModelDiscoveryServiceTestPromptTest.java index 05c742eb..18fd98b0 100644 --- a/mateclaw-server/src/test/java/vip/mate/llm/service/ModelDiscoveryServiceTestPromptTest.java +++ b/mateclaw-server/src/test/java/vip/mate/llm/service/ModelDiscoveryServiceTestPromptTest.java @@ -100,4 +100,13 @@ class ModelDiscoveryServiceTestPromptTest { Map requestBodyFromNull = ModelDiscoveryService.buildTestPromptRequestBody("gpt-4-turbo", null); assertEquals(requestBody, requestBodyFromNull); } + + @Test + @DisplayName("Kimi smoke test uses the same fixed temperature required by runtime") + void kimiForCoding_usesTemperatureOne() { + Map requestBody = ModelDiscoveryService.buildTestPromptRequestBody( + "kimi-for-coding", Map.of("temperature", 0.2)); + + assertEquals(1.0d, requestBody.get("temperature")); + } } diff --git a/mateclaw-server/src/test/java/vip/mate/team/controller/TeamControllerTest.java b/mateclaw-server/src/test/java/vip/mate/team/controller/TeamControllerTest.java index 00b56999..01b713e5 100644 --- a/mateclaw-server/src/test/java/vip/mate/team/controller/TeamControllerTest.java +++ b/mateclaw-server/src/test/java/vip/mate/team/controller/TeamControllerTest.java @@ -29,6 +29,7 @@ import vip.mate.team.service.TeamManualTaskService; import vip.mate.team.service.TeamRunService; import vip.mate.team.service.TeamService; import vip.mate.team.service.TeamTaskService; +import vip.mate.team.service.TeamWorkerInterventionService; import vip.mate.workspace.core.annotation.RequireWorkspaceRole; import vip.mate.workspace.core.service.WorkspaceService; @@ -69,6 +70,7 @@ class TeamControllerTest { @Mock private TeamDispatchService dispatchService; @Mock private TeamAnnounceService announceService; @Mock private TeamEventChannel eventChannel; + @Mock private TeamWorkerInterventionService workerInterventionService; @Mock private AgentMapper agentMapper; @Mock private WorkspaceService workspaceService; @Mock private AuthService authService; @@ -80,7 +82,7 @@ class TeamControllerTest { void setUp() { manualTaskService = new TeamManualTaskService(runService, taskService, events); controller = new TeamController(teamService, taskService, manualTaskService, dispatchService, - announceService, eventChannel, agentMapper); + announceService, eventChannel, workerInterventionService, agentMapper); AgentTeamEntity team = new AgentTeamEntity(); team.setId(TEAM_ID); team.setWorkspaceId(1L); @@ -112,6 +114,54 @@ class TeamControllerTest { return run; } + @Test + void workerApprovalEndpointUsesTaskScopedInterventionService() { + TeamTaskEntity waiting = task(TEAM_ID, TeamTaskStatus.AWAITING_APPROVAL); + waiting.setRunId(RUN_ID); + when(taskService.getTask(TASK_ID)).thenReturn(waiting); + when(workerInterventionService.approve(TEAM_ID, TASK_ID, "pending-42", "alice")) + .thenReturn(waiting); + TeamController.WorkerApprovalRequest request = new TeamController.WorkerApprovalRequest(); + request.setPendingId("pending-42"); + + R response = controller.approveWorkerTool( + TEAM_ID, TASK_ID, request, () -> "alice"); + + assertEquals(200, response.getCode()); + verify(workerInterventionService).approve(TEAM_ID, TASK_ID, "pending-42", "alice"); + } + + @Test + void workerFeedbackEndpointRejectsBlankContentBeforeRunningAgent() { + when(taskService.getTask(TASK_ID)).thenReturn(task(TEAM_ID, TeamTaskStatus.COMPLETED)); + TeamController.WorkerFeedbackRequest request = new TeamController.WorkerFeedbackRequest(); + request.setMessage(" "); + + R response = controller.feedbackWorker( + TEAM_ID, TASK_ID, request, () -> "alice"); + + assertEquals(400, response.getCode()); + verify(workerInterventionService, never()).feedback(any(), any(), any(), any()); + } + + @Test + void workerInterventionMapsMissingLinkAndBusyConversationToActionableCodes() { + TeamTaskEntity waiting = task(TEAM_ID, TeamTaskStatus.AWAITING_APPROVAL); + waiting.setRunId(RUN_ID); + when(taskService.getTask(TASK_ID)).thenReturn(waiting); + TeamController.WorkerApprovalRequest request = new TeamController.WorkerApprovalRequest(); + request.setPendingId("pending-42"); + when(workerInterventionService.approve(TEAM_ID, TASK_ID, "pending-42", "alice")) + .thenThrow(new IllegalArgumentException("worker conversation not found for this task")); + when(workerInterventionService.deny(TEAM_ID, TASK_ID, "pending-42", "alice")) + .thenThrow(new IllegalStateException("worker conversation is already running")); + + assertEquals(404, controller.approveWorkerTool( + TEAM_ID, TASK_ID, request, () -> "alice").getCode()); + assertEquals(409, controller.denyWorkerTool( + TEAM_ID, TASK_ID, request, () -> "alice").getCode()); + } + // ==================== team / membership ==================== @Test diff --git a/mateclaw-server/src/test/java/vip/mate/team/service/TeamDispatchServiceEventTest.java b/mateclaw-server/src/test/java/vip/mate/team/service/TeamDispatchServiceEventTest.java index 88b2de7a..4985798f 100644 --- a/mateclaw-server/src/test/java/vip/mate/team/service/TeamDispatchServiceEventTest.java +++ b/mateclaw-server/src/test/java/vip/mate/team/service/TeamDispatchServiceEventTest.java @@ -12,6 +12,7 @@ import org.springframework.transaction.support.AbstractPlatformTransactionManage import org.springframework.transaction.support.DefaultTransactionStatus; import org.springframework.transaction.support.TransactionTemplate; import vip.mate.agent.AgentService; +import vip.mate.approval.ApprovalWorkflowService; import vip.mate.channel.web.ChatStreamTracker; import vip.mate.team.event.TeamTasksDelegatedEvent; import vip.mate.workspace.conversation.ConversationService; @@ -65,7 +66,8 @@ class TeamDispatchServiceEventTest { return new TeamDispatchService( mock(TeamService.class), taskService, mock(AgentService.class), mock(ConversationService.class), mock(ChatStreamTracker.class), - mock(TeamAnnounceService.class), mock(TeamEventChannel.class)); + mock(TeamAnnounceService.class), mock(TeamEventChannel.class), + mock(ApprovalWorkflowService.class)); } } diff --git a/mateclaw-server/src/test/java/vip/mate/team/service/TeamDispatchServiceTest.java b/mateclaw-server/src/test/java/vip/mate/team/service/TeamDispatchServiceTest.java index 35c71f2e..84356046 100644 --- a/mateclaw-server/src/test/java/vip/mate/team/service/TeamDispatchServiceTest.java +++ b/mateclaw-server/src/test/java/vip/mate/team/service/TeamDispatchServiceTest.java @@ -4,6 +4,8 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; import vip.mate.agent.AgentService; +import vip.mate.approval.ApprovalWorkflowService; +import vip.mate.approval.PendingApproval; import vip.mate.channel.web.ChatStreamTracker; import vip.mate.team.model.AgentTeamEntity; import vip.mate.team.model.TeamTaskEntity; @@ -43,6 +45,7 @@ class TeamDispatchServiceTest { private ChatStreamTracker streamTracker; private TeamAnnounceService announceService; private TeamEventChannel eventChannel; + private ApprovalWorkflowService approvalService; private TeamDispatchService service; @BeforeEach @@ -54,8 +57,9 @@ class TeamDispatchServiceTest { streamTracker = mock(ChatStreamTracker.class); announceService = mock(TeamAnnounceService.class); eventChannel = mock(TeamEventChannel.class); + approvalService = mock(ApprovalWorkflowService.class); service = new TeamDispatchService(teamService, taskService, agentService, - conversationService, streamTracker, announceService, eventChannel); + conversationService, streamTracker, announceService, eventChannel, approvalService); } private TeamTaskEntity task(Long id, Long assignee) { @@ -365,6 +369,28 @@ class TeamDispatchServiceTest { verify(conversationService).saveMessage(startsWith("team-task-"), eq("assistant"), eq("all done")); } + @Test + @DisplayName("a worker tool approval parks the task instead of completing or retrying it") + void runTaskParksPendingToolApproval() { + TeamTaskEntity assigned = task(1L, MEMBER_A); + assigned.setStatus(TeamTaskStatus.IN_PROGRESS); + PendingApproval pending = new PendingApproval("pending-42", "worker", "system", + "execute_shell_command", "{}", "shell command requires approval"); + pending.setSummary("shell command requires approval"); + when(agentService.chatWithUsage(eq(MEMBER_A), anyString(), anyString())) + .thenReturn(AgentService.ChatResult.contentOnly("I need permission first.")); + when(approvalService.findPendingByConversation(startsWith("team-task-"))).thenReturn(pending); + when(taskService.parkForToolApproval(1L, "pending-42", "shell command requires approval")) + .thenReturn(true); + + service.runTask(TEAM_ID, assigned); + + verify(taskService).parkForToolApproval(1L, "pending-42", "shell command requires approval"); + verify(taskService, never()).completeTask(any(), any(), anyString()); + verify(taskService, never()).requeueUnusableResult(any(), anyString()); + verify(eventChannel).publishTaskEvent(any(), eq("team_task_awaiting_approval"), any()); + } + @Test @DisplayName("member child conversation inherits the team's workspace") void runTaskCreatesChildConversationInTeamWorkspace() { diff --git a/mateclaw-server/src/test/java/vip/mate/team/service/TeamRunStateMachineTest.java b/mateclaw-server/src/test/java/vip/mate/team/service/TeamRunStateMachineTest.java index 7f1d9fb4..c9f9212b 100644 --- a/mateclaw-server/src/test/java/vip/mate/team/service/TeamRunStateMachineTest.java +++ b/mateclaw-server/src/test/java/vip/mate/team/service/TeamRunStateMachineTest.java @@ -47,6 +47,8 @@ class TeamRunStateMachineTest { TeamRunStatus.RUNNING, null, 1, 0, 0, 50), Arguments.of("blocked tasks are active", TeamRunStatus.AWAITING_REVIEW, tasks(TeamTaskStatus.BLOCKED), TeamRunStatus.RUNNING, null, 0, 0, 0, 0), + Arguments.of("tool approval waits are active", TeamRunStatus.FINALIZING, + tasks(TeamTaskStatus.AWAITING_APPROVAL), TeamRunStatus.RUNNING, null, 0, 0, 0, 0), Arguments.of("review only", TeamRunStatus.RUNNING, tasks(TeamTaskStatus.COMPLETED, TeamTaskStatus.IN_REVIEW), TeamRunStatus.AWAITING_REVIEW, null, 1, 0, 1, 50), diff --git a/mateclaw-server/src/test/java/vip/mate/team/service/TeamRunViewFactoryTest.java b/mateclaw-server/src/test/java/vip/mate/team/service/TeamRunViewFactoryTest.java index 7f6ed5df..8a6e1df7 100644 --- a/mateclaw-server/src/test/java/vip/mate/team/service/TeamRunViewFactoryTest.java +++ b/mateclaw-server/src/test/java/vip/mate/team/service/TeamRunViewFactoryTest.java @@ -81,6 +81,39 @@ class TeamRunViewFactoryTest { assertEquals(List.of(), view.attentionItems()); } + @Test + void awaitingToolApprovalCreatesHighestPriorityAttentionItem() { + TeamTaskEntity task = task(101L, 201L, + "{\"toolApproval\":{\"pendingId\":\"pending-42\"," + + "\"summary\":\"shell command requires approval\"}}"); + task.setStatus(TeamTaskStatus.AWAITING_APPROVAL); + task.setReason("shell command requires approval"); + + TeamRunView view = project(run("{}"), List.of(task)); + + assertEquals(1, view.attentionItems().size()); + TeamRunView.AttentionItem item = view.attentionItems().getFirst(); + assertEquals("approval", item.type()); + assertEquals("action", item.severity()); + assertEquals(0, item.priority()); + assertEquals(101L, item.taskId()); + assertEquals("shell command requires approval", item.message()); + } + + @Test + void uncertainReplayProjectsOnlyTheSafeRecoveryAttentionType() { + TeamTaskEntity task = task(101L, 201L, + "{\"toolApproval\":{\"pendingId\":\"pending-42\"," + + "\"replayOutcomeUncertain\":true}}"); + task.setStatus(TeamTaskStatus.AWAITING_APPROVAL); + task.setReason("verify the tool outcome manually"); + + TeamRunView view = project(run("{}"), List.of(task)); + + assertEquals("replay_uncertain", view.attentionItems().getFirst().type()); + assertEquals("action", view.attentionItems().getFirst().severity()); + } + @Test void aggregatesRunOnlyDeliverables() { TeamRunEntity run = run("{\"deliverables\":[{\"name\":\"report.pdf\"," diff --git a/mateclaw-server/src/test/java/vip/mate/team/service/TeamTaskServiceTest.java b/mateclaw-server/src/test/java/vip/mate/team/service/TeamTaskServiceTest.java index a8839b16..6b4c97fe 100644 --- a/mateclaw-server/src/test/java/vip/mate/team/service/TeamTaskServiceTest.java +++ b/mateclaw-server/src/test/java/vip/mate/team/service/TeamTaskServiceTest.java @@ -390,6 +390,157 @@ class TeamTaskServiceTest { verify(projectionScheduler).scheduleTask(5L); } + @Test + @DisplayName("tool approval parks an in-progress task without releasing dependents") + void toolApprovalParksTask() { + TeamTaskEntity running = runTask(5L, TeamTaskStatus.IN_PROGRESS); + running.setMetadata("{\"deliverableRequired\":true}"); + when(taskMapper.selectById(5L)).thenReturn(running); + when(taskMapper.update(isNull(), any())).thenReturn(1); + + assertTrue(service.parkForToolApproval(5L, "pending-42", "shell command requires approval")); + + ArgumentCaptor> captor = + ArgumentCaptor.forClass(LambdaUpdateWrapper.class); + verify(taskMapper).update(isNull(), captor.capture()); + var values = captor.getValue().getParamNameValuePairs().values(); + assertTrue(values.contains("awaiting_approval")); + assertTrue(values.stream().anyMatch(value -> String.valueOf(value).contains("pending-42"))); + assertTrue(values.stream().anyMatch(value -> String.valueOf(value).contains("deliverableRequired")), + "parking must merge approval context into existing metadata"); + verify(taskMapper, never()).selectList(any()); + verify(projectionScheduler).scheduleTask(5L); + } + + @Test + @DisplayName("tool approval resume requires the exact pending id and records the replay lease") + void toolApprovalResumeUsesExactPendingId() { + TeamTaskEntity waiting = runTask(5L, TeamTaskStatus.AWAITING_APPROVAL); + waiting.setAssigneeAgentId(MEMBER_ID); + waiting.setMetadata("{\"deliverableRequired\":true,\"toolApproval\":{\"pendingId\":\"pending-42\"}}"); + when(taskMapper.selectById(5L)).thenReturn(waiting); + when(taskMapper.update(isNull(), any())).thenReturn(1); + + assertTrue(service.resumeAfterToolApproval(5L, "pending-42")); + + ArgumentCaptor> captor = + ArgumentCaptor.forClass(LambdaUpdateWrapper.class); + verify(taskMapper).update(isNull(), captor.capture()); + var values = captor.getValue().getParamNameValuePairs().values(); + assertTrue(values.contains(TeamTaskStatus.IN_PROGRESS)); + assertTrue(values.stream().anyMatch(value -> String.valueOf(value).contains("deliverableRequired"))); + assertTrue(values.stream().anyMatch(value -> String.valueOf(value).contains("replayInProgress"))); + verify(projectionScheduler).scheduleTask(5L); + } + + @Test + @DisplayName("a stale pending id cannot resume a worker task") + void staleToolApprovalCannotResumeTask() { + TeamTaskEntity waiting = runTask(5L, TeamTaskStatus.AWAITING_APPROVAL); + waiting.setMetadata("{\"toolApproval\":{\"pendingId\":\"pending-new\"}}"); + when(taskMapper.selectById(5L)).thenReturn(waiting); + + IllegalStateException error = assertThrows(IllegalStateException.class, + () -> service.resumeAfterToolApproval(5L, "pending-old")); + + assertTrue(error.getMessage().contains("no longer current")); + verify(taskMapper, never()).update(isNull(), any()); + } + + @Test + @DisplayName("a replay result is staged under the exact approval before it is consumed") + void stagesToolReplayResultForCrashSafeFinalization() { + TeamTaskEntity running = runTask(5L, TeamTaskStatus.IN_PROGRESS); + running.setMetadata("{\"deliverableRequired\":true," + + "\"toolApproval\":{\"pendingId\":\"pending-42\"}}"); + when(taskMapper.selectById(5L)).thenReturn(running); + when(taskMapper.update(isNull(), any())).thenReturn(1); + + assertTrue(service.stageToolReplayResult(5L, "pending-42", "tool completed")); + + ArgumentCaptor> captor = + ArgumentCaptor.forClass(LambdaUpdateWrapper.class); + verify(taskMapper).update(isNull(), captor.capture()); + var values = captor.getValue().getParamNameValuePairs().values(); + assertTrue(values.contains(TeamTaskStatus.AWAITING_APPROVAL)); + assertTrue(values.stream().anyMatch(value -> String.valueOf(value).contains("pending-42"))); + assertTrue(values.stream().anyMatch(value -> String.valueOf(value).contains("tool completed"))); + assertTrue(values.stream().anyMatch(value -> String.valueOf(value).contains("deliverableRequired"))); + verify(projectionScheduler).scheduleTask(5L); + } + + @Test + @DisplayName("a claimed replay can be stopped after failure but not after a result is staged") + void abortClaimedReplayHasExplicitGuard() { + TeamTaskEntity waiting = runTask(5L, TeamTaskStatus.AWAITING_APPROVAL); + waiting.setMetadata("{\"toolApproval\":{\"pendingId\":\"pending-42\"}}"); + when(taskMapper.selectById(5L)).thenReturn(waiting); + when(taskMapper.update(isNull(), any())).thenReturn(1); + + assertTrue(service.abortClaimedToolReplay(5L, "pending-42", "alice")); + + waiting.setMetadata("{\"toolApproval\":{\"pendingId\":\"pending-42\"," + + "\"replayResult\":\"done\"}}"); + assertThrows(IllegalStateException.class, + () -> service.abortClaimedToolReplay(5L, "pending-42", "alice")); + } + + @Test + @DisplayName("an expired replay lease becomes uncertain instead of an ordinary retryable stale task") + void expiredReplayLeaseRequiresManualResolution() { + TeamTaskEntity replay = runTask(5L, TeamTaskStatus.IN_PROGRESS); + replay.setMetadata("{\"toolApproval\":{\"pendingId\":\"pending-42\"," + + "\"replayInProgress\":true}}"); + when(taskMapper.selectList(any())).thenReturn(List.of(replay)); + when(taskMapper.update(isNull(), any())).thenReturn(1); + + service.recoverStaleTasks(); + + ArgumentCaptor> captor = + ArgumentCaptor.forClass(LambdaUpdateWrapper.class); + verify(taskMapper).update(isNull(), captor.capture()); + var values = captor.getValue().getParamNameValuePairs().values(); + assertTrue(values.contains(TeamTaskStatus.AWAITING_APPROVAL)); + assertTrue(values.stream().anyMatch(value -> String.valueOf(value) + .contains("replayOutcomeUncertain"))); + assertFalse(values.contains(TeamTaskStatus.STALE)); + } + + @Test + @DisplayName("a replay exception is immediately parked as outcome-uncertain") + void replayExceptionCannotBeAutomaticallyRetried() { + TeamTaskEntity replay = runTask(5L, TeamTaskStatus.IN_PROGRESS); + replay.setMetadata("{\"toolApproval\":{\"pendingId\":\"pending-42\"," + + "\"replayInProgress\":true}}"); + when(taskMapper.selectById(5L)).thenReturn(replay); + when(taskMapper.update(isNull(), any())).thenReturn(1); + + assertTrue(service.parkToolReplayUncertain( + 5L, "pending-42", "provider failed; outcome uncertain")); + + ArgumentCaptor> captor = + ArgumentCaptor.forClass(LambdaUpdateWrapper.class); + verify(taskMapper).update(isNull(), captor.capture()); + var values = captor.getValue().getParamNameValuePairs().values(); + assertTrue(values.contains(TeamTaskStatus.AWAITING_APPROVAL)); + assertTrue(values.stream().anyMatch(value -> String.valueOf(value) + .contains("replayOutcomeUncertain"))); + } + + @Test + @DisplayName("feedback reopens a settled worker task but not an active or approval-blocked task") + void feedbackResumeHasExplicitStateGuard() { + TeamTaskEntity completed = runTask(5L, TeamTaskStatus.COMPLETED); + completed.setAssigneeAgentId(MEMBER_ID); + when(taskMapper.selectById(5L)).thenReturn(completed); + when(taskMapper.update(isNull(), any())).thenReturn(1); + + assertTrue(service.resumeForWorkerFeedback(5L)); + + completed.setStatus(TeamTaskStatus.AWAITING_APPROVAL); + assertThrows(IllegalStateException.class, () -> service.resumeForWorkerFeedback(5L)); + } + // ==================== blocker comment ==================== @Test diff --git a/mateclaw-server/src/test/java/vip/mate/team/service/TeamWorkerInterventionServiceTest.java b/mateclaw-server/src/test/java/vip/mate/team/service/TeamWorkerInterventionServiceTest.java new file mode 100644 index 00000000..f5e8934f --- /dev/null +++ b/mateclaw-server/src/test/java/vip/mate/team/service/TeamWorkerInterventionServiceTest.java @@ -0,0 +1,284 @@ +package vip.mate.team.service; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import vip.mate.agent.AgentService; +import vip.mate.agent.context.ChatOrigin; +import vip.mate.agent.runtime.ConversationTurnGate; +import vip.mate.approval.ApprovalWorkflowService; +import vip.mate.approval.PendingApproval; +import vip.mate.approval.ResolveOutcome; +import vip.mate.channel.web.ChatStreamTracker; +import vip.mate.team.model.TeamTaskEntity; +import vip.mate.team.model.TeamTaskStatus; +import vip.mate.workspace.conversation.ConversationService; + +import java.util.Optional; + +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +class TeamWorkerInterventionServiceTest { + + private final TeamTaskService taskService = mock(TeamTaskService.class); + private final TeamWorkerConversationGovernanceService governance = + mock(TeamWorkerConversationGovernanceService.class); + private final ApprovalWorkflowService approvalService = mock(ApprovalWorkflowService.class); + private final AgentService agentService = mock(AgentService.class); + private final ConversationService conversationService = mock(ConversationService.class); + private final ChatStreamTracker streamTracker = mock(ChatStreamTracker.class); + private final TeamDispatchService dispatchService = mock(TeamDispatchService.class); + private final TeamAnnounceService announceService = mock(TeamAnnounceService.class); + private final TeamEventChannel eventChannel = mock(TeamEventChannel.class); + private final TeamWorkerReplayPersistenceService replayPersistenceService = + mock(TeamWorkerReplayPersistenceService.class); + private TeamWorkerInterventionService service; + + @BeforeEach + void setUp() { + service = new TeamWorkerInterventionService(taskService, governance, approvalService, + agentService, conversationService, new ConversationTurnGate(), streamTracker, + dispatchService, announceService, eventChannel, replayPersistenceService); + } + + @Test + void approvalReplaysInCanonicalConversationAndSettlesOriginalTask() { + TeamTaskEntity task = task(TeamTaskStatus.AWAITING_APPROVAL); + PendingApproval pending = pending("pending-42"); + when(taskService.getTask(101L)).thenReturn(task); + when(governance.resolve("worker-101", 11L, 101L)).thenReturn(Optional.of(context())); + when(approvalService.getPending("pending-42")).thenReturn(Optional.of(pending)); + when(approvalService.claimForReplay("pending-42", "alice")) + .thenReturn(ResolveOutcome.resolved(pending, "approved", true, 1)); + when(approvalService.getReplayClaim("pending-42")).thenReturn(Optional.of(pending)); + when(approvalService.consumeReplayClaim("pending-42", "alice")) + .thenReturn(ResolveOutcome.consumed(pending, true, 1)); + when(taskService.resumeAfterToolApproval(101L, "pending-42")).thenReturn(true); + when(taskService.stageToolReplayResult(101L, "pending-42", "tool completed")).thenReturn(true); + when(approvalService.restoreChatOrigin(null)).thenReturn(ChatOrigin.EMPTY); + when(agentService.chatWithReplayWithUsage(eq(201L), any(), eq("worker-101"), + eq("{\"name\":\"shell\"}"), eq(ChatOrigin.EMPTY))) + .thenReturn(AgentService.ChatResult.contentOnly("tool completed")); + + service.approve(7L, 101L, "pending-42", "alice"); + + verify(conversationService).removeApprovalPlaceholders("worker-101"); + verify(replayPersistenceService).persist(101L, "pending-42", "worker-101", + "tool completed", AgentService.ChatResult.contentOnly("tool completed")); + verify(dispatchService).settleOutcome(task, "tool completed"); + verify(approvalService).consumeReplayClaim("pending-42", "alice"); + } + + @Test + void canonicalLinkMismatchRejectsBeforeApprovalMutation() { + TeamTaskEntity task = task(TeamTaskStatus.AWAITING_APPROVAL); + when(taskService.getTask(101L)).thenReturn(task); + when(governance.resolve("worker-101", 11L, 101L)).thenReturn(Optional.empty()); + + IllegalArgumentException error = assertThrows(IllegalArgumentException.class, + () -> service.approve(7L, 101L, "pending-42", "alice")); + + assertTrue(error.getMessage().contains("worker conversation")); + verify(approvalService, never()).consumeReplayClaim(any(), any()); + } + + @Test + void denialSettlesWithoutExecutingTheTool() { + TeamTaskEntity task = task(TeamTaskStatus.AWAITING_APPROVAL); + PendingApproval pending = pending("pending-42"); + when(taskService.getTask(101L)).thenReturn(task); + when(governance.resolve("worker-101", 11L, 101L)).thenReturn(Optional.of(context())); + when(approvalService.getPending("pending-42")).thenReturn(Optional.of(pending)); + when(approvalService.resolve("pending-42", "alice", "denied")) + .thenReturn(ResolveOutcome.resolved(pending, "denied", true, 1)); + when(taskService.denyToolApproval(101L, "pending-42", "alice")).thenReturn(true); + + service.deny(7L, 101L, "pending-42", "alice"); + + verify(agentService, never()).chatWithReplayWithUsage(any(), any(), any(), any(), any()); + verify(taskService).denyToolApproval(101L, "pending-42", "alice"); + verify(announceService).announceTaskSettled(task); + } + + @Test + void feedbackContinuesOriginalConversationAndCannotBypassPendingApproval() { + TeamTaskEntity completed = task(TeamTaskStatus.COMPLETED); + when(taskService.getTask(101L)).thenReturn(completed); + when(governance.resolve("worker-101", 11L, 101L)).thenReturn(Optional.of(context())); + when(taskService.resumeForWorkerFeedback(101L)).thenReturn(true); + when(agentService.chatWithUsage(eq(201L), eq("tighten the summary"), eq("worker-101"), any())) + .thenReturn(AgentService.ChatResult.contentOnly("revised summary")); + + service.feedback(7L, 101L, "tighten the summary", "alice"); + + verify(conversationService).saveMessage("worker-101", "user", "tighten the summary"); + verify(conversationService).saveMessage("worker-101", "assistant", "revised summary", + null, "completed", 0, 0, null, null); + verify(dispatchService).settleOutcome(completed, "revised summary"); + + when(approvalService.findPendingByConversation("worker-101")) + .thenReturn(pending("pending-next")); + assertThrows(IllegalStateException.class, + () -> service.feedback(7L, 101L, "run another command", "alice")); + } + + @Test + void replayFailureReparksApprovedPayloadForSafeRetry() { + TeamTaskEntity task = task(TeamTaskStatus.AWAITING_APPROVAL); + PendingApproval pending = pending("pending-42"); + when(taskService.getTask(101L)).thenReturn(task); + when(governance.resolve("worker-101", 11L, 101L)).thenReturn(Optional.of(context())); + when(approvalService.getPending("pending-42")).thenReturn(Optional.of(pending)); + when(approvalService.claimForReplay("pending-42", "alice")) + .thenReturn(ResolveOutcome.resolved(pending, "approved", true, 1)); + when(approvalService.getReplayClaim("pending-42")).thenReturn(Optional.of(pending)); + when(taskService.resumeAfterToolApproval(101L, "pending-42")).thenReturn(true); + when(approvalService.restoreChatOrigin(null)).thenReturn(ChatOrigin.EMPTY); + when(agentService.chatWithReplayWithUsage(any(), any(), any(), any(), any())) + .thenThrow(new IllegalStateException("provider timeout")); + + assertThrows(IllegalStateException.class, + () -> service.approve(7L, 101L, "pending-42", "alice")); + + verify(taskService).parkToolReplayUncertain(eq(101L), eq("pending-42"), + org.mockito.ArgumentMatchers.contains("provider timeout")); + verify(taskService, never()).failTask(eq(101L), any()); + verify(approvalService, never()).consumeReplayClaim(any(), any()); + } + + @Test + void onlyTheInstanceHoldingTheTaskReplayLeaseExecutesTheTool() { + TeamTaskEntity task = task(TeamTaskStatus.AWAITING_APPROVAL); + PendingApproval pending = pending("pending-42"); + when(taskService.getTask(101L)).thenReturn(task); + when(governance.resolve("worker-101", 11L, 101L)).thenReturn(Optional.of(context())); + when(approvalService.getPending("pending-42")).thenReturn(Optional.of(pending)); + when(approvalService.claimForReplay("pending-42", "alice")) + .thenReturn(ResolveOutcome.resolved(pending, "approved", true, 1)); + when(approvalService.getReplayClaim("pending-42")).thenReturn(Optional.of(pending)); + when(taskService.resumeAfterToolApproval(101L, "pending-42")).thenReturn(false); + + assertThrows(IllegalStateException.class, + () -> service.approve(7L, 101L, "pending-42", "alice")); + + verify(agentService, never()).chatWithReplayWithUsage(any(), any(), any(), any(), any()); + } + + @Test + void uncertainReplayOutcomeCannotBeExecutedAgainAutomatically() { + TeamTaskEntity task = task(TeamTaskStatus.AWAITING_APPROVAL); + when(taskService.getTask(101L)).thenReturn(task); + when(governance.resolve("worker-101", 11L, 101L)).thenReturn(Optional.of(context())); + when(taskService.isToolReplayOutcomeUncertain(task)).thenReturn(true); + + assertThrows(IllegalStateException.class, + () -> service.approve(7L, 101L, "pending-42", "alice")); + + verify(approvalService, never()).claimForReplay(any(), any()); + verify(agentService, never()).chatWithReplayWithUsage(any(), any(), any(), any(), any()); + } + + @Test + void stagedReplayFinalizationDoesNotExecuteToolAgain() { + TeamTaskEntity task = task(TeamTaskStatus.AWAITING_APPROVAL); + task.setMetadata("{\"toolApproval\":{\"pendingId\":\"pending-42\"," + + "\"replayResult\":\"tool completed\",\"messagePersisted\":true}}"); + PendingApproval pending = pending("pending-42"); + when(taskService.getTask(101L)).thenReturn(task); + when(governance.resolve("worker-101", 11L, 101L)).thenReturn(Optional.of(context())); + when(approvalService.getPending("pending-42")).thenReturn(Optional.of(pending)); + when(taskService.stagedToolReplayResult(task)).thenReturn("tool completed"); + when(taskService.isToolReplayMessagePersisted(task)).thenReturn(true); + when(approvalService.consumeReplayClaim("pending-42", "alice")) + .thenReturn(ResolveOutcome.consumed(pending, true, 1)); + when(taskService.resumeAfterToolApproval(101L, "pending-42")).thenReturn(true); + + service.approve(7L, 101L, "pending-42", "alice"); + + verify(agentService, never()).chatWithReplayWithUsage(any(), any(), any(), any(), any()); + verify(conversationService, never()).removeApprovalPlaceholders(anyString()); + verify(dispatchService).settleOutcome(task, "tool completed"); + } + + @Test + void stagedReplayCannotBeReportedAsDeniedAfterToolExecution() { + TeamTaskEntity task = task(TeamTaskStatus.AWAITING_APPROVAL); + task.setMetadata("{\"toolApproval\":{\"pendingId\":\"pending-42\"," + + "\"replayResult\":\"tool completed\"}}"); + when(taskService.getTask(101L)).thenReturn(task); + when(governance.resolve("worker-101", 11L, 101L)).thenReturn(Optional.of(context())); + when(taskService.stagedToolReplayResult(task)).thenReturn("tool completed"); + + assertThrows(IllegalStateException.class, + () -> service.deny(7L, 101L, "pending-42", "alice")); + + verify(approvalService, never()).resolve(any(), any(), anyString()); + verify(taskService, never()).denyToolApproval(any(), any(), any()); + } + + @Test + void claimedReplayCanBeStoppedAfterExecutionFailure() { + TeamTaskEntity task = task(TeamTaskStatus.AWAITING_APPROVAL); + PendingApproval claimed = pending("pending-42"); + claimed.setStatus("approved"); + when(taskService.getTask(101L)).thenReturn(task); + when(governance.resolve("worker-101", 11L, 101L)).thenReturn(Optional.of(context())); + when(approvalService.getPending("pending-42")).thenReturn(Optional.of(claimed)); + when(approvalService.consumeReplayClaim("pending-42", "alice")) + .thenReturn(ResolveOutcome.consumed(claimed, true, 1)); + when(taskService.abortClaimedToolReplay(101L, "pending-42", "alice")) + .thenReturn(true); + + service.deny(7L, 101L, "pending-42", "alice"); + + verify(taskService).abortClaimedToolReplay(101L, "pending-42", "alice"); + verify(taskService, never()).denyToolApproval(any(), any(), any()); + verify(eventChannel).publishTaskEvent(task, "team_task_tool_replay_aborted", + java.util.Map.of("pendingId", "pending-42")); + } + + @Test + void duplicateDecisionReturnsCurrentTaskProjection() { + TeamTaskEntity completed = task(TeamTaskStatus.COMPLETED); + when(taskService.getTask(101L)).thenReturn(completed); + when(governance.resolve("worker-101", 11L, 101L)).thenReturn(Optional.of(context())); + + service.approve(7L, 101L, "pending-42", "alice"); + service.deny(7L, 101L, "pending-42", "alice"); + + verify(approvalService, never()).getPending(any()); + } + + private static TeamTaskEntity task(String status) { + TeamTaskEntity task = new TeamTaskEntity(); + task.setId(101L); + task.setTeamId(7L); + task.setRunId(11L); + task.setTaskNumber(3); + task.setStatus(status); + task.setAssigneeAgentId(201L); + task.setConversationId("worker-101"); + task.setMetadata("{\"toolApproval\":{\"pendingId\":\"pending-42\"}}"); + return task; + } + + private static TeamWorkerConversationContext context() { + return new TeamWorkerConversationContext(true, "team_worker", "worker-101", + 11L, 101L, 7L, "lead-11", 201L); + } + + private static PendingApproval pending(String id) { + PendingApproval pending = new PendingApproval(id, "worker-101", "owner", + "shell", "{}", "needs approval"); + pending.setAgentId("201"); + pending.setToolCallPayload("{\"name\":\"shell\"}"); + return pending; + } +} diff --git a/mateclaw-server/src/test/java/vip/mate/team/service/TeamWorkerReplayPersistenceServiceTest.java b/mateclaw-server/src/test/java/vip/mate/team/service/TeamWorkerReplayPersistenceServiceTest.java new file mode 100644 index 00000000..927da4e0 --- /dev/null +++ b/mateclaw-server/src/test/java/vip/mate/team/service/TeamWorkerReplayPersistenceServiceTest.java @@ -0,0 +1,30 @@ +package vip.mate.team.service; + +import org.junit.jupiter.api.Test; +import vip.mate.agent.AgentService; +import vip.mate.workspace.conversation.ConversationService; + +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +class TeamWorkerReplayPersistenceServiceTest { + + @Test + void messageAndIdempotencyMarkerMustBothSucceed() { + ConversationService conversations = mock(ConversationService.class); + TeamTaskService tasks = mock(TeamTaskService.class); + TeamWorkerReplayPersistenceService service = + new TeamWorkerReplayPersistenceService(conversations, tasks); + AgentService.ChatResult result = AgentService.ChatResult.contentOnly("done"); + when(tasks.markToolReplayMessagePersisted(101L, "pending-42")).thenReturn(false); + + assertThrows(IllegalStateException.class, + () -> service.persist(101L, "pending-42", "worker-101", "done", result)); + + verify(conversations).saveMessage("worker-101", "assistant", "done", + null, "completed", 0, 0, null, null); + verify(tasks).markToolReplayMessagePersisted(101L, "pending-42"); + } +} diff --git a/mateclaw-ui/src/api/index.ts b/mateclaw-ui/src/api/index.ts index 0844dee2..3efd9708 100644 --- a/mateclaw-ui/src/api/index.ts +++ b/mateclaw-ui/src/api/index.ts @@ -1013,6 +1013,12 @@ export const teamApi = { ) => http.post(`/teams/${id}/tasks`, data), listTaskEvents: (id: string, taskId: string) => http.get(`/teams/${id}/tasks/${taskId}/events`), approveTask: (id: string, taskId: string) => http.post(`/teams/${id}/tasks/${taskId}/approve`), + approveWorkerTool: (id: string, taskId: string, pendingId: string) => + http.post(`/teams/${id}/tasks/${taskId}/worker/approve`, { pendingId }), + denyWorkerTool: (id: string, taskId: string, pendingId: string) => + http.post(`/teams/${id}/tasks/${taskId}/worker/deny`, { pendingId }), + feedbackWorker: (id: string, taskId: string, message: string) => + http.post(`/teams/${id}/tasks/${taskId}/worker/feedback`, { message }), rejectTask: (id: string, taskId: string, reason?: string) => http.post(`/teams/${id}/tasks/${taskId}/reject`, { reason }), retryTask: (id: string, taskId: string) => http.post(`/teams/${id}/tasks/${taskId}/retry`), diff --git a/mateclaw-ui/src/components/team-run/TeamRunAttention.vue b/mateclaw-ui/src/components/team-run/TeamRunAttention.vue index 9e3b1623..73acad30 100644 --- a/mateclaw-ui/src/components/team-run/TeamRunAttention.vue +++ b/mateclaw-ui/src/components/team-run/TeamRunAttention.vue @@ -15,13 +15,17 @@ const emit = defineEmits<{ 'view-task': [taskId: string] 'retry-task': [taskId: string] 'approve-task': [taskId: string] + 'approve-tool': [taskId: string] + 'deny-tool': [taskId: string] }>() const { t } = useI18n() const items = computed(() => runAttention(props.run)) const retryable = (type: string) => ['failed', 'failure', 'stale'].includes(type.toLowerCase()) const reviewable = (type: string) => ['review', 'in_review'].includes(type.toLowerCase()) -const isPending = (taskId: string, action: 'retry' | 'approve') => props.pendingActions.includes(`${taskId}:${action}`) +const toolApproval = (type: string) => ['approval', 'replay_uncertain'].includes(type.toLowerCase()) +const toolApprovalCanExecute = (type: string) => type.toLowerCase() === 'approval' +const isPending = (taskId: string, action: 'retry' | 'approve' | 'approve-tool' | 'deny-tool') => props.pendingActions.includes(`${taskId}:${action}`)