Browse Source

refactor and add some comment.

master
yangxiaodong 7 years ago
parent
commit
c9cc55913a
17 changed files with 41 additions and 32 deletions
  1. +2
    -4
      src/DotNetCore.CAP.Kafka/CAP.KafkaCapOptionsExtension.cs
  2. +9
    -0
      src/DotNetCore.CAP.Kafka/CAP.Options.Extensions.cs
  3. +3
    -0
      src/DotNetCore.CAP.Kafka/CAP.SubscribeAttribute.cs
  4. +1
    -1
      src/DotNetCore.CAP.Kafka/KafkaConsumerClient.cs
  5. +3
    -3
      src/DotNetCore.CAP.Kafka/KafkaConsumerClientFactory.cs
  6. +3
    -4
      src/DotNetCore.CAP.Kafka/PublishQueueExecutor.cs
  7. +1
    -1
      src/DotNetCore.CAP.MySql/CAP.MySqlCapOptionsExtension.cs
  8. +1
    -1
      src/DotNetCore.CAP.MySql/FetchedMessage.cs
  9. +1
    -1
      src/DotNetCore.CAP.MySql/IAdditionalProcessor.Default.cs
  10. +1
    -1
      src/DotNetCore.CAP.RabbitMQ/CAP.Options.Extensions.cs
  11. +4
    -7
      src/DotNetCore.CAP.RabbitMQ/CAP.RabbitMQCapOptionsExtension.cs
  12. +3
    -0
      src/DotNetCore.CAP.RabbitMQ/CAP.SubscribeAttribute.cs
  13. +3
    -3
      src/DotNetCore.CAP.RabbitMQ/PublishQueueExecutor.cs
  14. +1
    -1
      src/DotNetCore.CAP.RabbitMQ/RabbitMQConsumerClient.cs
  15. +3
    -3
      src/DotNetCore.CAP.RabbitMQ/RabbitMQConsumerClientFactory.cs
  16. +1
    -1
      src/DotNetCore.CAP.SqlServer/CAP.SqlServerCapOptionsExtension.cs
  17. +1
    -1
      src/DotNetCore.CAP.SqlServer/FetchedMessage.cs

+ 2
- 4
src/DotNetCore.CAP.Kafka/CAP.KafkaCapOptionsExtension.cs View File

@@ -5,7 +5,7 @@ using Microsoft.Extensions.DependencyInjection;
// ReSharper disable once CheckNamespace
namespace DotNetCore.CAP
{
public class KafkaCapOptionsExtension : ICapOptionsExtension
internal sealed class KafkaCapOptionsExtension : ICapOptionsExtension
{
private readonly Action<KafkaOptions> _configure;

@@ -16,10 +16,8 @@ namespace DotNetCore.CAP

public void AddServices(IServiceCollection services)
{
services.Configure(_configure);

var kafkaOptions = new KafkaOptions();
_configure(kafkaOptions);
_configure?.Invoke(kafkaOptions);
services.AddSingleton(kafkaOptions);

services.AddSingleton<IConsumerClientFactory, KafkaConsumerClientFactory>();


+ 9
- 0
src/DotNetCore.CAP.Kafka/CAP.Options.Extensions.cs View File

@@ -6,6 +6,10 @@ namespace Microsoft.Extensions.DependencyInjection
{
public static class CapOptionsExtensions
{
/// <summary>
/// Configuration to use kafka in CAP.
/// </summary>
/// <param name="bootstrapServers">Kafka bootstrap server urls.</param>
public static CapOptions UseKafka(this CapOptions options, string bootstrapServers)
{
return options.UseKafka(opt =>
@@ -14,6 +18,11 @@ namespace Microsoft.Extensions.DependencyInjection
});
}

/// <summary>
/// Configuration to use kafka in CAP.
/// </summary>
/// <param name="configure">Provides programmatic configuration for the kafka .</param>
/// <returns></returns>
public static CapOptions UseKafka(this CapOptions options, Action<KafkaOptions> configure)
{
if (configure == null) throw new ArgumentNullException(nameof(configure));


+ 3
- 0
src/DotNetCore.CAP.Kafka/CAP.SubscribeAttribute.cs View File

@@ -3,6 +3,9 @@
// ReSharper disable once CheckNamespace
namespace DotNetCore.CAP
{
/// <summary>
/// An attribute for subscribe Kafka messages.
/// </summary>
public class CapSubscribeAttribute : TopicAttribute
{
public CapSubscribeAttribute(string name)


+ 1
- 1
src/DotNetCore.CAP.Kafka/KafkaConsumerClient.cs View File

@@ -6,7 +6,7 @@ using Confluent.Kafka.Serialization;

namespace DotNetCore.CAP.Kafka
{
public class KafkaConsumerClient : IConsumerClient
internal sealed class KafkaConsumerClient : IConsumerClient
{
private readonly string _groupId;
private readonly KafkaOptions _kafkaOptions;


+ 3
- 3
src/DotNetCore.CAP.Kafka/KafkaConsumerClientFactory.cs View File

@@ -3,13 +3,13 @@ using Microsoft.Extensions.Options;

namespace DotNetCore.CAP.Kafka
{
public class KafkaConsumerClientFactory : IConsumerClientFactory
internal sealed class KafkaConsumerClientFactory : IConsumerClientFactory
{
private readonly KafkaOptions _kafkaOptions;

public KafkaConsumerClientFactory(IOptions<KafkaOptions> kafkaOptions)
public KafkaConsumerClientFactory(KafkaOptions kafkaOptions)
{
_kafkaOptions = kafkaOptions?.Value ?? throw new ArgumentNullException(nameof(kafkaOptions));
_kafkaOptions = kafkaOptions;
}

public IConsumerClient Create(string groupId)


+ 3
- 4
src/DotNetCore.CAP.Kafka/PublishQueueExecutor.cs View File

@@ -4,22 +4,21 @@ using System.Threading.Tasks;
using Confluent.Kafka;
using DotNetCore.CAP.Processor.States;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;

namespace DotNetCore.CAP.Kafka
{
public class PublishQueueExecutor : BasePublishQueueExecutor
internal class PublishQueueExecutor : BasePublishQueueExecutor
{
private readonly ILogger _logger;
private readonly KafkaOptions _kafkaOptions;

public PublishQueueExecutor(IStateChanger stateChanger,
IOptions<KafkaOptions> options,
KafkaOptions options,
ILogger<PublishQueueExecutor> logger)
: base(stateChanger, logger)
{
_logger = logger;
_kafkaOptions = options.Value;
_kafkaOptions = options;
}

public override Task<OperateResult> PublishAsync(string keyName, string content)


+ 1
- 1
src/DotNetCore.CAP.MySql/CAP.MySqlCapOptionsExtension.cs View File

@@ -7,7 +7,7 @@ using Microsoft.Extensions.DependencyInjection;
// ReSharper disable once CheckNamespace
namespace DotNetCore.CAP
{
public class MySqlCapOptionsExtension : ICapOptionsExtension
internal class MySqlCapOptionsExtension : ICapOptionsExtension
{
private readonly Action<MySqlOptions> _configure;



+ 1
- 1
src/DotNetCore.CAP.MySql/FetchedMessage.cs View File

@@ -2,7 +2,7 @@

namespace DotNetCore.CAP.MySql
{
public class FetchedMessage
internal class FetchedMessage
{
public int MessageId { get; set; }



+ 1
- 1
src/DotNetCore.CAP.MySql/IAdditionalProcessor.Default.cs View File

@@ -7,7 +7,7 @@ using MySql.Data.MySqlClient;

namespace DotNetCore.CAP.MySql
{
public class DefaultAdditionalProcessor : IAdditionalProcessor
internal class DefaultAdditionalProcessor : IAdditionalProcessor
{
private readonly IServiceProvider _provider;
private readonly ILogger _logger;


+ 1
- 1
src/DotNetCore.CAP.RabbitMQ/CAP.Options.Extensions.cs View File

@@ -4,7 +4,7 @@ using DotNetCore.CAP;
// ReSharper disable once CheckNamespace
namespace Microsoft.Extensions.DependencyInjection
{
public static class CapOptionsExtensions
internal static class CapOptionsExtensions
{
public static CapOptions UseRabbitMQ(this CapOptions options, string hostName)
{


+ 4
- 7
src/DotNetCore.CAP.RabbitMQ/CAP.RabbitMQCapOptionsExtension.cs View File

@@ -5,7 +5,7 @@ using Microsoft.Extensions.DependencyInjection;
// ReSharper disable once CheckNamespace
namespace DotNetCore.CAP
{
public class RabbitMQCapOptionsExtension : ICapOptionsExtension
internal sealed class RabbitMQCapOptionsExtension : ICapOptionsExtension
{
private readonly Action<RabbitMQOptions> _configure;

@@ -16,12 +16,9 @@ namespace DotNetCore.CAP

public void AddServices(IServiceCollection services)
{
services.Configure(_configure);

var rabbitMQOptions = new RabbitMQOptions();
_configure(rabbitMQOptions);

services.AddSingleton(rabbitMQOptions);
var options = new RabbitMQOptions();
_configure?.Invoke(options);
services.AddSingleton(options);

services.AddSingleton<IConsumerClientFactory, RabbitMQConsumerClientFactory>();
services.AddTransient<IQueueExecutor, PublishQueueExecutor>();


+ 3
- 0
src/DotNetCore.CAP.RabbitMQ/CAP.SubscribeAttribute.cs View File

@@ -3,6 +3,9 @@
// ReSharper disable once CheckNamespace
namespace DotNetCore.CAP
{
/// <summary>
/// An attribute for subscribe RabbitMQ messages.
/// </summary>
public class CapSubscribeAttribute : TopicAttribute
{
public CapSubscribeAttribute(string name) : base(name)


+ 3
- 3
src/DotNetCore.CAP.RabbitMQ/PublishQueueExecutor.cs View File

@@ -8,18 +8,18 @@ using RabbitMQ.Client;

namespace DotNetCore.CAP.RabbitMQ
{
public class PublishQueueExecutor : BasePublishQueueExecutor
internal sealed class PublishQueueExecutor : BasePublishQueueExecutor
{
private readonly ILogger _logger;
private readonly RabbitMQOptions _rabbitMQOptions;

public PublishQueueExecutor(IStateChanger stateChanger,
IOptions<RabbitMQOptions> options,
RabbitMQOptions options,
ILogger<PublishQueueExecutor> logger)
: base(stateChanger, logger)
{
_logger = logger;
_rabbitMQOptions = options.Value;
_rabbitMQOptions = options;
}

public override Task<OperateResult> PublishAsync(string keyName, string content)


+ 1
- 1
src/DotNetCore.CAP.RabbitMQ/RabbitMQConsumerClient.cs View File

@@ -7,7 +7,7 @@ using RabbitMQ.Client.Events;

namespace DotNetCore.CAP.RabbitMQ
{
public class RabbitMQConsumerClient : IConsumerClient
internal sealed class RabbitMQConsumerClient : IConsumerClient
{
private readonly string _exchageName;
private readonly string _queueName;


+ 3
- 3
src/DotNetCore.CAP.RabbitMQ/RabbitMQConsumerClientFactory.cs View File

@@ -2,13 +2,13 @@

namespace DotNetCore.CAP.RabbitMQ
{
public class RabbitMQConsumerClientFactory : IConsumerClientFactory
internal sealed class RabbitMQConsumerClientFactory : IConsumerClientFactory
{
private readonly RabbitMQOptions _rabbitMQOptions;

public RabbitMQConsumerClientFactory(IOptions<RabbitMQOptions> rabbitMQOptions)
public RabbitMQConsumerClientFactory(RabbitMQOptions rabbitMQOptions)
{
_rabbitMQOptions = rabbitMQOptions.Value;
_rabbitMQOptions = rabbitMQOptions;
}

public IConsumerClient Create(string groupId)


+ 1
- 1
src/DotNetCore.CAP.SqlServer/CAP.SqlServerCapOptionsExtension.cs View File

@@ -7,7 +7,7 @@ using Microsoft.Extensions.DependencyInjection;
// ReSharper disable once CheckNamespace
namespace DotNetCore.CAP
{
public class SqlServerCapOptionsExtension : ICapOptionsExtension
internal class SqlServerCapOptionsExtension : ICapOptionsExtension
{
private readonly Action<SqlServerOptions> _configure;



+ 1
- 1
src/DotNetCore.CAP.SqlServer/FetchedMessage.cs View File

@@ -2,7 +2,7 @@

namespace DotNetCore.CAP.SqlServer
{
public class FetchedMessage
internal class FetchedMessage
{
public int MessageId { get; set; }



Loading…
Cancel
Save