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:

ServiceDescription
ConversationsServiceCRUD operations for conversations
EntriesServiceList, append, and sync entries
ConversationMembershipsServiceSharing and membership management
OwnershipTransfersServiceOwnership transfer operations
SearchServiceSemantic search and indexing
ResponseRecorderServiceStreaming response recording and resumption
SystemServiceHealth 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