diff --git a/doc/代码审核/网关代码审核20260604.md b/doc/代码审核/网关代码审核20260604.md new file mode 100644 index 0000000..fc084a2 --- /dev/null +++ b/doc/代码审核/网关代码审核20260604.md @@ -0,0 +1,535 @@ +# 网关项目深度代码审核报告 + +> 日期: 2026-06-04 | 项目: gateway/ | 扫描: 38 源文件, ~2800 行 | 问题: 30 项 + +--- + +## 1. Core/Infrastructure/AdapterRegistry.cs — 2 项 + +### AR1 [🟡] `GetOnlineAdapters()` 方法名误导 + +**位置**: 第 49-50 行 +```csharp +public IReadOnlyList GetOnlineAdapters() + => _adapters.AsReadOnly(); +``` +返回所有已注册适配器,**不做任何在线/离线判断**。调用者会误以为返回的是在线适配器。 + +**修复**: +```csharp +public IReadOnlyList GetAllAdapters() + => _adapters.AsReadOnly(); +``` +同时 grep 全项目无任何代码调用此方法,可直接删除(或保留为兼容性方法并标注 `[Obsolete]`)。 + +### AR2 [🟡] `FindByCode` O(n) 查找无缓存 + +每次 B 路由请求都执行 `_adapters.FirstOrDefault(...)`,适配器数量少时影响可忽略(<10 个),但如果未来扩展到 50+ 适配器会有性能问题。 + +**修复** — 加字典缓存: +```csharp +private readonly Dictionary _byCode = new(); +public void Register(IGatewayAdapter adapter) +{ + _adapters.Add(adapter); + _byCode[adapter.AdapterCode] = adapter; +} +public T? FindByCode(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` 不再需要此字段。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> 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> GetRecordingsAsync(RecordingQuery query); +``` +**跨文件影响**: `OwlAdapter.cs` GetRecordingsAsync 实现 + `Program.cs` B 路由。 + +--- + +## 6. Core/Infrastructure/GatewayClientFactory.cs — 2 项 + +### GF1 [🟡] `JsonDocument?` 返回类型不透明 + +**位置**: 第 33-42 行 +```csharp +public async Task RegisterAsync(GatewayRegisterRequest req) { ... } +public async Task SyncDevicesAsync(...) { ... } +``` +返回 `JsonDocument?` 使调用者必须知道 VolPro 响应的 JSON 结构。 + +**修复** — 定义响应 DTO: +```csharp +public class RegisterResponse { public int NodeId { get; set; } public List Devices { get; set; } = new(); } +public Task 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 +/// 类型: "DEVICE"(NVR) | "CHANNEL"(摄像头) +public string? Type { get; set; } +/// 设备 ID(通道记录指向其所属设备) +public string? Did { get; set; } +/// 云台类型: 0=无, 1=方向, 2=预置位 +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` 中添加 `true`。 + +### 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>()` 可能在配置格式错误时返回 null,直接 `foreach` 正常运作但无错误提示。 + +**修复** — 加空值检查和警告: +```csharp +var owlList = app.Configuration.GetSection("Owl").Get>(); +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`: +```csharp +public class OwlAdapter : ... +{ + private readonly ILogger _logger; + public OwlAdapter(..., ILogger logger) { _logger = logger; } + + public async Task 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 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` 无 `` 配置 + +三个适配器项目均无 XML 文档生成配置。 + +**修复** — 在各 `.csproj` 中添加 `true`。 + +--- + +## 统计 + +| 级别 | 数量 | 位置 | +|:--:|:--:|------| +| 🟠 严重 | 4 | RateLimiter(Owl.GetDevices size,MC4 MinValue,Program SyncAll) | +| 🟡 改善 | 21 | 命名/拆分/日志/HttpClient/配置 | +| ⚪ 低优 | 5 | 缩进/tools.json/注释 | + +## 总评 + +网关架构设计优秀——适配器模式隔离清晰、限流/认证/错误处理一致、接口粒度过细(19 条 B 路由)而非不足。最大改进空间在于:HttpClient 连接池复用、日志抽象统一、Program.cs 路由拆分。 diff --git a/gateway/src/IntegrationGateway.Adapters.Kms/KmsAdapter.cs b/gateway/src/IntegrationGateway.Adapters.Kms/KmsAdapter.cs index be5e66f..f4c31f5 100644 --- a/gateway/src/IntegrationGateway.Adapters.Kms/KmsAdapter.cs +++ b/gateway/src/IntegrationGateway.Adapters.Kms/KmsAdapter.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 _logger; /// 适配器编码,格式 "KMS:{实例名}" public string AdapterCode { get; } @@ -38,11 +40,12 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi /// KMS 服务地址 /// KMS 客户端 ID /// KMS 客户端密钥 - 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 logger = null!) { AdapterCode = adapterCode; _http = http; - _auth = new KmsAuthHelper(http, baseUrl, clientId, clientSecret); + _logger = logger; + _auth = new KmsAuthHelper(http, baseUrl, clientId, clientSecret, logger); } /// 初始化适配器:获取 KMS Token @@ -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()!; + _logger.LogDebug("[{Code}] POST /prod-api/getOpenerList → {Count} lockers", AdapterCode, data.Rows?.Count ?? 0); var devices = new List(); 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()!; + _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()!; + _logger.LogDebug("[{Code}] POST /prod-api/getRecordList → {Count} records", AdapterCode, data.Rows?.Count ?? 0); return new PagedResult { 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()!; + _logger.LogDebug("[{Code}] POST /prod-api/getPermissionList → {Count} permissions", AdapterCode, data.Rows?.Count ?? 0); return new PagedResult { Items = data.Rows ?? new(), Total = data.Total }; } diff --git a/gateway/src/IntegrationGateway.Adapters.Kms/KmsAuthHelper.cs b/gateway/src/IntegrationGateway.Adapters.Kms/KmsAuthHelper.cs index 8f39c9c..adbda69 100644 --- a/gateway/src/IntegrationGateway.Adapters.Kms/KmsAuthHelper.cs +++ b/gateway/src/IntegrationGateway.Adapters.Kms/KmsAuthHelper.cs @@ -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 /// KMS 服务地址 /// KMS 客户端 ID /// KMS 客户端密钥 - 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; } /// @@ -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; } diff --git a/gateway/src/IntegrationGateway.Adapters.MC4/Mc4Adapter.cs b/gateway/src/IntegrationGateway.Adapters.MC4/Mc4Adapter.cs index 2eedb64..a27eda7 100644 --- a/gateway/src/IntegrationGateway.Adapters.MC4/Mc4Adapter.cs +++ b/gateway/src/IntegrationGateway.Adapters.MC4/Mc4Adapter.cs @@ -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; /// 令牌桶限流器(2 QPS) private readonly RateLimiter _limiter = new(2); + private readonly ILogger _logger; /// 适配器编码,格式 "MC4:实例名" public string AdapterCode { get; } @@ -38,11 +40,12 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms /// 适配器编码 /// HttpClient 实例 /// MC4.0 服务地址 - 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 logger = null!) { AdapterCode = adapterCode; _http = http; - _auth = new Mc4AuthHelper(http, baseUrl, account, password); + _logger = logger; + _auth = new Mc4AuthHelper(http, baseUrl, account, password, logger); } /// 初始化适配器:获取 MC4.0 Token @@ -76,6 +79,7 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms resp.EnsureSuccessStatusCode(); var json = await resp.Content.ReadAsStringAsync(); var tree = JsonSerializer.Deserialize>(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>(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(json)!; + _logger.LogDebug("[{Code}] POST /api/central/alarm/query → {Len} bytes, {Count} alarms", AdapterCode, json.Length, result.List?.Count ?? 0); return new PagedResult { 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(); } /// 结束告警(同时写回 MC4.0) @@ -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(); } /// MC4.0 告警等级数字 → 中文映射 @@ -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(json)!; + _logger.LogDebug("[{Code}] POST /api/central/alarm/query → {Len} bytes, {Count} alarms", AdapterCode, json.Length, result.List?.Count ?? 0); return new PagedResult { Items = (result.List ?? new()).Select(MapAlarmItem).ToList(), diff --git a/gateway/src/IntegrationGateway.Adapters.MC4/Mc4AuthHelper.cs b/gateway/src/IntegrationGateway.Adapters.MC4/Mc4AuthHelper.cs index c527e2a..f095bcf 100644 --- a/gateway/src/IntegrationGateway.Adapters.MC4/Mc4AuthHelper.cs +++ b/gateway/src/IntegrationGateway.Adapters.MC4/Mc4AuthHelper.cs @@ -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 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; } diff --git a/gateway/src/IntegrationGateway.Adapters.Owl/OwlAdapter.cs b/gateway/src/IntegrationGateway.Adapters.Owl/OwlAdapter.cs index 68e747e..c90fbd9 100644 --- a/gateway/src/IntegrationGateway.Adapters.Owl/OwlAdapter.cs +++ b/gateway/src/IntegrationGateway.Adapters.Owl/OwlAdapter.cs @@ -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 _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 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>(json)!; var devices = new List(); @@ -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(json)!; return MapStreamUrls(play); } diff --git a/gateway/src/IntegrationGateway.Adapters.Owl/OwlAuthHelper.cs b/gateway/src/IntegrationGateway.Adapters.Owl/OwlAuthHelper.cs index 37aa8be..b37c2a0 100644 --- a/gateway/src/IntegrationGateway.Adapters.Owl/OwlAuthHelper.cs +++ b/gateway/src/IntegrationGateway.Adapters.Owl/OwlAuthHelper.cs @@ -3,6 +3,8 @@ using System.Text; using System.Text.Json; using System.Net.Http.Json; +using Microsoft.Extensions.Logging; + namespace IntegrationGateway.Adapters.Owl; /// @@ -21,6 +23,7 @@ public class OwlAuthHelper private readonly string _baseUrl; private readonly string _username; private readonly string _password; + private readonly ILogger _logger; /// 缓存的 JWT Token private string? _token; /// Token 过期时间(UTC) @@ -31,10 +34,11 @@ public class OwlAuthHelper /// Owl 服务地址,如 http://localhost:15123 /// Owl 登录用户名 /// Owl 登录密码 - 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; } /// @@ -62,7 +66,8 @@ public class OwlAuthHelper resp.EnsureSuccessStatusCode(); var loginResult = await resp.Content.ReadFromJsonAsync(); _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; } diff --git a/gateway/src/IntegrationGateway.Core/Infrastructure/AdapterRegistry.cs b/gateway/src/IntegrationGateway.Core/Infrastructure/AdapterRegistry.cs index f70ec92..ea25a08 100644 --- a/gateway/src/IntegrationGateway.Core/Infrastructure/AdapterRegistry.cs +++ b/gateway/src/IntegrationGateway.Core/Infrastructure/AdapterRegistry.cs @@ -46,7 +46,7 @@ public class AdapterRegistry public IGatewayAdapter? FindByCode(string adapterCode) => _adapters.FirstOrDefault(a => a.AdapterCode == adapterCode); - /// 获取所有在线适配器 - public IReadOnlyList GetOnlineAdapters() + /// 获取所有已注册适配器(含在线和离线) + public IReadOnlyList GetAllAdapters() => _adapters.AsReadOnly(); } diff --git a/gateway/src/IntegrationGateway.Core/Infrastructure/RateLimiter.cs b/gateway/src/IntegrationGateway.Core/Infrastructure/RateLimiter.cs index bdb3d54..72649a0 100644 --- a/gateway/src/IntegrationGateway.Core/Infrastructure/RateLimiter.cs +++ b/gateway/src/IntegrationGateway.Core/Infrastructure/RateLimiter.cs @@ -4,13 +4,12 @@ namespace IntegrationGateway.Core.Infrastructure; /// 令牌桶限流器。控制对第三方子系统的请求频率,防止超出 API 配额。 /// 每个适配器实例持有独立的限流器。 /// -/// 算法:启动时桶内有 tokensPerSecond 个令牌,每次请求消耗一个令牌, -/// 令牌按 (1000/tokensPerSecond) 毫秒的速率补充。 +/// 使用单个后台 PeriodicTimer 持续补充令牌,避免每次 WaitAsync 创建新 Task。 /// public class RateLimiter { private readonly SemaphoreSlim _semaphore; - private readonly int _intervalMs; + private readonly int _maxTokens; /// /// 创建限流器 @@ -18,22 +17,23 @@ public class RateLimiter /// 每秒允许的请求数(QPS) 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; } + } + }); } /// /// 等待获取一个令牌。如果当前没有可用令牌,阻塞直到有令牌被释放。 /// - /// 取消令牌 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); } diff --git a/gateway/src/IntegrationGateway.Core/Models/AdapterCapabilities.cs b/gateway/src/IntegrationGateway.Core/Models/AdapterCapabilities.cs index a3acecb..42180c1 100644 --- a/gateway/src/IntegrationGateway.Core/Models/AdapterCapabilities.cs +++ b/gateway/src/IntegrationGateway.Core/Models/AdapterCapabilities.cs @@ -15,12 +15,14 @@ public class AdapterCapabilities /// 是否支持视频取流 public bool HasStreams { get; set; } /// 是否支持云台控制(PTZ) + [Obsolete("Use HasStreams for PTZ capability check")] public bool HasPtz { get; set; } /// 是否支持录像回放 public bool HasRecordings { get; set; } /// 是否支持告警查询与处理 public bool HasAlarms { get; set; } /// 是否接受反向控制(点位写值) + [Obsolete("Check IAcceptsControl interface directly")] public bool AcceptsControl { get; set; } /// 是否接受元数据回写(如设备改名) public bool AcceptsMetadataPush { get; set; } diff --git a/gateway/src/IntegrationGateway.Core/Models/StandardDevice.cs b/gateway/src/IntegrationGateway.Core/Models/StandardDevice.cs index ca2f43a..444d290 100644 --- a/gateway/src/IntegrationGateway.Core/Models/StandardDevice.cs +++ b/gateway/src/IntegrationGateway.Core/Models/StandardDevice.cs @@ -8,7 +8,7 @@ namespace IntegrationGateway.Core.Models; public class StandardDevice { /// Vol.Pro 侧主键(同步后由 Vol.Pro 回填) - public int DeviceId { get; set; } + public int? DeviceId { get; set; } /// 来源适配器标识,格式 "类型:实例",如 "Owl:main" public string AdapterCode { get; set; } = ""; /// 子系统原始设备 ID(GB28181 编码 / MC4 sid) diff --git a/gateway/src/IntegrationGateway.Host/Program.cs b/gateway/src/IntegrationGateway.Host/Program.cs index 3571917..7f0bf79 100644 --- a/gateway/src/IntegrationGateway.Host/Program.cs +++ b/gateway/src/IntegrationGateway.Host/Program.cs @@ -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().CreateClient("VolPro"), - o.BaseUrl, o.Username, o.Password); + o.BaseUrl, o.Username, o.Password, + app.Services.GetRequiredService>()); 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().CreateClient("VolPro"), - k.BaseUrl, k.ClientId, k.ClientSecret); + k.BaseUrl, k.ClientId, k.ClientSecret, + app.Services.GetRequiredService>()); 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().CreateClient("VolPro"), - m.BaseUrl, m.Username, m.Password); + m.BaseUrl, m.Username, m.Password, + app.Services.GetRequiredService>()); registry.Register(a); } diff --git a/gateway/src/IntegrationGateway.Host/appsettings.json b/gateway/src/IntegrationGateway.Host/appsettings.json index 6dc3cf2..9da4f30 100644 --- a/gateway/src/IntegrationGateway.Host/appsettings.json +++ b/gateway/src/IntegrationGateway.Host/appsettings.json @@ -2,7 +2,8 @@ "Logging": { "LogLevel": { "Default": "Information", - "Microsoft.AspNetCore": "Warning" + "Microsoft.AspNetCore": "Warning", + "IntegrationGateway": "Debug" } }, "Owl": [