using Hangfire;
using QYZH.InteractiveMagazine.Infrastructure.OSS;
using QYZH.InteractiveMagazine.Models.Entity;
using QYZH.InteractiveMagazine.Models.Enum;
using QYZH.InteractiveMagazine.WorkService.Consumers;
using SqlSugar;
using System.Net.Http.Headers;
using System.Text;
using System.Text.Json;
using System.Text.Json.Serialization;
namespace QYZH.InteractiveMagazine.WorkService.Jobs;
///
/// 期刊任务AI批改定时任务。
///
public class JournalTaskAiScoreJob(
ILogger logger,
IConfiguration configuration,
IServiceScopeFactory scopeFactory,
IHttpClientFactory httpClientFactory,
OssService ossService)
{
private const int DefaultPendingAnswerMinutes = 30;
private const int DefaultPendingAnswerBatchSize = 100;
private const int DefaultAiScoreMaxRetryCount = 3;
private const int DefaultAiScoreRetryDelayMilliseconds = 1000;
private const int DefaultAiMaxConcurrency = 1;
private const long DefaultMaxImageBytes = 10 * 1024 * 1024;
private const string InvalidAnswerResult = "作答内容不符合要求";
private const string AiProcessingMessage = "AI批阅中,请稍后";
private static readonly object AiSemaphoreLock = new();
private static readonly object RunWindowLock = new();
private static SemaphoreSlim? aiSemaphore;
private static int aiSemaphoreLimit;
private static DateTime? lastRunAt;
///
/// 执行AI批改。
///
[DisableConcurrentExecution(1800)]
public void Execute()
{
ExecuteAsync().GetAwaiter().GetResult();
}
///
/// 异步执行AI批改。
///
public async Task ExecuteAsync()
{
using var scope = scopeFactory.CreateScope();
var client = scope.ServiceProvider.GetRequiredService();
var windowEnd = DateTime.Now;
var windowStart = GetWindowStart(windowEnd);
var batchSize = GetPendingAnswerBatchSize();
var pendingAnswers = await client.Queryable()
.Where(a => !a.IsDeleted
&& a.Status == (int)UserAnswerStatusEnum.Processing
&& ((a.UpdatedAt != null && a.UpdatedAt > windowStart && a.UpdatedAt <= windowEnd)
|| (a.UpdatedAt == null && a.CreatedAt > windowStart && a.CreatedAt <= windowEnd)))
.OrderBy(a => a.UpdatedAt, OrderByType.Asc)
.OrderBy(a => a.CreatedAt, OrderByType.Asc)
.Take(batchSize)
.ToListAsync();
if (pendingAnswers.Count == 0)
{
SetLastRunAt(windowEnd);
logger.LogInformation("期刊AI批改任务没有待处理答案,WindowStart: {WindowStart}, WindowEnd: {WindowEnd}", windowStart, windowEnd);
return;
}
logger.LogInformation("期刊AI批改任务开始,Count: {Count}, WindowStart: {WindowStart}, WindowEnd: {WindowEnd}", pendingAnswers.Count, windowStart, windowEnd);
foreach (var pendingAnswer in pendingAnswers)
{
try
{
await ProcessPendingAnswerAsync(client, pendingAnswer.Id, windowStart, windowEnd, CancellationToken.None);
}
catch (Exception ex)
{
logger.LogError(ex, "期刊AI批改答案失败,AnswerId: {AnswerId}", pendingAnswer.Id);
}
}
SetLastRunAt(windowEnd);
}
private async Task ProcessPendingAnswerAsync(
ISqlSugarClient client,
long answerId,
DateTime windowStart,
DateTime windowEnd,
CancellationToken cancellationToken)
{
var answer = await client.Queryable()
.Where(a => a.Id == answerId && !a.IsDeleted)
.FirstAsync(cancellationToken);
var answerTime = answer == null ? default : GetPendingTime(answer);
if (answer == null
|| answer.Status != (int)UserAnswerStatusEnum.Processing
|| answerTime <= windowStart
|| answerTime > windowEnd)
{
return;
}
var task = await client.Queryable()
.Where(t => t.Id == answer.JournalPageTaskId && !t.IsDeleted)
.FirstAsync(cancellationToken);
if (task == null)
{
logger.LogWarning("期刊AI批改未找到任务,AnswerId: {AnswerId}, TaskId: {TaskId}", answer.Id, answer.JournalPageTaskId);
return;
}
if (!task.NeedAiProcess)
{
logger.LogInformation("期刊任务配置为人工批改,跳过AI批改,AnswerId: {AnswerId}, TaskId: {TaskId}", answer.Id, task.Id);
return;
}
var originalUpdatedAt = answer.UpdatedAt;
var originalAnswerUrl = answer.AnswerUrl ?? string.Empty;
try
{
var scoreUnit = await BuildScoreUnitAsync(client, answer, task, cancellationToken);
var referenceAnswerMap = await BuildReferenceAnswerMapAsync(client, scoreUnit, cancellationToken);
var answerQuestion = BuildQuestionFromAnswer(answer);
var scoreResult = await ScoreQuestionAsync(scoreUnit, answerQuestion, referenceAnswerMap, cancellationToken);
var normalizedResult = NormalizeScoreResult(scoreResult, task);
await CompleteAnswerAsync(client, answer.Id, originalUpdatedAt, originalAnswerUrl, normalizedResult, cancellationToken);
}
catch (OperationCanceledException)
{
throw;
}
catch (Exception ex)
{
logger.LogWarning(ex, "期刊AI批改调用失败,保持处理中等待下次重试,AnswerId: {AnswerId}", answer.Id);
await client.Updateable()
.SetColumns(a => new JournalPageTaskUserAnswer
{
Result = AiProcessingMessage,
UpdatedBy = "AI",
UpdatedAt = DateTime.Now
})
.Where(a => a.Id == answer.Id && a.Status == (int)UserAnswerStatusEnum.Processing)
.ExecuteCommandAsync(cancellationToken);
}
}
private async Task CompleteAnswerAsync(
ISqlSugarClient client,
long answerId,
DateTime? originalUpdatedAt,
string originalAnswerUrl,
JournalAnswerScoreResult scoreResult,
CancellationToken cancellationToken)
{
client.Ado.BeginTran();
try
{
var latest = await client.Queryable()
.Where(a => a.Id == answerId && !a.IsDeleted)
.FirstAsync(cancellationToken);
if (latest == null
|| latest.Status != (int)UserAnswerStatusEnum.Processing
|| latest.UpdatedAt != originalUpdatedAt
|| !string.Equals(latest.AnswerUrl ?? string.Empty, originalAnswerUrl, StringComparison.Ordinal))
{
client.Ado.CommitTran();
logger.LogInformation("期刊AI批改结果已过期,跳过写入,AnswerId: {AnswerId}", answerId);
return;
}
await client.Insertable(BuildAnswerSnapshot(latest)).ExecuteCommandAsync(cancellationToken);
latest.Result = scoreResult.Result;
latest.Points = Math.Max(0, scoreResult.Points);
latest.GrowthPoints = Math.Max(0, scoreResult.GrowthPoint);
latest.Score = Math.Max(0, scoreResult.Score);
latest.Comprehension = Math.Max(0, scoreResult.Comprehension);
latest.Judgment = Math.Max(0, scoreResult.Judgment);
latest.Expression = Math.Max(0, scoreResult.Expression);
latest.Persuasiveness = Math.Max(0, scoreResult.Persuasiveness);
latest.Status = (int)UserAnswerStatusEnum.Complete;
latest.UpdatedBy = "AI";
latest.UpdatedAt = DateTime.Now;
await client.Updateable(latest)
.IgnoreColumns(a => new { a.CreatedBy, a.CreatedAt })
.Where(a => a.Id == latest.Id)
.ExecuteCommandAsync(cancellationToken);
client.Ado.CommitTran();
logger.LogInformation("期刊AI批改完成,AnswerId: {AnswerId}, Score: {Score}", latest.Id, latest.Score);
}
catch
{
client.Ado.RollbackTran();
throw;
}
}
private async Task BuildScoreUnitAsync(
ISqlSugarClient client,
JournalPageTaskUserAnswer answer,
JournalPageTask task,
CancellationToken cancellationToken)
{
var groupId = task.GroupId > 0 ? task.GroupId : task.Id;
var groupTasks = task.GroupId > 0
? await client.Queryable()
.Where(t => t.GroupId == task.GroupId && !t.IsDeleted)
.ToListAsync(cancellationToken)
: [task];
var pageIds = groupTasks.Select(t => t.JournalPageId).Distinct().ToList();
var pages = await client.Queryable()
.Where(p => pageIds.Contains(p.Id) && !p.IsDeleted)
.ToListAsync(cancellationToken);
var pageMap = pages.ToDictionary(p => p.Id);
var question = BuildQuestionFromAnswer(answer);
var contexts = groupTasks
.OrderBy(t => pageMap.TryGetValue(t.JournalPageId, out var page) ? page.PageNum : 0)
.ThenBy(t => ParseTaskNo(t.No))
.ThenBy(t => t.Id)
.Select(t =>
{
pageMap.TryGetValue(t.JournalPageId, out var page);
return new JournalQuestionContext(question, t, page);
})
.ToList();
return new JournalScoreUnit(groupId, contexts);
}
private static async Task>> BuildReferenceAnswerMapAsync(
ISqlSugarClient client,
JournalScoreUnit scoreUnit,
CancellationToken cancellationToken)
{
var taskIds = scoreUnit.Questions.Select(q => q.Task.Id).Distinct().ToList();
var referenceAnswers = await client.Queryable()
.Where(a => taskIds.Contains(a.JournalPageTaskId) && !a.IsDeleted)
.ToListAsync(cancellationToken);
return referenceAnswers
.GroupBy(a => a.JournalPageTaskId)
.ToDictionary(g => g.Key, g => g.ToList());
}
private async Task ScoreQuestionAsync(
JournalScoreUnit scoreUnit,
Question answerQuestion,
Dictionary> referenceAnswerMap,
CancellationToken cancellationToken)
{
var apiKey = configuration["AiChat:ApiKey"];
var baseUrl = configuration["AiChat:BaseUrl"];
var model = configuration["AiChat:Model"];
var timeoutSeconds = configuration.GetValue("AiChat:TimeoutSeconds");
var maxTokens = configuration.GetValue("AiChat:MaxTokens");
var temperature = configuration.GetValue("AiChat:Temperature");
if (string.IsNullOrWhiteSpace(apiKey) || string.IsNullOrWhiteSpace(baseUrl) || string.IsNullOrWhiteSpace(model))
{
throw new InvalidOperationException("AI聊天服务配置不完整,请检查 AiChat 配置节点");
}
var content = new List