1- using System . Text . Json ;
1+ using Google . Protobuf . WellKnownTypes ;
2+ using MessagePack ;
23using Microsoft . AspNetCore . SignalR ;
34using Microsoft . EntityFrameworkCore ;
45using OpenShock . Common . Hubs ;
56using OpenShock . Common . Models . WebSocket ;
67using OpenShock . Common . OpenShockDb ;
78using OpenShock . Common . Redis ;
9+ using OpenShock . Common . Redis . PubSub ;
810using OpenShock . Common . Services . RedisPubSub ;
911using OpenShock . Common . Utils ;
1012using Redis . OM . Contracts ;
1113using StackExchange . Redis ;
14+ using System . Text . Json ;
15+ using static System . Runtime . InteropServices . JavaScript . JSType ;
1216
1317namespace OpenShock . API . Realtime ;
1418
@@ -21,6 +25,7 @@ public sealed class RedisSubscriberService : IHostedService, IAsyncDisposable
2125 private readonly IDbContextFactory < OpenShockContext > _dbContextFactory ;
2226 private readonly IRedisConnectionProvider _redisConnectionProvider ;
2327 private readonly ISubscriber _subscriber ;
28+ private readonly ILogger < RedisSubscriberService > _logger ;
2429
2530 /// <summary>
2631 /// DI Constructor
@@ -29,34 +34,78 @@ public sealed class RedisSubscriberService : IHostedService, IAsyncDisposable
2934 /// <param name="hubContext"></param>
3035 /// <param name="dbContextFactory"></param>
3136 /// <param name="redisConnectionProvider"></param>
37+ /// <param name="logger"></param>
3238 public RedisSubscriberService (
3339 IConnectionMultiplexer connectionMultiplexer ,
3440 IHubContext < UserHub , IUserHub > hubContext ,
3541 IDbContextFactory < OpenShockContext > dbContextFactory ,
36- IRedisConnectionProvider redisConnectionProvider )
42+ IRedisConnectionProvider redisConnectionProvider ,
43+ ILogger < RedisSubscriberService > logger
44+ )
3745 {
3846 _hubContext = hubContext ;
3947 _dbContextFactory = dbContextFactory ;
4048 _redisConnectionProvider = redisConnectionProvider ;
4149 _subscriber = connectionMultiplexer . GetSubscriber ( ) ;
50+ _logger = logger ;
4251 }
4352
4453 /// <inheritdoc />
4554 public async Task StartAsync ( CancellationToken cancellationToken )
4655 {
4756 await _subscriber . SubscribeAsync ( RedisChannels . KeyEventExpired , ( _ , message ) => OsTask . Run ( ( ) => HandleKeyExpired ( message ) ) ) ;
48- await _subscriber . SubscribeAsync ( RedisChannels . DeviceStatus , ( _ , message ) => OsTask . Run ( ( ) => HandleDeviceStatus ( message ) ) ) ;
57+ await _subscriber . SubscribeAsync ( RedisChannels . DeviceStatus , ProcessDeviceStatusEvent ) ;
4958 }
5059
51- private async Task HandleDeviceStatus ( RedisValue message )
60+ private void ProcessDeviceStatusEvent ( RedisChannel _ , RedisValue value )
5261 {
53- if ( ! message . HasValue ) return ;
54- var data = JsonSerializer . Deserialize < DeviceUpdatedMessage > ( message . ToString ( ) ) ;
55- if ( data is null ) return ;
62+ if ( ! value . HasValue ) return ;
63+
64+ DeviceStatus message ;
65+ try
66+ {
67+ message = MessagePackSerializer . Deserialize < DeviceStatus > ( ( ReadOnlyMemory < byte > ) value ) ;
68+ if ( message is null ) return ;
69+ }
70+ catch ( Exception e )
71+ {
72+ _logger . LogError ( e , "Failed to deserialize redis message" ) ;
73+ return ;
74+ }
5675
57- await LogicDeviceOnlineStatus ( data . Id ) ;
76+ OsTask . Run ( ( ) => HandleDeviceStatusMessage ( message ) ) ;
5877 }
59-
78+
79+ private async Task HandleDeviceStatusMessage ( DeviceStatus message )
80+ {
81+ switch ( message . Payload )
82+ {
83+ case DeviceBoolStatePayload boolState :
84+ await HandleDeviceBoolState ( message . DeviceId , boolState ) ;
85+ break ;
86+ default :
87+ _logger . LogError ( "Got DeviceStatus with unknown payload type: {PayloadType}" , message . Payload ? . GetType ( ) . Name ) ;
88+ break ;
89+ }
90+
91+ }
92+
93+ private async Task HandleDeviceBoolState ( Guid deviceId , DeviceBoolStatePayload state )
94+ {
95+ switch ( state . Type )
96+ {
97+ case DeviceBoolStateType . Online :
98+ await LogicDeviceOnlineStatus ( deviceId ) ; // TODO: Handle device offline messages too
99+ break ;
100+ case DeviceBoolStateType . EStopped :
101+ _logger . LogWarning ( "This is not yet implemented" ) ;
102+ break ;
103+ default :
104+ _logger . LogError ( "Unknown DeviceBoolStateType: {StateType}" , state . Type ) ;
105+ break ;
106+ }
107+ }
108+
60109 private async Task HandleKeyExpired ( RedisValue message )
61110 {
62111 if ( ! message . HasValue ) return ;
0 commit comments