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=trueis 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
- 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