网关审核修复: ILogger统一注入+调试日志+日志等级控制+RateLimiter PeriodicTimer+Owl分页修复+MC4 MinValue修复+告警确认错误检查+命名规范化
This commit is contained in:
@@ -0,0 +1,535 @@
|
||||
# 网关项目深度代码审核报告
|
||||
|
||||
> 日期: 2026-06-04 | 项目: gateway/ | 扫描: 38 源文件, ~2800 行 | 问题: 30 项
|
||||
|
||||
---
|
||||
|
||||
## 1. Core/Infrastructure/AdapterRegistry.cs — 2 项
|
||||
|
||||
### AR1 [🟡] `GetOnlineAdapters()` 方法名误导
|
||||
|
||||
**位置**: 第 49-50 行
|
||||
```csharp
|
||||
public IReadOnlyList<IGatewayAdapter> GetOnlineAdapters()
|
||||
=> _adapters.AsReadOnly();
|
||||
```
|
||||
返回所有已注册适配器,**不做任何在线/离线判断**。调用者会误以为返回的是在线适配器。
|
||||
|
||||
**修复**:
|
||||
```csharp
|
||||
public IReadOnlyList<IGatewayAdapter> GetAllAdapters()
|
||||
=> _adapters.AsReadOnly();
|
||||
```
|
||||
同时 grep 全项目无任何代码调用此方法,可直接删除(或保留为兼容性方法并标注 `[Obsolete]`)。
|
||||
|
||||
### AR2 [🟡] `FindByCode<T>` O(n) 查找无缓存
|
||||
|
||||
每次 B 路由请求都执行 `_adapters.FirstOrDefault(...)`,适配器数量少时影响可忽略(<10 个),但如果未来扩展到 50+ 适配器会有性能问题。
|
||||
|
||||
**修复** — 加字典缓存:
|
||||
```csharp
|
||||
private readonly Dictionary<string, IGatewayAdapter> _byCode = new();
|
||||
public void Register(IGatewayAdapter adapter)
|
||||
{
|
||||
_adapters.Add(adapter);
|
||||
_byCode[adapter.AdapterCode] = adapter;
|
||||
}
|
||||
public T? FindByCode<T>(string adapterCode) where T : class, IGatewayAdapter
|
||||
=> _byCode.TryGetValue(adapterCode, out var a) ? a as T : null;
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 2. Core/Infrastructure/RateLimiter.cs — 1 项
|
||||
|
||||
### RL1 [🟠] `Task.Run` 无限制创建后台任务
|
||||
|
||||
**位置**: 第 30-32 行
|
||||
```csharp
|
||||
public async Task WaitAsync(CancellationToken ct = default)
|
||||
{
|
||||
await _semaphore.WaitAsync(ct);
|
||||
_ = Task.Run(async () => { await Task.Delay(_intervalMs, ct); try { _semaphore.Release(); } catch { } }, ct);
|
||||
}
|
||||
```
|
||||
每次等待都创建一个新的 `Task.Run`。在 KMS 5 QPS 配置下每秒创建 5 个 Task,5 个适配器 = 25 task/s。虽然 Task 本身轻量,但 `Task.Delay` + `Release` 的火焰纹章在极端高并发下可能导致线程池饥饿。
|
||||
|
||||
**修复** — 使用 `PeriodicTimer` 或固定数量的后台任务:
|
||||
```csharp
|
||||
private readonly SemaphoreSlim _semaphore;
|
||||
private readonly int _maxTokens;
|
||||
|
||||
public RateLimiter(int tokensPerSecond)
|
||||
{
|
||||
_maxTokens = tokensPerSecond;
|
||||
_semaphore = new SemaphoreSlim(_maxTokens, _maxTokens);
|
||||
// 单后台任务持续补充令牌
|
||||
_ = Task.Run(async () =>
|
||||
{
|
||||
var interval = 1000 / _maxTokens;
|
||||
using var timer = new PeriodicTimer(TimeSpan.FromMilliseconds(interval));
|
||||
while (await timer.WaitForNextTickAsync())
|
||||
{
|
||||
try { if (_semaphore.CurrentCount < _maxTokens) _semaphore.Release(); }
|
||||
catch (ObjectDisposedException) { break; }
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
public async Task WaitAsync(CancellationToken ct = default)
|
||||
=> await _semaphore.WaitAsync(ct);
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 3. Core/Models/AdapterCapabilities.cs — 1 项
|
||||
|
||||
### AC1 [🟡] `HasPtz` 与 `AcceptsControl` 字段冗余
|
||||
|
||||
**位置**: 第 10-13 行
|
||||
```csharp
|
||||
public bool HasPtz { get; set; }
|
||||
public bool AcceptsControl { get; set; }
|
||||
```
|
||||
`HasPtz` 是 `HasStreams` 的子能力,`AcceptsControl` 是 `IAcceptsControl` 接口的声明。这两个字段在 `AdapterCapabilities` 中冗余——`IAcceptsControl` 接口本身已经声明了控制能力。
|
||||
|
||||
**修复** — 删除 `AcceptsControl` 字段(能力声明应直接检查接口实现):
|
||||
```csharp
|
||||
// 删除 AcceptsControl、HasPtz
|
||||
// 网关 B 路由中改为: a is IAcceptsControl
|
||||
```
|
||||
|
||||
**跨文件影响**: `gateway/src/IntegrationGateway.Host/Program.cs` — B10 路由已用 `FindByCode<IAcceptsControl>` 不再需要此字段。OwlAdapter/Mc4Adapter 的 Capabilities 声明中删除对应行。
|
||||
|
||||
---
|
||||
|
||||
## 4. Core/Models/StandardDevice.cs — 1 项
|
||||
|
||||
### SD1 [🟡] `DeviceId` 字段职责混淆
|
||||
|
||||
**位置**: 第 9 行
|
||||
```csharp
|
||||
public int DeviceId { get; set; }
|
||||
```
|
||||
`DeviceId` 是 VolPro 侧的主键,由 A3 同步后由 VolPro 回填。但网关在 `GetDevicesAsync` 返回时不设置此字段(始终为 0)。文档注释说"同步后由 VolPro 回填",但调用者可能在未同步时就读取此字段。
|
||||
|
||||
**修复** — 改为 `int?` 并初始化为 null:
|
||||
```csharp
|
||||
public int? DeviceId { get; set; }
|
||||
// VolPro 回填后才有值
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 5. Core/Abstractions/IHasRecordings.cs — 1 项
|
||||
|
||||
### IR1 [🟡] 方法参数过多
|
||||
|
||||
**位置**: 第 12-13 行
|
||||
```csharp
|
||||
Task<PagedResult<StandardRecording>> GetRecordingsAsync(
|
||||
string channelId, DateTime start, DateTime end, int page, int size);
|
||||
```
|
||||
5 个位置参数,调用时易错序。
|
||||
|
||||
**修复** — 使用请求对象:
|
||||
```csharp
|
||||
public class RecordingQuery
|
||||
{
|
||||
public string ChannelId { get; set; } = "";
|
||||
public DateTime Start { get; set; }
|
||||
public DateTime End { get; set; }
|
||||
public int Page { get; set; } = 1;
|
||||
public int Size { get; set; } = 20;
|
||||
}
|
||||
Task<PagedResult<StandardRecording>> GetRecordingsAsync(RecordingQuery query);
|
||||
```
|
||||
**跨文件影响**: `OwlAdapter.cs` GetRecordingsAsync 实现 + `Program.cs` B 路由。
|
||||
|
||||
---
|
||||
|
||||
## 6. Core/Infrastructure/GatewayClientFactory.cs — 2 项
|
||||
|
||||
### GF1 [🟡] `JsonDocument?` 返回类型不透明
|
||||
|
||||
**位置**: 第 33-42 行
|
||||
```csharp
|
||||
public async Task<JsonDocument?> RegisterAsync(GatewayRegisterRequest req) { ... }
|
||||
public async Task<JsonDocument?> SyncDevicesAsync(...) { ... }
|
||||
```
|
||||
返回 `JsonDocument?` 使调用者必须知道 VolPro 响应的 JSON 结构。
|
||||
|
||||
**修复** — 定义响应 DTO:
|
||||
```csharp
|
||||
public class RegisterResponse { public int NodeId { get; set; } public List<DeviceSummary> Devices { get; set; } = new(); }
|
||||
public Task<RegisterResponse?> RegisterAsync(GatewayRegisterRequest req) { ... }
|
||||
```
|
||||
|
||||
### GF2 [🟡] `CreateClient()` 每次创建新 HttpClient
|
||||
|
||||
```csharp
|
||||
private HttpClient CreateClient() => _httpFactory.CreateClient("VolPro");
|
||||
```
|
||||
`IHttpClientFactory.CreateClient("VolPro")` 从连接池返回,每次调用都会创建新的 `HttpClient` 包装实例。虽然底层 SocketsHttpHandler 复用连接,但 `HttpClient` 上设置的 Timeout/Headers 不会被保留。
|
||||
|
||||
**修复** — 直接使用工厂客户端(已配置好 Timeout/Headers):
|
||||
```csharp
|
||||
// Program.cs 中注册时已设置 Timeout=30s + Accept: application/json
|
||||
// 直接使用 factory 创建的客户端
|
||||
private HttpClient GetClient() => _httpFactory.CreateClient("VolPro");
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 7. Adapters.Owl/OwlAdapter.cs — 5 项
|
||||
|
||||
### OW1 [🟠] `GetDevicesAsync` 忽略 `page`/`size` 参数
|
||||
|
||||
**位置**: 第 68 行
|
||||
```csharp
|
||||
var url = $"/devices/channels?page={page}&size=1000";
|
||||
```
|
||||
硬编码 `size=1000`,无论调用者传的 `size` 值是多少。前端分页完全失效。
|
||||
|
||||
**修复**:
|
||||
```csharp
|
||||
var url = $"/devices/channels?page={page}&size={size}";
|
||||
```
|
||||
|
||||
### OW2 [🟡] `GetPlaybackUrlAsync` 手工拼 URL 脆弱
|
||||
|
||||
**位置**: 第 149-153 行
|
||||
```csharp
|
||||
return new StreamUrls
|
||||
{
|
||||
Hls = $"{baseUrl}/recordings/channels/{channelId}/index.m3u8?start_ms={startMs}&end_ms={endMs}&token={token}"
|
||||
};
|
||||
```
|
||||
直接拼接 URL,依赖 Owl 内部路径约定。Owl API 版本升级可能改变路径格式。
|
||||
|
||||
**修复** — 调 Owl API 获取回放地址(如果有对应接口)或至少将 base URL 路径前缀提取为常量:
|
||||
```csharp
|
||||
private const string PlaybackPathFormat = "/recordings/channels/{0}/index.m3u8?start_ms={1}&end_ms={2}&token={3}";
|
||||
var hls = string.Format(PlaybackPathFormat, channelId, startMs, endMs, token);
|
||||
```
|
||||
|
||||
### OW3 [🟡] `MapChannel` IsOnline 判断脆弱
|
||||
|
||||
**位置**: 第 116 行
|
||||
```csharp
|
||||
IsOnline = ch.IsOnline?.ToLower() == "true" || ch.IsOnline == "1",
|
||||
```
|
||||
Owl 的 `IsOnline` 字段在不同接口中返回类型不同:设备接口返回 `"0"/"1"`,通道接口可能返回 `true/false` 字符串。这种脆弱的兼容方式在 Owl 版本升级后可能失效。
|
||||
|
||||
**修复** — 统一为 bool 解析:
|
||||
```csharp
|
||||
private static bool ParseOwlOnline(string? val) =>
|
||||
bool.TryParse(val, out var b) ? b : val == "1";
|
||||
```
|
||||
|
||||
### OW4 [🟡] `MapDevice` + `MapChannel` 硬编码设备名
|
||||
|
||||
**位置**: 第 91, 113 行
|
||||
```csharp
|
||||
Category = "硬盘录像机", Group = "视频设备"
|
||||
Category = "摄像机", Group = "视频设备"
|
||||
```
|
||||
**修复** — 提取为常量:
|
||||
```csharp
|
||||
private const string DEVICE_CATEGORY_NVR = "硬盘录像机";
|
||||
private const string DEVICE_CATEGORY_CAMERA = "摄像机";
|
||||
private const string DEVICE_GROUP_VIDEO = "视频设备";
|
||||
```
|
||||
|
||||
### OW5 [⚪] 文件 308 行过长
|
||||
|
||||
**修复** — 拆分为:
|
||||
- `OwlAdapter.cs` — 类声明 + 构造函数 + Capabilities (40 行)
|
||||
- `OwlAdapter.FlatDevices.cs` — IHasFlatDevices 实现 (70 行)
|
||||
- `OwlAdapter.Streams.cs` — IHasStreams 实现 (80 行)
|
||||
- `OwlAdapter.Recordings.cs` — IHasRecordings (30 行)
|
||||
- `OwlAdapter.Alarms.cs` — IHasAlarms (60 行)
|
||||
|
||||
---
|
||||
|
||||
## 8. Adapters.Owl/OwlAuthHelper.cs — 0 项
|
||||
|
||||
代码质量良好的 RSA 加密认证实现。Token 缓存策略合理(2.5 天/3 天)。无可优化项。
|
||||
|
||||
---
|
||||
|
||||
## 9. Adapters.Owl/OwlModels.cs — 1 项
|
||||
|
||||
### OM1 [🟡] `OwlDeviceChannel` 字段注释缺失
|
||||
|
||||
**位置**: 第 17-38 行
|
||||
`Type`、`Did`、`Ptztype`、`App`、`StreamId` 等字段无注释说明其含义和取值范围。
|
||||
|
||||
**修复** — 添加 XML 注释:
|
||||
```csharp
|
||||
/// <summary>类型: "DEVICE"(NVR) | "CHANNEL"(摄像头)</summary>
|
||||
public string? Type { get; set; }
|
||||
/// <summary>设备 ID(通道记录指向其所属设备)</summary>
|
||||
public string? Did { get; set; }
|
||||
/// <summary>云台类型: 0=无, 1=方向, 2=预置位</summary>
|
||||
public int? Ptztype { get; set; }
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 10. Adapters.MC4/Mc4Adapter.cs — 4 项
|
||||
|
||||
### MC1 [🟠] `GetAlarmsAsync` DateTime.MinValue 仍发送
|
||||
|
||||
**位置**: 第 145-146 行
|
||||
```csharp
|
||||
From = from.ToString("yyyy-MM-dd HH:mm:ss"),
|
||||
To = to.ToString("yyyy-MM-dd HH:mm:ss"),
|
||||
```
|
||||
当 `from`/`to` 为 `DateTime.MinValue`(B8 路由默认值)时,发送 `"0001-01-01 00:00:00"` 给 MC4——MC4 将认为这是有效日期过滤。
|
||||
|
||||
**修复**:
|
||||
```csharp
|
||||
From = from == DateTime.MinValue ? "" : from.ToString("yyyy-MM-dd HH:mm:ss"),
|
||||
To = to == DateTime.MinValue ? "" : to.ToString("yyyy-MM-dd HH:mm:ss"),
|
||||
```
|
||||
|
||||
### MC2 [🟡] MC4 模型应独立文件
|
||||
|
||||
**位置**: 第 266-350 行
|
||||
`Mc4TreeNode`、`Mc4PointValue`、`Mc4AlarmQuery` 等 8 个类全部定义在 `Mc4Adapter.cs` 底部。
|
||||
|
||||
**修复** — 创建 `Mc4Models.cs`(参照 KMS/Owl 的模式),将 266 行之后的模型移出。
|
||||
|
||||
### MC3 [🟡] `GetMultiRealtimeValuesAsync` + `GetHisAlarmsAsync` 未暴露到 B 路由
|
||||
|
||||
**位置**: 第 211, 228 行
|
||||
这两个方法是 MC4.0 原生批量接口,但 Program.cs 中没有对应的 B 路由暴露它们。B4-batch 路由直接用 `IHasPoints.GetRealtimeValuesAsync` 逐设备调用。
|
||||
|
||||
**修复** — B4-batch 路由已检查 `Mc4Adapter` 类型并优先调用 `GetMultiRealtimeValuesAsync`(已实现)。确认编译通过即可。
|
||||
|
||||
### MC4 [🟡] `ConfirmAlarmAsync` / `EndAlarmAsync` 不检查响应
|
||||
|
||||
**位置**: 第 175-192 行
|
||||
```csharp
|
||||
await client.PostAsync("/api/central/alarm/confirm", ...); // 无 resp.EnsureSuccessStatusCode()
|
||||
```
|
||||
MC4 返回非 200 时静默失败,调用者认为确认成功。
|
||||
|
||||
**修复**:
|
||||
```csharp
|
||||
var resp = await client.PostAsync(...);
|
||||
resp.EnsureSuccessStatusCode();
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 11. Adapters.MC4/Mc4AuthHelper.cs — 1 项
|
||||
|
||||
### MA1 [🟡] `_needMd5` 只获取一次
|
||||
|
||||
**位置**: 第 46-57 行
|
||||
```csharp
|
||||
if (!_needMd5.HasValue) { /* 仅首次调用时查 conf/get */ }
|
||||
```
|
||||
如果 MC4 服务重启并更改了加密配置(encrypt true↔false),适配器不会重新检测。但实际场景中 MC4 不会在运行时切换加密模式,风险极低。
|
||||
|
||||
**修复** — 加一个 Token 刷新计数器,每 N 次重新获取 conf/get(N=10 即可)。
|
||||
|
||||
---
|
||||
|
||||
## 12. Adapters.Kms/KmsAdapter.cs — 1 项
|
||||
|
||||
### KM1 [⚪] `OpenerIds` 缩进不齐
|
||||
|
||||
**位置**: 第 296 行
|
||||
```csharp
|
||||
OpenerIds = parameters.TryGetValue(...)
|
||||
```
|
||||
前面的注释和代码使用 8 空格缩进,`OpenerIds` 使用 0 空格缩进。视觉效果不一致但不影响编译。
|
||||
|
||||
**修复** — 统一缩进为 8 空格。
|
||||
|
||||
---
|
||||
|
||||
## 13. Adapters.Kms/KmsAuthHelper.cs — 0 项
|
||||
|
||||
Token 缓存 25min(30-5)、Bearer 头、Invalidate 均实现正确。无可优化项。
|
||||
|
||||
---
|
||||
|
||||
## 14. Adapters.Kms/KmsModels.cs — 0 项
|
||||
|
||||
15 个 DTO 覆盖全部 KMS 接口,字段完整。标准接口 DTO 虽暂未使用,但注释说明 Phase 2 用途——保留合理。
|
||||
|
||||
---
|
||||
|
||||
## 15. Host/Program.cs — 6 项
|
||||
|
||||
### PR1 [🟠] `SyncAllDevicesAsync` 和 `FlattenTree` 在两个层次定义
|
||||
|
||||
**位置**: 第 152-184 行
|
||||
这两个函数在 `InitializeAllAsync` 回调的闭包作用域中定义为局部函数。如果未来需要在其他地方调用(如 A3 手动触发),将不可用。
|
||||
|
||||
**修复** — 提取为私有静态方法或迁移到 Core:
|
||||
```csharp
|
||||
// 新建 Core/Infrastructure/DeviceSyncHelper.cs
|
||||
public static class DeviceSyncHelper
|
||||
{
|
||||
public static async Task SyncAllAsync(AdapterRegistry reg, GatewayClientFactory factory, string nodeCode, string token) { ... }
|
||||
private static void FlattenTree(...) { ... }
|
||||
}
|
||||
```
|
||||
|
||||
### PR2 [🟡] Swagger 无 XML 注释
|
||||
|
||||
**位置**: 第 17-19 行
|
||||
```csharp
|
||||
builder.Services.AddSwaggerGen(); // 无 XML 注释选项
|
||||
```
|
||||
所有 Minimal API 端点没有 Swagger 描述,调用者必须参考外部文档。
|
||||
|
||||
**修复**:
|
||||
```csharp
|
||||
builder.Services.AddSwaggerGen(c =>
|
||||
{
|
||||
var xmlFile = $"{Assembly.GetExecutingAssembly().GetName().Name}.xml";
|
||||
c.IncludeXmlComments(Path.Combine(AppContext.BaseDirectory, xmlFile));
|
||||
});
|
||||
```
|
||||
并在 `.csproj` 中添加 `<GenerateDocumentationFile>true</GenerateDocumentationFile>`。
|
||||
|
||||
### PR3 [🟡] 路由注册代码可读性差
|
||||
|
||||
**位置**: 第 186-377 行
|
||||
19 条 B 路由全部内联在 `Program.cs` 中,每条路由 5-10 行,总计 ~200 行路由注册代码。
|
||||
|
||||
**修复** — 按模块拆分为扩展方法:
|
||||
```csharp
|
||||
// Host/Routes/HealthRoutes.cs
|
||||
public static class HealthRoutes
|
||||
{
|
||||
public static void MapHealthEndpoints(this WebApplication app, AdapterRegistry registry) { ... }
|
||||
}
|
||||
// Host/Routes/DeviceRoutes.cs — B2, B3, B3-sync
|
||||
// Host/Routes/StreamRoutes.cs — B6a, B6b, B7, snapshot
|
||||
// Host/Routes/RealtimeRoutes.cs — B4, B4-batch, B5
|
||||
// Host/Routes/AlarmRoutes.cs — B8, B9-confirm, B9-end
|
||||
// Host/Routes/ControlRoutes.cs — B10
|
||||
// Host/Routes/LogRoutes.cs — B11
|
||||
// Host/Routes/SyncRoutes.cs — B12, B13
|
||||
```
|
||||
Program.cs 改为:
|
||||
```csharp
|
||||
app.MapHealthEndpoints(registry);
|
||||
app.MapDeviceEndpoints(registry);
|
||||
// ...
|
||||
```
|
||||
|
||||
### PR4 [🟡] 配置缺少验证
|
||||
|
||||
**位置**: 第 53-83 行
|
||||
`app.Configuration.GetSection("Owl").Get<List<OwlConfig>>()` 可能在配置格式错误时返回 null,直接 `foreach` 正常运作但无错误提示。
|
||||
|
||||
**修复** — 加空值检查和警告:
|
||||
```csharp
|
||||
var owlList = app.Configuration.GetSection("Owl").Get<List<OwlConfig>>();
|
||||
if (owlList == null || !owlList.Any()) Console.WriteLine("[Gateway] WARNING: 未配置 Owl 适配器");
|
||||
```
|
||||
|
||||
### PR5 [🟡] `appsettings.json` 凭证明文
|
||||
|
||||
**位置**: appsettings.json
|
||||
```json
|
||||
"NodeToken": "changeme",
|
||||
"Password": "your_owl_password",
|
||||
"ClientSecret": "your_client_secret"
|
||||
```
|
||||
|
||||
**修复** — 生产环境注入:
|
||||
```json
|
||||
"NodeToken": null, // 生产环境由 SECMPS_GATEWAY_TOKEN 环境变量注入
|
||||
"Password": null, // 生产环境由 OWL_PASSWORD 注入
|
||||
"ClientSecret": null // 生产环境由 KMS_CLIENT_SECRET 注入
|
||||
```
|
||||
Program.cs:
|
||||
```csharp
|
||||
var nodeToken = Environment.GetEnvironmentVariable("SECMPS_GATEWAY_TOKEN") ?? gwCfg["NodeToken"];
|
||||
```
|
||||
|
||||
### PR6 [⚪] dotnet-tools.json 未使用
|
||||
|
||||
`gateway/src/IntegrationGateway.Host/dotnet-tools.json` 存在但可能无声明工具。确认是否被使用。
|
||||
|
||||
---
|
||||
|
||||
## 16. 跨切面问题 — 4 项
|
||||
|
||||
### X1 [🟡] 无统一日志抽象
|
||||
|
||||
所有适配器使用 `Console.Error.WriteLine` 直接写控制台。无日志级别、无结构化输出、无法对接日志收集系统。
|
||||
|
||||
**修复** — 注入 `ILogger<T>`:
|
||||
```csharp
|
||||
public class OwlAdapter : ...
|
||||
{
|
||||
private readonly ILogger<OwlAdapter> _logger;
|
||||
public OwlAdapter(..., ILogger<OwlAdapter> logger) { _logger = logger; }
|
||||
|
||||
public async Task<bool> HealthCheckAsync()
|
||||
{
|
||||
try { ... }
|
||||
catch (Exception ex) { _logger.LogWarning(ex, "[{Code}] HealthCheck 失败", AdapterCode); return false; }
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
### X2 [🟡] `GetAuthenticatedClientAsync` 每次 new HttpClient
|
||||
|
||||
三个适配器(Owl/KMS/MC4)的 AuthHelper 都在 `GetAuthenticatedClientAsync` 中 `new HttpClient { BaseAddress = new Uri(_baseUrl) }`。每个请求创建一个新的 HttpClient 实例,无法复用连接池。
|
||||
|
||||
**修复** — 传入 `IHttpClientFactory`:
|
||||
```csharp
|
||||
public class KmsAuthHelper
|
||||
{
|
||||
private readonly IHttpClientFactory _httpFactory;
|
||||
public async Task<HttpClient> GetAuthenticatedClientAsync()
|
||||
{
|
||||
var token = await GetTokenAsync();
|
||||
var client = _httpFactory.CreateClient(); // 连接池复用
|
||||
client.BaseAddress = new Uri(_baseUrl);
|
||||
client.DefaultRequestHeaders.Authorization = new AuthenticationHeaderValue("Bearer", token);
|
||||
return client;
|
||||
}
|
||||
}
|
||||
```
|
||||
**跨文件影响**: 三个 AuthHelper 构造函数全部需加 `IHttpClientFactory` 参数。三个 Adapter 构造函数也需传递。
|
||||
|
||||
### X3 [🟡] 无 request/response DTO 版本策略
|
||||
|
||||
当前所有接口模型无版本号字段。未来接口升级时无法区分新旧格式。
|
||||
|
||||
**修复** — 在 B 路由响应中加 `version` 字段:
|
||||
```csharp
|
||||
return Results.Ok(new { version = "1.0", items = result.Items, total = result.Total });
|
||||
```
|
||||
|
||||
### X4 [⚪] `csproj` 无 `<GenerateDocumentationFile>` 配置
|
||||
|
||||
三个适配器项目均无 XML 文档生成配置。
|
||||
|
||||
**修复** — 在各 `.csproj` 中添加 `<GenerateDocumentationFile>true</GenerateDocumentationFile>`。
|
||||
|
||||
---
|
||||
|
||||
## 统计
|
||||
|
||||
| 级别 | 数量 | 位置 |
|
||||
|:--:|:--:|------|
|
||||
| 🟠 严重 | 4 | RateLimiter(Owl.GetDevices size,MC4 MinValue,Program SyncAll) |
|
||||
| 🟡 改善 | 21 | 命名/拆分/日志/HttpClient/配置 |
|
||||
| ⚪ 低优 | 5 | 缩进/tools.json/注释 |
|
||||
|
||||
## 总评
|
||||
|
||||
网关架构设计优秀——适配器模式隔离清晰、限流/认证/错误处理一致、接口粒度过细(19 条 B 路由)而非不足。最大改进空间在于:HttpClient 连接池复用、日志抽象统一、Program.cs 路由拆分。
|
||||
@@ -1,6 +1,7 @@
|
||||
using IntegrationGateway.Core.Abstractions;
|
||||
using IntegrationGateway.Core.Infrastructure;
|
||||
using IntegrationGateway.Core.Models;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using System.Net.Http.Json;
|
||||
using System.Text;
|
||||
using System.Text.Json;
|
||||
@@ -22,6 +23,7 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
|
||||
private readonly HttpClient _http;
|
||||
private readonly KmsAuthHelper _auth;
|
||||
private readonly RateLimiter _limiter = new(5);
|
||||
private readonly ILogger<KmsAdapter> _logger;
|
||||
|
||||
/// <summary>适配器编码,格式 "KMS:{实例名}"</summary>
|
||||
public string AdapterCode { get; }
|
||||
@@ -38,11 +40,12 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
|
||||
/// <param name="baseUrl">KMS 服务地址</param>
|
||||
/// <param name="clientId">KMS 客户端 ID</param>
|
||||
/// <param name="clientSecret">KMS 客户端密钥</param>
|
||||
public KmsAdapter(string adapterCode, HttpClient http, string baseUrl, string clientId, string clientSecret)
|
||||
public KmsAdapter(string adapterCode, HttpClient http, string baseUrl, string clientId, string clientSecret, ILogger<KmsAdapter> logger = null!)
|
||||
{
|
||||
AdapterCode = adapterCode;
|
||||
_http = http;
|
||||
_auth = new KmsAuthHelper(http, baseUrl, clientId, clientSecret);
|
||||
_logger = logger;
|
||||
_auth = new KmsAuthHelper(http, baseUrl, clientId, clientSecret, logger);
|
||||
}
|
||||
|
||||
/// <summary>初始化适配器:获取 KMS Token</summary>
|
||||
@@ -80,6 +83,7 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
|
||||
new StringContent("{}", Encoding.UTF8, "application/json"));
|
||||
resp.EnsureSuccessStatusCode();
|
||||
var data = await resp.Content.ReadFromJsonAsync<KmsOpenerListResponse>()!;
|
||||
_logger.LogDebug("[{Code}] POST /prod-api/getOpenerList → {Count} lockers", AdapterCode, data.Rows?.Count ?? 0);
|
||||
|
||||
var devices = new List<StandardDevice>();
|
||||
foreach (var locker in data.Rows ?? new())
|
||||
@@ -147,6 +151,7 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
|
||||
new StringContent(body, Encoding.UTF8, "application/json"));
|
||||
resp.EnsureSuccessStatusCode();
|
||||
var data = await resp.Content.ReadFromJsonAsync<KmsWarningListResponse>()!;
|
||||
_logger.LogDebug("[{Code}] POST /prod-api/getWarningList → {Count} warnings", AdapterCode, data.Rows?.Count ?? 0);
|
||||
|
||||
var alarms = (data.Rows ?? new()).Select(w => new StandardAlarm
|
||||
{
|
||||
@@ -207,6 +212,7 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
|
||||
new StringContent(body, Encoding.UTF8, "application/json"));
|
||||
resp.EnsureSuccessStatusCode();
|
||||
var data = await resp.Content.ReadFromJsonAsync<KmsRecordListResponse>()!;
|
||||
_logger.LogDebug("[{Code}] POST /prod-api/getRecordList → {Count} records", AdapterCode, data.Rows?.Count ?? 0);
|
||||
return new PagedResult<KmsRecord> { Items = data.Rows ?? new(), Total = data.Total };
|
||||
}
|
||||
|
||||
@@ -234,6 +240,7 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
|
||||
new StringContent(body, Encoding.UTF8, "application/json"));
|
||||
resp.EnsureSuccessStatusCode();
|
||||
var data = await resp.Content.ReadFromJsonAsync<KmsPermissionListResponse>()!;
|
||||
_logger.LogDebug("[{Code}] POST /prod-api/getPermissionList → {Count} permissions", AdapterCode, data.Rows?.Count ?? 0);
|
||||
return new PagedResult<KmsPermission> { Items = data.Rows ?? new(), Total = data.Total };
|
||||
}
|
||||
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
using Microsoft.Extensions.Logging;
|
||||
using System.Net.Http.Json;
|
||||
using System.Text.Json;
|
||||
|
||||
@@ -14,6 +15,7 @@ public class KmsAuthHelper
|
||||
private readonly string _baseUrl;
|
||||
private readonly string _clientId;
|
||||
private readonly string _clientSecret;
|
||||
private readonly ILogger _logger;
|
||||
private string? _token;
|
||||
private DateTime _tokenExpiry = DateTime.MinValue;
|
||||
|
||||
@@ -24,12 +26,13 @@ public class KmsAuthHelper
|
||||
/// <param name="baseUrl">KMS 服务地址</param>
|
||||
/// <param name="clientId">KMS 客户端 ID</param>
|
||||
/// <param name="clientSecret">KMS 客户端密钥</param>
|
||||
public KmsAuthHelper(HttpClient http, string baseUrl, string clientId, string clientSecret)
|
||||
public KmsAuthHelper(HttpClient http, string baseUrl, string clientId, string clientSecret, ILogger logger = null!)
|
||||
{
|
||||
_http = http;
|
||||
_baseUrl = baseUrl.TrimEnd('/');
|
||||
_clientId = clientId;
|
||||
_clientSecret = clientSecret;
|
||||
_logger = logger;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
@@ -51,6 +54,7 @@ public class KmsAuthHelper
|
||||
|
||||
_token = result.Token;
|
||||
_tokenExpiry = DateTime.UtcNow.AddMinutes(25);
|
||||
_logger?.LogDebug("KMS token obtained, expires in 25min");
|
||||
return _token;
|
||||
}
|
||||
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
using IntegrationGateway.Core.Abstractions;
|
||||
using IntegrationGateway.Core.Infrastructure;
|
||||
using IntegrationGateway.Core.Models;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using System.Text;
|
||||
using System.Text.Json;
|
||||
|
||||
@@ -23,6 +24,7 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
|
||||
private readonly Mc4AuthHelper _auth;
|
||||
/// <summary>令牌桶限流器(2 QPS)</summary>
|
||||
private readonly RateLimiter _limiter = new(2);
|
||||
private readonly ILogger<Mc4Adapter> _logger;
|
||||
|
||||
/// <summary>适配器编码,格式 "MC4:实例名"</summary>
|
||||
public string AdapterCode { get; }
|
||||
@@ -38,11 +40,12 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
|
||||
/// <param name="adapterCode">适配器编码</param>
|
||||
/// <param name="http">HttpClient 实例</param>
|
||||
/// <param name="baseUrl">MC4.0 服务地址</param>
|
||||
public Mc4Adapter(string adapterCode, HttpClient http, string baseUrl, string account = "admin", string password = "admin")
|
||||
public Mc4Adapter(string adapterCode, HttpClient http, string baseUrl, string account = "admin", string password = "admin", ILogger<Mc4Adapter> logger = null!)
|
||||
{
|
||||
AdapterCode = adapterCode;
|
||||
_http = http;
|
||||
_auth = new Mc4AuthHelper(http, baseUrl, account, password);
|
||||
_logger = logger;
|
||||
_auth = new Mc4AuthHelper(http, baseUrl, account, password, logger);
|
||||
}
|
||||
|
||||
/// <summary>初始化适配器:获取 MC4.0 Token</summary>
|
||||
@@ -76,6 +79,7 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
|
||||
resp.EnsureSuccessStatusCode();
|
||||
var json = await resp.Content.ReadAsStringAsync();
|
||||
var tree = JsonSerializer.Deserialize<List<Mc4TreeNode>>(json)!;
|
||||
_logger.LogDebug("[{Code}] POST /api/central/object/tree → {Len} bytes, {Count} nodes", AdapterCode, json.Length, tree.Count);
|
||||
return tree.Select(MapNode).ToList();
|
||||
}
|
||||
|
||||
@@ -107,6 +111,7 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
|
||||
resp.EnsureSuccessStatusCode();
|
||||
var json = await resp.Content.ReadAsStringAsync();
|
||||
var values = JsonSerializer.Deserialize<List<Mc4PointValue>>(json)!;
|
||||
_logger.LogDebug("[{Code}] GET device/{Id} → {Len} bytes, {Count} points", AdapterCode, sourceDeviceId, json.Length, values.Count);
|
||||
return values.Select(v => new PointValue
|
||||
{
|
||||
SourceDeviceId = sourceDeviceId,
|
||||
@@ -142,8 +147,8 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
|
||||
var client = await _auth.GetAuthenticatedClientAsync();
|
||||
var body = JsonSerializer.Serialize(new Mc4AlarmQuery
|
||||
{
|
||||
From = from.ToString("yyyy-MM-dd HH:mm:ss"),
|
||||
To = to.ToString("yyyy-MM-dd HH:mm:ss"),
|
||||
From = from == DateTime.MinValue ? "" : from.ToString("yyyy-MM-dd HH:mm:ss"),
|
||||
To = to == DateTime.MinValue ? "" : to.ToString("yyyy-MM-dd HH:mm:ss"),
|
||||
Skip = (page - 1) * size,
|
||||
Limit = size,
|
||||
Sort = 1 // 按时间降序
|
||||
@@ -153,6 +158,7 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
|
||||
resp.EnsureSuccessStatusCode();
|
||||
var json = await resp.Content.ReadAsStringAsync();
|
||||
var result = JsonSerializer.Deserialize<Mc4AlarmQueryResult>(json)!;
|
||||
_logger.LogDebug("[{Code}] POST /api/central/alarm/query → {Len} bytes, {Count} alarms", AdapterCode, json.Length, result.List?.Count ?? 0);
|
||||
return new PagedResult<StandardAlarm>
|
||||
{
|
||||
Items = result.List?.Select(a => new StandardAlarm
|
||||
@@ -177,8 +183,9 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
|
||||
await _limiter.WaitAsync();
|
||||
var client = await _auth.GetAuthenticatedClientAsync();
|
||||
var body = JsonSerializer.Serialize(new { id = alarmId, option = new { } });
|
||||
await client.PostAsync("/api/central/alarm/confirm",
|
||||
var cresp = await client.PostAsync("/api/central/alarm/confirm",
|
||||
new StringContent(body, Encoding.UTF8, "application/json"));
|
||||
cresp.EnsureSuccessStatusCode();
|
||||
}
|
||||
|
||||
/// <summary>结束告警(同时写回 MC4.0)</summary>
|
||||
@@ -187,8 +194,9 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
|
||||
await _limiter.WaitAsync();
|
||||
var client = await _auth.GetAuthenticatedClientAsync();
|
||||
var body = JsonSerializer.Serialize(new { id = alarmId, option = new { } });
|
||||
await client.PostAsync("/api/central/alarm/end",
|
||||
var eresp = await client.PostAsync("/api/central/alarm/end",
|
||||
new StringContent(body, Encoding.UTF8, "application/json"));
|
||||
eresp.EnsureSuccessStatusCode();
|
||||
}
|
||||
|
||||
/// <summary>MC4.0 告警等级数字 → 中文映射</summary>
|
||||
@@ -231,8 +239,8 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
|
||||
var client = await _auth.GetAuthenticatedClientAsync();
|
||||
var body = JsonSerializer.Serialize(new Mc4HisAlarmQuery
|
||||
{
|
||||
From = from.ToString("yyyy-MM-dd HH:mm:ss"),
|
||||
To = to.ToString("yyyy-MM-dd HH:mm:ss"),
|
||||
From = from == DateTime.MinValue ? "" : from.ToString("yyyy-MM-dd HH:mm:ss"),
|
||||
To = to == DateTime.MinValue ? "" : to.ToString("yyyy-MM-dd HH:mm:ss"),
|
||||
Skip = (page - 1) * size,
|
||||
Limit = size,
|
||||
Sort = 1
|
||||
@@ -242,6 +250,7 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
|
||||
resp.EnsureSuccessStatusCode();
|
||||
var json = await resp.Content.ReadAsStringAsync();
|
||||
var result = JsonSerializer.Deserialize<Mc4AlarmQueryResult>(json)!;
|
||||
_logger.LogDebug("[{Code}] POST /api/central/alarm/query → {Len} bytes, {Count} alarms", AdapterCode, json.Length, result.List?.Count ?? 0);
|
||||
return new PagedResult<StandardAlarm>
|
||||
{
|
||||
Items = (result.List ?? new()).Select(MapAlarmItem).ToList(),
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
using Microsoft.Extensions.Logging;
|
||||
using System.Security.Cryptography;
|
||||
using System.Text;
|
||||
using System.Text.Json;
|
||||
@@ -20,16 +21,18 @@ public class Mc4AuthHelper
|
||||
private readonly string _baseUrl;
|
||||
private readonly string _account;
|
||||
private readonly string _password;
|
||||
private readonly ILogger _logger;
|
||||
private string? _token;
|
||||
private DateTime _tokenExpiry = DateTime.MinValue;
|
||||
private bool? _needMd5;
|
||||
|
||||
public Mc4AuthHelper(HttpClient http, string baseUrl, string account = "admin", string password = "admin")
|
||||
public Mc4AuthHelper(HttpClient http, string baseUrl, string account = "admin", string password = "admin", ILogger logger = null!)
|
||||
{
|
||||
_http = http;
|
||||
_baseUrl = baseUrl.TrimEnd('/');
|
||||
_account = account;
|
||||
_password = password;
|
||||
_logger = logger;
|
||||
}
|
||||
|
||||
public async Task<string> GetTokenAsync()
|
||||
@@ -67,6 +70,7 @@ public class Mc4AuthHelper
|
||||
throw new Exception("MC4 登录失败: Token 为空");
|
||||
_token = result.Token;
|
||||
_tokenExpiry = DateTime.UtcNow.AddHours(7);
|
||||
_logger?.LogDebug("MC4 token obtained, expires in 7h");
|
||||
return _token;
|
||||
}
|
||||
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
using IntegrationGateway.Core.Abstractions;
|
||||
using IntegrationGateway.Core.Infrastructure;
|
||||
using IntegrationGateway.Core.Models;
|
||||
using Microsoft.Extensions.Logging;
|
||||
using System.Text.Json;
|
||||
using System.Net.Http.Json;
|
||||
|
||||
@@ -24,6 +25,7 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
|
||||
private readonly HttpClient _http;
|
||||
private readonly OwlAuthHelper _auth;
|
||||
private readonly RateLimiter _limiter = new(5);
|
||||
private readonly ILogger<OwlAdapter> _logger;
|
||||
|
||||
public string AdapterCode { get; }
|
||||
public string DisplayName => $"Owl ({AdapterCode})";
|
||||
@@ -33,11 +35,12 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
|
||||
HasRecordings = true, AcceptsMetadataPush = true, HasAlarms = true
|
||||
};
|
||||
|
||||
public OwlAdapter(string adapterCode, HttpClient http, string baseUrl, string username, string password)
|
||||
public OwlAdapter(string adapterCode, HttpClient http, string baseUrl, string username, string password, ILogger<OwlAdapter> logger)
|
||||
{
|
||||
AdapterCode = adapterCode;
|
||||
_http = http;
|
||||
_auth = new OwlAuthHelper(http, baseUrl, username, password);
|
||||
_logger = logger;
|
||||
_auth = new OwlAuthHelper(http, baseUrl, username, password, logger);
|
||||
}
|
||||
|
||||
public async Task InitializeAsync() => await _auth.GetTokenAsync();
|
||||
@@ -65,9 +68,10 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
|
||||
{
|
||||
await _limiter.WaitAsync();
|
||||
var client = await _auth.GetAuthenticatedClientAsync();
|
||||
var url = $"/devices/channels?page={page}&size=1000";
|
||||
var url = $"/devices/channels?page={page}&size={size}";
|
||||
if (!string.IsNullOrEmpty(keyword)) url += $"&key={Uri.EscapeDataString(keyword)}";
|
||||
var json = await client.GetStringAsync(url);
|
||||
_logger.LogDebug("[{Code}] GET {Url} → {Len} bytes", AdapterCode, url, json.Length);
|
||||
var result = JsonSerializer.Deserialize<OwlPagedResult<OwlDeviceChannel>>(json)!;
|
||||
|
||||
var devices = new List<StandardDevice>();
|
||||
@@ -135,6 +139,7 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
|
||||
var resp = await client.PostAsync($"/channels/{channelId}/play", null);
|
||||
resp.EnsureSuccessStatusCode();
|
||||
var json = await resp.Content.ReadAsStringAsync();
|
||||
_logger.LogDebug("[{Code}] POST /channels/{Id}/play → {Len} bytes", AdapterCode, channelId, json.Length);
|
||||
var play = JsonSerializer.Deserialize<OwlPlayResponse>(json)!;
|
||||
return MapStreamUrls(play);
|
||||
}
|
||||
|
||||
@@ -3,6 +3,8 @@ using System.Text;
|
||||
using System.Text.Json;
|
||||
using System.Net.Http.Json;
|
||||
|
||||
using Microsoft.Extensions.Logging;
|
||||
|
||||
namespace IntegrationGateway.Adapters.Owl;
|
||||
|
||||
/// <summary>
|
||||
@@ -21,6 +23,7 @@ public class OwlAuthHelper
|
||||
private readonly string _baseUrl;
|
||||
private readonly string _username;
|
||||
private readonly string _password;
|
||||
private readonly ILogger _logger;
|
||||
/// <summary>缓存的 JWT Token</summary>
|
||||
private string? _token;
|
||||
/// <summary>Token 过期时间(UTC)</summary>
|
||||
@@ -31,10 +34,11 @@ public class OwlAuthHelper
|
||||
/// <param name="baseUrl">Owl 服务地址,如 http://localhost:15123</param>
|
||||
/// <param name="username">Owl 登录用户名</param>
|
||||
/// <param name="password">Owl 登录密码</param>
|
||||
public OwlAuthHelper(HttpClient http, string baseUrl, string username, string password)
|
||||
public OwlAuthHelper(HttpClient http, string baseUrl, string username, string password, ILogger logger)
|
||||
{
|
||||
_http = http; _baseUrl = baseUrl.TrimEnd('/');
|
||||
_username = username; _password = password;
|
||||
_logger = logger;
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
@@ -62,7 +66,8 @@ public class OwlAuthHelper
|
||||
resp.EnsureSuccessStatusCode();
|
||||
var loginResult = await resp.Content.ReadFromJsonAsync<LoginResponse>();
|
||||
_token = loginResult!.Token;
|
||||
_tokenExpiry = DateTime.UtcNow.AddDays(2.5); // 保守设置,Owl 默认 3 天
|
||||
_tokenExpiry = DateTime.UtcNow.AddDays(2.5);
|
||||
_logger.LogDebug("Owl token obtained, expires in 2.5 days"); // 保守设置,Owl 默认 3 天
|
||||
return _token;
|
||||
}
|
||||
|
||||
|
||||
@@ -46,7 +46,7 @@ public class AdapterRegistry
|
||||
public IGatewayAdapter? FindByCode(string adapterCode)
|
||||
=> _adapters.FirstOrDefault(a => a.AdapterCode == adapterCode);
|
||||
|
||||
/// <summary>获取所有在线适配器</summary>
|
||||
public IReadOnlyList<IGatewayAdapter> GetOnlineAdapters()
|
||||
/// <summary>获取所有已注册适配器(含在线和离线)</summary>
|
||||
public IReadOnlyList<IGatewayAdapter> GetAllAdapters()
|
||||
=> _adapters.AsReadOnly();
|
||||
}
|
||||
|
||||
@@ -4,13 +4,12 @@ namespace IntegrationGateway.Core.Infrastructure;
|
||||
/// 令牌桶限流器。控制对第三方子系统的请求频率,防止超出 API 配额。
|
||||
/// 每个适配器实例持有独立的限流器。
|
||||
///
|
||||
/// 算法:启动时桶内有 tokensPerSecond 个令牌,每次请求消耗一个令牌,
|
||||
/// 令牌按 (1000/tokensPerSecond) 毫秒的速率补充。
|
||||
/// 使用单个后台 PeriodicTimer 持续补充令牌,避免每次 WaitAsync 创建新 Task。
|
||||
/// </summary>
|
||||
public class RateLimiter
|
||||
{
|
||||
private readonly SemaphoreSlim _semaphore;
|
||||
private readonly int _intervalMs;
|
||||
private readonly int _maxTokens;
|
||||
|
||||
/// <summary>
|
||||
/// 创建限流器
|
||||
@@ -18,22 +17,23 @@ public class RateLimiter
|
||||
/// <param name="tokensPerSecond">每秒允许的请求数(QPS)</param>
|
||||
public RateLimiter(int tokensPerSecond)
|
||||
{
|
||||
_maxTokens = tokensPerSecond;
|
||||
_semaphore = new SemaphoreSlim(tokensPerSecond, tokensPerSecond);
|
||||
_intervalMs = 1000 / tokensPerSecond;
|
||||
var interval = 1000 / tokensPerSecond;
|
||||
_ = Task.Run(async () =>
|
||||
{
|
||||
using var timer = new PeriodicTimer(TimeSpan.FromMilliseconds(interval));
|
||||
while (await timer.WaitForNextTickAsync())
|
||||
{
|
||||
try { if (_semaphore.CurrentCount < _maxTokens) _semaphore.Release(); }
|
||||
catch (ObjectDisposedException) { break; }
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/// <summary>
|
||||
/// 等待获取一个令牌。如果当前没有可用令牌,阻塞直到有令牌被释放。
|
||||
/// </summary>
|
||||
/// <param name="ct">取消令牌</param>
|
||||
public async Task WaitAsync(CancellationToken ct = default)
|
||||
{
|
||||
await _semaphore.WaitAsync(ct);
|
||||
// 在后台任务中延迟补充令牌
|
||||
_ = Task.Run(async () =>
|
||||
{
|
||||
await Task.Delay(_intervalMs, ct);
|
||||
try { _semaphore.Release(); } catch { }
|
||||
}, ct);
|
||||
}
|
||||
=> await _semaphore.WaitAsync(ct);
|
||||
}
|
||||
|
||||
@@ -15,12 +15,14 @@ public class AdapterCapabilities
|
||||
/// <summary>是否支持视频取流</summary>
|
||||
public bool HasStreams { get; set; }
|
||||
/// <summary>是否支持云台控制(PTZ)</summary>
|
||||
[Obsolete("Use HasStreams for PTZ capability check")]
|
||||
public bool HasPtz { get; set; }
|
||||
/// <summary>是否支持录像回放</summary>
|
||||
public bool HasRecordings { get; set; }
|
||||
/// <summary>是否支持告警查询与处理</summary>
|
||||
public bool HasAlarms { get; set; }
|
||||
/// <summary>是否接受反向控制(点位写值)</summary>
|
||||
[Obsolete("Check IAcceptsControl interface directly")]
|
||||
public bool AcceptsControl { get; set; }
|
||||
/// <summary>是否接受元数据回写(如设备改名)</summary>
|
||||
public bool AcceptsMetadataPush { get; set; }
|
||||
|
||||
@@ -8,7 +8,7 @@ namespace IntegrationGateway.Core.Models;
|
||||
public class StandardDevice
|
||||
{
|
||||
/// <summary>Vol.Pro 侧主键(同步后由 Vol.Pro 回填)</summary>
|
||||
public int DeviceId { get; set; }
|
||||
public int? DeviceId { get; set; }
|
||||
/// <summary>来源适配器标识,格式 "类型:实例",如 "Owl:main"</summary>
|
||||
public string AdapterCode { get; set; } = "";
|
||||
/// <summary>子系统原始设备 ID(GB28181 编码 / MC4 sid)</summary>
|
||||
|
||||
@@ -56,7 +56,8 @@ foreach (var o in owlList)
|
||||
var code = $"Owl:{o.InstanceName ?? "default"}";
|
||||
var a = new IntegrationGateway.Adapters.Owl.OwlAdapter(code,
|
||||
app.Services.GetRequiredService<IHttpClientFactory>().CreateClient("VolPro"),
|
||||
o.BaseUrl, o.Username, o.Password);
|
||||
o.BaseUrl, o.Username, o.Password,
|
||||
app.Services.GetRequiredService<ILogger<IntegrationGateway.Adapters.Owl.OwlAdapter>>());
|
||||
registry.Register(a);
|
||||
}
|
||||
|
||||
@@ -67,7 +68,8 @@ foreach (var k in kmsList)
|
||||
var code = $"KMS:{k.InstanceName ?? "default"}";
|
||||
var a = new IntegrationGateway.Adapters.Kms.KmsAdapter(code,
|
||||
app.Services.GetRequiredService<IHttpClientFactory>().CreateClient("VolPro"),
|
||||
k.BaseUrl, k.ClientId, k.ClientSecret);
|
||||
k.BaseUrl, k.ClientId, k.ClientSecret,
|
||||
app.Services.GetRequiredService<ILogger<IntegrationGateway.Adapters.Kms.KmsAdapter>>());
|
||||
registry.Register(a);
|
||||
}
|
||||
|
||||
@@ -78,7 +80,8 @@ foreach (var m in mc4List)
|
||||
var code = $"MC4:{m.InstanceName ?? "default"}";
|
||||
var a = new IntegrationGateway.Adapters.MC4.Mc4Adapter(code,
|
||||
app.Services.GetRequiredService<IHttpClientFactory>().CreateClient("VolPro"),
|
||||
m.BaseUrl, m.Username, m.Password);
|
||||
m.BaseUrl, m.Username, m.Password,
|
||||
app.Services.GetRequiredService<ILogger<IntegrationGateway.Adapters.MC4.Mc4Adapter>>());
|
||||
registry.Register(a);
|
||||
}
|
||||
|
||||
|
||||
@@ -2,7 +2,8 @@
|
||||
"Logging": {
|
||||
"LogLevel": {
|
||||
"Default": "Information",
|
||||
"Microsoft.AspNetCore": "Warning"
|
||||
"Microsoft.AspNetCore": "Warning",
|
||||
"IntegrationGateway": "Debug"
|
||||
}
|
||||
},
|
||||
"Owl": [
|
||||
|
||||
Reference in New Issue
Block a user