feat(harness): improve trace fallback and reasoning audit
This commit is contained in:
@@ -47,6 +47,7 @@ public final class DiagnosisAgentFactory {
|
||||
new HarnessToolInterceptor(context, evidenceTools, objectMapper))
|
||||
.hooks(auditHooks)
|
||||
.outputSchema(new DiagnosisDraftOutputSchema(objectMapper).getFormat())
|
||||
.returnReasoningContents(true)
|
||||
.parallelToolExecution(false)
|
||||
.releaseThread(true)
|
||||
.build();
|
||||
|
||||
@@ -3,6 +3,8 @@ package com.superbiz.agent.harness.application;
|
||||
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.superbiz.agent.harness.application.persistence.ChatRunStore;
|
||||
import com.superbiz.agent.harness.audit.DiagnosisTraceRecorder;
|
||||
import com.superbiz.agent.harness.audit.TraceAuditEvents;
|
||||
import com.superbiz.agent.harness.application.persistence.RoutingHistory;
|
||||
import com.superbiz.agent.harness.application.routing.IntentRouterInput;
|
||||
import com.superbiz.agent.harness.application.routing.IntentRoutingException;
|
||||
@@ -33,6 +35,7 @@ public final class ChatApplicationUseCase {
|
||||
private final KnowledgeQueryOperation knowledgeQuery;
|
||||
private final DiagnosisOperation diagnosis;
|
||||
private final ObjectMapper objectMapper;
|
||||
private final DiagnosisTraceRecorder traceRecorder;
|
||||
|
||||
public ChatApplicationUseCase(DiagnosisHarnessCore core,
|
||||
Supplier<String> sessionIdSupplier,
|
||||
@@ -42,6 +45,19 @@ public final class ChatApplicationUseCase {
|
||||
KnowledgeQueryOperation knowledgeQuery,
|
||||
DiagnosisOperation diagnosis,
|
||||
ObjectMapper objectMapper) {
|
||||
this(core, sessionIdSupplier, runStore, router, systemChat, knowledgeQuery,
|
||||
diagnosis, objectMapper, DiagnosisTraceRecorder.noop());
|
||||
}
|
||||
|
||||
public ChatApplicationUseCase(DiagnosisHarnessCore core,
|
||||
Supplier<String> sessionIdSupplier,
|
||||
ChatRunStore runStore,
|
||||
IntentRouting router,
|
||||
SystemChatOperation systemChat,
|
||||
KnowledgeQueryOperation knowledgeQuery,
|
||||
DiagnosisOperation diagnosis,
|
||||
ObjectMapper objectMapper,
|
||||
DiagnosisTraceRecorder traceRecorder) {
|
||||
this.core = Objects.requireNonNull(core, "core must not be null");
|
||||
this.sessionIdSupplier = Objects.requireNonNull(
|
||||
sessionIdSupplier, "sessionIdSupplier must not be null");
|
||||
@@ -51,6 +67,7 @@ public final class ChatApplicationUseCase {
|
||||
this.knowledgeQuery = Objects.requireNonNull(knowledgeQuery, "knowledgeQuery must not be null");
|
||||
this.diagnosis = Objects.requireNonNull(diagnosis, "diagnosis must not be null");
|
||||
this.objectMapper = Objects.requireNonNull(objectMapper, "objectMapper must not be null");
|
||||
this.traceRecorder = Objects.requireNonNull(traceRecorder, "traceRecorder must not be null");
|
||||
}
|
||||
|
||||
public ChatApplicationResult execute(ChatApplicationRequest request) {
|
||||
@@ -78,12 +95,14 @@ public final class ChatApplicationUseCase {
|
||||
IntentType intent = null;
|
||||
try {
|
||||
persistStart(context, request.query());
|
||||
traceRecorder.record(TraceAuditEvents.runStarted(context));
|
||||
observer.onStarted(new CoreRunControl(core, context));
|
||||
observer.onStatus(ChatApplicationStatus.ROUTING);
|
||||
intent = router.route(context, new IntentRouterInput(
|
||||
request.query(),
|
||||
history.map(RoutingHistory::intent).orElse(null),
|
||||
history.map(RoutingHistory::userQuery).orElse(null)));
|
||||
traceRecorder.record(TraceAuditEvents.routingDecision(context, intent));
|
||||
persistIntent(context.runId(), intent);
|
||||
|
||||
PathResult path = executePath(
|
||||
@@ -93,6 +112,8 @@ public final class ChatApplicationUseCase {
|
||||
String safeJson = write(path.content());
|
||||
persistFinish(context, intent, path.outcome(), safeJson,
|
||||
path.publishedResult(), durationMillis(startedNanos));
|
||||
traceRecorder.record(TraceAuditEvents.runFinished(
|
||||
context, intent, path.outcome(), durationMillis(startedNanos)));
|
||||
return new ChatApplicationResult(
|
||||
context.sessionId(), context.runId(), intent, path.outcome(),
|
||||
path.content().contentType(), path.content());
|
||||
@@ -104,6 +125,8 @@ public final class ChatApplicationUseCase {
|
||||
} catch (RuntimeException persistenceFailure) {
|
||||
exception.addSuppressed(persistenceFailure);
|
||||
}
|
||||
traceRecorder.record(TraceAuditEvents.runFinished(
|
||||
context, intent, terminal, durationMillis(startedNanos)));
|
||||
throw safeFailure(exception, terminal);
|
||||
}
|
||||
}
|
||||
|
||||
+7
-2
@@ -24,6 +24,7 @@ import com.superbiz.agent.harness.tool.contract.RagToolResult;
|
||||
import org.springframework.ai.chat.messages.SystemMessage;
|
||||
import org.springframework.ai.chat.messages.UserMessage;
|
||||
import org.springframework.ai.chat.prompt.Prompt;
|
||||
import org.springframework.ai.converter.BeanOutputConverter;
|
||||
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.ArrayList;
|
||||
@@ -39,7 +40,8 @@ public final class KnowledgeQueryExecutor implements KnowledgeQueryOperation {
|
||||
|
||||
private static final String SYSTEM_PROMPT = """
|
||||
你只根据提供的有界知识库证据回答原始问题。每个 answer item 必须引用 exact tool_call_id 和实际 document_ids。
|
||||
不要使用外部知识,不要改写问题,不要输出 markdown fence。返回严格 KnowledgeAnswerDraft JSON。
|
||||
不要使用外部知识,不要改写问题。answer_items 必须是非空数组;每个 item 必须包含非空 text、
|
||||
exact tool_call_id 和至少一个 actual document_id。limitations 必须是字符串数组。
|
||||
""";
|
||||
|
||||
private final DiagnosisHarnessCore core;
|
||||
@@ -48,6 +50,7 @@ public final class KnowledgeQueryExecutor implements KnowledgeQueryOperation {
|
||||
private final ObjectMapper objectMapper;
|
||||
private final ObjectReader ragReader;
|
||||
private final ObjectReader answerReader;
|
||||
private final String prompt;
|
||||
private final Supplier<String> callIdSupplier;
|
||||
private final KnowledgeQueryLimits limits;
|
||||
|
||||
@@ -67,6 +70,8 @@ public final class KnowledgeQueryExecutor implements KnowledgeQueryOperation {
|
||||
this.answerReader = objectMapper.readerFor(KnowledgeAnswerDraft.class)
|
||||
.with(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES)
|
||||
.with(DeserializationFeature.FAIL_ON_TRAILING_TOKENS);
|
||||
this.prompt = SYSTEM_PROMPT + System.lineSeparator()
|
||||
+ new BeanOutputConverter<>(KnowledgeAnswerDraft.class, objectMapper).getFormat();
|
||||
this.callIdSupplier = Objects.requireNonNull(callIdSupplier, "callIdSupplier must not be null");
|
||||
this.limits = Objects.requireNonNull(limits, "limits must not be null");
|
||||
}
|
||||
@@ -102,7 +107,7 @@ public final class KnowledgeQueryExecutor implements KnowledgeQueryOperation {
|
||||
core.reserveRunBytes(context, bytes);
|
||||
String output = modelCall.call(
|
||||
context,
|
||||
new Prompt(List.of(new SystemMessage(SYSTEM_PROMPT), new UserMessage(modelInput))),
|
||||
new Prompt(List.of(new SystemMessage(prompt), new UserMessage(modelInput))),
|
||||
limits.modelTimeout(),
|
||||
limits.maxModelOutputBytes());
|
||||
KnowledgeAnswerDraft draft = readAnswer(output);
|
||||
|
||||
@@ -4,6 +4,8 @@ import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.superbiz.agent.harness.contract.IntentType;
|
||||
import com.superbiz.agent.harness.audit.DiagnosisTraceRecorder;
|
||||
import com.superbiz.agent.harness.audit.TraceAuditEvents;
|
||||
import com.superbiz.agent.harness.application.IntentRouting;
|
||||
import com.superbiz.agent.harness.core.DiagnosisHarnessCore;
|
||||
import com.superbiz.agent.harness.core.RunContext;
|
||||
@@ -36,6 +38,7 @@ public final class IntentRouter implements IntentRouting {
|
||||
private final ObjectMapper objectMapper;
|
||||
private final IntentRouterLimits limits;
|
||||
private final Consumer<RetryAttempt> attemptRecorder;
|
||||
private final DiagnosisTraceRecorder traceRecorder;
|
||||
private final String systemPrompt;
|
||||
|
||||
public IntentRouter(DiagnosisHarnessCore core,
|
||||
@@ -44,12 +47,24 @@ public final class IntentRouter implements IntentRouting {
|
||||
ObjectMapper objectMapper,
|
||||
IntentRouterLimits limits,
|
||||
Consumer<RetryAttempt> attemptRecorder) {
|
||||
this(core, retryExecutor, modelCall, objectMapper, limits,
|
||||
attemptRecorder, DiagnosisTraceRecorder.noop());
|
||||
}
|
||||
|
||||
public IntentRouter(DiagnosisHarnessCore core,
|
||||
HarnessRetryExecutor retryExecutor,
|
||||
GuardModelCall modelCall,
|
||||
ObjectMapper objectMapper,
|
||||
IntentRouterLimits limits,
|
||||
Consumer<RetryAttempt> attemptRecorder,
|
||||
DiagnosisTraceRecorder traceRecorder) {
|
||||
this.core = Objects.requireNonNull(core, "core must not be null");
|
||||
this.retryExecutor = Objects.requireNonNull(retryExecutor, "retryExecutor must not be null");
|
||||
this.modelCall = Objects.requireNonNull(modelCall, "modelCall must not be null");
|
||||
this.objectMapper = Objects.requireNonNull(objectMapper, "objectMapper must not be null");
|
||||
this.limits = Objects.requireNonNull(limits, "limits must not be null");
|
||||
this.attemptRecorder = Objects.requireNonNull(attemptRecorder, "attemptRecorder must not be null");
|
||||
this.traceRecorder = Objects.requireNonNull(traceRecorder, "traceRecorder must not be null");
|
||||
this.systemPrompt = IntentRouterPrompt.load();
|
||||
}
|
||||
|
||||
@@ -73,7 +88,10 @@ public final class IntentRouter implements IntentRouting {
|
||||
() -> parse(modelCall.call(
|
||||
context, prompt, remaining(started), limits.maxOutputBytes())),
|
||||
this::classify,
|
||||
attemptRecorder);
|
||||
attempt -> {
|
||||
attemptRecorder.accept(attempt);
|
||||
traceRecorder.record(TraceAuditEvents.routingAttempt(context, attempt));
|
||||
});
|
||||
} catch (RetryExecutionException exception) {
|
||||
if (exception.failure() == RetryFailure.CANCELLED
|
||||
|| exception.failure() == RetryFailure.BUDGET_EXHAUSTED) {
|
||||
|
||||
@@ -0,0 +1,37 @@
|
||||
package com.superbiz.agent.harness.audit;
|
||||
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
|
||||
public record DiagnosisTraceAuditEvent(
|
||||
String sessionId,
|
||||
String runId,
|
||||
TracePhase phase,
|
||||
TraceEventType eventType,
|
||||
TraceEventStatus status,
|
||||
Integer attemptNo,
|
||||
Integer durationMs,
|
||||
Map<String, Object> details) {
|
||||
|
||||
public DiagnosisTraceAuditEvent {
|
||||
requireText(sessionId, "sessionId");
|
||||
requireText(runId, "runId");
|
||||
Objects.requireNonNull(phase, "phase must not be null");
|
||||
Objects.requireNonNull(eventType, "eventType must not be null");
|
||||
Objects.requireNonNull(status, "status must not be null");
|
||||
if (attemptNo != null && attemptNo < 1) {
|
||||
throw new IllegalArgumentException("attemptNo must be positive");
|
||||
}
|
||||
if (durationMs != null && durationMs < 0) {
|
||||
throw new IllegalArgumentException("durationMs must not be negative");
|
||||
}
|
||||
details = details == null ? Map.of() : Map.copyOf(new LinkedHashMap<>(details));
|
||||
}
|
||||
|
||||
private static void requireText(String value, String name) {
|
||||
if (value == null || value.isBlank()) {
|
||||
throw new IllegalArgumentException(name + " must not be blank");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
package com.superbiz.agent.harness.audit;
|
||||
|
||||
@FunctionalInterface
|
||||
public interface DiagnosisTraceRecorder {
|
||||
|
||||
void record(DiagnosisTraceAuditEvent event);
|
||||
|
||||
static DiagnosisTraceRecorder noop() {
|
||||
return ignored -> {
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -8,6 +8,8 @@ import com.alibaba.cloud.ai.graph.agent.hook.messages.MessagesModelHook;
|
||||
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.superbiz.agent.domain.entity.AgentStep;
|
||||
import com.superbiz.agent.domain.entity.AgentReasoningAudit;
|
||||
import com.superbiz.agent.repository.AgentReasoningAuditRepository;
|
||||
import com.superbiz.agent.repository.AgentStepRepository;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
@@ -17,6 +19,7 @@ import org.springframework.ai.chat.messages.Message;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
|
||||
@@ -28,16 +31,31 @@ public final class HarnessAgentAuditHook extends MessagesModelHook {
|
||||
private final AgentStepRepository repository;
|
||||
private final ObjectMapper objectMapper;
|
||||
private final String agentName;
|
||||
private final DiagnosisTraceRecorder traceRecorder;
|
||||
private final AgentReasoningAuditRepository reasoningRepository;
|
||||
private final ConcurrentHashMap<String, Integer> stepCounters = new ConcurrentHashMap<>();
|
||||
private final ConcurrentHashMap<String, PendingStep> pendingSteps = new ConcurrentHashMap<>();
|
||||
|
||||
public HarnessAgentAuditHook(AgentStepRepository repository, ObjectMapper objectMapper, String agentName) {
|
||||
this(repository, objectMapper, agentName, DiagnosisTraceRecorder.noop());
|
||||
}
|
||||
|
||||
public HarnessAgentAuditHook(AgentStepRepository repository, ObjectMapper objectMapper,
|
||||
String agentName, DiagnosisTraceRecorder traceRecorder) {
|
||||
this(repository, objectMapper, agentName, traceRecorder, null);
|
||||
}
|
||||
|
||||
public HarnessAgentAuditHook(AgentStepRepository repository, ObjectMapper objectMapper,
|
||||
String agentName, DiagnosisTraceRecorder traceRecorder,
|
||||
AgentReasoningAuditRepository reasoningRepository) {
|
||||
this.repository = Objects.requireNonNull(repository, "repository must not be null");
|
||||
this.objectMapper = Objects.requireNonNull(objectMapper, "objectMapper must not be null");
|
||||
if (agentName == null || agentName.isBlank()) {
|
||||
throw new IllegalArgumentException("agentName must not be blank");
|
||||
}
|
||||
this.agentName = agentName;
|
||||
this.traceRecorder = Objects.requireNonNull(traceRecorder, "traceRecorder must not be null");
|
||||
this.reasoningRepository = reasoningRepository;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -53,6 +71,7 @@ public final class HarnessAgentAuditHook extends MessagesModelHook {
|
||||
return new AgentCommand(messages);
|
||||
}
|
||||
int stepIndex = stepCounters.merge(identity.runId(), 0, (current, ignored) -> current + 1);
|
||||
PendingStep pending = new PendingStep(null, System.nanoTime(), inputMetadata(messages));
|
||||
try {
|
||||
AgentStep saved = repository.save(AgentStep.builder()
|
||||
.sessionId(identity.sessionId())
|
||||
@@ -62,11 +81,11 @@ public final class HarnessAgentAuditHook extends MessagesModelHook {
|
||||
.modelInput(write(inputMetadata(messages)))
|
||||
.hasToolCall(false)
|
||||
.build());
|
||||
pendingSteps.put(stepKey(identity.runId(), stepIndex),
|
||||
new PendingStep(saved.getId(), System.nanoTime()));
|
||||
pending = new PendingStep(saved.getId(), pending.startedNanos(), pending.input());
|
||||
} catch (RuntimeException exception) {
|
||||
log.warn("Failed to persist AgentStep audit before model: agent={}", agentName);
|
||||
}
|
||||
pendingSteps.put(stepKey(identity.runId(), stepIndex), pending);
|
||||
return new AgentCommand(messages);
|
||||
}
|
||||
|
||||
@@ -79,26 +98,33 @@ public final class HarnessAgentAuditHook extends MessagesModelHook {
|
||||
}
|
||||
int stepIndex = stepCounters.getOrDefault(identity.runId(), 0);
|
||||
PendingStep pending = pendingSteps.remove(stepKey(identity.runId(), stepIndex));
|
||||
if (pending == null || pending.id() == null) {
|
||||
if (pending == null) {
|
||||
return new AgentCommand(messages);
|
||||
}
|
||||
AssistantMessage assistant = lastAssistant(messages);
|
||||
List<String> toolNames = assistant == null || assistant.getToolCalls() == null
|
||||
? List.of()
|
||||
: assistant.getToolCalls().stream().map(AssistantMessage.ToolCall::name)
|
||||
.distinct().sorted().toList();
|
||||
ReasoningContent reasoning = reasoningContent(assistant);
|
||||
Map<String, Object> output = outputMetadata(assistant, toolNames, reasoning);
|
||||
int durationMs = durationMillis(pending.startedNanos());
|
||||
try {
|
||||
AgentStep step = repository.findById(pending.id()).orElse(null);
|
||||
AgentStep step = pending.id() == null ? null : repository.findById(pending.id()).orElse(null);
|
||||
if (step != null) {
|
||||
AssistantMessage assistant = lastAssistant(messages);
|
||||
List<String> toolNames = assistant == null || assistant.getToolCalls() == null
|
||||
? List.of()
|
||||
: assistant.getToolCalls().stream().map(AssistantMessage.ToolCall::name)
|
||||
.distinct().sorted().toList();
|
||||
step.setModelOutput(write(outputMetadata(assistant, toolNames)));
|
||||
step.setModelOutput(write(output));
|
||||
step.setThought(null);
|
||||
step.setHasToolCall(!toolNames.isEmpty());
|
||||
step.setDurationMs(durationMillis(pending.startedNanos()));
|
||||
step.setDurationMs(durationMs);
|
||||
repository.save(step);
|
||||
}
|
||||
} catch (RuntimeException exception) {
|
||||
log.warn("Failed to complete AgentStep audit: agent={}", agentName);
|
||||
}
|
||||
persistReasoning(identity, stepIndex, reasoning);
|
||||
traceRecorder.record(TraceAuditEvents.agentModelStep(
|
||||
identity.sessionId(), identity.runId(), agentName, stepIndex,
|
||||
durationMs, pending.input(), output));
|
||||
return new AgentCommand(messages);
|
||||
}
|
||||
|
||||
@@ -112,14 +138,56 @@ public final class HarnessAgentAuditHook extends MessagesModelHook {
|
||||
return metadata;
|
||||
}
|
||||
|
||||
private Map<String, Object> outputMetadata(AssistantMessage assistant, List<String> toolNames) {
|
||||
private Map<String, Object> outputMetadata(AssistantMessage assistant, List<String> toolNames,
|
||||
ReasoningContent reasoning) {
|
||||
Map<String, Object> metadata = new LinkedHashMap<>();
|
||||
metadata.put("has_text", assistant != null
|
||||
&& assistant.getText() != null && !assistant.getText().isBlank());
|
||||
metadata.put("tool_names", toolNames);
|
||||
metadata.put("reasoning_available", reasoning.available());
|
||||
metadata.put("reasoning_bytes", reasoning.bytes());
|
||||
return metadata;
|
||||
}
|
||||
|
||||
private ReasoningContent reasoningContent(AssistantMessage assistant) {
|
||||
if (assistant == null || assistant.getMetadata() == null) {
|
||||
return ReasoningContent.empty();
|
||||
}
|
||||
for (String key : List.of("reasoning_content", "reasoningContent", "reasoning", "thinking")) {
|
||||
Object value = assistant.getMetadata().get(key);
|
||||
if (value instanceof CharSequence text && !text.toString().isBlank()) {
|
||||
String content = bounded(text.toString());
|
||||
return new ReasoningContent(content, true,
|
||||
content.getBytes(StandardCharsets.UTF_8).length);
|
||||
}
|
||||
}
|
||||
return ReasoningContent.empty();
|
||||
}
|
||||
|
||||
private void persistReasoning(AuditIdentity identity, int stepIndex, ReasoningContent reasoning) {
|
||||
if (reasoningRepository == null) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
reasoningRepository.save(AgentReasoningAudit.builder()
|
||||
.sessionId(identity.sessionId())
|
||||
.runId(identity.runId())
|
||||
.stepIndex(stepIndex)
|
||||
.agentName(agentName)
|
||||
.reasoningAvailable(reasoning.available())
|
||||
.reasoningContent(reasoning.content())
|
||||
.contentBytes(reasoning.bytes())
|
||||
.build());
|
||||
} catch (RuntimeException exception) {
|
||||
log.warn("Failed to persist reasoning audit: agent={}, step={}", agentName, stepIndex);
|
||||
}
|
||||
}
|
||||
|
||||
private static String bounded(String value) {
|
||||
int maxChars = 32_000;
|
||||
return value.length() <= maxChars ? value : value.substring(0, maxChars);
|
||||
}
|
||||
|
||||
private AuditIdentity identity(RunnableConfig config) {
|
||||
if (config == null) {
|
||||
return null;
|
||||
@@ -165,6 +233,12 @@ public final class HarnessAgentAuditHook extends MessagesModelHook {
|
||||
private record AuditIdentity(String sessionId, String runId) {
|
||||
}
|
||||
|
||||
private record PendingStep(Long id, long startedNanos) {
|
||||
private record ReasoningContent(String content, boolean available, int bytes) {
|
||||
private static ReasoningContent empty() {
|
||||
return new ReasoningContent(null, false, 0);
|
||||
}
|
||||
}
|
||||
|
||||
private record PendingStep(Long id, long startedNanos, Map<String, Object> input) {
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
package com.superbiz.agent.harness.audit;
|
||||
|
||||
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.superbiz.agent.domain.entity.DiagnosisTraceEvent;
|
||||
import com.superbiz.agent.repository.DiagnosisTraceEventRepository;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.util.Objects;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
@Component
|
||||
public final class JpaDiagnosisTraceRecorder implements DiagnosisTraceRecorder {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(JpaDiagnosisTraceRecorder.class);
|
||||
|
||||
private final DiagnosisTraceEventRepository repository;
|
||||
private final ObjectMapper objectMapper;
|
||||
private final ConcurrentHashMap<String, AtomicInteger> sequences = new ConcurrentHashMap<>();
|
||||
|
||||
public JpaDiagnosisTraceRecorder(DiagnosisTraceEventRepository repository, ObjectMapper objectMapper) {
|
||||
this.repository = Objects.requireNonNull(repository, "repository must not be null");
|
||||
this.objectMapper = Objects.requireNonNull(objectMapper, "objectMapper must not be null");
|
||||
}
|
||||
|
||||
@Override
|
||||
public void record(DiagnosisTraceAuditEvent event) {
|
||||
Objects.requireNonNull(event, "event must not be null");
|
||||
try {
|
||||
int sequence = sequences.computeIfAbsent(event.runId(), this::loadSequence)
|
||||
.incrementAndGet();
|
||||
repository.save(DiagnosisTraceEvent.builder()
|
||||
.sessionId(event.sessionId())
|
||||
.runId(event.runId())
|
||||
.sequenceNo(sequence)
|
||||
.phase(event.phase().name())
|
||||
.eventType(event.eventType().name())
|
||||
.status(event.status().name())
|
||||
.attemptNo(event.attemptNo())
|
||||
.durationMs(event.durationMs())
|
||||
.details(writeDetails(event))
|
||||
.build());
|
||||
} catch (RuntimeException exception) {
|
||||
log.warn("Failed to persist diagnosis trace event: phase={}, type={}",
|
||||
event.phase(), event.eventType());
|
||||
} finally {
|
||||
if (event.eventType() == TraceEventType.RUN_FINISHED) {
|
||||
sequences.remove(event.runId());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private AtomicInteger loadSequence(String runId) {
|
||||
return new AtomicInteger(repository.findMaxSequenceNoByRunId(runId));
|
||||
}
|
||||
|
||||
private String writeDetails(DiagnosisTraceAuditEvent event) {
|
||||
try {
|
||||
return objectMapper.writeValueAsString(event.details());
|
||||
} catch (JsonProcessingException exception) {
|
||||
throw new IllegalStateException("Trace event details are not serializable", exception);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -5,6 +5,7 @@ import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.superbiz.agent.domain.entity.ToolInvocation;
|
||||
import com.superbiz.agent.harness.contract.InvocationStatus;
|
||||
import com.superbiz.agent.repository.ToolInvocationRepository;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.util.LinkedHashMap;
|
||||
@@ -16,15 +17,24 @@ public final class JpaToolInvocationAuditSink implements ToolInvocationAuditSink
|
||||
|
||||
private final ToolInvocationRepository repository;
|
||||
private final ObjectMapper objectMapper;
|
||||
private final DiagnosisTraceRecorder traceRecorder;
|
||||
|
||||
public JpaToolInvocationAuditSink(ToolInvocationRepository repository, ObjectMapper objectMapper) {
|
||||
this(repository, objectMapper, DiagnosisTraceRecorder.noop());
|
||||
}
|
||||
|
||||
@Autowired
|
||||
public JpaToolInvocationAuditSink(ToolInvocationRepository repository, ObjectMapper objectMapper,
|
||||
DiagnosisTraceRecorder traceRecorder) {
|
||||
this.repository = Objects.requireNonNull(repository, "repository must not be null");
|
||||
this.objectMapper = Objects.requireNonNull(objectMapper, "objectMapper must not be null");
|
||||
this.traceRecorder = Objects.requireNonNull(traceRecorder, "traceRecorder must not be null");
|
||||
}
|
||||
|
||||
@Override
|
||||
public void record(ToolInvocationAuditEvent event) {
|
||||
Objects.requireNonNull(event, "event must not be null");
|
||||
traceRecorder.record(TraceAuditEvents.toolInvocation(event));
|
||||
repository.save(ToolInvocation.builder()
|
||||
.sessionId(event.sessionId())
|
||||
.runId(event.runId())
|
||||
|
||||
@@ -0,0 +1,203 @@
|
||||
package com.superbiz.agent.harness.audit;
|
||||
|
||||
import com.superbiz.agent.harness.contract.DiagnosisDraft;
|
||||
import com.superbiz.agent.harness.contract.FallbackType;
|
||||
import com.superbiz.agent.harness.contract.IntentType;
|
||||
import com.superbiz.agent.harness.contract.ReleaseOutcome;
|
||||
import com.superbiz.agent.harness.contract.SemanticVerdict;
|
||||
import com.superbiz.agent.harness.core.RunContext;
|
||||
import com.superbiz.agent.harness.guard.evidence.EvidenceGuardResult;
|
||||
import com.superbiz.agent.harness.guard.evidence.EvidenceViolation;
|
||||
import com.superbiz.agent.harness.retry.RetryAttempt;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
public final class TraceAuditEvents {
|
||||
|
||||
private static final Pattern SAFE_TOOL_CALL_ID =
|
||||
Pattern.compile("[A-Za-z0-9][A-Za-z0-9._-]{0,127}");
|
||||
|
||||
private TraceAuditEvents() {
|
||||
}
|
||||
|
||||
public static DiagnosisTraceAuditEvent runStarted(RunContext context) {
|
||||
return event(context, TracePhase.RUN, TraceEventType.RUN_STARTED,
|
||||
TraceEventStatus.STARTED, null, null, Map.of());
|
||||
}
|
||||
|
||||
public static DiagnosisTraceAuditEvent runFinished(
|
||||
RunContext context, IntentType intent, ReleaseOutcome outcome, int durationMs) {
|
||||
Map<String, Object> details = new LinkedHashMap<>();
|
||||
if (intent != null) {
|
||||
details.put("intent", intent.name());
|
||||
}
|
||||
details.put("release_outcome", outcome.name());
|
||||
return event(context, TracePhase.RUN, TraceEventType.RUN_FINISHED,
|
||||
terminalStatus(outcome), null, durationMs, details);
|
||||
}
|
||||
|
||||
public static DiagnosisTraceAuditEvent routingAttempt(
|
||||
RunContext context, RetryAttempt attempt) {
|
||||
return retry(context, TracePhase.ROUTING, TraceEventType.ROUTING_ATTEMPT, attempt);
|
||||
}
|
||||
|
||||
public static DiagnosisTraceAuditEvent routingDecision(
|
||||
RunContext context, IntentType intent) {
|
||||
return event(context, TracePhase.ROUTING, TraceEventType.ROUTING_DECISION,
|
||||
TraceEventStatus.SUCCEEDED, null, null, Map.of("intent", intent.name()));
|
||||
}
|
||||
|
||||
public static DiagnosisTraceAuditEvent agentModelStep(
|
||||
String sessionId, String runId, String agentName, int stepIndex,
|
||||
int durationMs, Map<String, Object> input, Map<String, Object> output) {
|
||||
Map<String, Object> details = new LinkedHashMap<>();
|
||||
details.put("agent_name", agentName);
|
||||
details.put("step_index", stepIndex);
|
||||
details.put("input", input);
|
||||
details.put("output", output);
|
||||
return new DiagnosisTraceAuditEvent(
|
||||
sessionId, runId, TracePhase.AGENT, TraceEventType.AGENT_MODEL_STEP,
|
||||
TraceEventStatus.SUCCEEDED, stepIndex + 1, durationMs, details);
|
||||
}
|
||||
|
||||
public static DiagnosisTraceAuditEvent toolInvocation(ToolInvocationAuditEvent tool) {
|
||||
Map<String, Object> details = new LinkedHashMap<>();
|
||||
details.put("tool_call_id", tool.toolCallId());
|
||||
details.put("tool_name", tool.toolName());
|
||||
details.put("invocation_status", tool.status().name());
|
||||
details.put("evidence_status", tool.evidenceStatus().name());
|
||||
details.put("request_bytes", tool.requestBytes());
|
||||
details.put("agent_result_bytes", tool.agentResultBytes());
|
||||
if (tool.errorCode() != null) {
|
||||
details.put("error_code", tool.errorCode());
|
||||
}
|
||||
return new DiagnosisTraceAuditEvent(
|
||||
tool.sessionId(), tool.runId(), TracePhase.TOOL, TraceEventType.TOOL_INVOCATION,
|
||||
tool.status() == com.superbiz.agent.harness.contract.InvocationStatus.READY
|
||||
? TraceEventStatus.READY : TraceEventStatus.ERROR,
|
||||
null, tool.durationMs(), details);
|
||||
}
|
||||
|
||||
public static DiagnosisTraceAuditEvent evidenceValidation(
|
||||
RunContext context, TraceEventType type,
|
||||
EvidenceGuardResult result, DiagnosisDraft draft) {
|
||||
Map<String, Object> details = new LinkedHashMap<>();
|
||||
details.put("violation_count", result.violations().size());
|
||||
details.put("violations", violations(result.violations()));
|
||||
ToolReferences references = toolReferences(draft);
|
||||
details.put("referenced_tool_call_ids", references.safeIds());
|
||||
details.put("invalid_tool_reference_count", references.invalidCount());
|
||||
result.verifiedSnapshot().ifPresent(snapshot -> {
|
||||
details.put("verified_analysis_count", snapshot.analyses().size());
|
||||
details.put("verified_source_count", snapshot.verifiedSources().size());
|
||||
});
|
||||
return event(context, TracePhase.EVIDENCE, type,
|
||||
result.valid() ? TraceEventStatus.PASSED : TraceEventStatus.REJECTED,
|
||||
null, null, details);
|
||||
}
|
||||
|
||||
public static DiagnosisTraceAuditEvent evidenceRepairAttempt(
|
||||
RunContext context, RetryAttempt attempt) {
|
||||
return retry(context, TracePhase.EVIDENCE,
|
||||
TraceEventType.EVIDENCE_REPAIR_ATTEMPT, attempt);
|
||||
}
|
||||
|
||||
public static DiagnosisTraceAuditEvent semanticAttempt(
|
||||
RunContext context, RetryAttempt attempt) {
|
||||
return retry(context, TracePhase.SEMANTIC,
|
||||
TraceEventType.SEMANTIC_GUARD_ATTEMPT, attempt);
|
||||
}
|
||||
|
||||
public static DiagnosisTraceAuditEvent semanticDecision(
|
||||
RunContext context, SemanticVerdict verdict) {
|
||||
TraceEventStatus status = verdict == SemanticVerdict.SUPPORTED
|
||||
? TraceEventStatus.SUPPORTED : TraceEventStatus.UNSUPPORTED;
|
||||
return event(context, TracePhase.SEMANTIC, TraceEventType.SEMANTIC_GUARD_DECISION,
|
||||
status, null, null, Map.of("verdict", verdict.name()));
|
||||
}
|
||||
|
||||
public static DiagnosisTraceAuditEvent semanticUnavailable(RunContext context) {
|
||||
return event(context, TracePhase.SEMANTIC, TraceEventType.SEMANTIC_GUARD_DECISION,
|
||||
TraceEventStatus.UNAVAILABLE, null, null, Map.of());
|
||||
}
|
||||
|
||||
public static DiagnosisTraceAuditEvent releaseDecision(
|
||||
RunContext context, ReleaseOutcome outcome, FallbackType fallbackType) {
|
||||
Map<String, Object> details = new LinkedHashMap<>();
|
||||
details.put("release_outcome", outcome.name());
|
||||
if (fallbackType != null) {
|
||||
details.put("fallback_type", fallbackType.name());
|
||||
}
|
||||
return event(context, TracePhase.RELEASE, TraceEventType.RELEASE_DECISION,
|
||||
terminalStatus(outcome), null, null, details);
|
||||
}
|
||||
|
||||
private static DiagnosisTraceAuditEvent retry(
|
||||
RunContext context, TracePhase phase, TraceEventType type, RetryAttempt attempt) {
|
||||
Map<String, Object> details = new LinkedHashMap<>();
|
||||
if (attempt.failure() != null) {
|
||||
details.put("failure", attempt.failure().name());
|
||||
}
|
||||
return event(context, phase, type,
|
||||
attempt.success() ? TraceEventStatus.SUCCEEDED : TraceEventStatus.FAILED,
|
||||
attempt.attemptNumber(), null, details);
|
||||
}
|
||||
|
||||
private static DiagnosisTraceAuditEvent event(
|
||||
RunContext context, TracePhase phase, TraceEventType type, TraceEventStatus status,
|
||||
Integer attemptNo, Integer durationMs, Map<String, Object> details) {
|
||||
return new DiagnosisTraceAuditEvent(
|
||||
context.sessionId(), context.runId(), phase, type, status,
|
||||
attemptNo, durationMs, details);
|
||||
}
|
||||
|
||||
private static TraceEventStatus terminalStatus(ReleaseOutcome outcome) {
|
||||
return switch (outcome) {
|
||||
case SUCCESS -> TraceEventStatus.SUCCEEDED;
|
||||
case FALLBACK -> TraceEventStatus.FALLBACK;
|
||||
case FAILED -> TraceEventStatus.FAILED;
|
||||
case CANCELLED -> TraceEventStatus.CANCELLED;
|
||||
};
|
||||
}
|
||||
|
||||
private static List<Map<String, String>> violations(List<EvidenceViolation> violations) {
|
||||
List<Map<String, String>> result = new ArrayList<>();
|
||||
for (EvidenceViolation violation : violations) {
|
||||
Map<String, String> item = new LinkedHashMap<>();
|
||||
item.put("code", violation.code().name());
|
||||
item.put("target", violation.target());
|
||||
result.add(Map.copyOf(item));
|
||||
}
|
||||
return List.copyOf(result);
|
||||
}
|
||||
|
||||
private static ToolReferences toolReferences(DiagnosisDraft draft) {
|
||||
if (draft == null) {
|
||||
return new ToolReferences(List.of(), 0);
|
||||
}
|
||||
Set<String> safeIds = new LinkedHashSet<>();
|
||||
int invalid = 0;
|
||||
for (DiagnosisDraft.AnalysisItem item : draft.analysis()) {
|
||||
if (item == null) {
|
||||
continue;
|
||||
}
|
||||
for (String id : item.toolCallIds()) {
|
||||
if (id != null && SAFE_TOOL_CALL_ID.matcher(id).matches()) {
|
||||
safeIds.add(id);
|
||||
} else {
|
||||
invalid++;
|
||||
}
|
||||
}
|
||||
}
|
||||
return new ToolReferences(List.copyOf(safeIds), invalid);
|
||||
}
|
||||
|
||||
private record ToolReferences(List<String> safeIds, int invalidCount) {
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
package com.superbiz.agent.harness.audit;
|
||||
|
||||
public enum TraceEventStatus {
|
||||
STARTED,
|
||||
SUCCEEDED,
|
||||
FAILED,
|
||||
READY,
|
||||
ERROR,
|
||||
PASSED,
|
||||
REJECTED,
|
||||
SUPPORTED,
|
||||
UNSUPPORTED,
|
||||
UNAVAILABLE,
|
||||
FALLBACK,
|
||||
CANCELLED
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
package com.superbiz.agent.harness.audit;
|
||||
|
||||
public enum TraceEventType {
|
||||
RUN_STARTED,
|
||||
ROUTING_ATTEMPT,
|
||||
ROUTING_DECISION,
|
||||
AGENT_MODEL_STEP,
|
||||
TOOL_INVOCATION,
|
||||
EVIDENCE_GUARD_INITIAL,
|
||||
EVIDENCE_REPAIR_ATTEMPT,
|
||||
EVIDENCE_GUARD_RECHECK,
|
||||
SEMANTIC_GUARD_ATTEMPT,
|
||||
SEMANTIC_GUARD_DECISION,
|
||||
RELEASE_DECISION,
|
||||
RUN_FINISHED
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
package com.superbiz.agent.harness.audit;
|
||||
|
||||
public enum TracePhase {
|
||||
RUN,
|
||||
ROUTING,
|
||||
AGENT,
|
||||
TOOL,
|
||||
EVIDENCE,
|
||||
SEMANTIC,
|
||||
RELEASE
|
||||
}
|
||||
@@ -10,12 +10,27 @@ public record SafeFallback(
|
||||
@JsonProperty("message") String message,
|
||||
@JsonProperty("verified_sources") List<VerifiedSource> verifiedSources,
|
||||
@JsonProperty("limitations") List<String> limitations,
|
||||
@JsonProperty("next_steps") List<String> nextSteps) {
|
||||
@JsonProperty("next_steps") List<String> nextSteps,
|
||||
@JsonProperty("failure_stage") String failureStage,
|
||||
@JsonProperty("observed_facts") List<ObservedFact> observedFacts,
|
||||
@JsonProperty("validation_issues") List<ValidationIssue> validationIssues) {
|
||||
|
||||
public SafeFallback {
|
||||
verifiedSources = ContractCollections.immutable(verifiedSources);
|
||||
limitations = ContractCollections.immutable(limitations);
|
||||
nextSteps = ContractCollections.immutable(nextSteps);
|
||||
observedFacts = ContractCollections.immutable(observedFacts);
|
||||
validationIssues = ContractCollections.immutable(validationIssues);
|
||||
}
|
||||
|
||||
public SafeFallback(FallbackType type,
|
||||
String conclusion,
|
||||
String message,
|
||||
List<VerifiedSource> verifiedSources,
|
||||
List<String> limitations,
|
||||
List<String> nextSteps) {
|
||||
this(type, conclusion, message, verifiedSources, limitations, nextSteps,
|
||||
null, List.of(), List.of());
|
||||
}
|
||||
|
||||
public record VerifiedSource(
|
||||
@@ -23,4 +38,16 @@ public record SafeFallback(
|
||||
@JsonProperty("source") String source,
|
||||
@JsonProperty("scope") String scope) {
|
||||
}
|
||||
|
||||
public record ObservedFact(
|
||||
@JsonProperty("source_type") String sourceType,
|
||||
@JsonProperty("source") String source,
|
||||
@JsonProperty("scope") String scope,
|
||||
@JsonProperty("summary") String summary) {
|
||||
}
|
||||
|
||||
public record ValidationIssue(
|
||||
@JsonProperty("code") String code,
|
||||
@JsonProperty("target") String target) {
|
||||
}
|
||||
}
|
||||
|
||||
@@ -4,6 +4,8 @@ import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.superbiz.agent.harness.contract.SemanticVerdict;
|
||||
import com.superbiz.agent.harness.audit.DiagnosisTraceRecorder;
|
||||
import com.superbiz.agent.harness.audit.TraceAuditEvents;
|
||||
import com.superbiz.agent.harness.core.DiagnosisHarnessCore;
|
||||
import com.superbiz.agent.harness.core.RunContext;
|
||||
import com.superbiz.agent.harness.retry.HarnessRetryExecutor;
|
||||
@@ -31,6 +33,7 @@ public final class SemanticGuard {
|
||||
private final ObjectMapper objectMapper;
|
||||
private final SemanticGuardLimits limits;
|
||||
private final Consumer<RetryAttempt> attemptRecorder;
|
||||
private final DiagnosisTraceRecorder traceRecorder;
|
||||
private final String prompt;
|
||||
|
||||
public SemanticGuard(DiagnosisHarnessCore core,
|
||||
@@ -39,6 +42,17 @@ public final class SemanticGuard {
|
||||
ObjectMapper objectMapper,
|
||||
SemanticGuardLimits limits,
|
||||
Consumer<RetryAttempt> attemptRecorder) {
|
||||
this(core, retryExecutor, modelCall, objectMapper, limits,
|
||||
attemptRecorder, DiagnosisTraceRecorder.noop());
|
||||
}
|
||||
|
||||
public SemanticGuard(DiagnosisHarnessCore core,
|
||||
HarnessRetryExecutor retryExecutor,
|
||||
GuardModelCall modelCall,
|
||||
ObjectMapper objectMapper,
|
||||
SemanticGuardLimits limits,
|
||||
Consumer<RetryAttempt> attemptRecorder,
|
||||
DiagnosisTraceRecorder traceRecorder) {
|
||||
this.core = Objects.requireNonNull(core, "core must not be null");
|
||||
this.retryExecutor = Objects.requireNonNull(retryExecutor, "retryExecutor must not be null");
|
||||
this.modelCall = Objects.requireNonNull(modelCall, "modelCall must not be null");
|
||||
@@ -46,6 +60,7 @@ public final class SemanticGuard {
|
||||
this.limits = Objects.requireNonNull(limits, "limits must not be null");
|
||||
this.attemptRecorder = Objects.requireNonNull(
|
||||
attemptRecorder, "attemptRecorder must not be null");
|
||||
this.traceRecorder = Objects.requireNonNull(traceRecorder, "traceRecorder must not be null");
|
||||
this.prompt = SemanticGuardPrompt.load();
|
||||
}
|
||||
|
||||
@@ -68,7 +83,10 @@ public final class SemanticGuard {
|
||||
() -> parse(modelCall.call(
|
||||
context, modelPrompt, remainingTimeout(startedNanos), limits.maxOutputBytes())),
|
||||
this::classify,
|
||||
attemptRecorder);
|
||||
attempt -> {
|
||||
attemptRecorder.accept(attempt);
|
||||
traceRecorder.record(TraceAuditEvents.semanticAttempt(context, attempt));
|
||||
});
|
||||
}
|
||||
|
||||
private SemanticGuardDecision parse(String output) {
|
||||
|
||||
@@ -1,6 +1,10 @@
|
||||
package com.superbiz.agent.harness.release;
|
||||
|
||||
import com.superbiz.agent.harness.contract.DiagnosisDraft;
|
||||
import com.superbiz.agent.harness.audit.DiagnosisTraceRecorder;
|
||||
import com.superbiz.agent.harness.audit.TraceAuditEvents;
|
||||
import com.superbiz.agent.harness.audit.TraceEventType;
|
||||
import com.superbiz.agent.harness.contract.FallbackType;
|
||||
import com.superbiz.agent.harness.contract.SemanticVerdict;
|
||||
import com.superbiz.agent.harness.core.BudgetExceededException;
|
||||
import com.superbiz.agent.harness.core.RunAbortedException;
|
||||
@@ -22,16 +26,27 @@ public final class DiagnosisReleaseUseCase {
|
||||
private final EvidenceRepair evidenceRepair;
|
||||
private final SemanticGuard semanticGuard;
|
||||
private final SafeFallbackFactory fallbackFactory;
|
||||
private final DiagnosisTraceRecorder traceRecorder;
|
||||
|
||||
public DiagnosisReleaseUseCase(EvidenceGuard evidenceGuard,
|
||||
EvidenceRepair evidenceRepair,
|
||||
SemanticGuard semanticGuard,
|
||||
SafeFallbackFactory fallbackFactory) {
|
||||
this(evidenceGuard, evidenceRepair, semanticGuard, fallbackFactory,
|
||||
DiagnosisTraceRecorder.noop());
|
||||
}
|
||||
|
||||
public DiagnosisReleaseUseCase(EvidenceGuard evidenceGuard,
|
||||
EvidenceRepair evidenceRepair,
|
||||
SemanticGuard semanticGuard,
|
||||
SafeFallbackFactory fallbackFactory,
|
||||
DiagnosisTraceRecorder traceRecorder) {
|
||||
this.evidenceGuard = Objects.requireNonNull(evidenceGuard, "evidenceGuard must not be null");
|
||||
this.evidenceRepair = Objects.requireNonNull(evidenceRepair, "evidenceRepair must not be null");
|
||||
this.semanticGuard = Objects.requireNonNull(semanticGuard, "semanticGuard must not be null");
|
||||
this.fallbackFactory = Objects.requireNonNull(
|
||||
fallbackFactory, "fallbackFactory must not be null");
|
||||
this.traceRecorder = Objects.requireNonNull(traceRecorder, "traceRecorder must not be null");
|
||||
}
|
||||
|
||||
public DiagnosisReleaseResult execute(RunContext context, String query, DiagnosisDraft draft) {
|
||||
@@ -43,16 +58,20 @@ public final class DiagnosisReleaseUseCase {
|
||||
|
||||
DiagnosisDraft candidate = draft;
|
||||
EvidenceGuardResult evidence = evidenceGuard.validate(context, candidate);
|
||||
traceRecorder.record(TraceAuditEvents.evidenceValidation(
|
||||
context, TraceEventType.EVIDENCE_GUARD_INITIAL, evidence, candidate));
|
||||
if (!evidence.valid()) {
|
||||
try {
|
||||
candidate = evidenceRepair.repair(context, query, draft, evidence.violations());
|
||||
evidence = evidenceGuard.validate(context, candidate);
|
||||
traceRecorder.record(TraceAuditEvents.evidenceValidation(
|
||||
context, TraceEventType.EVIDENCE_GUARD_RECHECK, evidence, candidate));
|
||||
} catch (RuntimeException exception) {
|
||||
propagateTerminal(exception);
|
||||
return evidenceFailure();
|
||||
return evidenceFailure(context, evidence);
|
||||
}
|
||||
if (!evidence.valid()) {
|
||||
return evidenceFailure();
|
||||
return evidenceFailure(context, evidence);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -63,18 +82,32 @@ public final class DiagnosisReleaseUseCase {
|
||||
context, SemanticGuardInput.from(query, candidate, snapshot));
|
||||
} catch (RuntimeException exception) {
|
||||
propagateTerminal(exception);
|
||||
traceRecorder.record(TraceAuditEvents.semanticUnavailable(context));
|
||||
traceRecorder.record(TraceAuditEvents.releaseDecision(
|
||||
context, com.superbiz.agent.harness.contract.ReleaseOutcome.FALLBACK,
|
||||
FallbackType.SEMANTIC_UNAVAILABLE));
|
||||
return DiagnosisReleaseResult.fallback(
|
||||
fallbackFactory.semanticUnavailable(snapshot));
|
||||
}
|
||||
return decision.verdict() == SemanticVerdict.SUPPORTED
|
||||
? DiagnosisReleaseResult.success(candidate, snapshot)
|
||||
: DiagnosisReleaseResult.fallback(
|
||||
fallbackFactory.semanticUnsupported(snapshot));
|
||||
traceRecorder.record(TraceAuditEvents.semanticDecision(context, decision.verdict()));
|
||||
if (decision.verdict() == SemanticVerdict.SUPPORTED) {
|
||||
traceRecorder.record(TraceAuditEvents.releaseDecision(
|
||||
context, com.superbiz.agent.harness.contract.ReleaseOutcome.SUCCESS, null));
|
||||
return DiagnosisReleaseResult.success(candidate, snapshot);
|
||||
}
|
||||
traceRecorder.record(TraceAuditEvents.releaseDecision(
|
||||
context, com.superbiz.agent.harness.contract.ReleaseOutcome.FALLBACK,
|
||||
FallbackType.SEMANTIC_UNSUPPORTED));
|
||||
return DiagnosisReleaseResult.fallback(fallbackFactory.semanticUnsupported(snapshot));
|
||||
}
|
||||
|
||||
private DiagnosisReleaseResult evidenceFailure() {
|
||||
private DiagnosisReleaseResult evidenceFailure(
|
||||
RunContext context, EvidenceGuardResult evidence) {
|
||||
traceRecorder.record(TraceAuditEvents.releaseDecision(
|
||||
context, com.superbiz.agent.harness.contract.ReleaseOutcome.FALLBACK,
|
||||
FallbackType.EVIDENCE_VALIDATION_FAILED));
|
||||
return DiagnosisReleaseResult.fallback(
|
||||
fallbackFactory.evidenceValidationFailed());
|
||||
fallbackFactory.evidenceValidationFailed(evidence.violations()));
|
||||
}
|
||||
|
||||
private void propagateTerminal(RuntimeException exception) {
|
||||
|
||||
@@ -5,6 +5,8 @@ import com.fasterxml.jackson.databind.DeserializationFeature;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.fasterxml.jackson.databind.ObjectReader;
|
||||
import com.superbiz.agent.harness.contract.DiagnosisDraft;
|
||||
import com.superbiz.agent.harness.audit.DiagnosisTraceRecorder;
|
||||
import com.superbiz.agent.harness.audit.TraceAuditEvents;
|
||||
import com.superbiz.agent.harness.core.DiagnosisHarnessCore;
|
||||
import com.superbiz.agent.harness.core.RunContext;
|
||||
import com.superbiz.agent.harness.guard.evidence.EvidenceViolation;
|
||||
@@ -32,6 +34,7 @@ public final class EvidenceRepair {
|
||||
private final ObjectReader draftReader;
|
||||
private final EvidenceRepairLimits limits;
|
||||
private final Consumer<RetryAttempt> attemptRecorder;
|
||||
private final DiagnosisTraceRecorder traceRecorder;
|
||||
private final String prompt;
|
||||
|
||||
public EvidenceRepair(DiagnosisHarnessCore core,
|
||||
@@ -40,6 +43,17 @@ public final class EvidenceRepair {
|
||||
ObjectMapper objectMapper,
|
||||
EvidenceRepairLimits limits,
|
||||
Consumer<RetryAttempt> attemptRecorder) {
|
||||
this(core, retryExecutor, modelCall, objectMapper, limits,
|
||||
attemptRecorder, DiagnosisTraceRecorder.noop());
|
||||
}
|
||||
|
||||
public EvidenceRepair(DiagnosisHarnessCore core,
|
||||
HarnessRetryExecutor retryExecutor,
|
||||
GuardModelCall modelCall,
|
||||
ObjectMapper objectMapper,
|
||||
EvidenceRepairLimits limits,
|
||||
Consumer<RetryAttempt> attemptRecorder,
|
||||
DiagnosisTraceRecorder traceRecorder) {
|
||||
this.core = Objects.requireNonNull(core, "core must not be null");
|
||||
this.retryExecutor = Objects.requireNonNull(retryExecutor, "retryExecutor must not be null");
|
||||
this.modelCall = Objects.requireNonNull(modelCall, "modelCall must not be null");
|
||||
@@ -50,6 +64,7 @@ public final class EvidenceRepair {
|
||||
this.limits = Objects.requireNonNull(limits, "limits must not be null");
|
||||
this.attemptRecorder = Objects.requireNonNull(
|
||||
attemptRecorder, "attemptRecorder must not be null");
|
||||
this.traceRecorder = Objects.requireNonNull(traceRecorder, "traceRecorder must not be null");
|
||||
this.prompt = EvidenceRepairPrompt.load();
|
||||
}
|
||||
|
||||
@@ -87,7 +102,10 @@ public final class EvidenceRepair {
|
||||
return repaired;
|
||||
},
|
||||
this::classify,
|
||||
attemptRecorder);
|
||||
attempt -> {
|
||||
attemptRecorder.accept(attempt);
|
||||
traceRecorder.record(TraceAuditEvents.evidenceRepairAttempt(context, attempt));
|
||||
});
|
||||
}
|
||||
|
||||
private DiagnosisDraft parse(String output) {
|
||||
|
||||
@@ -2,50 +2,109 @@ package com.superbiz.agent.harness.release;
|
||||
|
||||
import com.superbiz.agent.harness.contract.FallbackType;
|
||||
import com.superbiz.agent.harness.contract.SafeFallback;
|
||||
import com.superbiz.agent.harness.guard.evidence.EvidenceViolation;
|
||||
import com.superbiz.agent.harness.guard.evidence.VerifiedEvidence;
|
||||
import com.superbiz.agent.harness.guard.evidence.VerifiedEvidenceSnapshot;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
|
||||
public final class SafeFallbackFactory {
|
||||
|
||||
public SafeFallback evidenceValidationFailed() {
|
||||
private static final int MAX_OBSERVED_FACTS = 12;
|
||||
private static final int MAX_SUMMARY_CHARS = 320;
|
||||
|
||||
public SafeFallback evidenceValidationFailed(List<EvidenceViolation> violations) {
|
||||
return fallback(
|
||||
FallbackType.EVIDENCE_VALIDATION_FAILED,
|
||||
List.of(),
|
||||
"当前证据无法完成真实性校验,无法确认根因",
|
||||
"证据引用校验未通过",
|
||||
"重新收集当前诊断范围内的证据后再发起诊断");
|
||||
"已完成证据引用检查,但当前报告无法通过真实性校验,因此不能确认根因",
|
||||
List.of("证据引用或报告结构校验未通过"),
|
||||
List.of("根据 validation_issues 修正报告结构,或补充对应范围的证据后重试"),
|
||||
"EVIDENCE_VALIDATION",
|
||||
List.of(),
|
||||
issues(violations));
|
||||
}
|
||||
|
||||
public SafeFallback semanticUnsupported(VerifiedEvidenceSnapshot snapshot) {
|
||||
return fallback(
|
||||
FallbackType.SEMANTIC_UNSUPPORTED,
|
||||
sources(snapshot),
|
||||
"当前证据不足,无法确认根因",
|
||||
"语义校验未通过",
|
||||
"补充当前缺失的数据后重新发起诊断");
|
||||
"已收集到可验证事实,但这些事实不足以支持当前根因结论",
|
||||
List.of("语义校验未通过,已确认事实仍可用于后续排查"),
|
||||
List.of("围绕 observed_facts 补充缺失的实时日志、指标或数据库证据后重试"),
|
||||
"SEMANTIC_VALIDATION",
|
||||
facts(snapshot),
|
||||
List.of());
|
||||
}
|
||||
|
||||
public SafeFallback semanticUnavailable(VerifiedEvidenceSnapshot snapshot) {
|
||||
return fallback(
|
||||
FallbackType.SEMANTIC_UNAVAILABLE,
|
||||
sources(snapshot),
|
||||
"当前证据暂时无法完成语义校验,无法确认根因",
|
||||
"语义校验暂不可用",
|
||||
"稍后重试或补充当前缺失的数据");
|
||||
"已收集到可验证事实,但当前无法完成语义校验,因此暂不发布根因结论",
|
||||
List.of("语义校验暂不可用"),
|
||||
List.of("稍后重试;已确认事实可继续用于人工排查"),
|
||||
"SEMANTIC_VALIDATION",
|
||||
facts(snapshot),
|
||||
List.of());
|
||||
}
|
||||
|
||||
private SafeFallback fallback(FallbackType type,
|
||||
List<SafeFallback.VerifiedSource> sources,
|
||||
String message,
|
||||
String limitation,
|
||||
String nextStep) {
|
||||
List<String> limitations,
|
||||
List<String> nextSteps,
|
||||
String failureStage,
|
||||
List<SafeFallback.ObservedFact> observedFacts,
|
||||
List<SafeFallback.ValidationIssue> validationIssues) {
|
||||
return new SafeFallback(
|
||||
type, null, message, sources, List.of(limitation), List.of(nextStep));
|
||||
type, null, message, sources, limitations, nextSteps,
|
||||
failureStage, observedFacts, validationIssues);
|
||||
}
|
||||
|
||||
private List<SafeFallback.VerifiedSource> sources(VerifiedEvidenceSnapshot snapshot) {
|
||||
return Objects.requireNonNull(snapshot, "snapshot must not be null").verifiedSources();
|
||||
}
|
||||
|
||||
private List<SafeFallback.ObservedFact> facts(VerifiedEvidenceSnapshot snapshot) {
|
||||
Objects.requireNonNull(snapshot, "snapshot must not be null");
|
||||
Map<String, SafeFallback.ObservedFact> unique = new LinkedHashMap<>();
|
||||
for (var analysis : snapshot.analyses()) {
|
||||
for (VerifiedEvidence evidence : analysis.evidence()) {
|
||||
if (unique.size() >= MAX_OBSERVED_FACTS) {
|
||||
return List.copyOf(unique.values());
|
||||
}
|
||||
String summary = bounded(evidence.excerpt());
|
||||
if (summary == null) {
|
||||
summary = "已验证结构化证据,具体值保留在受限审计记录中";
|
||||
}
|
||||
String key = evidence.sourceType() + '\u0000' + evidence.source()
|
||||
+ '\u0000' + evidence.scope() + '\u0000' + summary;
|
||||
unique.putIfAbsent(key, new SafeFallback.ObservedFact(
|
||||
evidence.sourceType(), evidence.source(), evidence.scope(), summary));
|
||||
}
|
||||
}
|
||||
return List.copyOf(unique.values());
|
||||
}
|
||||
|
||||
private List<SafeFallback.ValidationIssue> issues(List<EvidenceViolation> violations) {
|
||||
List<SafeFallback.ValidationIssue> result = new ArrayList<>();
|
||||
for (EvidenceViolation violation : violations == null ? List.<EvidenceViolation>of() : violations) {
|
||||
result.add(new SafeFallback.ValidationIssue(
|
||||
violation.code().name(), violation.target()));
|
||||
}
|
||||
return List.copyOf(result);
|
||||
}
|
||||
|
||||
private String bounded(String value) {
|
||||
if (value == null || value.isBlank()) {
|
||||
return null;
|
||||
}
|
||||
return value.length() <= MAX_SUMMARY_CHARS
|
||||
? value : value.substring(0, MAX_SUMMARY_CHARS);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user