Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/website.yml
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ jobs:
- name: Setup Deno
uses: denoland/setup-deno@v2
with:
deno-version: 2.5.7
deno-version: 2.9.7
cache: true

- name: Require telemetry configuration for deployment
Expand Down
8 changes: 4 additions & 4 deletions docs/serialization-architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -109,11 +109,11 @@ Messages: `EntityChanged`, `PlanChanged`, `UserMembershipChanged`, `ReleaseNotif

Cached values: Organizations, Projects, Stacks, Tokens, Users. Cache is ephemeral — keys expire. No migration needed; old cached values are simply evicted.

### Path 6: WebSocket Messages
### Path 6: Real-Time Push Messages (SSE and WebSocket)

**Config:** `WebSocketConnectionManager` resolves `ITextSerializer` from DI.
**Config:** `SseConnectionManager` and the temporary rollout-compatible `WebSocketConnectionManager` resolve `ITextSerializer` from DI.

Messages sent to browser clients use `serializer.SerializeToString(message)` with the standard snake_case options. The JavaScript/TypeScript frontend expects snake_case.
Messages sent to browser clients use `serializer.SerializeToString(message)` with the standard snake_case options. The Svelte client receives the payload over SSE, while the legacy Angular client temporarily remains on WebSocket. Both clients consume the same `TypedMessage` JSON contract and expect snake_case.

### Path 7: Webhook Payloads

Expand Down Expand Up @@ -445,7 +445,7 @@ Every code path that calls `GetValue<T>()`, mutates the result, and needs to per
| Existing ES documents | Read fine — `PropertyNameCaseInsensitive` handles any key format |
| In-flight queue messages | STJ reads Newtonsoft output (case-insensitive DTOs) |
| Cached values | Ephemeral with TTL, auto-expire |
| WebSocket messages | Consumers expect snake_case (unchanged) |
| SSE and WebSocket push messages | Both rollout transports share the same snake_case `TypedMessage` contract |
| Webhook consumers | snake_case output (unchanged from Newtonsoft era) |
| Old client submissions | Event upgrader handles V1/V2 format |
| Custom data keys in Event.Data | Preserved as-is (no DictionaryKeyPolicy in main serializer) |
Expand Down
19 changes: 18 additions & 1 deletion k8s/exceptionless/templates/api.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,11 @@ spec:
annotations:
checksum/config: {{ include (print $.Template.BasePath "/config.yaml") . | sha256sum }}
spec:
# SSE connections are long-lived; give the pod enough time to drain before SIGTERM.
# The preStop sleep lets the ALB/ingress controller deregister the pod before traffic stops,
# then the remaining window allows ASP.NET Core to cancel RequestAborted tokens and clean up.
# When push is eventually enabled behind a Gateway API RoutePolicy, revisit this value.
terminationGracePeriodSeconds: 60
topologySpreadConstraints:
- maxSkew: 1
topologyKey: kubernetes.io/hostname
Expand All @@ -39,6 +44,13 @@ spec:
- name: {{ template "exceptionless.name" . }}-api
image: "{{ .Values.api.image.repository }}:{{ .Values.version }}"
imagePullPolicy: {{ .Values.api.image.pullPolicy }}
lifecycle:
preStop:
# Give the ALB ~15s to deregister this pod before SIGTERM fires.
# The total graceful window is terminationGracePeriodSeconds (60s) minus
# this sleep, leaving 45s. ASP.NET Core uses 40s and keeps 5s of safety margin.
exec:
command: ["sleep", "15"]
livenessProbe:
httpGet:
path: /health
Expand Down Expand Up @@ -82,7 +94,10 @@ spec:
{{- include "exceptionless.otel-env" . | indent 12 }}
- name: RunJobsInProcess
value: 'false'
- name: EnableWebSockets
# SSE rollout prerequisite: Azure Application Gateway for Containers Ingress API
# does not support the routeTimeout=0s override required for long-lived SSE streams.
# Keep push disabled here until this route moves to Gateway API + RoutePolicy.
- name: EnablePush
value: 'false'
{{- if (empty .Values.storage.connectionString) }}
volumeMounts:
Expand Down Expand Up @@ -163,6 +178,8 @@ metadata:
alb.networking.azure.io/alb-namespace: {{ .Values.ingress.albNamespace }}
alb.networking.azure.io/alb-frontend: {{ template "exceptionless.fullname" . }}-fe
cert-manager.io/cluster-issuer: {{ .Values.ingress.clusterIssuer }}
# SSE is not safe to enable behind the current AGC Ingress API path.
# Migrate to Gateway API and attach a RoutePolicy with routeTimeout: 0s before enabling push.
spec:
ingressClassName: azure-alb-external
tls:
Expand Down
9 changes: 5 additions & 4 deletions src/Exceptionless.Core/Bootstrapper.cs
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,7 @@ public static void RegisterServices(IServiceCollection services, AppOptions appO
services.TryAddEnumerable(ServiceDescriptor.Singleton<IQueueBehavior<WorkItemData>, WorkItemDuplicateDetectionQueueBehavior>());

services.AddSingleton<IConnectionMapping, ConnectionMapping>();
services.AddSingleton<IConnectionLeaseStore, ConnectionLeaseStore>();
services.AddSingleton<MessageService>();
services.AddStartupAction<MessageService>();
services.AddSingleton<IMessageBus>(s => new InMemoryMessageBus(new InMemoryMessageBusOptions
Expand Down Expand Up @@ -282,8 +283,8 @@ public static void LogConfiguration(IServiceProvider serviceProvider, AppOptions
if (String.IsNullOrEmpty(appOptions.StorageOptions.Provider))
logger.LogWarning("Distributed storage is NOT enabled on {MachineName}", Environment.MachineName);

if (!appOptions.EnableWebSockets)
logger.LogWarning("Web Sockets is NOT enabled on {MachineName}", Environment.MachineName);
if (!appOptions.EnablePush)
logger.LogWarning("Real-time push (SSE) is NOT enabled on {MachineName}", Environment.MachineName);

if (String.IsNullOrEmpty(appOptions.EmailOptions.SmtpHost))
logger.LogWarning("Emails will NOT be sent until the SmtpHost is configured on {MachineName}", Environment.MachineName);
Expand Down Expand Up @@ -328,9 +329,9 @@ private static void LogConfigurationSummary(IServiceProvider serviceProvider, Ap
GetProviderEndpoint(options.StorageOptions.Data, options.StorageOptions.ConnectionString));

logger.LogInformation(
"Startup services: event submission {EventSubmission}; WebSockets {WebSockets}; jobs in process {JobsInProcess}; email {Email}; account creation {AccountCreation}; index configuration {IndexConfiguration}",
"Startup services: event submission {EventSubmission}; push {Push}; jobs in process {JobsInProcess}; email {Email}; account creation {AccountCreation}; index configuration {IndexConfiguration}",
GetStatus(!options.EventSubmissionDisabled),
GetStatus(options.EnableWebSockets),
GetStatus(options.EnablePush),
GetStatus(options.RunJobsInProcess),
GetStatus(!String.IsNullOrWhiteSpace(options.EmailOptions.SmtpHost)),
GetStatus(options.AuthOptions.EnableAccountCreation),
Expand Down
16 changes: 14 additions & 2 deletions src/Exceptionless.Core/Configuration/AppOptions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,18 @@ public class AppOptions

public bool EnableRepositoryNotifications { get; internal set; }

public bool EnableWebSockets { get; internal set; }
/// <summary>
/// Controls whether real-time push (SSE) is enabled. Reads from either 'EnablePush'
/// or legacy 'EnableWebSockets' config key for backward compatibility.
/// </summary>
public bool EnablePush { get; internal set; }

[Obsolete("Use EnablePush instead. This alias remains for source compatibility during the push transport rollout.")]
public bool EnableWebSockets
{
get => EnablePush;
internal set => EnablePush = value;
}

public string? Version { get; internal set; }

Expand Down Expand Up @@ -114,7 +125,8 @@ public static AppOptions ReadFromConfiguration(IConfiguration config)
options.BulkBatchSize = config.GetValue(nameof(options.BulkBatchSize), 1000);

options.EnableRepositoryNotifications = config.GetValue(nameof(options.EnableRepositoryNotifications), true);
options.EnableWebSockets = config.GetValue(nameof(options.EnableWebSockets), true);
// Support both new 'EnablePush' and legacy 'EnableWebSockets' config keys
options.EnablePush = config.GetValue(nameof(options.EnablePush), config.GetValue("EnableWebSockets", true));
Comment thread
niemyjski marked this conversation as resolved.

try
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,5 +51,6 @@ public class KnownKeys
public const string StackId = nameof(StackId);
public const string UserId = nameof(UserId);
public const string IsAuthenticationToken = nameof(IsAuthenticationToken);
public const string IsTokenRevoked = nameof(IsTokenRevoked);
}
}
11 changes: 11 additions & 0 deletions src/Exceptionless.Core/Repositories/OAuthTokenRepository.cs
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
using Exceptionless.Core.Messaging.Models;
using Exceptionless.Core.Models;
using Exceptionless.Core.Repositories.Configuration;
using Exceptionless.Core.Validation;
Expand Down Expand Up @@ -93,4 +94,14 @@ public Task<long> RemoveAllByUserIdAsync(string userId, CommandOptionsDescriptor
{
return RemoveAllAsync(q => q.FieldEquals(t => t.UserId, userId), options);
}

protected override Task PublishChangeTypeMessageAsync(ChangeType changeType, OAuthToken? document, IDictionary<string, object?>? data = null, TimeSpan? delay = null)
{
var items = new Dictionary<string, object?>(data ?? new Dictionary<string, object?>())
{
[ExtendedEntityChanged.KnownKeys.IsTokenRevoked] = document is { IsDisabled: true } or { IsSuspended: true }
};

return base.PublishChangeTypeMessageAsync(changeType, document, items, delay);
}
}
29 changes: 10 additions & 19 deletions src/Exceptionless.Core/Services/MessageService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,6 @@
using Foundatio.Repositories.Elasticsearch;
using Foundatio.Repositories.Models;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Logging.Abstractions;

namespace Exceptionless.Core.Services;

Expand All @@ -19,9 +18,7 @@ public class MessageService : IDisposable, IStartupAction
private readonly IEventRepository _eventRepository;
private readonly ITokenRepository _tokenRepository;
private readonly IWebHookRepository _webHookRepository;
private readonly IConnectionMapping _connectionMapping;
private readonly AppOptions _options;
private readonly ILogger _logger;
private readonly List<Action> _disposeActions = [];

public MessageService(
Expand All @@ -43,9 +40,12 @@ public MessageService(
_eventRepository = eventRepository;
_tokenRepository = tokenRepository;
_webHookRepository = webHookRepository;
_connectionMapping = connectionMapping;
_options = options;
_logger = loggerFactory.CreateLogger<MessageService>() ?? NullLogger<MessageService>.Instance;

// Preserve the public constructor shape for existing composition roots while push
// publication no longer depends on distributed connection state or trace logging.
_ = connectionMapping;
_ = loggerFactory;
}

public Task RunAsync(CancellationToken shutdownToken = default)
Expand Down Expand Up @@ -74,22 +74,13 @@ public Task RunAsync(CancellationToken shutdownToken = default)
_disposeActions.Add(() => repo.BeforePublishEntityChanged.RemoveHandler(handler));
}

private async Task OnBeforePublishEntityChangedAsync<T>(object sender, BeforePublishEntityChangedEventArgs<T> args)
private Task OnBeforePublishEntityChangedAsync<T>(object sender, BeforePublishEntityChangedEventArgs<T> args)
where T : class, IIdentity, new()
{
int listenerCount = await GetNumberOfListeners(args.Message);
args.Cancel = listenerCount == 0;
if (args.Cancel)
_logger.LogTrace("Cancelled {EntityType} Entity Changed Message: {@Message}", typeof(T).Name, args.Message);
}

private Task<int> GetNumberOfListeners(EntityChanged message)
{
var entityChanged = ExtendedEntityChanged.Create(message, false);
if (String.IsNullOrEmpty(entityChanged.OrganizationId))
return Task.FromResult(1); // Return 1 as we have no idea if people are listening.

return _connectionMapping.GetGroupConnectionCountAsync(entityChanged.OrganizationId);
// Push routing is replica-local. Publishing must not depend on immortal distributed
// connection indexes because another replica may own the interested client.
args.Cancel = false;
return Task.CompletedTask;
}

public void Dispose()
Expand Down
5 changes: 5 additions & 0 deletions src/Exceptionless.Core/Utility/AppDiagnostics.cs
Original file line number Diff line number Diff line change
Expand Up @@ -132,6 +132,11 @@ public GaugeInfo(Meter meter, string name)

internal static readonly Counter<int> SavedViewsSize = Meter.CreateCounter<int>("ex.savedviews.size", description: "Size of user saved views");
internal static readonly Counter<int> SavedViewsViewTypeSize = Meter.CreateCounter<int>("ex.savedviews.viewtype.size", description: "Size of user saved views by view type");

internal static readonly Counter<int> PushSseConnectionsOpened = Meter.CreateCounter<int>("ex.push.connections.sse.opened", description: "SSE push connections opened");
internal static readonly Counter<int> PushSseConnectionsClosed = Meter.CreateCounter<int>("ex.push.connections.sse.closed", description: "SSE push connections closed");
internal static readonly Counter<int> PushWebSocketConnectionsOpened = Meter.CreateCounter<int>("ex.push.connections.websocket.opened", description: "WebSocket push connections opened");
internal static readonly Counter<int> PushWebSocketConnectionsClosed = Meter.CreateCounter<int>("ex.push.connections.websocket.closed", description: "WebSocket push connections closed");
}

public static class MetricsClientExtensions
Expand Down
77 changes: 77 additions & 0 deletions src/Exceptionless.Core/Utility/IConnectionLeaseStore.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
namespace Exceptionless.Core.Utility;

public interface IConnectionLeaseStore
{
Task<bool> TryAcquireAsync(string userId, string connectionId, int maxConnections, TimeSpan leaseDuration);
Task<bool> RenewAsync(string userId, string connectionId, TimeSpan leaseDuration);
Task ReleaseAsync(string userId, string connectionId);
}

public sealed class ConnectionLeaseStore(TimeProvider timeProvider) : IConnectionLeaseStore
{
private readonly Dictionary<string, Dictionary<string, DateTimeOffset>> _leases = [];
private readonly object _lock = new();

public Task<bool> TryAcquireAsync(string userId, string connectionId, int maxConnections, TimeSpan leaseDuration)
{
if (maxConnections <= 0)
return Task.FromResult(false);

lock (_lock)
{
var now = timeProvider.GetUtcNow();
if (!_leases.TryGetValue(userId, out var userLeases))
{
userLeases = [];
_leases[userId] = userLeases;
}

RemoveExpired(userLeases, now);
if (!userLeases.ContainsKey(connectionId) && userLeases.Count >= maxConnections)
return Task.FromResult(false);

userLeases[connectionId] = now + leaseDuration;
return Task.FromResult(true);
}
}

public Task<bool> RenewAsync(string userId, string connectionId, TimeSpan leaseDuration)
{
lock (_lock)
{
var now = timeProvider.GetUtcNow();
if (!_leases.TryGetValue(userId, out var userLeases))
return Task.FromResult(false);

RemoveExpired(userLeases, now);
if (!userLeases.ContainsKey(connectionId))
return Task.FromResult(false);

userLeases[connectionId] = now + leaseDuration;
return Task.FromResult(true);
}
}

public Task ReleaseAsync(string userId, string connectionId)
{
lock (_lock)
{
if (_leases.TryGetValue(userId, out var userLeases))
{
userLeases.Remove(connectionId);
if (userLeases.Count is 0)
_leases.Remove(userId);
}
}

return Task.CompletedTask;
}

private static void RemoveExpired(Dictionary<string, DateTimeOffset> leases, DateTimeOffset now)
{
foreach (string connectionId in leases.Where(pair => pair.Value <= now).Select(pair => pair.Key).ToArray())
leases.Remove(connectionId);
}
}

public sealed class ConnectionLeaseStoreException(string message, Exception innerException) : Exception(message, innerException);
2 changes: 2 additions & 0 deletions src/Exceptionless.Insulation/Bootstrapper.cs
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,7 @@ private static void RegisterCache(IServiceCollection container, CacheOptions opt
container.ReplaceSingleton<ICacheClient>(CreateRedisCacheClient);

container.ReplaceSingleton<IConnectionMapping, RedisConnectionMapping>();
container.ReplaceSingleton<IConnectionLeaseStore, RedisConnectionLeaseStore>();
}
}

Expand All @@ -120,6 +121,7 @@ private static void RegisterMessageBus(IServiceCollection container, MessageBusO
if (String.Equals(options.Provider, "redis"))
{
container.ReplaceSingleton(s => GetRedisConnection(options.ConnectionString!, s.GetRequiredService<ILoggerFactory>()));
container.ReplaceSingleton<IConnectionLeaseStore, RedisConnectionLeaseStore>();

container.ReplaceSingleton<IMessageBus>(s => new RedisMessageBus(new RedisMessageBusOptions
{
Expand Down
Loading
Loading