Browse Source

Refactored background jobs.

pull/395/head
Halil ibrahim Kalkan 8 years ago
parent
commit
6b5ee29f74
  1. 119
      framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobExecuter.cs
  2. 4
      framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobInfo.cs
  3. 117
      framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobManager.cs
  4. 26
      framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobManagerExtensions.cs
  5. 30
      framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobNameAttribute.cs
  6. 19
      framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobOptions.cs
  7. 7
      framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobExecuter.cs
  8. 19
      framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobManager.cs
  9. 7
      framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobNameProvider.cs
  10. 3
      framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/JsonBackgroundJobSerializer.cs
  11. 1
      framework/test/Volo.Abp.BackgroundJobs.Tests/Volo.Abp.BackgroundJobs.Tests.csproj
  12. 2
      framework/test/Volo.Abp.BackgroundJobs.Tests/Volo.Abp.BackgroundJobs.Tests.csproj.DotSettings
  13. 16
      framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpBackgroundJobsTestModule.cs
  14. 42
      framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobManager_Tests.cs
  15. 7
      framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobsTestBase.cs
  16. 7
      framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/Class1.cs

119
framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobExecuter.cs

@ -0,0 +1,119 @@
using System;
using System.Diagnostics;
using Microsoft.Extensions.DependencyInjection;
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
{
public class BackgroundJobExecuter : IBackgroundJobExecuter, ITransientDependency
{
public ILogger<BackgroundJobExecuter> 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<BackgroundJobOptions> options)
{
ServiceProvider = serviceProvider;
Clock = clock;
Serializer = serializer;
Options = options.Value;
Store = store;
Logger = NullLogger<BackgroundJobExecuter>.Instance;
}
public void Execute(BackgroundJobInfo jobInfo)
{
try
{
jobInfo.TryCount++;
jobInfo.LastTryTime = Clock.Now;
var jobType = Options.GetJobType(jobInfo.JobName);
using (var scope = ServiceProvider.CreateScope())
{
var job = scope.ServiceProvider.GetService(jobType);
if (job == null)
{
throw new AbpException("JobName is not registered: " + jobType);
}
//TODO: Type check for the job object
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);
try
{
jobExecuteMethod.Invoke(job, new[] { argsObj });
AsyncHelper.RunSync(() => Store.DeleteAsync(jobInfo));
}
catch (Exception ex)
{
Logger.LogException(ex);
var nextTryTime = jobInfo.CalculateNextTryTime(Clock);
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)!
}
}
}
catch (Exception ex)
{
Logger.LogException(ex);
jobInfo.IsAbandoned = true;
TryUpdate(jobInfo);
}
}
private void TryUpdate(BackgroundJobInfo jobInfo)
{
try
{
Store.UpdateAsync(jobInfo);
}
catch (Exception updateEx)
{
Logger.LogException(updateEx);
}
}
}
}

4
framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobInfo.cs

@ -10,7 +10,7 @@ namespace Volo.Abp.BackgroundJobs
public class BackgroundJobInfo
{
/// <summary>
/// Maximum length of <see cref="JobType"/>.
/// Maximum length of <see cref="JobName"/>.
/// Value: 512.
/// </summary>
public const int MaxJobTypeLength = 512;
@ -48,7 +48,7 @@ namespace Volo.Abp.BackgroundJobs
/// </summary>
[Required]
[StringLength(MaxJobTypeLength)]
public virtual string JobType { get; set; }
public virtual string JobName { get; set; }
/// <summary>
/// Job arguments as JSON string.

117
framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobManager.cs

@ -1,8 +1,5 @@
using System;
using System.Diagnostics;
using System.Threading.Tasks;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Volo.Abp.BackgroundWorkers;
using Volo.Abp.DependencyInjection;
using Volo.Abp.Guids;
@ -11,6 +8,8 @@ using Volo.Abp.Timing;
namespace Volo.Abp.BackgroundJobs
{
//TODO: Split enqueueing & background worker!
/// <summary>
/// Default implementation of <see cref="IBackgroundJobManager"/>.
/// </summary>
@ -22,10 +21,10 @@ namespace Volo.Abp.BackgroundJobs
/// </summary>
public static int JobPollPeriod { get; set; } //TODO: Move to options
protected IServiceProvider ServiceProvider { get; }
protected IClock Clock { get; }
protected IBackgroundJobSerializer Serializer { get; }
protected IGuidGenerator GuidGenerator { get; }
protected IBackgroundJobExecuter JobExecuter { get; }
protected IBackgroundJobStore Store { get; }
static BackgroundJobManager()
@ -37,30 +36,35 @@ namespace Volo.Abp.BackgroundJobs
/// Initializes a new instance of the <see cref="BackgroundJobManager"/> class.
/// </summary>
public BackgroundJobManager(
IServiceProvider serviceProvider,
IClock clock,
IBackgroundJobSerializer serializer,
IBackgroundJobStore store,
IGuidGenerator guidGenerator,
AbpTimer timer)
AbpTimer timer,
IBackgroundJobExecuter jobExecuter)
: base(timer)
{
ServiceProvider = serviceProvider;
Clock = clock;
Serializer = serializer;
GuidGenerator = guidGenerator;
JobExecuter = jobExecuter;
Store = store;
Timer.Period = JobPollPeriod;
}
public async Task<Guid> EnqueueAsync<TJob, TArgs>(TArgs args, BackgroundJobPriority priority = BackgroundJobPriority.Normal, TimeSpan? delay = null)
where TJob : IBackgroundJob<TArgs>
public Task<Guid> EnqueueAsync<TArgs>(TArgs args, BackgroundJobPriority priority = BackgroundJobPriority.Normal, TimeSpan? delay = null)
{
var jobName = BackgroundJobNameAttribute.GetNameOrNull<TArgs>();
return EnqueueAsync(jobName, args, priority, delay);
}
public async Task<Guid> EnqueueAsync(string jobName, object args, BackgroundJobPriority priority = BackgroundJobPriority.Normal, TimeSpan? delay = null)
{
var jobInfo = new BackgroundJobInfo
{
Id = GuidGenerator.Create(),
JobType = typeof(TJob).AssemblyQualifiedName,
JobName = jobName,
JobArgs = Serializer.Serialize(args),
Priority = priority,
CreationTime = Clock.Now,
@ -77,104 +81,13 @@ namespace Volo.Abp.BackgroundJobs
return jobInfo.Id;
}
public async Task<bool> DeleteAsync(Guid jobId)
{
var jobInfo = await Store.FindAsync(jobId);
if (jobInfo == null)
{
return false;
}
await Store.DeleteAsync(jobInfo);
return true;
}
protected override void DoWork()
{
var waitingJobs = AsyncHelper.RunSync(() => Store.GetWaitingJobsAsync(1000));
foreach (var job in waitingJobs)
{
TryProcessJob(job);
}
}
private void TryProcessJob(BackgroundJobInfo jobInfo)
{
try
{
jobInfo.TryCount++;
jobInfo.LastTryTime = Clock.Now;
var jobType = Type.GetType(jobInfo.JobType);
using (var scope = ServiceProvider.CreateScope())
{
var job = scope.ServiceProvider.GetService(jobType);
if (job == null)
{
throw new AbpException("JobType is not registered: " + jobType);
}
//TODO: Type check for the job object
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);
try
{
jobExecuteMethod.Invoke(job, new[] { argsObj });
AsyncHelper.RunSync(() => Store.DeleteAsync(jobInfo));
}
catch (Exception ex)
{
Logger.LogException(ex);
var nextTryTime = jobInfo.CalculateNextTryTime(Clock);
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)!
}
}
}
catch (Exception ex)
{
Logger.LogException(ex);
jobInfo.IsAbandoned = true;
TryUpdate(jobInfo);
}
}
private void TryUpdate(BackgroundJobInfo jobInfo)
{
try
{
Store.UpdateAsync(jobInfo);
}
catch (Exception updateEx)
{
Logger.LogException(updateEx);
JobExecuter.Execute(job);
}
}
}

26
framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobManagerExtensions.cs

@ -1,4 +1,5 @@
using System;
using System.Threading.Tasks;
using Volo.Abp.Threading;
namespace Volo.Abp.BackgroundJobs
@ -11,16 +12,33 @@ namespace Volo.Abp.BackgroundJobs
/// <summary>
/// Enqueues a job to be executed.
/// </summary>
/// <typeparam name="TJob">Type of the job.</typeparam>
/// <typeparam name="TArgs">Type of the arguments of job.</typeparam>
/// <param name="backgroundJobManager">Background job manager reference</param>
/// <param name="args">Job arguments.</param>
/// <param name="priority">Job priority.</param>
/// <param name="delay">Job delay (wait duration before first try).</param>
public static void Enqueue<TJob, TArgs>(this IBackgroundJobManager backgroundJobManager, TArgs args, BackgroundJobPriority priority = BackgroundJobPriority.Normal, TimeSpan? delay = null)
where TJob : IBackgroundJob<TArgs>
public static void Enqueue<TArgs>(this IBackgroundJobManager backgroundJobManager, TArgs args, BackgroundJobPriority priority = BackgroundJobPriority.Normal, TimeSpan? delay = null)
{
AsyncHelper.RunSync(() => backgroundJobManager.EnqueueAsync<TJob, TArgs>(args, priority, delay));
AsyncHelper.RunSync(() => backgroundJobManager.EnqueueAsync<TArgs>(args, priority, delay));
}
/// <summary>
/// Enqueues a job to be executed.
/// </summary>
/// <param name="backgroundJobManager">Background job manager reference</param>
/// <param name="jobName">Job name.</param>
/// <param name="args">Job arguments.</param>
/// <param name="priority">Job priority.</param>
/// <param name="delay">Job delay (wait duration before first try).</param>
/// <returns>Unique identifier of a background job.</returns>
public static Guid EnqueueAsync(
this IBackgroundJobManager backgroundJobManager,
string jobName,
object args,
BackgroundJobPriority priority = BackgroundJobPriority.Normal,
TimeSpan? delay = null)
{
return AsyncHelper.RunSync(() => backgroundJobManager.EnqueueAsync(jobName, args, priority, delay));
}
}
}

30
framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobNameAttribute.cs

@ -0,0 +1,30 @@
using System;
using System.Linq;
using JetBrains.Annotations;
namespace Volo.Abp.BackgroundJobs
{
public class BackgroundJobNameAttribute : Attribute, IBackgroundJobNameProvider
{
public string Name { get; }
public BackgroundJobNameAttribute([NotNull] string name)
{
Name = Check.NotNullOrWhiteSpace(name, nameof(name));
}
public static string GetNameOrNull<TJobArgs>()
{
return GetNameOrNull(typeof(TJobArgs));
}
public static string GetNameOrNull(Type jobArgsType)
{
return jobArgsType
.GetCustomAttributes(true)
.OfType<IBackgroundJobNameProvider>()
.FirstOrDefault()
?.Name;
}
}
}

19
framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/BackgroundJobOptions.cs

@ -1,12 +1,29 @@
namespace Volo.Abp.BackgroundJobs
using System;
using System.Collections.Generic;
namespace Volo.Abp.BackgroundJobs
{
public class BackgroundJobOptions
{
public Dictionary<string, Type> JobTypes { get; }
public bool IsJobExecutionEnabled { get; set; }
public BackgroundJobOptions()
{
IsJobExecutionEnabled = true;
JobTypes = new Dictionary<string, Type>();
}
public Type GetJobType(string jobName)
{
var jobType = JobTypes.GetOrDefault(jobName);
if (jobType == null)
{
throw new AbpException("Undefined background job type for the job name: " + jobName);
}
return jobType;
}
}
}

7
framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobExecuter.cs

@ -0,0 +1,7 @@
namespace Volo.Abp.BackgroundJobs
{
public interface IBackgroundJobExecuter
{
void Execute(BackgroundJobInfo jobInfo);
}
}

19
framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobManager.cs

@ -13,20 +13,25 @@ namespace Volo.Abp.BackgroundJobs
/// <summary>
/// Enqueues a job to be executed.
/// </summary>
/// <typeparam name="TJob">Type of the job.</typeparam>
/// <typeparam name="TArgs">Type of the arguments of job.</typeparam>
/// <param name="args">Job arguments.</param>
/// <param name="priority">Job priority.</param>
/// <param name="delay">Job delay (wait duration before first try).</param>
/// <returns>Unique identifier of a background job.</returns>
Task<Guid> EnqueueAsync<TJob, TArgs>(TArgs args, BackgroundJobPriority priority = BackgroundJobPriority.Normal, TimeSpan? delay = null)
where TJob : IBackgroundJob<TArgs>;
Task<Guid> EnqueueAsync<TArgs>(
TArgs args,
BackgroundJobPriority priority = BackgroundJobPriority.Normal,
TimeSpan? delay = null
);
/// <summary>
/// Deletes a job with the specified jobId.
/// Enqueues a job to be executed.
/// </summary>
/// <param name="jobId">The Job Unique Identifier.</param>
/// <returns><c>True</c> on a successfull state transition, <c>false</c> otherwise.</returns>
Task<bool> DeleteAsync(Guid jobId);
/// <param name="jobName">Job name.</param>
/// <param name="args">Job arguments.</param>
/// <param name="priority">Job priority.</param>
/// <param name="delay">Job delay (wait duration before first try).</param>
/// <returns>Unique identifier of a background job.</returns>
Task<Guid> EnqueueAsync(string jobName, object args, BackgroundJobPriority priority = BackgroundJobPriority.Normal, TimeSpan? delay = null);
}
}

7
framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/IBackgroundJobNameProvider.cs

@ -0,0 +1,7 @@
namespace Volo.Abp.BackgroundJobs
{
public interface IBackgroundJobNameProvider
{
string Name { get; }
}
}

3
framework/src/Volo.Abp.BackgroundJobs/Volo/Abp/BackgroundJobs/JsonBackgroundJobSerializer.cs

@ -1,9 +1,10 @@
using System;
using Volo.Abp.DependencyInjection;
using Volo.Abp.Json;
namespace Volo.Abp.BackgroundJobs
{
public class JsonBackgroundJobSerializer : IBackgroundJobSerializer
public class JsonBackgroundJobSerializer : IBackgroundJobSerializer, ITransientDependency
{
private readonly IJsonSerializer _jsonSerializer;

1
framework/test/Volo.Abp.BackgroundJobs.Tests/Volo.Abp.BackgroundJobs.Tests.csproj

@ -2,6 +2,7 @@
<PropertyGroup>
<TargetFramework>netcoreapp2.1</TargetFramework>
<LangVersion>latest</LangVersion>
<AssemblyName>Volo.Abp.BackgroundJobs.Tests</AssemblyName>
<PackageId>Volo.Abp.BackgroundJobs.Tests</PackageId>
<GenerateRuntimeConfigurationFiles>true</GenerateRuntimeConfigurationFiles>

2
framework/test/Volo.Abp.BackgroundJobs.Tests/Volo.Abp.BackgroundJobs.Tests.csproj.DotSettings

@ -0,0 +1,2 @@
<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
<s:String x:Key="/Default/CodeInspection/CSharpLanguageProject/LanguageLevel/@EntryValue">CSharp71</s:String></wpf:ResourceDictionary>

16
framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/AbpBackgroundJobsTestModule.cs

@ -0,0 +1,16 @@
using Microsoft.Extensions.DependencyInjection;
using Volo.Abp.Modularity;
namespace Volo.Abp.BackgroundJobs
{
[DependsOn(
typeof(AbpBackgroundJobsModule)
)]
public class AbpBackgroundJobsTestModule : AbpModule
{
public override void ConfigureServices(ServiceConfigurationContext context)
{
context.Services.AddAssemblyOf<AbpBackgroundJobsTestModule>();
}
}
}

42
framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobManager_Tests.cs

@ -0,0 +1,42 @@
using System.Threading.Tasks;
using Shouldly;
using Xunit;
namespace Volo.Abp.BackgroundJobs
{
public class BackgroundJobManager_Tests : BackgroundJobsTestBase
{
private readonly IBackgroundJobManager _backgroundJobManager;
private readonly IBackgroundJobStore _backgroundJobStore;
public BackgroundJobManager_Tests()
{
_backgroundJobManager = GetRequiredService<IBackgroundJobManager>();
_backgroundJobStore = GetRequiredService<IBackgroundJobStore>();
}
[Fact]
public async Task Should_Store_Jobs()
{
var jobId = await _backgroundJobManager.EnqueueAsync(new MyJobArgs("42"));
jobId.ShouldNotBe(default);
(await _backgroundJobStore.FindAsync(jobId)).ShouldNotBeNull();
}
[BackgroundJobName("TestJobs.MyJob")]
private class MyJobArgs
{
public string Value { get; set; }
public MyJobArgs()
{
}
public MyJobArgs(string value)
{
Value = value;
}
}
}
}

7
framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/BackgroundJobsTestBase.cs

@ -0,0 +1,7 @@
namespace Volo.Abp.BackgroundJobs
{
public abstract class BackgroundJobsTestBase : AbpIntegratedTest<AbpBackgroundJobsTestModule>
{
}
}

7
framework/test/Volo.Abp.BackgroundJobs.Tests/Volo/Abp/BackgroundJobs/Class1.cs

@ -1,7 +0,0 @@
namespace Volo.Abp.BackgroundJobs
{
public class Class1
{
}
}
Loading…
Cancel
Save