Quarkus gRPC Client
The Quarkus extension provides gRPC clients for high-performance communication with Memory Service.
Setup
Add the gRPC client dependency:
<dependency>
<groupId>io.github.chirino.memory-service</groupId>
<artifactId>memory-service-proto-quarkus</artifactId>
<version>999-SNAPSHOT</version>
</dependency>
Configuration
# Configure the gRPC client for each service you need
quarkus.grpc.clients.conversations.host=localhost
quarkus.grpc.clients.conversations.port=9000
quarkus.grpc.clients.entries.host=localhost
quarkus.grpc.clients.entries.port=9000
quarkus.grpc.clients.search.host=localhost
quarkus.grpc.clients.search.port=9000
Calling as an asserted user
Use lowercase x-user-id metadata together with the API key or bearer credential. Derive an immutable stub for each request; do not attach changing user metadata to a shared injected stub:
import io.grpc.Metadata;
import io.grpc.stub.MetadataUtils;
Metadata requestMetadata = new Metadata();
requestMetadata.put(
Metadata.Key.of("x-api-key", Metadata.ASCII_STRING_MARSHALLER),
serviceApiKey);
requestMetadata.put(
Metadata.Key.of("x-user-id", Metadata.ASCII_STRING_MARSHALLER),
userId);
ConversationsServiceGrpc.ConversationsServiceBlockingStub requestStub =
conversationsClient.withInterceptors(
MetadataUtils.newAttachHeadersInterceptor(requestMetadata));
Conversation conversation = requestStub.createConversation(request);
The server honors x-user-id only when the authenticated client is listed in MEMORY_SERVICE_TRUSTED_USER_ID_CLIENTS. The metadata works across all normal user services and authorized event streams; admin and system methods ignore it. Identity is call metadata, so request messages do not contain RequestActor or another identity field.
Available Services
The Memory Service gRPC API is split into multiple services:
| Service | Description |
|---|---|
ConversationsService | CRUD operations for conversations |
EntriesService | List, append, and sync entries |
ConversationMembershipsService | Sharing and membership management |
OwnershipTransfersService | Ownership transfer operations |
SearchService | Semantic search and indexing |
ResponseRecorderService | Streaming response recording and resumption |
SystemService | Health checks |
Injecting Clients
ConversationsService
import io.github.chirino.memory.grpc.v1.*;
import io.quarkus.grpc.GrpcClient;
@GrpcClient("conversations")
ConversationsServiceGrpc.ConversationsServiceBlockingStub conversationsClient;
EntriesService
@GrpcClient("entries")
EntriesServiceGrpc.EntriesServiceBlockingStub entriesClient;
SearchService
@GrpcClient("search")
SearchServiceGrpc.SearchServiceBlockingStub searchClient;
Conversations API
List Conversations
import io.github.chirino.memory.grpc.v1.*;
ListConversationsResponse response = conversationsClient.listConversations(
ListConversationsRequest.newBuilder()
.setMode(ConversationListMode.LATEST_FORK)
.setPage(PageRequest.newBuilder()
.setPageSize(20)
.build())
.build()
);
for (ConversationSummary conv : response.getConversationsList()) {
System.out.println(conv.getTitle() + " - " + conv.getId());
}
Filtering by metadata
The gRPC API accepts typed ConversationMetadataFilter predicates. Because old servers silently ignore metadata_filters, check capabilities before sending them.
Inject a SystemService stub. The "conversations" client name reuses the same channel configuration as the conversations stub:
import io.github.chirino.memory.grpc.v1.*;
import io.quarkus.grpc.GrpcClient;
@GrpcClient("conversations")
SystemServiceGrpc.SystemServiceBlockingStub systemClient;
Then check capabilities and build the request:
import com.google.protobuf.Empty;
import io.github.chirino.memory.grpc.v1.*;
import java.util.List;
CapabilitiesResponse caps = systemClient.getCapabilities(Empty.getDefaultInstance());
CapabilitiesFeatures features = caps.getFeatures();
int filterVersion = features.getConversationMetadataFilterVersion();
int maxFilters = features.getMaxConversationMetadataFilters();
List<ConversationMetadataFilter> predicates = List.of(
ConversationMetadataFilter.newBuilder()
.setKey("status")
.setComparison(ConversationMetadataComparison.CONVERSATION_METADATA_COMPARISON_EQUAL)
.setValue("waiting")
.build(),
ConversationMetadataFilter.newBuilder()
.setKey("agent-id")
.setComparison(ConversationMetadataComparison.CONVERSATION_METADATA_COMPARISON_NOT_EQUAL)
.setValue("worker-2")
.build()
);
ListConversationsRequest.Builder req = ListConversationsRequest.newBuilder()
.setPage(PageRequest.newBuilder().setPageSize(20).build());
if (filterVersion >= 1 && maxFilters >= predicates.size()) {
req.addAllMetadataFilters(predicates);
} else if (predicates.size() == 1
&& predicates.get(0).getComparison()
== ConversationMetadataComparison.CONVERSATION_METADATA_COMPARISON_EQUAL) {
req.setMetadataFilterKey(predicates.get(0).getKey());
req.setMetadataFilterValue(predicates.get(0).getValue());
} else {
throw new IllegalStateException(
"server does not support repeatable metadata filters; upgrade the server or simplify the request to a single equality predicate");
}
ListConversationsResponse response = conversationsClient.listConversations(req.build());
The capability fields conversation_metadata_filter_version and max_conversation_metadata_filters decode to 0 on servers that predate repeatable filters. Do not send metadata_filters unless you have observed version >= 1 and a max large enough for your predicate count from the same server deployment.
Get Conversation
String conversationId = "d8e8a8c2-3767-4213-a590-c3a72bcbfdda";
Conversation conv = conversationsClient.getConversation(
GetConversationRequest.newBuilder()
.setConversationId(conversationId)
.build()
);
Create Conversation
import com.google.protobuf.Struct;
import com.google.protobuf.Value;
Conversation conv = conversationsClient.createConversation(
CreateConversationRequest.newBuilder()
.setTitle("My Conversation")
.setMetadata(Struct.newBuilder()
.putFields("topic", Value.newBuilder().setStringValue("support").build())
.build())
.build()
);
Delete Conversation
conversationsClient.deleteConversation(
DeleteConversationRequest.newBuilder()
.setConversationId(conversationId)
.build()
);
List Forks
ListForksResponse forks = conversationsClient.listForks(
ListForksRequest.newBuilder()
.setConversationId(conversationId)
.build()
);
for (String groupConversationId : forks.getConversationIdsList()) {
System.out.println("Group conversation: " + groupConversationId);
}
for (ConversationForkPoint point : forks.getForkPointsList()) {
System.out.println("Fork selector entry bytes: " + point.getEntryId());
for (ConversationForkOption option : point.getOptionsList()) {
System.out.println(" Option: " + option.getTitle());
}
}
The response is a complete navigation snapshot. Page the requested conversation’s entries separately with ListEntries.
Entries API
List Entries
ListEntriesResponse response = entriesClient.listEntries(
ListEntriesRequest.newBuilder()
.setConversationId(conversationId)
.setChannel(Channel.HISTORY)
.setEpochFilter("latest")
.setForks("none")
.setPage(PageRequest.newBuilder()
.setPageSize(50)
.build())
.build()
);
for (Entry entry : response.getEntriesList()) {
System.out.println(entry.getContentType() + ": " + entry.getContentList());
}
To open at the newest entries, set tail and then use
pageInfo.previousPageToken to request the adjacent older page:
ListEntriesResponse newest = entriesClient.listEntries(
ListEntriesRequest.newBuilder()
.setConversationId(conversationId)
.setChannel(Channel.HISTORY)
.setTail(true)
.setPage(PageRequest.newBuilder().setPageSize(50).build())
.build()
);
String olderToken = newest.getPageInfo().getPreviousPageToken();
if (!olderToken.isEmpty()) {
ListEntriesResponse older = entriesClient.listEntries(
ListEntriesRequest.newBuilder()
.setConversationId(conversationId)
.setChannel(Channel.HISTORY)
.setBeforePageToken(olderToken)
.setPage(PageRequest.newBuilder().setPageSize(50).build())
.build()
);
}
previousPageToken maps to REST beforeCursor; nextPageToken maps to REST
afterCursor. page.pageToken, beforePageToken, and tail=true are mutually
exclusive.
Append Entry
import com.google.protobuf.Value;
import com.google.protobuf.Struct;
Entry entry = entriesClient.appendEntry(
AppendEntryRequest.newBuilder()
.setConversationId(conversationId)
.setEntry(CreateEntryRequest.newBuilder()
.setUserId("user123")
.setChannel(Channel.HISTORY)
.setContentType("history")
.addContent(Value.newBuilder()
.setStructValue(Struct.newBuilder()
.putFields("text", Value.newBuilder()
.setStringValue("Hello, how can I help?").build())
.putFields("role", Value.newBuilder()
.setStringValue("USER").build())
.build())
.build())
.build())
.build()
);
Sync Agent Memory
SyncEntriesResponse response = entriesClient.syncEntries(
SyncEntriesRequest.newBuilder()
.setConversationId(conversationId)
.setEntry(CreateEntryRequest.newBuilder()
.setChannel(Channel.CONTEXT)
.setContentType("LC4J")
.addAllContent(memoryContent)
.build())
.build()
);
if (response.getNoOp()) {
System.out.println("Memory already up to date");
} else if (response.getEpochIncremented()) {
System.out.println("New epoch started: " + response.getEpoch());
}
Search API
Search Conversations
SearchEntriesResponse response = searchClient.searchConversations(
SearchEntriesRequest.newBuilder()
.setQuery("authentication configuration")
.setLimit(20)
.setIncludeEntry(true)
.build()
);
for (SearchResult result : response.getResultsList()) {
System.out.println("Score: " + result.getScore() +
" - " + result.getConversationTitle());
}
Response Recorder (Streaming)
The ResponseRecorderService handles streaming response recording and resumption for interrupted connections:
@GrpcClient("responserecorder")
ResponseRecorderServiceGrpc.ResponseRecorderServiceStub recorderClient;
// Check if response recorder is enabled
IsEnabledResponse enabled = recorderClient.isEnabled(Empty.getDefaultInstance());
// Check which conversations have recordings in progress
CheckRecordingsResponse check = recorderClient.checkRecordings(
CheckRecordingsRequest.newBuilder()
.addConversationIds(conversationId)
.build()
);
// Replay response tokens
recorderClient.replay(
ReplayRequest.newBuilder()
.setConversationId(conversationId)
.build(),
new StreamObserver<ReplayResponse>() {
@Override
public void onNext(ReplayResponse response) {
System.out.print(response.getContent());
}
@Override
public void onError(Throwable t) {
t.printStackTrace();
}
@Override
public void onCompleted() {
System.out.println("\nReplay completed");
}
}
);
Error Handling
import io.grpc.StatusRuntimeException;
try {
Conversation conv = conversationsClient.getConversation(request);
} catch (StatusRuntimeException e) {
switch (e.getStatus().getCode()) {
case NOT_FOUND:
// Conversation not found
break;
case UNAUTHENTICATED:
// Missing or invalid credentials
break;
case PERMISSION_DENIED:
// Access denied
break;
default:
throw e;
}
}
TLS Configuration
For production, enable TLS:
quarkus.grpc.clients.conversations.tls.enabled=true
quarkus.grpc.clients.conversations.tls.trust-certificate-pem.certs=ca.pem
When to Use gRPC
Choose gRPC over REST when you need:
- Streaming - Real-time response recording and resumption with
ResponseRecorderService - High throughput - Binary protocol is more efficient
- Strong typing - Generated stubs catch errors at compile time
- Bi-directional communication - Full duplex streaming for response recording
Next Steps
- REST Client - For simpler use cases
- Conversation Forking - Branch conversations to explore alternative paths
- Response Recording and Resumption - Streaming responses with resume and cancel support