using FluentFTP; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Configuration; using RMuseum.DbContext; using RMuseum.Models.ExternalFTPUpload; using RSecurityBackend.Models.Generic; using RSecurityBackend.Services; using RSecurityBackend.Services.Implementation; using System; using System.IO; using System.Linq; using System.Threading.Tasks; namespace RMuseum.Services.Implementation { /// /// Queued FTP Upload Service /// public class QueuedFTPUploadService { /// /// add upload (you should call ProcessQueue manually) /// /// /// /// /// public async Task> AddAsync(string localFilePath, string remoteFilePath, bool deleteFileAfterUpload) { try { var q = new QueuedFTPUpload() { LocalFilePath = localFilePath, RemoteFilePath = remoteFilePath, DeleteFileAfterUpload = deleteFileAfterUpload, QueueDate = DateTime.Now, Processing = false, }; _context.QueuedFTPUploads.Add(q); await _context.SaveChangesAsync(); return new RServiceResult(q); } catch (Exception exp) { return new RServiceResult(null, exp.ToString()); } } /// /// 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()); } } /// /// get queued ftp uploads /// /// /// public async Task> GetQueuedFTPUploadsAsync(PagingParameterModel paging) { try { var source = _context.QueuedFTPUploads.AsNoTracking() .OrderBy(t => t.QueueDate) .AsQueryable(); (PaginationMetadata PagingMeta, QueuedFTPUpload[] Items) paginatedResult = await QueryablePaginator.Paginate(source, paging); return new RServiceResult<(PaginationMetadata PagingMeta, QueuedFTPUpload[] Items)>(paginatedResult); } catch (Exception exp) { return new RServiceResult<(PaginationMetadata PagingMeta, QueuedFTPUpload[] Items)>((null, null), 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, IConfiguration configuration, IBackgroundTaskQueue backgroundTaskQueue) { _context = context; Configuration = configuration; _backgroundTaskQueue = backgroundTaskQueue; } } }