mirror of
https://github.com/zoriya/Kyoo.git
synced 2025-06-03 05:34:23 -04:00
Change rabbit channel from fanout to topic based
This commit is contained in:
parent
f1d72cb480
commit
cbb05ac977
@ -32,43 +32,31 @@ public class RabbitProducer
|
|||||||
{
|
{
|
||||||
_channel = rabbitConnection.CreateModel();
|
_channel = rabbitConnection.CreateModel();
|
||||||
|
|
||||||
_channel.ExchangeDeclare(exchange: "events.resource.collection", type: ExchangeType.Fanout);
|
_channel.ExchangeDeclare("events.resource", ExchangeType.Topic);
|
||||||
IRepository<Collection>.OnCreated += _Publish<Collection>("events.resource.collection", "created");
|
_ListenResourceEvents<Collection>("events.resource");
|
||||||
IRepository<Collection>.OnEdited += _Publish<Collection>("events.resource.collection", "edited");
|
_ListenResourceEvents<Movie>("events.resource");
|
||||||
IRepository<Collection>.OnDeleted += _Publish<Collection>("events.resource.collection", "deleted");
|
_ListenResourceEvents<Show>("events.resource");
|
||||||
|
_ListenResourceEvents<Season>("events.resource");
|
||||||
_channel.ExchangeDeclare(exchange: "events.resource.movie", type: ExchangeType.Fanout);
|
_ListenResourceEvents<Episode>("events.resource");
|
||||||
IRepository<Movie>.OnCreated += _Publish<Movie>("events.resource.movie", "created");
|
_ListenResourceEvents<Studio>("events.resource");
|
||||||
IRepository<Movie>.OnEdited += _Publish<Movie>("events.resource.movie", "edited");
|
_ListenResourceEvents<User>("events.resource");
|
||||||
IRepository<Movie>.OnDeleted += _Publish<Movie>("events.resource.movie", "deleted");
|
|
||||||
|
|
||||||
_channel.ExchangeDeclare(exchange: "events.resource.show", type: ExchangeType.Fanout);
|
|
||||||
IRepository<Show>.OnCreated += _Publish<Show>("events.resource.show", "created");
|
|
||||||
IRepository<Show>.OnEdited += _Publish<Show>("events.resource.show", "edited");
|
|
||||||
IRepository<Show>.OnDeleted += _Publish<Show>("events.resource.show", "deleted");
|
|
||||||
|
|
||||||
_channel.ExchangeDeclare(exchange: "events.resource.season", type: ExchangeType.Fanout);
|
|
||||||
IRepository<Season>.OnCreated += _Publish<Season>("events.resource.season", "created");
|
|
||||||
IRepository<Season>.OnEdited += _Publish<Season>("events.resource.season", "edited");
|
|
||||||
IRepository<Season>.OnDeleted += _Publish<Season>("events.resource.season", "deleted");
|
|
||||||
|
|
||||||
_channel.ExchangeDeclare(exchange: "events.resource.episode", type: ExchangeType.Fanout);
|
|
||||||
IRepository<Episode>.OnCreated += _Publish<Episode>("events.resource.episode", "created");
|
|
||||||
IRepository<Episode>.OnEdited += _Publish<Episode>("events.resource.episode", "edited");
|
|
||||||
IRepository<Episode>.OnDeleted += _Publish<Episode>("events.resource.episode", "deleted");
|
|
||||||
|
|
||||||
_channel.ExchangeDeclare(exchange: "events.resource.studio", type: ExchangeType.Fanout);
|
|
||||||
IRepository<Studio>.OnCreated += _Publish<Studio>("events.resource.studio", "created");
|
|
||||||
IRepository<Studio>.OnEdited += _Publish<Studio>("events.resource.studio", "edited");
|
|
||||||
IRepository<Studio>.OnDeleted += _Publish<Studio>("events.resource.studio", "deleted");
|
|
||||||
|
|
||||||
_channel.ExchangeDeclare(exchange: "events.resource.user", type: ExchangeType.Fanout);
|
|
||||||
IRepository<User>.OnCreated += _Publish<User>("events.resource.user", "created");
|
|
||||||
IRepository<User>.OnEdited += _Publish<User>("events.resource.user", "edited");
|
|
||||||
IRepository<User>.OnDeleted += _Publish<User>("events.resource.user", "deleted");
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private IRepository<T>.ResourceEventHandler _Publish<T>(string exchange, string action)
|
private void _ListenResourceEvents<T>(string exchange)
|
||||||
|
where T : IResource, IQuery
|
||||||
|
{
|
||||||
|
string type = typeof(T).Name.ToLowerInvariant();
|
||||||
|
|
||||||
|
IRepository<T>.OnCreated += _Publish<T>(exchange, type, "created");
|
||||||
|
IRepository<T>.OnEdited += _Publish<T>(exchange, type, "edited");
|
||||||
|
IRepository<T>.OnDeleted += _Publish<T>(exchange, type, "deleted");
|
||||||
|
}
|
||||||
|
|
||||||
|
private IRepository<T>.ResourceEventHandler _Publish<T>(
|
||||||
|
string exchange,
|
||||||
|
string type,
|
||||||
|
string action
|
||||||
|
)
|
||||||
where T : IResource, IQuery
|
where T : IResource, IQuery
|
||||||
{
|
{
|
||||||
return (T resource) =>
|
return (T resource) =>
|
||||||
@ -76,12 +64,12 @@ public class RabbitProducer
|
|||||||
var message = new
|
var message = new
|
||||||
{
|
{
|
||||||
Action = action,
|
Action = action,
|
||||||
Type = typeof(T).Name.ToLowerInvariant(),
|
Type = type,
|
||||||
Value = resource,
|
Resource = resource,
|
||||||
};
|
};
|
||||||
_channel.BasicPublish(
|
_channel.BasicPublish(
|
||||||
exchange,
|
exchange,
|
||||||
routingKey: string.Empty,
|
routingKey: $"{type}.{action}",
|
||||||
body: Encoding.UTF8.GetBytes(JsonSerializer.Serialize(message))
|
body: Encoding.UTF8.GetBytes(JsonSerializer.Serialize(message))
|
||||||
);
|
);
|
||||||
return Task.CompletedTask;
|
return Task.CompletedTask;
|
||||||
|
Loading…
x
Reference in New Issue
Block a user