V 1.1.2 新增了生图包 以及各种转发和接口
This commit is contained in:
@@ -1,62 +1,100 @@
|
||||
using LMS.Tools.TaskScheduler;
|
||||
using LMS.Tools.MJPackage;
|
||||
using LMS.Tools.TaskScheduler;
|
||||
using Quartz;
|
||||
namespace LMS.service.Configuration
|
||||
|
||||
public static class QuartzTaskSchedulerConfig
|
||||
{
|
||||
public static class QuartzTaskSchedulerConfig
|
||||
public static void AddQuartzTaskSchedulerService(this IServiceCollection services)
|
||||
{
|
||||
public static void AddQuartzTaskSchedulerService(this IServiceCollection services)
|
||||
services.AddQuartz(q =>
|
||||
{
|
||||
// 注册 Quartz 服务
|
||||
services.AddQuartz(q =>
|
||||
// 时区配置
|
||||
var chinaTimeZone = GetChinaTimeZone();
|
||||
|
||||
// 每月任务配置
|
||||
ConfigureMonthlyTask(q, chinaTimeZone);
|
||||
|
||||
// 每日任务配置
|
||||
ConfigureDailyTask(q, chinaTimeZone);
|
||||
|
||||
// 每30秒任务配置
|
||||
ConfigureThirtySecondTask(q, chinaTimeZone);
|
||||
|
||||
ConfigureFiveMinuteTask(q, chinaTimeZone);
|
||||
});
|
||||
|
||||
services.AddQuartzHostedService(options =>
|
||||
{
|
||||
options.WaitForJobsToComplete = true;
|
||||
});
|
||||
|
||||
// 注册作业类
|
||||
services.AddTransient<ResetUserFreeCount>();
|
||||
services.AddTransient<TokenResetService>();
|
||||
services.AddTransient<TokenSyncService>();
|
||||
services.AddTransient<TaskStatusCheckService>();
|
||||
}
|
||||
|
||||
private static TimeZoneInfo GetChinaTimeZone()
|
||||
{
|
||||
try
|
||||
{
|
||||
return TimeZoneInfo.FindSystemTimeZoneById("China Standard Time");
|
||||
}
|
||||
catch
|
||||
{
|
||||
try
|
||||
{
|
||||
|
||||
// 配置作业
|
||||
var jobKey = new JobKey("MonthlyTask", "DefaultGroup");
|
||||
|
||||
// 方法1:通过配置属性设置时区
|
||||
// 获取中国时区
|
||||
TimeZoneInfo chinaTimeZone;
|
||||
try
|
||||
{
|
||||
// 尝试获取 Windows 时区名称
|
||||
chinaTimeZone = TimeZoneInfo.FindSystemTimeZoneById("China Standard Time");
|
||||
}
|
||||
catch
|
||||
{
|
||||
try
|
||||
{
|
||||
// 尝试获取 Linux 时区名称
|
||||
chinaTimeZone = TimeZoneInfo.FindSystemTimeZoneById("Asia/Shanghai");
|
||||
}
|
||||
catch
|
||||
{
|
||||
// 如果都不可用,使用 UTC+8
|
||||
chinaTimeZone = TimeZoneInfo.CreateCustomTimeZone(
|
||||
"China_Custom",
|
||||
new TimeSpan(8, 0, 0),
|
||||
"China Custom Time",
|
||||
"China Standard Time");
|
||||
}
|
||||
}
|
||||
|
||||
// 添加作业
|
||||
q.AddJob<ResetUserFreeCount>(opts => opts.WithIdentity(jobKey));
|
||||
|
||||
// 添加触发器 - 每月1号凌晨0点执行
|
||||
q.AddTrigger(opts => opts
|
||||
.ForJob(jobKey)
|
||||
.WithIdentity("MonthlyTaskTrigger", "DefaultGroup")
|
||||
.WithCronSchedule("0 0 0 1 * ?")); // 每月1号凌晨0点
|
||||
});
|
||||
|
||||
// 添加 Quartz 托管服务
|
||||
services.AddQuartzHostedService(options =>
|
||||
return TimeZoneInfo.FindSystemTimeZoneById("Asia/Shanghai");
|
||||
}
|
||||
catch
|
||||
{
|
||||
options.WaitForJobsToComplete = true;
|
||||
});
|
||||
|
||||
// 注册作业
|
||||
services.AddTransient<ResetUserFreeCount>();
|
||||
return TimeZoneInfo.CreateCustomTimeZone(
|
||||
"China_Custom",
|
||||
new TimeSpan(8, 0, 0),
|
||||
"China Custom Time",
|
||||
"China Standard Time");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private static void ConfigureMonthlyTask(IServiceCollectionQuartzConfigurator q, TimeZoneInfo timeZone)
|
||||
{
|
||||
var jobKey = new JobKey("MonthlyTask", "DefaultGroup");
|
||||
q.AddJob<ResetUserFreeCount>(opts => opts.WithIdentity(jobKey));
|
||||
q.AddTrigger(opts => opts
|
||||
.ForJob(jobKey)
|
||||
.WithIdentity("MonthlyTaskTrigger", "DefaultGroup")
|
||||
.WithCronSchedule("0 0 0 1 * ?", x => x.InTimeZone(timeZone)));
|
||||
}
|
||||
|
||||
private static void ConfigureDailyTask(IServiceCollectionQuartzConfigurator q, TimeZoneInfo timeZone)
|
||||
{
|
||||
var jobKey = new JobKey("DailyTask", "DefaultGroup");
|
||||
q.AddJob<TokenResetService>(opts => opts.WithIdentity(jobKey));
|
||||
q.AddTrigger(opts => opts
|
||||
.ForJob(jobKey)
|
||||
.WithIdentity("DailyTaskTrigger", "DefaultGroup")
|
||||
.WithCronSchedule("0 10 0 * * ?", x => x.InTimeZone(timeZone))); // 每天凌晨0点10分执行
|
||||
}
|
||||
|
||||
private static void ConfigureThirtySecondTask(IServiceCollectionQuartzConfigurator q, TimeZoneInfo timeZone)
|
||||
{
|
||||
var jobKey = new JobKey("ThirtySecondTask", "DefaultGroup");
|
||||
q.AddJob<TokenSyncService>(opts => opts.WithIdentity(jobKey));
|
||||
q.AddTrigger(opts => opts
|
||||
.ForJob(jobKey)
|
||||
.WithIdentity("ThirtySecondTaskTrigger", "DefaultGroup")
|
||||
.WithCronSchedule("*/30 * * * * ?", x => x.InTimeZone(timeZone))); // 每30秒执行一次
|
||||
}
|
||||
|
||||
private static void ConfigureFiveMinuteTask(IServiceCollectionQuartzConfigurator q, TimeZoneInfo timeZone)
|
||||
{
|
||||
var jobKey = new JobKey("FiveMinuteTask", "DefaultGroup");
|
||||
q.AddJob<TaskStatusCheckService>(opts => opts.WithIdentity(jobKey));
|
||||
q.AddTrigger(opts => opts
|
||||
.ForJob(jobKey)
|
||||
.WithIdentity("FiveMinuteTaskTrigger", "DefaultGroup")
|
||||
.WithCronSchedule("0 */5 * * * ?", x => x.InTimeZone(timeZone))); // 每5分钟执行一次
|
||||
}
|
||||
}
|
||||
@@ -5,12 +5,14 @@ using LMS.DAO.UserDAO;
|
||||
using LMS.service.Configuration.InitConfiguration;
|
||||
using LMS.service.Extensions.Mail;
|
||||
using LMS.service.Service;
|
||||
using LMS.service.Service.MJPackage;
|
||||
using LMS.service.Service.Other;
|
||||
using LMS.service.Service.PermissionService;
|
||||
using LMS.service.Service.PromptService;
|
||||
using LMS.service.Service.RoleService;
|
||||
using LMS.service.Service.SoftwareService;
|
||||
using LMS.service.Service.UserService;
|
||||
using LMS.Tools.MJPackage;
|
||||
|
||||
namespace Lai_server.Configuration
|
||||
{
|
||||
@@ -38,6 +40,8 @@ namespace Lai_server.Configuration
|
||||
services.AddScoped<SoftwareService>();
|
||||
services.AddScoped<MachineAuthorizationService>();
|
||||
services.AddScoped<DataInfoService>();
|
||||
services.AddScoped<ITokenManagementService, TokenManagementService>();
|
||||
services.AddScoped<IMJPackageService, MJPackageService>();
|
||||
|
||||
// 注入 DAO
|
||||
services.AddScoped<UserBasicDao>();
|
||||
@@ -53,6 +57,20 @@ namespace Lai_server.Configuration
|
||||
|
||||
// 添加分布式缓存(用于存储验证码)
|
||||
services.AddDistributedMemoryCache();
|
||||
|
||||
// 注册自定义服务
|
||||
services.AddSingleton<TokenUsageTracker>();
|
||||
services.AddScoped<ITokenService, TokenService>();
|
||||
services.AddScoped<ITaskConcurrencyManager, TaskConcurrencyManager>();
|
||||
services.AddScoped<ITaskService, TaskService>();
|
||||
|
||||
|
||||
// 注册后台服务
|
||||
services.AddHostedService<RsaConfigurattions>();
|
||||
services.AddHostedService<DatabaseConfiguration>();
|
||||
|
||||
//services.AddHostedService<TokenSyncService>();
|
||||
//services.AddHostedService<DailyResetService>();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,115 @@
|
||||
using Betalgo.Ranul.OpenAI.ObjectModels.RealtimeModels;
|
||||
using LMS.Repository.MJPackage;
|
||||
using LMS.service.Extensions.Attributes;
|
||||
using LMS.service.Service.MJPackage;
|
||||
using LMS.Tools.MJPackage;
|
||||
using Microsoft.AspNetCore.Mvc;
|
||||
using Newtonsoft.Json;
|
||||
using System.Text;
|
||||
using System.Text.Json;
|
||||
|
||||
namespace LMS.service.Controllers
|
||||
{
|
||||
[Route("api/[controller]")]
|
||||
[ApiController]
|
||||
public class MJPackageController(TokenUsageTracker usageTracker, ITaskConcurrencyManager taskConcurrencyManager, ILogger<MJPackageController> logger, IMJPackageService mJPackageService) : ControllerBase
|
||||
{
|
||||
private readonly TokenUsageTracker _usageTracker = usageTracker;
|
||||
private readonly ILogger<MJPackageController> _logger = logger;
|
||||
private readonly ITaskConcurrencyManager _taskConcurrencyManager = taskConcurrencyManager;
|
||||
private readonly IMJPackageService _mJPackageService = mJPackageService;
|
||||
|
||||
[HttpPost("mj/submit/imagine")]
|
||||
[RateLimit]
|
||||
public async Task<IActionResult> Imagine([FromBody] MJSubmitImageModel model)
|
||||
{
|
||||
|
||||
string token = (string)(HttpContext.Items["UseToken"] ?? string.Empty);
|
||||
string requestToken = (string)(HttpContext.Items["RequestToken"] ?? string.Empty);
|
||||
if (string.IsNullOrWhiteSpace(token))
|
||||
{
|
||||
return Unauthorized("API token is empty");
|
||||
}
|
||||
if (string.IsNullOrWhiteSpace(requestToken))
|
||||
{
|
||||
return Unauthorized("API token is empty");
|
||||
}
|
||||
|
||||
using HttpClient client = new HttpClient();
|
||||
client.DefaultRequestHeaders.Add("Authorization", "Bearer sk-" + token);
|
||||
|
||||
model.NotifyHook = "https://lms.laitool.cn/api/MJPackage/mj/mj-notify-hook";
|
||||
string body = JsonConvert.SerializeObject(model);
|
||||
|
||||
client.Timeout = Timeout.InfiniteTimeSpan;
|
||||
string mjUrl = "https://api.laitool.cc/mj/submit/imagine";
|
||||
var response = await client.PostAsync(mjUrl, new StringContent(body, Encoding.UTF8, "application/json"));
|
||||
// 读取响应内容
|
||||
string content = await response.Content.ReadAsStringAsync();
|
||||
|
||||
// 直接返回原始响应,包括状态码和内容
|
||||
var res = new ContentResult
|
||||
{
|
||||
Content = content,
|
||||
ContentType = response.Content.Headers.ContentType?.ToString() ?? "application/json",
|
||||
StatusCode = (int)response.StatusCode
|
||||
};
|
||||
// 判断请求任务的状态,判断是不是需要释放
|
||||
if (res.StatusCode != 200 && res.StatusCode != 201)
|
||||
{
|
||||
// 释放
|
||||
_usageTracker.ReleaseConcurrencyPermit(requestToken);
|
||||
_logger.LogInformation($"请求失败,Token并发许可已释放: {token}, 状态码: {res.StatusCode}");
|
||||
return res;
|
||||
}
|
||||
|
||||
// 判断是不是提交成功,并且记录请求返回的ID
|
||||
// 把 Content 转换为 匿名 对象
|
||||
var result = JsonConvert.DeserializeAnonymousType(content, new
|
||||
{
|
||||
code = 1,
|
||||
description = "提交成功",
|
||||
result = 1320098173412546
|
||||
});
|
||||
if (result == null)
|
||||
{
|
||||
// 失败
|
||||
_usageTracker.ReleaseConcurrencyPermit(requestToken);
|
||||
_logger.LogInformation($"请求失败,返回的请求体为空,Token并发许可已释放: {token}, 状态码: {res.StatusCode}");
|
||||
return res;
|
||||
}
|
||||
else
|
||||
{
|
||||
if (result.code != 1 && result.code != 22)
|
||||
{
|
||||
_usageTracker.ReleaseConcurrencyPermit(requestToken);
|
||||
_logger.LogInformation($"请求失败,code: {result.code},Token并发许可已释放: {token}, 状态码: {res.StatusCode}");
|
||||
return res;
|
||||
}
|
||||
}
|
||||
// 开始写入任务
|
||||
await _taskConcurrencyManager.CreateTaskAsync(requestToken, result.result.ToString());
|
||||
|
||||
return res;
|
||||
}
|
||||
|
||||
[HttpGet("mj/task/{id}/fetch")]
|
||||
public async Task<ActionResult<object>> Fetch(string id)
|
||||
{
|
||||
string token = HttpContext.Request.Headers.Authorization.ToString()
|
||||
.Replace("Bearer ", "", StringComparison.OrdinalIgnoreCase)
|
||||
.Trim();
|
||||
|
||||
var res = await _mJPackageService.FetchTaskAsync(id, token);
|
||||
return res;
|
||||
}
|
||||
|
||||
[HttpPost("mj/mj-notify-hook")]
|
||||
//[Route("mjPackage/mj/mj-notify-hook")]
|
||||
public async Task<IActionResult> MJNotifyHook([FromBody] JsonElement model)
|
||||
{
|
||||
return await _mJPackageService.MJNotifyHookAsync(model);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,212 @@
|
||||
using LMS.Common.Extensions;
|
||||
using LMS.Repository.DB;
|
||||
using LMS.Repository.DTO;
|
||||
using LMS.Repository.MJPackage;
|
||||
using LMS.service.Service.MJPackage;
|
||||
using Microsoft.AspNetCore.Authorization;
|
||||
using Microsoft.AspNetCore.Mvc;
|
||||
using System.ComponentModel.DataAnnotations;
|
||||
|
||||
namespace LMS.service.Controllers
|
||||
{
|
||||
[Route("api/[controller]/[action]")]
|
||||
[ApiController]
|
||||
public class TokenManagementController : ControllerBase
|
||||
{
|
||||
private readonly ITokenManagementService _tokenManagementService;
|
||||
private readonly ILogger<TokenManagementController> _logger;
|
||||
|
||||
public TokenManagementController(
|
||||
ITokenManagementService tokenManagementService,
|
||||
ILogger<TokenManagementController> logger)
|
||||
{
|
||||
_tokenManagementService = tokenManagementService;
|
||||
_logger = logger;
|
||||
}
|
||||
|
||||
#region 用户-使用Token查询对应的任务
|
||||
|
||||
[HttpGet("{token}")]
|
||||
public async Task<ActionResult<APIResponseModel<CollectionResponse<TokenAndTaskCollection>>>> QueryTokenTaskCollection(string token, [Required] int page, [Required] int pageSize, string? thirdPartyTaskId)
|
||||
{
|
||||
return await _tokenManagementService.QueryTokenTaskCollection(token, page, pageSize, thirdPartyTaskId);
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 用户-查询Token是不是存在,返回简单数据
|
||||
|
||||
[HttpGet("{token}")]
|
||||
public async Task<ActionResult<APIResponseModel<TokenCacheItem>>> GetTokenItem(string token)
|
||||
{
|
||||
return await _tokenManagementService.GetTokenItem(token);
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
|
||||
#region 管理员-新增Token
|
||||
|
||||
[HttpPost]
|
||||
[Authorize]
|
||||
public async Task<ActionResult<APIResponseModel<string>>> AddToken([FromBody] AddOrModifyTokenModel model)
|
||||
{
|
||||
long requestUserId = ConvertExtension.ObjectToLong(HttpContext.Items["UserId"] ?? 0);
|
||||
return await _tokenManagementService.AddToken(requestUserId, model);
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 管理员-修改token
|
||||
|
||||
[HttpPost("{tokenId}")]
|
||||
[Authorize]
|
||||
public async Task<ActionResult<APIResponseModel<string>>> ModifyToken(long tokenId, [FromBody] AddOrModifyTokenModel model)
|
||||
{
|
||||
long requestUserId = ConvertExtension.ObjectToLong(HttpContext.Items["UserId"] ?? 0);
|
||||
return await _tokenManagementService.ModifyToken(requestUserId, tokenId, model);
|
||||
}
|
||||
|
||||
|
||||
#endregion
|
||||
|
||||
#region 管理员-删除token
|
||||
|
||||
[HttpDelete("{tokenId}")]
|
||||
[Authorize]
|
||||
public async Task<ActionResult<APIResponseModel<string>>> DeleteToken(long tokenId)
|
||||
{
|
||||
long requestUserId = ConvertExtension.ObjectToLong(HttpContext.Items["UserId"] ?? 0);
|
||||
return await _tokenManagementService.DeleteToken(requestUserId, tokenId);
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 管理员-获取所有Token
|
||||
|
||||
[HttpGet]
|
||||
[Authorize]
|
||||
public async Task<ActionResult<APIResponseModel<CollectionResponse<TokenCacheItem>>>> QueryTokenCollection([Required] int page, [Required] int pageSize, string? token, long? tokenId, bool? efficient)
|
||||
{
|
||||
long requestUserId = ConvertExtension.ObjectToLong(HttpContext.Items["UserId"] ?? 0);
|
||||
return await _tokenManagementService.QueryTokenCollection(page, pageSize, token, tokenId, efficient, requestUserId);
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 管理员-获取指定ID得Token
|
||||
|
||||
[HttpGet("{tokenId}")]
|
||||
[Authorize]
|
||||
public async Task<ActionResult<APIResponseModel<MJApiTokens>>> QueryTokenById(long tokenId)
|
||||
{
|
||||
long requestUserId = ConvertExtension.ObjectToLong(HttpContext.Items["UserId"] ?? 0);
|
||||
return await _tokenManagementService.QueryTokenById(tokenId, requestUserId);
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 管理员-获取所有的任务
|
||||
|
||||
[HttpGet]
|
||||
[Authorize]
|
||||
public async Task<ActionResult<APIResponseModel<CollectionResponse<MJApiTaskCollection>>>> QueryTaskCollection([Required] int page, [Required] int pageSize, string? thirdPartyTaskId, string? token, long? tokenId)
|
||||
{
|
||||
long requestUserId = ConvertExtension.ObjectToLong(HttpContext.Items["UserId"] ?? 0);
|
||||
return await _tokenManagementService.QueryTaskCollection(requestUserId, page, pageSize, thirdPartyTaskId, token, tokenId);
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 管理员-删除小于指定时间的任务
|
||||
|
||||
[HttpDelete("{timestamp}")]
|
||||
[Authorize]
|
||||
public async Task<ActionResult<APIResponseModel<string>>> DeleteMJTaskEarlyTimestamp(long timestamp)
|
||||
{
|
||||
long requestUserId = ConvertExtension.ObjectToLong(HttpContext.Items["UserId"] ?? 0);
|
||||
return await _tokenManagementService.DeleteMJTaskEarlyTimestamp(requestUserId, timestamp);
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 管理员-删除指定ID的任务
|
||||
|
||||
[HttpDelete("{taskId}")]
|
||||
[Authorize]
|
||||
public async Task<ActionResult<APIResponseModel<string>>> DeleteMJTask(string taskId)
|
||||
{
|
||||
long requestUserId = ConvertExtension.ObjectToLong(HttpContext.Items["UserId"] ?? 0);
|
||||
return await _tokenManagementService.DeleteMJTask(requestUserId, taskId);
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 管理员-手动刷新token缓存
|
||||
|
||||
/// <summary>
|
||||
/// 手动刷新Token缓存
|
||||
/// </summary>
|
||||
/// <param name="token">Token字符串</param>
|
||||
/// <returns>刷新结果</returns>
|
||||
[HttpPost("{token}")]
|
||||
[Authorize]
|
||||
public async Task<ActionResult<APIResponseModel<string>>> RefreshToken(string token)
|
||||
{
|
||||
long requestUserId = ConvertExtension.ObjectToLong(HttpContext.Items["UserId"] ?? 0);
|
||||
return await _tokenManagementService.RefreshToken(requestUserId, token);
|
||||
}
|
||||
#endregion
|
||||
|
||||
#region 管理员-获取活跃的token
|
||||
|
||||
[HttpGet]
|
||||
[Authorize]
|
||||
public async Task<ActionResult<APIResponseModel<CollectionResponse<TokenCacheItem>>>> GetActiveTokens([FromQuery] int minutes = 5)
|
||||
{
|
||||
long requestUserId = ConvertExtension.ObjectToLong(HttpContext.Items["UserId"] ?? 0);
|
||||
return await _tokenManagementService.GetActiveTokens(requestUserId, minutes);
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 管理员-移除不活跃的token
|
||||
|
||||
[HttpGet]
|
||||
[Authorize]
|
||||
public async Task<ActionResult<APIResponseModel<string>>> RemoveNotActiveTokens([FromQuery] int minutes = 5)
|
||||
{
|
||||
long requestUserId = ConvertExtension.ObjectToLong(HttpContext.Items["UserId"] ?? 0);
|
||||
return await _tokenManagementService.RemoveNotActiveTokens(requestUserId, minutes);
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 内存统计接口
|
||||
|
||||
/// <summary>
|
||||
/// 健康检查端点
|
||||
/// </summary>
|
||||
/// <returns>系统健康状态</returns>
|
||||
[HttpGet]
|
||||
[Authorize]
|
||||
public async Task<ActionResult<APIResponseModel<object>>> GetHealth()
|
||||
{
|
||||
long requestUserId = ConvertExtension.ObjectToLong(HttpContext.Items["UserId"] ?? 0);
|
||||
return await _tokenManagementService.GetHealth(requestUserId);
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 任务统计
|
||||
[HttpGet]
|
||||
[Authorize]
|
||||
public async Task<ActionResult<APIResponseModel<TaskStatistics>>> GetDayTaskStatistics()
|
||||
{
|
||||
long requestUserId = ConvertExtension.ObjectToLong(HttpContext.Items["UserId"] ?? 0);
|
||||
return await _tokenManagementService.GetDayTaskStatistics(requestUserId);
|
||||
}
|
||||
|
||||
#endregion
|
||||
}
|
||||
}
|
||||
@@ -80,7 +80,7 @@ namespace LMS.service.Controllers
|
||||
HttpOnly = true,
|
||||
Secure = true, // 如果使用 HTTPS
|
||||
SameSite = SameSiteMode.None,
|
||||
Expires = DateTime.UtcNow.AddDays(7),
|
||||
Expires = BeijingTimeExtension.GetBeijingTime().AddDays(7),
|
||||
};
|
||||
Response.Cookies.Append("refreshToken", ((LoginResponse)apiResponse.Data).RefreshToken, cookieOptions);
|
||||
return apiResponse;
|
||||
|
||||
@@ -0,0 +1,176 @@
|
||||
using LMS.Common.Extensions;
|
||||
using LMS.Tools.MJPackage;
|
||||
using Microsoft.AspNetCore.Mvc;
|
||||
using Microsoft.AspNetCore.Mvc.Filters;
|
||||
|
||||
namespace LMS.service.Extensions.Attributes
|
||||
{
|
||||
[AttributeUsage(AttributeTargets.Method)]
|
||||
public class RateLimitAttribute : ActionFilterAttribute, IAsyncActionFilter
|
||||
{
|
||||
private string _token;
|
||||
private DateTime _startTime;
|
||||
private bool _concurrencyAcquired = false;
|
||||
|
||||
public RateLimitAttribute()
|
||||
{
|
||||
}
|
||||
|
||||
public override async Task OnActionExecutionAsync(ActionExecutingContext context, ActionExecutionDelegate next)
|
||||
{
|
||||
_startTime = BeijingTimeExtension.GetBeijingTime();
|
||||
|
||||
// 从服务容器获取需要的服务
|
||||
var serviceProvider = context.HttpContext.RequestServices;
|
||||
var logger = serviceProvider.GetRequiredService<ILogger<RateLimitAttribute>>();
|
||||
var tokenService = serviceProvider.GetRequiredService<ITokenService>();
|
||||
var usageTracker = serviceProvider.GetRequiredService<TokenUsageTracker>();
|
||||
|
||||
try
|
||||
{
|
||||
// 1. 获取Token
|
||||
_token = context.HttpContext.Request.Headers.Authorization.ToString()
|
||||
.Replace("Bearer ", "", StringComparison.OrdinalIgnoreCase)
|
||||
.Trim();
|
||||
|
||||
if (string.IsNullOrEmpty(_token))
|
||||
{
|
||||
logger.LogWarning($"请求缺少Token, IP: {context.HttpContext.Connection.RemoteIpAddress}");
|
||||
context.Result = new UnauthorizedObjectResult("Missing API token");
|
||||
return;
|
||||
}
|
||||
|
||||
// 2. 获取Token配置
|
||||
var tokenConfig = await tokenService.GetTokenAsync(_token);
|
||||
if (tokenConfig == null)
|
||||
{
|
||||
context.Result = new ObjectResult("Invalid token")
|
||||
{
|
||||
StatusCode = StatusCodes.Status403Forbidden
|
||||
};
|
||||
return;
|
||||
}
|
||||
|
||||
// 3. 检查Token是不是到期
|
||||
if (tokenConfig.ExpiresAt != null && tokenConfig.ExpiresAt < BeijingTimeExtension.GetBeijingTime())
|
||||
{
|
||||
context.Result = new ObjectResult("expired token")
|
||||
{
|
||||
StatusCode = StatusCodes.Status403Forbidden
|
||||
};
|
||||
return;
|
||||
}
|
||||
|
||||
// 4. 判断当前Token得上一次使用时间是否超过了10分钟,超过了重新从数据库获取
|
||||
if (tokenConfig.LastActivityTime < BeijingTimeExtension.GetBeijingTime().AddMinutes(-10))
|
||||
{
|
||||
logger.LogInformation($"Token {_token} 上次活动时间超过10分钟,重新从数据库获取配置");
|
||||
tokenConfig = await tokenService.GetDatabaseTokenAsync(_token);
|
||||
if (tokenConfig == null)
|
||||
{
|
||||
context.Result = new ObjectResult("Invalid or expired token")
|
||||
{
|
||||
StatusCode = StatusCodes.Status403Forbidden
|
||||
};
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
// 5. 检查日使用限制
|
||||
if (tokenConfig.DailyLimit > 0 && tokenConfig.DailyUsage >= tokenConfig.DailyLimit)
|
||||
{
|
||||
logger.LogWarning($"Token日限制已达上限: {_token}, 当前使用: {tokenConfig.DailyUsage}, 限制: {tokenConfig.DailyLimit}");
|
||||
context.Result = new ObjectResult("Daily limit exceeded")
|
||||
{
|
||||
StatusCode = StatusCodes.Status403Forbidden
|
||||
};
|
||||
return;
|
||||
}
|
||||
|
||||
// 6. 检查总使用限制
|
||||
if (tokenConfig.TotalLimit > 0 && tokenConfig.TotalUsage >= tokenConfig.TotalLimit)
|
||||
{
|
||||
logger.LogWarning($"Token总限制已达上限: {_token}, 当前使用: {tokenConfig.TotalUsage}, 限制: {tokenConfig.TotalLimit}");
|
||||
context.Result = new ObjectResult("Total limit exceeded")
|
||||
{
|
||||
StatusCode = StatusCodes.Status403Forbidden
|
||||
};
|
||||
return;
|
||||
}
|
||||
|
||||
// 7. 并发控制
|
||||
var (maxCount, currentlyExecuting, available) = usageTracker.GetConcurrencyStatus(_token);
|
||||
logger.LogInformation($"Token并发状态: {_token}, 最大: {maxCount}, 执行中: {currentlyExecuting}, 可用: {available}");
|
||||
|
||||
// 等待获取并发许可
|
||||
_concurrencyAcquired = await usageTracker.WaitForConcurrencyPermitAsync(_token);
|
||||
if (!_concurrencyAcquired)
|
||||
{
|
||||
logger.LogInformation($"Token并发限制超出: {_token}, 并发限制: {tokenConfig.ConcurrencyLimit}");
|
||||
context.Result = new ObjectResult($"Concurrency limit exceeded (max: {tokenConfig.ConcurrencyLimit})")
|
||||
{
|
||||
StatusCode = StatusCodes.Status429TooManyRequests
|
||||
};
|
||||
return;
|
||||
}
|
||||
|
||||
logger.LogInformation($"Token验证成功,开始处理请求: {_token}, 并发限制: {tokenConfig.ConcurrencyLimit}");
|
||||
|
||||
if (string.IsNullOrWhiteSpace(tokenConfig.UseToken))
|
||||
{
|
||||
context.Result = new ObjectResult($"Token Error")
|
||||
{
|
||||
StatusCode = StatusCodes.Status401Unauthorized
|
||||
};
|
||||
return;
|
||||
}
|
||||
|
||||
// 将新token存储在HttpContext.Items中
|
||||
context.HttpContext.Items["UseToken"] = tokenConfig.UseToken;
|
||||
context.HttpContext.Items["RequestToken"] = _token;
|
||||
|
||||
// 执行 Action
|
||||
var executedContext = await next();
|
||||
|
||||
// 6. 更新使用计数 (仅成功请求)
|
||||
if (executedContext.HttpContext.Response.StatusCode < 400)
|
||||
{
|
||||
tokenService.IncrementUsage(_token);
|
||||
|
||||
var duration = BeijingTimeExtension.GetBeijingTime() - _startTime;
|
||||
logger.LogInformation($"请求处理成功: Token={_token}, 状态码={executedContext.HttpContext.Response.StatusCode}, 耗时={duration.TotalMilliseconds}ms");
|
||||
}
|
||||
else
|
||||
{
|
||||
logger.LogWarning($"请求处理失败: Token={_token}, 状态码={executedContext.HttpContext.Response.StatusCode}");
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
// 在异常情况下也要释放并发许可
|
||||
if (_concurrencyAcquired)
|
||||
{
|
||||
|
||||
usageTracker.ReleaseConcurrencyPermit(_token);
|
||||
|
||||
}
|
||||
logger.LogError(ex, $"处理Token请求时发生错误: {_token},已释放Token许可!");
|
||||
context.Result = new ObjectResult("Internal server error")
|
||||
{
|
||||
StatusCode = StatusCodes.Status500InternalServerError
|
||||
};
|
||||
}
|
||||
finally
|
||||
{
|
||||
// 7. 释放并发许可
|
||||
//if (_concurrencyAcquired)
|
||||
//{
|
||||
// usageTracker.ReleaseConcurrencyPermit(_token);
|
||||
|
||||
// var newStatus = usageTracker.GetConcurrencyStatus(_token);
|
||||
// logger.LogInformation($"Token并发许可已释放: {_token}, 执行中: {newStatus.currentlyExecuting}, 可用: {newStatus.available}");
|
||||
//}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -5,6 +5,7 @@ using LMS.Repository.Models.DB;
|
||||
using LMS.service.Configuration;
|
||||
using LMS.service.Configuration.InitConfiguration;
|
||||
using LMS.service.Extensions.Middleware;
|
||||
using LMS.Tools.MJPackage;
|
||||
using Microsoft.AspNetCore.Identity;
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using Serilog;
|
||||
@@ -89,9 +90,6 @@ builder.Services.AddEndpointsApiExplorer();
|
||||
// 注入Swagger
|
||||
builder.Services.AddSwaggerService();
|
||||
|
||||
builder.Services.AddHostedService<RsaConfigurattions>();
|
||||
builder.Services.AddHostedService<DatabaseConfiguration>();
|
||||
|
||||
|
||||
var app = builder.Build();
|
||||
|
||||
@@ -120,6 +118,7 @@ app.MapControllers();
|
||||
// 在管道中使用IP速率限制中间件
|
||||
app.UseIpRateLimiting();
|
||||
|
||||
// 添加动态权限的中间件
|
||||
app.UseMiddleware<DynamicPermissionMiddleware>();
|
||||
app.UseEndpoints(endpoints =>
|
||||
{
|
||||
|
||||
@@ -115,6 +115,16 @@ public class ForwardWordService(ApplicationDbContext context)
|
||||
throw new Exception("参数错误");
|
||||
}
|
||||
|
||||
// 判断请求的url是不是满足条件
|
||||
if (string.IsNullOrEmpty(request.GptUrl))
|
||||
{
|
||||
throw new Exception("请求的url为空");
|
||||
}
|
||||
if (!request.GptUrl.StartsWith("https://ark.cn-beijing.volces.com") && !request.GptUrl.StartsWith("https://api.moonshot.cn") && !request.GptUrl.StartsWith("https://laitool.net") && !request.GptUrl.StartsWith("https://api.laitool.cc") && !request.GptUrl.StartsWith("https://laitool.cc"))
|
||||
{
|
||||
throw new Exception("请求的url不合法");
|
||||
}
|
||||
|
||||
// 获取提示词预设
|
||||
Prompt? prompt = await _context.Prompt.FirstOrDefaultAsync(x => x.PromptTypeId == request.PromptTypeId && x.Id == request.PromptId);
|
||||
if (prompt == null)
|
||||
|
||||
@@ -0,0 +1,18 @@
|
||||
using Microsoft.AspNetCore.Mvc;
|
||||
using System.Text.Json;
|
||||
|
||||
namespace LMS.service.Service.MJPackage
|
||||
{
|
||||
public interface IMJPackageService
|
||||
{
|
||||
/// <summary>
|
||||
/// mj 得 /mj/task/{id}/fetch 转发接口得实现
|
||||
/// 内含重试请求,全部重试失败 会尝试充数据库中读取消息
|
||||
/// </summary>
|
||||
/// <param name="id"></param>
|
||||
/// <param name="token"></param>
|
||||
/// <returns></returns>
|
||||
Task<ActionResult<object>> FetchTaskAsync(string id, string token);
|
||||
Task<ActionResult> MJNotifyHookAsync(JsonElement model);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
using LMS.Repository.DB;
|
||||
using LMS.Repository.DTO;
|
||||
using LMS.Repository.MJPackage;
|
||||
using Microsoft.AspNetCore.Mvc;
|
||||
|
||||
namespace LMS.service.Service.MJPackage
|
||||
{
|
||||
public interface ITokenManagementService
|
||||
{
|
||||
Task<ActionResult<APIResponseModel<CollectionResponse<TokenAndTaskCollection>>>> QueryTokenTaskCollection(string token, int page, int pageSize, string? thirdPartyTaskId);
|
||||
Task<ActionResult<APIResponseModel<TokenCacheItem>>> GetTokenItem(string token);
|
||||
Task<ActionResult<APIResponseModel<string>>> AddToken(long requestUserId, AddOrModifyTokenModel model);
|
||||
Task<ActionResult<APIResponseModel<string>>> ModifyToken(long requestUserId, long tokenId, AddOrModifyTokenModel model);
|
||||
Task<ActionResult<APIResponseModel<string>>> DeleteToken(long requestUserId, long tokenId);
|
||||
Task<ActionResult<APIResponseModel<CollectionResponse<MJApiTaskCollection>>>> QueryTaskCollection(long requestUserId, int page, int pageSize, string? thirdPartyTaskId, string? token, long? tokenId);
|
||||
Task<ActionResult<APIResponseModel<CollectionResponse<TokenCacheItem>>>> QueryTokenCollection(int page, int pageSize, string? token, long? tokenId, bool? efficient, long requestUserId);
|
||||
Task<ActionResult<APIResponseModel<string>>> DeleteMJTaskEarlyTimestamp(long requestUserId, long timestamp);
|
||||
Task<ActionResult<APIResponseModel<string>>> DeleteMJTask(long requestUserId, string taskId);
|
||||
Task<ActionResult<APIResponseModel<string>>> RefreshToken(long requestUserId, string token);
|
||||
Task<ActionResult<APIResponseModel<CollectionResponse<TokenCacheItem>>>> GetActiveTokens(long requestUserId, int minutes);
|
||||
Task<ActionResult<APIResponseModel<string>>> RemoveNotActiveTokens(long requestUserId, int minutes);
|
||||
Task<ActionResult<APIResponseModel<object>>> GetHealth(long requestUserId);
|
||||
Task<ActionResult<APIResponseModel<MJApiTokens>>> QueryTokenById(long tokenId, long requestUserId);
|
||||
Task<ActionResult<APIResponseModel<TaskStatistics>>> GetDayTaskStatistics(long requestUserId);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,295 @@
|
||||
using LMS.Common.Extensions;
|
||||
using LMS.DAO;
|
||||
using LMS.Repository.DB;
|
||||
using LMS.Tools.MJPackage;
|
||||
using Microsoft.AspNetCore.Mvc;
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using System.Net.Sockets;
|
||||
using System.Text.Json;
|
||||
using static Betalgo.Ranul.OpenAI.ObjectModels.StaticValues.AssistantsStatics.MessageStatics;
|
||||
|
||||
namespace LMS.service.Service.MJPackage
|
||||
{
|
||||
public class MJPackageService(ILogger<MJPackageService> logger, ApplicationDbContext dbContext, TokenUsageTracker usageTracker, ITaskConcurrencyManager taskConcurrencyManager, ITokenService tokenService) : IMJPackageService
|
||||
{
|
||||
private readonly ILogger<MJPackageService> _logger = logger;
|
||||
private readonly ApplicationDbContext _dbContext = dbContext;
|
||||
private readonly TokenUsageTracker _usageTracker = usageTracker;
|
||||
private readonly ITaskConcurrencyManager _taskConcurrencyManager = taskConcurrencyManager;
|
||||
private readonly ITokenService _tokenService = tokenService;
|
||||
|
||||
public async Task<ActionResult<object>> FetchTaskAsync(string id, string token)
|
||||
{
|
||||
|
||||
// 参数验证
|
||||
var validationResult = ValidateParameters(id, token);
|
||||
if (validationResult != null) return validationResult;
|
||||
|
||||
// 获取UseToken
|
||||
// 2. 获取Token配置
|
||||
var tokenConfig = await _tokenService.GetTokenAsync(token);
|
||||
if (tokenConfig == null || string.IsNullOrWhiteSpace(tokenConfig.UseToken))
|
||||
{
|
||||
return new UnauthorizedObjectResult(new { error = "无效或过期的Token" });
|
||||
}
|
||||
|
||||
// 尝试原始API
|
||||
var originResult = await TryOriginApiAsync(id);
|
||||
if (originResult != null) return originResult;
|
||||
|
||||
// 尝试备用API
|
||||
var backupResult = await TryBackupApiAsync(id, tokenConfig.UseToken);
|
||||
if (backupResult != null) return backupResult;
|
||||
|
||||
// 从数据库获取缓存数据
|
||||
return await GetTaskFromDatabaseAsync(id, token);
|
||||
}
|
||||
|
||||
|
||||
#region MJ Package 的回调处理
|
||||
|
||||
public async Task<ActionResult> MJNotifyHookAsync(JsonElement model)
|
||||
{
|
||||
try
|
||||
{
|
||||
string rawJson = model.GetRawText();
|
||||
|
||||
// 尝试获取ID字段
|
||||
string mjId = string.Empty;
|
||||
if (model.TryGetProperty("id", out var idElement))
|
||||
{
|
||||
mjId = idElement.ToString();
|
||||
}
|
||||
else if (model.TryGetProperty("Id", out var idElementCap))
|
||||
{
|
||||
mjId = idElementCap.ToString();
|
||||
}
|
||||
|
||||
if (string.IsNullOrWhiteSpace(mjId))
|
||||
{
|
||||
_logger.LogWarning("MJNotifyHook: 接收到的回调数据中缺少ID");
|
||||
return new BadRequestObjectResult("缺少ID");
|
||||
}
|
||||
|
||||
// 获取任务
|
||||
var mjTask = await _taskConcurrencyManager.GetTaskInfoByThirdPartyIdAsync(mjId);
|
||||
if (mjTask == null)
|
||||
{
|
||||
return new NotFoundObjectResult($"未找到ID为 {mjId} 的任务");
|
||||
}
|
||||
|
||||
// 尝试获取状态字段
|
||||
string status = MJTaskStatus.SUBMITTED;
|
||||
|
||||
if (model.TryGetProperty("status", out var statusElement))
|
||||
{
|
||||
status = statusElement.GetString() ?? MJTaskStatus.SUBMITTED;
|
||||
}
|
||||
else if (model.TryGetProperty("Status", out var statusElementCap))
|
||||
{
|
||||
status = statusElementCap.GetString() ?? MJTaskStatus.SUBMITTED;
|
||||
}
|
||||
|
||||
MJApiTasks mJApiTasks = new()
|
||||
{
|
||||
TaskId = mjTask.TaskId,
|
||||
Token = mjTask.Token,
|
||||
Status = status,
|
||||
StartTime = mjTask.StartTime,
|
||||
EndTime = null,
|
||||
ThirdPartyTaskId = mjId,
|
||||
Properties = rawJson // 或者直接存储 model
|
||||
};
|
||||
|
||||
if (mjTask.Status == MJTaskStatus.SUCCESS || mjTask.Status == MJTaskStatus.FAILURE || mjTask.Status == MJTaskStatus.CANCEL)
|
||||
{
|
||||
// 当前任务已经被释放过了
|
||||
// 开始修改数据
|
||||
mJApiTasks.EndTime = BeijingTimeExtension.GetBeijingTime();
|
||||
await _taskConcurrencyManager.UpdateTaskInDatabase(mJApiTasks);
|
||||
|
||||
return new OkObjectResult(null);
|
||||
}
|
||||
|
||||
if (status == MJTaskStatus.SUCCESS || status == MJTaskStatus.FAILURE || status == MJTaskStatus.CANCEL)
|
||||
{
|
||||
mJApiTasks.EndTime = BeijingTimeExtension.GetBeijingTime();
|
||||
_usageTracker.ReleaseConcurrencyPermit(mjTask.Token);
|
||||
}
|
||||
|
||||
// 开始修改数据
|
||||
await _taskConcurrencyManager.UpdateTaskInDatabase(mJApiTasks);
|
||||
|
||||
return new OkObjectResult(null);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.LogError(ex, "MJNotifyHook 处理回调数据时发生异常");
|
||||
return new StatusCodeResult(StatusCodes.Status500InternalServerError);
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 私有辅助方法
|
||||
|
||||
private static ActionResult<object>? ValidateParameters(string id, string token)
|
||||
{
|
||||
if (string.IsNullOrWhiteSpace(id))
|
||||
return new BadRequestObjectResult(new { error = "任务ID不能为空" });
|
||||
|
||||
if (string.IsNullOrWhiteSpace(token))
|
||||
return new UnauthorizedObjectResult(new { error = "缺少授权Token" });
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
private async Task<ActionResult<object>?> TryOriginApiAsync(string id)
|
||||
{
|
||||
const string originUrl = "https://mjapi.bzu.cn/mj/task/{0}/fetch";
|
||||
|
||||
// 判断 原始token 不存在 直接 返回空
|
||||
string orginToken = await _tokenService.GetOriginToken();
|
||||
if (string.IsNullOrWhiteSpace(orginToken))
|
||||
{
|
||||
return null;
|
||||
}
|
||||
|
||||
try
|
||||
{
|
||||
using var client = CreateHttpClient(orginToken, false);
|
||||
var response = await client.GetAsync(string.Format(originUrl, id));
|
||||
var content = await response.Content.ReadAsStringAsync();
|
||||
if ((int)response.StatusCode == 204 || string.IsNullOrWhiteSpace(content))
|
||||
{
|
||||
throw new Exception("原始API返回204 No Content或空内容");
|
||||
}
|
||||
|
||||
return new ContentResult
|
||||
{
|
||||
Content = content,
|
||||
ContentType = response.Content.Headers.ContentType?.ToString() ?? "application/json",
|
||||
StatusCode = (int)response.StatusCode
|
||||
};
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.LogWarning(ex, "原始API调用失败,TaskId: {TaskId},准备尝试备用API", id);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
private static bool IsRetriableException(Exception ex)
|
||||
{
|
||||
return ex is HttpRequestException ||
|
||||
ex is TaskCanceledException ||
|
||||
ex is SocketException;
|
||||
}
|
||||
|
||||
private async Task<ActionResult<object>?> TryBackupApiAsync(string id, string useToken)
|
||||
{
|
||||
const string backupUrlTemplate = "https://api.laitool.cc/mj/task/{0}/fetch";
|
||||
const int maxRetries = 3;
|
||||
const int baseDelayMs = 1000;
|
||||
|
||||
using var client = CreateHttpClient($"Bearer sk-{useToken}", true);
|
||||
var backupUrl = string.Format(backupUrlTemplate, id);
|
||||
|
||||
for (int attempt = 1; attempt <= maxRetries; attempt++)
|
||||
{
|
||||
try
|
||||
{
|
||||
var response = await client.GetAsync(backupUrl);
|
||||
var content = await response.Content.ReadAsStringAsync();
|
||||
|
||||
_logger.LogInformation("备用API调用成功,TaskId: {TaskId}, Attempt: {Attempt}, StatusCode: {StatusCode}",
|
||||
id, attempt, response.StatusCode);
|
||||
|
||||
return new ContentResult
|
||||
{
|
||||
Content = content,
|
||||
ContentType = response.Content.Headers.ContentType?.ToString() ?? "application/json",
|
||||
StatusCode = (int)response.StatusCode
|
||||
};
|
||||
}
|
||||
catch (Exception ex) when (IsRetriableException(ex))
|
||||
{
|
||||
if (attempt < maxRetries)
|
||||
{
|
||||
var delay = baseDelayMs * (int)Math.Pow(2, attempt - 1);
|
||||
_logger.LogWarning(ex, "备用API调用失败,TaskId: {TaskId}, Attempt: {Attempt}, 将在{Delay}ms后重试",
|
||||
id, attempt, delay);
|
||||
await Task.Delay(delay);
|
||||
}
|
||||
else
|
||||
{
|
||||
_logger.LogError(ex, "备用API调用最终失败,TaskId: {TaskId}, MaxAttempts: {MaxAttempts}",
|
||||
id, maxRetries);
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.LogError(ex, "备用API调用发生不可重试异常,TaskId: {TaskId}, Attempt: {Attempt}",
|
||||
id, attempt);
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
|
||||
private async Task<ActionResult<object>> GetTaskFromDatabaseAsync(string id, string token)
|
||||
{
|
||||
try
|
||||
{
|
||||
// 这里需要根据您的数据库结构调整
|
||||
// 假设您有一个任务表存储MJ任务的状态信息
|
||||
var taskFromDb = await _dbContext.MJApiTasks // 替换为您实际的表名
|
||||
.Where(x => x.ThirdPartyTaskId == id && x.Token == token) // 确保用户只能访问自己的任务
|
||||
.FirstOrDefaultAsync();
|
||||
|
||||
if (taskFromDb != null && !string.IsNullOrWhiteSpace(taskFromDb.Properties))
|
||||
{
|
||||
return new OkObjectResult(taskFromDb.Properties);
|
||||
}
|
||||
else
|
||||
{
|
||||
return new ObjectResult(new
|
||||
{
|
||||
error = "服务暂时不可用",
|
||||
message = "无法获取任务状态,MJ服务连接失败且本地无缓存数据",
|
||||
ThirdPartyTaskId = id,
|
||||
})
|
||||
{
|
||||
StatusCode = 502
|
||||
};
|
||||
}
|
||||
}
|
||||
catch (Exception dbEx)
|
||||
{
|
||||
_logger.LogError(dbEx, $"从数据库获取任务数据时发生异常,TaskId: {id}");
|
||||
|
||||
return new ObjectResult(new
|
||||
{
|
||||
error = "系统异常",
|
||||
message = "获取任务状态失败,请稍后重试",
|
||||
ThirdPartyTaskId = id
|
||||
})
|
||||
{
|
||||
StatusCode = 500
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
private static HttpClient CreateHttpClient(string authorization, bool isBearerToken)
|
||||
{
|
||||
var client = new HttpClient();
|
||||
client.DefaultRequestHeaders.Add("Authorization", authorization);
|
||||
client.Timeout = TimeSpan.FromSeconds(30);
|
||||
return client;
|
||||
}
|
||||
|
||||
#endregion
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,878 @@
|
||||
using LMS.Common.Extensions;
|
||||
using LMS.DAO;
|
||||
using LMS.DAO.UserDAO;
|
||||
using LMS.Repository.DB;
|
||||
using LMS.Repository.DTO;
|
||||
using LMS.Repository.MJPackage;
|
||||
using LMS.Tools.MJPackage;
|
||||
using Microsoft.AspNetCore.Mvc;
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using System.Data;
|
||||
using System.Diagnostics;
|
||||
using static LMS.Common.Enums.ResponseCodeEnum;
|
||||
|
||||
namespace LMS.service.Service.MJPackage
|
||||
{
|
||||
public class TokenManagementService(
|
||||
ApplicationDbContext dbContext,
|
||||
TokenUsageTracker usageTracker,
|
||||
ITokenService tokenService,
|
||||
UserBasicDao userBasicDao,
|
||||
ILogger<TokenManagementService> logger) : ITokenManagementService
|
||||
{
|
||||
|
||||
private readonly ApplicationDbContext _dbContext = dbContext;
|
||||
private readonly TokenUsageTracker _usageTracker = usageTracker;
|
||||
private readonly ITokenService _tokenService = tokenService;
|
||||
private readonly ILogger<TokenManagementService> _logger = logger;
|
||||
private readonly UserBasicDao _userBasicDao = userBasicDao;
|
||||
|
||||
#region 使用Token查询对应的任务-用户
|
||||
/// <summary>
|
||||
/// 查询任务集合 通过参数
|
||||
/// </summary>
|
||||
/// <param name="token"></param>
|
||||
/// <param name="page"></param>
|
||||
/// <param name="pageSize"></param>
|
||||
/// <param name="thirdPartyTaskId"></param>
|
||||
/// <returns></returns>
|
||||
public async Task<ActionResult<APIResponseModel<CollectionResponse<TokenAndTaskCollection>>>> QueryTokenTaskCollection(string token, int page, int pageSize, string? thirdPartyTaskId)
|
||||
{
|
||||
try
|
||||
{
|
||||
if (string.IsNullOrWhiteSpace(token))
|
||||
{
|
||||
return APIResponseModel<CollectionResponse<TokenAndTaskCollection>>.CreateErrorResponseModel(ResponseCode.ParameterError, "Token不能为空");
|
||||
}
|
||||
|
||||
TokenCacheItem? tokenCache = await _tokenService.GetDatabaseTokenAsync(token, true);
|
||||
if (tokenCache == null)
|
||||
{
|
||||
return APIResponseModel<CollectionResponse<TokenAndTaskCollection>>.CreateErrorResponseModel(ResponseCode.ParameterError, "Token不存在或已过期");
|
||||
}
|
||||
|
||||
// 处理token数据
|
||||
TokenAndTaskCollection tokenCacheItem = new()
|
||||
{
|
||||
Id = tokenCache.Id,
|
||||
Token = tokenCache.Token,
|
||||
DailyLimit = tokenCache.DailyLimit,
|
||||
TotalLimit = tokenCache.TotalLimit,
|
||||
ConcurrencyLimit = tokenCache.ConcurrencyLimit,
|
||||
CreatedAt = tokenCache.CreatedAt,
|
||||
ExpiresAt = tokenCache.ExpiresAt,
|
||||
DailyUsage = tokenCache.DailyUsage,
|
||||
TotalUsage = tokenCache.TotalUsage,
|
||||
LastActivityTime = tokenCache.LastActivityTime
|
||||
};
|
||||
|
||||
|
||||
var (maxCount, currentlyExecuting, available) = _usageTracker.GetConcurrencyStatus(tokenCache.Token);
|
||||
tokenCacheItem.CurrentlyExecuting = currentlyExecuting;
|
||||
// 开始处理 task 数据
|
||||
IQueryable<MJApiTasks> query = _dbContext.MJApiTasks.Where(x => x.TokenId == tokenCacheItem.Id);
|
||||
|
||||
if (!string.IsNullOrWhiteSpace(thirdPartyTaskId))
|
||||
{
|
||||
query = query.Where(x => x.ThirdPartyTaskId == thirdPartyTaskId);
|
||||
}
|
||||
|
||||
int total = await query.CountAsync();
|
||||
List<MJApiTasks> mJApiTasks = await query.OrderByDescending(x => x.StartTime).Skip((page - 1) * pageSize).Take(pageSize).ToListAsync();
|
||||
|
||||
// 处理任务集合
|
||||
// 将某个属性设置为空值
|
||||
foreach (var task in mJApiTasks)
|
||||
{
|
||||
task.Token = "****";
|
||||
}
|
||||
|
||||
tokenCacheItem.TaskCollections = mJApiTasks;
|
||||
|
||||
return APIResponseModel<CollectionResponse<TokenAndTaskCollection>>.CreateSuccessResponseModel(ResponseCode.Success, new CollectionResponse<TokenAndTaskCollection>
|
||||
{
|
||||
Total = total,
|
||||
Collection = [tokenCacheItem],
|
||||
Current = page
|
||||
});
|
||||
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
return APIResponseModel<CollectionResponse<TokenAndTaskCollection>>.CreateErrorResponseModel(Common.Enums.ResponseCodeEnum.ResponseCode.SystemError, ex.Message);
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 查询Token是不是存在-用户
|
||||
/// <summary>
|
||||
/// 获取Token的数据
|
||||
/// </summary>
|
||||
/// <param name="token"></param>
|
||||
/// <returns></returns>
|
||||
public async Task<ActionResult<APIResponseModel<TokenCacheItem>>> GetTokenItem(string token)
|
||||
{
|
||||
try
|
||||
{
|
||||
if (string.IsNullOrWhiteSpace(token))
|
||||
{
|
||||
return APIResponseModel<TokenCacheItem>.CreateErrorResponseModel(ResponseCode.ParameterError, "Token不能为空");
|
||||
}
|
||||
|
||||
TokenCacheItem? tokenCache = await _tokenService.GetDatabaseTokenAsync(token, false);
|
||||
if (tokenCache == null)
|
||||
{
|
||||
return APIResponseModel<TokenCacheItem>.CreateErrorResponseModel(ResponseCode.ParameterError, "Token不存在");
|
||||
}
|
||||
|
||||
var (maxCount, currentlyExecuting, available) = _usageTracker.GetConcurrencyStatus(tokenCache.Token);
|
||||
tokenCache.CurrentlyExecuting = currentlyExecuting;
|
||||
tokenCache.UseToken = "*******************";
|
||||
return APIResponseModel<TokenCacheItem>.CreateSuccessResponseModel(ResponseCode.Success, tokenCache);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
return APIResponseModel<TokenCacheItem>.CreateErrorResponseModel(ResponseCode.SystemError, ex.Message);
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 新增Token
|
||||
|
||||
public async Task<ActionResult<APIResponseModel<string>>> AddToken(long requestUserId, AddOrModifyTokenModel model)
|
||||
{
|
||||
try
|
||||
{
|
||||
bool isSuperAdmin = await _userBasicDao.CheckUserIsSuperAdmin(requestUserId);
|
||||
if (!isSuperAdmin)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.NotPermissionAction);
|
||||
}
|
||||
|
||||
// 是 超级管理员 直接添加数据就行
|
||||
if (string.IsNullOrWhiteSpace(model.Token))
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "Token不能为空");
|
||||
}
|
||||
|
||||
if (string.IsNullOrWhiteSpace(model.UseToken))
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "使用Token不能为空");
|
||||
}
|
||||
|
||||
// 判断token是不是存在
|
||||
MJApiTokens? exitMJApiTokens = await _dbContext.MJApiTokens
|
||||
.AsNoTracking()
|
||||
.FirstOrDefaultAsync(x => x.Token == model.Token || x.UseToken == model.UseToken);
|
||||
if (exitMJApiTokens != null)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "Token或使用Token已存在,请检查数据!");
|
||||
}
|
||||
|
||||
if (model.DailyLimit < 0 || model.TotalLimit < 0 || model.ConcurrencyLimit < 0 || model.UseDayCount < 0)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "限制参数不能小于0");
|
||||
}
|
||||
|
||||
MJApiTokens mJApiTokens = new()
|
||||
{
|
||||
Token = model.Token,
|
||||
UseToken = model.UseToken,
|
||||
DailyLimit = model.DailyLimit,
|
||||
TotalLimit = model.TotalLimit,
|
||||
ConcurrencyLimit = model.ConcurrencyLimit,
|
||||
CreatedAt = BeijingTimeExtension.GetBeijingTime(),
|
||||
ExpiresAt = BeijingTimeExtension.GetBeijingTime().AddDays(model.UseDayCount)
|
||||
};
|
||||
|
||||
// 开始新增
|
||||
await _dbContext.MJApiTokens.AddAsync(mJApiTokens);
|
||||
await _dbContext.SaveChangesAsync();
|
||||
string message = $"Token添加成功: {model.Token}, 日限制: {model.DailyLimit}, 总限制: {model.TotalLimit}, 并发限制: {model.ConcurrencyLimit},有效期: {mJApiTokens.ExpiresAt}";
|
||||
return APIResponseModel<string>.CreateSuccessResponseModel(ResponseCode.Success, message);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.SystemError, ex.Message);
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 修改Token
|
||||
|
||||
public async Task<ActionResult<APIResponseModel<string>>> ModifyToken(long requestUserId, long tokenId, AddOrModifyTokenModel model)
|
||||
{
|
||||
try
|
||||
{
|
||||
if (string.IsNullOrWhiteSpace(model.Token))
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "Token不能为空");
|
||||
}
|
||||
if (string.IsNullOrWhiteSpace(model.UseToken))
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "使用Token不能为空");
|
||||
}
|
||||
|
||||
bool isSuperAdmin = await _userBasicDao.CheckUserIsSuperAdmin(requestUserId);
|
||||
if (!isSuperAdmin)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.NotPermissionAction);
|
||||
}
|
||||
|
||||
// 是 超级管理员 直接修改数据就行
|
||||
MJApiTokens? mJApiTokens = await _dbContext.MJApiTokens
|
||||
.Where(x => x.Id == tokenId).FirstOrDefaultAsync();
|
||||
|
||||
if (mJApiTokens == null)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "Token不存在");
|
||||
}
|
||||
|
||||
// 判断当前传入的token 是不是已经再别的地方使用了
|
||||
MJApiTokens? exitMJApiTokens = await _dbContext.MJApiTokens
|
||||
.AsNoTracking()
|
||||
.FirstOrDefaultAsync(x => (x.Token == model.Token || x.UseToken == model.UseToken) && x.Id != tokenId);
|
||||
if (exitMJApiTokens != null)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "Token或者是使用Token已存在,请检查数据");
|
||||
}
|
||||
|
||||
|
||||
// 开始修改
|
||||
mJApiTokens.Token = model.Token;
|
||||
mJApiTokens.UseToken = model.UseToken;
|
||||
|
||||
if (model.DailyLimit > 0)
|
||||
{
|
||||
mJApiTokens.DailyLimit = model.DailyLimit;
|
||||
}
|
||||
if (model.TotalLimit > 0)
|
||||
{
|
||||
mJApiTokens.TotalLimit = model.TotalLimit;
|
||||
}
|
||||
if (model.ConcurrencyLimit > 0)
|
||||
{
|
||||
mJApiTokens.ConcurrencyLimit = model.ConcurrencyLimit;
|
||||
}
|
||||
|
||||
if (model.UseDayCount > 0)
|
||||
{
|
||||
// 不是 -1 就需要重新设置到期时间
|
||||
mJApiTokens.ExpiresAt = mJApiTokens.CreatedAt.AddDays(model.UseDayCount);
|
||||
}
|
||||
|
||||
_dbContext.MJApiTokens.Update(mJApiTokens);
|
||||
await _dbContext.SaveChangesAsync();
|
||||
|
||||
// 刷新一下内存中的限制
|
||||
await RefreshTokenFromDatabaseAsync(model.Token);
|
||||
|
||||
string message = $"Token修改成功: {model.Token}, 日限制: {mJApiTokens.DailyLimit}, 总限制: {mJApiTokens.TotalLimit}, 并发限制: {mJApiTokens.ConcurrencyLimit},有效期: {mJApiTokens.ExpiresAt},请注意:修改Token后,之前的Token将失效。";
|
||||
|
||||
return APIResponseModel<string>.CreateSuccessResponseModel(ResponseCode.Success, message);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.SystemError, ex.Message);
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 删除Token
|
||||
|
||||
public async Task<ActionResult<APIResponseModel<string>>> DeleteToken(long requestUserId, long tokenId)
|
||||
{
|
||||
var transaction = await _dbContext.Database.BeginTransactionAsync();
|
||||
try
|
||||
{
|
||||
bool isSuperAdmin = await _userBasicDao.CheckUserIsSuperAdmin(requestUserId);
|
||||
if (!isSuperAdmin)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.NotPermissionAction);
|
||||
}
|
||||
// 是 超级管理员 直接删除数据就行
|
||||
// 先删除 使用数据
|
||||
MJApiTokenUsage? tokenUsage = await _dbContext.MJApiTokenUsage
|
||||
.Where(x => x.TokenId == tokenId)
|
||||
.FirstOrDefaultAsync();
|
||||
if (tokenUsage != null)
|
||||
{
|
||||
_dbContext.MJApiTokenUsage.Remove(tokenUsage);
|
||||
}
|
||||
|
||||
// 再删除 Token 数据
|
||||
MJApiTokens? mJApiTokens = await _dbContext.MJApiTokens
|
||||
.Where(x => x.Id == tokenId).FirstOrDefaultAsync();
|
||||
if (mJApiTokens == null)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "Token不存在");
|
||||
}
|
||||
_dbContext.MJApiTokens.Remove(mJApiTokens);
|
||||
await _dbContext.SaveChangesAsync();
|
||||
await transaction.CommitAsync();
|
||||
// 将内存中的token移除
|
||||
|
||||
_usageTracker.RemoveToken(mJApiTokens.Token);
|
||||
_logger.LogInformation($"Token删除成功: {mJApiTokens.Token},ID: {tokenId}");
|
||||
string message = $"Token删除成功: {mJApiTokens.Token},ID: {tokenId},请注意:删除Token后,之前的Token将失效。";
|
||||
return APIResponseModel<string>.CreateSuccessResponseModel(ResponseCode.Success, message);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
await transaction.RollbackAsync();
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.SystemError, ex.Message);
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 查询所有的任务,可以带参数
|
||||
|
||||
public async Task<ActionResult<APIResponseModel<CollectionResponse<MJApiTaskCollection>>>> QueryTaskCollection(long requestUserId, int page, int pageSize, string? thirdPartyTaskId, string? token, long? tokenId)
|
||||
{
|
||||
try
|
||||
{
|
||||
// 权限检查
|
||||
bool isSuperAdmin = await _userBasicDao.CheckUserIsSuperAdmin(requestUserId);
|
||||
if (!isSuperAdmin)
|
||||
{
|
||||
return APIResponseModel<CollectionResponse<MJApiTaskCollection>>.CreateErrorResponseModel(ResponseCode.NotPermissionAction);
|
||||
}
|
||||
|
||||
// 构建查询条件
|
||||
IQueryable<MJApiTasks> query = _dbContext.MJApiTasks.AsQueryable();
|
||||
|
||||
// 根据不同参数组合构建查询
|
||||
if (tokenId.HasValue)
|
||||
{
|
||||
query = query.Where(x => x.TokenId == tokenId);
|
||||
}
|
||||
|
||||
if (!string.IsNullOrWhiteSpace(token))
|
||||
{
|
||||
// 如果只提供了token字符串
|
||||
query = query.Where(x => x.Token.Contains(token));
|
||||
}
|
||||
|
||||
// 添加第三方任务ID过滤
|
||||
if (!string.IsNullOrWhiteSpace(thirdPartyTaskId))
|
||||
{
|
||||
query = query.Where(x => x.ThirdPartyTaskId == thirdPartyTaskId);
|
||||
}
|
||||
|
||||
// 获取总数
|
||||
int total = await query.CountAsync();
|
||||
|
||||
if (total == 0)
|
||||
{
|
||||
return APIResponseModel<CollectionResponse<MJApiTaskCollection>>.CreateSuccessResponseModel(
|
||||
ResponseCode.Success, new CollectionResponse<MJApiTaskCollection>
|
||||
{
|
||||
Total = 0,
|
||||
Collection = new List<MJApiTaskCollection>(),
|
||||
Current = page
|
||||
});
|
||||
}
|
||||
|
||||
// 分页查询任务数据
|
||||
var tasks = await query
|
||||
.OrderByDescending(x => x.StartTime)
|
||||
.Skip((page - 1) * pageSize)
|
||||
.Take(pageSize)
|
||||
.ToListAsync();
|
||||
|
||||
// 构建返回结果
|
||||
var result = tasks.Select(task => new MJApiTaskCollection
|
||||
{
|
||||
TaskId = task.TaskId,
|
||||
Token = task.Token,
|
||||
TokenId = task.TokenId,
|
||||
StartTime = task.StartTime,
|
||||
EndTime = task.EndTime,
|
||||
Status = task.Status,
|
||||
ThirdPartyTaskId = task.ThirdPartyTaskId,
|
||||
Properties = task.Properties,
|
||||
}).ToList();
|
||||
|
||||
return APIResponseModel<CollectionResponse<MJApiTaskCollection>>.CreateSuccessResponseModel(
|
||||
ResponseCode.Success, new CollectionResponse<MJApiTaskCollection>
|
||||
{
|
||||
Total = total,
|
||||
Collection = result, // 直接赋值,不要用数组包装
|
||||
Current = page
|
||||
});
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger?.LogError(ex, "查询任务集合时发生错误 - 用户ID: {UserId}, Token: {Token}, TokenId: {TokenId}",
|
||||
requestUserId, token, tokenId);
|
||||
|
||||
return APIResponseModel<CollectionResponse<MJApiTaskCollection>>.CreateErrorResponseModel(
|
||||
ResponseCode.SystemError, "查询失败,请稍后重试");
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 查询所有的Token, 可以带参数
|
||||
public async Task<ActionResult<APIResponseModel<CollectionResponse<TokenCacheItem>>>> QueryTokenCollection(int page, int pageSize, string? token, long? tokenId, bool? efficient, long requestUserId)
|
||||
{
|
||||
try
|
||||
{
|
||||
// 权限检查
|
||||
bool isSuperAdmin = await _userBasicDao.CheckUserIsSuperAdmin(requestUserId);
|
||||
if (!isSuperAdmin)
|
||||
{
|
||||
return APIResponseModel<CollectionResponse<TokenCacheItem>>.CreateErrorResponseModel(ResponseCode.NotPermissionAction);
|
||||
}
|
||||
|
||||
var currentUtcTime = BeijingTimeExtension.GetBeijingTime(); // 2025-06-09 12:20:40
|
||||
|
||||
// 构建基础查询
|
||||
var query = from t in _dbContext.MJApiTokens
|
||||
join u in _dbContext.MJApiTokenUsage on t.Id equals u.TokenId into tokenUsage
|
||||
from usage in tokenUsage.DefaultIfEmpty()
|
||||
select new
|
||||
{
|
||||
Token = t,
|
||||
Usage = usage
|
||||
};
|
||||
|
||||
// 应用过滤条件
|
||||
if (!string.IsNullOrWhiteSpace(token))
|
||||
{
|
||||
query = query.Where(x => x.Token.Token == token);
|
||||
}
|
||||
|
||||
if (tokenId.HasValue)
|
||||
{
|
||||
query = query.Where(x => x.Token.Id == tokenId.Value);
|
||||
}
|
||||
|
||||
if (efficient.HasValue)
|
||||
{
|
||||
if (efficient.Value)
|
||||
{
|
||||
// 查询有效得token
|
||||
query = query.Where(x =>
|
||||
(x.Token.ExpiresAt == null || x.Token.ExpiresAt > currentUtcTime));
|
||||
}
|
||||
else
|
||||
{
|
||||
// 查询无效或不活跃的Token
|
||||
query = query.Where(x =>
|
||||
(x.Token.ExpiresAt != null && x.Token.ExpiresAt <= currentUtcTime));
|
||||
}
|
||||
}
|
||||
|
||||
// 获取总数
|
||||
var total = await query.CountAsync();
|
||||
|
||||
if (total == 0)
|
||||
{
|
||||
return APIResponseModel<CollectionResponse<TokenCacheItem>>.CreateSuccessResponseModel(
|
||||
ResponseCode.Success, new CollectionResponse<TokenCacheItem>
|
||||
{
|
||||
Total = 0,
|
||||
Collection = [],
|
||||
Current = page
|
||||
});
|
||||
}
|
||||
|
||||
// 分页查询并投影到TokenCacheItem
|
||||
var tokenItems = await query
|
||||
.OrderByDescending(x => x.Token.CreatedAt)
|
||||
.Skip((page - 1) * pageSize)
|
||||
.Take(pageSize)
|
||||
.Select(x => new TokenCacheItem
|
||||
{
|
||||
Id = x.Token.Id,
|
||||
Token = x.Token.Token,
|
||||
UseToken = x.Token.UseToken,
|
||||
DailyLimit = x.Token.DailyLimit,
|
||||
TotalLimit = x.Token.TotalLimit,
|
||||
ConcurrencyLimit = x.Token.ConcurrencyLimit,
|
||||
CreatedAt = x.Token.CreatedAt,
|
||||
ExpiresAt = x.Token.ExpiresAt,
|
||||
DailyUsage = x.Usage != null ? x.Usage.DailyUsage : 0,
|
||||
TotalUsage = x.Usage != null ? x.Usage.TotalUsage : 0,
|
||||
LastActivityTime = x.Usage != null ? x.Usage.LastActivityAt : x.Token.CreatedAt,
|
||||
HistoryUse = x.Usage != null ? x.Usage.HistoryUse : ""
|
||||
})
|
||||
.ToListAsync();
|
||||
|
||||
for (int i = 0; i < tokenItems.Count; i++)
|
||||
{
|
||||
var tokenItem = tokenItems[i];
|
||||
var (maxCount, currentlyExecuting, available) = usageTracker.GetConcurrencyStatus(tokenItem.Token);
|
||||
tokenItems[i].CurrentlyExecuting = currentlyExecuting;
|
||||
}
|
||||
|
||||
_logger?.LogInformation($"✅ Token查询完成, 总数: {total}, 当前页: {page}, 页大小: {pageSize}, 返回: {tokenItems.Count} 条");
|
||||
|
||||
return APIResponseModel<CollectionResponse<TokenCacheItem>>.CreateSuccessResponseModel(
|
||||
ResponseCode.Success, new CollectionResponse<TokenCacheItem>
|
||||
{
|
||||
Total = total,
|
||||
Collection = tokenItems,
|
||||
Current = page
|
||||
});
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger?.LogError(ex, "❌ 查询Token集合时发生错误 - 用户: qiang-lo, UTC时间: 2025-06-09 12:20:40");
|
||||
return APIResponseModel<CollectionResponse<TokenCacheItem>>.CreateErrorResponseModel(
|
||||
ResponseCode.SystemError, "查询失败,请稍后重试");
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
|
||||
#region 删除小于指定时间的任务
|
||||
|
||||
public async Task<ActionResult<APIResponseModel<string>>> DeleteMJTaskEarlyTimestamp(long requestUserId, long timestamp)
|
||||
{
|
||||
const string operationName = "删除早期MJ任务";
|
||||
try
|
||||
{
|
||||
// 1. 参数验证
|
||||
if (!IsValidTimestamp(timestamp))
|
||||
{
|
||||
_logger.LogWarning($"{operationName} - 无效的时间戳: {timestamp}");
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "时间戳格式错误");
|
||||
}
|
||||
|
||||
// 2. 权限检查
|
||||
if (!await _userBasicDao.CheckUserIsSuperAdmin(requestUserId))
|
||||
{
|
||||
_logger.LogWarning($"{operationName} - 用户 {requestUserId} 无权限执行此操作");
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.NotPermissionAction);
|
||||
}
|
||||
// 3. 时间戳转换
|
||||
DateTime targetDateTime = DateTimeOffset.FromUnixTimeMilliseconds(timestamp).DateTime;
|
||||
_logger.LogInformation($"{operationName} - 开始删除早于 {targetDateTime:yyyy-MM-dd HH:mm:ss} 的任务");
|
||||
|
||||
DateTime beijingTargetDateTime = targetDateTime.AddHours(8); // 转换为北京时间
|
||||
// 4. 使用批量删除,避免加载到内存
|
||||
int deletedCount = await _dbContext.MJApiTasks
|
||||
.Where(x => x.StartTime <= beijingTargetDateTime)
|
||||
.ExecuteDeleteAsync(); // 使用 EF Core 7+ 的批量删除
|
||||
|
||||
// 5. 记录结果
|
||||
string resultMessage = deletedCount == 0
|
||||
? "没有找到符合条件的任务"
|
||||
: $"成功删除 {deletedCount} 个任务";
|
||||
|
||||
_logger.LogInformation($"{operationName} - {resultMessage},目标时间: {targetDateTime:yyyy-MM-dd HH:mm:ss}");
|
||||
|
||||
return APIResponseModel<string>.CreateSuccessResponseModel(ResponseCode.Success, resultMessage);
|
||||
}
|
||||
catch (ArgumentOutOfRangeException ex)
|
||||
{
|
||||
_logger.LogError(ex, $"{operationName} - 时间戳转换失败: {timestamp}");
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "时间戳格式错误");
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.LogError(ex, $"{operationName} - 执行失败,时间戳: {timestamp}");
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.SystemError, "系统错误,请稍后重试");
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 删除指定ID的任务
|
||||
|
||||
public async Task<ActionResult<APIResponseModel<string>>> DeleteMJTask(long requestUserId, string taskId)
|
||||
{
|
||||
try
|
||||
{
|
||||
bool isSuperAdmin = await _userBasicDao.CheckUserIsSuperAdmin(requestUserId);
|
||||
if (!isSuperAdmin)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.NotPermissionAction);
|
||||
}
|
||||
// 1. 参数验证
|
||||
if (string.IsNullOrWhiteSpace(taskId))
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "任务ID不能为空");
|
||||
}
|
||||
// 2. 查找任务
|
||||
var task = await _dbContext.MJApiTasks
|
||||
.Where(x => x.TaskId == taskId)
|
||||
.FirstOrDefaultAsync();
|
||||
if (task == null)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "任务不存在");
|
||||
}
|
||||
// 3. 删除任务
|
||||
_dbContext.MJApiTasks.Remove(task);
|
||||
await _dbContext.SaveChangesAsync();
|
||||
_logger.LogInformation($"删除指定ID的任务成功,任务ID: {taskId}");
|
||||
return APIResponseModel<string>.CreateSuccessResponseModel(ResponseCode.Success, "任务删除成功");
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.LogError(ex, $"删除指定ID的任务失败,任务ID: {taskId}, 失败原因:{ex.Message}");
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.SystemError, "删除任务失败,请稍后重试");
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 手动刷新Token
|
||||
|
||||
public async Task<ActionResult<APIResponseModel<string>>> RefreshToken(long requestUserId, string token)
|
||||
{
|
||||
try
|
||||
{
|
||||
bool isSuperAdmin = await _userBasicDao.CheckUserIsSuperAdmin(requestUserId);
|
||||
if (!isSuperAdmin)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.NotPermissionAction);
|
||||
}
|
||||
// 1. 参数验证
|
||||
if (string.IsNullOrWhiteSpace(token))
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "Token不能为空");
|
||||
}
|
||||
TokenCacheItem tokenCache = await RefreshTokenFromDatabaseAsync(token);
|
||||
if (tokenCache == null)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.ParameterError, "Token不存在");
|
||||
}
|
||||
return APIResponseModel<string>.CreateSuccessResponseModel(ResponseCode.Success, $"Token刷新成功: {token}, 并发限制: {tokenCache.ConcurrencyLimit},日出图限制: {tokenCache.DailyLimit}, 当前日出图总数: {tokenCache.DailyUsage}, 总出图量: {tokenCache.TotalUsage}");
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
string message = $"手动刷新Token失败,Token: {token}, 失败原因:{ex.Message}";
|
||||
_logger.LogError(ex, message);
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.SystemError, message);
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 获取活跃的Token
|
||||
|
||||
public async Task<ActionResult<APIResponseModel<CollectionResponse<TokenCacheItem>>>> GetActiveTokens(long requestUserId, int minutes)
|
||||
{
|
||||
try
|
||||
{
|
||||
bool isSuperAdmin = await _userBasicDao.CheckUserIsSuperAdmin(requestUserId);
|
||||
if (!isSuperAdmin)
|
||||
{
|
||||
return APIResponseModel<CollectionResponse<TokenCacheItem>>.CreateErrorResponseModel(ResponseCode.NotPermissionAction);
|
||||
}
|
||||
var threshold = TimeSpan.FromMinutes(minutes);
|
||||
List<TokenCacheItem> activeTokens = _usageTracker.GetActiveTokens(threshold);
|
||||
return APIResponseModel<CollectionResponse<TokenCacheItem>>.CreateSuccessResponseModel(ResponseCode.Success, new CollectionResponse<TokenCacheItem>
|
||||
{
|
||||
Total = activeTokens.Count,
|
||||
Collection = activeTokens,
|
||||
Current = 1 // 这里可以设置为1,因为我们只返回一页
|
||||
});
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
string message = $"获取活跃的Token失败,错误信息:{ex.Message}";
|
||||
_logger.LogError(ex, message);
|
||||
return APIResponseModel<CollectionResponse<TokenCacheItem>>.CreateErrorResponseModel(ResponseCode.SystemError, message);
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 移除不活跃的Token
|
||||
|
||||
public async Task<ActionResult<APIResponseModel<string>>> RemoveNotActiveTokens(long requestUserId, int minutes)
|
||||
{
|
||||
try
|
||||
{
|
||||
bool isSuperAdmin = await _userBasicDao.CheckUserIsSuperAdmin(requestUserId);
|
||||
if (!isSuperAdmin)
|
||||
{
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.NotPermissionAction);
|
||||
}
|
||||
var threshold = TimeSpan.FromMinutes(minutes);
|
||||
var (activateTokenCount, notActivateTokenCount) = _usageTracker.RemoveNotActiveTokens(threshold);
|
||||
string message = $"删除不活跃得 Token 数: {notActivateTokenCount},阈值: {minutes}分钟";
|
||||
_logger.LogInformation(message);
|
||||
return APIResponseModel<string>.CreateSuccessResponseModel(ResponseCode.Success, message);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
string message = $"移除不活跃的Token失败,错误信息:{ex.Message}";
|
||||
_logger.LogError(ex, message);
|
||||
return APIResponseModel<string>.CreateErrorResponseModel(ResponseCode.SystemError, message);
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 健康检查和统计接口
|
||||
public async Task<ActionResult<APIResponseModel<object>>> GetHealth(long requestUserId)
|
||||
{
|
||||
try
|
||||
{
|
||||
bool isSuperAdmin = await _userBasicDao.CheckUserIsSuperAdmin(requestUserId);
|
||||
if (!isSuperAdmin)
|
||||
{
|
||||
return APIResponseModel<object>.CreateErrorResponseModel(ResponseCode.NotPermissionAction);
|
||||
}
|
||||
|
||||
var stats = _usageTracker.GetCacheStats();
|
||||
var now = BeijingTimeExtension.GetBeijingTime();
|
||||
return APIResponseModel<object>.CreateSuccessResponseModel(ResponseCode.Success, new
|
||||
{
|
||||
Status = "Healthy",
|
||||
Timestamp = now,
|
||||
CacheStats = stats,
|
||||
Uptime = now - Process.GetCurrentProcess().StartTime.ToUniversalTime()
|
||||
});
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.LogError(ex, "获取系统健康状态失败");
|
||||
return APIResponseModel<object>.CreateErrorResponseModel(ResponseCode.SystemError, "获取系统健康状态失败");
|
||||
}
|
||||
}
|
||||
#endregion
|
||||
|
||||
#region 获取指定ID得Token
|
||||
|
||||
/// <summary>
|
||||
/// 获取指定ID得Token
|
||||
/// </summary>
|
||||
/// <param name="tokenId"></param>
|
||||
/// <param name="requestUserId"></param>
|
||||
/// <returns></returns>
|
||||
public async Task<ActionResult<APIResponseModel<MJApiTokens>>> QueryTokenById(long tokenId, long requestUserId)
|
||||
{
|
||||
try
|
||||
{
|
||||
|
||||
bool isSuperAdmin = await _userBasicDao.CheckUserIsSuperAdmin(requestUserId);
|
||||
if (!isSuperAdmin)
|
||||
{
|
||||
return APIResponseModel<MJApiTokens>.CreateErrorResponseModel(ResponseCode.NotPermissionAction);
|
||||
}
|
||||
|
||||
MJApiTokens? mjApiTokens = await _dbContext.MJApiTokens
|
||||
.AsNoTracking()
|
||||
.Where(x => x.Id == tokenId)
|
||||
.FirstOrDefaultAsync();
|
||||
if (mjApiTokens == null)
|
||||
{
|
||||
return APIResponseModel<MJApiTokens>.CreateErrorResponseModel(ResponseCode.ParameterError, "Token不存在");
|
||||
}
|
||||
return APIResponseModel<MJApiTokens>.CreateSuccessResponseModel(ResponseCode.Success, mjApiTokens);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.LogError(ex, "查询Token失败,TokenId: {TokenId}, 错误信息: {Message}", tokenId, ex.Message);
|
||||
return APIResponseModel<MJApiTokens>.CreateErrorResponseModel(ResponseCode.SystemError, "查询Token失败,请稍后重试");
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
#region 获取日统计数据
|
||||
|
||||
public async Task<ActionResult<APIResponseModel<TaskStatistics>>> GetDayTaskStatistics(long requestUserId)
|
||||
{
|
||||
try
|
||||
{
|
||||
bool isSuperAdmin = await _userBasicDao.CheckUserIsSuperAdmin(requestUserId);
|
||||
if (!isSuperAdmin)
|
||||
{
|
||||
return APIResponseModel<TaskStatistics>.CreateErrorResponseModel(ResponseCode.NotPermissionAction);
|
||||
}
|
||||
// 获取当前日期
|
||||
DateTime today = BeijingTimeExtension.GetBeijingTime().Date; // 获取今天的日期,不包含时间部分
|
||||
// 查询今天的任务
|
||||
var tasks = await _dbContext.MJApiTasks
|
||||
.Where(x => x.StartTime.Date == today)
|
||||
.OrderByDescending(x => x.StartTime)
|
||||
.ToListAsync();
|
||||
if (tasks.Count == 0)
|
||||
{
|
||||
return APIResponseModel<TaskStatistics>.CreateSuccessResponseModel(
|
||||
ResponseCode.Success, new TaskStatistics());
|
||||
}
|
||||
|
||||
// 统计任务数量
|
||||
int totalTasks = tasks.Count;
|
||||
// 统计成功任务数量
|
||||
int successfulTasks = tasks.Count(x => x.Status == MJTaskStatus.SUCCESS);
|
||||
// 统计失败任务数量
|
||||
int failedTasks = tasks.Count(x => x.Status == MJTaskStatus.FAILURE || x.Status == MJTaskStatus.CANCEL);
|
||||
|
||||
// 剩下的都是再执行的
|
||||
int inProgressTasks = totalTasks - successfulTasks - failedTasks;
|
||||
return APIResponseModel<TaskStatistics>.CreateSuccessResponseModel(ResponseCode.Success, new TaskStatistics
|
||||
{
|
||||
TotalTasks = totalTasks,
|
||||
CompletedTasks = successfulTasks,
|
||||
FailedTasks = failedTasks,
|
||||
InProgressTasks = inProgressTasks
|
||||
});
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.LogError(ex, "获取日统计数据失败,错误信息: {Message}", ex.Message);
|
||||
return APIResponseModel<TaskStatistics>.CreateErrorResponseModel(ResponseCode.SystemError, "获取日统计数据失败,请稍后重试");
|
||||
}
|
||||
}
|
||||
|
||||
#endregion
|
||||
|
||||
|
||||
/// <summary>
|
||||
/// 从数据库重新加载Token到内存缓存(平滑更新)
|
||||
/// </summary>
|
||||
private async Task<TokenCacheItem> RefreshTokenFromDatabaseAsync(string token)
|
||||
{
|
||||
_logger.LogDebug($"从数据库平滑刷新Token到内存: {token}");
|
||||
|
||||
try
|
||||
{
|
||||
TokenCacheItem? tokenItem = await _tokenService.GetDatabaseTokenAsync(token);
|
||||
if (tokenItem == null)
|
||||
{
|
||||
// 将内存的中的这个token删掉
|
||||
_usageTracker.RemoveToken(token);
|
||||
throw new Exception($"Token不存在");
|
||||
}
|
||||
|
||||
// 3. 平滑更新到内存缓存(这里会自动处理并发限制的平滑调整)
|
||||
_usageTracker.AddOrUpdateTokenAsync(tokenItem);
|
||||
|
||||
_logger.LogInformation($"Token平滑刷新成功: {token}, 并发限制: {tokenItem.ConcurrencyLimit}");
|
||||
return tokenItem;
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.LogError(ex, $"平滑刷新Token失败: {token} ," + ex.Message);
|
||||
throw;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
// 辅助方法:验证时间戳
|
||||
private static bool IsValidTimestamp(long timestamp)
|
||||
{
|
||||
// 检查时间戳是否在合理范围内
|
||||
// Unix时间戳最小值(1970-01-01)和最大值(约2038年或更远)
|
||||
const long minTimestamp = 0;
|
||||
const long maxTimestamp = 253402300799999; // 9999-12-31的毫秒时间戳
|
||||
|
||||
return timestamp >= minTimestamp && timestamp <= maxTimestamp;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -9,8 +9,6 @@ using LMS.Repository.Models.DB;
|
||||
using Microsoft.AspNetCore.Identity;
|
||||
using Microsoft.AspNetCore.Mvc;
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using System.Collections.Generic;
|
||||
using System.Linq;
|
||||
using static LMS.Common.Enums.MachineEnum;
|
||||
using static LMS.Common.Enums.ResponseCodeEnum;
|
||||
using static LMS.Repository.DTO.MachineDto;
|
||||
|
||||
@@ -1,25 +1,22 @@
|
||||
using AutoMapper;
|
||||
using LinqKit;
|
||||
using LMS.Common.Dictionary;
|
||||
using LMS.Common.Enums;
|
||||
using LMS.Common.Extensions;
|
||||
using LMS.Common.Templates;
|
||||
using LMS.DAO;
|
||||
using LMS.DAO.UserDAO;
|
||||
using LMS.Repository.DB;
|
||||
using LMS.Repository.DTO;
|
||||
using LMS.Repository.DTO.OptionDto;
|
||||
using LMS.Repository.Models.DB;
|
||||
using LMS.Repository.Options;
|
||||
using LMS.service.Extensions.Mail;
|
||||
using Microsoft.AspNetCore.Identity;
|
||||
using Microsoft.AspNetCore.Mvc;
|
||||
using Microsoft.EntityFrameworkCore;
|
||||
using Microsoft.Extensions.Options;
|
||||
using System.Linq;
|
||||
using System.Linq.Dynamic.Core;
|
||||
using static LMS.Common.Enums.ResponseCodeEnum;
|
||||
using Options = LMS.Repository.DB.Options;
|
||||
using System.Linq.Dynamic.Core;
|
||||
using LinqKit;
|
||||
using LMS.Repository.DTO.OptionDto;
|
||||
using LMS.Common.Extensions;
|
||||
|
||||
namespace LMS.service.Service
|
||||
{
|
||||
|
||||
@@ -68,6 +68,6 @@
|
||||
}
|
||||
]
|
||||
},
|
||||
"Version": "1.1.1",
|
||||
"Version": "1.1.2",
|
||||
"AllowedHosts": "*"
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user