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
15 changes: 15 additions & 0 deletions .tegami/2026-08-05-retrieval-admission-control.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
---
packages:
orgmemory: patch
subject: Faster, fairer assistant retrieval under load
---

## Improvements

Assistant knowledge retrieval now admits snapshot queries through one fair
process-wide limit instead of per-request batches, so concurrent
conversations can no longer exhaust the database connection pool and stall at
the turn timeout. The API connection pool is right-sized for the production
host, retrieval breadth returns to the upstream LightRAG default, and new
payload-free timing stages make the previously unattributed portion of
time-to-first-token observable.
1 change: 1 addition & 0 deletions apps/api/build.gradle.kts
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ dependencies {
// not have.
testImplementation("io.micrometer:micrometer-registry-otlp")
testImplementation("io.opentelemetry:opentelemetry-sdk")
testImplementation("io.opentelemetry:opentelemetry-sdk-testing")
testRuntimeOnly("org.junit.platform:junit-platform-launcher")
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
import com.orgmemory.core.assistant.AssistantAssetToolService;
import com.orgmemory.core.assistant.AssistantAssetTraceRecorder;
import com.orgmemory.core.assistant.AssistantService;
import com.orgmemory.core.assistant.observability.AssistantStageEventSink;
import com.orgmemory.core.assistant.observability.AssistantTurnEvent;
import com.orgmemory.core.assistant.observability.AssistantTurnMeterObservationHandler;
import com.orgmemory.core.assetregistry.AssetRegistryService;
Expand All @@ -15,7 +16,9 @@
import com.orgmemory.core.knowledge.search.PermissionAwareKnowledgeSearch;
import io.micrometer.core.instrument.MeterRegistry;
import io.micrometer.observation.ObservationRegistry;
import io.opentelemetry.api.OpenTelemetry;
import java.time.Clock;
import java.util.List;
import org.springframework.ai.chat.memory.ChatMemory;
import org.springframework.ai.chat.memory.ChatMemoryRepository;
import org.springframework.ai.chat.memory.MessageWindowChatMemory;
Expand All @@ -35,11 +38,18 @@ Clock assistantClock() {
}

@Bean
ChatMemory assistantChatMemory(ChatMemoryRepository repository) {
return MessageWindowChatMemory.builder()
ChatMemory assistantChatMemory(
ChatMemoryRepository repository,
AssistantStageEventSink stages,
AssistantProperties properties) {
ChatMemory memory = MessageWindowChatMemory.builder()
.chatMemoryRepository(repository)
.maxMessages(20)
.build();
return new ObservedChatMemory(
memory,
stages,
observedEngine(properties));
}

@Bean
Expand All @@ -62,9 +72,26 @@ AssistantService assistantService(
PermissionAwareKnowledgeSearch retrieval,
ChatModelPort chat,
ObservationRegistry observations,
AssistantProperties properties) {
AssistantProperties properties,
AssistantStageEventSink stages) {
return new AssistantService(
retrieval, chat, observations, observedEngine(properties));
retrieval,
chat,
observations,
observedEngine(properties),
stages);
}

@Bean
AssistantStageEventSink assistantStageEventSink(
OpenTelemetry openTelemetry,
MeterRegistry meters) {
return AssistantStageEventSink.composite(List.of(
AssistantStageEventSink.failureTolerant(
new OpenTelemetryAssistantStageEventSink(
openTelemetry)),
AssistantStageEventSink.failureTolerant(
new MicrometerAssistantStageEventSink(meters))));
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ record GraphRagQueryRuntimeProperties(
Integer maximumEmbeddingBatchSize,
Integer maximumKnowledgeSpaces,
Integer maximumConcurrentSpaces,
Integer retrievalAdmissionPermits,
Integer topK,
Integer chunkTopK,
Integer relatedChunkNumber,
Expand Down Expand Up @@ -45,7 +46,11 @@ record GraphRagQueryRuntimeProperties(
throw new IllegalArgumentException(
"maximumConcurrentSpaces must not exceed maximumKnowledgeSpaces");
}
topK = positive(topK, 60, "topK");
retrievalAdmissionPermits = positive(
retrievalAdmissionPermits,
4,
"retrievalAdmissionPermits");
topK = positive(topK, 40, "topK");
chunkTopK = positive(chunkTopK, 20, "chunkTopK");
relatedChunkNumber = positive(
relatedChunkNumber,
Expand Down Expand Up @@ -98,6 +103,7 @@ GraphRagRetrievalPolicy toPolicy() {
return new GraphRagRetrievalPolicy(
maximumKnowledgeSpaces,
maximumConcurrentSpaces,
retrievalAdmissionPermits,
topK,
chunkTopK,
relatedChunkNumber,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
package com.orgmemory.api.assistant;

import com.orgmemory.core.assistant.observability.AssistantStageEventSink;
import io.micrometer.core.instrument.MeterRegistry;
import io.micrometer.core.instrument.Timer;
import java.util.Locale;
import java.util.Objects;

/** Bounded-cardinality timer for assistant latency attribution stages. */
final class MicrometerAssistantStageEventSink
implements AssistantStageEventSink {

static final String STAGE_TIMER = "orgmemory.assistant.stage";

private final MeterRegistry registry;

MicrometerAssistantStageEventSink(MeterRegistry registry) {
this.registry = Objects.requireNonNull(registry, "registry");
}

@Override
public void emit(AssistantStageEvent event) {
Objects.requireNonNull(event, "event");
Timer.builder(STAGE_TIMER)
.description(
"Assistant latency stages above permission-aware retrieval")
.tag("engine", value(event.engine()))
.tag("stage", value(event.stage()))
.tag("outcome", value(event.outcome()))
.register(registry)
.record(event.duration());
}

private static String value(Enum<?> value) {
return value.name().toLowerCase(Locale.ROOT);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
package com.orgmemory.api.assistant;

import com.orgmemory.core.assistant.observability.AssistantStageEventSink;
import com.orgmemory.core.assistant.observability.AssistantStageEventSink.AssistantStageEvent;
import com.orgmemory.core.assistant.observability.AssistantTurnEvent;
import java.time.Duration;
import java.time.Instant;
import java.util.List;
import java.util.Objects;
import org.springframework.ai.chat.memory.ChatMemory;
import org.springframework.ai.chat.messages.Message;

/** Measures the history read without exposing its conversation or messages. */
final class ObservedChatMemory implements ChatMemory {

private final ChatMemory delegate;
private final AssistantStageEventSink events;
private final AssistantTurnEvent.RetrievalEngine engine;

ObservedChatMemory(
ChatMemory delegate,
AssistantStageEventSink events,
AssistantTurnEvent.RetrievalEngine engine) {
this.delegate = Objects.requireNonNull(delegate, "delegate");
this.events = Objects.requireNonNull(events, "events");
this.engine = Objects.requireNonNull(engine, "engine");
}

@Override
public void add(String conversationId, List<Message> messages) {
delegate.add(conversationId, messages);
}

@Override
public List<Message> get(String conversationId) {
long startedAt = System.nanoTime();
try {
List<Message> messages = delegate.get(conversationId);
emit(
AssistantStageEventSink.Outcome.SUCCEEDED,
startedAt,
null);
return messages;
} catch (RuntimeException | Error failure) {
emit(
AssistantStageEventSink.Outcome.FAILED,
startedAt,
"history_load_failed");
throw failure;
}
}

@Override
public void clear(String conversationId) {
delegate.clear(conversationId);
}

private void emit(
AssistantStageEventSink.Outcome outcome,
long startedAt,
String failureCode) {
events.emit(new AssistantStageEvent(
engine,
AssistantStageEventSink.Stage.CONVERSATION_HISTORY_LOAD,
outcome,
Duration.ofNanos(Math.max(
0L,
System.nanoTime() - startedAt)),
failureCode,
Instant.now()));
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
package com.orgmemory.api.assistant;

import com.orgmemory.core.assistant.observability.AssistantStageEventSink;
import io.opentelemetry.api.OpenTelemetry;
import io.opentelemetry.api.common.AttributeKey;
import io.opentelemetry.api.trace.Span;
import io.opentelemetry.api.trace.SpanKind;
import io.opentelemetry.api.trace.StatusCode;
import io.opentelemetry.api.trace.Tracer;
import java.time.Instant;
import java.util.Locale;
import java.util.Objects;
import java.util.concurrent.TimeUnit;

/** Payload-free OpenTelemetry adapter for assistant latency attribution. */
final class OpenTelemetryAssistantStageEventSink
implements AssistantStageEventSink {

static final String INSTRUMENTATION_SCOPE =
"com.orgmemory.assistant";
static final AttributeKey<String> ENGINE =
AttributeKey.stringKey("orgmemory.assistant.engine");
static final AttributeKey<String> STAGE =
AttributeKey.stringKey("orgmemory.assistant.stage");
static final AttributeKey<String> OUTCOME =
AttributeKey.stringKey("orgmemory.assistant.outcome");
static final AttributeKey<Long> DURATION_NANOS =
AttributeKey.longKey("orgmemory.assistant.duration_nanos");
static final AttributeKey<String> FAILURE_CODE =
AttributeKey.stringKey("orgmemory.assistant.failure_code");

private final Tracer tracer;

OpenTelemetryAssistantStageEventSink(OpenTelemetry openTelemetry) {
tracer = Objects.requireNonNull(openTelemetry, "openTelemetry")
.getTracer(INSTRUMENTATION_SCOPE);
}

@Override
public void emit(AssistantStageEvent event) {
Objects.requireNonNull(event, "event");
long endEpochNanos = epochNanos(event.occurredAt());
long startEpochNanos = Math.subtractExact(
endEpochNanos,
event.duration().toNanos());
Span span = tracer.spanBuilder(
"orgmemory.assistant." + value(event.stage()))
.setSpanKind(SpanKind.INTERNAL)
.setStartTimestamp(
startEpochNanos,
TimeUnit.NANOSECONDS)
.startSpan();
span.setAttribute(ENGINE, value(event.engine()));
span.setAttribute(STAGE, value(event.stage()));
span.setAttribute(OUTCOME, value(event.outcome()));
span.setAttribute(
DURATION_NANOS,
event.duration().toNanos());
if (event.failureCode() != null) {
span.setAttribute(FAILURE_CODE, event.failureCode());
}
if (event.outcome() == Outcome.FAILED) {
span.setStatus(StatusCode.ERROR);
}
span.end(endEpochNanos, TimeUnit.NANOSECONDS);
}

private static String value(Enum<?> value) {
return value.name().toLowerCase(Locale.ROOT);
}

private static long epochNanos(Instant instant) {
return Math.addExact(
Math.multiplyExact(
instant.getEpochSecond(),
1_000_000_000L),
instant.getNano());
}
}
4 changes: 2 additions & 2 deletions apps/api/src/main/resources/application-prod.yml
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,8 @@ spring:
password: ${ORGMEMORY_DB_PASSWORD}
hikari:
pool-name: orgmemory-api
maximum-pool-size: ${ORGMEMORY_API_DB_POOL_MAXIMUM_SIZE:12}
minimum-idle: ${ORGMEMORY_API_DB_POOL_MINIMUM_IDLE:2}
maximum-pool-size: ${ORGMEMORY_API_DB_POOL_MAXIMUM_SIZE:8}
minimum-idle: ${ORGMEMORY_API_DB_POOL_MINIMUM_IDLE:8}
connection-timeout: ${ORGMEMORY_DB_CONNECTION_TIMEOUT_MS:10000}
validation-timeout: ${ORGMEMORY_DB_VALIDATION_TIMEOUT_MS:5000}
max-lifetime: ${ORGMEMORY_DB_MAX_LIFETIME_MS:1800000}
Expand Down
3 changes: 2 additions & 1 deletion apps/api/src/main/resources/application.yml
Original file line number Diff line number Diff line change
Expand Up @@ -150,7 +150,8 @@ orgmemory:
keyword-cache-ttl: ${ORGMEMORY_GRAPH_QUERY_KEYWORD_CACHE_TTL:24h}
maximum-knowledge-spaces: ${ORGMEMORY_GRAPH_QUERY_MAXIMUM_SPACES:20}
maximum-concurrent-spaces: ${ORGMEMORY_GRAPH_QUERY_MAXIMUM_CONCURRENT_SPACES:4}
top-k: ${ORGMEMORY_GRAPH_QUERY_TOP_K:60}
retrieval-admission-permits: ${ORGMEMORY_GRAPH_QUERY_ADMISSION_PERMITS:4}
top-k: ${ORGMEMORY_GRAPH_QUERY_TOP_K:40}
chunk-top-k: ${ORGMEMORY_GRAPH_QUERY_CHUNK_TOP_K:20}
related-chunk-number: ${ORGMEMORY_GRAPH_QUERY_RELATED_CHUNKS:5}
maximum-graph-depth: ${ORGMEMORY_GRAPH_QUERY_MAXIMUM_DEPTH:1}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@ class MetricsDistributionTests {
"jvm.gc.pause",
"orgmemory.assistant.turn",
"orgmemory.assistant.time_to_first_token",
"orgmemory.assistant.stage",
"orgmemory.graph_rag.stage",
"gen_ai.client.operation"
})
Expand All @@ -63,6 +64,7 @@ void meterChartedAsAQuantilePublishesAHistogram(String name) {
strings = {
"http.server.requests",
"orgmemory.assistant.turn",
"orgmemory.assistant.stage",
"orgmemory.graph_rag.stage",
"gen_ai.client.operation"
})
Expand Down
Loading