♻️ Refactor NATS message handling

This commit is contained in:
2025-08-21 18:47:23 +08:00
parent 4d1972bc99
commit d92220b4bc
2 changed files with 20 additions and 8 deletions

View File

@@ -1,4 +1,5 @@
using System.Text.Json;
using DysonNetwork.Shared.Proto;
using NATS.Client.Core;
namespace DysonNetwork.Pusher.Services;
@@ -20,7 +21,8 @@ public class QueueService(INatsConnection nats)
Body = body
})
};
await nats.PublishAsync(QueueName, message);
var rawMessage = GrpcTypeHelper.ConvertObjectToByteString(message).ToByteArray();
await nats.PublishAsync(QueueName, rawMessage);
}
public async Task EnqueuePushNotification(Notification.Notification notification, Guid userId, bool isSavable = false)
@@ -34,8 +36,8 @@ public class QueueService(INatsConnection nats)
TargetId = userId.ToString(),
Data = JsonSerializer.Serialize(notification)
};
await nats.PublishAsync(QueueName, message);
var rawMessage = GrpcTypeHelper.ConvertObjectToByteString(message).ToByteArray();
await nats.PublishAsync(QueueName, rawMessage);
}
}