226 lines
12 KiB
C#
226 lines
12 KiB
C#
using dy.net.model.dto;
|
|
using dy.net.model.entity;
|
|
using dy.net.storage;
|
|
using Microsoft.Extensions.DependencyInjection;
|
|
using SqlSugar;
|
|
using MediaStorageType = dy.net.model.dto.StorageType;
|
|
|
|
namespace dy.net.service
|
|
{
|
|
/// <summary>
|
|
/// Replaces only the main media after a complete new download. The old row and file remain valid on failure.
|
|
/// </summary>
|
|
public class VideoRedownloadWorker : BackgroundService
|
|
{
|
|
private readonly IServiceScopeFactory _scopeFactory;
|
|
|
|
public VideoRedownloadWorker(IServiceScopeFactory scopeFactory) => _scopeFactory = scopeFactory;
|
|
|
|
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
|
{
|
|
while (!stoppingToken.IsCancellationRequested)
|
|
{
|
|
try
|
|
{
|
|
using var scope = _scopeFactory.CreateScope();
|
|
var db = scope.ServiceProvider.GetRequiredService<ISqlSugarClient>();
|
|
var waitingTaskIds = await db.Queryable<VideoDownloadTask>()
|
|
.Where(x => x.Status == VideoTaskStatus.WaitingForStorage)
|
|
.Select(x => x.Id).ToListAsync();
|
|
var pendingJobs = await db.Queryable<DouyinReDownload>()
|
|
.Where(x => x.Status == 0).OrderBy(x => x.CreateTime).ToListAsync();
|
|
var job = pendingJobs.FirstOrDefault(x => string.IsNullOrWhiteSpace(x.TaskId)
|
|
|| !waitingTaskIds.Contains(x.TaskId));
|
|
if (job == null)
|
|
{
|
|
await Task.Delay(TimeSpan.FromSeconds(5), stoppingToken);
|
|
continue;
|
|
}
|
|
await ProcessAsync(scope.ServiceProvider, job, stoppingToken);
|
|
}
|
|
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { }
|
|
catch (Exception ex)
|
|
{
|
|
Serilog.Log.Error(ex, "安全重新下载后台任务异常");
|
|
await Task.Delay(TimeSpan.FromSeconds(3), stoppingToken);
|
|
}
|
|
}
|
|
}
|
|
|
|
private static async Task ProcessAsync(IServiceProvider services, DouyinReDownload job, CancellationToken cancellationToken)
|
|
{
|
|
var db = services.GetRequiredService<ISqlSugarClient>();
|
|
var tasks = services.GetRequiredService<VideoTaskService>();
|
|
if (!string.IsNullOrWhiteSpace(job.TaskId))
|
|
{
|
|
var parent = await db.Queryable<VideoDownloadTask>().InSingleAsync(job.TaskId);
|
|
if (parent?.Status is VideoTaskStatus.WaitingForStorage or VideoTaskStatus.WaitingForSource)
|
|
{
|
|
if (parent.Status == VideoTaskStatus.WaitingForStorage) return;
|
|
}
|
|
try
|
|
{
|
|
await tasks.EnsureTaskStorageAvailableAsync(job.TaskId);
|
|
}
|
|
catch (InvalidOperationException ex)
|
|
{
|
|
if (!string.IsNullOrWhiteSpace(job.TaskItemId))
|
|
{
|
|
var blockedItem = await db.Queryable<VideoDownloadTaskItem>().InSingleAsync(job.TaskItemId);
|
|
if (blockedItem != null)
|
|
{
|
|
blockedItem.Stage = VideoTaskItemStage.WaitingForStorage;
|
|
blockedItem.ErrorType = VideoTaskErrorType.StorageUnavailable;
|
|
blockedItem.ErrorMessage = ex.Message;
|
|
blockedItem.UpdatedAt = DateTime.Now;
|
|
await db.Updateable(blockedItem).ExecuteCommandAsync();
|
|
}
|
|
}
|
|
await tasks.BlockTaskForStorageAsync(job.TaskId, ex.Message);
|
|
await tasks.RefreshCountsAsync(job.TaskId);
|
|
return;
|
|
}
|
|
}
|
|
var http = services.GetRequiredService<DouyinHttpClientService>();
|
|
var resolver = services.GetRequiredService<DouyinMigrationSourceResolver>();
|
|
var video = !string.IsNullOrWhiteSpace(job.VideoRecordId)
|
|
? await db.Queryable<DouyinVideo>().InSingleAsync(job.VideoRecordId)
|
|
: await db.Queryable<DouyinVideo>().Where(x => x.AwemeId == job.ViedoId).FirstAsync();
|
|
if (video == null)
|
|
{
|
|
await FinishAsync(db, job, false, "原视频记录不存在");
|
|
if (!string.IsNullOrWhiteSpace(job.TaskItemId))
|
|
await tasks.MarkFailedAsync(job.TaskItemId, new InvalidOperationException("原视频记录不存在"), VideoTaskErrorType.SourceUnavailable);
|
|
await CompleteTaskWhenIdleAsync(db, tasks, job.TaskId);
|
|
return;
|
|
}
|
|
var cookie = await db.Queryable<DouyinCookie>().InSingleAsync(video.CookieId);
|
|
if (!string.IsNullOrWhiteSpace(job.TaskItemId))
|
|
{
|
|
var tracked = await db.Queryable<VideoDownloadTaskItem>().InSingleAsync(job.TaskItemId);
|
|
var access = tracked?.Stage == VideoTaskItemStage.WaitingForSource
|
|
? await tasks.TryAcquireSourceAsync(video.CookieId, tracked.Id)
|
|
: await tasks.GetSourceAccessAsync(video.CookieId);
|
|
if (!access.Allowed)
|
|
{
|
|
await tasks.WaitItemForSourceAsync(job.TaskItemId, video.CookieId, access.Message);
|
|
return;
|
|
}
|
|
if (tracked?.Stage == VideoTaskItemStage.WaitingForSource)
|
|
{
|
|
tracked.Stage = VideoTaskItemStage.Pending;
|
|
tracked.ErrorMessage = null;
|
|
tracked.ErrorType = VideoTaskErrorType.None;
|
|
tracked.UpdatedAt = DateTime.Now;
|
|
await db.Updateable(tracked).ExecuteCommandAsync();
|
|
}
|
|
}
|
|
if (!string.IsNullOrWhiteSpace(job.TaskId)) await tasks.StartTaskAsync(job.TaskId);
|
|
if (!string.IsNullOrWhiteSpace(job.TaskItemId))
|
|
await tasks.SetItemStageAsync(job.TaskItemId, VideoTaskItemStage.Downloading, job.SavePath);
|
|
var urls = new List<string>();
|
|
if (!string.IsNullOrWhiteSpace(video.VideoUrl)) urls.Add(video.VideoUrl);
|
|
foreach (var url in await resolver.ResolveAsync(video, cookie, cancellationToken))
|
|
if (!urls.Contains(url)) urls.Add(url);
|
|
if (urls.Count == 0)
|
|
{
|
|
await FinishAsync(db, job, false, "记录 URL 与抖音全量列表均无可用地址");
|
|
if (!string.IsNullOrWhiteSpace(job.TaskItemId))
|
|
await tasks.MarkFailedAsync(job.TaskItemId, new InvalidOperationException("记录 URL 与抖音全量列表均无可用地址"), VideoTaskErrorType.SourceUnavailable);
|
|
await CompleteTaskWhenIdleAsync(db, tasks, job.TaskId);
|
|
return;
|
|
}
|
|
|
|
job.Attempts++;
|
|
job.UpdateTime = DateTime.UtcNow;
|
|
await db.Updateable(job).ExecuteCommandAsync();
|
|
try
|
|
{
|
|
long length;
|
|
if (video.StorageType == MediaStorageType.Local)
|
|
{
|
|
var target = Path.GetFullPath(video.VideoSavePath);
|
|
var directory = Path.GetDirectoryName(target) ?? throw new InvalidOperationException("原视频目录无效");
|
|
Directory.CreateDirectory(directory);
|
|
var temporary = Path.Combine(directory, "." + Path.GetFileName(target) + ".redownload-" + Guid.NewGuid().ToString("N"));
|
|
try
|
|
{
|
|
var downloaded = await http.DownloadAsync(urls[0], temporary, cookie?.Cookies, urls.Skip(1).ToList(), cancellationToken, maxRetryCount: Math.Min(12, urls.Count));
|
|
if (!downloaded.Success)
|
|
{
|
|
var disposition = SourceFailureDisposition.Continue;
|
|
if (!string.IsNullOrWhiteSpace(job.TaskItemId))
|
|
disposition = await tasks.MarkDownloadFailedAsync(job.TaskItemId, video.CookieId, downloaded);
|
|
if (disposition == SourceFailureDisposition.Continue)
|
|
await FinishAsync(db, job, false, downloaded.Message);
|
|
await CompleteTaskWhenIdleAsync(db, tasks, job.TaskId);
|
|
return;
|
|
}
|
|
await tasks.RecordSourceSuccessAsync(video.CookieId);
|
|
if (!File.Exists(downloaded.ActualSavePath)) throw new MediaStorageException("重新下载完成后临时文件不存在。");
|
|
length = new FileInfo(downloaded.ActualSavePath).Length;
|
|
if (length <= 0) throw new IOException("重新下载文件为空");
|
|
File.Move(downloaded.ActualSavePath, target, true);
|
|
}
|
|
finally
|
|
{
|
|
if (File.Exists(temporary)) File.Delete(temporary);
|
|
}
|
|
}
|
|
else
|
|
{
|
|
var downloaded = await http.DownloadToStorageAsync(video.StorageType, urls[0], video.VideoSavePath,
|
|
cookie?.Cookies, urls.Skip(1).ToList(), cancellationToken, Math.Min(12, urls.Count));
|
|
if (!downloaded.Success)
|
|
{
|
|
var disposition = SourceFailureDisposition.Continue;
|
|
if (!string.IsNullOrWhiteSpace(job.TaskItemId))
|
|
disposition = await tasks.MarkDownloadFailedAsync(job.TaskItemId, video.CookieId, downloaded);
|
|
if (disposition == SourceFailureDisposition.Continue)
|
|
await FinishAsync(db, job, false, downloaded.Message);
|
|
await CompleteTaskWhenIdleAsync(db, tasks, job.TaskId);
|
|
return;
|
|
}
|
|
await tasks.RecordSourceSuccessAsync(video.CookieId);
|
|
length = await services.GetRequiredService<MediaStorageRouter>()
|
|
.Resolve(video.StorageType).GetLengthAsync(video.VideoSavePath, cancellationToken) ?? 0;
|
|
if (length <= 0) throw new IOException($"{video.StorageType} 替换后的文件为空");
|
|
}
|
|
if (!string.IsNullOrWhiteSpace(job.TaskItemId))
|
|
await tasks.SetItemStageAsync(job.TaskItemId, VideoTaskItemStage.Verifying, video.VideoSavePath);
|
|
if (length <= 0) throw new IOException("重新下载后的文件为空");
|
|
if (!string.IsNullOrWhiteSpace(job.TaskItemId))
|
|
await tasks.SetItemStageAsync(job.TaskItemId, VideoTaskItemStage.Committing, video.VideoSavePath);
|
|
video.FileSize = length;
|
|
video.SyncTime = DateTime.Now;
|
|
await db.Updateable(video).ExecuteCommandAsync();
|
|
await FinishAsync(db, job, true, null);
|
|
if (!string.IsNullOrWhiteSpace(job.TaskItemId)) await tasks.MarkSucceededAsync(job.TaskItemId, length, video.Id);
|
|
}
|
|
catch (Exception ex)
|
|
{
|
|
Serilog.Log.Warning(ex, "安全重新下载失败,旧记录和旧文件保持不变:{VideoId}", video.Id);
|
|
await FinishAsync(db, job, false, ex.Message);
|
|
if (!string.IsNullOrWhiteSpace(job.TaskItemId)) await tasks.MarkFailedAsync(job.TaskItemId, ex);
|
|
}
|
|
await CompleteTaskWhenIdleAsync(db, tasks, job.TaskId);
|
|
}
|
|
|
|
private static async Task CompleteTaskWhenIdleAsync(ISqlSugarClient db, VideoTaskService tasks, string taskId)
|
|
{
|
|
if (string.IsNullOrWhiteSpace(taskId)) return;
|
|
var pending = await db.Queryable<DouyinReDownload>().Where(x => x.TaskId == taskId && x.Status == 0).CountAsync();
|
|
if (pending == 0) await tasks.CompleteTaskAsync(taskId);
|
|
}
|
|
|
|
private static async Task FinishAsync(ISqlSugarClient db, DouyinReDownload job, bool success, string error)
|
|
{
|
|
job.Status = success ? 1 : 2;
|
|
job.ErrorMessage = error;
|
|
job.UpdateTime = DateTime.UtcNow;
|
|
if (success) job.DownTime = DateTime.UtcNow;
|
|
await db.Updateable(job).ExecuteCommandAsync();
|
|
}
|
|
}
|
|
}
|