diff --git a/infra/docker/.env.example b/infra/docker/.env.example index b802cd92..944a1790 100644 --- a/infra/docker/.env.example +++ b/infra/docker/.env.example @@ -20,11 +20,9 @@ COMPOSE_PATH_SEPARATOR=: # Infra dependencies — pinned for reproducibility. CLICKHOUSE_IMAGE_NAME=clickhouse/clickhouse-server:24.3.15.72-alpine -KAFKA_IMAGE_NAME=confluentinc/cp-kafka:7.7.0 POSTGRES_IMAGE_NAME=ankane/pgvector:v0.5.1 # OTEL collector image vars removed — HOL-21 dropped the collector container. # Backend hosts the OTLP receiver at /otel/v1/{logs,traces,metrics}. -ZOOKEEPER_IMAGE_NAME=confluentinc/cp-zookeeper:7.7.0 # ── ClickHouse (analytics) ────────────────────────────────────────── CLICKHOUSE_DATABASE=default @@ -39,9 +37,9 @@ PSQL_PORT=5432 PSQL_USER=postgres # ── Kafka ─────────────────────────────────────────────────────────── -KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092 -KAFKA_SERVERS=kafka:9092 -KAFKA_TOPIC=dev +# Kafka removed in HOL-23 — backend uses an in-process Channel-backed +# message bus. Re-introduce these vars only if swapping a Kafka-backed +# IMessageBus implementation back in for multi-node deployments. # ── Redis ─────────────────────────────────────────────────────────── # Removed in HOL-22 — backend uses IMemoryCache + Postgres instead. diff --git a/infra/docker/compose.hobby-dotnet.yml b/infra/docker/compose.hobby-dotnet.yml index 9f9b7986..e01e89dc 100644 --- a/infra/docker/compose.hobby-dotnet.yml +++ b/infra/docker/compose.hobby-dotnet.yml @@ -21,7 +21,6 @@ services: - ClickHouse__Password=${CLICKHOUSE_PASSWORD:-} - ClickHouse__Migrations__Path=/app/clickhouse-migrations - ClickHouse__Migrations__Disabled=${CLICKHOUSE_MIGRATIONS_DISABLED:-false} - - Kafka__BootstrapServers=${KAFKA_SERVERS:-kafka:9092} - Storage__Type=filesystem - Storage__FilesystemRoot=/highlight-data - Frontend__Uri=${REACT_APP_FRONTEND_URI:-http://localhost:3000} diff --git a/infra/docker/compose.yml b/infra/docker/compose.yml index 958542b2..9a823607 100644 --- a/infra/docker/compose.yml +++ b/infra/docker/compose.yml @@ -3,43 +3,10 @@ x-local-logging: &local-logging # HoldFast services for the development deployment. services: - zookeeper: - logging: *local-logging - image: ${ZOOKEEPER_IMAGE_NAME} - container_name: zookeeper - restart: on-failure - volumes: - - zoo-data:/var/lib/zookeeper/data - - zoo-log:/var/lib/zookeeper/log - environment: - ZOOKEEPER_CLIENT_PORT: 2181 - ZOOKEEPER_TICK_TIME: 2000 - - kafka: - logging: *local-logging - image: ${KAFKA_IMAGE_NAME} - container_name: kafka - volumes: - - kafka-data:/var/lib/kafka/data - ports: - - '0.0.0.0:9092:9092' - restart: on-failure - depends_on: - - zookeeper - environment: - KAFKA_ADVERTISED_LISTENERS: ${KAFKA_ADVERTISED_LISTENERS} - KAFKA_BROKER_ID: 1 - KAFKA_CONSUMER_MAX_PARTITION_FETCH_BYTES: 268435456 - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT - KAFKA_LOG_RETENTION_HOURS: 1 - KAFKA_LOG_SEGMENT_BYTES: 268435456 - KAFKA_MESSAGE_MAX_BYTES: 268435456 - KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 - KAFKA_PRODUCER_MAX_REQUEST_SIZE: 268435456 - KAFKA_REPLICA_FETCH_MAX_BYTES: 268435456 - KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 - KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 - KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' + # Kafka + Zookeeper removed in HOL-23 — backend uses an in-process + # Channel-backed message bus (HoldFast.Shared.Messaging.InProcessMessageBus). + # Producers and consumers live in the same .NET host (all-in-one runtime + # mode), so no broker is required at hobby/lean self-hosted scale. # Redis removed in HOL-22 — RedisService was registered in DI but no # code path actually injected it. Self-hosted scale doesn't need a @@ -88,6 +55,3 @@ volumes: postgres-data: clickhouse-data: clickhouse-logs: - kafka-data: - zoo-log: - zoo-data: diff --git a/infra/docker/env.sh b/infra/docker/env.sh index 25766697..bb71b97c 100644 --- a/infra/docker/env.sh +++ b/infra/docker/env.sh @@ -26,8 +26,6 @@ fi if [[ "$*" == *"--go-docker"* ]]; then export CLICKHOUSE_ADDRESS=clickhouse:9000 export IN_DOCKER_GO=true - export KAFKA_ADVERTISED_LISTENERS="PLAINTEXT://kafka:9092" - export KAFKA_SERVERS=kafka:9092 export ON_PREM=true # HOL-21: collector container removed; backend hosts the OTLP receiver. export OTLP_DOGFOOD_ENDPOINT=http://backend:8082/otel @@ -36,8 +34,6 @@ if [[ "$*" == *"--go-docker"* ]]; then echo "Using docker-internal infra." else export CLICKHOUSE_ADDRESS=localhost:9000 - export KAFKA_ADVERTISED_LISTENERS="PLAINTEXT://localhost:9092" - export KAFKA_SERVERS=localhost:9092 export OTLP_DOGFOOD_ENDPOINT=http://localhost:8082/otel export OTLP_ENDPOINT=http://localhost:8082/otel export PSQL_HOST=localhost diff --git a/infra/docker/start-infra.sh b/infra/docker/start-infra.sh index a437efe9..26a7a301 100644 --- a/infra/docker/start-infra.sh +++ b/infra/docker/start-infra.sh @@ -4,12 +4,9 @@ source env.sh # startup the infra -SERVICES="clickhouse kafka postgres redis zookeeper" -BUILD_ARGS="--build-arg OTEL_COLLECTOR_ALPINE_IMAGE_NAME=${OTEL_COLLECTOR_ALPINE_IMAGE_NAME} \ - --build-arg OTEL_COLLECTOR_IMAGE_NAME=${OTEL_COLLECTOR_IMAGE_NAME}" +SERVICES="clickhouse postgres" docker compose pull $SERVICES -docker compose build --pull $BUILD_ARGS $SERVICES docker compose up --detach --wait --remove-orphans $SERVICES if [[ "$*" != *"--go-docker"* ]]; then diff --git a/src/dotnet/src/HoldFast.Api/Program.cs b/src/dotnet/src/HoldFast.Api/Program.cs index e1fffcf1..4e77e747 100644 --- a/src/dotnet/src/HoldFast.Api/Program.cs +++ b/src/dotnet/src/HoldFast.Api/Program.cs @@ -10,6 +10,7 @@ using HoldFast.Shared.Notifications; using HoldFast.Shared.ErrorGrouping; using HoldFast.Shared.Kafka; +using HoldFast.Shared.Messaging; using HoldFast.Shared.Runtime; using HoldFast.Shared.SessionProcessing; using HoldFast.Storage; @@ -137,20 +138,15 @@ req.RequestUri is null || }); builder.Services.AddScoped(); -// ── Kafka ───────────────────────────────────────────────────────────── -builder.Services.Configure( - builder.Configuration.GetSection("Kafka")); -builder.Services.AddSingleton(); +// ── Message bus ────────────────────────────────────────────────────── +// HOL-23: in-process Channel-backed bus replaces Kafka for hobby/lean +// self-hosted deployments. Producer and consumers both run in the same +// .NET host (all-in-one runtime mode); see HoldFast.Shared.Messaging. +// To run the worker in a separate process from the API, swap in a +// Kafka-backed IMessageBus implementation that honors the same shape. +builder.Services.AddSingleton(); builder.Services.AddSingleton(); -// Topic bootstrap — pre-create required topics so consumers don't crash on -// subscribe (Confluent.Kafka auto-create only triggers on producer writes). -// Disable via Kafka__TopicBootstrap__Disabled=true when topics are managed -// externally (Strimzi KafkaTopic resources, Helm pre-jobs, etc). -builder.Services.Configure( - builder.Configuration.GetSection("Kafka:TopicBootstrap")); -builder.Services.AddHostedService(); - // ── ClickHouse ──────────────────────────────────────────────────────── builder.Services.Configure( builder.Configuration.GetSection("ClickHouse")); diff --git a/src/dotnet/src/HoldFast.GraphQL.Public/KafkaProducerAdapter.cs b/src/dotnet/src/HoldFast.GraphQL.Public/KafkaProducerAdapter.cs index 58fe218d..6dd569f6 100644 --- a/src/dotnet/src/HoldFast.GraphQL.Public/KafkaProducerAdapter.cs +++ b/src/dotnet/src/HoldFast.GraphQL.Public/KafkaProducerAdapter.cs @@ -1,17 +1,21 @@ using System.Text.Json; using HoldFast.GraphQL.Public.InputTypes; using HoldFast.Shared.Kafka; +using HoldFast.Shared.Messaging; namespace HoldFast.GraphQL.Public; /// -/// Implements IKafkaProducer by forwarding to the shared KafkaProducerService. +/// Implements IKafkaProducer by forwarding to the shared . +/// (HOL-23: the bus is now an in-process Channel-backed implementation by default, +/// not Confluent.Kafka. The class name is kept for now to keep the diff narrow — +/// rename pending in a follow-up.) /// public class KafkaProducerAdapter : IKafkaProducer { - private readonly KafkaProducerService _producer; + private readonly IMessageBus _producer; - public KafkaProducerAdapter(KafkaProducerService producer) + public KafkaProducerAdapter(IMessageBus producer) { _producer = producer; } @@ -19,7 +23,7 @@ public KafkaProducerAdapter(KafkaProducerService producer) public Task ProduceSessionEventsAsync(string sessionSecureId, long payloadId, string data, CancellationToken ct) { var message = new { SessionSecureId = sessionSecureId, PayloadId = payloadId, Data = data }; - return _producer.ProduceAsync(KafkaTopics.SessionEvents, sessionSecureId, message, ct); + return _producer.PublishAsync(KafkaTopics.SessionEvents, sessionSecureId, message, ct); } public async Task ProducePushPayloadAsync(string sessionSecureId, long payloadId, string events, @@ -37,7 +41,7 @@ public async Task ProducePushPayloadAsync(string sessionSecureId, long payloadId PayloadId = payloadId, Data = events, }; - await _producer.ProduceAsync(KafkaTopics.SessionEvents, sessionSecureId, sessionMsg, ct); + await _producer.PublishAsync(KafkaTopics.SessionEvents, sessionSecureId, sessionMsg, ct); // Each frontend error → FrontendErrors topic so FrontendErrorsConsumer can // group them into ClickHouse error_objects. Errors are dropped on the @@ -58,7 +62,7 @@ public async Task ProducePushPayloadAsync(string sessionSecureId, long payloadId Timestamp = err.Timestamp, Payload = err.Payload, }; - await _producer.ProduceAsync(KafkaTopics.FrontendErrors, sessionSecureId, frontendMsg, ct); + await _producer.PublishAsync(KafkaTopics.FrontendErrors, sessionSecureId, frontendMsg, ct); } // _ = messages; _ = resources; _ = webSocketEvents; _ = isBeacon; @@ -71,21 +75,21 @@ public async Task ProducePushPayloadAsync(string sessionSecureId, long payloadId public Task ProduceBackendErrorAsync(string? projectId, BackendErrorObjectInput error, CancellationToken ct) { - return _producer.ProduceAsync(KafkaTopics.BackendErrors, projectId ?? "unknown", error, ct); + return _producer.PublishAsync(KafkaTopics.BackendErrors, projectId ?? "unknown", error, ct); } public Task ProduceMetricAsync(MetricInput metric, CancellationToken ct) { - return _producer.ProduceAsync(KafkaTopics.Metrics, metric.SessionSecureId, metric, ct); + return _producer.PublishAsync(KafkaTopics.Metrics, metric.SessionSecureId, metric, ct); } public Task ProduceLogAsync(LogInput log, CancellationToken ct) { - return _producer.ProduceAsync(KafkaTopics.Logs, log.TraceId, log, ct); + return _producer.PublishAsync(KafkaTopics.Logs, log.TraceId, log, ct); } public Task ProduceTraceAsync(TraceInput trace, CancellationToken ct) { - return _producer.ProduceAsync(KafkaTopics.Traces, trace.TraceId, trace, ct); + return _producer.PublishAsync(KafkaTopics.Traces, trace.TraceId, trace, ct); } } diff --git a/src/dotnet/src/HoldFast.Shared/HoldFast.Shared.csproj b/src/dotnet/src/HoldFast.Shared/HoldFast.Shared.csproj index 4e6509e0..9e602dfa 100644 --- a/src/dotnet/src/HoldFast.Shared/HoldFast.Shared.csproj +++ b/src/dotnet/src/HoldFast.Shared/HoldFast.Shared.csproj @@ -6,7 +6,6 @@ - diff --git a/src/dotnet/src/HoldFast.Shared/Kafka/KafkaConsumerService.cs b/src/dotnet/src/HoldFast.Shared/Kafka/KafkaConsumerService.cs deleted file mode 100644 index 53bca21f..00000000 --- a/src/dotnet/src/HoldFast.Shared/Kafka/KafkaConsumerService.cs +++ /dev/null @@ -1,88 +0,0 @@ -using System.Text.Json; -using Confluent.Kafka; -using Microsoft.Extensions.Logging; -using Microsoft.Extensions.Options; - -namespace HoldFast.Shared.Kafka; - -/// -/// Base Kafka consumer that handles deserialization and error handling. -/// Worker BackgroundServices inherit from this. -/// -public abstract class KafkaConsumerService where T : class -{ - private readonly IConsumer _consumer; - private readonly ILogger _logger; - private readonly string _topic; - - protected KafkaConsumerService( - IOptions options, - string topic, - string groupId, - ILogger logger) - { - _topic = topic; - _logger = logger; - - var config = new ConsumerConfig - { - BootstrapServers = options.Value.BootstrapServers, - GroupId = groupId, - AutoOffsetReset = AutoOffsetReset.Earliest, - EnableAutoCommit = false, - MaxPollIntervalMs = 300000, - }; - _consumer = new ConsumerBuilder(config).Build(); - } - - /// - /// Process a single consumed message. Implementations define the business logic. - /// - protected abstract Task ProcessAsync(string key, T value, CancellationToken ct); - - /// - /// Run the consume loop. Call from BackgroundService.ExecuteAsync. - /// - public async Task ConsumeLoopAsync(CancellationToken ct) - { - _consumer.Subscribe(_topic); - _logger.LogInformation("Kafka consumer started for topic {Topic}", _topic); - - try - { - while (!ct.IsCancellationRequested) - { - var result = _consumer.Consume(ct); - if (result?.Message?.Value == null) continue; - - try - { - var value = JsonSerializer.Deserialize(result.Message.Value); - if (value != null) - { - await ProcessAsync(result.Message.Key, value, ct); - } - _consumer.Commit(result); - } - catch (JsonException ex) - { - _logger.LogError(ex, "Failed to deserialize message from {Topic}", _topic); - _consumer.Commit(result); // Skip bad messages - } - catch (Exception ex) - { - _logger.LogError(ex, "Error processing message from {Topic}", _topic); - // Don't commit — message will be redelivered - } - } - } - catch (OperationCanceledException) - { - // Normal shutdown - } - finally - { - _consumer.Close(); - } - } -} diff --git a/src/dotnet/src/HoldFast.Shared/Kafka/KafkaProducerService.cs b/src/dotnet/src/HoldFast.Shared/Kafka/KafkaProducerService.cs deleted file mode 100644 index bb61c74a..00000000 --- a/src/dotnet/src/HoldFast.Shared/Kafka/KafkaProducerService.cs +++ /dev/null @@ -1,59 +0,0 @@ -using System.Text.Json; -using Confluent.Kafka; -using Microsoft.Extensions.Logging; -using Microsoft.Extensions.Options; - -namespace HoldFast.Shared.Kafka; - -/// -/// Configuration for Kafka connections. BootstrapServers is comma-separated broker addresses. -/// -public class KafkaOptions -{ - public string BootstrapServers { get; set; } = "localhost:9092"; -} - -/// -/// Confluent.Kafka producer wrapper. Produces JSON-serialized messages. -/// -public class KafkaProducerService : IDisposable -{ - private readonly IProducer _producer; - private readonly ILogger _logger; - - public KafkaProducerService(IOptions options, ILogger logger) - { - _logger = logger; - var config = new ProducerConfig - { - BootstrapServers = options.Value.BootstrapServers, - Acks = Acks.All, - EnableIdempotence = true, - LingerMs = 5, - BatchNumMessages = 1000, - }; - _producer = new ProducerBuilder(config).Build(); - } - - public async Task ProduceAsync(string topic, string key, T value, CancellationToken ct) - { - var json = JsonSerializer.Serialize(value); - var message = new Message { Key = key, Value = json }; - - try - { - await _producer.ProduceAsync(topic, message, ct); - } - catch (ProduceException ex) - { - _logger.LogError(ex, "Failed to produce message to {Topic} with key {Key}", topic, key); - throw; - } - } - - public void Dispose() - { - _producer.Flush(TimeSpan.FromSeconds(5)); - _producer.Dispose(); - } -} diff --git a/src/dotnet/src/HoldFast.Shared/Kafka/KafkaTopicBootstrapService.cs b/src/dotnet/src/HoldFast.Shared/Kafka/KafkaTopicBootstrapService.cs deleted file mode 100644 index dbbfcb89..00000000 --- a/src/dotnet/src/HoldFast.Shared/Kafka/KafkaTopicBootstrapService.cs +++ /dev/null @@ -1,149 +0,0 @@ -using Confluent.Kafka; -using Confluent.Kafka.Admin; -using Microsoft.Extensions.Hosting; -using Microsoft.Extensions.Logging; -using Microsoft.Extensions.Options; - -namespace HoldFast.Shared.Kafka; - -/// -/// Configuration for the Kafka topic bootstrap runner. -/// -public class KafkaTopicBootstrapOptions -{ - /// - /// Number of partitions for newly-created topics. Hobby/dev defaults to 1; - /// production deployments should set this higher (typically 3-12 depending - /// on consumer parallelism). - /// - public int Partitions { get; set; } = 1; - - /// - /// Replication factor for newly-created topics. Hobby/dev defaults to 1 - /// since the cluster has a single broker; production should be 3. - /// - public short ReplicationFactor { get; set; } = 1; - - /// - /// Timeout for the create-topics admin call. - /// - public TimeSpan CreateTimeout { get; set; } = TimeSpan.FromSeconds(30); - - /// - /// Skip topic creation. Use when topics are managed externally (Helm - /// chart pre-job, Strimzi KafkaTopic resources, ops automation). - /// - public bool Disabled { get; set; } -} - -/// -/// Pre-creates the Kafka topics consumers will subscribe to, on backend startup. -/// -/// Why: Confluent.Kafka's auto-create only triggers on producer writes, not on -/// consumer subscribe. With the .NET backend's default `BackgroundServiceException -/// Behavior = StopHost`, a single consumer failing to subscribe to a missing -/// topic kills the entire process. Net effect on a fresh stack: backend starts, -/// SessionEventsConsumer fails to subscribe to "session-events", host stops, -/// container restarts, repeat forever. Pre-creating topics breaks the loop. -/// -/// See HOL-12. -/// -public class KafkaTopicBootstrapService : IHostedService -{ - /// - /// Topics required by the backend's worker hosted services. Order doesn't - /// matter; the admin call is batched. - /// - private static readonly string[] RequiredTopics = - [ - KafkaTopics.SessionEvents, - KafkaTopics.BackendErrors, - KafkaTopics.FrontendErrors, - KafkaTopics.Metrics, - KafkaTopics.Logs, - KafkaTopics.Traces, - ]; - - private readonly KafkaOptions _kafka; - private readonly KafkaTopicBootstrapOptions _options; - private readonly ILogger _logger; - - public KafkaTopicBootstrapService( - IOptions kafka, - IOptions options, - ILogger logger) - { - _kafka = kafka.Value; - _options = options.Value; - _logger = logger; - } - - public async Task StartAsync(CancellationToken cancellationToken) - { - if (_options.Disabled) - { - _logger.LogInformation("Kafka topic bootstrap: disabled by configuration, skipping"); - return; - } - - var config = new AdminClientConfig { BootstrapServers = _kafka.BootstrapServers }; - using var admin = new AdminClientBuilder(config).Build(); - - // Find which topics actually exist so we only create the missing ones — - // CreateTopicsAsync returns one Error per topic, but plumbing it through - // is awkward; querying first keeps the log output clean. - var meta = admin.GetMetadata(TimeSpan.FromSeconds(15)); - var existing = meta.Topics - .Where(t => t.Error.Code == ErrorCode.NoError) - .Select(t => t.Topic) - .ToHashSet(); - - var missing = RequiredTopics.Where(t => !existing.Contains(t)).ToList(); - if (missing.Count == 0) - { - _logger.LogInformation("Kafka topic bootstrap: all {Count} topics already exist", - RequiredTopics.Length); - return; - } - - var specs = missing.Select(name => new TopicSpecification - { - Name = name, - NumPartitions = _options.Partitions, - ReplicationFactor = _options.ReplicationFactor, - }).ToList(); - - try - { - await admin.CreateTopicsAsync(specs, new CreateTopicsOptions - { - OperationTimeout = _options.CreateTimeout, - }); - _logger.LogInformation( - "Kafka topic bootstrap: created {Created} of {Required} topics ({Names})", - missing.Count, RequiredTopics.Length, string.Join(", ", missing)); - } - catch (CreateTopicsException ex) - { - // Distinguish "raced with another instance / hobby restart" from real failures - var failed = ex.Results - .Where(r => r.Error.Code != ErrorCode.NoError && - r.Error.Code != ErrorCode.TopicAlreadyExists) - .ToList(); - if (failed.Count > 0) - { - foreach (var f in failed) - _logger.LogError( - "Kafka topic bootstrap: failed to create {Topic} — {Error}", - f.Topic, f.Error.Reason); - throw; - } - _logger.LogInformation( - "Kafka topic bootstrap: {Existed} of {Created} topics already existed (race-safe)", - ex.Results.Count(r => r.Error.Code == ErrorCode.TopicAlreadyExists), - missing.Count); - } - } - - public Task StopAsync(CancellationToken cancellationToken) => Task.CompletedTask; -} diff --git a/src/dotnet/src/HoldFast.Shared/Messaging/IMessageBus.cs b/src/dotnet/src/HoldFast.Shared/Messaging/IMessageBus.cs new file mode 100644 index 00000000..b29980e6 --- /dev/null +++ b/src/dotnet/src/HoldFast.Shared/Messaging/IMessageBus.cs @@ -0,0 +1,31 @@ +namespace HoldFast.Shared.Messaging; + +/// +/// Pub/sub message bus the worker hosted services use to decouple ingest from +/// processing. Replaces Kafka in the hobby/lean architecture (HOL-23) — the +/// in-process Channel-backed implementation handles single-node self-hosted +/// scale without a broker, JVM, or zookeeper container. +/// +/// The interface intentionally mirrors the small slice of Kafka semantics we +/// were actually using: per-topic JSON messages with a string key, no +/// partition control, no consumer-group rebalancing. If you outgrow the +/// in-process implementation, swap in a Kafka-backed implementation that +/// honors the same shape. +/// +public interface IMessageBus +{ + /// + /// Publish a value to a topic. Value is JSON-serialized. + /// In the in-process implementation this is non-blocking (channel write). + /// + Task PublishAsync(string topic, string key, T value, CancellationToken ct); + + /// + /// Subscribe to a topic and receive messages as (key, json-body) pairs. + /// Multiple consumers on the same topic will fan in (each message goes + /// to exactly one of them) — matches Kafka consumer-group semantics with + /// a single-partition topic. + /// + IAsyncEnumerable<(string Key, string Value)> SubscribeAsync( + string topic, CancellationToken ct); +} diff --git a/src/dotnet/src/HoldFast.Shared/Messaging/InProcessMessageBus.cs b/src/dotnet/src/HoldFast.Shared/Messaging/InProcessMessageBus.cs new file mode 100644 index 00000000..78d6b0ac --- /dev/null +++ b/src/dotnet/src/HoldFast.Shared/Messaging/InProcessMessageBus.cs @@ -0,0 +1,69 @@ +using System.Collections.Concurrent; +using System.Runtime.CompilerServices; +using System.Text.Json; +using System.Threading.Channels; +using Microsoft.Extensions.Logging; + +namespace HoldFast.Shared.Messaging; + +/// +/// Single-process message bus backed by . Replaces +/// Kafka for hobby/lean self-hosted deployments (HOL-23). +/// +/// One unbounded channel per topic, lazily created on first publish or +/// subscribe. Producers write (key, json) pairs; consumers read them in +/// FIFO order. Multiple subscribers on the same topic compete for messages +/// (fan-in), so the same SessionEventsConsumer-as-singleton + worker pattern +/// the Kafka version used keeps working unchanged. +/// +/// Tradeoffs vs. Kafka: +/// - No durability — messages live in memory; a backend restart drops the +/// in-flight queue. Acceptable at hobby scale; the API responses to +/// pushPayload are still 200, so SDK-side retries cover transient loss. +/// - No replay — no offsets, no consumer groups, no rebalancing. +/// - Single-node only — producer and consumer must be in the same process. +/// The .NET backend's "all-in-one" runtime mode satisfies this. +/// +/// If a deployment outgrows this, swap an alternative IMessageBus +/// implementation back in (Kafka, Redis Streams, etc). +/// +public class InProcessMessageBus : IMessageBus +{ + private readonly ConcurrentDictionary> _channels = new(); + private readonly ILogger _logger; + + public InProcessMessageBus(ILogger logger) + { + _logger = logger; + } + + private Channel<(string Key, string Value)> GetOrCreate(string topic) => + _channels.GetOrAdd(topic, _ => + { + _logger.LogDebug("In-process message bus: creating channel for topic {Topic}", topic); + // Unbounded so producers never block — at hobby scale the queue + // depth stays trivial. If you observe unbounded growth here, that's + // a sign you're outgrowing the in-process bus. + return Channel.CreateUnbounded<(string, string)>(new UnboundedChannelOptions + { + SingleReader = false, // multiple consumer instances OK + SingleWriter = false, // multiple producer call sites OK + }); + }); + + public async Task PublishAsync(string topic, string key, T value, CancellationToken ct) + { + var json = JsonSerializer.Serialize(value); + var ch = GetOrCreate(topic); + await ch.Writer.WriteAsync((key, json), ct); + } + + public async IAsyncEnumerable<(string Key, string Value)> SubscribeAsync( + string topic, + [EnumeratorCancellation] CancellationToken ct) + { + var ch = GetOrCreate(topic); + await foreach (var item in ch.Reader.ReadAllAsync(ct)) + yield return item; + } +} diff --git a/src/dotnet/src/HoldFast.Shared/Messaging/MessageConsumerBase.cs b/src/dotnet/src/HoldFast.Shared/Messaging/MessageConsumerBase.cs new file mode 100644 index 00000000..e8adabf5 --- /dev/null +++ b/src/dotnet/src/HoldFast.Shared/Messaging/MessageConsumerBase.cs @@ -0,0 +1,77 @@ +using System.Text.Json; +using Microsoft.Extensions.Logging; + +namespace HoldFast.Shared.Messaging; + +/// +/// Base consumer that reads JSON-serialized messages from an +/// topic and dispatches them to a +/// strongly-typed implementation. +/// +/// Replaces the Kafka-specific KafkaConsumerService base class (HOL-23). +/// Subclasses just specify the topic and group identifier; the consume loop +/// + JSON deserialization is shared. +/// +public abstract class MessageConsumerBase where T : class +{ + private readonly IMessageBus _bus; + private readonly ILogger _logger; + private readonly string _topic; + private readonly string _groupId; + + protected MessageConsumerBase( + IMessageBus bus, + string topic, + string groupId, + ILogger logger) + { + _bus = bus; + _topic = topic; + _groupId = groupId; + _logger = logger; + } + + /// + /// Process a single consumed message. Implementations define the business logic. + /// + protected abstract Task ProcessAsync(string key, T value, CancellationToken ct); + + /// + /// Run the consume loop. Call from BackgroundService.ExecuteAsync. + /// + public async Task ConsumeLoopAsync(CancellationToken ct) + { + _logger.LogInformation( + "In-process consumer started for topic {Topic} (group {Group})", + _topic, _groupId); + + try + { + await foreach (var (key, body) in _bus.SubscribeAsync(_topic, ct)) + { + try + { + var value = JsonSerializer.Deserialize(body); + if (value != null) + await ProcessAsync(key, value, ct); + } + catch (JsonException ex) + { + _logger.LogError(ex, "Failed to deserialize message from {Topic}", _topic); + } + catch (Exception ex) + { + _logger.LogError(ex, "Error processing message from {Topic}", _topic); + // No commit semantics — message is already consumed. Loss + // here is the same tradeoff documented in + // InProcessMessageBus (acceptable at hobby scale; SDK-side + // retry covers transient loss). + } + } + } + catch (OperationCanceledException) + { + // Normal shutdown + } + } +} diff --git a/src/dotnet/src/HoldFast.Worker/ErrorGroupingWorker.cs b/src/dotnet/src/HoldFast.Worker/ErrorGroupingWorker.cs index c2bc3cae..bfc21622 100644 --- a/src/dotnet/src/HoldFast.Worker/ErrorGroupingWorker.cs +++ b/src/dotnet/src/HoldFast.Worker/ErrorGroupingWorker.cs @@ -3,6 +3,7 @@ using HoldFast.Shared.AlertEvaluation; using HoldFast.Shared.ErrorGrouping; using HoldFast.Shared.Kafka; +using HoldFast.Shared.Messaging; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; @@ -34,16 +35,16 @@ public record BackendErrorMessage( /// Consumes backend errors from Kafka, groups them, and stores to database. /// Replaces the Go worker's processBackendPayloadImpl handler. /// -public class ErrorGroupingConsumer : KafkaConsumerService +public class ErrorGroupingConsumer : MessageConsumerBase { private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; public ErrorGroupingConsumer( - IOptions options, + IMessageBus bus, IServiceScopeFactory scopeFactory, ILogger logger) - : base(options, KafkaTopics.BackendErrors, "error-grouping-worker", logger) + : base(bus, KafkaTopics.BackendErrors, "error-grouping-worker", logger) { _scopeFactory = scopeFactory; _logger = logger; diff --git a/src/dotnet/src/HoldFast.Worker/FrontendErrorsWorker.cs b/src/dotnet/src/HoldFast.Worker/FrontendErrorsWorker.cs index 586dd352..28e0c83d 100644 --- a/src/dotnet/src/HoldFast.Worker/FrontendErrorsWorker.cs +++ b/src/dotnet/src/HoldFast.Worker/FrontendErrorsWorker.cs @@ -5,6 +5,7 @@ using HoldFast.Shared.AlertEvaluation; using HoldFast.Shared.ErrorGrouping; using HoldFast.Shared.Kafka; +using HoldFast.Shared.Messaging; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; @@ -41,16 +42,16 @@ public record FrontendErrorMessage( /// IErrorGroupingService. Mirrors how ErrorGroupingConsumer handles backend /// errors, but resolves projectId via the session. /// -public class FrontendErrorsConsumer : KafkaConsumerService +public class FrontendErrorsConsumer : MessageConsumerBase { private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; public FrontendErrorsConsumer( - IOptions options, + IMessageBus bus, IServiceScopeFactory scopeFactory, ILogger logger) - : base(options, KafkaTopics.FrontendErrors, "frontend-errors-worker", logger) + : base(bus, KafkaTopics.FrontendErrors, "frontend-errors-worker", logger) { _scopeFactory = scopeFactory; _logger = logger; diff --git a/src/dotnet/src/HoldFast.Worker/LogIngestionWorker.cs b/src/dotnet/src/HoldFast.Worker/LogIngestionWorker.cs index df37141d..31866858 100644 --- a/src/dotnet/src/HoldFast.Worker/LogIngestionWorker.cs +++ b/src/dotnet/src/HoldFast.Worker/LogIngestionWorker.cs @@ -1,6 +1,7 @@ using HoldFast.Data.ClickHouse; using HoldFast.Data.ClickHouse.Models; using HoldFast.Shared.Kafka; +using HoldFast.Shared.Messaging; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; @@ -29,16 +30,16 @@ public record LogIngestionMessage( /// /// Consumes log rows from Kafka and writes to ClickHouse. /// -public class LogIngestionConsumer : KafkaConsumerService +public class LogIngestionConsumer : MessageConsumerBase { private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; public LogIngestionConsumer( - IOptions options, + IMessageBus bus, IServiceScopeFactory scopeFactory, ILogger logger) - : base(options, KafkaTopics.Logs, "log-ingestion-worker", logger) + : base(bus, KafkaTopics.Logs, "log-ingestion-worker", logger) { _scopeFactory = scopeFactory; _logger = logger; diff --git a/src/dotnet/src/HoldFast.Worker/MetricsWorker.cs b/src/dotnet/src/HoldFast.Worker/MetricsWorker.cs index 82149753..44db855e 100644 --- a/src/dotnet/src/HoldFast.Worker/MetricsWorker.cs +++ b/src/dotnet/src/HoldFast.Worker/MetricsWorker.cs @@ -1,5 +1,6 @@ using HoldFast.Data.ClickHouse; using HoldFast.Shared.Kafka; +using HoldFast.Shared.Messaging; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; @@ -22,16 +23,16 @@ public record MetricsMessage( /// Consumes metrics from Kafka and writes to ClickHouse. /// Replaces the Go worker's pushMetrics handler. /// -public class MetricsConsumer : KafkaConsumerService +public class MetricsConsumer : MessageConsumerBase { private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; public MetricsConsumer( - IOptions options, + IMessageBus bus, IServiceScopeFactory scopeFactory, ILogger logger) - : base(options, KafkaTopics.Metrics, "metrics-worker", logger) + : base(bus, KafkaTopics.Metrics, "metrics-worker", logger) { _scopeFactory = scopeFactory; _logger = logger; diff --git a/src/dotnet/src/HoldFast.Worker/SessionEventsWorker.cs b/src/dotnet/src/HoldFast.Worker/SessionEventsWorker.cs index e4f31837..94cbd9cb 100644 --- a/src/dotnet/src/HoldFast.Worker/SessionEventsWorker.cs +++ b/src/dotnet/src/HoldFast.Worker/SessionEventsWorker.cs @@ -1,4 +1,5 @@ using HoldFast.Shared.Kafka; +using HoldFast.Shared.Messaging; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; @@ -18,16 +19,16 @@ public record SessionEventsMessage( /// Consumes session events from Kafka and processes them. /// Replaces the Go worker's processPublicWorkerMessage handler. /// -public class SessionEventsConsumer : KafkaConsumerService +public class SessionEventsConsumer : MessageConsumerBase { private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; public SessionEventsConsumer( - IOptions options, + IMessageBus bus, IServiceScopeFactory scopeFactory, ILogger logger) - : base(options, KafkaTopics.SessionEvents, "session-events-worker", logger) + : base(bus, KafkaTopics.SessionEvents, "session-events-worker", logger) { _scopeFactory = scopeFactory; _logger = logger; diff --git a/src/dotnet/src/HoldFast.Worker/TraceIngestionWorker.cs b/src/dotnet/src/HoldFast.Worker/TraceIngestionWorker.cs index b846ab6b..d6e13b7f 100644 --- a/src/dotnet/src/HoldFast.Worker/TraceIngestionWorker.cs +++ b/src/dotnet/src/HoldFast.Worker/TraceIngestionWorker.cs @@ -1,6 +1,7 @@ using HoldFast.Data.ClickHouse; using HoldFast.Data.ClickHouse.Models; using HoldFast.Shared.Kafka; +using HoldFast.Shared.Messaging; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; @@ -32,16 +33,16 @@ public record TraceIngestionMessage( /// /// Consumes trace spans from Kafka and writes to ClickHouse. /// -public class TraceIngestionConsumer : KafkaConsumerService +public class TraceIngestionConsumer : MessageConsumerBase { private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; public TraceIngestionConsumer( - IOptions options, + IMessageBus bus, IServiceScopeFactory scopeFactory, ILogger logger) - : base(options, KafkaTopics.Traces, "trace-ingestion-worker", logger) + : base(bus, KafkaTopics.Traces, "trace-ingestion-worker", logger) { _scopeFactory = scopeFactory; _logger = logger; diff --git a/src/dotnet/tests/HoldFast.Shared.Tests/Kafka/KafkaOptionsTests.cs b/src/dotnet/tests/HoldFast.Shared.Tests/Kafka/KafkaOptionsTests.cs deleted file mode 100644 index 4afc0873..00000000 --- a/src/dotnet/tests/HoldFast.Shared.Tests/Kafka/KafkaOptionsTests.cs +++ /dev/null @@ -1,39 +0,0 @@ -using HoldFast.Shared.Kafka; -using Xunit; - -namespace HoldFast.Shared.Tests.Kafka; - -/// -/// Tests for KafkaOptions configuration defaults and overrides. -/// -public class KafkaOptionsTests -{ - [Fact] - public void KafkaOptions_DefaultBootstrapServers() - { - var options = new KafkaOptions(); - Assert.Equal("localhost:9092", options.BootstrapServers); - } - - [Fact] - public void KafkaOptions_SetBootstrapServers() - { - var options = new KafkaOptions { BootstrapServers = "kafka1:9092,kafka2:9092" }; - Assert.Equal("kafka1:9092,kafka2:9092", options.BootstrapServers); - } - - [Fact] - public void KafkaOptions_EmptyBootstrapServers() - { - var options = new KafkaOptions { BootstrapServers = "" }; - Assert.Equal("", options.BootstrapServers); - } - - [Fact] - public void KafkaOptions_MultipleServers() - { - var options = new KafkaOptions { BootstrapServers = "host1:9092,host2:9092,host3:9092" }; - var servers = options.BootstrapServers.Split(','); - Assert.Equal(3, servers.Length); - } -} diff --git a/src/dotnet/tests/HoldFast.Worker.Tests/ConsumerProcessAsyncTests.cs b/src/dotnet/tests/HoldFast.Worker.Tests/ConsumerProcessAsyncTests.cs index a1a443ec..c5826fa2 100644 --- a/src/dotnet/tests/HoldFast.Worker.Tests/ConsumerProcessAsyncTests.cs +++ b/src/dotnet/tests/HoldFast.Worker.Tests/ConsumerProcessAsyncTests.cs @@ -5,6 +5,7 @@ using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Options; using HoldFast.Shared.Kafka; +using HoldFast.Shared.Messaging; using Xunit; namespace HoldFast.Worker.Tests; @@ -39,7 +40,7 @@ public ConsumerProcessAsyncTests() private class TestableMetricsConsumer : MetricsConsumer { public TestableMetricsConsumer(IServiceScopeFactory sf) - : base(Options.Create(new KafkaOptions { BootstrapServers = "test:9092" }), + : base(new InProcessMessageBus(NullLogger.Instance), sf, NullLogger.Instance) { } public new Task ProcessAsync(string key, MetricsMessage value, CancellationToken ct) @@ -117,7 +118,7 @@ public async Task MetricsConsumer_ManyTags() private class TestableLogConsumer : LogIngestionConsumer { public TestableLogConsumer(IServiceScopeFactory sf) - : base(Options.Create(new KafkaOptions { BootstrapServers = "test:9092" }), + : base(new InProcessMessageBus(NullLogger.Instance), sf, NullLogger.Instance) { } public new Task ProcessAsync(string key, LogIngestionMessage value, CancellationToken ct) @@ -193,7 +194,7 @@ public async Task LogConsumer_HighSeverityNumber() private class TestableTraceConsumer : TraceIngestionConsumer { public TestableTraceConsumer(IServiceScopeFactory sf) - : base(Options.Create(new KafkaOptions { BootstrapServers = "test:9092" }), + : base(new InProcessMessageBus(NullLogger.Instance), sf, NullLogger.Instance) { } public new Task ProcessAsync(string key, TraceIngestionMessage value, CancellationToken ct)