Files
Print_server/Inspectron.Epson/Queue/PrinterQueue.cs
2026-02-12 15:45:49 +01:00

144 lines
5.4 KiB
C#

using System.Collections.Concurrent;
using Microsoft.Extensions.Logging;
namespace Inspectron.Epson.Queue;
public class PrinterQueue
{
public string PrinterIp { get; }
private readonly ConcurrentQueue<PrintJob> _queue;
private readonly ConcurrentStack<PrintJob> _priorityQueue; // For failed jobs
private readonly SemaphoreSlim _signal;
private readonly CancellationTokenSource _cancellationTokenSource;
private Task _processingTask;
private readonly IPrintService _printService;
private readonly ILogger _logger;
private readonly PrintServer _printServer;
private readonly IJobStatusReporter _statusReporter;
public bool IsProcessing { get; private set; }
public int QueueLength => _queue.Count + _priorityQueue.Count;
public PrinterQueue(string printerIp, IPrintService printService, ILogger logger, PrintServer printServer, IJobStatusReporter statusReporter)
{
PrinterIp = printerIp;
_queue = new ConcurrentQueue<PrintJob>();
_priorityQueue = new ConcurrentStack<PrintJob>();
_signal = new SemaphoreSlim(0);
_cancellationTokenSource = new CancellationTokenSource();
_printService = printService;
_logger = logger;
_printServer = printServer;
_statusReporter = statusReporter;
}
public void Enqueue(PrintJob job)
{
_queue.Enqueue(job);
_signal.Release(); // Signal that there's work to do
_statusReporter.ReportStatusAsync(job, PrintJobStatus.Received).GetAwaiter().GetResult();
_logger.LogInformation("Job queued for printer {PrinterId}. Queue length: {QueueLength}", PrinterIp, QueueLength);
}
public void Start()
{
if (_processingTask != null)
return;
IsProcessing = true;
_processingTask = Task.Run(() => ProcessQueueAsync(_cancellationTokenSource.Token));
_logger.LogInformation("Printer queue {PrinterId} started", PrinterIp);
}
private void EnqueuePriority(PrintJob job)
{
_priorityQueue.Push(job);
_signal.Release();
_logger.LogInformation("Job priority queued for printer {PrinterId}. Queue length: {QueueLength}", PrinterIp, QueueLength);
}
private async Task ProcessQueueAsync(CancellationToken cancellationToken)
{
while (!cancellationToken.IsCancellationRequested)
{
try
{
// Wait for signal that there's work or cancellation
await _signal.WaitAsync(cancellationToken);
PrintJob job = null;
if (!_priorityQueue.TryPop(out job))
{
_queue.TryDequeue(out job);
}
if (job!=null)
{
_logger.LogInformation("Processing job on printer {PrinterId}", PrinterIp);
_printServer.EnterPrintLock();
try
{
// Call the actual print function
var result = await _printService.PrintAsync(PrinterIp, job);
if (!result.Success)
{
if (result.ErrorType == PrintErrorType.ConversionError)
{
_logger.LogWarning("Job failed due to conversion error, not re-queuing");
await _statusReporter.ReportStatusAsync(job, PrintJobStatus.Failed);
}
else
{
// Re-queue with retry logic
job.RetryCount++;
if (job.RetryCount == 1)
{
await _statusReporter.ReportStatusAsync(job, PrintJobStatus.OutOfPaper);
}
_logger.LogWarning("Job failed, retrying ({RetryCount}/3)", job.RetryCount);
await Task.Delay(5000, cancellationToken); // Wait before retry
EnqueuePriority(job);
}
}
else
{
_logger.LogInformation("Job completed successfully on printer {PrinterId}", PrinterIp);
await _statusReporter.ReportStatusAsync(job, PrintJobStatus.Completed);
}
}
catch (Exception ex)
{
_logger.LogError(ex, "Error processing job");
await _statusReporter.ReportStatusAsync(job, PrintJobStatus.Failed);
}
finally
{
_printServer.ExitPrintLock();
}
}
}
catch (OperationCanceledException)
{
break;
}
}
IsProcessing = false;
_logger.LogInformation("Printer queue {PrinterId} stopped", PrinterIp);
}
public async Task StopAsync()
{
_cancellationTokenSource.Cancel();
_signal.Release(); // Release to unblock the wait
if (_processingTask != null)
{
await _processingTask;
}
}
public List<PrintJob> GetPendingJobs()
{
return new List<PrintJob>(_queue);
}
}