Spring gRPC Client

The Spring Boot starter provides a gRPC client for high-performance communication with Memory Service. gRPC is particularly useful for streaming scenarios with response recording and resumption.

When to Use gRPC

Choose gRPC over REST when you need:

  • Response Recording and Resumption - Required for resuming interrupted streaming responses
  • High throughput - Binary protocol is more efficient than JSON
  • Real-time streaming - Server-side streaming for message updates
  • Strong typing - Generated stubs catch errors at compile time

Setup

The gRPC client is included in the starter. gRPC is auto-configured when:

  • memory-service.grpc.enabled=true is set, OR
  • gRPC target is auto-derived from the REST client base URL
<dependency>
  <groupId>io.github.chirino.memory-service</groupId>
  <artifactId>memory-service-spring-boot-starter</artifactId>
  <version>999-SNAPSHOT</version>
</dependency>

Configuration

Configure gRPC in application.properties:

# Enable gRPC explicitly
memory-service.grpc.enabled=true
memory-service.grpc.target=localhost:6565
memory-service.grpc.plaintext=true

# Or let it auto-derive from REST client baseUrl
memory-service.client.url=http://localhost:8082
# gRPC will be auto-configured to use the same host/port

When memory-service.grpc.enabled is not explicitly set to true, the starter will attempt to auto-derive gRPC settings from the REST client’s base-url. This allows you to configure both REST and gRPC with a single URL.

Calling as an asserted user

Attach lowercase x-user-id metadata together with the service credential, and derive a stub for each request so user identity cannot leak between concurrent calls:

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 =
    stubs.conversationsService().withInterceptors(
        MetadataUtils.newAttachHeadersInterceptor(requestMetadata));

Conversation conversation = requestStub.createConversation(request);

The server honors the asserted user only for clients listed in MEMORY_SERVICE_TRUSTED_USER_ID_CLIENTS and only on normal user APIs. Admin and system methods ignore it. Request protobuf messages intentionally contain no identity wrapper or RequestActor field.

Injecting Stubs

The starter auto-configures a ManagedChannel and MemoryServiceGrpcClients.MemoryServiceStubs bean. Inject the stubs:

import io.github.chirino.memoryservice.grpc.MemoryServiceGrpcClients.MemoryServiceStubs;
import io.grpc.ManagedChannel;
import org.springframework.stereotype.Service;

@Service
public class MyService {
    private final MemoryServiceStubs stubs;

    public MyService(MemoryServiceStubs stubs) {
        this.stubs = stubs;
    }
}

The MemoryServiceStubs provides access to all service stubs:

  • systemService() - System health and status (blocking)
  • conversationsService() - Conversation operations (blocking)
  • membershipsService() - Conversation membership operations (blocking)
  • messagesService() - Message operations (blocking)
  • searchService() - Search operations (blocking)
  • responseRecorderService() - Response recording and resumption (async, streaming)

Example Operations

Get Conversation (Blocking)

import io.github.chirino.memory.grpc.v1.ConversationsServiceGrpc;
import io.github.chirino.memory.grpc.v1.GetConversationRequest;

ConversationsServiceGrpc.ConversationsServiceBlockingStub conversations =
    stubs.conversationsService();

GetConversationRequest request = GetConversationRequest.newBuilder()
    .setId("my-conversation")
    .build();

Conversation conversation = conversations.getConversation(request);
System.out.println(conversation.getTitle());

Create Conversation (Blocking)

import io.github.chirino.memory.grpc.v1.CreateConversationRequest;

CreateConversationRequest request = CreateConversationRequest.newBuilder()
    .setId("my-conversation")
    .putMetadata("topic", "support")
    .build();

Conversation conversation = conversations.createConversation(request);

List conversations (blocking)

import io.github.chirino.memory.grpc.v1.*;

ListConversationsResponse response = conversations.listConversations(
    ListConversationsRequest.newBuilder()
        .setPage(PageRequest.newBuilder().setPageSize(20).build())
        .build());
response.getConversationsList().forEach(conv -> {
    System.out.println(conv.getId());
});

Filtering by metadata

The gRPC API accepts typed ConversationMetadataFilter predicates. Because old servers silently ignore metadata_filters, check capabilities before sending them.

stubs.systemService() returns a SystemServiceGrpc.SystemServiceBlockingStub directly:

import com.google.protobuf.Empty;
import io.github.chirino.memory.grpc.v1.*;
import java.util.List;

CapabilitiesResponse caps = stubs.systemService().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 = conversations.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.

Streaming Messages

Server Streaming - Get All Messages

import io.github.chirino.memory.grpc.v1.MessagesServiceGrpc;
import io.github.chirino.memory.grpc.v1.GetMessagesRequest;
import io.grpc.stub.StreamObserver;

MessagesServiceGrpc.MessagesServiceStub messages =
    MessagesServiceGrpc.newStub(stubs.conversationsService().getChannel());

GetMessagesRequest request = GetMessagesRequest.newBuilder()
    .setConversationId("my-conversation")
    .build();

messages.getMessages(request, new StreamObserver<Message>() {
    @Override
    public void onNext(Message message) {
        System.out.println("Received: " + message.getContent());
    }

    @Override
    public void onError(Throwable t) {
        t.printStackTrace();
    }

    @Override
    public void onCompleted() {
        System.out.println("Stream completed");
    }
});

Real-time Message Updates

Subscribe to new messages as they arrive:

import io.github.chirino.memory.grpc.v1.StreamMessagesRequest;

StreamMessagesRequest request = StreamMessagesRequest.newBuilder()
    .setConversationId("my-conversation")
    .setFromSequence(lastKnownSequence)
    .build();

messages.streamMessages(request, new StreamObserver<Message>() {
    @Override
    public void onNext(Message message) {
        // Handle new message in real-time
        updateUI(message);
    }

    @Override
    public void onError(Throwable t) {
        // Handle reconnection
        reconnect();
    }

    @Override
    public void onCompleted() {
        // Stream ended
    }
});

Response Recording and Resumption

The responseRecorderService() provides async streaming for response recording and resumption. This is used internally by the ResponseRecordingManager bean, but you can also use it directly:

import io.github.chirino.memory.grpc.v1.ResponseRecorderServiceGrpc;
import io.github.chirino.memory.grpc.v1.ReplayRequest;
import reactor.core.publisher.Flux;

ResponseRecorderServiceGrpc.ResponseRecorderServiceStub recorder =
    stubs.responseRecorderService();

ReplayRequest request = ReplayRequest.newBuilder()
    .setConversationId(conversationId)
    .build();

// Convert gRPC stream to Reactor Flux
Flux<String> responseFlux = Flux.create(sink -> {
    recorder.replay(request, new StreamObserver<ReplayResponse>() {
        @Override
        public void onNext(ReplayResponse response) {
            sink.next(response.getContent());
        }

        @Override
        public void onError(Throwable t) {
            sink.error(t);
        }

        @Override
        public void onCompleted() {
            sink.complete();
        }
    });
});

Error Handling

import io.grpc.StatusRuntimeException;

try {
    Conversation conv = conversations.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;
    }
}

Authentication

gRPC uses the same authentication as REST. Configure API keys or bearer tokens via headers:

# API key is automatically added to gRPC requests
memory-service.client.api-key=agent-api-key-1

# Or use bearer token
memory-service.client.bearer-token=${BEARER_TOKEN:}

For OAuth2, the bearer token is obtained from the OAuth2 context and passed via gRPC metadata.

TLS Configuration

For production, enable TLS:

memory-service.grpc.plaintext=false
memory-service.grpc.target=memory-service.example.com:443

Configuration Properties

# Enable gRPC (or auto-derive from REST baseUrl)
memory-service.grpc.enabled=true

# gRPC target (host:port)
memory-service.grpc.target=localhost:6565

# Use plaintext (false for TLS)
memory-service.grpc.plaintext=true

# Keep-alive settings
memory-service.grpc.keep-alive-time=30s
memory-service.grpc.keep-alive-timeout=5s

# Custom headers
memory-service.grpc.headers.X-Custom-Header=value

Next Steps