feat(harness): add rag and log projections

This commit is contained in:
zhuyongxin
2026-07-21 20:12:18 +08:00
parent 0dbdd7d8d3
commit 3e602781d6
20 changed files with 1140 additions and 0 deletions
@@ -0,0 +1,94 @@
package com.superbiz.agent.harness.tool.adapter;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.superbiz.agent.harness.core.RunContext;
import com.superbiz.agent.harness.tool.boundary.ToolBoundary;
import com.superbiz.agent.harness.tool.boundary.ToolBoundaryErrorCode;
import com.superbiz.agent.harness.tool.boundary.ToolBoundaryResult;
import com.superbiz.agent.harness.tool.boundary.ToolCallRequestEnvelope;
import com.superbiz.agent.harness.tool.contract.LogQueryScope;
import com.superbiz.agent.harness.tool.contract.LogTopic;
import com.superbiz.agent.harness.tool.contract.QueryLogsRequest;
import com.superbiz.agent.harness.tool.projection.QueryLogsResultProjector;
import java.time.Clock;
import java.time.Duration;
import java.time.Instant;
import java.time.format.DateTimeFormatter;
import java.util.Objects;
/** Bridges logical query-log requests and the existing Mock tool through ToolBoundary. */
public final class QueryLogsToolAdapter {
@FunctionalInterface
public interface LegacyExecutor {
String execute(String region, String legacyTopic, String query, Integer limit) throws Exception;
}
private static final String DEFAULT_REGION = "ap-guangzhou";
private static final int DEFAULT_LOOKBACK_MINUTES = 30;
private static final int DEFAULT_LEGACY_LIMIT = 100;
private final ToolBoundary boundary;
private final ObjectMapper objectMapper;
private final QueryLogsResultProjector projector;
private final LegacyExecutor legacyExecutor;
private final Clock clock;
private final String region;
private final int legacyLimit;
public QueryLogsToolAdapter(ToolBoundary boundary, ObjectMapper objectMapper,
QueryLogsResultProjector projector, LegacyExecutor legacyExecutor,
Clock clock) {
this(boundary, objectMapper, projector, legacyExecutor, clock, DEFAULT_REGION, DEFAULT_LEGACY_LIMIT);
}
public QueryLogsToolAdapter(ToolBoundary boundary, ObjectMapper objectMapper,
QueryLogsResultProjector projector, LegacyExecutor legacyExecutor,
Clock clock, String region, int legacyLimit) {
this.boundary = Objects.requireNonNull(boundary, "boundary must not be null");
this.objectMapper = Objects.requireNonNull(objectMapper, "objectMapper must not be null");
this.projector = Objects.requireNonNull(projector, "projector must not be null");
this.legacyExecutor = Objects.requireNonNull(legacyExecutor, "legacyExecutor must not be null");
this.clock = Objects.requireNonNull(clock, "clock must not be null");
if (region == null || region.isBlank() || legacyLimit <= 0) {
throw new IllegalArgumentException("legacy region and limit must be valid");
}
this.region = region;
this.legacyLimit = legacyLimit;
}
public ToolBoundaryResult execute(RunContext context, ToolCallRequestEnvelope envelope) {
try {
QueryLogsRequest request = objectMapper.readValue(envelope.requestJson(), QueryLogsRequest.class);
if (request.topic() == null || request.query() == null || request.query().isBlank()) {
return ToolBoundaryResult.error(envelope.toolCallId(), ToolBoundaryErrorCode.INVALID_REQUEST);
}
int lookback = request.lookbackMinutes() == null
? DEFAULT_LOOKBACK_MINUTES : request.lookbackMinutes();
if (lookback <= 0 || lookback > 24 * 60) {
return ToolBoundaryResult.error(envelope.toolCallId(), ToolBoundaryErrorCode.INVALID_REQUEST);
}
Instant end = clock.instant();
LogQueryScope scope = new LogQueryScope(
request.topic(), request.query(),
DateTimeFormatter.ISO_INSTANT.format(end.minus(Duration.ofMinutes(lookback))),
DateTimeFormatter.ISO_INSTANT.format(end));
String legacyTopic = legacyTopic(request.topic());
return boundary.execute(context, envelope,
ignored -> legacyExecutor.execute(region, legacyTopic, request.query(), legacyLimit),
raw -> projector.project(request, envelope.toolCallId(), scope, raw));
} catch (Exception e) {
return ToolBoundaryResult.error(envelope == null ? null : envelope.toolCallId(),
ToolBoundaryErrorCode.INVALID_REQUEST);
}
}
private static String legacyTopic(LogTopic topic) {
return switch (topic) {
case APPLICATION -> "application-logs";
case DATABASE_SLOW_QUERY -> "database-slow-query";
case SYSTEM_EVENTS -> "system-events";
};
}
}
@@ -0,0 +1,49 @@
package com.superbiz.agent.harness.tool.adapter;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.superbiz.agent.harness.tool.boundary.ToolBoundary;
import com.superbiz.agent.harness.tool.boundary.ToolBoundaryErrorCode;
import com.superbiz.agent.harness.tool.boundary.ToolBoundaryResult;
import com.superbiz.agent.harness.tool.boundary.ToolCallRequestEnvelope;
import com.superbiz.agent.harness.core.RunContext;
import com.superbiz.agent.harness.tool.contract.RagToolRequest;
import com.superbiz.agent.harness.tool.projection.RagResultProjector;
import java.util.Objects;
/** Bridges a typed RAG request and a legacy knowledge executor through ToolBoundary. */
public final class RagToolAdapter {
@FunctionalInterface
public interface LegacyExecutor {
Object execute(String query) throws Exception;
}
private final ToolBoundary boundary;
private final ObjectMapper objectMapper;
private final RagResultProjector projector;
private final LegacyExecutor legacyExecutor;
public RagToolAdapter(ToolBoundary boundary, ObjectMapper objectMapper,
RagResultProjector projector, LegacyExecutor legacyExecutor) {
this.boundary = Objects.requireNonNull(boundary, "boundary must not be null");
this.objectMapper = Objects.requireNonNull(objectMapper, "objectMapper must not be null");
this.projector = Objects.requireNonNull(projector, "projector must not be null");
this.legacyExecutor = Objects.requireNonNull(legacyExecutor, "legacyExecutor must not be null");
}
public ToolBoundaryResult execute(RunContext context, ToolCallRequestEnvelope envelope) {
try {
RagToolRequest request = objectMapper.readValue(envelope.requestJson(), RagToolRequest.class);
if (request.query() == null || request.query().isBlank()) {
return ToolBoundaryResult.error(envelope.toolCallId(), ToolBoundaryErrorCode.INVALID_REQUEST);
}
return boundary.execute(context, envelope,
ignored -> objectMapper.writeValueAsString(legacyExecutor.execute(request.query())),
raw -> projector.project(request, envelope.toolCallId(), raw));
} catch (Exception e) {
return ToolBoundaryResult.error(envelope == null ? null : envelope.toolCallId(),
ToolBoundaryErrorCode.INVALID_REQUEST);
}
}
}
@@ -0,0 +1,208 @@
package com.superbiz.agent.harness.tool.projection;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.superbiz.agent.harness.contract.EvidenceStatus;
import com.superbiz.agent.harness.tool.boundary.ProjectedToolResult;
import com.superbiz.agent.harness.tool.contract.LogEvent;
import com.superbiz.agent.harness.tool.contract.LogPattern;
import com.superbiz.agent.harness.tool.contract.LogQueryScope;
import com.superbiz.agent.harness.tool.contract.LogSourceKind;
import com.superbiz.agent.harness.tool.contract.QueryLogsRequest;
import com.superbiz.agent.harness.tool.contract.QueryLogsToolResult;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.regex.Pattern;
/** Projects legacy Mock log JSON into the frozen query_logs contract. */
public final class QueryLogsResultProjector {
private static final Pattern SECRET = Pattern.compile(
"(?i)(password|token|secret|api[_-]?key)\\s*[:=]\\s*[^\\s,;]+" );
private static final Pattern POD = Pattern.compile("\\bpod-[A-Za-z0-9-]+\\b");
private static final Pattern HOST = Pattern.compile(
"\\b(?:redis|mysql|kafka|rabbitmq|node|host)-[A-Za-z0-9._-]+(?::\\d+)?\\b",
Pattern.CASE_INSENSITIVE);
private static final Pattern PID = Pattern.compile("(?i)\\bPID\\s*:\\s*\\d+");
private static final Pattern IP = Pattern.compile("\\b(?:\\d{1,3}\\.){3}\\d{1,3}(?::\\d+)?\\b");
private static final Pattern SQL_LITERAL = Pattern.compile("'[^']*'");
private static final Pattern STACK_SUFFIX = Pattern.compile("\\s+at\\s+[\\w.$]+\\([^)]*\\).*$", Pattern.CASE_INSENSITIVE);
private final ObjectMapper objectMapper;
private final ToolProjectionLimits limits;
public QueryLogsResultProjector(ObjectMapper objectMapper, ToolProjectionLimits limits) {
this.objectMapper = objectMapper;
this.limits = limits;
}
public ProjectedToolResult project(QueryLogsRequest request, String toolCallId,
LogQueryScope scope, String rawResponse) throws Exception {
if (request == null || toolCallId == null || toolCallId.isBlank() || scope == null) {
throw new IllegalArgumentException("request, scope and tool call ID are required");
}
JsonNode root = objectMapper.readTree(rawResponse);
if (root == null || !root.isObject() || !root.path("success").isBoolean()
|| !root.path("success").asBoolean()) {
throw new IllegalArgumentException("log response indicates failure");
}
JsonNode logs = root.get("logs");
if (logs == null || !logs.isArray()) {
throw new IllegalArgumentException("log response has no logs array");
}
boolean truncated = logs.size() > limits.maxEvents();
List<LogEvent> allEvents = new ArrayList<>();
Map<String, PatternAccumulator> aggregates = new LinkedHashMap<>();
for (JsonNode log : logs) {
String timestamp = bounded(text(log, "timestamp"), limits.maxMessageChars());
String level = bounded(text(log, "level"), 32);
String service = bounded(text(log, "service"), 128);
String message = sanitize(text(log, "message"));
message = bounded(message, limits.maxMessageChars());
if (message.isBlank()) {
continue;
}
String example = message;
allEvents.add(new LogEvent(nullable(timestamp), nullable(level), nullable(service), message));
String patternKey = level + "\u0000" + service + "\u0000" + normalizePattern(message);
aggregates.computeIfAbsent(patternKey, ignored -> new PatternAccumulator(level, service, example))
.add(timestamp);
}
List<LogEvent> events = sample(allEvents, limits.maxEvents());
truncated |= events.size() < allEvents.size();
List<LogPattern> patterns = aggregates.values().stream()
.sorted(Comparator.comparingLong(PatternAccumulator::count).reversed()
.thenComparing(PatternAccumulator::example))
.limit(limits.maxPatterns())
.map(PatternAccumulator::toPattern)
.toList();
truncated |= aggregates.size() > patterns.size();
long matchCount = root.path("total").canConvertToLong() ? root.path("total").asLong() : allEvents.size();
QueryLogsToolResult result = new QueryLogsToolResult(
allEvents.isEmpty() ? EvidenceStatus.NO_EVIDENCE : EvidenceStatus.EVIDENCE_FOUND,
toolCallId,
LogSourceKind.MOCK,
scope,
Math.max(0, matchCount),
events.size(),
patterns,
events,
truncated);
result = fitBudget(result);
return new ProjectedToolResult(objectMapper.writeValueAsString(result), result.evidenceStatus());
}
private QueryLogsToolResult fitBudget(QueryLogsToolResult result) throws Exception {
QueryLogsToolResult current = result;
while (bytes(objectMapper.writeValueAsString(current)) > limits.maxAgentUtf8Bytes()
&& (!current.events().isEmpty() || !current.patterns().isEmpty())) {
List<LogEvent> events = new ArrayList<>(current.events());
List<LogPattern> patterns = new ArrayList<>(current.patterns());
if (!events.isEmpty()) {
events.remove(events.size() - 1);
} else {
patterns.remove(patterns.size() - 1);
}
current = new QueryLogsToolResult(
current.evidenceStatus(), current.toolCallId(), current.sourceKind(), current.scope(),
current.matchCount(), events.size(), patterns, events, true);
}
if (bytes(objectMapper.writeValueAsString(current)) > limits.maxAgentUtf8Bytes()) {
throw new IllegalArgumentException("log projection exceeds total budget");
}
return current;
}
private static List<LogEvent> sample(List<LogEvent> events, int max) {
if (events.size() <= max) {
return List.copyOf(events);
}
List<LogEvent> sampled = new ArrayList<>(max);
for (int i = 0; i < max; i++) {
int index = max == 1 ? 0 : Math.round((float) i * (events.size() - 1) / (max - 1));
sampled.add(events.get(index));
}
return sampled;
}
private static String sanitize(String message) {
String value = SECRET.matcher(message).replaceAll("$1=[REDACTED]");
value = POD.matcher(value).replaceAll("[REDACTED_POD]");
value = HOST.matcher(value).replaceAll("[REDACTED_HOST]");
value = PID.matcher(value).replaceAll("PID:[REDACTED]");
value = IP.matcher(value).replaceAll("[REDACTED_IP]");
value = SQL_LITERAL.matcher(value).replaceAll("'[REDACTED_LITERAL]'");
return STACK_SUFFIX.matcher(value).replaceAll("").trim();
}
private static String normalizePattern(String message) {
return message.replaceAll("\\b\\d+(?:\\.\\d+)?\\b", "<n>")
.replaceAll("\\s+", " ").trim();
}
private static String text(JsonNode node, String field) {
JsonNode value = node == null ? null : node.get(field);
return value == null || value.isNull() ? "" : value.asText("");
}
private static String bounded(String value, int max) {
if (value == null) {
return "";
}
return value.length() <= max ? value : value.substring(0, max);
}
private static String nullable(String value) {
return value == null || value.isBlank() ? null : value;
}
private static int bytes(String value) {
return value.getBytes(StandardCharsets.UTF_8).length;
}
private static final class PatternAccumulator {
private final String level;
private final String service;
private final String example;
private long count;
private String firstSeen;
private String lastSeen;
private PatternAccumulator(String level, String service, String example) {
this.level = level;
this.service = service;
this.example = example;
}
private PatternAccumulator add(String timestamp) {
count++;
if (firstSeen == null || timestamp.compareTo(firstSeen) < 0) {
firstSeen = timestamp;
}
if (lastSeen == null || timestamp.compareTo(lastSeen) > 0) {
lastSeen = timestamp;
}
return this;
}
private long count() {
return count;
}
private String example() {
return example;
}
private LogPattern toPattern() {
return new LogPattern(count, firstSeen, lastSeen, nullable(level), nullable(service), example);
}
}
}
@@ -0,0 +1,133 @@
package com.superbiz.agent.harness.tool.projection;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.superbiz.agent.harness.contract.EvidenceStatus;
import com.superbiz.agent.harness.tool.boundary.ProjectedToolResult;
import com.superbiz.agent.harness.tool.contract.RagEvidence;
import com.superbiz.agent.harness.tool.contract.RagToolRequest;
import com.superbiz.agent.harness.tool.contract.RagToolResult;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
/** Projects legacy knowledge retrieval JSON into the frozen RAG contract. */
public final class RagResultProjector {
private final ObjectMapper objectMapper;
private final ToolProjectionLimits limits;
public RagResultProjector(ObjectMapper objectMapper, ToolProjectionLimits limits) {
this.objectMapper = objectMapper;
this.limits = limits;
}
public ProjectedToolResult project(RagToolRequest request, String toolCallId,
String rawResponse) throws Exception {
if (request == null || toolCallId == null || toolCallId.isBlank()) {
throw new IllegalArgumentException("request and tool call ID are required");
}
JsonNode root = objectMapper.readTree(rawResponse);
if (root == null || !root.isObject()) {
throw new IllegalArgumentException("RAG response must be a JSON object");
}
boolean truncated = false;
String query = bounded(request.query(), limits.maxQueryChars());
truncated = !query.equals(request.query());
List<RagEvidence> evidence = new ArrayList<>();
Set<String> documentIds = new HashSet<>();
JsonNode blocks = root.has("evidenceBlocks") ? root.get("evidenceBlocks") : root.get("evidence_blocks");
if (blocks != null && blocks.isArray()) {
int ordinal = 0;
for (JsonNode block : blocks) {
ordinal++;
if (evidence.size() >= limits.maxEvidence()) {
truncated = true;
break;
}
String excerpt = text(block, "content");
if (excerpt.isBlank()) {
excerpt = text(block, "excerpt");
}
if (excerpt.isBlank()) {
continue;
}
String source = text(block, "source");
String title = text(block, "title");
String documentId = firstNonBlank(text(block, "document_id"), source, title,
"legacy-document-" + ordinal);
if (!documentIds.add(documentId)) {
truncated = true;
continue;
}
String boundedExcerpt = bounded(excerpt, limits.maxExcerptChars());
truncated |= !boundedExcerpt.equals(excerpt);
evidence.add(new RagEvidence(
documentId,
nullable(source),
nullable(title),
nullable(text(block, "breadcrumb")),
boundedExcerpt));
}
if (blocks.size() > limits.maxEvidence()) {
truncated = true;
}
}
EvidenceStatus status = evidence.isEmpty()
? EvidenceStatus.NO_EVIDENCE
: EvidenceStatus.EVIDENCE_FOUND;
RagToolResult result = new RagToolResult(status, toolCallId, query, evidence, evidence.size(), truncated);
result = fitBudget(result, truncated);
return new ProjectedToolResult(objectMapper.writeValueAsString(result), result.evidenceStatus());
}
private RagToolResult fitBudget(RagToolResult result, boolean truncated) throws Exception {
RagToolResult current = result;
while (bytes(objectMapper.writeValueAsString(current)) > limits.maxAgentUtf8Bytes()
&& !current.evidence().isEmpty()) {
List<RagEvidence> reduced = new ArrayList<>(current.evidence());
reduced.remove(reduced.size() - 1);
current = new RagToolResult(
reduced.isEmpty() ? EvidenceStatus.NO_EVIDENCE : EvidenceStatus.EVIDENCE_FOUND,
current.toolCallId(), current.query(), reduced, reduced.size(), true);
}
if (bytes(objectMapper.writeValueAsString(current)) > limits.maxAgentUtf8Bytes()) {
throw new IllegalArgumentException("RAG projection exceeds total budget");
}
return current;
}
private static String text(JsonNode node, String field) {
JsonNode value = node == null ? null : node.get(field);
return value == null || value.isNull() ? "" : value.asText("");
}
private static String firstNonBlank(String... values) {
for (String value : values) {
if (value != null && !value.isBlank()) {
return value;
}
}
return "unknown-document";
}
private static String nullable(String value) {
return value == null || value.isBlank() ? null : value;
}
private static String bounded(String value, int max) {
if (value == null) {
return "";
}
return value.length() <= max ? value : value.substring(0, max);
}
private static int bytes(String value) {
return value.getBytes(StandardCharsets.UTF_8).length;
}
}
@@ -0,0 +1,24 @@
package com.superbiz.agent.harness.tool.projection;
/** Bounds applied to Agent-facing projections. */
public record ToolProjectionLimits(
int maxEvidence,
int maxExcerptChars,
int maxPatterns,
int maxEvents,
int maxMessageChars,
int maxQueryChars,
int maxAgentUtf8Bytes) {
public ToolProjectionLimits {
if (maxEvidence <= 0 || maxExcerptChars <= 0 || maxPatterns <= 0
|| maxEvents <= 0 || maxMessageChars <= 0 || maxQueryChars <= 0
|| maxAgentUtf8Bytes <= 0) {
throw new IllegalArgumentException("projection limits must be positive");
}
}
public static ToolProjectionLimits defaults() {
return new ToolProjectionLimits(8, 1200, 12, 30, 1000, 500, 16_384);
}
}