Files
QYZH.InteractiveMagazine/QYZH.InteractiveMagazine.WorkService/Consumers/JournalTaskReceiveConsumer.cs

1308 lines
48 KiB
C#
Raw Normal View History

using QYZH.InteractiveMagazine.Infrastructure.OSS;
using QYZH.InteractiveMagazine.Models.Entity;
using QYZH.InteractiveMagazine.Models.Enum;
using SqlSugar;
using System.Net.Http.Headers;
using System.Text;
using System.Text.Json;
using System.Text.Json.Serialization;
namespace QYZH.InteractiveMagazine.WorkService.Consumers;
/// <summary>
/// 期刊任务接收消费者
/// </summary>
public class JournalTaskReceiveConsumer(
ILogger<JournalTaskReceiveConsumer> logger,
IConfiguration configuration,
IServiceScopeFactory scopeFactory,
IHttpClientFactory httpClientFactory,
OssService ossService) : IQueueConsumer
{
private const int DefaultAiScoreMaxRetryCount = 3;
private const int DefaultAiScoreRetryDelayMilliseconds = 1000;
private const int DefaultAiMaxConcurrency = 1;
private const long DefaultMaxImageBytes = 10 * 1024 * 1024;
private const float DefaultCompletionThreshold = 80;
private const float DefaultCommunityScoreThreshold = 90;
private const string InvalidAnswerResult = "作答内容不符合要求";
private const string AiProcessingMessage = "AI批阅中请稍后";
private static readonly object AiSemaphoreLock = new();
private static SemaphoreSlim? aiSemaphore;
private static int aiSemaphoreLimit;
public string Exchange => "ex.journal";
public string QueueName => "mq.journal.task.receive";
public string RoutingKey => "rk.journal.task.receive";
public async Task HandleAsync(byte[] body, CancellationToken cancellationToken = default)
{
var message = Encoding.UTF8.GetString(body);
logger.LogInformation("收到期刊任务消息: {Message}", message);
using var scope = scopeFactory.CreateScope();
var client = scope.ServiceProvider.GetRequiredService<ISqlSugarClient>();
var data = JsonSerializer.Deserialize<QuestionData>(message, new JsonSerializerOptions
{
PropertyNameCaseInsensitive = true
}) ?? throw new InvalidOperationException("期刊任务消息内容为空");
data.Normalize();
if (data.Questions == null || data.Questions.Length == 0)
{
logger.LogWarning("期刊任务消息没有题目UserId: {UserId}, JournalId: {JournalId}, PageId: {PageId}", data.UserId, data.JournalId, data.PageId);
return;
}
var taskIds = data.Questions.Select(q => q.Id).Distinct().ToList();
var tasks = await client.Queryable<JournalPageTask>()
.Where(t => taskIds.Contains(t.Id) && !t.IsDeleted)
.ToListAsync(cancellationToken);
var taskMap = tasks.ToDictionary(t => t.Id);
var groupIds = tasks.Select(t => t.GroupId > 0 ? t.GroupId : t.Id).Distinct().ToList();
var groupTasks = await client.Queryable<JournalPageTask>()
.Where(t => groupIds.Contains(t.GroupId) && !t.IsDeleted)
.ToListAsync(cancellationToken);
var persistedTaskMap = groupTasks
.GroupBy(t => t.GroupId > 0 ? t.GroupId : t.Id)
.ToDictionary(
g => g.Key,
g => g.FirstOrDefault(t => t.Id == g.Key)
?? g.OrderBy(t => t.Id).First());
var groupTaskIds = groupTasks.Select(t => t.Id).Distinct().ToList();
var referenceAnswers = await client.Queryable<JournalPageTaskAnswer>()
.Where(a => groupTaskIds.Contains(a.JournalPageTaskId) && !a.IsDeleted)
.ToListAsync(cancellationToken);
var referenceAnswerMap = referenceAnswers
.GroupBy(a => a.JournalPageTaskId)
.ToDictionary(g => g.Key, g => g.ToList());
var pageIds = groupTasks.Select(t => t.JournalPageId).Append(data.PageId).Distinct().ToList();
var pages = await client.Queryable<JournalPage>()
.Where(p => pageIds.Contains(p.Id) && !p.IsDeleted)
.ToListAsync(cancellationToken);
var pageMap = pages.ToDictionary(p => p.Id);
var existingAnswers = await client.Queryable<JournalPageTaskUserAnswer>()
.Where(a => a.UserId == data.UserId && groupTaskIds.Contains(a.JournalPageTaskId) && !a.IsDeleted)
.ToListAsync(cancellationToken);
var existingAnswerMap = existingAnswers
.GroupBy(a => a.JournalPageTaskId)
.ToDictionary(g => g.Key, g => g.OrderByDescending(a => a.UpdatedAt ?? a.CreatedAt).First());
var questionContexts = new List<JournalQuestionContext>();
var answerContexts = new List<JournalAnswerContext>();
foreach (var question in data.Questions)
{
if (!taskMap.TryGetValue(question.Id, out var task))
{
logger.LogWarning("未找到期刊任务TaskId: {TaskId}, UserId: {UserId}", question.Id, data.UserId);
continue;
}
if (!task.NeedAiProcess)
{
logger.LogInformation("期刊任务配置为人工批改跳过AI评分TaskId: {TaskId}, UserId: {UserId}", task.Id, data.UserId);
continue;
}
pageMap.TryGetValue(task.JournalPageId, out var taskPage);
questionContexts.Add(new JournalQuestionContext(question, task, taskPage));
}
foreach (var scoreUnit in BuildScoreUnits(questionContexts))
{
var persistedContext = BuildPersistedContext(scoreUnit, persistedTaskMap, pageMap);
if (persistedContext == null)
{
logger.LogWarning("期刊任务评分单元缺少可入库任务UserId: {UserId}, GroupId: {GroupId}",
data.UserId, scoreUnit.GroupId);
continue;
}
if (IsCompletedSameAnswer(scoreUnit, persistedContext.Task, data, existingAnswerMap))
{
logger.LogInformation("期刊任务作答已完成且内容未变化跳过AI评分UserId: {UserId}, GroupId: {GroupId}, TaskIds: {TaskIds}",
data.UserId, scoreUnit.GroupId, string.Join(",", scoreUnit.Questions.Select(q => q.Task.Id)));
continue;
}
try
{
var scoreResult = await ScoreQuestionAsync(scoreUnit, referenceAnswerMap, cancellationToken);
var normalizedResult = NormalizeScoreResult(scoreResult, persistedContext.Task);
answerContexts.Add(new JournalAnswerContext(
BuildAnswerEntity(data, scoreUnit, persistedContext.Question, persistedContext.Task, persistedContext.Page, normalizedResult, GetCompletionThreshold()),
persistedContext.Question,
persistedContext.Task,
normalizedResult));
LogSkippedGroupTasks(scoreUnit, persistedContext.Task);
}
catch (OperationCanceledException)
{
throw;
}
catch (Exception ex)
{
logger.LogWarning(ex, "期刊任务AI评分失败已标记为处理中UserId: {UserId}, GroupId: {GroupId}, TaskIds: {TaskIds}",
data.UserId, scoreUnit.GroupId, string.Join(",", scoreUnit.Questions.Select(q => q.Task.Id)));
var failureResult = BuildFailureScoreResult();
answerContexts.Add(new JournalAnswerContext(
BuildAnswerEntity(data, scoreUnit, persistedContext.Question, persistedContext.Task, persistedContext.Page, failureResult, GetCompletionThreshold()),
persistedContext.Question,
persistedContext.Task,
failureResult));
LogSkippedGroupTasks(scoreUnit, persistedContext.Task);
}
}
if (answerContexts.Count == 0)
{
logger.LogWarning("期刊任务消息没有可入库的答题记录UserId: {UserId}, JournalId: {JournalId}, PageId: {PageId}", data.UserId, data.JournalId, data.PageId);
return;
}
client.Ado.BeginTran();
try
{
foreach (var context in answerContexts)
{
var answer = context.Answer;
var existing = await client.Queryable<JournalPageTaskUserAnswer>()
.Where(a => a.UserId == answer.UserId && a.JournalPageTaskId == answer.JournalPageTaskId && !a.IsDeleted)
.FirstAsync(cancellationToken);
if (existing == null)
{
await client.Insertable(answer).ExecuteCommandAsync(cancellationToken);
}
else
{
await client.Insertable(BuildAnswerSnapshot(existing)).ExecuteCommandAsync(cancellationToken);
answer.Id = existing.Id;
answer.CreatedBy = existing.CreatedBy;
answer.CreatedAt = existing.CreatedAt;
answer.UpdatedBy = answer.UserId.ToString();
answer.UpdatedAt = DateTime.Now;
await client.Updateable(answer)
.IgnoreColumns(a => new { a.CreatedBy, a.CreatedAt })
.Where(a => a.Id == existing.Id)
.ExecuteCommandAsync(cancellationToken);
}
await InsertCommunityMessageIfNeededAsync(client, context, cancellationToken);
}
client.Ado.CommitTran();
}
catch
{
client.Ado.RollbackTran();
throw;
}
logger.LogInformation("期刊任务答题记录保存完成UserId: {UserId}, JournalId: {JournalId}, PageId: {PageId}, Count: {Count}",
data.UserId, data.JournalId, data.PageId, answerContexts.Count);
}
public Task OnErrorAsync(byte[] body, Exception exception)
{
var message = Encoding.UTF8.GetString(body);
logger.LogError(exception, "处理期刊任务消息失败: {Message}", message);
return Task.CompletedTask;
}
private async Task<JournalAnswerScoreResult> ScoreQuestionAsync(
JournalScoreUnit scoreUnit,
Dictionary<long, List<JournalPageTaskAnswer>> referenceAnswerMap,
CancellationToken cancellationToken)
{
var apiKey = configuration["AiChat:ApiKey"];
var baseUrl = configuration["AiChat:BaseUrl"];
var model = configuration["AiChat:Model"];
var timeoutSeconds = configuration.GetValue<int>("AiChat:TimeoutSeconds");
var maxTokens = configuration.GetValue<int>("AiChat:MaxTokens");
var temperature = configuration.GetValue<double?>("AiChat:Temperature");
if (string.IsNullOrWhiteSpace(apiKey) || string.IsNullOrWhiteSpace(baseUrl) || string.IsNullOrWhiteSpace(model))
{
throw new InvalidOperationException("AI聊天服务配置不完整请检查 AiChat 配置节点");
}
var content = new List<object>
{
new
{
type = "text",
text = BuildScorePrompt(scoreUnit, referenceAnswerMap)
}
};
var totalAnswerImageCount = 0;
foreach (var context in scoreUnit.Questions)
{
var answerImages = await BuildAnswerImageContentsAsync(context.Question, cancellationToken);
totalAnswerImageCount += answerImages.Count;
content.Add(new
{
type = "text",
text = $"以下为任务 {context.Task.Id} 的学生作答图片,共 {answerImages.Count} 张。"
});
foreach (var answerImage in answerImages)
{
content.Add(new
{
type = "image_url",
image_url = new { url = answerImage.DataUrl }
});
}
}
if (totalAnswerImageCount == 0)
{
throw new InvalidOperationException($"评分单元 {scoreUnit.GroupId} 缺少答案图片");
}
foreach (var context in scoreUnit.Questions)
{
referenceAnswerMap.TryGetValue(context.Task.Id, out var referenceAnswers);
var referenceAnswerImages = await BuildReferenceAnswerImageContentsAsync(context.Task.Id, referenceAnswers ?? [], cancellationToken);
if (referenceAnswerImages.Count == 0)
{
continue;
}
content.Add(new
{
type = "text",
text = $"以下为任务 {context.Task.Id} 的参考答案图片,共 {referenceAnswerImages.Count} 张。参考答案不是必有评分时以题目Prompt和学生答案为主。"
});
foreach (var referenceAnswerImage in referenceAnswerImages)
{
content.Add(new
{
type = "image_url",
image_url = new { url = referenceAnswerImage.DataUrl }
});
}
}
var requestBody = new
{
model,
messages = new object[]
{
new
{
role = "system",
content = "你是专业的学生作答评分助手。必须只返回合法 JSON不要返回 Markdown、解释或代码块。"
},
new
{
role = "user",
content
}
},
max_tokens = maxTokens > 0 ? maxTokens : 2000,
temperature = temperature is >= 0 ? temperature.Value : 0.1,
stream = false,
response_format = new { type = "json_object" }
};
var requestJson = JsonSerializer.Serialize(requestBody, new JsonSerializerOptions
{
DefaultIgnoreCondition = JsonIgnoreCondition.WhenWritingNull
});
var semaphore = GetAiSemaphore();
await semaphore.WaitAsync(cancellationToken);
try
{
var scoreResult = await SendAiScoreRequestWithRetryAsync(
scoreUnit.GroupId,
$"{baseUrl.TrimEnd('/')}/chat/completions",
apiKey,
requestJson,
timeoutSeconds > 0 ? timeoutSeconds : 300,
cancellationToken);
scoreResult.Result = TrimResult(scoreResult.Result);
return scoreResult;
}
finally
{
semaphore.Release();
}
}
private static string BuildScorePrompt(
JournalScoreUnit scoreUnit,
Dictionary<long, List<JournalPageTaskAnswer>> referenceAnswerMap)
{
var referenceAnswerTexts = scoreUnit.Questions
.SelectMany(q =>
{
referenceAnswerMap.TryGetValue(q.Task.Id, out var answers);
return answers ?? [];
})
.Select(a => a.Answer?.Trim())
.Where(a => !string.IsNullOrWhiteSpace(a))
.Distinct()
.ToList();
var prompt = new StringBuilder();
prompt.AppendLine(scoreUnit.Questions.Count > 1
? "这是同一跨页题目的多页作答,请综合全部学生作答图片评分。"
: "这是单页题目作答,请根据学生作答图片评分。");
prompt.AppendLine("以题目 Prompt 和学生答案为准评分,若有参考答案请结合参考答案。");
if (referenceAnswerTexts.Count > 0)
{
prompt.AppendLine("参考答案文本:");
for (var i = 0; i < referenceAnswerTexts.Count; i++)
{
prompt.AppendLine($"{i + 1}. {referenceAnswerTexts[i]}");
}
}
prompt.AppendLine();
prompt.AppendLine("题目信息:");
foreach (var context in scoreUnit.Questions)
{
prompt.AppendLine($" 题目内容:{context.Task.Task}");
prompt.AppendLine($" 评分Prompt{context.Task.Prompt}");
prompt.AppendLine($" 成长值上限:{context.Task.GrowthPoint}");
prompt.AppendLine($" 积分上限:{context.Task.Points}");
prompt.AppendLine($" 理解力上限:{context.Task.Comprehension}");
prompt.AppendLine($" 判断力上限:{context.Task.Judgment}");
prompt.AppendLine($" 表达力上限:{context.Task.Expression}");
prompt.AppendLine($" 说服力上限:{context.Task.Persuasiveness}");
}
prompt.AppendLine();
prompt.AppendLine("评分要求:");
prompt.AppendLine("- 不得超过题目配置中的各项上限。");
prompt.AppendLine("- 看不清、缺页、无法识别或答案明显不完整时,降低 Completion不要猜测高分。");
prompt.AppendLine("- Completion 表示作答完整度,范围 0-100。");
prompt.AppendLine("- Result 返回 50 字内中文评语。");
prompt.AppendLine("若作答内容与题目要求不符、答非所问、空白、仅抄题或无法形成有效答案Score/GrowthPoint/Points/Comprehension/Judgment/Expression/Persuasiveness 均返回 0。");
prompt.AppendLine("若作答内容不符合要求Result 必须且只能返回:作答内容不符合要求。不要补充原因、建议或其他文字。");
prompt.AppendLine();
prompt.AppendLine("只返回如下 JSON 字段:");
prompt.AppendLine("{");
prompt.AppendLine(" \"Score\": 0,");
prompt.AppendLine(" \"GrowthPoint\": 0,");
prompt.AppendLine(" \"Points\": 0,");
prompt.AppendLine(" \"Comprehension\": 0,");
prompt.AppendLine(" \"Judgment\": 0,");
prompt.AppendLine(" \"Expression\": 0,");
prompt.AppendLine(" \"Persuasiveness\": 0,");
prompt.AppendLine(" \"Completion\": 100,");
prompt.AppendLine(" \"Result\": \"50字内的中文评语\"");
prompt.AppendLine("}");
return prompt.ToString();
}
private async Task<List<AnswerImageContent>> BuildAnswerImageContentsAsync(Question question, CancellationToken cancellationToken)
{
var imageUrls = question.AnswerUrl?.Where(url => !string.IsNullOrWhiteSpace(url)).Distinct().ToList() ?? [];
var maxImageBytes = configuration.GetValue<long>("AiChat:MaxImageBytes");
if (maxImageBytes <= 0)
{
maxImageBytes = DefaultMaxImageBytes;
}
var result = new List<AnswerImageContent>();
foreach (var imageUrl in imageUrls)
{
cancellationToken.ThrowIfCancellationRequested();
await using var imageStream = await GetImageStreamAsync(imageUrl, cancellationToken);
if (imageStream == null)
{
throw new InvalidOperationException($"答案图片读取失败TaskId: {question.Id}, Url: {imageUrl}");
}
using var memoryStream = new MemoryStream();
await imageStream.CopyToAsync(memoryStream, cancellationToken);
if (memoryStream.Length == 0)
{
throw new InvalidOperationException($"答案图片内容为空TaskId: {question.Id}, Url: {imageUrl}");
}
if (memoryStream.Length > maxImageBytes)
{
throw new InvalidOperationException($"答案图片超过大小限制TaskId: {question.Id}, Url: {imageUrl}, Size: {memoryStream.Length}");
}
var imageBytes = memoryStream.ToArray();
var mimeType = GetImageMimeType(imageUrl, imageBytes);
result.Add(new AnswerImageContent($"data:{mimeType};base64,{Convert.ToBase64String(imageBytes)}"));
}
return result;
}
private async Task<List<AnswerImageContent>> BuildReferenceAnswerImageContentsAsync(
long taskId,
List<JournalPageTaskAnswer> referenceAnswers,
CancellationToken cancellationToken)
{
var imageUrls = referenceAnswers
.Select(a => a.AnswerUrL)
.Where(url => !string.IsNullOrWhiteSpace(url))
.Distinct()
.Select(url => url!)
.ToList();
var maxImageBytes = configuration.GetValue<long>("AiChat:MaxImageBytes");
if (maxImageBytes <= 0)
{
maxImageBytes = DefaultMaxImageBytes;
}
var result = new List<AnswerImageContent>();
foreach (var imageUrl in imageUrls)
{
cancellationToken.ThrowIfCancellationRequested();
try
{
await using var imageStream = await GetImageStreamAsync(imageUrl, cancellationToken);
if (imageStream == null)
{
throw new InvalidOperationException($"参考答案图片读取失败TaskId: {taskId}, Url: {imageUrl}");
}
using var memoryStream = new MemoryStream();
await imageStream.CopyToAsync(memoryStream, cancellationToken);
if (memoryStream.Length == 0)
{
throw new InvalidOperationException($"参考答案图片内容为空TaskId: {taskId}, Url: {imageUrl}");
}
if (memoryStream.Length > maxImageBytes)
{
throw new InvalidOperationException($"参考答案图片超过大小限制TaskId: {taskId}, Url: {imageUrl}, Size: {memoryStream.Length}");
}
var imageBytes = memoryStream.ToArray();
var mimeType = GetImageMimeType(imageUrl, imageBytes);
result.Add(new AnswerImageContent($"data:{mimeType};base64,{Convert.ToBase64String(imageBytes)}"));
}
catch (OperationCanceledException)
{
throw;
}
catch (Exception ex)
{
logger.LogWarning(ex, "参考答案图片读取失败已跳过TaskId: {TaskId}, Url: {Url}", taskId, imageUrl);
}
}
return result;
}
private async Task<JournalAnswerScoreResult> SendAiScoreRequestWithRetryAsync(
long taskId,
string requestUrl,
string apiKey,
string requestJson,
int timeoutSeconds,
CancellationToken cancellationToken)
{
var maxRetryCount = configuration.GetValue<int>("AiChat:ScoreMaxRetryCount");
if (maxRetryCount <= 0)
{
maxRetryCount = DefaultAiScoreMaxRetryCount;
}
var retryDelayMilliseconds = configuration.GetValue<int>("AiChat:ScoreRetryDelayMilliseconds");
if (retryDelayMilliseconds <= 0)
{
retryDelayMilliseconds = DefaultAiScoreRetryDelayMilliseconds;
}
Exception? lastException = null;
for (var attempt = 1; attempt <= maxRetryCount; attempt++)
{
cancellationToken.ThrowIfCancellationRequested();
try
{
var httpClient = httpClientFactory.CreateClient();
httpClient.Timeout = TimeSpan.FromSeconds(timeoutSeconds);
using var request = new HttpRequestMessage(HttpMethod.Post, requestUrl);
request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", apiKey);
request.Content = new StringContent(requestJson, Encoding.UTF8, "application/json");
using var response = await httpClient.SendAsync(request, cancellationToken);
var responseContent = await response.Content.ReadAsStringAsync(cancellationToken);
if (response.IsSuccessStatusCode)
{
var resultJson = ExtractAssistantContent(responseContent);
return ParseScoreResult(resultJson);
}
logger.LogWarning("AI评分调用失败TaskId: {TaskId}, Attempt: {Attempt}/{MaxRetryCount}, StatusCode: {StatusCode}, Response: {Response}",
taskId, attempt, maxRetryCount, response.StatusCode, responseContent);
if (!ShouldRetry(response.StatusCode))
{
throw new NonRetryAiScoreException($"AI评分调用失败{response.StatusCode},响应:{responseContent}");
}
if (attempt == maxRetryCount)
{
throw new InvalidOperationException($"AI评分调用失败{response.StatusCode},响应:{responseContent}");
}
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
throw;
}
catch (OperationCanceledException ex) when (attempt < maxRetryCount)
{
lastException = ex;
logger.LogWarning(ex, "AI评分调用超时准备重试TaskId: {TaskId}, Attempt: {Attempt}/{MaxRetryCount}", taskId, attempt, maxRetryCount);
}
catch (OperationCanceledException ex)
{
throw new InvalidOperationException($"AI评分调用超时TaskId: {taskId}", ex);
}
catch (Exception ex) when (attempt < maxRetryCount && ex is not NonRetryAiScoreException)
{
lastException = ex;
logger.LogWarning(ex, "AI评分调用异常准备重试TaskId: {TaskId}, Attempt: {Attempt}/{MaxRetryCount}", taskId, attempt, maxRetryCount);
}
var delay = TimeSpan.FromMilliseconds(retryDelayMilliseconds * attempt);
await Task.Delay(delay, cancellationToken);
}
throw new InvalidOperationException($"AI评分调用失败TaskId: {taskId}", lastException);
}
private async Task<Stream?> GetImageStreamAsync(string imageUrl, CancellationToken cancellationToken)
{
if (Uri.TryCreate(imageUrl, UriKind.Absolute, out var uri)
&& (uri.Scheme == Uri.UriSchemeHttp || uri.Scheme == Uri.UriSchemeHttps))
{
var httpClient = httpClientFactory.CreateClient();
return await httpClient.GetStreamAsync(uri, cancellationToken);
}
return ossService.GetObjectStream(imageUrl);
}
private static bool ShouldRetry(System.Net.HttpStatusCode statusCode)
{
var status = (int)statusCode;
return status == 408 || status == 429 || status >= 500;
}
private static string GetImageMimeType(string imageUrl, byte[] imageBytes)
{
if (imageBytes.Length >= 4)
{
if (imageBytes[0] == 0x89 && imageBytes[1] == 0x50 && imageBytes[2] == 0x4E && imageBytes[3] == 0x47)
{
return "image/png";
}
if (imageBytes[0] == 0xFF && imageBytes[1] == 0xD8)
{
return "image/jpeg";
}
if (imageBytes[0] == 0x47 && imageBytes[1] == 0x49 && imageBytes[2] == 0x46)
{
return "image/gif";
}
if (imageBytes[0] == 0x52 && imageBytes[1] == 0x49 && imageBytes[2] == 0x46 && imageBytes[3] == 0x46)
{
return "image/webp";
}
}
var extension = Path.GetExtension(imageUrl).ToLowerInvariant();
return extension switch
{
".png" => "image/png",
".jpg" or ".jpeg" => "image/jpeg",
".gif" => "image/gif",
".webp" => "image/webp",
_ => "image/jpeg"
};
}
private static string ExtractAssistantContent(string responseContent)
{
using var document = JsonDocument.Parse(responseContent);
var message = document.RootElement.GetProperty("choices")[0].GetProperty("message");
if (!message.TryGetProperty("content", out var contentElement))
{
throw new InvalidOperationException("AI评分响应缺少 content");
}
if (contentElement.ValueKind == JsonValueKind.String)
{
return contentElement.GetString() ?? string.Empty;
}
return contentElement.GetRawText();
}
private static JournalAnswerScoreResult ParseScoreResult(string resultJson)
{
var cleanedJson = CleanJsonContent(resultJson);
EnsureRequiredScoreFields(cleanedJson);
var result = JsonSerializer.Deserialize<JournalAnswerScoreResult>(cleanedJson, new JsonSerializerOptions
{
PropertyNameCaseInsensitive = true
});
return result ?? throw new InvalidOperationException("AI评分结果解析失败");
}
private static string CleanJsonContent(string content)
{
var text = content.Trim();
if (text.StartsWith("```", StringComparison.Ordinal))
{
var firstLineEnd = text.IndexOf('\n');
if (firstLineEnd >= 0)
{
text = text[(firstLineEnd + 1)..];
}
var fenceIndex = text.LastIndexOf("```", StringComparison.Ordinal);
if (fenceIndex >= 0)
{
text = text[..fenceIndex];
}
}
var start = text.IndexOf('{');
var end = text.LastIndexOf('}');
if (start >= 0 && end > start)
{
text = text[start..(end + 1)];
}
return text.Trim();
}
private static void EnsureRequiredScoreFields(string json)
{
using var document = JsonDocument.Parse(json);
var requiredFields = new[]
{
nameof(JournalAnswerScoreResult.Score),
nameof(JournalAnswerScoreResult.GrowthPoint),
nameof(JournalAnswerScoreResult.Points),
nameof(JournalAnswerScoreResult.Comprehension),
nameof(JournalAnswerScoreResult.Judgment),
nameof(JournalAnswerScoreResult.Expression),
nameof(JournalAnswerScoreResult.Persuasiveness),
nameof(JournalAnswerScoreResult.Completion),
nameof(JournalAnswerScoreResult.Result)
};
foreach (var field in requiredFields)
{
if (!document.RootElement.EnumerateObject().Any(p => string.Equals(p.Name, field, StringComparison.OrdinalIgnoreCase)))
{
throw new InvalidOperationException($"AI评分结果缺少字段{field}");
}
}
}
private static JournalAnswerScoreResult NormalizeScoreResult(JournalAnswerScoreResult scoreResult, JournalPageTask task)
{
var scoreMax = task.Comprehension + task.Judgment + task.Expression + task.Persuasiveness;
if (scoreMax <= 0)
{
scoreMax = 100;
}
return new JournalAnswerScoreResult
{
Score = Clamp(scoreResult.Score, 0, scoreMax),
GrowthPoint = (int)Clamp(scoreResult.GrowthPoint, 0, task.GrowthPoint),
Points = (int)Clamp(scoreResult.Points, 0, task.Points),
Comprehension = Clamp(scoreResult.Comprehension, 0, task.Comprehension),
Judgment = Clamp(scoreResult.Judgment, 0, task.Judgment),
Expression = Clamp(scoreResult.Expression, 0, task.Expression),
Persuasiveness = Clamp(scoreResult.Persuasiveness, 0, task.Persuasiveness),
Completion = Clamp(scoreResult.Completion, 0, 100),
Result = NormalizeResult(scoreResult.Result)
};
}
private static string NormalizeResult(string? result)
{
var trimmedResult = TrimResult(result);
if (trimmedResult.Contains("不符合要求", StringComparison.Ordinal) ||
trimmedResult.Contains("不符", StringComparison.Ordinal))
{
return InvalidAnswerResult;
}
return trimmedResult;
}
private static float Clamp(float value, float min, float max)
{
if (float.IsNaN(value) || float.IsInfinity(value))
{
return min;
}
if (max < min)
{
max = min;
}
return Math.Min(Math.Max(value, min), max);
}
private static JournalAnswerScoreResult BuildFailureScoreResult()
{
return new JournalAnswerScoreResult
{
Completion = 0,
Result = AiProcessingMessage
};
}
private static List<JournalScoreUnit> BuildScoreUnits(List<JournalQuestionContext> contexts)
{
return contexts
.GroupBy(c => c.Task.GroupId > 0 ? c.Task.GroupId : c.Task.Id)
.Select(g => new JournalScoreUnit(
g.Key,
g.OrderBy(c => c.Page?.PageNum ?? 0)
.ThenBy(c => ParseTaskNo(c.Task.No))
.ThenBy(c => c.Task.Id)
.ToList()))
.ToList();
}
private static JournalQuestionContext? BuildPersistedContext(
JournalScoreUnit scoreUnit,
Dictionary<long, JournalPageTask> persistedTaskMap,
Dictionary<long, JournalPage> pageMap)
{
if (scoreUnit.Questions.Count == 0)
{
return null;
}
var sourceContext = scoreUnit.Questions.FirstOrDefault(q => q.Task.Id == scoreUnit.GroupId)
?? scoreUnit.Questions.First();
var persistedTask = persistedTaskMap.TryGetValue(scoreUnit.GroupId, out var task)
? task
: sourceContext.Task;
pageMap.TryGetValue(persistedTask.JournalPageId, out var persistedPage);
return new JournalQuestionContext(sourceContext.Question, persistedTask, persistedPage);
}
private static int ParseTaskNo(string? taskNo)
{
return int.TryParse(taskNo, out var no) ? no : int.MaxValue;
}
private static bool IsCompletedSameAnswer(
JournalScoreUnit scoreUnit,
JournalPageTask persistedTask,
QuestionData data,
Dictionary<long, JournalPageTaskUserAnswer> existingAnswerMap)
{
return scoreUnit.Questions.Count > 0
&& existingAnswerMap.TryGetValue(persistedTask.Id, out var existing)
&& existing.Status == (int)UserAnswerStatusEnum.Complete
&& string.Equals(existing.PageAnswerUrl ?? string.Empty, data.PageAnswerUrl ?? string.Empty, StringComparison.Ordinal)
&& string.Equals(existing.AnswerUrl ?? string.Empty, SerializeAnswerUrls(scoreUnit), StringComparison.Ordinal)
&& existing.AnswerStartTime == GetAnswerStartTime(scoreUnit)
&& existing.AnswerEndTime == GetAnswerEndTime(scoreUnit);
}
private SemaphoreSlim GetAiSemaphore()
{
var maxConcurrency = configuration.GetValue<int>("AiChat:MaxConcurrency");
if (maxConcurrency <= 0)
{
maxConcurrency = DefaultAiMaxConcurrency;
}
lock (AiSemaphoreLock)
{
if (aiSemaphore == null || aiSemaphoreLimit != maxConcurrency)
{
aiSemaphore = new SemaphoreSlim(maxConcurrency, maxConcurrency);
aiSemaphoreLimit = maxConcurrency;
}
return aiSemaphore;
}
}
private void LogSkippedGroupTasks(JournalScoreUnit scoreUnit, JournalPageTask persistedTask)
{
var skippedTaskIds = scoreUnit.Questions
.Select(q => q.Task.Id)
.Where(id => id != persistedTask.Id)
.Distinct()
.ToList();
if (skippedTaskIds.Count == 0)
{
return;
}
logger.LogInformation("期刊跨页题评分结果仅保存到主任务GroupId: {GroupId}, PersistedTaskId: {PersistedTaskId}, SkippedTaskIds: {SkippedTaskIds}",
scoreUnit.GroupId, persistedTask.Id, string.Join(",", skippedTaskIds));
}
private static string SerializeAnswerUrls(JournalScoreUnit scoreUnit)
{
var answerUrls = scoreUnit.Questions
.SelectMany(q => q.Question.AnswerUrl ?? [])
.Where(url => !string.IsNullOrWhiteSpace(url))
.Distinct()
.ToArray();
return JsonSerializer.Serialize(answerUrls);
}
private static DateTime GetAnswerStartTime(JournalScoreUnit scoreUnit)
{
var startTimes = scoreUnit.Questions
.Select(q => q.Question.AnswerStartTime)
.Where(t => t != default)
.ToList();
return startTimes.Count == 0 ? default : startTimes.Min();
}
private static DateTime GetAnswerEndTime(JournalScoreUnit scoreUnit)
{
var endTimes = scoreUnit.Questions
.Select(q => q.Question.AnswerEndTime)
.Where(t => t != default)
.ToList();
return endTimes.Count == 0 ? default : endTimes.Max();
}
private static int GetAnswerSeconds(JournalScoreUnit scoreUnit)
{
return scoreUnit.Questions.Sum(q => Math.Max(0, q.Question.AnswerTime));
}
private static int GetBreakCount(JournalScoreUnit scoreUnit)
{
return scoreUnit.Questions.Sum(q => Math.Max(0, q.Question.BreakCount));
}
private static string SerializeBreakTimes(JournalScoreUnit scoreUnit)
{
var breakTimes = scoreUnit.Questions
.SelectMany(q => q.Question.BreakTimes ?? [])
.ToList();
return JsonSerializer.Serialize(breakTimes);
}
private static JournalPageTaskUserAnswer BuildAnswerEntity(
QuestionData data,
JournalScoreUnit scoreUnit,
Question question,
JournalPageTask task,
JournalPage? page,
JournalAnswerScoreResult scoreResult,
float completionThreshold)
{
var now = DateTime.Now;
var growthPoint = Math.Max(0, scoreResult.GrowthPoint);
var points = Math.Max(0, scoreResult.Points);
var answerStatus = scoreResult.Completion >= completionThreshold
? UserAnswerStatusEnum.Complete
: UserAnswerStatusEnum.Processing;
return new JournalPageTaskUserAnswer
{
JournalId = data.JournalId,
JournalPageId = data.PageId,
JournalPageTaskId = task.Id,
JournalPageTaskGroupId = task.GroupId,
UserId = data.UserId,
Result = scoreResult.Result,
Points = points,
GrowthPoints = growthPoint,
Score = Math.Max(0, scoreResult.Score),
Comprehension = Math.Max(0, scoreResult.Comprehension),
Judgment = Math.Max(0, scoreResult.Judgment),
Expression = Math.Max(0, scoreResult.Expression),
Persuasiveness = Math.Max(0, scoreResult.Persuasiveness),
QuestionAnswerUrl = question.Url,
AnswerUrl = SerializeAnswerUrls(scoreUnit),
PageAnswerUrl = data.PageAnswerUrl,
2026-07-01 09:20:07 +08:00
Revision = 0,
AnswerStartTime = GetAnswerStartTime(scoreUnit),
AnswerEndTime = GetAnswerEndTime(scoreUnit),
AnswerSeconds = GetAnswerSeconds(scoreUnit),
ImageRecognition = 0,
JournalPageNum = page?.PageNum ?? 0,
Modify = 0,
LastTag = 0,
DotPageNum = page?.PageNum ?? 0,
PageResultUrl = string.Empty,
Type = ((int)task.Type).ToString(),
DotPageNo = page?.PageNo ?? string.Empty,
PageAnswerDotUrl = string.Empty,
BreakCount = GetBreakCount(scoreUnit),
BreakTimes = SerializeBreakTimes(scoreUnit),
AssignmentStatus = (int)AssignmentStatusEnum.UnAssignmented,
Status = (int)answerStatus,
CreatedBy = data.UserId.ToString(),
CreatedAt = data.CreatedTime == default ? now : data.CreatedTime,
UpdatedBy = data.UserId.ToString(),
UpdatedAt = now
};
}
private static JournalPageTaskUserAnswerSnapshot BuildAnswerSnapshot(JournalPageTaskUserAnswer answer)
{
var now = DateTime.Now;
return new JournalPageTaskUserAnswerSnapshot
{
JournalPageTaskUserAnswerId = answer.Id,
JournalId = answer.JournalId,
JournalPageId = answer.JournalPageId,
JournalPageTaskId = answer.JournalPageTaskId,
JournalPageTaskGroupId = answer.JournalPageTaskGroupId,
UserId = answer.UserId,
Result = answer.Result,
Points = answer.Points,
GrowthPoints = answer.GrowthPoints,
Score = answer.Score,
Comprehension = answer.Comprehension,
Judgment = answer.Judgment,
Expression = answer.Expression,
Persuasiveness = answer.Persuasiveness,
QuestionAnswerUrl = answer.QuestionAnswerUrl,
AnswerUrl = answer.AnswerUrl,
PageAnswerUrl = answer.PageAnswerUrl,
2026-07-01 09:20:07 +08:00
Revision = answer.Revision,
AnswerStartTime = answer.AnswerStartTime,
AnswerEndTime = answer.AnswerEndTime,
AnswerSeconds = answer.AnswerSeconds,
ImageRecognition = answer.ImageRecognition,
JournalPageNum = answer.JournalPageNum,
Modify = answer.Modify,
LastTag = answer.LastTag,
DotPageNum = answer.DotPageNum,
PageResultUrl = answer.PageResultUrl,
Type = answer.Type,
DotPageNo = answer.DotPageNo,
PageAnswerDotUrl = answer.PageAnswerDotUrl,
BreakCount = answer.BreakCount,
BreakTimes = answer.BreakTimes,
AssignmentStatus = answer.AssignmentStatus,
Status = answer.Status,
CreatedBy = answer.UpdatedBy ?? answer.CreatedBy ?? string.Empty,
CreatedAt = now,
UpdatedBy = answer.UpdatedBy ?? string.Empty,
UpdatedAt = now
};
}
private async Task InsertCommunityMessageIfNeededAsync(
ISqlSugarClient client,
JournalAnswerContext context,
CancellationToken cancellationToken)
{
var scoreThreshold = configuration.GetValue<float>("AiChat:CommunityScoreThreshold");
if (scoreThreshold <= 0)
{
scoreThreshold = DefaultCommunityScoreThreshold;
}
if (context.Answer.Status != (int)UserAnswerStatusEnum.Complete || context.ScoreResult.Score < scoreThreshold)
{
return;
}
var exists = await client.Queryable<CommunityMessage>()
.Where(m => m.JournalTaskAnswerId == context.Answer.Id && !m.IsDeleted)
.AnyAsync(cancellationToken);
if (exists)
{
return;
}
var userJournal = await client.Queryable<UserJournal>()
.Where(uj => uj.UserId == context.Answer.UserId && uj.JournalId == context.Answer.JournalId && !uj.IsDeleted)
.FirstAsync(cancellationToken);
if (userJournal == null)
{
logger.LogWarning("高分答案未找到用户期刊关系跳过社区消息写入UserId: {UserId}, JournalId: {JournalId}, AnswerId: {AnswerId}",
context.Answer.UserId, context.Answer.JournalId, context.Answer.Id);
return;
}
var now = DateTime.Now;
var communityMessage = new CommunityMessage
{
JournalId = context.Answer.JournalId,
UserId = context.Answer.UserId,
UserJournalId = userJournal.Id,
JournalTaskId = context.Answer.JournalPageTaskId,
JournalTaskAnswerId = context.Answer.Id,
Content = context.ScoreResult.Result,
ImageUrl = context.Question.AnswerUrl.FirstOrDefault() ?? string.Empty,
SortOrder = 0,
IsActive = true,
Type = MessageTypeEnum.Message,
LikeCount = 0,
IsFeatured = 0,
Status = 1,
CreatedBy = context.Answer.UserId.ToString(),
CreatedAt = now,
UpdatedBy = context.Answer.UserId.ToString(),
UpdatedAt = now
};
await client.Insertable(communityMessage).ExecuteCommandAsync(cancellationToken);
}
private float GetCompletionThreshold()
{
var threshold = configuration.GetValue<float>("AiChat:CompletionThreshold");
if (threshold <= 0)
{
threshold = DefaultCompletionThreshold;
}
return threshold;
}
private static string TrimResult(string? result)
{
if (string.IsNullOrWhiteSpace(result))
{
return string.Empty;
}
return result.Length <= 50 ? result : result[..50];
}
}
/// <summary>
/// 期刊题目作答消息
/// </summary>
public class QuestionData
{
/// <summary>
/// 用户ID
/// </summary>
public long UserId { get; set; }
/// <summary>
/// 学生ID
/// </summary>
public long StudentId { get; set; }
/// <summary>
/// 期刊ID
/// </summary>
public long JournalId { get; set; }
/// <summary>
/// 作业ID
/// </summary>
public long HomeworkId { get; set; }
/// <summary>
/// 书籍ID
/// </summary>
public long BookId { get; set; }
/// <summary>
/// 页ID
/// </summary>
public long PageId { get; set; }
/// <summary>
/// 页面作答图片地址
/// </summary>
public string PageAnswerUrl { get; set; } = string.Empty;
/// <summary>
/// 题目作答列表
/// </summary>
public Question[] Questions { get; set; } = [];
/// <summary>
/// 创建时间
/// </summary>
public DateTime CreatedTime { get; set; }
/// <summary>
/// 兼容新版作答消息字段
/// </summary>
public void Normalize()
{
if (UserId <= 0)
{
UserId = StudentId;
}
if (JournalId <= 0)
{
JournalId = BookId > 0 ? BookId : HomeworkId;
}
}
}
/// <summary>
/// 题目作答数据
/// </summary>
public class Question
{
/// <summary>
/// 题目ID
/// </summary>
public long Id { get; set; }
/// <summary>
/// 题目图片地址
/// </summary>
public string Url { get; set; } = string.Empty;
/// <summary>
/// 答案图片地址
/// </summary>
public string[] AnswerUrl { get; set; } = [];
/// <summary>
/// 作答开始时间
/// </summary>
public DateTime AnswerStartTime { get; set; }
/// <summary>
/// 作答结束时间
/// </summary>
public DateTime AnswerEndTime { get; set; }
/// <summary>
/// 作答耗时秒数
/// </summary>
public int AnswerTime { get; set; }
/// <summary>
/// 中断次数
/// </summary>
public int BreakCount { get; set; }
/// <summary>
/// 中断记录
/// </summary>
public List<BreakTime> BreakTimes { get; set; } = [];
}
/// <summary>
/// 中断时间记录
/// </summary>
public class BreakTime
{
/// <summary>
/// 中断时间
/// </summary>
public DateTime Time { get; set; }
/// <summary>
/// 等待时间
/// </summary>
public long WaitTime { get; set; }
}
/// <summary>
/// AI评分结果
/// </summary>
public class JournalAnswerScoreResult
{
/// <summary>
/// 题目得分
/// </summary>
public float Score { get; set; }
/// <summary>
/// 成长值
/// </summary>
public int GrowthPoint { get; set; }
/// <summary>
/// 积分
/// </summary>
public int Points { get; set; }
/// <summary>
/// 理解力评分
/// </summary>
public float Comprehension { get; set; }
/// <summary>
/// 判断力评分
/// </summary>
public float Judgment { get; set; }
/// <summary>
/// 表达力评分
/// </summary>
public float Expression { get; set; }
/// <summary>
/// 说服力评分
/// </summary>
public float Persuasiveness { get; set; }
/// <summary>
/// 瀹屾垚搴?
/// </summary>
public float Completion { get; set; }
/// <summary>
/// 50字内评语
/// </summary>
public string Result { get; set; } = string.Empty;
}
public record JournalAnswerContext(
JournalPageTaskUserAnswer Answer,
Question Question,
JournalPageTask Task,
JournalAnswerScoreResult ScoreResult);
public record JournalQuestionContext(
Question Question,
JournalPageTask Task,
JournalPage? Page);
public record JournalScoreUnit(
long GroupId,
List<JournalQuestionContext> Questions);
public record AnswerImageContent(string DataUrl);
/// <summary>
/// AI评分不可重试异常。
/// </summary>
public class NonRetryAiScoreException(string message) : Exception(message);