diff --git a/packages/Realtime/Realtime.Tests/PostgresChanges/PostgresChangesDeliveryTests.cs b/packages/Realtime/Realtime.Tests/PostgresChanges/PostgresChangesDeliveryTests.cs index b19211a1..75e16d3e 100644 --- a/packages/Realtime/Realtime.Tests/PostgresChanges/PostgresChangesDeliveryTests.cs +++ b/packages/Realtime/Realtime.Tests/PostgresChanges/PostgresChangesDeliveryTests.cs @@ -1,11 +1,11 @@ using System; -using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; using Microsoft.VisualStudio.TestTools.UnitTesting; using Realtime.Tests.Models; using Supabase.Postgrest.Interfaces; using Supabase.Realtime; +using Supabase.Realtime.Exceptions; using Supabase.Realtime.Interfaces; using Supabase.Realtime.PostgresChanges; using static Supabase.Realtime.Constants; @@ -28,38 +28,70 @@ public class PostgresChangesDeliveryTests [TestInitialize] public async Task InitializeTest() { - restClient = Helpers.RestClient(); - socketClient = Helpers.SocketClient(); - await socketClient.ConnectAsync(); + this.restClient = Helpers.RestClient(); + this.socketClient = Helpers.SocketClient(); + await this.socketClient.ConnectAsync(); } [TestCleanup] - public void CleanupTest() => socketClient.Disconnect(); + public void CleanupTest() => this.socketClient.Disconnect(); [TestMethod] public async Task OnPostgresChange_ShouldModelPayload() { var tsc = new TaskCompletionSource(); - var channel = socketClient.Channel("example"); + var channel = this.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? ✅" }); + await this.restClient.Table().Insert(new Todo { UserId = 1, Details = "Client Models a response? ✅" }); Assert.IsTrue(await tsc.Task); } + [TestMethod] + public async Task OnPostgresChange_ShouldThrowError_GivenOnChangeAfterSubscribe() + { + var tsc = new TaskCompletionSource(); + var channel = this.socketClient.Channel("example"); + channel.OnPostgresChange((_, changes) => + { + var model = changes.Model(); + tsc.SetResult(model != null); + }, ListenType.Inserts, new PostgresChangesFilter { Table = "*" }); + + await channel.Subscribe(); + + var act = () => channel.OnPostgresChange((_, changes) => + { + var model = changes.Model(); + tsc.SetResult(model != null); + }, ListenType.Inserts, new PostgresChangesFilter { Table = "*" }); + + Assert.Throws(act); + } + + [TestMethod] + public async Task OnPostgresChange_ShouldThrowError_GivenRegisterChangesAfterSubscribe() + { + var channel = this.socketClient.Channel("example"); + await channel.Subscribe(); + + var act = () => channel.RegisterPostgresChangesOptions(new PostgresChangesOptions("example")); + Assert.Throws(act); + } + [TestMethod] public async Task OnPostgresChange_ShouldReceiveInsert() { var tsc = new TaskCompletionSource(); - var channel = socketClient.Channel("realtime", "public", "todos"); + var channel = this.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? ✅" }); + await this.restClient.Table().Insert(new Todo { UserId = 1, Details = "Client receives insert callback? ✅" }); Assert.IsTrue(await tsc.Task); } @@ -67,7 +99,7 @@ public async Task OnPostgresChange_ShouldReceiveInsert() public async Task OnPostgresChange_ShouldReceiveFilteredInsert() { var tsc = new TaskCompletionSource(); - var channel = socketClient.Channel("realtime", "public", "todos"); + var channel = this.socketClient.Channel("realtime", "public", "todos"); channel.OnPostgresChange((_, changes) => { Assert.AreEqual("Client receives filtered insert callback? ✅", changes.Model()?.Details); @@ -75,8 +107,8 @@ public async Task OnPostgresChange_ShouldReceiveFilteredInsert() }, 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? ✅" }); + await this.restClient.Table().Insert(new Todo { UserId = 1, Details = "Client receives insert callback? ✅" }); + await this.restClient.Table().Insert(new Todo { UserId = 2, Details = "Client receives filtered insert callback? ✅" }); Assert.IsTrue(await tsc.Task); } @@ -84,14 +116,14 @@ public async Task OnPostgresChange_ShouldReceiveFilteredInsert() public async Task OnPostgresChange_ShouldReceiveUpdateAndFilteredInsert() { var tsc = new TaskCompletionSource(); - var response = await restClient.Table() + var response = await this.restClient.Table() .Insert(new Todo { UserId = 1, Details = "Client receives insert callback? ✅" }); - await restClient.Table() + await this.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"); + var channel = this.socketClient.Channel("realtime", "public", "todos"); channel.OnPostgresChange((_, changes) => { Assert.AreEqual(oldDetails, changes.OldModel()?.Details); @@ -111,7 +143,7 @@ await restClient.Table() 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(); + await this.restClient.Table().Set(x => x.Details!, newDetails).Match(model).Update(); Assert.IsTrue(await tsc.Task); } @@ -119,12 +151,12 @@ await restClient.Table() public async Task OnPostgresChange_ShouldReceiveUpdate() { var tsc = new TaskCompletionSource(); - var response = await restClient.Table() + var response = await this.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"); + var channel = this.socketClient.Channel("realtime", "public", "todos"); channel.OnPostgresChange((_, changes) => { Assert.AreEqual(oldDetails, changes.OldModel()?.Details); @@ -138,7 +170,7 @@ public async Task OnPostgresChange_ShouldReceiveUpdate() tsc.SetResult(true); }, ListenType.Updates, new PostgresChangesFilter { Table = "todos" }); await channel.Subscribe(); - await restClient.Table().Set(x => x.Details!, newDetails).Match(model).Update(); + await this.restClient.Table().Set(x => x.Details!, newDetails).Match(model).Update(); Assert.IsTrue(await tsc.Task); } @@ -146,13 +178,13 @@ public async Task OnPostgresChange_ShouldReceiveUpdate() public async Task OnPostgresChange_ShouldReceiveDelete() { var tsc = new TaskCompletionSource(); - var channel = socketClient.Channel("realtime", "public", "todos"); + var channel = this.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 result = await this.restClient.Table().Get(); var model = result.Models.Last(); - await restClient.Table().Match(model).Delete(); + await this.restClient.Table().Match(model).Delete(); Assert.IsTrue(await tsc.Task); } @@ -160,10 +192,10 @@ public async Task OnPostgresChange_ShouldReceiveDelete() 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? ✅" }); + var channel = this.socketClient.Channel("realtime", "public", "todos"); + var todo1 = await this.restClient.Table().Insert(new Todo { UserId = 1, Details = "Client receives callbacks 1? ✅" }); + var todo2 = await this.restClient.Table().Insert(new Todo { UserId = 2, Details = "Client receives callbacks 2? ✅" }); + await this.restClient.Table().Insert(new Todo { UserId = 3, Details = "Client receives callbacks 3? ✅" }); channel.OnPostgresChange((_, removed) => { var result = removed.OldModel(); @@ -173,8 +205,8 @@ public async Task OnPostgresChange_ShouldReceiveFilteredDelete() }, 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(); + await this.restClient.Table().Match(todo1.Models.First()).Delete(); + await this.restClient.Table().Match(todo2.Models.First()).Delete(); Assert.IsTrue(await tsc.Task); } @@ -184,7 +216,7 @@ 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"); + var channel = this.socketClient.Channel("realtime", "public", "todos"); channel.OnPostgresChange((_, changes) => { switch (changes.Payload?.Data?.Type) @@ -195,10 +227,10 @@ public async Task OnPostgresChange_ShouldReceiveAllEvents_GivenWildcard() } }, 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 inserted = await this.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 this.restClient.Table().Set(x => x.Details!, "And edits.").Match(newModel).Update(); + await this.restClient.Table().Match(newModel).Delete(); await Task.WhenAll(insertTsc.Task, updateTsc.Task, deleteTsc.Task); Assert.IsTrue(insertTsc.Task.Result); Assert.IsTrue(updateTsc.Task.Result); @@ -213,7 +245,7 @@ public async Task OnPostgresChange_ShouldFanOutToMultipleInsertListeners() 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 channel = this.socketClient.Channel("realtime", "public", "todos"); var count = 0; channel.OnPostgresChange((_, _) => { @@ -227,9 +259,9 @@ public async Task OnPostgresChange_ShouldFanOutToMultipleInsertListeners() 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 this.restClient.Table().Insert(new Todo { UserId = 1, Details = "Client receives wildcard callbacks? ✅" }); + await this.restClient.Table().Insert(new Todo { UserId = 1, Details = filter1 }); + await this.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); @@ -240,11 +272,11 @@ public async Task OnPostgresChange_ShouldFanOutToMultipleInsertListeners() public async Task OnPostgresChange_ShouldRegisterAndDeliver_GivenChainedSubscribe() { var tsc = new TaskCompletionSource(); - await socketClient.Channel("public:todos") + await this.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? ✅" }); + await this.restClient.Table().Insert(new Todo { UserId = 1, Details = "OnPostgresChange receives insert? ✅" }); Assert.IsTrue(await WithinTimeout(tsc.Task)); } diff --git a/packages/Realtime/Realtime/Exceptions/FailureHint.cs b/packages/Realtime/Realtime/Exceptions/FailureHint.cs index 385cb75c..229d4bdd 100644 --- a/packages/Realtime/Realtime/Exceptions/FailureHint.cs +++ b/packages/Realtime/Realtime/Exceptions/FailureHint.cs @@ -1,4 +1,3 @@ -using System; using Websocket.Client; namespace Supabase.Realtime.Exceptions; @@ -48,6 +47,11 @@ public enum Reason /// If seen, please open an issue. /// ConnectionStale, + + /// + /// Cannot make changes after subscribe + /// + StateInvalid, } /// @@ -63,7 +67,7 @@ public static Reason Parse(DisconnectionInfo info) DisconnectionType.NoMessageReceived => Reason.ConnectionStale, DisconnectionType.Lost => Reason.ConnectionLost, DisconnectionType.ByServer => Reason.Unknown, - _ => Reason.Unknown + _ => Reason.Unknown, }; } -} \ No newline at end of file +} diff --git a/packages/Realtime/Realtime/PublicAPI.Unshipped.txt b/packages/Realtime/Realtime/PublicAPI.Unshipped.txt index 5c55f925..10e94093 100644 --- a/packages/Realtime/Realtime/PublicAPI.Unshipped.txt +++ b/packages/Realtime/Realtime/PublicAPI.Unshipped.txt @@ -10,6 +10,7 @@ Supabase.Realtime.Channel.ChannelOptions.SerializerSettings.get -> System.Text.J Supabase.Realtime.Client.SerializerSettings.get -> System.Text.Json.JsonSerializerOptions! Supabase.Realtime.ClientOptions.WebSocketFactory.get -> Supabase.Realtime.Sockets.IWebSocketFactory? Supabase.Realtime.ClientOptions.WebSocketFactory.set -> void +Supabase.Realtime.Exceptions.FailureHint.Reason.StateInvalid = 7 -> Supabase.Realtime.Exceptions.FailureHint.Reason Supabase.Realtime.Interfaces.IRealtimeClient.SerializerSettings.get -> System.Text.Json.JsonSerializerOptions! Supabase.Realtime.PostgresChanges.PostgresChangesResponse.PostgresChangesResponse() -> void Supabase.Realtime.PostgresChanges.PostgresChangesResponse.PostgresChangesResponse(System.Text.Json.JsonSerializerOptions! serializerSettings) -> void diff --git a/packages/Realtime/Realtime/RealtimeChannel.cs b/packages/Realtime/Realtime/RealtimeChannel.cs index f12a3083..dc51fbf0 100644 --- a/packages/Realtime/Realtime/RealtimeChannel.cs +++ b/packages/Realtime/Realtime/RealtimeChannel.cs @@ -443,6 +443,13 @@ public IRealtimeChannel Register(PostgresChangesOptions postgresChangesOptions) /// internal void RegisterPostgresChangesOptions(PostgresChangesOptions postgresChangesOptions) { + if (this.IsJoined || this.IsJoining) + throw new RealtimeException( + $"Cannot add `postgres_changes` callbacks for {this.Topic} after `Subscribe()`.") + { + Reason = FailureHint.Reason.StateInvalid, + }; + this.PostgresChangesOptions.Add(postgresChangesOptions); this.BindPostgresChangesOptions(postgresChangesOptions); } @@ -534,7 +541,7 @@ public IRealtimeChannel Unsubscribe() /// /// Sends a `Push` request under this channel. - /// + /// /// Maintains a buffer in the event push is called prior to the channel being joined. /// /// @@ -885,7 +892,7 @@ private void BindIdPostgresChanges(PhoenixPostgresChangeResponse joinResponse) } /// - /// Try to invoke the handler properly based on event type and socket response + /// Try to invoke the handler properly based on event type and socket response /// /// ///