Compare commits
1
Commits
0.5.0
..
5b52677c28
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
5b52677c28
|
@@ -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
@@ -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);
|
||||||
|
|
||||||
|
|||||||
+7
-23
@@ -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)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
+11
-26
@@ -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) {
|
||||||
|
|||||||
Reference in New Issue
Block a user