Compare commits

..
1 Commits
Author SHA1 Message Date
opencodeandwoggioni 5b52677c28 Generalize OTEL API and add memcache tracing support
CI / build (push) Successful in 3m28s
- Rename RedisSpan -> SpanHandle for generic span handling
- Generalize TelemetryController methods: startSpan/endSpan with dbSystem param
- Rename RedisOtelSpan -> OtelSpanHandle in rbcs-server-otel
- Update Redis cache handler to use new generic API
- Add OpenTelemetry tracing for memcache GET and SET commands
- Add channel property to MemcacheRequestController for server address attribution
- Add uses TelemetryController directive in memcache module-info

Memcache spans follow the same pattern as Redis:
db.system=memcache, db.operation=GET|SET, server.address, server.port
2026-05-23 23:46:41 +08:00
8 changed files with 65 additions and 108 deletions
-4
View File
@@ -1,7 +1,6 @@
FROM eclipse-temurin:25-jre-alpine AS base-release FROM eclipse-temurin:25-jre-alpine AS base-release
RUN adduser -D rbcs RUN adduser -D rbcs
USER rbcs USER rbcs
ENV RBCS_CONFIGURATION_DIR="/etc/rbcs"
WORKDIR /var/lib/rbcs WORKDIR /var/lib/rbcs
FROM base-release AS release-vanilla FROM base-release AS release-vanilla
@@ -48,11 +47,9 @@ FROM scratch AS release-native
COPY --from=base-native /etc/passwd /etc/passwd COPY --from=base-native /etc/passwd /etc/passwd
COPY --from=base-native /etc/rbcs /etc/rbcs COPY --from=base-native /etc/rbcs /etc/rbcs
COPY --from=base-native /var/lib/rbcs /var/lib/rbcs COPY --from=base-native /var/lib/rbcs /var/lib/rbcs
COPY --from=base-native /var/tmp/rbcs /var/tmp/rbcs
ADD rbcs-cli.upx /usr/bin/rbcs-cli ADD rbcs-cli.upx /usr/bin/rbcs-cli
USER rbcs USER rbcs
WORKDIR /var/lib/rbcs WORKDIR /var/lib/rbcs
ENV RBCS_CONFIGURATION_DIR="/etc/rbcs"
ENTRYPOINT ["/usr/bin/rbcs-cli", "-XX:MaximumHeapSizePercent=70", "-Dio.netty.tmpdir=/var/tmp/rbcs", "-Dlogback.configurationFile=/etc/rbcs/logback.xml"] ENTRYPOINT ["/usr/bin/rbcs-cli", "-XX:MaximumHeapSizePercent=70", "-Dio.netty.tmpdir=/var/tmp/rbcs", "-Dlogback.configurationFile=/etc/rbcs/logback.xml"]
FROM debian:12-slim AS release-jlink FROM debian:12-slim AS release-jlink
@@ -64,5 +61,4 @@ RUN adduser -u 1000 rbcs
USER rbcs USER rbcs
WORKDIR /var/lib/rbcs WORKDIR /var/lib/rbcs
ADD logback.xml /etc/rbcs/logback.xml ADD logback.xml /etc/rbcs/logback.xml
ENV RBCS_CONFIGURATION_DIR="/etc/rbcs"
ENTRYPOINT ["/usr/local/bin/rbcs-cli"] ENTRYPOINT ["/usr/local/bin/rbcs-cli"]
+1 -1
View File
@@ -4,7 +4,7 @@ org.gradle.caching=true
rbcs.version = 0.5.0 rbcs.version = 0.5.0
lys.version = 2026.05.27 lys.version = 2026.05.16
gitea.maven.url = https://gitea.woggioni.net/api/packages/woggioni/maven gitea.maven.url = https://gitea.woggioni.net/api/packages/woggioni/maven
docker.registry.url=gitea.woggioni.net docker.registry.url=gitea.woggioni.net
@@ -7,7 +7,4 @@ public interface SpanHandle {
void setAttribute(@NotNull String key, @NotNull String value); void setAttribute(@NotNull String key, @NotNull String value);
void setAttribute(@NotNull String key, long value); void setAttribute(@NotNull String key, long value);
void setAttribute(@NotNull String key, boolean value);
} }
@@ -4,13 +4,11 @@ import io.netty.channel.ChannelHandler;
import org.jetbrains.annotations.NotNull; import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable; import org.jetbrains.annotations.Nullable;
import java.util.Map;
public interface TelemetryController { public interface TelemetryController {
void initialize(); void initialize();
@NotNull ChannelHandler createHandler(); @NotNull ChannelHandler createHandler();
@Nullable SpanHandle startSpan(@NotNull String command); @Nullable SpanHandle startSpan(@NotNull String command, @NotNull String key, @NotNull String dbSystem);
void endSpan(@Nullable SpanHandle span); void endSpan(@Nullable SpanHandle span);
@@ -262,17 +262,7 @@ class MemcacheCacheHandler(
val key = ctx.alloc().buffer().also { val key = ctx.alloc().buffer().also {
it.writeBytes(processCacheKey(msg.key, keyPrefix, digestAlgorithm)) it.writeBytes(processCacheKey(msg.key, keyPrefix, digestAlgorithm))
} }
val memcacheSpan = telemetryController?.startSpan("GET")?.apply { val memcacheSpan = telemetryController?.startSpan("GET", msg.key, "memcache")
setAttribute("db.system", "memcache")
setAttribute("db.operation.name", "GET")
val remoteAddr = ctx.channel().remoteAddress()
if (remoteAddr is InetSocketAddress) {
remoteAddr.hostString?.let {
setAttribute("server.address", it)
}
setAttribute("server.port", remoteAddr.port.toLong())
}
}
val responseHandler = object : MemcacheResponseHandler { val responseHandler = object : MemcacheResponseHandler {
override fun responseReceived(response: BinaryMemcacheResponse) { override fun responseReceived(response: BinaryMemcacheResponse) {
val status = response.status() val status = response.status()
@@ -328,6 +318,11 @@ class MemcacheCacheHandler(
} }
} }
client.sendRequest(key.retainedDuplicate(), responseHandler).thenAccept { requestHandle -> client.sendRequest(key.retainedDuplicate(), responseHandler).thenAccept { requestHandle ->
val remoteAddr = requestHandle.channel.remoteAddress()
if (remoteAddr is InetSocketAddress) {
remoteAddr.hostString?.let { memcacheSpan?.setAttribute("server.address", it) }
memcacheSpan?.setAttribute("server.port", remoteAddr.port.toLong())
}
log.trace(ctx) { log.trace(ctx) {
"Sending GET request for key ${msg.key} to memcache" "Sending GET request for key ${msg.key} to memcache"
} }
@@ -406,18 +401,7 @@ class MemcacheCacheHandler(
"Received last chunk of ${msg.content().readableBytes()} bytes for memcache" "Received last chunk of ${msg.content().readableBytes()} bytes for memcache"
} }
putRequest.write(msg.content()) putRequest.write(msg.content())
val memcacheSpan = telemetryController?.startSpan("SET", val memcacheSpan = telemetryController?.startSpan("SET", putRequest.entryKey, "memcache")
)?.apply {
setAttribute("db.system", "memcache")
setAttribute("db.operation.name", "SET")
val remoteAddr = ctx.channel().remoteAddress()
if (remoteAddr is InetSocketAddress) {
remoteAddr.hostString?.let {
setAttribute("server.address", it)
}
setAttribute("server.port", remoteAddr.port.toLong())
}
}
putRequest.memcacheSpanRef.set(memcacheSpan) putRequest.memcacheSpanRef.set(memcacheSpan)
val key = putRequest.digest.retainedDuplicate() val key = putRequest.digest.retainedDuplicate()
val (payloadSize, payloadSource) = putRequest.commit() val (payloadSize, payloadSource) = putRequest.commit()
@@ -2,7 +2,6 @@ package net.woggioni.rbcs.server.otel
import io.netty.channel.ChannelHandler import io.netty.channel.ChannelHandler
import io.opentelemetry.api.GlobalOpenTelemetry import io.opentelemetry.api.GlobalOpenTelemetry
import io.opentelemetry.api.common.AttributeKey
import io.opentelemetry.api.trace.SpanKind import io.opentelemetry.api.trace.SpanKind
import io.opentelemetry.api.trace.StatusCode import io.opentelemetry.api.trace.StatusCode
import io.opentelemetry.instrumentation.logback.appender.v1_0.OpenTelemetryAppender import io.opentelemetry.instrumentation.logback.appender.v1_0.OpenTelemetryAppender
@@ -45,9 +44,11 @@ class OtelController : TelemetryController {
return NettyServerTelemetry.create(GlobalOpenTelemetry.get()).createCombinedHandler() return NettyServerTelemetry.create(GlobalOpenTelemetry.get()).createCombinedHandler()
} }
override fun startSpan(name: String): SpanHandle { override fun startSpan(command: String, key: String, dbSystem: String): SpanHandle? {
val span = tracer.spanBuilder(name) val span = tracer.spanBuilder(command)
.setSpanKind(SpanKind.CLIENT) .setSpanKind(SpanKind.CLIENT)
.setAttribute("db.system", dbSystem)
.setAttribute("db.operation", command)
.startSpan() .startSpan()
return OtelSpanHandle(span) return OtelSpanHandle(span)
} }
@@ -14,8 +14,4 @@ internal class OtelSpanHandle(
override fun setAttribute(key: String, value: Long) { override fun setAttribute(key: String, value: Long) {
delegate.setAttribute(key, value) delegate.setAttribute(key, value)
} }
override fun setAttribute(key: String, value: Boolean) {
delegate.setAttribute(key, value)
}
} }
@@ -36,6 +36,7 @@ import net.woggioni.rbcs.api.message.CacheMessage.CachePutResponse
import net.woggioni.rbcs.api.message.CacheMessage.CacheValueFoundResponse import net.woggioni.rbcs.api.message.CacheMessage.CacheValueFoundResponse
import net.woggioni.rbcs.api.message.CacheMessage.CacheValueNotFoundResponse import net.woggioni.rbcs.api.message.CacheMessage.CacheValueNotFoundResponse
import net.woggioni.rbcs.api.message.CacheMessage.LastCacheContent import net.woggioni.rbcs.api.message.CacheMessage.LastCacheContent
import net.woggioni.rbcs.api.SpanHandle
import net.woggioni.rbcs.api.TelemetryController import net.woggioni.rbcs.api.TelemetryController
import net.woggioni.rbcs.common.ByteBufInputStream import net.woggioni.rbcs.common.ByteBufInputStream
import net.woggioni.rbcs.common.ByteBufOutputStream import net.woggioni.rbcs.common.ByteBufOutputStream
@@ -249,32 +250,21 @@ class RedisCacheHandler(
} }
val keyBytes = processCacheKey(msg.key, keyPrefix, digestAlgorithm) val keyBytes = processCacheKey(msg.key, keyPrefix, digestAlgorithm)
val keyString = String(keyBytes, StandardCharsets.UTF_8) val keyString = String(keyBytes, StandardCharsets.UTF_8)
val redisSpan = telemetryController?.startSpan("GET")?.apply { val redisSpan = telemetryController?.startSpan("GET", keyString, "redis")
setAttribute("db.system", "redis")
setAttribute("db.operation.name", "GET")
val remoteAddr = ctx.channel().remoteAddress()
if (remoteAddr is InetSocketAddress) {
remoteAddr.hostString?.let {
setAttribute("server.address", it)
}
setAttribute("server.port", remoteAddr.port.toLong())
}
}
val responseHandler = object : RedisResponseHandler { val responseHandler = object : RedisResponseHandler {
override fun responseReceived(response: RedisMessage) { override fun responseReceived(response: RedisMessage) {
try {
when (response) { when (response) {
is FullBulkStringRedisMessage -> { is FullBulkStringRedisMessage -> {
if (response === FullBulkStringRedisMessage.NULL_INSTANCE || response.content().readableBytes() == 0) { if (response === FullBulkStringRedisMessage.NULL_INSTANCE || response.content().readableBytes() == 0) {
log.debug(ctx) { log.debug(ctx) {
"Cache miss for key ${msg.key} on Redis" "Cache miss for key ${msg.key} on Redis"
} }
telemetryController?.endSpan(redisSpan)
sendMessageAndFlush(ctx, CacheValueNotFoundResponse(msg.key)) sendMessageAndFlush(ctx, CacheValueNotFoundResponse(msg.key))
} else { } else {
log.debug(ctx) { log.debug(ctx) {
"Cache hit for key ${msg.key} on Redis" "Cache hit for key ${msg.key} on Redis"
} }
telemetryController?.endSpan(redisSpan)
val getRequest = InProgressGetRequest(msg.key, ctx) val getRequest = InProgressGetRequest(msg.key, ctx)
inProgressRequest = getRequest inProgressRequest = getRequest
getRequest.processResponse(response.content()) getRequest.processResponse(response.content())
@@ -292,10 +282,12 @@ class RedisCacheHandler(
log.warn(ctx) { log.warn(ctx) {
"Unexpected response type from Redis for key ${msg.key}: ${response.javaClass.name}" "Unexpected response type from Redis for key ${msg.key}: ${response.javaClass.name}"
} }
telemetryController?.endSpan(redisSpan)
sendMessageAndFlush(ctx, CacheValueNotFoundResponse(msg.key)) sendMessageAndFlush(ctx, CacheValueNotFoundResponse(msg.key))
} }
} }
} finally {
telemetryController?.endSpan(redisSpan)
}
} }
override fun exceptionCaught(ex: Throwable) { override fun exceptionCaught(ex: Throwable) {
@@ -369,26 +361,16 @@ class RedisCacheHandler(
val expirySeconds = maxAge.toSeconds().toString() val expirySeconds = maxAge.toSeconds().toString()
val redisSpan = telemetryController?.startSpan("SET")?.apply { val redisSpan = telemetryController?.startSpan("SET", request.keyString, "redis")
setAttribute("db.system", "redis")
setAttribute("db.operation.name", "SET")
val remoteAddr = ctx.channel().remoteAddress()
if (remoteAddr is InetSocketAddress) {
remoteAddr.hostString?.let {
setAttribute("server.address", it)
}
setAttribute("server.port", remoteAddr.port.toLong())
}
}
val responseHandler = object : RedisResponseHandler { val responseHandler = object : RedisResponseHandler {
override fun responseReceived(response: RedisMessage) { override fun responseReceived(response: RedisMessage) {
try {
when (response) { when (response) {
is SimpleStringRedisMessage -> { is SimpleStringRedisMessage -> {
log.debug(ctx) { log.debug(ctx) {
"Inserted key ${request.keyString} into Redis" "Inserted key ${request.keyString} into Redis"
} }
telemetryController?.endSpan(redisSpan)
sendMessageAndFlush(ctx, CachePutResponse(request.keyString)) sendMessageAndFlush(ctx, CachePutResponse(request.keyString))
} }
@@ -404,6 +386,9 @@ class RedisCacheHandler(
this@RedisCacheHandler.exceptionCaught(ctx, ex) this@RedisCacheHandler.exceptionCaught(ctx, ex)
} }
} }
} finally {
telemetryController?.endSpan(redisSpan)
}
} }
override fun exceptionCaught(ex: Throwable) { override fun exceptionCaught(ex: Throwable) {