Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions memind-core/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,15 @@
<name>Memind - Core</name>

<dependencies>
<dependency>
<groupId>io.micrometer</groupId>
<artifactId>micrometer-observation</artifactId>
</dependency>
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-core-micrometer</artifactId>
</dependency>

<!-- Reactor Core for Flux/Mono -->
<dependency>
<groupId>io.projectreactor</groupId>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,27 +16,20 @@
import com.openmemind.ai.memory.core.DefaultMemory;
import com.openmemind.ai.memory.core.Memory;
import com.openmemind.ai.memory.core.buffer.MemoryBuffer;
import com.openmemind.ai.memory.core.extraction.MemoryExtractor;
import com.openmemind.ai.memory.core.extraction.insight.tree.BubbleTrackerStore;
import com.openmemind.ai.memory.core.llm.ChatClientRegistry;
import com.openmemind.ai.memory.core.llm.ChatClientSlot;
import com.openmemind.ai.memory.core.llm.StructuredChatClient;
import com.openmemind.ai.memory.core.llm.rerank.NoopReranker;
import com.openmemind.ai.memory.core.llm.rerank.Reranker;
import com.openmemind.ai.memory.core.metrics.MemoryMetricsRecorder;
import com.openmemind.ai.memory.core.metrics.NoopMemoryMetricsRecorder;
import com.openmemind.ai.memory.core.plugin.RawDataPlugin;
import com.openmemind.ai.memory.core.prompt.PromptRegistry;
import com.openmemind.ai.memory.core.resource.ContentParserRegistry;
import com.openmemind.ai.memory.core.resource.ResourceFetcher;
import com.openmemind.ai.memory.core.retrieval.MemoryRetriever;
import com.openmemind.ai.memory.core.store.MemoryStore;
import com.openmemind.ai.memory.core.textsearch.MemoryTextSearch;
import com.openmemind.ai.memory.core.tracing.MemoryObserver;
import com.openmemind.ai.memory.core.tracing.NoopMemoryObserver;
import com.openmemind.ai.memory.core.tracing.decorator.TracingMemoryExtractor;
import com.openmemind.ai.memory.core.tracing.decorator.TracingMemoryRetriever;
import com.openmemind.ai.memory.core.vector.MemoryVector;
import io.micrometer.observation.ObservationRegistry;
import java.util.ArrayList;
import java.util.EnumMap;
import java.util.IdentityHashMap;
Expand Down Expand Up @@ -64,8 +57,7 @@ public final class DefaultMemoryBuilder implements MemoryBuilder {
private BubbleTrackerStore bubbleTrackerStore;
private final List<RawDataPlugin> rawDataPlugins = new ArrayList<>();
private MemoryBuildOptions options = MemoryBuildOptions.defaults();
private MemoryObserver memoryObserver = new NoopMemoryObserver();
private MemoryMetricsRecorder memoryMetricsRecorder = NoopMemoryMetricsRecorder.INSTANCE;
private ObservationRegistry observationRegistry = ObservationRegistry.NOOP;
private boolean externallyManaged;

@Override
Expand Down Expand Up @@ -155,14 +147,9 @@ public MemoryBuilder options(MemoryBuildOptions options) {
}

@Override
public MemoryBuilder memoryObserver(MemoryObserver observer) {
this.memoryObserver = Objects.requireNonNull(observer, "observer");
return this;
}

@Override
public MemoryBuilder memoryMetricsRecorder(MemoryMetricsRecorder recorder) {
this.memoryMetricsRecorder = Objects.requireNonNull(recorder, "recorder");
public MemoryBuilder observationRegistry(ObservationRegistry observationRegistry) {
this.observationRegistry =
Objects.requireNonNull(observationRegistry, "observationRegistry");
return this;
}

Expand Down Expand Up @@ -194,21 +181,12 @@ public Memory build() {
resourceFetcher,
List.copyOf(rawDataPlugins),
bubbleTrackerStore,
memoryObserver,
memoryMetricsRecorder,
observationRegistry,
sanitization.memoryThreadForcedDisableReason());
MemoryExtractionAssembly extractionAssembly =
new MemoryExtractionAssembler().assemble(context);
MemoryExtractor pipeline =
tracingExtractor(
extractionAssembly.pipeline(),
context.memoryObserver(),
context.memoryMetricsRecorder());
MemoryRetriever memoryRetriever =
tracingRetriever(
new MemoryRetrievalAssembler().assemble(context),
context.memoryObserver(),
context.memoryMetricsRecorder());
var pipeline = extractionAssembly.pipeline();
var memoryRetriever = new MemoryRetrievalAssembler().assemble(context);
AutoCloseable lifecycle =
externallyManaged
? lifecycle(extractionAssembly.lifecycle())
Expand All @@ -231,27 +209,6 @@ public Memory build() {
extractionAssembly.memoryThreadLayer());
}

private MemoryExtractor tracingExtractor(
MemoryExtractor extractor, MemoryObserver observer, MemoryMetricsRecorder recorder) {
if (!hasObservability(observer, recorder) || extractor instanceof TracingMemoryExtractor) {
return extractor;
}
return new TracingMemoryExtractor(extractor, observer, recorder);
}

private MemoryRetriever tracingRetriever(
MemoryRetriever retriever, MemoryObserver observer, MemoryMetricsRecorder recorder) {
if (!hasObservability(observer, recorder) || retriever instanceof TracingMemoryRetriever) {
return retriever;
}
return new TracingMemoryRetriever(retriever, observer);
}

private boolean hasObservability(MemoryObserver observer, MemoryMetricsRecorder recorder) {
return !(observer instanceof NoopMemoryObserver)
|| !(recorder instanceof NoopMemoryMetricsRecorder);
}

MemoryBuildOptions buildOptions() {
return options;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,17 +20,14 @@
import com.openmemind.ai.memory.core.extraction.insight.tree.BubbleTrackerStore;
import com.openmemind.ai.memory.core.llm.ChatClientRegistry;
import com.openmemind.ai.memory.core.llm.rerank.Reranker;
import com.openmemind.ai.memory.core.metrics.MemoryMetricsRecorder;
import com.openmemind.ai.memory.core.metrics.NoopMemoryMetricsRecorder;
import com.openmemind.ai.memory.core.plugin.RawDataPlugin;
import com.openmemind.ai.memory.core.prompt.PromptRegistry;
import com.openmemind.ai.memory.core.resource.ContentParserRegistry;
import com.openmemind.ai.memory.core.resource.ResourceFetcher;
import com.openmemind.ai.memory.core.store.MemoryStore;
import com.openmemind.ai.memory.core.textsearch.MemoryTextSearch;
import com.openmemind.ai.memory.core.tracing.MemoryObserver;
import com.openmemind.ai.memory.core.tracing.NoopMemoryObserver;
import com.openmemind.ai.memory.core.vector.MemoryVector;
import io.micrometer.observation.ObservationRegistry;
import java.util.List;
import java.util.Objects;
import java.util.Optional;
Expand All @@ -48,8 +45,7 @@ record MemoryAssemblyContext(
ResourceFetcher resourceFetcher,
List<RawDataPlugin> rawDataPlugins,
BubbleTrackerStore bubbleTrackerStore,
MemoryObserver memoryObserver,
MemoryMetricsRecorder memoryMetricsRecorder,
ObservationRegistry observationRegistry,
Optional<String> memoryThreadForcedDisableReason) {

MemoryAssemblyContext {
Expand All @@ -70,50 +66,14 @@ record MemoryAssemblyContext(
Objects.requireNonNull(promptRegistry, "promptRegistry");
Objects.requireNonNull(options, "options");
rawDataPlugins = List.copyOf(Objects.requireNonNull(rawDataPlugins, "rawDataPlugins"));
memoryObserver = memoryObserver != null ? memoryObserver : new NoopMemoryObserver();
memoryMetricsRecorder =
memoryMetricsRecorder != null
? memoryMetricsRecorder
: NoopMemoryMetricsRecorder.INSTANCE;
observationRegistry =
observationRegistry != null ? observationRegistry : ObservationRegistry.NOOP;
memoryThreadForcedDisableReason =
memoryThreadForcedDisableReason != null
? memoryThreadForcedDisableReason
: Optional.empty();
}

MemoryAssemblyContext(
ChatClientRegistry chatClientRegistry,
MemoryStore memoryStore,
MemoryBuffer memoryBuffer,
MemoryTextSearch textSearch,
MemoryVector memoryVector,
Reranker reranker,
PromptRegistry promptRegistry,
MemoryBuildOptions options,
ContentParserRegistry contentParserRegistry,
ResourceFetcher resourceFetcher,
List<RawDataPlugin> rawDataPlugins,
BubbleTrackerStore bubbleTrackerStore,
MemoryObserver memoryObserver,
Optional<String> memoryThreadForcedDisableReason) {
this(
chatClientRegistry,
memoryStore,
memoryBuffer,
textSearch,
memoryVector,
reranker,
promptRegistry,
options,
contentParserRegistry,
resourceFetcher,
rawDataPlugins,
bubbleTrackerStore,
memoryObserver,
NoopMemoryMetricsRecorder.INSTANCE,
memoryThreadForcedDisableReason);
}

InsightBuffer insightBuffer() {
return memoryBuffer.insightBuffer();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,15 +19,14 @@
import com.openmemind.ai.memory.core.llm.ChatClientSlot;
import com.openmemind.ai.memory.core.llm.StructuredChatClient;
import com.openmemind.ai.memory.core.llm.rerank.Reranker;
import com.openmemind.ai.memory.core.metrics.MemoryMetricsRecorder;
import com.openmemind.ai.memory.core.plugin.RawDataPlugin;
import com.openmemind.ai.memory.core.prompt.PromptRegistry;
import com.openmemind.ai.memory.core.resource.ContentParserRegistry;
import com.openmemind.ai.memory.core.resource.ResourceFetcher;
import com.openmemind.ai.memory.core.store.MemoryStore;
import com.openmemind.ai.memory.core.textsearch.MemoryTextSearch;
import com.openmemind.ai.memory.core.tracing.MemoryObserver;
import com.openmemind.ai.memory.core.vector.MemoryVector;
import io.micrometer.observation.ObservationRegistry;

/**
* Builds a {@link Memory} instance from runtime components.
Expand Down Expand Up @@ -60,11 +59,7 @@ public interface MemoryBuilder {

MemoryBuilder options(MemoryBuildOptions options);

default MemoryBuilder memoryObserver(MemoryObserver observer) {
return this;
}

default MemoryBuilder memoryMetricsRecorder(MemoryMetricsRecorder recorder) {
default MemoryBuilder observationRegistry(ObservationRegistry observationRegistry) {
return this;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -95,7 +95,6 @@
import com.openmemind.ai.memory.core.store.graph.NoOpItemGraphCommitOperations;
import com.openmemind.ai.memory.core.store.thread.NoOpThreadEnrichmentInputStore;
import com.openmemind.ai.memory.core.store.thread.ThreadEnrichmentInputStore;
import com.openmemind.ai.memory.core.tracing.decorator.TracingItemGraphMaterializer;
import com.openmemind.ai.memory.core.utils.IdUtils;
import java.util.ArrayList;
import java.util.HashMap;
Expand Down Expand Up @@ -141,13 +140,15 @@ MemoryExtractionAssembly assemble(MemoryAssemblyContext context) {
captionGenerator,
context.memoryStore(),
context.memoryVector(),
context.options().extraction().rawdata().vectorBatchSize());
context.options().extraction().rawdata().vectorBatchSize(),
context.observationRegistry());

MemoryItemExtractor itemExtractor =
createMemoryItemExtractor(registry, processors, context.promptRegistry());
MemoryItemDeduplicator deduplicator =
new CompositeDeduplicator(
List.of(new HashBasedDeduplicator(context.memoryStore())));
List.of(new HashBasedDeduplicator(context.memoryStore())),
context.observationRegistry());
ItemGraphMaterializer graphMaterializer = graphMaterializer(context);
MemoryItemLayer memoryItemLayer =
new MemoryItemLayer(
Expand All @@ -157,19 +158,22 @@ MemoryExtractionAssembly assemble(MemoryAssemblyContext context) {
context.memoryVector(),
IdUtils.snowflake(),
null,
graphMaterializer);
graphMaterializer,
context.observationRegistry());
MemoryItemExtractStep memoryItemStep = memoryItemLayer;
MemoryThreadLayer memoryThreadLayer = null;

InsightGenerator insightGenerator =
new LlmInsightGenerator(
registry.resolve(ChatClientSlot.INSIGHT_GENERATOR),
context.promptRegistry());
context.promptRegistry(),
context.observationRegistry());
InsightGraphAssistant insightGraphAssistant = insightGraphAssistant(context);
InsightGroupClassifier insightGroupClassifier =
new LlmInsightGroupClassifier(
registry.resolve(ChatClientSlot.INSIGHT_GROUP_CLASSIFIER),
context.promptRegistry());
context.promptRegistry(),
context.observationRegistry());
var identityManager = new InsightPointIdentityManager();
var evidenceNormalizer = new InsightPointEvidenceNormalizer();
BubbleTrackerStore bubbleTrackerStore =
Expand Down Expand Up @@ -201,12 +205,13 @@ MemoryExtractionAssembly assemble(MemoryAssemblyContext context) {
identityManager,
evidenceNormalizer,
insightGraphAssistant,
null);
context.observationRegistry());
InsightLayer insightLayer =
new InsightLayer(
context.memoryStore(),
insightBuildScheduler,
unsupportedInsightTypes(processors));
unsupportedInsightTypes(processors),
context.observationRegistry());

ContextCommitDetector contextCommitDetector =
new LlmContextCommitDetector(
Expand Down Expand Up @@ -320,7 +325,8 @@ MemoryExtractionAssembly assemble(MemoryAssemblyContext context) {
resolveResourceFetcher(context.resourceFetcher()),
ingestionPolicyRegistry,
context.options().extraction().rawdata(),
context.options().extraction().item());
context.options().extraction().item(),
context.observationRegistry());
return new MemoryExtractionAssembly(
pipeline, insightLayer, extractionLifecycle, memoryThreadLayer);
}
Expand Down Expand Up @@ -405,8 +411,9 @@ yield new ConservativeHeuristicEntityResolutionStrategy(
planner,
context.memoryStore().itemGraphCommitOperations(),
derivedMaintainer,
context.options().extraction().item().graph());
return new TracingItemGraphMaterializer(graphMaterializer, context.memoryObserver());
context.options().extraction().item().graph(),
context.observationRegistry());
return graphMaterializer;
}

private ConversationContentProcessor conversationProcessor(
Expand Down
Loading
Loading