♻️ To use CorePush to no longer depends on gorush
This commit is contained in:
@@ -91,7 +91,7 @@ public class RegistryProxyConfigProvider : IProxyConfigProvider, IDisposable
|
|||||||
{
|
{
|
||||||
{ "destination1", new DestinationConfig { Address = serviceUrl } }
|
{ "destination1", new DestinationConfig { Address = serviceUrl } }
|
||||||
},
|
},
|
||||||
HttpRequest = new ForwarderRequestConfig()
|
HttpRequest = new ForwarderRequestConfig
|
||||||
{
|
{
|
||||||
ActivityTimeout = directRoute.IsWebsocket ? TimeSpan.FromHours(24) : TimeSpan.FromMinutes(2)
|
ActivityTimeout = directRoute.IsWebsocket ? TimeSpan.FromHours(24) : TimeSpan.FromMinutes(2)
|
||||||
}
|
}
|
||||||
|
@@ -8,6 +8,7 @@
|
|||||||
</PropertyGroup>
|
</PropertyGroup>
|
||||||
|
|
||||||
<ItemGroup>
|
<ItemGroup>
|
||||||
|
<PackageReference Include="CorePush" Version="4.3.0" />
|
||||||
<PackageReference Include="EFCore.BulkExtensions" Version="9.0.1" />
|
<PackageReference Include="EFCore.BulkExtensions" Version="9.0.1" />
|
||||||
<PackageReference Include="EFCore.BulkExtensions.PostgreSql" Version="9.0.1" />
|
<PackageReference Include="EFCore.BulkExtensions.PostgreSql" Version="9.0.1" />
|
||||||
<PackageReference Include="EFCore.NamingConventions" Version="9.0.0" />
|
<PackageReference Include="EFCore.NamingConventions" Version="9.0.0" />
|
||||||
|
@@ -1,5 +1,6 @@
|
|||||||
using System.Text;
|
using CorePush.Apple;
|
||||||
using System.Text.Json;
|
using CorePush.Firebase;
|
||||||
|
using DysonNetwork.Pusher.Connection;
|
||||||
using DysonNetwork.Shared.Proto;
|
using DysonNetwork.Shared.Proto;
|
||||||
using EFCore.BulkExtensions;
|
using EFCore.BulkExtensions;
|
||||||
using Microsoft.EntityFrameworkCore;
|
using Microsoft.EntityFrameworkCore;
|
||||||
@@ -7,14 +8,55 @@ using NodaTime;
|
|||||||
|
|
||||||
namespace DysonNetwork.Pusher.Notification;
|
namespace DysonNetwork.Pusher.Notification;
|
||||||
|
|
||||||
public class PushService(IConfiguration config, AppDatabase db, IHttpClientFactory httpFactory)
|
public class PushService
|
||||||
{
|
{
|
||||||
private readonly string _notifyTopic = config["Notifications:Topic"]!;
|
private readonly AppDatabase _db;
|
||||||
private readonly Uri _notifyEndpoint = new(config["Notifications:Endpoint"]!);
|
private readonly WebSocketService _ws;
|
||||||
|
private readonly ILogger<PushService> _logger;
|
||||||
|
private readonly FirebaseSender? _fcm;
|
||||||
|
private readonly ApnSender? _apns;
|
||||||
|
private readonly string? _apnsTopic;
|
||||||
|
|
||||||
|
public PushService(
|
||||||
|
IConfiguration config,
|
||||||
|
AppDatabase db,
|
||||||
|
WebSocketService ws,
|
||||||
|
IHttpClientFactory httpFactory,
|
||||||
|
ILogger<PushService> logger
|
||||||
|
)
|
||||||
|
{
|
||||||
|
var cfgSection = config.GetSection("Notifications:Push");
|
||||||
|
|
||||||
|
// Set up Firebase Cloud Messaging
|
||||||
|
var fcmConfig = cfgSection.GetValue<string>("Google");
|
||||||
|
if (fcmConfig != null && File.Exists(fcmConfig))
|
||||||
|
_fcm = new FirebaseSender(File.ReadAllText(fcmConfig), httpFactory.CreateClient());
|
||||||
|
|
||||||
|
// Set up Apple Push Notification Service
|
||||||
|
var apnsKeyPath = cfgSection.GetValue<string>("Apple:PrivateKey");
|
||||||
|
if (apnsKeyPath != null && File.Exists(apnsKeyPath))
|
||||||
|
{
|
||||||
|
_apns = new ApnSender(new ApnSettings
|
||||||
|
{
|
||||||
|
P8PrivateKey = File.ReadAllText(apnsKeyPath),
|
||||||
|
P8PrivateKeyId = cfgSection.GetValue<string>("Apple:PrivateKeyId"),
|
||||||
|
TeamId = cfgSection.GetValue<string>("Apple:TeamId"),
|
||||||
|
AppBundleIdentifier = cfgSection.GetValue<string>("Apple:BundleIdentifier"),
|
||||||
|
ServerType = cfgSection.GetValue<bool>("Production")
|
||||||
|
? ApnServerType.Production
|
||||||
|
: ApnServerType.Development
|
||||||
|
}, httpFactory.CreateClient());
|
||||||
|
_apnsTopic = cfgSection.GetValue<string>("Apple:BundleIdentifier");
|
||||||
|
}
|
||||||
|
|
||||||
|
_db = db;
|
||||||
|
_ws = ws;
|
||||||
|
_logger = logger;
|
||||||
|
}
|
||||||
|
|
||||||
public async Task UnsubscribeDevice(string deviceId)
|
public async Task UnsubscribeDevice(string deviceId)
|
||||||
{
|
{
|
||||||
await db.PushSubscriptions
|
await _db.PushSubscriptions
|
||||||
.Where(s => s.DeviceId == deviceId)
|
.Where(s => s.DeviceId == deviceId)
|
||||||
.ExecuteDeleteAsync();
|
.ExecuteDeleteAsync();
|
||||||
}
|
}
|
||||||
@@ -27,41 +69,40 @@ public class PushService(IConfiguration config, AppDatabase db, IHttpClientFacto
|
|||||||
)
|
)
|
||||||
{
|
{
|
||||||
var now = SystemClock.Instance.GetCurrentInstant();
|
var now = SystemClock.Instance.GetCurrentInstant();
|
||||||
|
|
||||||
// First check if a matching subscription exists
|
|
||||||
var accountId = Guid.Parse(account.Id!);
|
var accountId = Guid.Parse(account.Id!);
|
||||||
var existingSubscription = await db.PushSubscriptions
|
|
||||||
|
// Check for existing subscription with same device ID or token
|
||||||
|
var existingSubscription = await _db.PushSubscriptions
|
||||||
.Where(s => s.AccountId == accountId)
|
.Where(s => s.AccountId == accountId)
|
||||||
.Where(s => s.DeviceId == deviceId || s.DeviceToken == deviceToken)
|
.Where(s => s.DeviceId == deviceId || s.DeviceToken == deviceToken)
|
||||||
.FirstOrDefaultAsync();
|
.FirstOrDefaultAsync();
|
||||||
|
|
||||||
if (existingSubscription is not null)
|
if (existingSubscription != null)
|
||||||
{
|
{
|
||||||
// Update the existing subscription directly in the database
|
// Update existing subscription
|
||||||
await db.PushSubscriptions
|
|
||||||
.Where(s => s.Id == existingSubscription.Id)
|
|
||||||
.ExecuteUpdateAsync(setters => setters
|
|
||||||
.SetProperty(s => s.DeviceId, deviceId)
|
|
||||||
.SetProperty(s => s.DeviceToken, deviceToken)
|
|
||||||
.SetProperty(s => s.UpdatedAt, now));
|
|
||||||
|
|
||||||
// Return the updated subscription
|
|
||||||
existingSubscription.DeviceId = deviceId;
|
existingSubscription.DeviceId = deviceId;
|
||||||
existingSubscription.DeviceToken = deviceToken;
|
existingSubscription.DeviceToken = deviceToken;
|
||||||
|
existingSubscription.Provider = provider;
|
||||||
existingSubscription.UpdatedAt = now;
|
existingSubscription.UpdatedAt = now;
|
||||||
|
|
||||||
|
_db.Update(existingSubscription);
|
||||||
|
await _db.SaveChangesAsync();
|
||||||
return existingSubscription;
|
return existingSubscription;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Create new subscription
|
||||||
var subscription = new PushSubscription
|
var subscription = new PushSubscription
|
||||||
{
|
{
|
||||||
DeviceId = deviceId,
|
DeviceId = deviceId,
|
||||||
DeviceToken = deviceToken,
|
DeviceToken = deviceToken,
|
||||||
Provider = provider,
|
Provider = provider,
|
||||||
AccountId = accountId,
|
AccountId = accountId,
|
||||||
|
CreatedAt = now,
|
||||||
|
UpdatedAt = now
|
||||||
};
|
};
|
||||||
|
|
||||||
db.PushSubscriptions.Add(subscription);
|
_db.PushSubscriptions.Add(subscription);
|
||||||
await db.SaveChangesAsync();
|
await _db.SaveChangesAsync();
|
||||||
|
|
||||||
return subscription;
|
return subscription;
|
||||||
}
|
}
|
||||||
@@ -94,8 +135,8 @@ public class PushService(IConfiguration config, AppDatabase db, IHttpClientFacto
|
|||||||
|
|
||||||
if (save)
|
if (save)
|
||||||
{
|
{
|
||||||
db.Add(notification);
|
_db.Add(notification);
|
||||||
await db.SaveChangesAsync();
|
await _db.SaveChangesAsync();
|
||||||
}
|
}
|
||||||
|
|
||||||
if (!isSilent) _ = DeliveryNotification(notification);
|
if (!isSilent) _ = DeliveryNotification(notification);
|
||||||
@@ -104,7 +145,7 @@ public class PushService(IConfiguration config, AppDatabase db, IHttpClientFacto
|
|||||||
public async Task DeliveryNotification(Notification notification)
|
public async Task DeliveryNotification(Notification notification)
|
||||||
{
|
{
|
||||||
// Pushing the notification
|
// Pushing the notification
|
||||||
var subscribers = await db.PushSubscriptions
|
var subscribers = await _db.PushSubscriptions
|
||||||
.Where(s => s.AccountId == notification.AccountId)
|
.Where(s => s.AccountId == notification.AccountId)
|
||||||
.ToListAsync();
|
.ToListAsync();
|
||||||
|
|
||||||
@@ -117,7 +158,7 @@ public class PushService(IConfiguration config, AppDatabase db, IHttpClientFacto
|
|||||||
var id = notifications.Where(n => n.ViewedAt == null).Select(n => n.Id).ToList();
|
var id = notifications.Where(n => n.ViewedAt == null).Select(n => n.Id).ToList();
|
||||||
if (id.Count == 0) return;
|
if (id.Count == 0) return;
|
||||||
|
|
||||||
await db.Notifications
|
await _db.Notifications
|
||||||
.Where(n => id.Contains(n.Id))
|
.Where(n => id.Contains(n.Id))
|
||||||
.ExecuteUpdateAsync(s => s.SetProperty(n => n.ViewedAt, now)
|
.ExecuteUpdateAsync(s => s.SetProperty(n => n.ViewedAt, now)
|
||||||
);
|
);
|
||||||
@@ -141,105 +182,113 @@ public class PushService(IConfiguration config, AppDatabase db, IHttpClientFacto
|
|||||||
};
|
};
|
||||||
return newNotification;
|
return newNotification;
|
||||||
}).ToList();
|
}).ToList();
|
||||||
await db.BulkInsertAsync(notifications);
|
await _db.BulkInsertAsync(notifications);
|
||||||
}
|
}
|
||||||
|
|
||||||
var subscribers = await db.PushSubscriptions
|
var subscribers = await _db.PushSubscriptions
|
||||||
.Where(s => accounts.Contains(s.AccountId))
|
.Where(s => accounts.Contains(s.AccountId))
|
||||||
.ToListAsync();
|
.ToListAsync();
|
||||||
await _PushNotification(notification, subscribers);
|
await _PushNotification(notification, subscribers);
|
||||||
}
|
}
|
||||||
|
|
||||||
private List<Dictionary<string, object>> _BuildNotificationPayload(Notification notification,
|
|
||||||
IEnumerable<PushSubscription> subscriptions)
|
|
||||||
{
|
|
||||||
var subDict = subscriptions
|
|
||||||
.GroupBy(x => x.Provider)
|
|
||||||
.ToDictionary(x => x.Key, x => x.ToList());
|
|
||||||
|
|
||||||
var notifications = subDict.Select(value =>
|
|
||||||
{
|
|
||||||
var platformCode = value.Key switch
|
|
||||||
{
|
|
||||||
PushProvider.Apple => 1,
|
|
||||||
PushProvider.Google => 2,
|
|
||||||
_ => throw new InvalidOperationException($"Unknown push provider: {value.Key}")
|
|
||||||
};
|
|
||||||
|
|
||||||
var tokens = value.Value.Select(x => x.DeviceToken).ToList();
|
|
||||||
return _BuildNotificationPayload(notification, platformCode, tokens);
|
|
||||||
}).ToList();
|
|
||||||
|
|
||||||
return notifications.ToList();
|
|
||||||
}
|
|
||||||
|
|
||||||
private Dictionary<string, object> _BuildNotificationPayload(Pusher.Notification.Notification notification,
|
|
||||||
int platformCode,
|
|
||||||
IEnumerable<string> deviceTokens)
|
|
||||||
{
|
|
||||||
var alertDict = new Dictionary<string, object>();
|
|
||||||
var dict = new Dictionary<string, object>
|
|
||||||
{
|
|
||||||
["notif_id"] = notification.Id.ToString(),
|
|
||||||
["apns_id"] = notification.Id.ToString(),
|
|
||||||
["topic"] = _notifyTopic,
|
|
||||||
["tokens"] = deviceTokens,
|
|
||||||
["data"] = new Dictionary<string, object>
|
|
||||||
{
|
|
||||||
["type"] = notification.Topic,
|
|
||||||
["meta"] = notification.Meta ?? new Dictionary<string, object>(),
|
|
||||||
},
|
|
||||||
["mutable_content"] = true,
|
|
||||||
["priority"] = notification.Priority >= 5 ? "high" : "normal",
|
|
||||||
};
|
|
||||||
|
|
||||||
if (!string.IsNullOrWhiteSpace(notification.Title))
|
|
||||||
{
|
|
||||||
dict["title"] = notification.Title;
|
|
||||||
alertDict["title"] = notification.Title;
|
|
||||||
}
|
|
||||||
|
|
||||||
if (!string.IsNullOrWhiteSpace(notification.Content))
|
|
||||||
{
|
|
||||||
dict["message"] = notification.Content;
|
|
||||||
alertDict["body"] = notification.Content;
|
|
||||||
}
|
|
||||||
|
|
||||||
if (!string.IsNullOrWhiteSpace(notification.Subtitle))
|
|
||||||
{
|
|
||||||
dict["message"] = $"{notification.Subtitle}\n{dict["message"]}";
|
|
||||||
alertDict["subtitle"] = notification.Subtitle;
|
|
||||||
}
|
|
||||||
|
|
||||||
if (notification.Priority >= 5)
|
|
||||||
dict["name"] = "default";
|
|
||||||
|
|
||||||
dict["platform"] = platformCode;
|
|
||||||
dict["alert"] = alertDict;
|
|
||||||
|
|
||||||
return dict;
|
|
||||||
}
|
|
||||||
|
|
||||||
private async Task _PushNotification(
|
private async Task _PushNotification(
|
||||||
Notification notification,
|
Notification notification,
|
||||||
IEnumerable<PushSubscription> subscriptions
|
IEnumerable<PushSubscription> subscriptions
|
||||||
)
|
)
|
||||||
{
|
{
|
||||||
var subList = subscriptions.ToList();
|
var tasks = subscriptions
|
||||||
if (subList.Count == 0) return;
|
.Select(subscription => _PushSingleNotification(notification, subscription))
|
||||||
|
.ToList();
|
||||||
|
|
||||||
var requestDict = new Dictionary<string, object>
|
await Task.WhenAll(tasks);
|
||||||
|
}
|
||||||
|
|
||||||
|
private async Task _PushSingleNotification(Notification notification, PushSubscription subscription)
|
||||||
|
{
|
||||||
|
try
|
||||||
{
|
{
|
||||||
["notifications"] = _BuildNotificationPayload(notification, subList)
|
_logger.LogDebug(
|
||||||
};
|
$"Pushing notification {notification.Topic} #{notification.Id} to device #{subscription.DeviceId}");
|
||||||
|
|
||||||
var client = httpFactory.CreateClient();
|
switch (subscription.Provider)
|
||||||
client.BaseAddress = _notifyEndpoint;
|
{
|
||||||
var request = await client.PostAsync("/push", new StringContent(
|
case PushProvider.Google:
|
||||||
JsonSerializer.Serialize(requestDict),
|
if (_fcm == null)
|
||||||
Encoding.UTF8,
|
throw new InvalidOperationException("Firebase Cloud Messaging is not initialized.");
|
||||||
"application/json"
|
|
||||||
));
|
var body = string.Empty;
|
||||||
request.EnsureSuccessStatusCode();
|
if (!string.IsNullOrEmpty(notification.Subtitle) || !string.IsNullOrEmpty(notification.Content))
|
||||||
|
{
|
||||||
|
body = string.Join("\n",
|
||||||
|
notification.Subtitle ?? string.Empty,
|
||||||
|
notification.Content ?? string.Empty).Trim();
|
||||||
|
}
|
||||||
|
|
||||||
|
await _fcm.SendAsync(new Dictionary<string, object>
|
||||||
|
{
|
||||||
|
["message"] = new Dictionary<string, object>
|
||||||
|
{
|
||||||
|
["token"] = subscription.DeviceToken,
|
||||||
|
["notification"] = new Dictionary<string, object>
|
||||||
|
{
|
||||||
|
["title"] = notification.Title ?? string.Empty,
|
||||||
|
["body"] = body
|
||||||
|
},
|
||||||
|
["data"] = new Dictionary<string, object>
|
||||||
|
{
|
||||||
|
["id"] = notification.Id,
|
||||||
|
["topic"] = notification.Topic,
|
||||||
|
["meta"] = notification.Meta ?? new Dictionary<string, object>()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
break;
|
||||||
|
|
||||||
|
case PushProvider.Apple:
|
||||||
|
if (_apns == null)
|
||||||
|
throw new InvalidOperationException("Apple Push Notification Service is not initialized.");
|
||||||
|
|
||||||
|
var alertDict = new Dictionary<string, object>();
|
||||||
|
if (!string.IsNullOrEmpty(notification.Title))
|
||||||
|
alertDict["title"] = notification.Title;
|
||||||
|
if (!string.IsNullOrEmpty(notification.Subtitle))
|
||||||
|
alertDict["subtitle"] = notification.Subtitle;
|
||||||
|
if (!string.IsNullOrEmpty(notification.Content))
|
||||||
|
alertDict["body"] = notification.Content;
|
||||||
|
|
||||||
|
var payload = new Dictionary<string, object?>
|
||||||
|
{
|
||||||
|
["topic"] = _apnsTopic,
|
||||||
|
["aps"] = new Dictionary<string, object?>
|
||||||
|
{
|
||||||
|
["alert"] = alertDict,
|
||||||
|
["sound"] = notification.Priority >= 5 ? "default" : null,
|
||||||
|
["mutable-content"] = 1
|
||||||
|
},
|
||||||
|
["meta"] = notification.Meta ?? new Dictionary<string, object>()
|
||||||
|
};
|
||||||
|
|
||||||
|
await _apns.SendAsync(
|
||||||
|
payload,
|
||||||
|
deviceToken: subscription.DeviceToken,
|
||||||
|
apnsId: notification.Id.ToString(),
|
||||||
|
apnsPriority: notification.Priority,
|
||||||
|
apnPushType: ApnPushType.Alert
|
||||||
|
);
|
||||||
|
break;
|
||||||
|
|
||||||
|
default:
|
||||||
|
throw new InvalidOperationException($"Push provider not supported: {subscription.Provider}");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
catch (Exception ex)
|
||||||
|
{
|
||||||
|
_logger.LogError(ex,
|
||||||
|
$"Failed to push notification #{notification.Id} to device {subscription.DeviceId}. {ex.Message}");
|
||||||
|
throw new Exception($"Failed to send notification to {subscription.Provider}: {ex.Message}", ex);
|
||||||
|
}
|
||||||
|
|
||||||
|
_logger.LogInformation(
|
||||||
|
$"Successfully pushed notification #{notification.Id} to device {subscription.DeviceId}");
|
||||||
}
|
}
|
||||||
}
|
}
|
@@ -23,6 +23,7 @@ builder.Services.AddAppFlushHandlers();
|
|||||||
|
|
||||||
// Add business services
|
// Add business services
|
||||||
builder.Services.AddAppBusinessServices();
|
builder.Services.AddAppBusinessServices();
|
||||||
|
builder.Services.AddPushServices(builder.Configuration);
|
||||||
|
|
||||||
// Add scheduled jobs
|
// Add scheduled jobs
|
||||||
builder.Services.AddAppScheduledJobs();
|
builder.Services.AddAppScheduledJobs();
|
||||||
|
@@ -1,5 +1,7 @@
|
|||||||
using System.Text.Json;
|
using System.Text.Json;
|
||||||
using System.Threading.RateLimiting;
|
using System.Threading.RateLimiting;
|
||||||
|
using CorePush.Apple;
|
||||||
|
using CorePush.Firebase;
|
||||||
using DysonNetwork.Pusher.Connection;
|
using DysonNetwork.Pusher.Connection;
|
||||||
using DysonNetwork.Pusher.Email;
|
using DysonNetwork.Pusher.Email;
|
||||||
using DysonNetwork.Pusher.Notification;
|
using DysonNetwork.Pusher.Notification;
|
||||||
@@ -137,4 +139,13 @@ public static class ServiceCollectionExtensions
|
|||||||
|
|
||||||
return services;
|
return services;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public static void AddPushServices(this IServiceCollection services, IConfiguration configuration)
|
||||||
|
{
|
||||||
|
services.Configure<ApnSettings>(configuration.GetSection("PushNotify:Apple"));
|
||||||
|
services.AddHttpClient<ApnSender>();
|
||||||
|
|
||||||
|
services.Configure<FirebaseSettings>(configuration.GetSection("PushNotify:Firebase"));
|
||||||
|
services.AddHttpClient<FirebaseSettings>();
|
||||||
|
}
|
||||||
}
|
}
|
@@ -14,8 +14,16 @@
|
|||||||
"Etcd": "etcd.orb.local:2379"
|
"Etcd": "etcd.orb.local:2379"
|
||||||
},
|
},
|
||||||
"Notifications": {
|
"Notifications": {
|
||||||
"Topic": "dev.solsynth.solian",
|
"Push": {
|
||||||
"Endpoint": "http://localhost:8088"
|
"Production": true,
|
||||||
|
"Google": "./Keys/Solian.json",
|
||||||
|
"Apple": {
|
||||||
|
"PrivateKey": "./Keys/Solian.p8",
|
||||||
|
"PrivateKeyId": "4US4KSX4W6",
|
||||||
|
"TeamId": "W7HPZ53V6B",
|
||||||
|
"BundleIdentifier": "dev.solsynth.solian"
|
||||||
|
}
|
||||||
|
}
|
||||||
},
|
},
|
||||||
"Email": {
|
"Email": {
|
||||||
"Server": "smtp4dev.orb.local",
|
"Server": "smtp4dev.orb.local",
|
||||||
|
Reference in New Issue
Block a user