diff --git a/RMuseum/RMuseum.xml b/RMuseum/RMuseum.xml
index df5db4b6..a6cf7ffe 100644
--- a/RMuseum/RMuseum.xml
+++ b/RMuseum/RMuseum.xml
@@ -18522,23 +18522,47 @@
- add upload
+ add upload (you should call ProcessQueue manually)
+
+
+ process queue
+
+
+
+
+
+ reset queue
+
+
+
Database Context
-
+
+
+ Configuration
+
+
+
+
+ Background Task Queue Instance
+
+
+
constructor
+
+
diff --git a/RMuseum/Services/Implementation/QueuedFTPUploadService.cs b/RMuseum/Services/Implementation/QueuedFTPUploadService.cs
index cb098d65..9f8ad851 100644
--- a/RMuseum/Services/Implementation/QueuedFTPUploadService.cs
+++ b/RMuseum/Services/Implementation/QueuedFTPUploadService.cs
@@ -1,7 +1,13 @@
-using RMuseum.DbContext;
+using FluentFTP;
+using Microsoft.EntityFrameworkCore;
+using Microsoft.Extensions.Configuration;
+using RMuseum.DbContext;
using RMuseum.Models.ExternalFTPUpload;
using RSecurityBackend.Models.Generic;
+using RSecurityBackend.Services;
using System;
+using System.IO;
+using System.Linq;
using System.Threading.Tasks;
namespace RMuseum.Services.Implementation
@@ -13,7 +19,7 @@ namespace RMuseum.Services.Implementation
public class QueuedFTPUploadService
{
///
- /// add upload
+ /// add upload (you should call ProcessQueue manually)
///
///
///
@@ -41,19 +47,159 @@ namespace RMuseum.Services.Implementation
}
}
+ ///
+ /// process queue
+ ///
+ ///
+ public async Task> ProcessQueue()
+ {
+ try
+ {
+ var processing = await _context.QueuedFTPUploads.AsNoTracking().Where(q => q.Processing).FirstOrDefaultAsync();
+ if(processing != null)
+ {
+ return new RServiceResult(false, $"already processing {processing.Id}.");
+ }
+
+ if (false == bool.Parse(Configuration.GetSection("ExternalFTPServer")["UploadEnabled"]))
+ {
+ return new RServiceResult(false, "ExternalFTPServer.UploadEnabled is not set to True.");
+ }
+
+ _backgroundTaskQueue.QueueBackgroundWorkItem
+ (
+ async token =>
+ {
+ using (RMuseumDbContext context = new RMuseumDbContext(new DbContextOptions()))
+ {
+ var next = await context.QueuedFTPUploads.Where(q => q.Processing == false).FirstOrDefaultAsync();
+ if (next == null) return;
+ var ftpClient = new AsyncFtpClient
+ (
+ Configuration.GetSection("ExternalFTPServer")["Host"],
+ Configuration.GetSection("ExternalFTPServer")["Username"],
+ Configuration.GetSection("ExternalFTPServer")["Password"]
+ );
+ ftpClient.ValidateCertificate += FtpClient_ValidateCertificate;
+ await ftpClient.AutoConnect();
+ ftpClient.Config.RetryAttempts = 3;
+ while (next != null)
+ {
+ try
+ {
+ next.Processing = true;
+ next.ProcessDate = DateTime.Now;
+ context.Update(next);
+ await context.SaveChangesAsync();
+ var status = await ftpClient.UploadFile(next.LocalFilePath, next.RemoteFilePath, createRemoteDir: true);
+ if(status != FtpStatus.Failed)
+ {
+ if(next.DeleteFileAfterUpload)
+ {
+ try
+ {
+ var dir = Path.GetDirectoryName(next.LocalFilePath);
+ File.Delete(next.LocalFilePath);
+ if(Directory.GetFiles(dir).Length == 0)
+ {
+ Directory.Delete(dir);
+ }
+ }
+ catch
+ {
+ //do nothing! not very important
+ }
+ }
+ context.Remove(next);
+ await context.SaveChangesAsync();
+ }
+ else
+ {
+ next.Error = "ftp client status is FtpStatus.Failed";
+ context.Update(next);
+ await context.SaveChangesAsync();
+ await ftpClient.Disconnect();
+ return;
+ }
+ }
+ catch (Exception e)
+ {
+ next.Error = e.ToString();
+ context.Update(next);
+ await context.SaveChangesAsync();
+ await ftpClient.Disconnect();
+ return;
+ }
+
+ next = await context.QueuedFTPUploads.Where(q => q.Processing == false).FirstOrDefaultAsync();
+ }
+ await ftpClient.Disconnect();
+ }
+ }
+ );
+
+ return new RServiceResult(true);
+ }
+ catch (Exception exp)
+ {
+ return new RServiceResult(false, exp.ToString());
+ }
+ }
+
+ private void FtpClient_ValidateCertificate(FluentFTP.Client.BaseClient.BaseFtpClient control, FtpSslValidationEventArgs e)
+ {
+ e.Accept = true;
+ }
+
+ ///
+ /// reset queue
+ ///
+ ///
+ public async Task> ResetQueue()
+ {
+ try
+ {
+ var queued = await _context.QueuedFTPUploads.Where(q => q.Processing).ToListAsync();
+ foreach (var queuedFTPUpload in queued)
+ {
+ queuedFTPUpload.Processing = false;
+ }
+ _context.UpdateRange(queued);
+ await _context.SaveChangesAsync();
+ return new RServiceResult(true);
+ }
+ catch (Exception exp)
+ {
+ return new RServiceResult(false, exp.ToString());
+ }
+ }
+
///
/// Database Context
///
protected readonly RMuseumDbContext _context;
+ ///
+ /// Configuration
+ ///
+ protected IConfiguration Configuration { get; }
+
+ ///
+ /// Background Task Queue Instance
+ ///
+ protected readonly IBackgroundTaskQueue _backgroundTaskQueue;
///
/// constructor
///
///
- public QueuedFTPUploadService(RMuseumDbContext context)
+ ///
+ ///
+ public QueuedFTPUploadService(RMuseumDbContext context, IConfiguration configuration, IBackgroundTaskQueue backgroundTaskQueue)
{
_context = context;
+ Configuration = configuration;
+ _backgroundTaskQueue = backgroundTaskQueue;
}
}
}