pendingExtraction(String bufferKey, Message message) {
+ recentConversationBuffer.append(bufferKey, message);
+ if (message.role() == Message.Role.ASSISTANT) {
+ appendToPendingBuffer(bufferKey, message);
+ return Optional.empty();
+ }
+
+ var snapshot = pendingConversationBuffer.load(bufferKey);
+ if (snapshot.isEmpty()) {
+ appendToPendingBuffer(bufferKey, message);
+ return Optional.empty();
+ }
+
+ var detectionInput =
+ new CommitDetectionInput(
+ snapshot, List.of(message), CommitDetectionContext.empty());
+ var decision =
+ contextCommitDetector
+ .shouldCommit(detectionInput)
+ .defaultIfEmpty(CommitDecision.hold())
+ .block();
+
+ if (decision == null || !decision.shouldSeal()) {
+ appendToPendingBuffer(bufferKey, message);
+ return Optional.empty();
+ }
+
+ log.debug(
+ "Boundary detection triggered sealing: memoryId={}, reason={}, confidence={}",
+ bufferKey,
+ decision.reason(),
+ decision.confidence());
+
+ var sealedMessages = List.copyOf(snapshot);
+ pendingConversationBuffer.clear(bufferKey);
+ pendingConversationBuffer.append(bufferKey, message);
+
+ return Optional.of(new PendingExtraction(sealedMessages, new HashMap<>()));
}
private void appendToPendingBuffer(String bufferKey, Message message) {
diff --git a/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/MemoryExtractor.java b/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/MemoryExtractor.java
index ce72f54a..5ea6c865 100644
--- a/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/MemoryExtractor.java
+++ b/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/MemoryExtractor.java
@@ -23,8 +23,8 @@
* Defines the memory extraction contract: batch extraction via {@link #extract(ExtractionRequest)}
* and context single-message extraction via {@link #addMessage(MemoryId, Message, ExtractionConfig)}.
*
- *
The primary implementation is {@link DefaultMemoryExtractor}. Decorators (e.g., tracing) wrap this
- * interface to add cross-cutting concerns.
+ *
The primary implementation is {@link DefaultMemoryExtractor}. Implementations publish observability signals at
+ * the relevant pipeline stages.
*/
public interface MemoryExtractor {
diff --git a/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/InsightLayer.java b/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/InsightLayer.java
index f8df7092..1a535c5b 100644
--- a/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/InsightLayer.java
+++ b/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/InsightLayer.java
@@ -18,11 +18,13 @@
import com.openmemind.ai.memory.core.data.MemoryItem;
import com.openmemind.ai.memory.core.data.enums.InsightAnalysisMode;
import com.openmemind.ai.memory.core.data.enums.MemoryItemType;
+import com.openmemind.ai.memory.core.extraction.insight.observation.InsightLayerObservation;
import com.openmemind.ai.memory.core.extraction.insight.scheduler.InsightBuildScheduler;
import com.openmemind.ai.memory.core.extraction.result.InsightResult;
import com.openmemind.ai.memory.core.extraction.result.MemoryItemResult;
import com.openmemind.ai.memory.core.extraction.step.InsightExtractStep;
import com.openmemind.ai.memory.core.store.MemoryStore;
+import io.micrometer.observation.ObservationRegistry;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
@@ -48,11 +50,19 @@ public class InsightLayer implements InsightExtractStep {
private final MemoryStore memoryStore;
private final InsightBuildScheduler scheduler;
private final Set unsupportedContentTypes;
+ private final ObservationRegistry observationRegistry;
public InsightLayer(MemoryStore memoryStore, InsightBuildScheduler scheduler) {
this(memoryStore, scheduler, Set.of());
}
+ public InsightLayer(
+ MemoryStore memoryStore,
+ InsightBuildScheduler scheduler,
+ ObservationRegistry observationRegistry) {
+ this(memoryStore, scheduler, Set.of(), observationRegistry);
+ }
+
/**
* @param memoryStore memory store
* @param scheduler insight build scheduler
@@ -63,10 +73,20 @@ public InsightLayer(
MemoryStore memoryStore,
InsightBuildScheduler scheduler,
Set unsupportedContentTypes) {
+ this(memoryStore, scheduler, unsupportedContentTypes, ObservationRegistry.NOOP);
+ }
+
+ public InsightLayer(
+ MemoryStore memoryStore,
+ InsightBuildScheduler scheduler,
+ Set unsupportedContentTypes,
+ ObservationRegistry observationRegistry) {
this.memoryStore = Objects.requireNonNull(memoryStore);
this.scheduler = Objects.requireNonNull(scheduler);
this.unsupportedContentTypes =
unsupportedContentTypes != null ? Set.copyOf(unsupportedContentTypes) : Set.of();
+ this.observationRegistry =
+ observationRegistry == null ? ObservationRegistry.NOOP : observationRegistry;
}
@Override
@@ -77,6 +97,14 @@ public Mono extract(MemoryId memoryId, MemoryItemResult memoryIte
@Override
public Mono extract(
MemoryId memoryId, MemoryItemResult memoryItemResult, String language) {
+ return InsightLayerObservation.observe(
+ observationRegistry,
+ memoryId,
+ () -> extractInternal(memoryId, memoryItemResult, language));
+ }
+
+ private Mono extractInternal(
+ MemoryId memoryId, MemoryItemResult memoryItemResult, String language) {
// Only take FACT type items whose content type supports insight building
var factItems =
memoryItemResult.newItems().stream()
diff --git a/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/generator/LlmInsightGenerator.java b/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/generator/LlmInsightGenerator.java
index 6dce10b3..9c4fe2fd 100644
--- a/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/generator/LlmInsightGenerator.java
+++ b/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/generator/LlmInsightGenerator.java
@@ -17,6 +17,8 @@
import com.openmemind.ai.memory.core.data.MemoryInsight;
import com.openmemind.ai.memory.core.data.MemoryInsightType;
import com.openmemind.ai.memory.core.data.MemoryItem;
+import com.openmemind.ai.memory.core.extraction.insight.generator.observation.LlmInsightGeneratorObservation;
+import com.openmemind.ai.memory.core.extraction.insight.generator.observation.LlmInsightGeneratorObservation.InsightGenerateDocument;
import com.openmemind.ai.memory.core.llm.ChatMessages;
import com.openmemind.ai.memory.core.llm.StructuredChatClient;
import com.openmemind.ai.memory.core.prompt.PromptRegistry;
@@ -24,6 +26,7 @@
import com.openmemind.ai.memory.core.prompt.extraction.insight.InsightLeafPrompts;
import com.openmemind.ai.memory.core.prompt.extraction.insight.InteractionGuideSynthesisPrompts;
import com.openmemind.ai.memory.core.prompt.extraction.insight.RootSynthesisPrompts;
+import io.micrometer.observation.ObservationRegistry;
import java.time.Duration;
import java.util.List;
import java.util.Objects;
@@ -45,18 +48,33 @@ public class LlmInsightGenerator implements InsightGenerator {
private final StructuredChatClient structuredChatClient;
private final PromptRegistry promptRegistry;
+ private final ObservationRegistry observationRegistry;
public LlmInsightGenerator(StructuredChatClient structuredChatClient) {
this(structuredChatClient, PromptRegistry.EMPTY);
}
+ public LlmInsightGenerator(
+ StructuredChatClient structuredChatClient, ObservationRegistry observationRegistry) {
+ this(structuredChatClient, PromptRegistry.EMPTY, observationRegistry);
+ }
+
public LlmInsightGenerator(
StructuredChatClient structuredChatClient, PromptRegistry promptRegistry) {
+ this(structuredChatClient, promptRegistry, ObservationRegistry.NOOP);
+ }
+
+ public LlmInsightGenerator(
+ StructuredChatClient structuredChatClient,
+ PromptRegistry promptRegistry,
+ ObservationRegistry observationRegistry) {
this.structuredChatClient =
Objects.requireNonNull(
structuredChatClient, "structuredChatClient must not be null");
this.promptRegistry =
Objects.requireNonNull(promptRegistry, "promptRegistry must not be null");
+ this.observationRegistry =
+ observationRegistry == null ? ObservationRegistry.NOOP : observationRegistry;
}
@Override
@@ -68,6 +86,29 @@ public Mono generatePoints(
int targetTokens,
String additionalContext,
String language) {
+ return LlmInsightGeneratorObservation.observeLeafPointGeneration(
+ observationRegistry,
+ insightType,
+ groupName,
+ () ->
+ generatePointsInternal(
+ insightType,
+ groupName,
+ existingPoints,
+ newItems,
+ targetTokens,
+ additionalContext,
+ language));
+ }
+
+ private Mono generatePointsInternal(
+ MemoryInsightType insightType,
+ String groupName,
+ List existingPoints,
+ List newItems,
+ int targetTokens,
+ String additionalContext,
+ String language) {
var template =
InsightLeafPrompts.build(
@@ -109,6 +150,29 @@ public Mono generateLeafPointOps(
int targetTokens,
String additionalContext,
String language) {
+ return LlmInsightGeneratorObservation.observeLeafPointOperations(
+ observationRegistry,
+ insightType,
+ groupName,
+ () ->
+ generateLeafPointOpsInternal(
+ insightType,
+ groupName,
+ existingPoints,
+ newItems,
+ targetTokens,
+ additionalContext,
+ language));
+ }
+
+ private Mono generateLeafPointOpsInternal(
+ MemoryInsightType insightType,
+ String groupName,
+ List existingPoints,
+ List newItems,
+ int targetTokens,
+ String additionalContext,
+ String language) {
var template =
InsightLeafPrompts.buildPointOps(
@@ -160,6 +224,28 @@ public Mono generateBranchSummary(
int targetTokens,
String additionalContext,
String language) {
+ return LlmInsightGeneratorObservation.observeAggregatePointGeneration(
+ observationRegistry,
+ InsightGenerateDocument.BRANCH,
+ insightType,
+ leafInsights,
+ () ->
+ generateBranchSummaryInternal(
+ insightType,
+ existingPoints,
+ leafInsights,
+ targetTokens,
+ additionalContext,
+ language));
+ }
+
+ private Mono generateBranchSummaryInternal(
+ MemoryInsightType insightType,
+ List existingPoints,
+ List leafInsights,
+ int targetTokens,
+ String additionalContext,
+ String language) {
var promptResult =
BranchAggregationPrompts.build(
@@ -209,6 +295,28 @@ public Mono generateBranchPointOps(
int targetTokens,
String additionalContext,
String language) {
+ return LlmInsightGeneratorObservation.observeAggregatePointOperations(
+ observationRegistry,
+ InsightGenerateDocument.BRANCH,
+ insightType,
+ leafInsights,
+ () ->
+ generateBranchPointOpsInternal(
+ insightType,
+ existingPoints,
+ leafInsights,
+ targetTokens,
+ additionalContext,
+ language));
+ }
+
+ private Mono generateBranchPointOpsInternal(
+ MemoryInsightType insightType,
+ List existingPoints,
+ List leafInsights,
+ int targetTokens,
+ String additionalContext,
+ String language) {
var promptResult =
BranchAggregationPrompts.buildPointOps(
@@ -258,6 +366,28 @@ public Mono generateRootSynthesis(
int targetTokens,
String additionalContext,
String language) {
+ return LlmInsightGeneratorObservation.observeAggregatePointGeneration(
+ observationRegistry,
+ InsightGenerateDocument.ROOT,
+ rootInsightType,
+ branchInsights,
+ () ->
+ generateRootSynthesisInternal(
+ rootInsightType,
+ existingPoints,
+ branchInsights,
+ targetTokens,
+ additionalContext,
+ language));
+ }
+
+ private Mono generateRootSynthesisInternal(
+ MemoryInsightType rootInsightType,
+ List existingPoints,
+ List branchInsights,
+ int targetTokens,
+ String additionalContext,
+ String language) {
var template =
switch (rootInsightType.name()) {
diff --git a/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/generator/observation/LlmInsightGeneratorObservation.java b/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/generator/observation/LlmInsightGeneratorObservation.java
new file mode 100644
index 00000000..970b5caf
--- /dev/null
+++ b/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/generator/observation/LlmInsightGeneratorObservation.java
@@ -0,0 +1,390 @@
+/*
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package com.openmemind.ai.memory.core.extraction.insight.generator.observation;
+
+import com.openmemind.ai.memory.core.data.MemoryInsight;
+import com.openmemind.ai.memory.core.data.MemoryInsightType;
+import com.openmemind.ai.memory.core.data.PointOperation;
+import com.openmemind.ai.memory.core.extraction.insight.generator.InsightPointGenerateResponse;
+import com.openmemind.ai.memory.core.extraction.insight.generator.InsightPointOpsResponse;
+import com.openmemind.ai.memory.core.observation.MemoryObservation;
+import io.micrometer.common.KeyValues;
+import io.micrometer.common.docs.KeyName;
+import io.micrometer.observation.Observation;
+import io.micrometer.observation.ObservationConvention;
+import io.micrometer.observation.ObservationRegistry;
+import io.micrometer.observation.docs.ObservationDocumentation;
+import java.util.List;
+import java.util.function.Function;
+import java.util.function.Supplier;
+import reactor.core.publisher.Mono;
+
+/** Observation contracts for LlmInsightGenerator. */
+public final class LlmInsightGeneratorObservation {
+
+ private LlmInsightGeneratorObservation() {}
+
+ public static Mono observeLeafPointGeneration(
+ ObservationRegistry observationRegistry,
+ MemoryInsightType insightType,
+ String groupName,
+ Supplier> operation) {
+ return observeLeafPointGeneration(
+ observationRegistry, insightType, groupName, ignored -> operation.get());
+ }
+
+ public static Mono observeLeafPointGeneration(
+ ObservationRegistry observationRegistry,
+ MemoryInsightType insightType,
+ String groupName,
+ Function>
+ operation) {
+ return observePointGeneration(
+ observationRegistry,
+ InsightGenerateDocument.LEAF,
+ () -> InsightGenerateObservationContext.leaf(insightType, groupName),
+ operation);
+ }
+
+ public static Mono observeLeafPointOperations(
+ ObservationRegistry observationRegistry,
+ MemoryInsightType insightType,
+ String groupName,
+ Supplier> operation) {
+ return observeLeafPointOperations(
+ observationRegistry, insightType, groupName, ignored -> operation.get());
+ }
+
+ public static Mono observeLeafPointOperations(
+ ObservationRegistry observationRegistry,
+ MemoryInsightType insightType,
+ String groupName,
+ Function> operation) {
+ return observePointOperations(
+ observationRegistry,
+ InsightGenerateDocument.LEAF,
+ () -> InsightGenerateObservationContext.leaf(insightType, groupName),
+ operation);
+ }
+
+ public static Mono observeAggregatePointGeneration(
+ ObservationRegistry observationRegistry,
+ InsightGenerateDocument document,
+ MemoryInsightType insightType,
+ List insights,
+ Supplier> operation) {
+ return observeAggregatePointGeneration(
+ observationRegistry, document, insightType, insights, ignored -> operation.get());
+ }
+
+ public static Mono observeAggregatePointGeneration(
+ ObservationRegistry observationRegistry,
+ InsightGenerateDocument document,
+ MemoryInsightType insightType,
+ List insights,
+ Function>
+ operation) {
+ return observePointGeneration(
+ observationRegistry,
+ document,
+ () -> InsightGenerateObservationContext.aggregate(document, insightType, insights),
+ operation);
+ }
+
+ public static Mono observeAggregatePointOperations(
+ ObservationRegistry observationRegistry,
+ InsightGenerateDocument document,
+ MemoryInsightType insightType,
+ List insights,
+ Supplier> operation) {
+ return observeAggregatePointOperations(
+ observationRegistry, document, insightType, insights, ignored -> operation.get());
+ }
+
+ public static Mono observeAggregatePointOperations(
+ ObservationRegistry observationRegistry,
+ InsightGenerateDocument document,
+ MemoryInsightType insightType,
+ List insights,
+ Function> operation) {
+ return observePointOperations(
+ observationRegistry,
+ document,
+ () -> InsightGenerateObservationContext.aggregate(document, insightType, insights),
+ operation);
+ }
+
+ private static Mono observePointGeneration(
+ ObservationRegistry observationRegistry,
+ InsightGenerateDocument document,
+ Supplier contextFactory,
+ Function>
+ operation) {
+ return MemoryObservation.mono(
+ observationRegistry,
+ document,
+ InsightGenerateConvention.of(document),
+ contextFactory,
+ context -> operation.apply(context).doOnNext(context::recordPointResponse));
+ }
+
+ private static Mono observePointOperations(
+ ObservationRegistry observationRegistry,
+ InsightGenerateDocument document,
+ Supplier contextFactory,
+ Function> operation) {
+ return MemoryObservation.mono(
+ observationRegistry,
+ document,
+ InsightGenerateConvention.of(document),
+ contextFactory,
+ context -> operation.apply(context).doOnNext(context::recordOpsResponse));
+ }
+
+ public enum InsightGenerateDocument implements ObservationDocumentation {
+ LEAF("memind.extraction.insight.generate.leaf"),
+ BRANCH("memind.extraction.insight.generate.branch"),
+ ROOT("memind.extraction.insight.generate.root");
+
+ private final String name;
+
+ InsightGenerateDocument(String name) {
+ this.name = name;
+ }
+
+ @Override
+ public String getName() {
+ return name;
+ }
+
+ @Override
+ public Class extends ObservationConvention extends Observation.Context>>
+ getDefaultConvention() {
+ return InsightGenerateConvention.class;
+ }
+
+ @Override
+ public KeyName[] getHighCardinalityKeyNames() {
+ return switch (this) {
+ case LEAF ->
+ new KeyName[] {
+ HighCardinalityKeyNames.INSIGHT_TYPE,
+ HighCardinalityKeyNames.GROUP_NAME,
+ HighCardinalityKeyNames.POINT_COUNT,
+ HighCardinalityKeyNames.ADD_COUNT,
+ HighCardinalityKeyNames.UPDATE_COUNT,
+ HighCardinalityKeyNames.DELETE_COUNT
+ };
+ case BRANCH ->
+ new KeyName[] {
+ HighCardinalityKeyNames.INSIGHT_TYPE,
+ HighCardinalityKeyNames.LEAF_COUNT,
+ HighCardinalityKeyNames.POINT_COUNT,
+ HighCardinalityKeyNames.ADD_COUNT,
+ HighCardinalityKeyNames.UPDATE_COUNT,
+ HighCardinalityKeyNames.DELETE_COUNT
+ };
+ case ROOT ->
+ new KeyName[] {
+ HighCardinalityKeyNames.INSIGHT_TYPE,
+ HighCardinalityKeyNames.LEAF_COUNT,
+ HighCardinalityKeyNames.POINT_COUNT
+ };
+ };
+ }
+ }
+
+ public enum HighCardinalityKeyNames implements KeyName {
+ INSIGHT_TYPE {
+ @Override
+ public String asString() {
+ return "memind.extraction.insight_type";
+ }
+ },
+ GROUP_NAME {
+ @Override
+ public String asString() {
+ return "memind.extraction.insight_group_name";
+ }
+ },
+ LEAF_COUNT {
+ @Override
+ public String asString() {
+ return "memind.extraction.insight_leaf_count";
+ }
+ },
+ POINT_COUNT {
+ @Override
+ public String asString() {
+ return "memind.extraction.insight_point_count";
+ }
+ },
+ ADD_COUNT {
+ @Override
+ public String asString() {
+ return "memind.extraction.insight_add_count";
+ }
+ },
+ UPDATE_COUNT {
+ @Override
+ public String asString() {
+ return "memind.extraction.insight_update_count";
+ }
+ },
+ DELETE_COUNT {
+ @Override
+ public String asString() {
+ return "memind.extraction.insight_delete_count";
+ }
+ };
+ }
+
+ public static final class InsightGenerateObservationContext extends Observation.Context {
+
+ private final InsightGenerateDocument document;
+ private final String insightTypeName;
+ private final String groupName;
+ private final Integer leafCount;
+ private Integer pointCount;
+ private Integer addCount;
+ private Integer updateCount;
+ private Integer deleteCount;
+
+ public InsightGenerateObservationContext(
+ InsightGenerateDocument document,
+ String insightTypeName,
+ String groupName,
+ Integer leafCount) {
+ this.document = document;
+ this.insightTypeName = insightTypeName;
+ this.groupName = groupName;
+ this.leafCount = leafCount;
+ }
+
+ public static InsightGenerateObservationContext leaf(
+ MemoryInsightType insightType, String groupName) {
+ return new InsightGenerateObservationContext(
+ InsightGenerateDocument.LEAF,
+ insightType == null ? "" : insightType.name(),
+ groupName == null ? "" : groupName,
+ null);
+ }
+
+ public static InsightGenerateObservationContext aggregate(
+ InsightGenerateDocument document,
+ MemoryInsightType insightType,
+ List insights) {
+ return new InsightGenerateObservationContext(
+ document,
+ insightType == null ? "" : insightType.name(),
+ null,
+ insights == null ? 0 : insights.size());
+ }
+
+ public void recordPointResponse(InsightPointGenerateResponse response) {
+ this.pointCount =
+ response == null || response.points() == null ? 0 : response.points().size();
+ }
+
+ public void recordOpsResponse(InsightPointOpsResponse response) {
+ var operations = response == null ? List.of() : response.operations();
+ this.addCount = countOperations(operations, PointOperation.OpType.ADD);
+ this.updateCount = countOperations(operations, PointOperation.OpType.UPDATE);
+ this.deleteCount = countOperations(operations, PointOperation.OpType.DELETE);
+ }
+
+ private static int countOperations(
+ List operations, PointOperation.OpType opType) {
+ return (int) operations.stream().filter(operation -> operation.op() == opType).count();
+ }
+ }
+
+ public static final class InsightGenerateConvention
+ implements ObservationConvention {
+
+ public static final InsightGenerateConvention LEAF =
+ new InsightGenerateConvention(InsightGenerateDocument.LEAF);
+ public static final InsightGenerateConvention BRANCH =
+ new InsightGenerateConvention(InsightGenerateDocument.BRANCH);
+ public static final InsightGenerateConvention ROOT =
+ new InsightGenerateConvention(InsightGenerateDocument.ROOT);
+
+ private final InsightGenerateDocument document;
+
+ public InsightGenerateConvention(InsightGenerateDocument document) {
+ this.document = document;
+ }
+
+ public static InsightGenerateConvention of(InsightGenerateDocument document) {
+ return switch (document) {
+ case LEAF -> LEAF;
+ case BRANCH -> BRANCH;
+ case ROOT -> ROOT;
+ };
+ }
+
+ @Override
+ public String getName() {
+ return document.getName();
+ }
+
+ @Override
+ public String getContextualName(InsightGenerateObservationContext context) {
+ return context.document.getName();
+ }
+
+ @Override
+ public KeyValues getHighCardinalityKeyValues(InsightGenerateObservationContext context) {
+ var values =
+ KeyValues.of(
+ HighCardinalityKeyNames.INSIGHT_TYPE.withValue(
+ context.insightTypeName));
+ if (context.groupName != null) {
+ values =
+ values.and(HighCardinalityKeyNames.GROUP_NAME.withValue(context.groupName));
+ }
+ if (context.leafCount != null) {
+ values =
+ values.and(
+ HighCardinalityKeyNames.LEAF_COUNT.withValue(
+ String.valueOf(context.leafCount)));
+ }
+ if (context.pointCount != null) {
+ values =
+ values.and(
+ HighCardinalityKeyNames.POINT_COUNT.withValue(
+ String.valueOf(context.pointCount)));
+ }
+ if (context.addCount != null) {
+ values =
+ values.and(
+ HighCardinalityKeyNames.ADD_COUNT.withValue(
+ String.valueOf(context.addCount)))
+ .and(
+ HighCardinalityKeyNames.UPDATE_COUNT.withValue(
+ String.valueOf(context.updateCount)))
+ .and(
+ HighCardinalityKeyNames.DELETE_COUNT.withValue(
+ String.valueOf(context.deleteCount)));
+ }
+ return values;
+ }
+
+ @Override
+ public boolean supportsContext(Observation.Context context) {
+ return context instanceof InsightGenerateObservationContext generateContext
+ && generateContext.document == document;
+ }
+ }
+}
diff --git a/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/group/LlmInsightGroupClassifier.java b/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/group/LlmInsightGroupClassifier.java
index c5eb0c86..fb0eea20 100644
--- a/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/group/LlmInsightGroupClassifier.java
+++ b/memind-core/src/main/java/com/openmemind/ai/memory/core/extraction/insight/group/LlmInsightGroupClassifier.java
@@ -15,10 +15,12 @@
import com.openmemind.ai.memory.core.data.MemoryInsightType;
import com.openmemind.ai.memory.core.data.MemoryItem;
+import com.openmemind.ai.memory.core.extraction.insight.group.observation.LlmInsightGroupClassifierObservation;
import com.openmemind.ai.memory.core.llm.ChatMessages;
import com.openmemind.ai.memory.core.llm.StructuredChatClient;
import com.openmemind.ai.memory.core.prompt.PromptRegistry;
import com.openmemind.ai.memory.core.prompt.extraction.insight.InsightGroupPrompts;
+import io.micrometer.observation.ObservationRegistry;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Comparator;
@@ -81,18 +83,33 @@ public class LlmInsightGroupClassifier implements InsightGroupClassifier {
private final StructuredChatClient structuredChatClient;
private final PromptRegistry promptRegistry;
+ private final ObservationRegistry observationRegistry;
public LlmInsightGroupClassifier(StructuredChatClient structuredChatClient) {
this(structuredChatClient, PromptRegistry.EMPTY);
}
+ public LlmInsightGroupClassifier(
+ StructuredChatClient structuredChatClient, ObservationRegistry observationRegistry) {
+ this(structuredChatClient, PromptRegistry.EMPTY, observationRegistry);
+ }
+
public LlmInsightGroupClassifier(
StructuredChatClient structuredChatClient, PromptRegistry promptRegistry) {
+ this(structuredChatClient, promptRegistry, ObservationRegistry.NOOP);
+ }
+
+ public LlmInsightGroupClassifier(
+ StructuredChatClient structuredChatClient,
+ PromptRegistry promptRegistry,
+ ObservationRegistry observationRegistry) {
this.structuredChatClient =
Objects.requireNonNull(
structuredChatClient, "structuredChatClient must not be null");
this.promptRegistry =
Objects.requireNonNull(promptRegistry, "promptRegistry must not be null");
+ this.observationRegistry =
+ observationRegistry == null ? ObservationRegistry.NOOP : observationRegistry;
}
@Override
@@ -119,6 +136,25 @@ public Mono