From 73fb92f436d0b0b20b0b745933a429c1dece95bc Mon Sep 17 00:00:00 2001 From: Halil ibrahim Kalkan Date: Tue, 24 Jul 2018 16:47:13 +0300 Subject: [PATCH] Refactored background job executer. --- framework/Volo.Abp.sln | 2 +- .../BackgroundJobs/BackgroundJobException.cs | 10 +- .../BackgroundJobs/BackgroundJobExecuter.cs | 114 ++++-------------- .../BackgroundJobs/BackgroundJobOptions.cs | 1 + .../Abp/BackgroundJobs/BackgroundJobWorker.cs | 77 +++++++++++- .../BackgroundJobs/IBackgroundJobExecuter.cs | 2 +- .../Abp/BackgroundJobs/JobExecutionContext.cs | 18 +++ .../Abp/BackgroundJobs/JobExecutionResult.cs | 8 ++ .../BackgroundJobExecuter_Tests.cs | 22 ++-- 9 files changed, 137 insertions(+), 117 deletions(-) create mode 100644 framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/JobExecutionContext.cs create mode 100644 framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/JobExecutionResult.cs diff --git a/framework/Volo.Abp.sln b/framework/Volo.Abp.sln index 9a5abbc32e..81a858e968 100644 --- a/framework/Volo.Abp.sln +++ b/framework/Volo.Abp.sln @@ -198,7 +198,7 @@ Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Volo.Abp.BackgroundJobs", " EndProject Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Volo.Abp.BackgroundWorkers", "src\Volo.Abp.BackgroundWorkers\Volo.Abp.BackgroundWorkers.csproj", "{6C3E76B8-C4DA-4E74-9F8B-A8BC4C831722}" EndProject -Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Volo.Abp.BackgroundJobs.Tests", "test\Volo.Abp.BackgroundJobs.Tests\Volo.Abp.BackgroundJobs.Tests.csproj", "{D86548EA-7047-4623-8824-F6285CD254AA}" +Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Volo.Abp.BackgroundJobs.Tests", "test\Volo.Abp.BackgroundJobs.Tests\Volo.Abp.BackgroundJobs.Tests.csproj", "{D86548EA-7047-4623-8824-F6285CD254AA}" EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobException.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobException.cs index 8e16d983bd..bc85df58a6 100644 --- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobException.cs +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobException.cs @@ -1,21 +1,15 @@ using System; using System.Runtime.Serialization; -using JetBrains.Annotations; namespace Volo.Abp.BackgroundJobs { [Serializable] public class BackgroundJobException : AbpException { - [CanBeNull] - public BackgroundJobInfo BackgroundJob { get; set; } + public string JobName { get; set; } - [CanBeNull] - public object JobObject { get; set; } + public string JobArgs { get; set; } - /// - /// Creates a new object. - /// public BackgroundJobException() { diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobExecuter.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobExecuter.cs index c95c1eb2aa..32e86ccb47 100644 --- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobExecuter.cs +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobExecuter.cs @@ -5,8 +5,6 @@ using Microsoft.Extensions.Logging; using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Options; using Volo.Abp.DependencyInjection; -using Volo.Abp.Threading; -using Volo.Abp.Timing; namespace Volo.Abp.BackgroundJobs { @@ -15,120 +13,58 @@ namespace Volo.Abp.BackgroundJobs public ILogger Logger { protected get; set; } protected IServiceProvider ServiceProvider { get; } - protected IClock Clock { get; } protected IBackgroundJobSerializer Serializer { get; } - protected IBackgroundJobStore Store { get; } protected BackgroundJobOptions Options { get; } public BackgroundJobExecuter( IServiceProvider serviceProvider, - IClock clock, IBackgroundJobSerializer serializer, - IBackgroundJobStore store, IOptions options) { ServiceProvider = serviceProvider; - Clock = clock; Serializer = serializer; Options = options.Value; - Store = store; Logger = NullLogger.Instance; } - public virtual void Execute(BackgroundJobInfo jobInfo) + public virtual void Execute(JobExecutionContext context) { //TODO: Refactor (split to multiple methods). - try + var jobType = Options.GetJobType(context.JobName); + + using (var scope = ServiceProvider.CreateScope()) { - jobInfo.TryCount++; - jobInfo.LastTryTime = Clock.Now; + var job = scope.ServiceProvider.GetService(jobType); + if (job == null) + { + throw new AbpException("The job type is not registered to DI: " + jobType); + } - var jobType = Options.GetJobType(jobInfo.JobName); + var jobExecuteMethod = job.GetType().GetMethod("Execute"); + Debug.Assert(jobExecuteMethod != null, nameof(jobExecuteMethod) + " != null"); + var argsType = jobExecuteMethod.GetParameters()[0].ParameterType; + var argsObj = Serializer.Deserialize(context.JobArgs, argsType); - using (var scope = ServiceProvider.CreateScope()) + try { - var job = scope.ServiceProvider.GetService(jobType); - if (job == null) - { - throw new AbpException("The job type is not registered to DI: " + jobType); - } + jobExecuteMethod.Invoke(job, new[] { argsObj }); + } + catch (Exception ex) + { + context.Result = JobExecutionResult.Failed; - var jobExecuteMethod = job.GetType().GetMethod("Execute"); - Debug.Assert(jobExecuteMethod != null, nameof(jobExecuteMethod) + " != null"); - var argsType = jobExecuteMethod.GetParameters()[0].ParameterType; - var argsObj = Serializer.Deserialize(jobInfo.JobArgs, argsType); + Logger.LogException(ex); - try + //TODO: Somehow trigger an event for the exception (may create an Volo.Abp.ExceptionHandling package)! + var backgroundJobException = new BackgroundJobException("A background job execution is failed. See inner exception for details.", ex) { - jobExecuteMethod.Invoke(job, new[] { argsObj }); - AsyncHelper.RunSync(() => Store.DeleteAsync(jobInfo.Id)); - } - catch (Exception ex) - { - Logger.LogException(ex); - - var nextTryTime = CalculateNextTryTime(jobInfo); - if (nextTryTime.HasValue) - { - jobInfo.NextTryTime = nextTryTime.Value; - } - else - { - jobInfo.IsAbandoned = true; - } - - TryUpdate(jobInfo); - - var backgroundJobException = new BackgroundJobException( - "A background job execution is failed. See inner exception for details. See BackgroundJob property to get information on the background job.", - ex - ) - { - BackgroundJob = jobInfo, - JobObject = job - }; - - //TODO: Somehow trigger an event for the exception (may create an Volo.Abp.ExceptionHandling package)! - } + JobName = context.JobName, + JobArgs = context.JobArgs + }; } } - catch (Exception ex) - { - Logger.LogException(ex); - - jobInfo.IsAbandoned = true; - - TryUpdate(jobInfo); - } - } - - protected virtual void TryUpdate(BackgroundJobInfo jobInfo) - { - try - { - Store.UpdateAsync(jobInfo); - } - catch (Exception updateEx) - { - Logger.LogException(updateEx); - } - } - - protected virtual DateTime? CalculateNextTryTime(BackgroundJobInfo jobInfo) //TODO: Move to another place to override easier - { - var nextWaitDuration = Options.DefaultFirstWaitDuration * (Math.Pow(Options.DefaultWaitFactor, jobInfo.TryCount - 1)); - var nextTryDate = jobInfo.LastTryTime.HasValue - ? jobInfo.LastTryTime.Value.AddSeconds(nextWaitDuration) - : Clock.Now.AddSeconds(nextWaitDuration); - - if (nextTryDate.Subtract(jobInfo.CreationTime).TotalSeconds > Options.DefaultTimeout) - { - return null; - } - - return nextTryDate; } } } \ No newline at end of file diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobOptions.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobOptions.cs index 4bcfeb56f8..173f194f0f 100644 --- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobOptions.cs +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobOptions.cs @@ -49,6 +49,7 @@ namespace Volo.Abp.BackgroundJobs { JobTypes = new Dictionary(); + JobPollPeriod = 5000; DefaultFirstWaitDuration = 60; DefaultTimeout = 172800; DefaultWaitFactor = 2.0; diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorker.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorker.cs index d8f2ad7367..07a52cbb50 100644 --- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorker.cs +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobWorker.cs @@ -1,7 +1,10 @@ -using Microsoft.Extensions.Options; +using System; +using Microsoft.Extensions.Logging; +using Microsoft.Extensions.Options; using Volo.Abp.BackgroundWorkers; using Volo.Abp.DependencyInjection; using Volo.Abp.Threading; +using Volo.Abp.Timing; namespace Volo.Abp.BackgroundJobs { @@ -10,7 +13,8 @@ namespace Volo.Abp.BackgroundJobs protected IBackgroundJobExecuter JobExecuter { get; } protected IBackgroundJobStore Store { get; } protected BackgroundJobOptions Options { get; } - + protected IClock Clock { get; } + /// /// Initializes a new instance of the class. /// @@ -18,10 +22,12 @@ namespace Volo.Abp.BackgroundJobs IBackgroundJobStore store, AbpTimer timer, IBackgroundJobExecuter jobExecuter, - IOptions options) + IOptions options, + IClock clock) : base(timer) { JobExecuter = jobExecuter; + Clock = clock; Store = store; Options = options.Value; Timer.Period = Options.JobPollPeriod; @@ -31,10 +37,71 @@ namespace Volo.Abp.BackgroundJobs { var waitingJobs = AsyncHelper.RunSync(() => Store.GetWaitingJobsAsync(Options.MaxJobFetchCount)); - foreach (var job in waitingJobs) + foreach (var jobInfo in waitingJobs) + { + jobInfo.TryCount++; + jobInfo.LastTryTime = Clock.Now; + + var context = new JobExecutionContext(jobInfo.JobName, jobInfo.JobArgs); + + try + { + JobExecuter.Execute(context); + + if (context.Result == JobExecutionResult.Success) + { + AsyncHelper.RunSync(() => Store.DeleteAsync(jobInfo.Id)); + } + else if (context.Result == JobExecutionResult.Failed) + { + + var nextTryTime = CalculateNextTryTime(jobInfo); + if (nextTryTime.HasValue) + { + jobInfo.NextTryTime = nextTryTime.Value; + } + else + { + jobInfo.IsAbandoned = true; + } + + TryUpdate(jobInfo); + } + } + catch (Exception ex) + { + Logger.LogException(ex); + jobInfo.IsAbandoned = true; + TryUpdate(jobInfo); + } + } + } + + protected virtual void TryUpdate(BackgroundJobInfo jobInfo) + { + try + { + Store.UpdateAsync(jobInfo); + } + catch (Exception updateEx) + { + Logger.LogException(updateEx); + } + } + + protected virtual DateTime? CalculateNextTryTime(BackgroundJobInfo jobInfo) //TODO: Move to another place to override easier + { + var nextWaitDuration = Options.DefaultFirstWaitDuration * (Math.Pow(Options.DefaultWaitFactor, jobInfo.TryCount - 1)); + var nextTryDate = jobInfo.LastTryTime.HasValue + ? jobInfo.LastTryTime.Value.AddSeconds(nextWaitDuration) + : Clock.Now.AddSeconds(nextWaitDuration); + + if (nextTryDate.Subtract(jobInfo.CreationTime).TotalSeconds > Options.DefaultTimeout) { - JobExecuter.Execute(job); + return null; } + + return nextTryDate; } } } \ No newline at end of file diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobExecuter.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobExecuter.cs index a0010952d3..2470b2945b 100644 --- a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobExecuter.cs +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobExecuter.cs @@ -2,6 +2,6 @@ { public interface IBackgroundJobExecuter { - void Execute(BackgroundJobInfo jobInfo); + void Execute(JobExecutionContext context); } } \ No newline at end of file diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/JobExecutionContext.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/JobExecutionContext.cs new file mode 100644 index 0000000000..6323c09cda --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/JobExecutionContext.cs @@ -0,0 +1,18 @@ +namespace Volo.Abp.BackgroundJobs +{ + public class JobExecutionContext + { + public string JobName { get; } + + public string JobArgs { get; } + + public JobExecutionResult Result { get; set; } + + public JobExecutionContext(string jobName, string jobArgs) + { + JobName = jobName; + JobArgs = jobArgs; + Result = JobExecutionResult.Success; + } + } +} \ No newline at end of file diff --git a/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/JobExecutionResult.cs b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/JobExecutionResult.cs new file mode 100644 index 0000000000..dca4732bb5 --- /dev/null +++ b/framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/JobExecutionResult.cs @@ -0,0 +1,8 @@ +namespace Volo.Abp.BackgroundJobs +{ + public enum JobExecutionResult + { + Success, + Failed + } +} \ No newline at end of file diff --git a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobExecuter_Tests.cs b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobExecuter_Tests.cs index f2c38d07c3..065368fdab 100644 --- a/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobExecuter_Tests.cs +++ b/framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobExecuter_Tests.cs @@ -1,5 +1,6 @@ using System.Threading.Tasks; using Shouldly; +using Volo.Abp.Json; using Xunit; namespace Volo.Abp.BackgroundJobs @@ -7,14 +8,12 @@ namespace Volo.Abp.BackgroundJobs public class BackgroundJobExecuter_Tests : BackgroundJobsTestBase { private readonly IBackgroundJobExecuter _backgroundJobExecuter; - private readonly IBackgroundJobManager _backgroundJobManager; - private readonly IBackgroundJobStore _backgroundJobStore; + private readonly IJsonSerializer _jsonSerializer; public BackgroundJobExecuter_Tests() { _backgroundJobExecuter = GetRequiredService(); - _backgroundJobManager = GetRequiredService(); - _backgroundJobStore = GetRequiredService(); + _jsonSerializer = GetRequiredService(); } [Fact] @@ -25,21 +24,18 @@ namespace Volo.Abp.BackgroundJobs var jobObject = GetRequiredService(); jobObject.ExecutedValues.ShouldBeEmpty(); - var jobId = await _backgroundJobManager.EnqueueAsync(new MyJobArgs("42")); - - var job = await _backgroundJobStore.FindAsync(jobId); - job.ShouldNotBeNull(); - //Act - _backgroundJobExecuter.Execute(job); + _backgroundJobExecuter.Execute( + new JobExecutionContext( + BackgroundJobNameAttribute.GetName(), + _jsonSerializer.Serialize(new MyJobArgs("42")) + ) + ); //Assert jobObject.ExecutedValues.ShouldContain("42"); - - job = await _backgroundJobStore.FindAsync(jobId); - job.ShouldBeNull(); //Because it's deleted after the execution } } } \ No newline at end of file