Skip to content
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 .github/workflows/nuget-publish.yml
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ jobs:
- name: Pack and Push NuGet Package
uses: simplify9/sw-workflows/actions/dotnet-pack-push@main
with:
projects: "SW.Bus/SW.Bus.csproj"
projects: "SW.Bus/SW.Bus.csproj,SW.Bus.RabbitMqExtensions/SW.Bus.RabbitMqExtensions.csproj"
configuration: "Release"
version: ${{ steps.semver.outputs.version }}
api-key: ${{ secrets.SWNUGETKEY }}
Expand Down
7 changes: 7 additions & 0 deletions SW.Bus.RabbitMqExtensions/ConsumerOptions.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
namespace SW.Bus.RabbitMqExtensions;

public class ConsumerOptions
{
public ushort? Prefetch { get; set; }
public int? Priority { get; set; }
}
8 changes: 8 additions & 0 deletions SW.Bus.RabbitMqExtensions/IConsumeExtended.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
using SW.PrimitiveTypes;

namespace SW.Bus.RabbitMqExtensions;

public interface IConsumeExtended : IConsume
{
Task<IDictionary<string,ConsumerOptions>> GetMessageTypeNamesWithOptions();
}
11 changes: 11 additions & 0 deletions SW.Bus.RabbitMqExtensions/QueueOptions.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
namespace SW.Bus.RabbitMqExtensions;

public class QueueOptions:ConsumerOptions
{
public int? RetryCount { get; set; }
public uint? RetryAfterSeconds { get; set; }
public IDictionary<string, object>? ConsumerArgs => Priority is null or 0 ? null : new Dictionary<string, object>
{
{ "x-priority", Priority},
};
}
27 changes: 27 additions & 0 deletions SW.Bus.RabbitMqExtensions/SW.Bus.RabbitMqExtensions.csproj
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
<Project Sdk="Microsoft.NET.Sdk">

<PropertyGroup>
<TargetFramework>net8.0</TargetFramework>
<ImplicitUsings>enable</ImplicitUsings>
<Nullable>enable</Nullable>
<PackageId>SimplyWorks.Bus.RabbitMqExtensions</PackageId>
<Product>SimplyWorks.Bus.RabbitMqExtensions</Product>
<Authors>Simplify9</Authors>
<Description>Extensions for SimplyWorks.Bus library to support RabbitMQ specific features.</Description>
<PackageTags>messagebus;rabbitmq;aspnetcore;messaging;pubsub;eventdriven;microservices;dotnet8</PackageTags>
<PackageLicenseExpression>MIT</PackageLicenseExpression>
<PackageProjectUrl>https://github.com/simplify9/SW-Bus</PackageProjectUrl>
<RepositoryUrl>https://github.com/simplify9/SW-Bus</RepositoryUrl>
<PackageReadmeFile>README.md</PackageReadmeFile>
<PackageIcon>icon.png</PackageIcon>
<RepositoryType>git</RepositoryType>
<Copyright>Copyright © 2020 Simplify9</Copyright>
<PackageReleaseNotes>See https://github.com/simplify9/SW-Bus/releases for release notes and changelog.</PackageReleaseNotes>

</PropertyGroup>

<ItemGroup>
<PackageReference Include="SimplyWorks.PrimitiveTypes" Version="8.1.2" />
</ItemGroup>

</Project>
2 changes: 1 addition & 1 deletion SW.Bus.SampleWeb/SW.Bus.SampleWeb.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
<ItemGroup>
<PackageReference Include="Microsoft.AspNetCore.Authentication.JwtBearer" Version="5.0.11" />
<PackageReference Include="SimplyWorks.CqApi" Version="2.0.38" />
<PackageReference Include="SimplyWorks.PrimitiveTypes" Version="6.0.5" />
<PackageReference Include="SimplyWorks.PrimitiveTypes" Version="8.1.2" />
</ItemGroup>

<ItemGroup>
Expand Down
6 changes: 6 additions & 0 deletions SW.Bus.sln
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "Solution Items", "Solution
.editorconfig = .editorconfig
EndProjectSection
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "SW.Bus.RabbitMqExtensions", "SW.Bus.RabbitMqExtensions\SW.Bus.RabbitMqExtensions.csproj", "{4C85A34C-FBDA-4BBF-938D-86A46F9E5593}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
Expand All @@ -32,6 +34,10 @@ Global
{2FB51BFB-82A8-40EF-8FEB-D9F12993F2E1}.Debug|Any CPU.Build.0 = Debug|Any CPU
{2FB51BFB-82A8-40EF-8FEB-D9F12993F2E1}.Release|Any CPU.ActiveCfg = Release|Any CPU
{2FB51BFB-82A8-40EF-8FEB-D9F12993F2E1}.Release|Any CPU.Build.0 = Release|Any CPU
{4C85A34C-FBDA-4BBF-938D-86A46F9E5593}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{4C85A34C-FBDA-4BBF-938D-86A46F9E5593}.Debug|Any CPU.Build.0 = Debug|Any CPU
{4C85A34C-FBDA-4BBF-938D-86A46F9E5593}.Release|Any CPU.ActiveCfg = Release|Any CPU
{4C85A34C-FBDA-4BBF-938D-86A46F9E5593}.Release|Any CPU.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
Expand Down
14 changes: 9 additions & 5 deletions SW.Bus/BasicPublisher.cs
Original file line number Diff line number Diff line change
Expand Up @@ -24,23 +24,23 @@ public BasicPublisher(IModel model, BusOptions busOptions, RequestContext reques
this.requestContext = requestContext;
}

async public Task Publish<TMessage>(TMessage message, string exchange)
async public Task Publish<TMessage>(TMessage message, string exchange, byte? priority = null)
{
var serializerOptions = new JsonSerializerOptions()
{
ReferenceHandler = ReferenceHandler.IgnoreCycles
};
var body = JsonSerializer.Serialize(message,message.GetType(), serializerOptions);

await Publish(message.GetType().Name, body,exchange);
await Publish(message.GetType().Name, body,exchange, priority);
}

public async Task Publish(string messageTypeName, string message,string exchange)
public async Task Publish(string messageTypeName, string message,string exchange, byte? priority = null)
{
try
{
var body = Encoding.UTF8.GetBytes(message);
await Publish(messageTypeName, body,exchange);
await Publish(messageTypeName, body,exchange, priority);
}
catch (Exception e)
{
Expand All @@ -49,11 +49,15 @@ public async Task Publish(string messageTypeName, string message,string exchange

}

public Task Publish(string messageTypeName, byte[] message,string exchange)
public Task Publish(string messageTypeName, byte[] message,string exchange, byte? priority = null)
{
IBasicProperties props = null;
props = model.CreateBasicProperties();
props.Headers = new Dictionary<string, object>();
if (priority.HasValue)
{
props.Priority = priority.Value;
}
if (requestContext.IsValid && busOptions.Token.IsValid)
{
var jwt = busOptions.Token.WriteJwt((ClaimsIdentity)requestContext.User.Identity);
Expand Down
2 changes: 2 additions & 0 deletions SW.Bus/BusOptions.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
using SW.HttpExtensions;
using System;
using System.Collections.Generic;
using SW.Bus.RabbitMqExtensions;

namespace SW.Bus
{
Expand Down Expand Up @@ -34,6 +35,7 @@ public BusOptions(string environment)
public ushort DefaultQueuePrefetch { get; set; }
public ushort DefaultRetryCount { get; set; }
public uint DefaultRetryAfter { get; set; }
public int DefaultMaxPriority { get; set; }
public string NodeId { get; set; }
public int ListenRetryCount { get; set; }
public ushort ListenRetryAfter { get; set; }
Expand Down
103 changes: 59 additions & 44 deletions SW.Bus/ConsumerDefinition.cs
Original file line number Diff line number Diff line change
@@ -1,58 +1,73 @@
using System;
using System.Collections.Generic;
using System.Reflection;
using RabbitMQ.Client.Events;
using SW.Bus.RabbitMqExtensions;

namespace SW.Bus
namespace SW.Bus;

public class ConsumerDefinition
{
public class ConsumerDefinition
private readonly string queueNamePrefix;
private readonly BusOptions busOptions;
private readonly QueueOptions queueOptions;

public ConsumerDefinition(string queueNamePrefix, BusOptions busOptions, string nakedQueueName)
{
private readonly string queueNamePrefix;
private readonly BusOptions busOptions;
private readonly QueueOptions queueOptions;
this.queueNamePrefix = queueNamePrefix;
this.busOptions = busOptions;
NakedQueueName = nakedQueueName;
busOptions.Options.TryGetValue(NakedQueueName, out queueOptions);
}

public ConsumerDefinition(string queueNamePrefix, BusOptions busOptions, string nakedQueueName)
{
this.queueNamePrefix = queueNamePrefix;
this.busOptions = busOptions;
NakedQueueName = nakedQueueName;
busOptions.Options.TryGetValue(NakedQueueName, out queueOptions);

}

public Type ServiceType { get; set; }
public Type MessageType { get; set; }
public string MessageTypeName { get; set; }
public MethodInfo Method { get; set; }
public int RetryCount => queueOptions?.RetryCount ?? busOptions.DefaultRetryCount;
public uint RetryAfter => queueOptions?.RetryAfterSeconds ?? busOptions.DefaultRetryAfter;
public ushort QueuePrefetch => queueOptions?.Prefetch ?? busOptions.DefaultQueuePrefetch;
public string NakedQueueName { get; private set; }
public string QueueName => $"{queueNamePrefix}.{NakedQueueName}".ToLower();
public string RoutingKey => MessageTypeName.ToLower();
public string RetryRoutingKey => $"{NakedQueueName}.retry".ToLower();
public string RetryQueueName => $"{queueNamePrefix}.{NakedQueueName}.retry".ToLower();
public string BadRoutingKey => $"{NakedQueueName}.bad".ToLower();
public string BadQueueName => $"{queueNamePrefix}.{NakedQueueName}.bad".ToLower();
public IDictionary<string, object> RetryArgs => RetryCount == 0 ? null : new Dictionary<string, object>
{
{ "x-dead-letter-exchange", busOptions.ProcessExchange},
{ "x-dead-letter-routing-key", RetryRoutingKey},
{ "x-message-ttl", RetryAfter == 0 ? 100 : RetryAfter * 1000 }
};
public ConsumerDefinition(string queueNamePrefix, BusOptions busOptions, ConsumerOptions consumerOptions,
string nakedQueueName): this(queueNamePrefix, busOptions, nakedQueueName)
{
ArgumentNullException.ThrowIfNull(consumerOptions);
queueOptions ??= new QueueOptions();
if (consumerOptions.Prefetch.HasValue)
queueOptions.Prefetch = consumerOptions.Prefetch;
if (consumerOptions.Priority.HasValue)
queueOptions.Priority = consumerOptions.Priority;
}

public Type ServiceType { get; set; }
public Type MessageType { get; set; }
public string MessageTypeName { get; set; }
public MethodInfo Method { get; set; }
public int RetryCount => queueOptions?.RetryCount ?? busOptions.DefaultRetryCount;
public uint RetryAfter => queueOptions?.RetryAfterSeconds ?? busOptions.DefaultRetryAfter;
public ushort QueuePrefetch => queueOptions?.Prefetch ?? busOptions.DefaultQueuePrefetch;
public string NakedQueueName { get; private set; }
public string QueueName => $"{queueNamePrefix}.{NakedQueueName}".ToLower();
public string RoutingKey => MessageTypeName.ToLower();
public string RetryRoutingKey => $"{NakedQueueName}.retry".ToLower();
public string RetryQueueName => $"{queueNamePrefix}.{NakedQueueName}.retry".ToLower();
public string BadRoutingKey => $"{NakedQueueName}.bad".ToLower();
public string BadQueueName => $"{queueNamePrefix}.{NakedQueueName}.bad".ToLower();

public IDictionary<string, object> ProcessArgs => new Dictionary<string, object>
public IDictionary<string, object> RetryArgs => RetryCount == 0
? null
: new Dictionary<string, object>
{
{ "x-dead-letter-exchange", busOptions.DeadLetterExchange },
{ "x-dead-letter-exchange", busOptions.ProcessExchange },
{ "x-dead-letter-routing-key", RetryRoutingKey },

{ "x-message-ttl", RetryAfter == 0 ? 100 : RetryAfter * 1000 }
};

public static IDictionary<string, object> BadArgs => new Dictionary<string, object>
{
{ "x-message-ttl", (uint)TimeSpan.FromDays(7).TotalMilliseconds }
};
public IDictionary<string, object> ConsumerArgs => queueOptions?.ConsumerArgs;
}
public IDictionary<string, object> ProcessArgs => new Dictionary<string, object>
{
{ "x-dead-letter-exchange", busOptions.DeadLetterExchange },
{ "x-dead-letter-routing-key", RetryRoutingKey },
};

}
public static IDictionary<string, object> BadArgs => new Dictionary<string, object>
{
{ "x-message-ttl", (uint)TimeSpan.FromDays(7).TotalMilliseconds }
};

public IDictionary<string, object> ConsumerArgs => queueOptions?.ConsumerArgs;
public int? ConsumerPriority => queueOptions?.Priority;
public string ConsumerTag { get; set; }
public AsyncEventingBasicConsumer ConsumerObject { get; set; }
}
Loading