diff --git a/.editorconfig b/.editorconfig
new file mode 100644
index 0000000..35161b2
--- /dev/null
+++ b/.editorconfig
@@ -0,0 +1,86 @@
+root = true
+
+[*]
+charset = utf-8
+end_of_line = lf
+trim_trailing_whitespace = true
+insert_final_newline = true
+indent_style = space
+indent_size = 4
+
+[*.cs]
+# --- Language style ---
+csharp_preferred_modifier_order = public, private, protected, internal, new, static, abstract, virtual, sealed, readonly, override, extern, unsafe, volatile, async, file, required:suggestion
+csharp_space_after_cast = true
+csharp_new_line_before_members_in_object_initializers = false
+csharp_style_var_for_built_in_types = true:suggestion
+csharp_style_var_when_type_is_apparent = true:suggestion
+csharp_style_var_elsewhere = true:suggestion
+
+dotnet_style_predefined_type_for_locals_parameters_members = true:suggestion
+dotnet_style_predefined_type_for_member_access = true:suggestion
+dotnet_style_qualification_for_field = true:suggestion
+dotnet_style_qualification_for_property = true:suggestion
+dotnet_style_qualification_for_method = true:suggestion
+dotnet_style_qualification_for_event = true:suggestion
+dotnet_style_require_accessibility_modifiers = for_non_interface_members:suggestion
+
+# --- Expression-bodied members: use an expression body when the member is a single statement ---
+csharp_style_expression_bodied_methods = when_on_single_line:suggestion
+csharp_style_expression_bodied_constructors = when_on_single_line:suggestion
+csharp_style_expression_bodied_operators = when_on_single_line:suggestion
+csharp_style_expression_bodied_local_functions = when_on_single_line:suggestion
+csharp_style_expression_bodied_lambdas = when_on_single_line:suggestion
+csharp_style_expression_bodied_properties = when_on_single_line:suggestion
+csharp_style_expression_bodied_indexers = when_on_single_line:suggestion
+csharp_style_expression_bodied_accessors = when_on_single_line:suggestion
+
+# --- Naming (warning-level) ---
+# private instance / static fields: camelCase, NO underscore prefix (CONVENTIONS §4)
+dotnet_naming_rule.private_fields_camel_case.severity = warning
+dotnet_naming_rule.private_fields_camel_case.symbols = private_fields
+dotnet_naming_rule.private_fields_camel_case.style = camel_case
+dotnet_naming_symbols.private_fields.applicable_kinds = field
+dotnet_naming_symbols.private_fields.applicable_accessibilities = private
+
+# private const: PascalCase (declared before the general rule so it wins)
+dotnet_naming_rule.private_const_pascal_case.severity = warning
+dotnet_naming_rule.private_const_pascal_case.symbols = private_const_fields
+dotnet_naming_rule.private_const_pascal_case.style = pascal_case
+dotnet_naming_symbols.private_const_fields.applicable_kinds = field
+dotnet_naming_symbols.private_const_fields.applicable_accessibilities = private
+dotnet_naming_symbols.private_const_fields.required_modifiers = const
+
+# private static readonly: PascalCase
+dotnet_naming_rule.private_static_readonly_pascal_case.severity = warning
+dotnet_naming_rule.private_static_readonly_pascal_case.symbols = private_static_readonly_fields
+dotnet_naming_rule.private_static_readonly_pascal_case.style = pascal_case
+dotnet_naming_symbols.private_static_readonly_fields.applicable_kinds = field
+dotnet_naming_symbols.private_static_readonly_fields.applicable_accessibilities = private
+dotnet_naming_symbols.private_static_readonly_fields.required_modifiers = static, readonly
+
+dotnet_naming_style.camel_case.capitalization = camel_case
+dotnet_naming_style.pascal_case.capitalization = pascal_case
+
+# --- ReSharper / Rider hints (IDE-only; not enforced by dotnet build) ---
+# Collapse blank lines in method bodies to 0 — backs the "no blank lines between arrange/act/assert"
+# rule. Rider-enforced; CI cannot enforce this via editorconfig alone.
+resharper_csharp_keep_blank_lines_in_code = 0
+resharper_csharp_keep_blank_lines_in_declarations = 1
+resharper_braces_for_ifelse = required
+resharper_braces_for_for = required
+resharper_braces_for_foreach = required
+resharper_braces_for_while = required
+resharper_trailing_comma_in_multiline_lists = true
+
+[*.{csproj,props,targets,nuspec}]
+indent_size = 2
+
+[*.{json,jsonc}]
+indent_size = 2
+
+[*.{yaml,yml}]
+indent_size = 2
+
+[*.{bash,sh,zsh}]
+indent_size = 2
diff --git a/RealtimeTests/Broadcast/BroadcastMessageTests.cs b/RealtimeTests/Broadcast/BroadcastMessageTests.cs
new file mode 100644
index 0000000..4e71033
--- /dev/null
+++ b/RealtimeTests/Broadcast/BroadcastMessageTests.cs
@@ -0,0 +1,51 @@
+using System;
+using FluentAssertions;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using RealtimeTests.Models;
+using RealtimeTests.Support;
+using Supabase.Realtime.Socket;
+
+namespace RealtimeTests.Broadcast;
+
+///
+/// What a channel does with an inbound broadcast frame: it routes it to the registered broadcast
+/// instance, which types the payload and hands it to listeners. Driven through the channel's public
+/// message entry point over a fake socket, so no live server is needed.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class BroadcastMessageTests
+{
+ private const string Frame =
+ "{\"topic\":\"realtime:example\",\"event\":\"broadcast\",\"ref\":\"srv-1\",\"payload\":{\"event\":\"user\",\"userId\":\"abc-123\"}}";
+
+ [TestMethod]
+ public void HandleSocketMessage_ShouldDeliverTypedBroadcastToListener()
+ {
+ var channel = Wire.Channel();
+ var broadcast = channel.Register();
+ BroadcastExample? received = null;
+ broadcast.AddBroadcastEventHandler((_, _) => received = broadcast.Current());
+ channel.HandleSocketMessage(Wire.Decode(Frame));
+ received.Should().NotBeNull();
+ received!.UserId.Should().Be("abc-123");
+ received.Event.Should().Be("user");
+ }
+
+ [TestMethod]
+ public void Current_ShouldReturnNull_GivenNoBroadcastReceived()
+ {
+ var channel = Wire.Channel();
+ var broadcast = channel.Register();
+ broadcast.Current().Should().BeNull();
+ }
+
+ [TestMethod]
+ public void TriggerReceived_ShouldThrow_GivenUnparsableResponse()
+ {
+ var channel = Wire.Channel();
+ var broadcast = channel.Register();
+ var act = () => broadcast.TriggerReceived(new SocketResponse(Wire.Settings()));
+ act.Should().Throw();
+ }
+}
diff --git a/RealtimeTests/Broadcast/BroadcastRegistrationTests.cs b/RealtimeTests/Broadcast/BroadcastRegistrationTests.cs
new file mode 100644
index 0000000..cd164ac
--- /dev/null
+++ b/RealtimeTests/Broadcast/BroadcastRegistrationTests.cs
@@ -0,0 +1,46 @@
+using System;
+using FluentAssertions;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using RealtimeTests.Models;
+using RealtimeTests.Support;
+using Supabase.Realtime.Broadcast;
+
+namespace RealtimeTests.Broadcast;
+
+///
+/// Client-side rules a channel enforces when registering for broadcast, before any server is involved:
+/// broadcast can only be registered once, and replay-from-history requires a private channel.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class BroadcastRegistrationTests
+{
+ private static BroadcastOptions WithReplay() => new()
+ {
+ Replay = new BroadcastOptions.ReplayOptions { Limit = 10, Since = 0 }
+ };
+
+ [TestMethod]
+ public void Register_ShouldThrow_GivenReplayOnPublicChannel()
+ {
+ var channel = Wire.Channel(options: Wire.PublicOptions());
+ var act = () => channel.Register(WithReplay());
+ act.Should().Throw();
+ }
+
+ [TestMethod]
+ public void Register_ShouldSucceed_GivenReplayOnPrivateChannel()
+ {
+ var channel = Wire.Channel(options: Wire.PrivateOptions());
+ channel.Register(WithReplay()).Should().NotBeNull();
+ }
+
+ [TestMethod]
+ public void Register_ShouldThrow_GivenCalledTwice()
+ {
+ var channel = Wire.Channel();
+ channel.Register();
+ var act = () => channel.Register();
+ act.Should().Throw();
+ }
+}
diff --git a/RealtimeTests/Broadcast/BroadcastRelayTests.cs b/RealtimeTests/Broadcast/BroadcastRelayTests.cs
new file mode 100644
index 0000000..9b82f8d
--- /dev/null
+++ b/RealtimeTests/Broadcast/BroadcastRelayTests.cs
@@ -0,0 +1,155 @@
+using System;
+using System.Collections.Generic;
+using System.Threading.Tasks;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using Newtonsoft.Json;
+using RealtimeTests.Models;
+using Supabase.Postgrest.Interfaces;
+using Supabase.Realtime;
+using Supabase.Realtime.Broadcast;
+using Supabase.Realtime.Channel;
+using Supabase.Realtime.Interfaces;
+
+namespace RealtimeTests.Broadcast;
+
+///
+/// Broadcast against the live stack: messages relay between two clients, work on private (RLS-authorized)
+/// channels, replay from history on a private channel, and a send resolves even when no server ack was
+/// configured (supabase-community/realtime-csharp#38).
+///
+[TestClass]
+[TestCategory("E2E")]
+public class BroadcastRelayTests
+{
+ private IPostgrestClient restClient = null!;
+ private IRealtimeClient socketClient = null!;
+
+ [TestInitialize]
+ public async Task InitializeTest()
+ {
+ restClient = Helpers.RestClient();
+ socketClient = Helpers.SocketClient();
+ await socketClient.ConnectAsync();
+ }
+
+ [TestCleanup]
+ public void CleanupTest() => socketClient.Disconnect();
+
+ [TestMethod]
+ public async Task Broadcast_ShouldRelayBetweenClients()
+ {
+ var tsc = new TaskCompletionSource();
+ var tsc2 = new TaskCompletionSource();
+ var guid1 = Guid.NewGuid().ToString();
+ var guid2 = Guid.NewGuid().ToString();
+ var channel1 = socketClient.Channel("online-users");
+ var broadcast1 = channel1.Register(true, true);
+ broadcast1.AddBroadcastEventHandler((_, _) =>
+ {
+ var broadcast = broadcast1.Current();
+ if (broadcast?.UserId != guid1 && broadcast?.Event == "user")
+ tsc.TrySetResult(true);
+ });
+ var client2 = Helpers.SocketClient();
+ await client2.ConnectAsync();
+ var channel2 = client2.Channel("online-users");
+ var broadcast2 = channel2.Register(true, true);
+ broadcast2.AddBroadcastEventHandler((_, _) =>
+ {
+ var broadcast = broadcast2.Current();
+ if (broadcast?.UserId != guid2 && broadcast?.Event == "user")
+ tsc2.TrySetResult(true);
+ });
+ await channel1.Subscribe();
+ await channel2.Subscribe();
+ await broadcast1.Send("user", new BroadcastExample { UserId = guid1 });
+ await broadcast2.Send("user", new BroadcastExample { UserId = guid2 });
+ await Task.WhenAll(tsc.Task, tsc2.Task);
+ }
+
+ [TestMethod]
+ public async Task Broadcast_ShouldRelayOnPrivateChannel()
+ {
+ var tsc = new TaskCompletionSource();
+ var tsc2 = new TaskCompletionSource();
+ var guid1 = Guid.NewGuid().ToString();
+ var guid2 = Guid.NewGuid().ToString();
+ var client1 = Helpers.PrivateSocketClient();
+ await client1.ConnectAsync();
+ var channel1 = client1.Channel("online-users",
+ ChannelOptions.Private(client1.Options, () => Helpers.ApiKey, new JsonSerializerSettings()));
+ var broadcast1 = channel1.Register(true, true);
+ broadcast1.AddBroadcastEventHandler((_, _) =>
+ {
+ var broadcast = broadcast1.Current();
+ if (broadcast?.UserId != guid1 && broadcast?.Event == "user")
+ tsc.TrySetResult(true);
+ });
+ var client2 = Helpers.PrivateSocketClient();
+ await client2.ConnectAsync();
+ var channel2 = client2.Channel("online-users",
+ ChannelOptions.Private(client2.Options, () => Helpers.ApiKey, new JsonSerializerSettings()));
+ var broadcast2 = channel2.Register(true, true);
+ broadcast2.AddBroadcastEventHandler((_, _) =>
+ {
+ var broadcast = broadcast2.Current();
+ if (broadcast?.UserId != guid2 && broadcast?.Event == "user")
+ tsc2.TrySetResult(true);
+ });
+ await channel1.Subscribe();
+ await channel2.Subscribe();
+ await broadcast1.Send("user", new BroadcastExample { UserId = guid1 });
+ await broadcast2.Send("user", new BroadcastExample { UserId = guid2 });
+ await Task.WhenAll(tsc.Task, tsc2.Task);
+ }
+
+ [TestMethod]
+ public async Task Send_ShouldResolve_GivenNoExplicitAck()
+ {
+ // Mirrors the most natural usage: Channel(name) -> Subscribe() -> Send(), with no
+ // Register(broadcastAck: true). See supabase-community/realtime-csharp#38.
+ var channel = socketClient.Channel("no-ack-broadcast");
+ await channel.Subscribe();
+ var sendTask = channel.Send(Constants.ChannelEventName.Broadcast, "test_event",
+ new BroadcastExample { UserId = Guid.NewGuid().ToString() });
+ var winner = await Task.WhenAny(sendTask, Task.Delay(TimeSpan.FromSeconds(5)));
+ Assert.AreSame(sendTask, winner, "channel.Send() did not complete within 5s (issue #38).");
+ Assert.IsTrue(await sendTask);
+ }
+
+ [TestMethod]
+ public async Task Broadcast_ShouldReplayHistory_GivenPrivateChannel()
+ {
+ var send = new Dictionary
+ {
+ { "event", "user" },
+ { "topic", "online-users" },
+ { "private", true }
+ };
+ await restClient.Rpc("send", send);
+ var tsc = new TaskCompletionSource();
+ var client1 = Helpers.PrivateSocketClient();
+ await client1.ConnectAsync();
+ var broadcastOptions = new BroadcastOptions
+ {
+ BroadcastAck = true,
+ BroadcastSelf = true,
+ Replay = new BroadcastOptions.ReplayOptions
+ {
+ Limit = 10,
+ Since = DateTimeOffset.UtcNow.AddDays(-3).ToUnixTimeMilliseconds()
+ }
+ };
+ var channel1 = client1.Channel("online-users",
+ ChannelOptions.Private(client1.Options, () => null, new JsonSerializerSettings()));
+ var broadcast1 = channel1.Register(broadcastOptions);
+ broadcast1.AddBroadcastEventHandler((_, _) =>
+ {
+ var broadcast = broadcast1.Current();
+ if (broadcast is { Event: "user", Meta.Replayed: true })
+ tsc.TrySetResult(true);
+ });
+ await channel1.Subscribe();
+ await tsc.Task;
+ }
+}
diff --git a/RealtimeTests/ChannelBroadcastReplayTests.cs b/RealtimeTests/ChannelBroadcastReplayTests.cs
deleted file mode 100644
index 58fef3b..0000000
--- a/RealtimeTests/ChannelBroadcastReplayTests.cs
+++ /dev/null
@@ -1,56 +0,0 @@
-using System;
-using Microsoft.VisualStudio.TestTools.UnitTesting;
-using Newtonsoft.Json;
-using Supabase.Realtime;
-using Supabase.Realtime.Broadcast;
-using Supabase.Realtime.Channel;
-
-namespace RealtimeTests;
-
-///
-/// Client-side validation of broadcast replay registration. These tests exercise the guard on
-/// without a live
-/// server: the channel is built directly against an unconnected socket, so no stack is required.
-///
-[TestClass]
-public class ChannelBroadcastReplayTests
-{
- [TestMethod(DisplayName = "Channel: Registering broadcast replay on a public channel throws")]
- public void ClientCannotRegisterReplayOnPublicChannel()
- {
- var channel = PublicChannel();
-
- Assert.Throws(
- () => channel.Register(WithReplay()));
- }
-
- [TestMethod(DisplayName = "Channel: Registering broadcast replay on a private channel is allowed")]
- public void ClientCanRegisterReplayOnPrivateChannel()
- {
- var channel = PrivateChannel();
-
- var broadcast = channel.Register(WithReplay());
-
- Assert.IsNotNull(broadcast);
- }
-
- private static BroadcastOptions WithReplay() => new()
- {
- Replay = new BroadcastOptions.ReplayOptions { Limit = 10, Since = 0 }
- };
-
- private static RealtimeChannel PublicChannel() =>
- Channel(ChannelOptions.Public(ClientOptions(), () => null, new JsonSerializerSettings()));
-
- private static RealtimeChannel PrivateChannel() =>
- Channel(ChannelOptions.Private(ClientOptions(), () => null, new JsonSerializerSettings()));
-
- private static ClientOptions ClientOptions() => new();
-
- private static RealtimeChannel Channel(ChannelOptions options)
- {
- var socket = new RealtimeSocket("ws://localhost:54321/realtime/v1", options.ClientOptions);
-
- return new RealtimeChannel(socket, "realtime:online-users", options);
- }
-}
diff --git a/RealtimeTests/ChannelBroadcastTests.cs b/RealtimeTests/ChannelBroadcastTests.cs
deleted file mode 100644
index 6c7e9c1..0000000
--- a/RealtimeTests/ChannelBroadcastTests.cs
+++ /dev/null
@@ -1,229 +0,0 @@
-using System;
-using System.Collections.Generic;
-using System.Threading.Tasks;
-using Microsoft.VisualStudio.TestTools.UnitTesting;
-using Newtonsoft.Json;
-using RealtimeTests.Models;
-using Supabase.Postgrest.Interfaces;
-using Supabase.Realtime;
-using Supabase.Realtime.Broadcast;
-using Supabase.Realtime.Channel;
-using Supabase.Realtime.Interfaces;
-using Supabase.Realtime.Models;
-using Supabase.Realtime.PostgresChanges;
-using static Supabase.Realtime.PostgresChanges.PostgresChangesOptions;
-
-namespace RealtimeTests;
-
-public class BroadcastExample : BaseBroadcast
-{
- [JsonProperty("userId")]
- public string? UserId { get; set; }
-}
-
-[TestClass]
-public class ChannelBroadcastTests
-{
- private IPostgrestClient? _restClient;
- private IRealtimeClient? _socketClient;
-
- [TestInitialize]
- public async Task InitializeTest()
- {
- _restClient = Helpers.RestClient();
- _socketClient = Helpers.SocketClient();
- await _socketClient!.ConnectAsync();
- }
-
- [TestCleanup]
- public void CleanupTest()
- {
- _socketClient!.Disconnect();
- }
-
- [TestMethod("Channel: Can listen for broadcast")]
- public async Task ClientCanListenForBroadcast()
- {
- var tsc = new TaskCompletionSource();
- var tsc2 = new TaskCompletionSource();
-
- var guid1 = Guid.NewGuid().ToString();
- var guid2 = Guid.NewGuid().ToString();
-
- var channel1 = _socketClient!.Channel("online-users");
- var broadcast1 = channel1.Register(true, true);
- broadcast1.AddBroadcastEventHandler(
- (_, _) =>
- {
- var broadcast = broadcast1.Current();
- if (broadcast?.UserId != guid1 && broadcast?.Event == "user")
- tsc.TrySetResult(true);
- }
- );
-
- var client2 = Helpers.SocketClient();
- await client2.ConnectAsync();
- var channel2 = client2.Channel("online-users");
- var broadcast2 = channel2.Register(true, true);
- broadcast2.AddBroadcastEventHandler(
- (_, _) =>
- {
- var broadcast = broadcast2.Current();
- if (broadcast?.UserId != guid2 && broadcast?.Event == "user")
- tsc2.TrySetResult(true);
- }
- );
-
- await channel1.Subscribe();
- await channel2.Subscribe();
-
- await broadcast1.Send("user", new BroadcastExample { UserId = guid1 });
- await broadcast2.Send("user", new BroadcastExample { UserId = guid2 });
-
- await Task.WhenAll(new[] { tsc.Task, tsc2.Task });
- }
-
- [TestMethod("Channel: Can listen for broadcast on a private channel")]
- public async Task ClientCanListenForBroadcastPrivate()
- {
- var tsc = new TaskCompletionSource();
- var tsc2 = new TaskCompletionSource();
-
- var guid1 = Guid.NewGuid().ToString();
- var guid2 = Guid.NewGuid().ToString();
-
- var client1 = Helpers.PrivateSocketClient();
- await client1.ConnectAsync();
- var options1 = ChannelOptions.Private(
- client1.Options,
- () => Helpers.ApiKey,
- new JsonSerializerSettings()
- );
- var channel1 = client1.Channel("online-users", options1);
- var broadcast1 = channel1.Register(true, true);
- broadcast1.AddBroadcastEventHandler(
- (_, _) =>
- {
- var broadcast = broadcast1.Current();
- if (broadcast?.UserId != guid1 && broadcast?.Event == "user")
- tsc.TrySetResult(true);
- }
- );
-
- var client2 = Helpers.PrivateSocketClient();
- await client2.ConnectAsync();
- var options2 = ChannelOptions.Private(
- client2.Options,
- () => Helpers.ApiKey,
- new JsonSerializerSettings()
- );
- var channel2 = client2.Channel("online-users", options2);
- var broadcast2 = channel2.Register(true, true);
- broadcast2.AddBroadcastEventHandler(
- (_, _) =>
- {
- var broadcast = broadcast2.Current();
- if (broadcast?.UserId != guid2 && broadcast?.Event == "user")
- tsc2.TrySetResult(true);
- }
- );
-
- await channel1.Subscribe();
- await channel2.Subscribe();
-
- await broadcast1.Send("user", new BroadcastExample { UserId = guid1 });
- await broadcast2.Send("user", new BroadcastExample { UserId = guid2 });
-
- await Task.WhenAll(new[] { tsc.Task, tsc2.Task });
- }
-
- [TestMethod("Channel: Send resolves when broadcast ack is not explicitly enabled")]
- public async Task ChannelSendResolvesWithoutExplicitAck()
- {
- // Mirrors the most natural usage: `client.Channel(name)` -> `Subscribe()` -> `Send()`,
- // with no call to `Register(broadcastAck: true)`. See supabase-community/realtime-csharp#38.
- var channel = _socketClient!.Channel("no-ack-broadcast");
- await channel.Subscribe();
-
- var sendTask = channel.Send(Supabase.Realtime.Constants.ChannelEventName.Broadcast, "test_event",
- new BroadcastExample { UserId = Guid.NewGuid().ToString() });
-
- var winner = await Task.WhenAny(sendTask, Task.Delay(TimeSpan.FromSeconds(5)));
-
- Assert.AreSame(sendTask, winner, "channel.Send() did not complete within 5s (issue #38: awaited broadcast send hangs).");
- Assert.IsTrue(await sendTask);
- }
-
- [TestMethod("Channel: Can listen history for private broadcast")]
- public async Task ClientCanListenHistoryForBroadcastPrivate()
- {
- var send = new Dictionary
- {
- { "event", "user" },
- { "topic", "online-users" },
- { "private", true }
- };
- await _restClient!.Rpc("send", send);
-
- var tsc = new TaskCompletionSource();
-
- var client1 = Helpers.PrivateSocketClient();
- await client1.ConnectAsync();
- var options1 = ChannelOptions.Private(
- client1.Options,
- () => null,
- new JsonSerializerSettings()
- );
- var broadcastOptions = new BroadcastOptions
- {
- BroadcastAck = true,
- BroadcastSelf = true,
- Replay = new BroadcastOptions.ReplayOptions
- {
- Limit = 10,
- Since = DateTimeOffset.UtcNow.AddDays(-3).ToUnixTimeMilliseconds()
- }
- };
- var channel1 = client1.Channel("online-users", options1);
- var broadcast1 = channel1.Register(broadcastOptions);
- broadcast1.AddBroadcastEventHandler(
- (_, _) =>
- {
-
- var broadcast = broadcast1.Current();
- if (broadcast is { Event: "user", Meta.Replayed: true })
- tsc.TrySetResult(true);
- }
- );
-
- await channel1.Subscribe();
-
- await Task.WhenAll(tsc.Task);
- }
-
- [TestMethod("Channel: Payload returns a modeled response (if possible)")]
- public async Task ChannelPayloadReturnsModel()
- {
- var tsc = new TaskCompletionSource();
-
- var channel = _socketClient!.Channel("example");
- channel.Register(new PostgresChangesOptions("public", "*"));
- channel.AddPostgresChangeHandler(
- ListenType.Inserts,
- (_, changes) =>
- {
- var model = changes.Model();
- tsc.SetResult(model != null);
- }
- );
-
- await channel.Subscribe();
-
- await _restClient!
- .Table()
- .Insert(new Todo { UserId = 1, Details = "Client Models a response? ✅" });
-
- var check = await tsc.Task;
- Assert.IsTrue(check);
- }
-}
diff --git a/RealtimeTests/ChannelPostgresChangesTests.cs b/RealtimeTests/ChannelPostgresChangesTests.cs
deleted file mode 100644
index b889a30..0000000
--- a/RealtimeTests/ChannelPostgresChangesTests.cs
+++ /dev/null
@@ -1,364 +0,0 @@
-using System;
-using System.Collections.Generic;
-using System.Diagnostics;
-using System.Linq;
-using System.Threading.Tasks;
-using Microsoft.VisualStudio.TestTools.UnitTesting;
-using Supabase.Postgrest.Interfaces;
-using RealtimeTests.Models;
-using Supabase.Realtime;
-using Supabase.Realtime.Interfaces;
-using Supabase.Realtime.PostgresChanges;
-using static Supabase.Realtime.Constants;
-using static Supabase.Realtime.PostgresChanges.PostgresChangesOptions;
-
-namespace RealtimeTests;
-
-[TestClass]
-public class ChannelPostgresChangesTests
-{
- private IPostgrestClient? _restClient;
- private IRealtimeClient? _socketClient;
-
- [TestInitialize]
- public async Task InitializeTest()
- {
- _restClient = Helpers.RestClient();
- _socketClient = Helpers.SocketClient();
- await _socketClient!.ConnectAsync();
- }
-
- [TestCleanup]
- public void CleanupTest()
- {
- _socketClient!.Disconnect();
- }
-
- [TestMethod(DisplayName = "Channel: Payload returns a modeled response (if possible)")]
- public async Task ChannelPayloadReturnsModel()
- {
- var tsc = new TaskCompletionSource();
-
- var channel = _socketClient!.Channel("example");
- channel.OnPostgresChange((_, changes) =>
- {
- var model = changes.Model();
- tsc.SetResult(model != null);
- }, ListenType.Inserts, new PostgresChangesFilter { Table = "*" });
-
- await channel.Subscribe();
-
- await _restClient!.Table().Insert(new Todo { UserId = 1, Details = "Client Models a response? ✅" });
-
- var check = await tsc.Task;
- Assert.IsTrue(check);
- }
-
- [TestMethod(DisplayName = "Channel: Receives Insert Callback")]
- public async Task ChannelReceivesInsertCallback()
- {
- var tsc = new TaskCompletionSource();
-
- var channel = _socketClient!.Channel("realtime", "public", "todos");
-
- channel.OnPostgresChange((_, _) => tsc.SetResult(true), ListenType.Inserts,
- new PostgresChangesFilter { Table = "todos" });
-
- await channel.Subscribe();
- await _restClient!.Table()
- .Insert(new Todo { UserId = 1, Details = "Client receives insert callback? ✅" });
-
- var check = await tsc.Task;
- Assert.IsTrue(check);
- }
-
- [TestMethod(DisplayName = "Channel: Receives Filtered Insert Callback")]
- public async Task ChannelReceivesInsertCallbackFiltered()
- {
- var tsc = new TaskCompletionSource();
-
- var channel = _socketClient!.Channel("realtime", "public", "todos");
-
- channel.OnPostgresChange((_, changes) =>
- {
- var oldModel = changes.Model();
-
- Assert.AreEqual("Client receives filtered insert callback? ✅", oldModel?.Details);
-
- tsc.SetResult(true);
- }, ListenType.Inserts,
- new PostgresChangesFilter { Table = "todos", Filter = "details=eq.Client receives filtered insert callback? ✅" });
-
- await channel.Subscribe();
- await _restClient!.Table()
- .Insert(new Todo { UserId = 1, Details = "Client receives insert callback? ✅" });
-
- await _restClient!.Table()
- .Insert(new Todo { UserId = 2, Details = "Client receives filtered insert callback? ✅" });
-
- var check = await tsc.Task;
- Assert.IsTrue(check);
- }
-
- [TestMethod(DisplayName = "Channel: Receives Filtered Two Callback")]
- public async Task ChannelReceivesTwoCallbacks()
- {
- var tsc = new TaskCompletionSource();
-
- var response = await _restClient!.Table()
- .Insert(new Todo { UserId = 1, Details = "Client receives insert callback? ✅" });
- await _restClient!.Table()
- .Insert(new Todo { UserId = 2, Details = "Client receives filtered insert callback? ✅" });
-
- var model = response.Models.First();
- var oldDetails = model.Details;
- var newDetails = $"I'm an updated item ✏️ - {DateTime.Now}";
-
- var channel = _socketClient!.Channel("realtime", "public", "todos");
- channel.OnPostgresChange((_, changes) =>
- {
- var oldModel = changes.OldModel();
-
- Assert.AreEqual(oldDetails, oldModel?.Details);
-
- var updated = changes.Model();
- Assert.AreEqual(newDetails, updated?.Details);
-
- if (updated != null)
- {
- Assert.AreEqual(model.Id, updated.Id);
- Assert.AreEqual(model.UserId, updated.UserId);
- }
-
- tsc.SetResult(true);
- }, ListenType.Updates, new PostgresChangesFilter { Table = "todos" });
-
- const string filter = "Client receives filtered insert callback? ✅";
- channel.OnPostgresChange((_, changes) =>
- {
- var insertedModel = changes.Model();
-
- Assert.AreEqual("Client receives filtered insert callback? ✅", insertedModel?.Details);
-
- tsc.SetResult(true);
- }, ListenType.Inserts, new PostgresChangesFilter { Table = "todos", Filter = $"details=eq.{filter}" });
-
- await channel.Subscribe();
-
- await _restClient.Table()
- .Set(x => x.Details!, newDetails)
- .Match(model)
- .Update();
-
- var check = await tsc.Task;
- Assert.IsTrue(check);
- }
-
- [TestMethod(DisplayName = "Channel: Receives Update Callback")]
- public async Task ChannelReceivesUpdateCallback()
- {
- var tsc = new TaskCompletionSource();
-
- var response = await _restClient!.Table()
- .Insert(new Todo { UserId = 1, Details = "Client receives insert callback? ✅" });
-
- var model = response.Models.First();
- var oldDetails = model.Details;
- var newDetails = $"I'm an updated item ✏️ - {DateTime.Now}";
-
- var channel = _socketClient!.Channel("realtime", "public", "todos");
-
- channel.OnPostgresChange((_, changes) =>
- {
- var oldModel = changes.OldModel();
-
- Assert.AreEqual(oldDetails, oldModel?.Details);
-
- var updated = changes.Model();
- Assert.AreEqual(newDetails, updated?.Details);
-
- if (updated != null)
- {
- Assert.AreEqual(model.Id, updated.Id);
- Assert.AreEqual(model.UserId, updated.UserId);
- }
-
- tsc.SetResult(true);
- }, ListenType.Updates, new PostgresChangesFilter { Table = "todos" });
-
- await channel.Subscribe();
-
- await _restClient.Table()
- .Set(x => x.Details!, newDetails)
- .Match(model)
- .Update();
-
- var check = await tsc.Task;
- Assert.IsTrue(check);
- }
-
- [TestMethod(DisplayName = "Channel: Receives Delete Callback")]
- public async Task ChannelReceivesDeleteCallback()
- {
- var tsc = new TaskCompletionSource();
-
- var channel = _socketClient!.Channel("realtime", "public", "todos");
-
- channel.OnPostgresChange((_, _) => tsc.SetResult(true), ListenType.Deletes,
- new PostgresChangesFilter { Table = "todos" });
-
- await channel.Subscribe();
-
- var result = await _restClient!.Table().Get();
- var model = result.Models.Last();
-
- await _restClient.Table().Match(model).Delete();
-
- var check = await tsc.Task;
- Assert.IsTrue(check);
- }
-
- [TestMethod(DisplayName = "Channel: Receives Delete Callback")]
- public async Task ChannelReceivesFilteredDeleteCallback()
- {
- var tsc = new TaskCompletionSource();
- var channel = _socketClient!.Channel("realtime", "public", "todos");
-
- var todo1 = await _restClient!.Table().Insert(new Todo
- { UserId = 1, Details = "Client receives callbacks 1? ✅" });
- var todo2 = await _restClient!.Table().Insert(new Todo
- { UserId = 2, Details = "Client receives callbacks 2? ✅" });
- await _restClient!.Table().Insert(new Todo
- { UserId = 3, Details = "Client receives callbacks 3? ✅" });
-
- channel.OnPostgresChange((_, removed) =>
- {
- var result = removed.OldModel();
- Assert.AreEqual(result?.Details, todo1.Model?.Details);
- Assert.AreNotEqual(result?.Details, todo2.Model?.Details);
-
- tsc.SetResult(true);
- }, ListenType.Deletes,
- new PostgresChangesFilter { Table = "todos", Filter = $"details=eq.{todo1.Model?.Details}" });
-
- await channel.Subscribe();
-
- await _restClient.Table().Match(todo1.Models.First()).Delete();
- await _restClient.Table().Match(todo2.Models.First()).Delete();
-
- var check = await tsc.Task;
- Assert.IsTrue(check);
- }
-
- [TestMethod(DisplayName = "Channel: Receives '*' Callback")]
- public async Task ChannelReceivesWildcardCallback()
- {
- var insertTsc = new TaskCompletionSource();
- var updateTsc = new TaskCompletionSource();
- var deleteTsc = new TaskCompletionSource();
-
- List tasks = new List { insertTsc.Task, updateTsc.Task, deleteTsc.Task };
-
- var channel = _socketClient!.Channel("realtime", "public", "todos");
-
- channel.OnPostgresChange((_, changes) =>
- {
- switch (changes.Payload?.Data?.Type)
- {
- case EventType.Insert:
- insertTsc.SetResult(true);
- break;
- case EventType.Update:
- updateTsc.SetResult(true);
- break;
- case EventType.Delete:
- deleteTsc.SetResult(true);
- break;
- }
- }, ListenType.All, new PostgresChangesFilter { Table = "todos" });
-
- await channel.Subscribe();
-
- var modeledResponse = await _restClient!.Table().Insert(new Todo
- { UserId = 1, Details = "Client receives wildcard callbacks? ✅" });
- var newModel = modeledResponse.Models.First();
-
- await _restClient.Table().Set(x => x.Details!, "And edits.").Match(newModel).Update();
- await _restClient.Table().Match(newModel).Delete();
-
- await Task.WhenAll(tasks);
-
- Assert.IsTrue(insertTsc.Task.Result);
- Assert.IsTrue(updateTsc.Task.Result);
- Assert.IsTrue(deleteTsc.Task.Result);
- }
-
- [TestMethod(DisplayName = "Channel: Receives Several Same Callback")]
- public async Task ChannelReceivesSeveralSameCallback()
- {
- var insertTask1 = new TaskCompletionSource();
- var insertTask2 = new TaskCompletionSource();
- var insertTask3 = new TaskCompletionSource();
- const string filter1 = "Client receives callbacks 1? ✅";
- const string filter2 = "Client receives callbacks 2? ✅";
-
- var channel = _socketClient!.Channel("realtime", "public", "todos");
-
- var count = 0;
- channel.OnPostgresChange((_, added) =>
- {
- count++;
- if (count == 3) insertTask1.TrySetResult(true);
- }, ListenType.Inserts, new PostgresChangesFilter { Table = "todos" });
-
- channel.OnPostgresChange((_, added) =>
- {
- var model = added.Model();
-
- insertTask2.SetResult(model?.Details == filter1);
- }, ListenType.Inserts, new PostgresChangesFilter { Table = "todos", Filter = $"details=eq.{filter1}" });
-
- channel.OnPostgresChange((_, added) =>
- {
- var model = added.Model();
-
- insertTask3.SetResult(model?.Details == filter2);
- }, ListenType.Inserts, new PostgresChangesFilter { Table = "todos", Filter = $"details=eq.{filter2}" });
-
- await channel.Subscribe();
-
- await _restClient!.Table().Insert(new Todo { UserId = 1, Details = "Client receives wildcard callbacks? ✅" });
- await _restClient!.Table().Insert(new Todo { UserId = 1, Details = filter1 });
- await _restClient!.Table().Insert(new Todo { UserId = 1, Details = filter2 });
-
- await Task.WhenAll(insertTask1.Task, insertTask2.Task, insertTask3.Task);
-
- Assert.IsTrue(insertTask1.Task.Result);
- Assert.IsTrue(insertTask2.Task.Result);
- Assert.IsTrue(insertTask3.Task.Result);
- }
-
- [TestMethod]
- public async Task OnPostgresChange_ShouldRegisterAndDeliverToHandler_GivenChainedSubscribe()
- {
- var tsc = new TaskCompletionSource();
-
- await _socketClient!.Channel("public:todos")
- .OnPostgresChange((_, changes) => tsc.TrySetResult(changes.Model() != null),
- ListenType.Inserts, new PostgresChangesFilter { Table = "todos" })
- .Subscribe();
-
- await _restClient!.Table()
- .Insert(new Todo { UserId = 1, Details = "OnPostgresChange receives insert? ✅" });
-
- Assert.IsTrue(await WithinTimeout(tsc.Task),
- "OnPostgresChange should register the option and bind the handler in one call, and return the channel so Subscribe can be chained");
- }
-
- private static async Task WithinTimeout(Task task, int timeoutMs = 15000)
- {
- var completed = await Task.WhenAny(task, Task.Delay(timeoutMs));
- return completed == task && task.Result;
- }
-
-}
diff --git a/RealtimeTests/ChannelTests.cs b/RealtimeTests/ChannelTests.cs
deleted file mode 100644
index ec4a44d..0000000
--- a/RealtimeTests/ChannelTests.cs
+++ /dev/null
@@ -1,135 +0,0 @@
-using System.Collections.Generic;
-using System.Threading.Tasks;
-using Microsoft.VisualStudio.TestTools.UnitTesting;
-using Newtonsoft.Json;
-using Supabase.Postgrest.Interfaces;
-using RealtimeTests.Models;
-using Supabase.Realtime;
-using Supabase.Realtime.Interfaces;
-using Supabase.Realtime.PostgresChanges;
-using static Supabase.Realtime.Constants;
-using static Supabase.Realtime.PostgresChanges.PostgresChangesOptions;
-
-namespace RealtimeTests;
-
-[TestClass]
-public class ChannelTests
-{
- private IPostgrestClient? _restClient;
- private IRealtimeClient? _socketClient;
-
- [TestInitialize]
- public async Task InitializeTest()
- {
- _restClient = Helpers.RestClient();
- _socketClient = Helpers.SocketClient();
- await _socketClient!.ConnectAsync();
- }
-
- [TestCleanup]
- public void CleanupTest()
- {
- _socketClient!.Disconnect();
- }
-
- [TestMethod("Channel: Close Event Handler")]
- public async Task ChannelCloseEventHandler()
- {
- var tsc = new TaskCompletionSource();
-
- var channel = _socketClient!.Channel("realtime", "public", "todos");
- channel.AddStateChangedHandler((_, state) =>
- {
- if (state == ChannelState.Closed)
- tsc.SetResult(true);
- });
-
- await channel.Subscribe();
- channel.Unsubscribe();
-
- var check = await tsc.Task;
- Assert.IsTrue(check);
- }
-
- [TestMethod("Channel: Supports WALRUS Array Changes")]
- public async Task ChannelSupportsWalrusArray()
- {
- Todo? result = null;
- var tsc = new TaskCompletionSource();
-
- var channel = _socketClient!.Channel("realtime", "public", "todos");
- var numbers = new List { 4, 5, 6 };
-
- await channel.Subscribe();
-
- channel.AddPostgresChangeHandler(ListenType.Inserts, (_, changes) =>
- {
- result = changes.Model();
- tsc.SetResult(true);
- });
-
- await _restClient!.Table().Insert(new Todo { UserId = 1, Numbers = numbers });
-
- await tsc.Task;
- CollectionAssert.AreEqual(numbers, result?.Numbers);
- }
-
- [TestMethod("Channel: Sends Join parameters")]
- public async Task ChannelSendsJoinParameters()
- {
- var parameters = new Dictionary { { "key", "value" } };
- var channel = _socketClient!.Channel("realtime", "public", "todos", parameters: parameters);
-
- await channel.Subscribe();
-
- var serialized = JsonConvert.SerializeObject(channel.JoinPush?.Payload);
- Assert.IsTrue(serialized.Contains("\"key\":\"value\""));
- }
-
- [TestMethod("Channel: Returns single subscription per unique topic.")]
- public async Task ChannelJoinsDuplicateSubscription()
- {
- var subscription1 = _socketClient!.Channel("realtime", "public", "todos");
- var subscription2 = _socketClient!.Channel("realtime", "public", "todos");
- var subscription3 = _socketClient!.Channel("realtime", "public", "todos", "user_id", "1");
-
- Assert.AreEqual(subscription1.Topic, subscription2.Topic);
-
- await subscription1.Subscribe();
-
- Assert.AreEqual(subscription1.HasJoinedOnce, subscription2.HasJoinedOnce);
- Assert.AreNotEqual(subscription1.HasJoinedOnce, subscription3.HasJoinedOnce);
-
- var subscription4 = _socketClient!.Channel("realtime", "public", "todos");
-
- Assert.AreEqual(subscription1.HasJoinedOnce, subscription4.HasJoinedOnce);
- }
-
- [TestMethod("Channel: Registers Handlers")]
- public async Task ChannelRegistersHandlers()
- {
- var channel = _socketClient!.Channel("test");
-
- IRealtimeChannel.StateChangedHandler stateHandler = (_, _) => Assert.Fail("State Handler was called");
- IRealtimeChannel.MessageReceivedHandler messageReceivedHandler =
- (_, _) => Assert.Fail("Message Handler was called");
- IRealtimeChannel.PostgresChangesHandler postgresChangesHandler =
- (_, _) => Assert.Fail("Postgres Changes Handler was called");
-
- channel.AddStateChangedHandler(stateHandler);
- channel.AddMessageReceivedHandler(messageReceivedHandler);
- channel.AddPostgresChangeHandler(ListenType.All, postgresChangesHandler);
-
- channel.Register(new PostgresChangesOptions("public", "todos"));
- channel.Register();
- channel.Register("user");
-
- channel.RemoveStateChangedHandler(stateHandler);
- channel.RemoveMessageReceivedHandler(messageReceivedHandler);
- channel.RemovePostgresChangeHandler(ListenType.All, postgresChangesHandler);
-
- await channel.Subscribe();
-
- await Task.Delay(500);
- }
-}
\ No newline at end of file
diff --git a/RealtimeTests/Channels/ChannelLifecycleTests.cs b/RealtimeTests/Channels/ChannelLifecycleTests.cs
new file mode 100644
index 0000000..83193da
--- /dev/null
+++ b/RealtimeTests/Channels/ChannelLifecycleTests.cs
@@ -0,0 +1,37 @@
+using System.Collections.Generic;
+using FluentAssertions;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using RealtimeTests.Support;
+using Supabase.Realtime.Exceptions;
+using static Supabase.Realtime.Constants;
+
+namespace RealtimeTests.Channels;
+
+///
+/// A channel's join lifecycle over an unconnected socket: pushing before a join is refused, and
+/// unsubscribing walks state through Leaving to Closed. The join handshake itself needs a
+/// server reply and is covered in the E2E tier.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class ChannelLifecycleTests
+{
+ [TestMethod]
+ public void Push_ShouldThrow_GivenNotJoined()
+ {
+ var channel = Wire.Channel();
+ var act = () => channel.Push("some_event");
+ act.Should().Throw().Which.Reason.Should().Be(FailureHint.Reason.ChannelNotOpen);
+ }
+
+ [TestMethod]
+ public void Unsubscribe_ShouldWalkStateToClosed()
+ {
+ var channel = Wire.Channel();
+ var states = new List();
+ channel.AddStateChangedHandler((_, state) => states.Add(state));
+ channel.Unsubscribe();
+ states.Should().ContainInOrder(ChannelState.Leaving, ChannelState.Closed);
+ channel.IsClosed.Should().BeTrue();
+ }
+}
diff --git a/RealtimeTests/Channels/ChannelMessagingTests.cs b/RealtimeTests/Channels/ChannelMessagingTests.cs
new file mode 100644
index 0000000..4226c82
--- /dev/null
+++ b/RealtimeTests/Channels/ChannelMessagingTests.cs
@@ -0,0 +1,115 @@
+using System.Collections.Generic;
+using System.Threading.Tasks;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using Newtonsoft.Json;
+using RealtimeTests.Models;
+using Supabase.Postgrest.Interfaces;
+using Supabase.Realtime;
+using Supabase.Realtime.Interfaces;
+using Supabase.Realtime.PostgresChanges;
+using static Supabase.Realtime.Constants;
+using static Supabase.Realtime.PostgresChanges.PostgresChangesOptions;
+
+namespace RealtimeTests.Channels;
+
+///
+/// A subscribed channel's messaging behavior against the live stack: close notification on unsubscribe,
+/// WALRUS array column delivery, join parameters on the wire, shared join state across duplicate topics,
+/// and handler registration/removal.
+///
+[TestClass]
+[TestCategory("E2E")]
+public class ChannelMessagingTests
+{
+ private IPostgrestClient restClient = null!;
+ private IRealtimeClient socketClient = null!;
+
+ [TestInitialize]
+ public async Task InitializeTest()
+ {
+ restClient = Helpers.RestClient();
+ socketClient = Helpers.SocketClient();
+ await socketClient.ConnectAsync();
+ }
+
+ [TestCleanup]
+ public void CleanupTest() => socketClient.Disconnect();
+
+ [TestMethod]
+ public async Task Unsubscribe_ShouldRaiseClosedState()
+ {
+ var tsc = new TaskCompletionSource();
+ var channel = socketClient.Channel("realtime", "public", "todos");
+ channel.AddStateChangedHandler((_, state) =>
+ {
+ if (state == ChannelState.Closed)
+ tsc.SetResult(true);
+ });
+ await channel.Subscribe();
+ channel.Unsubscribe();
+ Assert.IsTrue(await tsc.Task);
+ }
+
+ [TestMethod]
+ public async Task Subscribe_ShouldDeliverArrayColumn_GivenInsert()
+ {
+ Todo? result = null;
+ var tsc = new TaskCompletionSource();
+ var channel = socketClient.Channel("realtime", "public", "todos");
+ var numbers = new List { 4, 5, 6 };
+ channel.OnPostgresChange((_, changes) =>
+ {
+ result = changes.Model();
+ tsc.SetResult(true);
+ }, ListenType.Inserts, new PostgresChangesFilter { Table = "todos" });
+ await channel.Subscribe();
+ await restClient.Table().Insert(new Todo { UserId = 1, Numbers = numbers });
+ await tsc.Task;
+ CollectionAssert.AreEqual(numbers, result?.Numbers);
+ }
+
+ [TestMethod]
+ public async Task Subscribe_ShouldSendJoinParameters()
+ {
+ var parameters = new Dictionary { { "key", "value" } };
+ var channel = socketClient.Channel("realtime", "public", "todos", parameters: parameters);
+ await channel.Subscribe();
+ var serialized = JsonConvert.SerializeObject(channel.JoinPush?.Payload);
+ Assert.IsTrue(serialized.Contains("\"key\":\"value\""));
+ }
+
+ [TestMethod]
+ public async Task Channel_ShouldShareJoinState_GivenDuplicateTopics()
+ {
+ var first = socketClient.Channel("realtime", "public", "todos");
+ var second = socketClient.Channel("realtime", "public", "todos");
+ var filtered = socketClient.Channel("realtime", "public", "todos", "user_id", "1");
+ Assert.AreEqual(first.Topic, second.Topic);
+ await first.Subscribe();
+ Assert.AreEqual(first.HasJoinedOnce, second.HasJoinedOnce);
+ Assert.AreNotEqual(first.HasJoinedOnce, filtered.HasJoinedOnce);
+ var fourth = socketClient.Channel("realtime", "public", "todos");
+ Assert.AreEqual(first.HasJoinedOnce, fourth.HasJoinedOnce);
+ }
+
+ [TestMethod]
+ public async Task Channel_ShouldRegisterAndRemoveHandlers()
+ {
+ var channel = socketClient.Channel("test");
+ IRealtimeChannel.StateChangedHandler stateHandler = (_, _) => Assert.Fail("State Handler was called");
+ IRealtimeChannel.MessageReceivedHandler messageReceivedHandler =
+ (_, _) => Assert.Fail("Message Handler was called");
+ IRealtimeChannel.PostgresChangesHandler postgresChangesHandler =
+ (_, _) => Assert.Fail("Postgres Changes Handler was called");
+ channel.AddStateChangedHandler(stateHandler);
+ channel.AddMessageReceivedHandler(messageReceivedHandler);
+ channel.OnPostgresChange(postgresChangesHandler, ListenType.All, new PostgresChangesFilter { Table = "todos" });
+ channel.Register();
+ channel.Register("user");
+ channel.RemoveStateChangedHandler(stateHandler);
+ channel.RemoveMessageReceivedHandler(messageReceivedHandler);
+ channel.RemovePostgresChangeHandler(ListenType.All, postgresChangesHandler);
+ await channel.Subscribe();
+ await Task.Delay(500);
+ }
+}
diff --git a/RealtimeTests/Channels/ChannelOptionsTests.cs b/RealtimeTests/Channels/ChannelOptionsTests.cs
new file mode 100644
index 0000000..9dde109
--- /dev/null
+++ b/RealtimeTests/Channels/ChannelOptionsTests.cs
@@ -0,0 +1,44 @@
+using FluentAssertions;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using Newtonsoft.Json;
+using Supabase.Realtime;
+using Supabase.Realtime.Channel;
+
+namespace RealtimeTests.Channels;
+
+///
+/// The public/private distinction a channel carries: marks the
+/// channel as authorized against Row Level Security (required for broadcast replay), and the access-token
+/// accessor is surfaced back to callers.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class ChannelOptionsTests
+{
+ private static readonly JsonSerializerSettings Settings = new();
+
+ [TestMethod]
+ public void Public_ShouldNotBePrivate()
+ {
+ ChannelOptions.Public(new ClientOptions(), () => null, Settings).IsPrivate.Should().BeFalse();
+ }
+
+ [TestMethod]
+ public void Private_ShouldBePrivate()
+ {
+ ChannelOptions.Private(new ClientOptions(), () => null, Settings).IsPrivate.Should().BeTrue();
+ }
+
+ [TestMethod]
+ public void Constructor_ShouldDefaultToPublic()
+ {
+ new ChannelOptions(new ClientOptions(), () => null, Settings).IsPrivate.Should().BeFalse();
+ }
+
+ [TestMethod]
+ public void RetrieveAccessToken_ShouldReturnConfiguredToken()
+ {
+ var options = ChannelOptions.Public(new ClientOptions(), () => "jwt-123", Settings);
+ options.RetrieveAccessToken().Should().Be("jwt-123");
+ }
+}
diff --git a/RealtimeTests/Channels/ChannelSubscriptionTests.cs b/RealtimeTests/Channels/ChannelSubscriptionTests.cs
new file mode 100644
index 0000000..b32066e
--- /dev/null
+++ b/RealtimeTests/Channels/ChannelSubscriptionTests.cs
@@ -0,0 +1,107 @@
+using System.Threading.Tasks;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using Supabase.Realtime;
+using Supabase.Realtime.Exceptions;
+using static Supabase.Realtime.Constants;
+
+namespace RealtimeTests.Channels;
+
+///
+/// How the client resolves channels from a live connection: the topic each overload produces, that the
+/// same topic returns one shared instance, that an already-prefixed name is not double-prefixed, removal,
+/// and pushing the access token on SetAuth.
+///
+[TestClass]
+[TestCategory("E2E")]
+public class ChannelSubscriptionTests
+{
+ private Client client = null!;
+
+ [TestInitialize]
+ public async Task InitializeTest()
+ {
+ client = Helpers.SocketClient();
+ await client.ConnectAsync();
+ }
+
+ [TestCleanup]
+ public void CleanupTest() => client.Disconnect();
+
+ [TestMethod]
+ public async Task Channel_ShouldTopicByTable()
+ {
+ var channel = client.Channel(table: "todos");
+ await channel.Subscribe();
+ Assert.AreEqual("realtime:public:todos", channel.Topic);
+ }
+
+ [TestMethod]
+ public async Task Channel_ShouldTopicBySchemaWildcard()
+ {
+ var channel = client.Channel("realtime", "public", "*");
+ await channel.Subscribe();
+ Assert.AreEqual("realtime:public:*", channel.Topic);
+ }
+
+ [TestMethod]
+ public async Task Channel_ShouldThrowThenJoin_GivenUnpublishedThenPublishedTable()
+ {
+ var users = client.Channel("realtime", "public", "users");
+ await Assert.ThrowsAsync(() => users.Subscribe());
+ var todos = client.Channel("realtime", "public", "todos");
+ await todos.Subscribe();
+ Assert.AreEqual("realtime:public:todos", todos.Topic);
+ }
+
+ [TestMethod]
+ public async Task Channel_ShouldTopicByColumnFilter()
+ {
+ var channel = client.Channel("realtime", "public", "todos", "id", "1");
+ await channel.Subscribe();
+ Assert.AreEqual("realtime:public:todos:id=eq.1", channel.Topic);
+ }
+
+ [TestMethod]
+ public async Task Channel_ShouldReturnSameInstance_GivenSameTopic()
+ {
+ var first = client.Channel("realtime", "public", "todos");
+ await first.Subscribe();
+ var second = client.Channel("realtime", "public", "todos");
+ Assert.AreEqual(true, second.IsJoined);
+ }
+
+ [TestMethod]
+ public void Channel_ShouldNotDoublePrefixTopic_GivenNameAlreadyPrefixed()
+ {
+ var prefixed = client.Channel("realtime:public:todos");
+ Assert.AreEqual("realtime:public:todos", prefixed.Topic);
+ Assert.AreSame(prefixed, client.Channel("public:todos"));
+ }
+
+ [TestMethod]
+ public async Task Remove_ShouldClearStoredSubscription()
+ {
+ var channel = client.Channel("realtime", "public", "todos");
+ await channel.Subscribe();
+ client.Remove(channel);
+ var reopened = client.Channel("realtime", "public", "todos");
+ Assert.AreEqual(ChannelState.Closed, reopened.State);
+ }
+
+ [TestMethod]
+ public async Task SetAuth_ShouldPushAccessTokenToJoinedChannelsOnly()
+ {
+ var first = client.Channel("realtime", "public", "todos");
+ var second = client.Channel("realtime", "public", "todos");
+ const string token =
+ "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJzdWIiOiIxMjM0NTY3ODkwIiwibmFtZSI6IkpvaG4gRG9lIiwiaWF0IjoxNTE2MjM5MDIyfQ.C8oVtF5DICct_4HcdSKt8pdrxBFMQOAnPpbiiUbaXAY";
+ client.SetAuth(token);
+ foreach (var subscription in client.Subscriptions.Values)
+ Assert.IsNull(subscription.LastPush);
+ await first.Subscribe();
+ await second.Subscribe();
+ client.SetAuth(token);
+ foreach (var subscription in client.Subscriptions.Values)
+ Assert.IsTrue(subscription.LastPush?.EventName == ChannelAccessToken);
+ }
+}
diff --git a/RealtimeTests/Channels/ChannelTopicTests.cs b/RealtimeTests/Channels/ChannelTopicTests.cs
new file mode 100644
index 0000000..88cd0aa
--- /dev/null
+++ b/RealtimeTests/Channels/ChannelTopicTests.cs
@@ -0,0 +1,55 @@
+using System.Collections.Generic;
+using FluentAssertions;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using Supabase.Realtime;
+
+namespace RealtimeTests.Channels;
+
+///
+/// The topic string a channel is keyed under: joins the
+/// database/schema/table segments and appends a col=eq.val row filter, dropping empty segments.
+/// Also covers the connection query string builder.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class ChannelTopicTests
+{
+ [TestMethod]
+ public void GenerateChannelTopic_ShouldJoinSegments()
+ {
+ Utils.GenerateChannelTopic("realtime", "public", "todos", null, null)
+ .Should().Be("realtime:public:todos");
+ }
+
+ [TestMethod]
+ public void GenerateChannelTopic_ShouldAppendRowFilter_GivenColumnAndValue()
+ {
+ Utils.GenerateChannelTopic("realtime", "public", "todos", "id", "1")
+ .Should().Be("realtime:public:todos:id=eq.1");
+ }
+
+ [TestMethod]
+ public void GenerateChannelTopic_ShouldDropEmptySegments()
+ {
+ Utils.GenerateChannelTopic("realtime", "", "", null, null).Should().Be("realtime");
+ }
+
+ [TestMethod]
+ public void GenerateChannelTopic_ShouldIgnoreRowFilter_GivenMissingValue()
+ {
+ Utils.GenerateChannelTopic("realtime", "public", "todos", "id", null)
+ .Should().Be("realtime:public:todos");
+ }
+
+ [TestMethod]
+ public void QueryString_ShouldEmitProvidedPairsAndSkipEmpty()
+ {
+ var query = Utils.QueryString(new Dictionary
+ {
+ { "token", null },
+ { "apikey", "anon-key" },
+ { "vsn", "1.0.0" }
+ });
+ query.Should().Be("apikey=anon-key&vsn=1.0.0");
+ }
+}
diff --git a/RealtimeTests/ClientFailureTests.cs b/RealtimeTests/ClientFailureTests.cs
deleted file mode 100644
index 0c1e822..0000000
--- a/RealtimeTests/ClientFailureTests.cs
+++ /dev/null
@@ -1,47 +0,0 @@
-using System.Diagnostics;
-using System.Threading.Tasks;
-using Microsoft.VisualStudio.TestTools.UnitTesting;
-using Supabase.Realtime;
-using Supabase.Realtime.Exceptions;
-
-namespace RealtimeTests;
-
-[TestClass]
-public class ClientFailureTests
-{
- [TestMethod("Client throws exception when unable to initially connect.")]
- public async Task ClientThrowsExceptionOnInitialConnectionFailure()
- {
- var client = new Client("ws://localhost");
- var client2 = new Client("ws://localhost");
- client.AddDebugHandler((_, message, _) => Debug.WriteLine($"Client 1: {message}"));
- client2.AddDebugHandler((_, message, _) => Debug.WriteLine($"Client 2: {message}"));
-
- await Assert.ThrowsAsync(async () => { await client.ConnectAsync(); });
-
- await Assert.ThrowsAsync(() =>
- {
- var tsc = new TaskCompletionSource();
-
- client2.Connect((_, exception) =>
- {
- if (exception != null)
- tsc.SetException(exception);
- });
- return tsc.Task;
- });
- }
-
- [TestMethod("Client: Allows for multiple connection attempts after failure.")]
- public async Task ClientShouldAllowForMultipleSocketConnectionAttempts()
- {
- var client = new Client("ws://localhost");
- client.AddDebugHandler((_, message, _) => Debug.WriteLine(message));
-
- // Should throw first time and clear socket instance.
- await Assert.ThrowsAsync(client.ConnectAsync);
-
- // Should throw again, as socket is still cleared (as opposed to merely logging).
- await Assert.ThrowsAsync(client.ConnectAsync);
- }
-}
\ No newline at end of file
diff --git a/RealtimeTests/ClientTests.cs b/RealtimeTests/ClientTests.cs
deleted file mode 100644
index 0b734aa..0000000
--- a/RealtimeTests/ClientTests.cs
+++ /dev/null
@@ -1,152 +0,0 @@
-using System;
-using System.Collections.Generic;
-using System.Diagnostics;
-using System.Threading.Tasks;
-using Microsoft.VisualStudio.TestTools.UnitTesting;
-using Supabase.Realtime.Exceptions;
-using static Supabase.Realtime.Constants;
-
-namespace RealtimeTests;
-
-[TestClass]
-public class ClientTests
-{
- private Supabase.Realtime.Client? client;
-
- [TestInitialize]
- public async Task InitializeTest()
- {
- client = Helpers.SocketClient();
- client.AddDebugHandler((sender, message, exception) => Debug.WriteLine(message));
-
- await client!.ConnectAsync();
- }
-
- [TestCleanup]
- public void CleanupTest()
- {
- client?.Disconnect();
- }
-
-
- [TestMethod("Client: Join channels of format: {database}")]
- public async Task ClientJoinsChannel_DB()
- {
- var channel = client!.Channel(table: "todos");
- await channel.Subscribe();
-
- Assert.AreEqual("realtime:public:todos", channel.Topic);
- }
-
- [TestMethod("Client: Join channels of format: {database}:{schema}:*")]
- public async Task ClientJoinsChannel_DB_Schema()
- {
- var channel = client!.Channel("realtime", "public", "*");
- await channel.Subscribe();
-
- Assert.AreEqual("realtime:public:*", channel.Topic);
- }
-
- [TestMethod("Client: Join channels of format: {database}:{schema}:{table}")]
- public async Task ClientJoinsChannel_DB_Schema_Table()
- {
- var channel = client!.Channel("realtime", "public", "users");
- await Assert.ThrowsAsync(() => channel.Subscribe());
-
- var channel2 = client!.Channel("realtime", "public", "todos");
- await channel2.Subscribe();
-
- Assert.AreEqual("realtime:public:todos", channel2.Topic);
- }
-
- [TestMethod("Client: Join channels of format: {database}:{schema}:{table}:{col}=eq.{val}")]
- public async Task ClientJoinsChannel_DB_Schema_Table_Query()
- {
- var channel = client!.Channel("realtime", "public", "todos", "id", "1");
- await channel.Subscribe();
-
- Assert.AreEqual("realtime:public:todos:id=eq.1", channel.Topic);
- }
-
- [TestMethod("Client: Returns a single instance of a channel based on topic")]
- public async Task ClientReturnsSingleChannelInstance()
- {
- var channel1 = client!.Channel("realtime", "public", "todos");
-
- await channel1.Subscribe();
-
- // Client should return an instance of `realtime:public:todos` that is already joined.
- var channel2 = client!.Channel("realtime", "public", "todos");
-
- Assert.AreEqual(true, channel2.IsJoined);
- }
-
- [TestMethod]
- public void Channel_ShouldNotDoublePrefixTopic_GivenNameAlreadyPrefixed()
- {
- var prefixed = client!.Channel("realtime:public:todos");
- Assert.AreEqual("realtime:public:todos", prefixed.Topic);
- Assert.AreSame(prefixed, client!.Channel("public:todos"));
- }
-
- [TestMethod("Client: Removes Channel Subscriptions")]
- public async Task ClientCanRemoveChannelSubscription()
- {
- var channel1 = client!.Channel("realtime", "public", "todos");
- await channel1.Subscribe();
-
- // Removing channel should remove the stored instance, so a future instance would need
- // to resubscribe.
- client!.Remove(channel1);
-
- var channel2 = client!.Channel("realtime", "public", "todos");
- Assert.AreEqual(ChannelState.Closed, channel2.State);
- }
-
- [TestMethod("Client: SetsAuth")]
- public async Task ClientSetsAuth()
- {
- var channel = client!.Channel("realtime", "public", "todos");
- var channel2 = client!.Channel("realtime", "public", "todos");
-
- var token =
- @"eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJzdWIiOiIxMjM0NTY3ODkwIiwibmFtZSI6IkpvaG4gRG9lIiwiaWF0IjoxNTE2MjM5MDIyfQ.C8oVtF5DICct_4HcdSKt8pdrxBFMQOAnPpbiiUbaXAY";
-
- // No subscriptions should show a push
- client!.SetAuth(token);
- foreach (var subscription in client!.Subscriptions.Values)
- {
- Assert.IsNull(subscription.LastPush);
- }
-
- await channel.Subscribe();
- await channel2.Subscribe();
-
- client!.SetAuth(token);
- foreach (var subscription in client!.Subscriptions.Values)
- {
- Assert.IsTrue(subscription?.LastPush?.EventName == ChannelAccessToken);
- }
- }
-
- [TestMethod("Client: Can reconnect after programmatic disconnect")]
- public async Task ClientCanReconnectAfterProgrammaticDisconnect()
- {
- client!.Disconnect();
- await client!.ConnectAsync();
- }
-
- [TestMethod("Client: Sets headers")]
- public async Task ClientCanSetHeaders()
- {
- client!.Disconnect();
-
- client!.GetHeaders = () => new Dictionary() { { "testing", "123" } };
- await client.ConnectAsync();
-
- Assert.IsNotNull(client!);
- Assert.IsNotNull(client!.Socket);
- Assert.IsNotNull(client!.Socket.GetHeaders);
- Assert.AreEqual("123",client.Socket.GetHeaders()["testing"]);
- }
-}
\ No newline at end of file
diff --git a/RealtimeTests/Connection/ClientConnectionTests.cs b/RealtimeTests/Connection/ClientConnectionTests.cs
new file mode 100644
index 0000000..cdd9060
--- /dev/null
+++ b/RealtimeTests/Connection/ClientConnectionTests.cs
@@ -0,0 +1,74 @@
+using System.Collections.Generic;
+using System.Diagnostics;
+using System.Threading.Tasks;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using Supabase.Realtime;
+using Supabase.Realtime.Exceptions;
+
+namespace RealtimeTests.Connection;
+
+///
+/// The client's connection lifecycle against a live socket: failing loudly when the server is
+/// unreachable, staying failed on retry, reconnecting after a programmatic disconnect, and forwarding
+/// custom headers on connect.
+///
+[TestClass]
+[TestCategory("E2E")]
+public class ClientConnectionTests
+{
+ [TestMethod]
+ public async Task ConnectAsync_ShouldThrow_GivenUnreachableServer()
+ {
+ var client = new Client("ws://localhost");
+ client.AddDebugHandler((_, message, _) => Debug.WriteLine(message));
+ await Assert.ThrowsAsync(() => client.ConnectAsync());
+ }
+
+ [TestMethod]
+ public async Task Connect_ShouldReportError_GivenUnreachableServer()
+ {
+ var client = new Client("ws://localhost");
+ client.AddDebugHandler((_, message, _) => Debug.WriteLine(message));
+ await Assert.ThrowsAsync(() =>
+ {
+ var tsc = new TaskCompletionSource();
+#pragma warning disable CS0618 // pins the still-shipping obsolete callback Connect's error path
+ client.Connect((_, exception) =>
+ {
+ if (exception != null)
+ tsc.SetException(exception);
+ });
+#pragma warning restore CS0618
+ return tsc.Task;
+ });
+ }
+
+ [TestMethod]
+ public async Task ConnectAsync_ShouldStayFailed_GivenRepeatedAttempts()
+ {
+ var client = new Client("ws://localhost");
+ client.AddDebugHandler((_, message, _) => Debug.WriteLine(message));
+ await Assert.ThrowsAsync(client.ConnectAsync);
+ await Assert.ThrowsAsync(client.ConnectAsync);
+ }
+
+ [TestMethod]
+ public async Task ConnectAsync_ShouldReconnect_GivenProgrammaticDisconnect()
+ {
+ var client = Helpers.SocketClient();
+ await client.ConnectAsync();
+ client.Disconnect();
+ await client.ConnectAsync();
+ }
+
+ [TestMethod]
+ public async Task ConnectAsync_ShouldForwardCustomHeaders()
+ {
+ var client = Helpers.SocketClient();
+ client.GetHeaders = () => new Dictionary { { "testing", "123" } };
+ await client.ConnectAsync();
+ Assert.IsNotNull(client.Socket);
+ Assert.IsNotNull(client.Socket.GetHeaders);
+ Assert.AreEqual("123", client.Socket.GetHeaders!()["testing"]);
+ }
+}
diff --git a/RealtimeTests/ConverterTests.cs b/RealtimeTests/ConverterTests.cs
deleted file mode 100644
index 43d5481..0000000
--- a/RealtimeTests/ConverterTests.cs
+++ /dev/null
@@ -1,41 +0,0 @@
-using System.Collections.Generic;
-using Microsoft.VisualStudio.TestTools.UnitTesting;
-using Newtonsoft.Json;
-using Supabase.Realtime;
-using Supabase.Realtime.Converters;
-
-namespace RealtimeTests;
-
-internal class TestJson
-{
- [JsonProperty("intArray")] public List intArray { get; set; } = new();
-
- [JsonProperty("stringArray")] public List stringArray { get; set; } = new();
-}
-
-[TestClass]
-public class ConverterTests
-{
- [TestMethod("Support Array Conversions (WALRUS + Backwards Compat.)")]
- public void SupportArrayConversions()
- {
- var testJson = "{\"intArray\":[9999,99,99999], \"stringArray\": [\"testing\",\"1\",\"2\"]}";
-
- var parsed = JsonConvert.DeserializeObject(testJson, new JsonSerializerSettings
- {
- ContractResolver = new CustomContractResolver()
- });
-
- CollectionAssert.AreEqual(new List { 9999, 99, 99999 }, parsed?.intArray);
- CollectionAssert.AreEqual(new List { "testing", "1", "2" }, parsed?.stringArray);
-
- CollectionAssert.AreEqual(new List(), IntArrayConverter.Parse("{}"));
- CollectionAssert.AreEqual(new List { 1, 2, 3 }, IntArrayConverter.Parse("{1,2,3}"));
- CollectionAssert.AreEqual(new List { 1, 2, 3 }, IntArrayConverter.Parse("[1,2,3]"));
- CollectionAssert.AreEqual(new List { 99, 999, 9999, 999999 },
- IntArrayConverter.Parse("[99, 999, 9999, 999999]"));
-
- CollectionAssert.AreEqual(new List { "a", "b", "c" }, StringArrayConverter.Parse("{a,b,c}"));
- CollectionAssert.AreEqual(new List { "a", "b", "c" }, StringArrayConverter.Parse("[a,b,c]"));
- }
-}
\ No newline at end of file
diff --git a/RealtimeTests/Errors/FailureHintTests.cs b/RealtimeTests/Errors/FailureHintTests.cs
new file mode 100644
index 0000000..839cc31
--- /dev/null
+++ b/RealtimeTests/Errors/FailureHintTests.cs
@@ -0,0 +1,41 @@
+using System;
+using System.Net.WebSockets;
+using FluentAssertions;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using Supabase.Realtime.Exceptions;
+using Websocket.Client;
+
+namespace RealtimeTests.Errors;
+
+///
+/// How a transport disconnection is classified into a and surfaced on a
+/// , so a caller can branch on why a socket dropped.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class FailureHintTests
+{
+ private static DisconnectionInfo Info(DisconnectionType type, string? description = "dropped", Exception? exception = null) =>
+ new(type, WebSocketCloseStatus.NormalClosure, description, subProtocol: null, exception);
+
+ [TestMethod]
+ [DataRow(DisconnectionType.Error, FailureHint.Reason.SocketError)]
+ [DataRow(DisconnectionType.NoMessageReceived, FailureHint.Reason.ConnectionStale)]
+ [DataRow(DisconnectionType.Lost, FailureHint.Reason.ConnectionLost)]
+ [DataRow(DisconnectionType.ByServer, FailureHint.Reason.Unknown)]
+ [DataRow(DisconnectionType.ByUser, FailureHint.Reason.Unknown)]
+ public void Parse_ShouldClassifyDisconnectionType(DisconnectionType type, FailureHint.Reason expected)
+ {
+ FailureHint.Parse(Info(type)).Should().Be(expected);
+ }
+
+ [TestMethod]
+ public void FromDisconnectionInfo_ShouldCarryReasonMessageAndInnerException()
+ {
+ var inner = new InvalidOperationException("boom");
+ var exception = RealtimeException.FromDisconnectionInfo(Info(DisconnectionType.Error, "socket died", inner));
+ exception.Reason.Should().Be(FailureHint.Reason.SocketError);
+ exception.Message.Should().Be("socket died");
+ exception.InnerException.Should().BeSameAs(inner);
+ }
+}
diff --git a/RealtimeTests/Helpers.cs b/RealtimeTests/Helpers.cs
index d568300..5a5a4e7 100644
--- a/RealtimeTests/Helpers.cs
+++ b/RealtimeTests/Helpers.cs
@@ -1,4 +1,4 @@
-using System.Collections.Generic;
+using System.Collections.Generic;
using System.Diagnostics;
using Supabase.Realtime;
using Supabase.Realtime.Socket;
@@ -45,4 +45,4 @@ public static Client PrivateSocketClient()
return client;
}
-}
\ No newline at end of file
+}
diff --git a/RealtimeTests/Models/BroadcastExample.cs b/RealtimeTests/Models/BroadcastExample.cs
new file mode 100644
index 0000000..65d233b
--- /dev/null
+++ b/RealtimeTests/Models/BroadcastExample.cs
@@ -0,0 +1,12 @@
+using Newtonsoft.Json;
+using Supabase.Realtime.Models;
+
+namespace RealtimeTests.Models;
+
+///
+/// A modeled broadcast payload used across the broadcast tests.
+///
+public class BroadcastExample : BaseBroadcast
+{
+ [JsonProperty("userId")] public string? UserId { get; set; }
+}
diff --git a/RealtimeTests/Models/PresenceExample.cs b/RealtimeTests/Models/PresenceExample.cs
new file mode 100644
index 0000000..025a0c1
--- /dev/null
+++ b/RealtimeTests/Models/PresenceExample.cs
@@ -0,0 +1,13 @@
+using System;
+using Newtonsoft.Json;
+using Supabase.Realtime.Models;
+
+namespace RealtimeTests.Models;
+
+///
+/// A modeled presence payload used across the presence tests.
+///
+public class PresenceExample : BasePresence
+{
+ [JsonProperty("time")] public DateTime? Time { get; set; }
+}
diff --git a/RealtimeTests/Models/Todo.cs b/RealtimeTests/Models/Todo.cs
index faa385b..68659d8 100644
--- a/RealtimeTests/Models/Todo.cs
+++ b/RealtimeTests/Models/Todo.cs
@@ -1,4 +1,4 @@
-using System;
+using System;
using System.Collections.Generic;
using Supabase.Postgrest.Attributes;
using Supabase.Postgrest.Models;
@@ -17,4 +17,4 @@ public class Todo : BaseModel
[Column("numbers")] public List? Numbers { get; set; } = new List();
[Column("inserted_at")] public DateTime? InsertedAt { get; set; }
-}
\ No newline at end of file
+}
diff --git a/RealtimeTests/Models/User.cs b/RealtimeTests/Models/User.cs
index c7d6a87..7cd56b2 100644
--- a/RealtimeTests/Models/User.cs
+++ b/RealtimeTests/Models/User.cs
@@ -1,4 +1,4 @@
-using Supabase.Postgrest.Attributes;
+using Supabase.Postgrest.Attributes;
using Supabase.Postgrest.Models;
namespace RealtimeTests.Models;
@@ -11,4 +11,4 @@ public class User : BaseModel
[Column("name")]
public string? Name { get; set; }
-}
\ No newline at end of file
+}
diff --git a/RealtimeTests/PostgresChanges/PostgresChangesDeliveryTests.cs b/RealtimeTests/PostgresChanges/PostgresChangesDeliveryTests.cs
new file mode 100644
index 0000000..d028387
--- /dev/null
+++ b/RealtimeTests/PostgresChanges/PostgresChangesDeliveryTests.cs
@@ -0,0 +1,256 @@
+using System;
+using System.Collections.Generic;
+using System.Linq;
+using System.Threading.Tasks;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using RealtimeTests.Models;
+using Supabase.Postgrest.Interfaces;
+using Supabase.Realtime;
+using Supabase.Realtime.Interfaces;
+using Supabase.Realtime.PostgresChanges;
+using static Supabase.Realtime.Constants;
+using static Supabase.Realtime.PostgresChanges.PostgresChangesOptions;
+
+namespace RealtimeTests.PostgresChanges;
+
+///
+/// Postgres change delivery against the live stack: a channel that subscribed with an
+/// OnPostgresChange listener receives the matching insert/update/delete callbacks (filtered and
+/// wildcard), models the payload, and fans out to multiple listeners.
+///
+[TestClass]
+[TestCategory("E2E")]
+public class PostgresChangesDeliveryTests
+{
+ private IPostgrestClient restClient = null!;
+ private IRealtimeClient socketClient = null!;
+
+ [TestInitialize]
+ public async Task InitializeTest()
+ {
+ restClient = Helpers.RestClient();
+ socketClient = Helpers.SocketClient();
+ await socketClient.ConnectAsync();
+ }
+
+ [TestCleanup]
+ public void CleanupTest() => socketClient.Disconnect();
+
+ [TestMethod]
+ public async Task OnPostgresChange_ShouldModelPayload()
+ {
+ var tsc = new TaskCompletionSource();
+ var channel = socketClient.Channel("example");
+ channel.OnPostgresChange((_, changes) =>
+ {
+ var model = changes.Model();
+ tsc.SetResult(model != null);
+ }, ListenType.Inserts, new PostgresChangesFilter { Table = "*" });
+ await channel.Subscribe();
+ await restClient.Table().Insert(new Todo { UserId = 1, Details = "Client Models a response? ✅" });
+ Assert.IsTrue(await tsc.Task);
+ }
+
+ [TestMethod]
+ public async Task OnPostgresChange_ShouldReceiveInsert()
+ {
+ var tsc = new TaskCompletionSource();
+ var channel = socketClient.Channel("realtime", "public", "todos");
+ channel.OnPostgresChange((_, _) => tsc.SetResult(true), ListenType.Inserts,
+ new PostgresChangesFilter { Table = "todos" });
+ await channel.Subscribe();
+ await restClient.Table().Insert(new Todo { UserId = 1, Details = "Client receives insert callback? ✅" });
+ Assert.IsTrue(await tsc.Task);
+ }
+
+ [TestMethod]
+ public async Task OnPostgresChange_ShouldReceiveFilteredInsert()
+ {
+ var tsc = new TaskCompletionSource();
+ var channel = socketClient.Channel("realtime", "public", "todos");
+ channel.OnPostgresChange((_, changes) =>
+ {
+ Assert.AreEqual("Client receives filtered insert callback? ✅", changes.Model()?.Details);
+ tsc.SetResult(true);
+ }, ListenType.Inserts,
+ new PostgresChangesFilter { Table = "todos", Filter = "details=eq.Client receives filtered insert callback? ✅" });
+ await channel.Subscribe();
+ await restClient.Table().Insert(new Todo { UserId = 1, Details = "Client receives insert callback? ✅" });
+ await restClient.Table().Insert(new Todo { UserId = 2, Details = "Client receives filtered insert callback? ✅" });
+ Assert.IsTrue(await tsc.Task);
+ }
+
+ [TestMethod]
+ public async Task OnPostgresChange_ShouldReceiveUpdateAndFilteredInsert()
+ {
+ var tsc = new TaskCompletionSource();
+ var response = await restClient.Table()
+ .Insert(new Todo { UserId = 1, Details = "Client receives insert callback? ✅" });
+ await restClient.Table()
+ .Insert(new Todo { UserId = 2, Details = "Client receives filtered insert callback? ✅" });
+ var model = response.Models.First();
+ var oldDetails = model.Details;
+ var newDetails = $"I'm an updated item ✏️ - {DateTime.Now}";
+ var channel = socketClient.Channel("realtime", "public", "todos");
+ channel.OnPostgresChange((_, changes) =>
+ {
+ Assert.AreEqual(oldDetails, changes.OldModel()?.Details);
+ var updated = changes.Model();
+ Assert.AreEqual(newDetails, updated?.Details);
+ if (updated != null)
+ {
+ Assert.AreEqual(model.Id, updated.Id);
+ Assert.AreEqual(model.UserId, updated.UserId);
+ }
+ tsc.SetResult(true);
+ }, ListenType.Updates, new PostgresChangesFilter { Table = "todos" });
+ const string filter = "Client receives filtered insert callback? ✅";
+ channel.OnPostgresChange((_, changes) =>
+ {
+ Assert.AreEqual("Client receives filtered insert callback? ✅", changes.Model()?.Details);
+ tsc.SetResult(true);
+ }, ListenType.Inserts, new PostgresChangesFilter { Table = "todos", Filter = $"details=eq.{filter}" });
+ await channel.Subscribe();
+ await restClient.Table().Set(x => x.Details!, newDetails).Match(model).Update();
+ Assert.IsTrue(await tsc.Task);
+ }
+
+ [TestMethod]
+ public async Task OnPostgresChange_ShouldReceiveUpdate()
+ {
+ var tsc = new TaskCompletionSource();
+ var response = await restClient.Table()
+ .Insert(new Todo { UserId = 1, Details = "Client receives insert callback? ✅" });
+ var model = response.Models.First();
+ var oldDetails = model.Details;
+ var newDetails = $"I'm an updated item ✏️ - {DateTime.Now}";
+ var channel = socketClient.Channel("realtime", "public", "todos");
+ channel.OnPostgresChange((_, changes) =>
+ {
+ Assert.AreEqual(oldDetails, changes.OldModel()?.Details);
+ var updated = changes.Model();
+ Assert.AreEqual(newDetails, updated?.Details);
+ if (updated != null)
+ {
+ Assert.AreEqual(model.Id, updated.Id);
+ Assert.AreEqual(model.UserId, updated.UserId);
+ }
+ tsc.SetResult(true);
+ }, ListenType.Updates, new PostgresChangesFilter { Table = "todos" });
+ await channel.Subscribe();
+ await restClient.Table().Set(x => x.Details!, newDetails).Match(model).Update();
+ Assert.IsTrue(await tsc.Task);
+ }
+
+ [TestMethod]
+ public async Task OnPostgresChange_ShouldReceiveDelete()
+ {
+ var tsc = new TaskCompletionSource();
+ var channel = socketClient.Channel("realtime", "public", "todos");
+ channel.OnPostgresChange((_, _) => tsc.SetResult(true), ListenType.Deletes,
+ new PostgresChangesFilter { Table = "todos" });
+ await channel.Subscribe();
+ var result = await restClient.Table().Get();
+ var model = result.Models.Last();
+ await restClient.Table().Match(model).Delete();
+ Assert.IsTrue(await tsc.Task);
+ }
+
+ [TestMethod]
+ public async Task OnPostgresChange_ShouldReceiveFilteredDelete()
+ {
+ var tsc = new TaskCompletionSource();
+ var channel = socketClient.Channel("realtime", "public", "todos");
+ var todo1 = await restClient.Table().Insert(new Todo { UserId = 1, Details = "Client receives callbacks 1? ✅" });
+ var todo2 = await restClient.Table().Insert(new Todo { UserId = 2, Details = "Client receives callbacks 2? ✅" });
+ await restClient.Table().Insert(new Todo { UserId = 3, Details = "Client receives callbacks 3? ✅" });
+ channel.OnPostgresChange((_, removed) =>
+ {
+ var result = removed.OldModel();
+ Assert.AreEqual(result?.Details, todo1.Model?.Details);
+ Assert.AreNotEqual(result?.Details, todo2.Model?.Details);
+ tsc.SetResult(true);
+ }, ListenType.Deletes,
+ new PostgresChangesFilter { Table = "todos", Filter = $"details=eq.{todo1.Model?.Details}" });
+ await channel.Subscribe();
+ await restClient.Table().Match(todo1.Models.First()).Delete();
+ await restClient.Table().Match(todo2.Models.First()).Delete();
+ Assert.IsTrue(await tsc.Task);
+ }
+
+ [TestMethod]
+ public async Task OnPostgresChange_ShouldReceiveAllEvents_GivenWildcard()
+ {
+ var insertTsc = new TaskCompletionSource();
+ var updateTsc = new TaskCompletionSource();
+ var deleteTsc = new TaskCompletionSource();
+ var channel = socketClient.Channel("realtime", "public", "todos");
+ channel.OnPostgresChange((_, changes) =>
+ {
+ switch (changes.Payload?.Data?.Type)
+ {
+ case EventType.Insert: insertTsc.SetResult(true); break;
+ case EventType.Update: updateTsc.SetResult(true); break;
+ case EventType.Delete: deleteTsc.SetResult(true); break;
+ }
+ }, ListenType.All, new PostgresChangesFilter { Table = "todos" });
+ await channel.Subscribe();
+ var inserted = await restClient.Table().Insert(new Todo { UserId = 1, Details = "Client receives wildcard callbacks? ✅" });
+ var newModel = inserted.Models.First();
+ await restClient.Table().Set(x => x.Details!, "And edits.").Match(newModel).Update();
+ await restClient.Table().Match(newModel).Delete();
+ await Task.WhenAll(insertTsc.Task, updateTsc.Task, deleteTsc.Task);
+ Assert.IsTrue(insertTsc.Task.Result);
+ Assert.IsTrue(updateTsc.Task.Result);
+ Assert.IsTrue(deleteTsc.Task.Result);
+ }
+
+ [TestMethod]
+ public async Task OnPostgresChange_ShouldFanOutToMultipleInsertListeners()
+ {
+ var insertTask1 = new TaskCompletionSource();
+ var insertTask2 = new TaskCompletionSource();
+ var insertTask3 = new TaskCompletionSource();
+ const string filter1 = "Client receives callbacks 1? ✅";
+ const string filter2 = "Client receives callbacks 2? ✅";
+ var channel = socketClient.Channel("realtime", "public", "todos");
+ var count = 0;
+ channel.OnPostgresChange((_, _) =>
+ {
+ count++;
+ if (count == 3) insertTask1.TrySetResult(true);
+ }, ListenType.Inserts, new PostgresChangesFilter { Table = "todos" });
+ channel.OnPostgresChange((_, added) =>
+ insertTask2.SetResult(added.Model()?.Details == filter1), ListenType.Inserts,
+ new PostgresChangesFilter { Table = "todos", Filter = $"details=eq.{filter1}" });
+ channel.OnPostgresChange((_, added) =>
+ insertTask3.SetResult(added.Model()?.Details == filter2), ListenType.Inserts,
+ new PostgresChangesFilter { Table = "todos", Filter = $"details=eq.{filter2}" });
+ await channel.Subscribe();
+ await restClient.Table().Insert(new Todo { UserId = 1, Details = "Client receives wildcard callbacks? ✅" });
+ await restClient.Table().Insert(new Todo { UserId = 1, Details = filter1 });
+ await restClient.Table().Insert(new Todo { UserId = 1, Details = filter2 });
+ await Task.WhenAll(insertTask1.Task, insertTask2.Task, insertTask3.Task);
+ Assert.IsTrue(insertTask1.Task.Result);
+ Assert.IsTrue(insertTask2.Task.Result);
+ Assert.IsTrue(insertTask3.Task.Result);
+ }
+
+ [TestMethod]
+ public async Task OnPostgresChange_ShouldRegisterAndDeliver_GivenChainedSubscribe()
+ {
+ var tsc = new TaskCompletionSource();
+ await socketClient.Channel("public:todos")
+ .OnPostgresChange((_, changes) => tsc.TrySetResult(changes.Model() != null),
+ ListenType.Inserts, new PostgresChangesFilter { Table = "todos" })
+ .Subscribe();
+ await restClient.Table().Insert(new Todo { UserId = 1, Details = "OnPostgresChange receives insert? ✅" });
+ Assert.IsTrue(await WithinTimeout(tsc.Task));
+ }
+
+ private static async Task WithinTimeout(Task task, int timeoutMs = 15000)
+ {
+ var completed = await Task.WhenAny(task, Task.Delay(timeoutMs));
+ return completed == task && task.Result;
+ }
+}
diff --git a/RealtimeTests/PostgresChanges/PostgresChangesModelTests.cs b/RealtimeTests/PostgresChanges/PostgresChangesModelTests.cs
new file mode 100644
index 0000000..453ed96
--- /dev/null
+++ b/RealtimeTests/PostgresChanges/PostgresChangesModelTests.cs
@@ -0,0 +1,77 @@
+using System;
+using System.IO;
+using System.Threading.Tasks;
+using FluentAssertions;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using Newtonsoft.Json;
+using RealtimeTests.Models;
+using Supabase.Postgrest.Exceptions;
+using Supabase.Postgrest.Interfaces;
+using Supabase.Realtime.PostgresChanges;
+
+namespace RealtimeTests.PostgresChanges;
+
+///
+/// Hydration of a postgres_changes frame into a model (realtime-csharp#35): Model<T> and
+/// OldModel<T> deserialize the record, and when a PostgrestClient is configured they
+/// attach its context so Update/Delete can be called on the returned model. This reproduces
+/// the exact deserialization path RealtimeChannel.HandleSocketMessage uses, without a live socket.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class PostgresChangesModelTests
+{
+ private const string BaseUrl = "http://localhost:54321/rest/v1";
+
+ private static PostgresChangesResponse BuildResponse(IPostgrestClient? postgrestClient)
+ {
+ var json = File.ReadAllText(Path.Combine(AppContext.BaseDirectory, "Fixtures", "PostgresChangesUpdateEvent.json"));
+ var settings = Supabase.Postgrest.Client.SerializerSettings();
+ var response = JsonConvert.DeserializeObject(json, settings)!;
+ response.Json = json;
+ response.SerializerSettings = settings;
+ response.PostgrestClient = postgrestClient;
+ return response;
+ }
+
+ [TestMethod]
+ public void Model_ShouldDeserialiseRecord()
+ {
+ var model = BuildResponse(postgrestClient: null).Model();
+ model.Should().NotBeNull();
+ model!.Id.Should().Be(12);
+ model.Details.Should().Be("test...");
+ }
+
+ [TestMethod]
+ public void OldModel_ShouldDeserialiseOldRecord()
+ {
+ var model = BuildResponse(postgrestClient: null).OldModel();
+ model.Should().NotBeNull();
+ model!.Details.Should().Be("previous");
+ }
+
+ [TestMethod]
+ public void Model_ShouldLeaveClientContextNull_GivenNoPostgrestClient()
+ {
+ var model = BuildResponse(postgrestClient: null).Model()!;
+ model.BaseUrl.Should().BeNull();
+ model.RequestClientOptions.Should().BeNull();
+ }
+
+ [TestMethod]
+ public void Model_ShouldAttachClientContext_GivenPostgrestClient()
+ {
+ var model = BuildResponse(new Supabase.Postgrest.Client(BaseUrl)).Model()!;
+ model.BaseUrl.Should().Be(BaseUrl);
+ model.RequestClientOptions.Should().NotBeNull();
+ }
+
+ [TestMethod]
+ public async Task Update_ShouldThrowLoudly_GivenNoClientContext()
+ {
+ var model = BuildResponse(postgrestClient: null).Model()!;
+ var act = () => model.Update();
+ (await act.Should().ThrowAsync()).Which.Message.Should().Contain("BaseUrl");
+ }
+}
diff --git a/RealtimeTests/PostgresChanges/PostgresChangesOptionsTests.cs b/RealtimeTests/PostgresChanges/PostgresChangesOptionsTests.cs
new file mode 100644
index 0000000..53167dc
--- /dev/null
+++ b/RealtimeTests/PostgresChanges/PostgresChangesOptionsTests.cs
@@ -0,0 +1,89 @@
+using FluentAssertions;
+using FluentAssertions.Execution;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using Newtonsoft.Json;
+using Newtonsoft.Json.Linq;
+using Supabase.Realtime.PostgresChanges;
+using static Supabase.Realtime.PostgresChanges.PostgresChangesOptions;
+
+namespace RealtimeTests.PostgresChanges;
+
+///
+/// The shape takes inside a channel's join
+/// config.postgres_changes payload, plus the value-equality that lets a channel dedupe repeated
+/// registrations of the same listener.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class PostgresChangesOptionsTests
+{
+ [TestMethod]
+ public void Table_ShouldBeOmitted_GivenNull()
+ {
+ var json = JsonConvert.SerializeObject(new PostgresChangesOptions("public"));
+ JObject.Parse(json).ContainsKey("table").Should().BeFalse();
+ }
+
+ [TestMethod]
+ public void Table_ShouldBeSerialized_GivenProvided()
+ {
+ var json = JsonConvert.SerializeObject(new PostgresChangesOptions("public", "todos"));
+ JObject.Parse(json)["table"]?.Value().Should().Be("todos");
+ }
+
+ [TestMethod]
+ public void Filter_ShouldBeOmitted_GivenNull()
+ {
+ var json = JsonConvert.SerializeObject(new PostgresChangesOptions("public", "todos"));
+ JObject.Parse(json).ContainsKey("filter").Should().BeFalse();
+ }
+
+ [TestMethod]
+ [DataRow(ListenType.All, "*")]
+ [DataRow(ListenType.Inserts, "INSERT")]
+ [DataRow(ListenType.Updates, "UPDATE")]
+ [DataRow(ListenType.Deletes, "DELETE")]
+ public void Event_ShouldRenderListenType(ListenType listenType, string expected)
+ {
+ new PostgresChangesOptions("public", "todos", listenType).Event.Should().Be(expected);
+ }
+
+ [TestMethod]
+ public void Equals_ShouldHold_GivenAllDiscriminatorsMatch()
+ {
+ var a = new PostgresChangesOptions("public", "todos", ListenType.Inserts, "id=eq.1");
+ var b = new PostgresChangesOptions("public", "todos", ListenType.Inserts, "id=eq.1");
+ using (new AssertionScope())
+ {
+ a.Equals(b).Should().BeTrue();
+ a.GetHashCode().Should().Be(b.GetHashCode());
+ }
+ }
+
+ [TestMethod]
+ [DataRow("other", "todos", ListenType.Inserts, "id=eq.1")]
+ [DataRow("public", "users", ListenType.Inserts, "id=eq.1")]
+ [DataRow("public", "todos", ListenType.Updates, "id=eq.1")]
+ [DataRow("public", "todos", ListenType.Inserts, "id=eq.2")]
+ public void Equals_ShouldFail_GivenDiscriminatorDiffers(string schema, string table, ListenType listenType, string filter)
+ {
+ var reference = new PostgresChangesOptions("public", "todos", ListenType.Inserts, "id=eq.1");
+ reference.Equals(new PostgresChangesOptions(schema, table, listenType, filter)).Should().BeFalse();
+ }
+
+ [TestMethod]
+ public void Equals_ShouldFail_GivenNull()
+ {
+ var reference = new PostgresChangesOptions("public", "todos");
+ var result = reference.Equals((object?) null);
+ result.Should().BeFalse();
+ }
+
+ [TestMethod]
+ public void Equals_ShouldFail_GivenDifferentType()
+ {
+ var reference = new PostgresChangesOptions("public", "todos");
+ var result = reference.Equals(new object());
+ result.Should().BeFalse();
+ }
+}
diff --git a/RealtimeTests/PostgresChanges/PostgresChangesRegistrationTests.cs b/RealtimeTests/PostgresChanges/PostgresChangesRegistrationTests.cs
new file mode 100644
index 0000000..1ecbec4
--- /dev/null
+++ b/RealtimeTests/PostgresChanges/PostgresChangesRegistrationTests.cs
@@ -0,0 +1,57 @@
+using FluentAssertions;
+using FluentAssertions.Execution;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using RealtimeTests.Support;
+using Supabase.Realtime.PostgresChanges;
+using static Supabase.Realtime.PostgresChanges.PostgresChangesOptions;
+
+namespace RealtimeTests.PostgresChanges;
+
+///
+/// How a channel records a postgres_changes listener. OnPostgresChange is the one-call API: it
+/// stores the options built from the filter and returns the channel so Subscribe can be chained.
+/// The obsolete Register path feeds the same list.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class PostgresChangesRegistrationTests
+{
+ [TestMethod]
+ public void OnPostgresChange_ShouldStoreOptionsAndReturnChannel()
+ {
+ var channel = Wire.Channel();
+ var returned = channel.OnPostgresChange((_, _) => { }, ListenType.Inserts,
+ new PostgresChangesFilter { Schema = "public", Table = "todos", Filter = "id=eq.1" });
+ using (new AssertionScope())
+ {
+ returned.Should().BeSameAs(channel);
+ channel.PostgresChangesOptions.Should().ContainSingle();
+ var options = channel.PostgresChangesOptions[0];
+ options.Schema.Should().Be("public");
+ options.Table.Should().Be("todos");
+ options.Filter.Should().Be("id=eq.1");
+ options.Event.Should().Be("INSERT");
+ }
+ }
+
+ [TestMethod]
+ public void OnPostgresChange_ShouldDefaultToSchemaWideListener_GivenNoFilter()
+ {
+ var channel = Wire.Channel();
+ channel.OnPostgresChange((_, _) => { }, ListenType.All);
+ var options = channel.PostgresChangesOptions.Should().ContainSingle().Subject;
+ options.Schema.Should().Be("public");
+ options.Table.Should().BeNull();
+ }
+
+ [TestMethod]
+ public void Register_ShouldRecordOptions()
+ {
+ var channel = Wire.Channel();
+#pragma warning disable CS0618
+ var returned = channel.Register(new PostgresChangesOptions("public", "todos"));
+#pragma warning restore CS0618
+ returned.Should().BeSameAs(channel);
+ channel.PostgresChangesOptions.Should().ContainSingle().Which.Table.Should().Be("todos");
+ }
+}
diff --git a/RealtimeTests/PostgresChangesResponseAttachTests.cs b/RealtimeTests/PostgresChangesResponseAttachTests.cs
deleted file mode 100644
index 46ff8a5..0000000
--- a/RealtimeTests/PostgresChangesResponseAttachTests.cs
+++ /dev/null
@@ -1,122 +0,0 @@
-using System;
-using System.IO;
-using System.Threading.Tasks;
-using Microsoft.VisualStudio.TestTools.UnitTesting;
-using Newtonsoft.Json;
-using RealtimeTests.Models;
-using Supabase.Postgrest.Exceptions;
-using Supabase.Postgrest.Interfaces;
-using Supabase.Realtime.PostgresChanges;
-
-namespace RealtimeTests;
-
-///
-/// Feature under test: realtime-csharp#35 - a model returned by PostgresChangesResponse.Model<TModel>/
-/// OldModel<TModel> gets a configured 's
-/// context attached to it (via Attach<T>()), so that Update/Delete can be called
-/// directly on it. Reproduces the exact deserialization path RealtimeChannel.HandleSocketMessage() uses for
-/// postgres_changes events, without needing a live socket connection.
-///
-/// Split into two groups:
-/// - no PostgrestClient configured: the default/standalone case, which must stay backward-compatible.
-/// - a PostgrestClient configured: the case the Supabase umbrella client wires up automatically.
-///
-[TestClass]
-public class PostgresChangesResponseAttachTests
-{
- private const string BaseUrl = "http://localhost:54321/rest/v1";
-
- private static PostgresChangesResponse BuildResponse(IPostgrestClient? postgrestClient)
- {
- var json = File.ReadAllText(Path.Combine(AppContext.BaseDirectory, "Fixtures", "PostgresChangesUpdateEvent.json"));
- var serializerSettings = Supabase.Postgrest.Client.SerializerSettings();
-
- // Mirrors RealtimeChannel.HandleSocketMessage() exactly: deserialize, then stamp Json, SerializerSettings,
- // and PostgrestClient for the later Model() re-deserialization pass.
- var deserialized = JsonConvert.DeserializeObject(json, serializerSettings);
- deserialized!.Json = json;
- deserialized.SerializerSettings = serializerSettings;
- deserialized.PostgrestClient = postgrestClient;
- return deserialized;
- }
-
- private static void AssertDoesNotThrowBaseUrlException(Exception ex)
- {
- Assert.IsFalse(ex is PostgrestException { Message: var m } && m.Contains("should be set in the model"),
- $"Unexpectedly got the BaseUrl exception: {ex}");
- }
-
- [TestMethod(DisplayName = "Without a PostgrestClient, Model() leaves BaseUrl/RequestClientOptions null")]
- public void GivenNoPostgrestClient_Model_LeavesClientContextNull()
- {
- var response = BuildResponse(postgrestClient: null);
- var model = response.Model();
-
- Assert.IsNotNull(model, "Sanity check: the model itself should deserialize correctly.");
- Assert.AreEqual(12, model!.Id);
- Assert.IsNull(model.BaseUrl);
- Assert.IsNull(model.RequestClientOptions);
- }
-
- [TestMethod(DisplayName = "Without a PostgrestClient, Update() throws PostgrestException (loud, not silent)")]
- public async Task GivenNoPostgrestClient_Update_ThrowsPostgrestException()
- {
- var response = BuildResponse(postgrestClient: null);
- var model = response.Model()!;
-
- var ex = await Assert.ThrowsAsync(() => model.Update());
- StringAssert.Contains(ex.Message, "BaseUrl");
- }
-
- [TestMethod(DisplayName = "Without a PostgrestClient, Delete() throws PostgrestException (loud, not silent)")]
- public async Task GivenNoPostgrestClient_Delete_ThrowsPostgrestException()
- {
- var response = BuildResponse(postgrestClient: null);
- var model = response.Model()!;
-
- var ex = await Assert.ThrowsAsync(() => model.Delete());
- StringAssert.Contains(ex.Message, "BaseUrl");
- }
-
- [TestMethod(DisplayName = "With a PostgrestClient configured, Model() attaches BaseUrl/RequestClientOptions")]
- public void GivenPostgrestClient_Model_AttachesClientContext()
- {
- var response = BuildResponse(new Supabase.Postgrest.Client(BaseUrl));
- var model = response.Model()!;
-
- Assert.AreEqual(BaseUrl, model.BaseUrl);
- Assert.IsNotNull(model.RequestClientOptions);
- }
-
- [TestMethod(DisplayName = "With a PostgrestClient configured, Update() does not throw the BaseUrl exception")]
- public async Task GivenPostgrestClient_Update_DoesNotThrowBaseUrlException()
- {
- var response = BuildResponse(new Supabase.Postgrest.Client(BaseUrl));
- var model = response.Model()!;
-
- try
- {
- await model.Update();
- }
- catch (Exception ex)
- {
- AssertDoesNotThrowBaseUrlException(ex);
- }
- }
-
- [TestMethod(DisplayName = "With a PostgrestClient configured, Delete() does not throw the BaseUrl exception")]
- public async Task GivenPostgrestClient_Delete_DoesNotThrowBaseUrlException()
- {
- var response = BuildResponse(new Supabase.Postgrest.Client(BaseUrl));
- var model = response.Model()!;
-
- try
- {
- await model.Delete();
- }
- catch (Exception ex)
- {
- AssertDoesNotThrowBaseUrlException(ex);
- }
- }
-}
diff --git a/RealtimeTests/Presence/PresenceStateTests.cs b/RealtimeTests/Presence/PresenceStateTests.cs
new file mode 100644
index 0000000..5c660e5
--- /dev/null
+++ b/RealtimeTests/Presence/PresenceStateTests.cs
@@ -0,0 +1,77 @@
+using System;
+using FluentAssertions;
+using FluentAssertions.Execution;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using RealtimeTests.Models;
+using RealtimeTests.Support;
+using Supabase.Realtime;
+using Supabase.Realtime.Interfaces;
+using Supabase.Realtime.Socket;
+
+namespace RealtimeTests.Presence;
+
+///
+/// How a channel folds presence frames into shared state: a presence_state frame seeds the current
+/// state, and a presence_diff applies joins and leaves. Both are driven through the channel over a
+/// fake socket.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class PresenceStateTests
+{
+ private const string StateFrame =
+ "{\"topic\":\"realtime:online-users\",\"event\":\"presence_state\",\"ref\":\"s1\"," +
+ "\"payload\":{\"user-1\":{\"metas\":[{\"phx_ref\":\"r1\",\"time\":\"2023-09-01T10:00:00Z\"}]}}}";
+
+ private const string DiffFrame =
+ "{\"topic\":\"realtime:online-users\",\"event\":\"presence_diff\",\"ref\":\"s2\"," +
+ "\"payload\":{\"joins\":{\"user-2\":{\"metas\":[{\"phx_ref\":\"r2\",\"time\":\"2023-09-01T11:00:00Z\"}]}}," +
+ "\"leaves\":{\"user-1\":{\"metas\":[{\"phx_ref\":\"r1\"}]}}}}";
+
+ private static (RealtimeChannel channel, RealtimePresence presence) Presence()
+ {
+ var channel = Wire.Channel("realtime:online-users");
+ return (channel, channel.Register("user-1"));
+ }
+
+ [TestMethod]
+ public void HandleSocketMessage_ShouldSeedStateAndRaiseSync_GivenPresenceState()
+ {
+ var (channel, presence) = Presence();
+ var syncs = 0;
+ presence.AddPresenceEventHandler(IRealtimePresence.EventType.Sync, (_, _) => syncs++);
+ channel.HandleSocketMessage(Wire.Decode(StateFrame));
+ using (new AssertionScope())
+ {
+ presence.CurrentState.Should().ContainKey("user-1");
+ presence.CurrentState["user-1"].Should().ContainSingle().Which.Time.Should().NotBeNull();
+ syncs.Should().Be(1);
+ }
+ }
+
+ [TestMethod]
+ public void HandleSocketMessage_ShouldApplyJoinsAndLeaves_GivenPresenceDiff()
+ {
+ var (channel, presence) = Presence();
+ channel.HandleSocketMessage(Wire.Decode(StateFrame));
+ int joins = 0, leaves = 0;
+ presence.AddPresenceEventHandler(IRealtimePresence.EventType.Join, (_, _) => joins++);
+ presence.AddPresenceEventHandler(IRealtimePresence.EventType.Leave, (_, _) => leaves++);
+ channel.HandleSocketMessage(Wire.Decode(DiffFrame));
+ using (new AssertionScope())
+ {
+ presence.CurrentState.Should().ContainKey("user-2");
+ presence.CurrentState.Should().NotContainKey("user-1");
+ joins.Should().Be(1);
+ leaves.Should().Be(1);
+ }
+ }
+
+ [TestMethod]
+ public void TriggerDiff_ShouldThrow_GivenUnparsableResponse()
+ {
+ var (_, presence) = Presence();
+ var act = () => presence.TriggerDiff(new SocketResponse(Wire.Settings()));
+ act.Should().Throw();
+ }
+}
diff --git a/RealtimeTests/ChannelPresenceTests.cs b/RealtimeTests/Presence/PresenceTrackingTests.cs
similarity index 65%
rename from RealtimeTests/ChannelPresenceTests.cs
rename to RealtimeTests/Presence/PresenceTrackingTests.cs
index 1e0353b..0d55d26 100644
--- a/RealtimeTests/ChannelPresenceTests.cs
+++ b/RealtimeTests/Presence/PresenceTrackingTests.cs
@@ -1,50 +1,44 @@
-using System;
+using System;
using System.Linq;
using System.Threading.Tasks;
using Microsoft.VisualStudio.TestTools.UnitTesting;
-using Newtonsoft.Json;
+using RealtimeTests.Models;
using Supabase.Postgrest.Interfaces;
using Supabase.Realtime;
using Supabase.Realtime.Interfaces;
-using Supabase.Realtime.Models;
-namespace RealtimeTests;
-
-public class PresenceExample : BasePresence
-{
- [JsonProperty("time")] public DateTime? Time { get; set; }
-}
+namespace RealtimeTests.Presence;
+///
+/// Presence against the live stack: two clients tracking presence on the same channel each observe the
+/// other's state through the sync event.
+///
[TestClass]
-public class ChannelPresenceTests
+[TestCategory("E2E")]
+public class PresenceTrackingTests
{
- private IPostgrestClient? _restClient;
- private IRealtimeClient? _socketClient;
+ private IPostgrestClient restClient = null!;
+ private IRealtimeClient socketClient = null!;
[TestInitialize]
public async Task InitializeTest()
{
- _restClient = Helpers.RestClient();
- _socketClient = Helpers.SocketClient();
- await _socketClient!.ConnectAsync();
+ restClient = Helpers.RestClient();
+ socketClient = Helpers.SocketClient();
+ await socketClient.ConnectAsync();
}
[TestCleanup]
- public void CleanupTest()
- {
- _socketClient!.Disconnect();
- }
+ public void CleanupTest() => socketClient.Disconnect();
- [TestMethod("Channel: Can create presence")]
- public async Task ClientCanCreatePresence()
+ [TestMethod]
+ public async Task Track_ShouldSyncPresenceBetweenClients()
{
var tsc = new TaskCompletionSource();
var tsc2 = new TaskCompletionSource();
-
var guid1 = Guid.NewGuid().ToString();
var guid2 = Guid.NewGuid().ToString();
-
- var channel1 = _socketClient!.Channel("online-users");
+ var channel1 = socketClient.Channel("online-users");
var presence1 = channel1.Register(guid1);
presence1.AddPresenceEventHandler(IRealtimePresence.EventType.Sync, (_, _) =>
{
@@ -52,7 +46,6 @@ public async Task ClientCanCreatePresence()
if (state.ContainsKey(guid2) && state[guid2].First().Time != null)
tsc.TrySetResult(true);
});
-
var client2 = Helpers.SocketClient();
await client2.ConnectAsync();
var channel2 = client2.Channel("online-users");
@@ -63,15 +56,11 @@ public async Task ClientCanCreatePresence()
if (state.ContainsKey(guid1) && state[guid1].First().Time != null)
tsc2.TrySetResult(true);
});
-
await channel1.Subscribe();
await channel2.Subscribe();
-
await presence1.Track(new PresenceExample { Time = DateTime.Now });
await presence2.Track(new PresenceExample { Time = DateTime.Now });
-
await presence1.Untrack();
-
- await Task.WhenAll(new[] { tsc.Task, tsc2.Task });
+ await Task.WhenAll(tsc.Task, tsc2.Task);
}
-}
\ No newline at end of file
+}
diff --git a/RealtimeTests/RealtimeTests.csproj b/RealtimeTests/RealtimeTests.csproj
index 665f926..ba32653 100644
--- a/RealtimeTests/RealtimeTests.csproj
+++ b/RealtimeTests/RealtimeTests.csproj
@@ -5,6 +5,9 @@
enable
latest
CS8600;CS8602;CS8603
+ true
+ true
+ latest
@@ -15,12 +18,29 @@
runtime; build; native; contentfiles; analyzers; buildtransitive
all
+
+
+
+
+
+
+
+
+
+
+
PreserveNewest
diff --git a/RealtimeTests/Serialization/ArrayConverterTests.cs b/RealtimeTests/Serialization/ArrayConverterTests.cs
new file mode 100644
index 0000000..12dc287
--- /dev/null
+++ b/RealtimeTests/Serialization/ArrayConverterTests.cs
@@ -0,0 +1,72 @@
+using System.Collections.Generic;
+using FluentAssertions;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using Newtonsoft.Json;
+using Supabase.Realtime;
+using Supabase.Realtime.Converters;
+
+namespace RealtimeTests.Serialization;
+
+///
+/// WALRUS delivers Postgres array columns as strings in either {1,2,3} (curly) or [1,2,3]
+/// (bracket) form; and coerce both
+/// back into a typed list. These pin that coercion, including the empty and whitespace edges.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class ArrayConverterTests
+{
+ private class ArrayModel
+ {
+ [JsonProperty("intArray")] public List IntArray { get; set; } = new();
+ [JsonProperty("stringArray")] public List StringArray { get; set; } = new();
+ }
+
+ private static ArrayModel Coerce(string json) =>
+ JsonConvert.DeserializeObject(json,
+ new JsonSerializerSettings { ContractResolver = new CustomContractResolver() })!;
+
+ [TestMethod]
+ public void ContractResolver_ShouldCoerceArrayColumns_GivenJsonArrays()
+ {
+ var parsed = Coerce("{\"intArray\":[9999,99,99999], \"stringArray\": [\"testing\",\"1\",\"2\"]}");
+ parsed.IntArray.Should().Equal(9999, 99, 99999);
+ parsed.StringArray.Should().Equal("testing", "1", "2");
+ }
+
+ [TestMethod]
+ public void ContractResolver_ShouldCoerceArrayColumns_GivenWalrusEncodedStrings()
+ {
+ // WALRUS delivers array columns as a single Postgres-array string; the converter (not default
+ // deserialization) is what turns it back into a list, so this exercises the converter is engaged.
+ var parsed = Coerce("{\"intArray\":\"{9999,99,99999}\", \"stringArray\":\"{testing,1,2}\"}");
+ parsed.IntArray.Should().Equal(9999, 99, 99999);
+ parsed.StringArray.Should().Equal("testing", "1", "2");
+ }
+
+ [TestMethod]
+ public void IntArrayParse_ShouldReadCurlyForm()
+ {
+ IntArrayConverter.Parse("{1,2,3}").Should().Equal(1, 2, 3);
+ }
+
+ [TestMethod]
+ public void IntArrayParse_ShouldReadBracketForm_GivenWhitespace()
+ {
+ IntArrayConverter.Parse("[1,2,3]").Should().Equal(1, 2, 3);
+ IntArrayConverter.Parse("[99, 999, 9999, 999999]").Should().Equal(99, 999, 9999, 999999);
+ }
+
+ [TestMethod]
+ public void IntArrayParse_ShouldReturnEmpty_GivenEmptyBraces()
+ {
+ IntArrayConverter.Parse("{}").Should().BeEmpty();
+ }
+
+ [TestMethod]
+ public void StringArrayParse_ShouldReadBothForms()
+ {
+ StringArrayConverter.Parse("{a,b,c}").Should().Equal("a", "b", "c");
+ StringArrayConverter.Parse("[a,b,c]").Should().Equal("a", "b", "c");
+ }
+}
diff --git a/RealtimeTests/Serialization/DateTimeCoercionTests.cs b/RealtimeTests/Serialization/DateTimeCoercionTests.cs
new file mode 100644
index 0000000..6215b92
--- /dev/null
+++ b/RealtimeTests/Serialization/DateTimeCoercionTests.cs
@@ -0,0 +1,54 @@
+using System;
+using System.Collections.Generic;
+using FluentAssertions;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using Newtonsoft.Json;
+using Supabase.Realtime;
+
+namespace RealtimeTests.Serialization;
+
+///
+/// Date coercion the resolver wires onto columns: ISO timestamps round-trip, and
+/// Postgres' infinity/-infinity sentinels map onto /
+/// . Driven through , which is how the
+/// converter is actually applied.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class DateTimeCoercionTests
+{
+ private class DateModel
+ {
+ [JsonProperty("at")] public DateTime? At { get; set; }
+ [JsonProperty("many")] public List? Many { get; set; }
+ }
+
+ private static DateModel Parse(string json) =>
+ JsonConvert.DeserializeObject(json,
+ new JsonSerializerSettings { ContractResolver = new CustomContractResolver() })!;
+
+ [TestMethod]
+ public void At_ShouldParseIsoTimestamp()
+ {
+ Parse("{\"at\":\"2023-09-11T15:30:21Z\"}").At.Should().Be(new DateTime(2023, 9, 11, 15, 30, 21, DateTimeKind.Utc));
+ }
+
+ [TestMethod]
+ public void At_ShouldMapPositiveInfinityToMaxValue()
+ {
+ Parse("{\"at\":\"infinity\"}").At.Should().Be(DateTime.MaxValue);
+ }
+
+ [TestMethod]
+ public void At_ShouldMapNegativeInfinityToMinValue()
+ {
+ Parse("{\"at\":\"-infinity\"}").At.Should().Be(DateTime.MinValue);
+ }
+
+ [TestMethod]
+ public void Many_ShouldParseArrayOfTimestamps()
+ {
+ Parse("{\"many\":[\"2023-01-01T00:00:00Z\",\"2024-02-02T00:00:00Z\"]}")
+ .Many.Should().HaveCount(2);
+ }
+}
diff --git a/RealtimeTests/Serialization/JoinPushConfigTests.cs b/RealtimeTests/Serialization/JoinPushConfigTests.cs
new file mode 100644
index 0000000..91200b5
--- /dev/null
+++ b/RealtimeTests/Serialization/JoinPushConfigTests.cs
@@ -0,0 +1,68 @@
+using System.Collections.Generic;
+using FluentAssertions;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using Newtonsoft.Json;
+using Newtonsoft.Json.Linq;
+using Supabase.Realtime.Broadcast;
+using Supabase.Realtime.Channel;
+using Supabase.Realtime.PostgresChanges;
+using Supabase.Realtime.Presence;
+
+namespace RealtimeTests.Serialization;
+
+///
+/// The config block a channel sends when it joins: whether it is flagged private, which
+/// postgres_changes listeners it carries, and that absent broadcast/presence sections are omitted rather
+/// than sent as null.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class JoinPushConfigTests
+{
+ private static JObject Config(JoinPush joinPush) =>
+ (JObject) JObject.Parse(JsonConvert.SerializeObject(joinPush))["config"]!;
+
+ [TestMethod]
+ public void ForPrivateChannel_ShouldFlagConfigPrivate()
+ {
+ Config(JoinPush.ForPrivateChannel())["private"]!.Value().Should().BeTrue();
+ }
+
+ [TestMethod]
+ public void ForPublicChannel_ShouldNotFlagConfigPrivate()
+ {
+ Config(JoinPush.ForPublicChannel())["private"]!.Value().Should().BeFalse();
+ }
+
+ [TestMethod]
+ public void ForPublicChannel_ShouldOmitBroadcastAndPresence_GivenAbsent()
+ {
+ var config = Config(JoinPush.ForPublicChannel());
+ config.ContainsKey("broadcast").Should().BeFalse();
+ config.ContainsKey("presence").Should().BeFalse();
+ }
+
+ [TestMethod]
+ public void ForPublicChannel_ShouldCarryPostgresChangesListeners()
+ {
+ var options = new List { new("public", "todos") };
+ var config = Config(JoinPush.ForPublicChannel(postgresChangesOptions: options));
+ config["postgres_changes"]!.Should().HaveCount(1);
+ config["postgres_changes"]![0]!["table"]!.Value().Should().Be("todos");
+ }
+
+ [TestMethod]
+ public void ForPublicChannel_ShouldSerialiseBroadcast_GivenProvided()
+ {
+ var config = Config(JoinPush.ForPublicChannel(new BroadcastOptions(broadcastSelf: true, broadcastAck: true)));
+ config["broadcast"]!["self"]!.Value().Should().BeTrue();
+ config["broadcast"]!["ack"]!.Value().Should().BeTrue();
+ }
+
+ [TestMethod]
+ public void ForPublicChannel_ShouldSerialisePresenceKey_GivenProvided()
+ {
+ var config = Config(JoinPush.ForPublicChannel(presenceOptions: new PresenceOptions("client-1")));
+ config["presence"]!["key"]!.Value().Should().Be("client-1");
+ }
+}
diff --git a/RealtimeTests/Serialization/PostgresChangesOptionsTests.cs b/RealtimeTests/Serialization/PostgresChangesOptionsTests.cs
deleted file mode 100644
index a98c9f8..0000000
--- a/RealtimeTests/Serialization/PostgresChangesOptionsTests.cs
+++ /dev/null
@@ -1,29 +0,0 @@
-using Microsoft.VisualStudio.TestTools.UnitTesting;
-using Newtonsoft.Json;
-using Newtonsoft.Json.Linq;
-using Supabase.Realtime.PostgresChanges;
-
-namespace RealtimeTests;
-
-///
-/// Wire-shape of as it is serialized into a channel's
-/// join config.postgres_changes payload.
-///
-[TestClass]
-public class PostgresChangesOptionsTests
-{
- [TestMethod]
- public void Table_ShouldBeOmittedFromJson_WhenNull()
- {
- // A schema-wide listener has no table; the server expects the key absent, not `"table": null`.
- var json = JsonConvert.SerializeObject(new PostgresChangesOptions("public"));
- Assert.IsFalse(JObject.Parse(json).ContainsKey("table"));
- }
-
- [TestMethod]
- public void Table_ShouldBeSerialized_WhenProvided()
- {
- var json = JsonConvert.SerializeObject(new PostgresChangesOptions("public", "todos"));
- Assert.AreEqual("todos", JObject.Parse(json)["table"]?.Value());
- }
-}
diff --git a/RealtimeTests/Serialization/SocketResponseSerializationTests.cs b/RealtimeTests/Serialization/SocketResponseSerializationTests.cs
new file mode 100644
index 0000000..fec31aa
--- /dev/null
+++ b/RealtimeTests/Serialization/SocketResponseSerializationTests.cs
@@ -0,0 +1,73 @@
+using FluentAssertions;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+using Newtonsoft.Json;
+using RealtimeTests.Support;
+using Supabase.Realtime.Socket;
+using static Supabase.Realtime.Constants;
+
+namespace RealtimeTests.Serialization;
+
+///
+/// How a raw Phoenix frame maps onto the typed the channel switches
+/// on, and how a postgres_changes payload's type and errors fields deserialize. These pin
+/// the wire contract the whole dispatch path depends on.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class SocketResponseSerializationTests
+{
+ [TestMethod]
+ [DataRow(ChannelEventPresenceState, EventType.PresenceState)]
+ [DataRow(ChannelEventPresenceDiff, EventType.PresenceDiff)]
+ [DataRow(ChannelEventBroadcast, EventType.Broadcast)]
+ [DataRow(ChannelEventPostgresChanges, EventType.PostgresChanges)]
+ [DataRow(ChannelEventSystem, EventType.System)]
+ [DataRow(ChannelEventReply, EventType.PostgresChanges)]
+ public void Event_ShouldMapPhoenixEventNameToTypedEvent(string phoenixEvent, EventType expected)
+ {
+ var response = Wire.Decode($"{{\"topic\":\"realtime:x\",\"event\":\"{phoenixEvent}\",\"ref\":\"1\"}}");
+ response.Event.Should().Be(expected);
+ }
+
+ [TestMethod]
+ public void Event_ShouldFallBackToPayloadType_GivenUnrecognisedEventName()
+ {
+ var response = Wire.Decode("{\"topic\":\"realtime:x\",\"event\":\"unhandled\",\"payload\":{\"type\":\"INSERT\"},\"ref\":\"1\"}");
+ response.Event.Should().Be(EventType.Insert);
+ }
+
+ [TestMethod]
+ public void Event_ShouldBeUnknown_GivenNoEventAndNoPayloadType()
+ {
+ var response = Wire.Decode("{\"topic\":\"realtime:x\",\"event\":\"unhandled\",\"ref\":\"1\"}");
+ response.Event.Should().Be(EventType.Unknown);
+ }
+
+ [TestMethod]
+ [DataRow("INSERT", EventType.Insert)]
+ [DataRow("UPDATE", EventType.Update)]
+ [DataRow("DELETE", EventType.Delete)]
+ [DataRow("TRUNCATE", EventType.Unknown)]
+ public void PayloadType_ShouldMapActionString(string action, EventType expected)
+ {
+ var payload = JsonConvert.DeserializeObject($"{{\"type\":\"{action}\"}}");
+ payload!.Type.Should().Be(expected);
+ }
+
+ [TestMethod]
+ public void PayloadErrors_ShouldDeserialize_GivenPresent()
+ {
+ const string json =
+ "{\"schema\":\"public\",\"table\":\"todos\",\"type\":\"UPDATE\",\"errors\":[\"Error 413: Payload Too Large\"]}";
+ var payload = JsonConvert.DeserializeObject(json);
+ payload!.Errors.Should().ContainSingle().Which.Should().Be("Error 413: Payload Too Large");
+ }
+
+ [TestMethod]
+ public void PayloadErrors_ShouldBeNull_GivenAbsent()
+ {
+ const string json = "{\"schema\":\"public\",\"table\":\"todos\",\"type\":\"UPDATE\",\"errors\":null}";
+ var payload = JsonConvert.DeserializeObject(json);
+ payload!.Errors.Should().BeNull();
+ }
+}
diff --git a/RealtimeTests/SocketResponseTests.cs b/RealtimeTests/SocketResponseTests.cs
deleted file mode 100644
index 1e0483a..0000000
--- a/RealtimeTests/SocketResponseTests.cs
+++ /dev/null
@@ -1,25 +0,0 @@
-using Microsoft.VisualStudio.TestTools.UnitTesting;
-using Newtonsoft.Json;
-using Supabase.Realtime.Socket;
-
-namespace RealtimeTests;
-
-[TestClass]
-public class SocketResponseTests
-{
- [TestMethod("Error Response Included and Deserialized.")]
- public void SocketResponseIncludesError()
- {
- var responseWithError =
- "{\"columns\":[{\"name\":\"id\",\"type\":\"int8\"},{\"name\":\"details\",\"type\":\"text\"}],\"commit_timestamp\":\"2021-12-28T23:59:38.984538+00:00\",\"schema\":\"public\",\"table\":\"todos\",\"type\":\"UPDATE\",\"old_record\":{\"details\":\"previous test\",\"id\":12,\"user_id\":1},\"record\":{\"details\":\"test...\",\"id\":12,\"user_id\":1},\"errors\":[\"Error 413: Payload Too Large\"]}";
- var errorResponse = JsonConvert.DeserializeObject(responseWithError);
-
- CollectionAssert.Contains(errorResponse?.Errors, "Error 413: Payload Too Large");
-
- var responseWithoutError =
- "{\"columns\":[{\"name\":\"id\",\"type\":\"int8\"},{\"name\":\"details\",\"type\":\"text\"}],\"commit_timestamp\":\"2021-12-28T23:59:38.984538+00:00\",\"schema\":\"public\",\"table\":\"todos\",\"type\":\"UPDATE\",\"old_record\":{\"details\":\"previous test\",\"id\":12,\"user_id\":1},\"record\":{\"details\":\"test...\",\"id\":12,\"user_id\":1},\"errors\": null}";
- var successResponse = JsonConvert.DeserializeObject(responseWithoutError);
-
- Assert.IsNull(successResponse?.Errors);
- }
-}
\ No newline at end of file
diff --git a/RealtimeTests/Support/Wire.cs b/RealtimeTests/Support/Wire.cs
new file mode 100644
index 0000000..b5bc46a
--- /dev/null
+++ b/RealtimeTests/Support/Wire.cs
@@ -0,0 +1,51 @@
+using Newtonsoft.Json;
+using Newtonsoft.Json.Converters;
+using Supabase.Realtime;
+using Supabase.Realtime.Channel;
+using Supabase.Realtime.Interfaces;
+using Supabase.Realtime.Socket;
+
+namespace RealtimeTests.Support;
+
+///
+/// Shared builders for the hermetic channel tests. The seam is a real, unconnected
+/// : constructing it opens no connection, so a channel can be exercised in
+/// process without a live Phoenix server and without a hand-written stub. Inbound routing is driven by
+/// calling RealtimeChannel.HandleSocketMessage with a frame from ; the
+/// subscribe/send handshake (which needs the server to reply) stays in the E2E tier.
+///
+internal static class Wire
+{
+ private const string Endpoint = "ws://127.0.0.1:54321/realtime/v1";
+
+ public static JsonSerializerSettings Settings() => new()
+ {
+ ContractResolver = new CustomContractResolver(),
+ Converters = { new IsoDateTimeConverter { DateTimeFormat = @"yyyy'-'MM'-'dd' 'HH':'mm':'ss.FFFFFFK" } },
+ MissingMemberHandling = MissingMemberHandling.Ignore
+ };
+
+ public static ChannelOptions PublicOptions() =>
+ ChannelOptions.Public(new ClientOptions(), () => null, Settings());
+
+ public static ChannelOptions PrivateOptions() =>
+ ChannelOptions.Private(new ClientOptions(), () => null, Settings());
+
+ public static IRealtimeSocket Socket() => new RealtimeSocket(Endpoint, new ClientOptions());
+
+ public static RealtimeChannel Channel(string topic = "realtime:example", ChannelOptions? options = null) =>
+ new(Socket(), topic, options ?? PublicOptions());
+
+ ///
+ /// Reproduces the client's default decoder: populate a from the frame and
+ /// stamp the raw JSON, which the typed re-parsing in broadcast/presence/postgres_changes relies on.
+ ///
+ public static SocketResponse Decode(string json)
+ {
+ var settings = Settings();
+ var response = new SocketResponse(settings);
+ JsonConvert.PopulateObject(json, response, settings);
+ response.Json = json;
+ return response;
+ }
+}
diff --git a/RealtimeTests/TestConventions.cs b/RealtimeTests/TestConventions.cs
new file mode 100644
index 0000000..73f002d
--- /dev/null
+++ b/RealtimeTests/TestConventions.cs
@@ -0,0 +1,57 @@
+using System;
+using System.Linq;
+using System.Reflection;
+using System.Text.RegularExpressions;
+using FluentAssertions;
+using FluentAssertions.Execution;
+using Microsoft.VisualStudio.TestTools.UnitTesting;
+
+namespace RealtimeTests;
+
+///
+/// Mechanized guardrails for the suite itself: these fail the build when a test class or method drifts
+/// from the documented conventions (CONVENTIONS §5.6), so the rules are enforced deterministically
+/// instead of relying on review. The two rules below cannot be expressed in .editorconfig.
+///
+[TestClass]
+[TestCategory("Unit")]
+public class TestConventions
+{
+ private static readonly Regex NamePattern =
+ new("^[A-Za-z][A-Za-z0-9]*_Should[A-Za-z0-9]+(_Given[A-Za-z0-9]+)?$", RegexOptions.Compiled);
+
+ private static readonly Type[] TestClasses = typeof(TestConventions).Assembly.GetTypes()
+ .Where(type => type.GetCustomAttribute() != null)
+ .ToArray();
+
+ [TestMethod]
+ public void TestClass_ShouldDeclareExactlyOneTestCategory()
+ {
+ using (new AssertionScope())
+ {
+ var tierBearingClasses = TestClasses.Where(type =>
+ type.GetMethods().Any(method => method.GetCustomAttribute() != null));
+ foreach (var type in tierBearingClasses)
+ {
+ type.GetCustomAttributes().Should().ContainSingle(
+ $"{type.Name} must carry exactly one [TestCategory] tier (Unit/E2E)");
+ }
+ }
+ }
+
+ [TestMethod]
+ public void TestMethod_ShouldFollowNamingConvention()
+ {
+ using (new AssertionScope())
+ {
+ var testMethods = TestClasses
+ .SelectMany(type => type.GetMethods())
+ .Where(method => method.GetCustomAttribute() != null);
+ foreach (var method in testMethods)
+ {
+ method.Name.Should().MatchRegex(NamePattern,
+ $"{method.DeclaringType!.Name}.{method.Name} must read Sut_ShouldConsequence[_GivenScenario]");
+ }
+ }
+ }
+}
diff --git a/RealtimeTests/stryker-config.json b/RealtimeTests/stryker-config.json
new file mode 100644
index 0000000..a38dd7f
--- /dev/null
+++ b/RealtimeTests/stryker-config.json
@@ -0,0 +1,31 @@
+{
+ "stryker-config": {
+ "project": "Realtime.csproj",
+ "test-projects": [
+ "RealtimeTests.csproj"
+ ],
+ "target-framework": "net8.0",
+ "reporters": [
+ "json",
+ "progress"
+ ],
+ "mutate": [
+ "Utils.cs",
+ "CustomContractResolver.cs",
+ "Converters/IntArrayConverter.cs",
+ "Converters/StringArrayConverter.cs",
+ "Converters/DateTimeConverter.cs",
+ "Channel/JoinPush.cs",
+ "Channel/ChannelOptions.cs",
+ "Exceptions/FailureHint.cs",
+ "Exceptions/RealtimeException.cs",
+ "Broadcast/BroadcastOptions.cs",
+ "Presence/PresenceOptions.cs",
+ "PostgresChanges/PostgresChangesOptions.cs",
+ "PostgresChanges/PostgresChangesFilter.cs",
+ "PostgresChanges/PostgresChangesResponse.cs",
+ "Socket/SocketResponse.cs",
+ "Socket/SocketResponsePayload.cs"
+ ]
+ }
+}