Skip to content
This repository was archived by the owner on Aug 7, 2026. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Realtime/Client.cs
Original file line number Diff line number Diff line change
Expand Up @@ -404,7 +404,7 @@ public RealtimeChannel Channel(string database = "realtime", string schema = "pu
var options = ChannelOptions.Public(Options, () => AccessToken, SerializerSettings);

var subscription = new RealtimeChannel(Socket!, key, options);
subscription.Register(changesOptions);
subscription.RegisterPostgresChangesOptions(changesOptions);

_subscriptions.Add(key, subscription);

Expand Down
5 changes: 4 additions & 1 deletion Realtime/Interfaces/IRealtimeChannel.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using Supabase.Realtime.Broadcast;
using System;
using Supabase.Realtime.Broadcast;
using Supabase.Realtime.Channel;
using Supabase.Realtime.Models;
using Supabase.Realtime.PostgresChanges;
Expand Down Expand Up @@ -162,6 +163,7 @@ IRealtimeChannel OnPostgresChange(PostgresChangesHandler postgresChangeHandler,
/// </summary>
/// <param name="listenType"></param>
/// <param name="postgresChangeHandler"></param>
[Obsolete("Favor OnPostgresChange instead.")]
void AddPostgresChangeHandler(ListenType listenType, PostgresChangesHandler postgresChangeHandler);

/// <summary>
Expand Down Expand Up @@ -266,6 +268,7 @@ RealtimePresence<TPresenceResponse> Register<TPresenceResponse>(string presenceK
/// </summary>
/// <param name="postgresChangesOptions"></param>
/// <returns></returns>
[Obsolete("Favor OnPostgresChange instead.")]
IRealtimeChannel Register(PostgresChangesOptions postgresChangesOptions);

/// <summary>
Expand Down
4 changes: 3 additions & 1 deletion Realtime/RealtimeChannel.cs
Original file line number Diff line number Diff line change
Expand Up @@ -352,6 +352,7 @@ public IRealtimeChannel OnPostgresChange(PostgresChangesHandler postgresChangeHa
/// </summary>
/// <param name="listenType">The type of event this callback should process.</param>
/// <param name="postgresChangeHandler"></param>
[Obsolete("Favor OnPostgresChange instead.")]
public void AddPostgresChangeHandler(ListenType listenType, PostgresChangesHandler postgresChangeHandler)
{
BindPostgresChangesHandler(listenType, postgresChangeHandler);
Expand Down Expand Up @@ -437,6 +438,7 @@ private void NotifyPostgresChanges(EventType eventType, PostgresChangesResponse
/// </summary>
/// <param name="postgresChangesOptions"></param>
/// <returns></returns>
[Obsolete("Favor OnPostgresChange instead.")]
public IRealtimeChannel Register(PostgresChangesOptions postgresChangesOptions)
{
RegisterPostgresChangesOptions(postgresChangesOptions);
Expand All @@ -448,7 +450,7 @@ public IRealtimeChannel Register(PostgresChangesOptions postgresChangesOptions)
/// and <see cref="OnPostgresChange"/> share a single registration path.
/// </summary>
/// <param name="postgresChangesOptions"></param>
private void RegisterPostgresChangesOptions(PostgresChangesOptions postgresChangesOptions)
internal void RegisterPostgresChangesOptions(PostgresChangesOptions postgresChangesOptions)
{
PostgresChangesOptions.Add(postgresChangesOptions);
BindPostgresChangesOptions(postgresChangesOptions);
Expand Down
107 changes: 51 additions & 56 deletions RealtimeTests/ChannelPostgresChangesTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -34,18 +34,17 @@ public void CleanupTest()
_socketClient!.Disconnect();
}

[TestMethod("Channel: Payload returns a modeled response (if possible)")]
[TestMethod(DisplayName = "Channel: Payload returns a modeled response (if possible)")]
public async Task ChannelPayloadReturnsModel()
{
var tsc = new TaskCompletionSource<bool>();

var channel = _socketClient!.Channel("example");
channel.Register(new PostgresChangesOptions("public", "*"));
channel.AddPostgresChangeHandler(ListenType.Inserts, (_, changes) =>
channel.OnPostgresChange((_, changes) =>
{
var model = changes.Model<Todo>();
tsc.SetResult(model != null);
});
}, ListenType.Inserts, new PostgresChangesFilter { Table = "*" });

await channel.Subscribe();

Expand All @@ -55,14 +54,15 @@ public async Task ChannelPayloadReturnsModel()
Assert.IsTrue(check);
}

[TestMethod("Channel: Receives Insert Callback")]
[TestMethod(DisplayName = "Channel: Receives Insert Callback")]
public async Task ChannelReceivesInsertCallback()
{
var tsc = new TaskCompletionSource<bool>();

var channel = _socketClient!.Channel("realtime", "public", "todos");

channel.AddPostgresChangeHandler(ListenType.Inserts, (_, _) => tsc.SetResult(true));
channel.OnPostgresChange((_, _) => tsc.SetResult(true), ListenType.Inserts,
new PostgresChangesFilter { Table = "todos" });

await channel.Subscribe();
await _restClient!.Table<Todo>()
Expand All @@ -72,35 +72,35 @@ public async Task ChannelReceivesInsertCallback()
Assert.IsTrue(check);
}

[TestMethod("Channel: Receives Filtered Insert Callback")]
[TestMethod(DisplayName = "Channel: Receives Filtered Insert Callback")]
public async Task ChannelReceivesInsertCallbackFiltered()
{
var tsc = new TaskCompletionSource<bool>();

var channel = _socketClient!.Channel("realtime", "public", "todos", "details",
"Client receives filtered insert callback? ✅");

channel.AddPostgresChangeHandler(ListenType.Inserts, (_, changes) =>
var channel = _socketClient!.Channel("realtime", "public", "todos");

channel.OnPostgresChange((_, changes) =>
{
var oldModel = changes.Model<Todo>();

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<Todo>()
.Insert(new Todo { UserId = 1, Details = "Client receives insert callback? ✅" });

await _restClient!.Table<Todo>()
.Insert(new Todo { UserId = 2, Details = "Client receives filtered insert callback? ✅" });

var check = await tsc.Task;
Assert.IsTrue(check);
}

[TestMethod("Channel: Receives Filtered Two Callback")]
[TestMethod(DisplayName = "Channel: Receives Filtered Two Callback")]
public async Task ChannelReceivesTwoCallbacks()
{
var tsc = new TaskCompletionSource<bool>();
Expand All @@ -115,7 +115,7 @@ public async Task ChannelReceivesTwoCallbacks()
var newDetails = $"I'm an updated item ✏️ - {DateTime.Now}";

var channel = _socketClient!.Channel("realtime", "public", "todos");
channel.AddPostgresChangeHandler(ListenType.Updates, (_, changes) =>
channel.OnPostgresChange((_, changes) =>
{
var oldModel = changes.OldModel<Todo>();

Expand All @@ -131,19 +131,18 @@ public async Task ChannelReceivesTwoCallbacks()
}

tsc.SetResult(true);
});
}, ListenType.Updates, new PostgresChangesFilter { Table = "todos" });

const string filter = "Client receives filtered insert callback? ✅";
channel.Register(new PostgresChangesOptions("public", "todos", ListenType.Inserts, $"details=eq.{filter}"));
channel.AddPostgresChangeHandler(ListenType.Inserts, (_, changes) =>
channel.OnPostgresChange((_, changes) =>
{
var insertedModel = changes.Model<Todo>();

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<Todo>()
Expand All @@ -154,8 +153,8 @@ await _restClient.Table<Todo>()
var check = await tsc.Task;
Assert.IsTrue(check);
}
[TestMethod("Channel: Receives Update Callback")]

[TestMethod(DisplayName = "Channel: Receives Update Callback")]
public async Task ChannelReceivesUpdateCallback()
{
var tsc = new TaskCompletionSource<bool>();
Expand All @@ -169,7 +168,7 @@ public async Task ChannelReceivesUpdateCallback()

var channel = _socketClient!.Channel("realtime", "public", "todos");

channel.AddPostgresChangeHandler(ListenType.Updates, (_, changes) =>
channel.OnPostgresChange((_, changes) =>
{
var oldModel = changes.OldModel<Todo>();

Expand All @@ -185,7 +184,7 @@ public async Task ChannelReceivesUpdateCallback()
}

tsc.SetResult(true);
});
}, ListenType.Updates, new PostgresChangesFilter { Table = "todos" });

await channel.Subscribe();

Expand All @@ -198,14 +197,15 @@ await _restClient.Table<Todo>()
Assert.IsTrue(check);
}

[TestMethod("Channel: Receives Delete Callback")]
[TestMethod(DisplayName = "Channel: Receives Delete Callback")]
public async Task ChannelReceivesDeleteCallback()
{
var tsc = new TaskCompletionSource<bool>();

var channel = _socketClient!.Channel("realtime", "public", "todos");

channel.AddPostgresChangeHandler(ListenType.Deletes, (_, _) => tsc.SetResult(true));
channel.OnPostgresChange((_, _) => tsc.SetResult(true), ListenType.Deletes,
new PostgresChangesFilter { Table = "todos" });

await channel.Subscribe();

Expand All @@ -218,28 +218,28 @@ public async Task ChannelReceivesDeleteCallback()
Assert.IsTrue(check);
}

[TestMethod("Channel: Receives Delete Callback")]
[TestMethod(DisplayName = "Channel: Receives Delete Callback")]
public async Task ChannelReceivesFilteredDeleteCallback()
{
var tsc = new TaskCompletionSource<bool>();
var channel = _socketClient!.Channel("realtime", "public", "todos");

var todo1 = await _restClient!.Table<Todo>().Insert(new Todo
{ UserId = 1, Details = "Client receives callbacks 1? ✅" });
var todo2 = await _restClient!.Table<Todo>().Insert(new Todo
{ UserId = 2, Details = "Client receives callbacks 2? ✅" });
await _restClient!.Table<Todo>().Insert(new Todo
{ UserId = 3, Details = "Client receives callbacks 3? ✅" });

channel.Register(new PostgresChangesOptions("public", "todos", ListenType.Deletes, $"details=eq.{todo1.Model?.Details}"));
channel.AddPostgresChangeHandler(ListenType.Deletes, (_, removed) =>

channel.OnPostgresChange((_, removed) =>
{
var result = removed.OldModel<Todo>();
var result = removed.OldModel<Todo>();
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();

Expand All @@ -249,8 +249,8 @@ public async Task ChannelReceivesFilteredDeleteCallback()
var check = await tsc.Task;
Assert.IsTrue(check);
}
[TestMethod("Channel: Receives '*' Callback")]

[TestMethod(DisplayName = "Channel: Receives '*' Callback")]
public async Task ChannelReceivesWildcardCallback()
{
var insertTsc = new TaskCompletionSource<bool>();
Expand All @@ -261,7 +261,7 @@ public async Task ChannelReceivesWildcardCallback()

var channel = _socketClient!.Channel("realtime", "public", "todos");

channel.AddPostgresChangeHandler(ListenType.All, (_, changes) =>
channel.OnPostgresChange((_, changes) =>
{
switch (changes.Payload?.Data?.Type)
{
Expand All @@ -275,7 +275,7 @@ public async Task ChannelReceivesWildcardCallback()
deleteTsc.SetResult(true);
break;
}
});
}, ListenType.All, new PostgresChangesFilter { Table = "todos" });

await channel.Subscribe();

Expand All @@ -293,49 +293,44 @@ public async Task ChannelReceivesWildcardCallback()
Assert.IsTrue(deleteTsc.Task.Result);
}

[TestMethod("Channel: Receives Several Same Callback")]
[TestMethod(DisplayName = "Channel: Receives Several Same Callback")]
public async Task ChannelReceivesSeveralSameCallback()
{
var insertTask1 = new TaskCompletionSource<bool>();
var insertTask2 = new TaskCompletionSource<bool>();
var insertTask3 = new TaskCompletionSource<bool>();
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.Register(new PostgresChangesOptions("public", "todos", ListenType.Inserts));
channel.AddPostgresChangeHandler(ListenType.Inserts, (_, added) =>
channel.OnPostgresChange((_, added) =>
{
count++;
if (count == 3) insertTask1.TrySetResult(true);
});
}, ListenType.Inserts, new PostgresChangesFilter { Table = "todos" });

channel.Register(new PostgresChangesOptions("public", "todos", ListenType.Inserts, $"details=eq.{filter1}"));
channel.AddPostgresChangeHandler(ListenType.Inserts, (_, added) =>
channel.OnPostgresChange((_, added) =>
{
var model = added.Model<Todo>();

insertTask2.SetResult(model?.Details == filter1);
});

insertTask2.SetResult(model?.Details == filter1);
}, ListenType.Inserts, new PostgresChangesFilter { Table = "todos", Filter = $"details=eq.{filter1}" });

channel.Register(new PostgresChangesOptions("public", "todos", ListenType.Inserts, $"details=eq.{filter2}"));
channel.AddPostgresChangeHandler(ListenType.Inserts, (_, added) =>
channel.OnPostgresChange((_, added) =>
{
var model = added.Model<Todo>();

insertTask3.SetResult(model?.Details == filter2);
});

}, ListenType.Inserts, new PostgresChangesFilter { Table = "todos", Filter = $"details=eq.{filter2}" });

await channel.Subscribe();

await _restClient!.Table<Todo>().Insert(new Todo { UserId = 1, Details = "Client receives wildcard callbacks? ✅" });
await _restClient!.Table<Todo>().Insert(new Todo { UserId = 1, Details = filter1 });
await _restClient!.Table<Todo>().Insert(new Todo { UserId = 1, Details = filter2 });

await Task.WhenAll(insertTask1.Task, insertTask2.Task, insertTask3.Task);

Assert.IsTrue(insertTask1.Task.Result);
Expand Down Expand Up @@ -366,4 +361,4 @@ private static async Task<bool> WithinTimeout(Task<bool> task, int timeoutMs = 1
return completed == task && task.Result;
}

}
}
Loading