using Elsa.Common;
using Elsa.Mediator.Contracts;
using Elsa.Scheduling.Commands;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using Timer = System.Timers.Timer;
namespace Elsa.Scheduling.ScheduledTasks;
///
/// A task that is scheduled to execute at a specific instant.
///
public class ScheduledSpecificInstantTask : IScheduledTask, IDisposable
{
private readonly ITask _task;
private readonly ISystemClock _systemClock;
private readonly IServiceScopeFactory _scopeFactory;
private readonly ILogger _logger;
private readonly DateTimeOffset _startAt;
private readonly CancellationTokenSource _cancellationTokenSource;
private Timer? _timer;
///
/// Initializes a new instance of .
///
public ScheduledSpecificInstantTask(ITask task, DateTimeOffset startAt, ISystemClock systemClock, IServiceScopeFactory scopeFactory, ILogger logger)
{
_task = task;
_systemClock = systemClock;
_scopeFactory = scopeFactory;
_logger = logger;
_startAt = startAt;
_cancellationTokenSource = new CancellationTokenSource();
Schedule();
}
///
public void Cancel() => _timer?.Dispose();
private void Schedule()
{
var now = _systemClock.UtcNow;
var delay = _startAt - now;
if (delay.Milliseconds <= 0)
delay = TimeSpan.FromMilliseconds(1);
_timer = new Timer(delay.TotalMilliseconds) { Enabled = true };
_timer.Elapsed += async (_, _) =>
{
_timer?.Dispose();
_timer = null;
using var scope = _scopeFactory.CreateScope();
var commandSender = scope.ServiceProvider.GetRequiredService();
var cancellationToken = _cancellationTokenSource.Token;
if (!cancellationToken.IsCancellationRequested)
{
try
{
await commandSender.SendAsync(new RunScheduledTask(_task), cancellationToken);
}
catch (Exception e)
{
_logger.LogError(e, "Error scheduled task");
}
}
};
}
void IDisposable.Dispose()
{
_cancellationTokenSource.Dispose();
_timer?.Dispose();
}
}