Compare commits

..

2 Commits

22 changed files with 1301 additions and 291 deletions
@@ -87,6 +87,41 @@ namespace VolPro.Entity.DomainModels
[Editable(true)] [Editable(true)]
public int? RuleID { get; set; } public int? RuleID { get; set; }
/// <summary>
///数值型恢复阈值(如>28℃触发,≤26℃恢复)
/// </summary>
[Display(Name ="数值型恢复阈值(如>28℃触发,≤26℃恢复)")]
[DisplayFormat(DataFormatString="18,2")]
[Column(TypeName="decimal")]
[Editable(true)]
public decimal? RecoveryThreshold_Numeric { get; set; }
/// <summary>
///开关型恢复阈值(如开触发,关恢复)
/// </summary>
[Display(Name ="开关型恢复阈值(如开触发,关恢复)")]
[MaxLength(50)]
[Column(TypeName="nvarchar(50)")]
[Editable(true)]
public string RecoveryThreshold_Switch { get; set; }
/// <summary>
///该条件上次触发的时间
/// </summary>
[Display(Name ="该条件上次触发的时间")]
[Column(TypeName="datetime")]
[Editable(true)]
public DateTime? LastTriggered { get; set; }
/// <summary>
///该条件上次触发时的实际值
/// </summary>
[Display(Name ="该条件上次触发时的实际值")]
[DisplayFormat(DataFormatString="18,2")]
[Column(TypeName="decimal")]
[Editable(true)]
public decimal? LastTriggerValue { get; set; }
} }
} }
@@ -0,0 +1,47 @@
using Microsoft.AspNetCore.Authorization;
using Microsoft.AspNetCore.Mvc;
using Microsoft.Extensions.Options;
namespace VolPro.WebApi.Controllers.Warehouse;
/// <summary>
/// 文件服务。对外暴露 VolPro 文件系统中的静态文件(截图、导出等)。
/// 不走 VolPro JWT 认证体系——网关 B 组接口直接调用。
/// </summary>
[ApiController]
[AllowAnonymous]
public class FileServiceController : Controller
{
/// <summary>
/// 获取截图文件。
/// 文件存放于 VolPro.WebApi/Download/Screenshots/ 目录。
/// </summary>
/// <param name="filename">文件名(含扩展名,如 abc.png</param>
[HttpGet("api/gateway/screenshots/{filename}")]
public IActionResult GetScreenshot(string filename)
{
// 安全检查:禁止路径穿越(.., /, \)
if (string.IsNullOrWhiteSpace(filename) ||
filename.Contains("..") ||
filename.Contains('/') ||
filename.Contains('\\'))
return BadRequest(new { error = "非法文件名" });
var folder = Path.Combine(AppContext.BaseDirectory, "Download", "Screenshots");
var filePath = Path.Combine(folder, filename);
if (!System.IO.File.Exists(filePath))
return NotFound(new { error = "文件不存在" });
var ext = Path.GetExtension(filename).ToLowerInvariant();
var contentType = ext switch
{
".png" => "image/png",
".jpg" or ".jpeg" => "image/jpeg",
".gif" => "image/gif",
_ => "application/octet-stream"
};
return PhysicalFile(filePath, contentType);
}
}
@@ -0,0 +1,7 @@
中文提示 : 检测到你没有开启文件,AllowLoadLocalInfile=true加到自符串上,已自动执行 SET GLOBAL local_infile=1 在试一次
English Message : Loading local data is disabled; this must be enabled on both the client and server sides at SqlSugar.Check.ExceptionEasy(String enMessage, String cnMessage)
at SqlSugar.MySqlFastBuilder.ExecuteBulkCopyAsync(DataTable dt)
at SqlSugar.FastestProvider`1._BulkCopy(List`1 datas)
at SqlSugar.FastestProvider`1.BulkCopyAsync(List`1 datas)
at SqlSugar.FastestProvider`1.BulkCopy(List`1 datas)
at VolPro.Core.Services.Logger.Start() in D:\Code\SecMPS\api_sqlsugar\VolPro.Core\Services\Logger.cs:line 194SqlSugar
@@ -21,6 +21,7 @@ using System;
using System.Collections.Generic; using System.Collections.Generic;
using System.Threading.Tasks; using System.Threading.Tasks;
using System.Text.Json; using System.Text.Json;
using Microsoft.Extensions.Logging;
namespace Warehouse.Services namespace Warehouse.Services
{ {
@@ -31,16 +32,19 @@ namespace Warehouse.Services
{ {
private readonly IHttpContextAccessor _httpContextAccessor; private readonly IHttpContextAccessor _httpContextAccessor;
private readonly Igateway_nodesRepository _repository; private readonly Igateway_nodesRepository _repository;
private readonly ILogger<gateway_nodesService> _logger;
[ActivatorUtilitiesConstructor] [ActivatorUtilitiesConstructor]
public gateway_nodesService( public gateway_nodesService(
Igateway_nodesRepository dbRepository, Igateway_nodesRepository dbRepository,
IHttpContextAccessor httpContextAccessor IHttpContextAccessor httpContextAccessor,
ILogger<gateway_nodesService> logger
) )
: base(dbRepository) : base(dbRepository)
{ {
_httpContextAccessor = httpContextAccessor; _httpContextAccessor = httpContextAccessor;
_repository = dbRepository; _repository = dbRepository;
_logger = logger;
} }
/// <summary> /// <summary>
@@ -51,39 +55,53 @@ namespace Warehouse.Services
[Obsolete("由 A1 API Controller 自动调用,不建议手动调用")] [Obsolete("由 A1 API Controller 自动调用,不建议手动调用")]
public async Task<gateway_nodes> RegisterNodeAsync(string nodeCode, string token, string adapterTypes, string baseUrl) public async Task<gateway_nodes> RegisterNodeAsync(string nodeCode, string token, string adapterTypes, string baseUrl)
{ {
var existingList = await _repository.FindAsIQueryable(x => x.NodeCode == nodeCode).ToListAsync(); _logger.LogInformation("[A1] 网关注册: NodeCode={Node}, Adapters={Adapters}", nodeCode, adapterTypes);
var existing = existingList.FirstOrDefault(); try
gateway_nodes entity;
if (existing != null)
{ {
if (existing.NodeToken != token) var existingList = await _repository.FindAsIQueryable(x => x.NodeCode == nodeCode).ToListAsync();
throw new UnauthorizedAccessException("NodeToken 不匹配"); var existing = existingList.FirstOrDefault();
existing.AdapterTypes = adapterTypes; gateway_nodes entity;
existing.BaseUrl = baseUrl; if (existing != null)
existing.IsOnline = "在线";
existing.LastHeartbeat = DateTime.Now;
_repository.DbContext.Updateable(existing).ExecuteCommand();
entity = existing;
}
else
{
entity = new gateway_nodes
{ {
NodeCode = nodeCode, if (existing.NodeToken != token)
NodeName = nodeCode, {
NodeToken = token, _logger.LogWarning("[A1] 注册失败: NodeCode={Node} Token不匹配", nodeCode);
AdapterTypes = adapterTypes, throw new UnauthorizedAccessException("NodeToken 不匹配");
BaseUrl = baseUrl, }
IsOnline = "在线",
Enable = "启用", existing.AdapterTypes = adapterTypes;
LastHeartbeat = DateTime.Now, existing.BaseUrl = baseUrl;
CreateDate = DateTime.Now existing.IsOnline = "在线";
}; existing.LastHeartbeat = DateTime.Now;
_repository.DbContext.Insertable(entity).ExecuteCommand(); _repository.DbContext.Updateable(existing).ExecuteCommand();
entity = existing;
_logger.LogInformation("[A1] 网关注册(更新): NodeId={Id}, NodeCode={Node}", entity.NodeId, nodeCode);
}
else
{
entity = new gateway_nodes
{
NodeCode = nodeCode,
NodeName = nodeCode,
NodeToken = token,
AdapterTypes = adapterTypes,
BaseUrl = baseUrl,
IsOnline = "在线",
Enable = "启用",
LastHeartbeat = DateTime.Now,
CreateDate = DateTime.Now
};
_repository.DbContext.Insertable(entity).ExecuteCommand();
_logger.LogInformation("[A1] 网关注册(新增): NodeId={Id}, NodeCode={Node}", entity.NodeId, nodeCode);
}
return entity;
}
catch (Exception ex)
{
_logger.LogError(ex, "[A1] 注册异常: NodeCode={Node}", nodeCode);
throw;
} }
return entity;
} }
/// <summary> /// <summary>
@@ -92,14 +110,27 @@ namespace Warehouse.Services
[Obsolete("由 A2 API Controller 自动调用,不建议手动调用")] [Obsolete("由 A2 API Controller 自动调用,不建议手动调用")]
public async Task UpdateHeartbeatAsync(string nodeCode, string token) public async Task UpdateHeartbeatAsync(string nodeCode, string token)
{ {
var entityList = await _repository.FindAsIQueryable(x => x.NodeCode == nodeCode && x.NodeToken == token).ToListAsync(); try
var entity = entityList.FirstOrDefault(); {
if (entity == null) var entityList = await _repository.FindAsIQueryable(x => x.NodeCode == nodeCode && x.NodeToken == token).ToListAsync();
throw new UnauthorizedAccessException("认证失败:NodeCode 或 Token 无效"); var entity = entityList.FirstOrDefault();
if (entity == null)
{
_logger.LogWarning("[A2] 心跳认证失败: NodeCode={Node}", nodeCode);
throw new UnauthorizedAccessException("认证失败:NodeCode 或 Token 无效");
}
entity.IsOnline = "在线"; entity.IsOnline = "在线";
entity.LastHeartbeat = DateTime.Now; entity.LastHeartbeat = DateTime.Now;
_repository.DbContext.Updateable(entity).ExecuteCommand(); _repository.DbContext.Updateable(entity).ExecuteCommand();
_logger.LogDebug("[A2] 心跳更新: NodeCode={Node}", nodeCode);
}
catch (UnauthorizedAccessException) { throw; }
catch (Exception ex)
{
_logger.LogError(ex, "[A2] 心跳异常: NodeCode={Node}", nodeCode);
throw;
}
} }
/// <summary> /// <summary>
@@ -110,69 +141,83 @@ namespace Warehouse.Services
[Obsolete("由 A3 API Controller 自动调用,不建议手动调用")] [Obsolete("由 A3 API Controller 自动调用,不建议手动调用")]
public async Task<(int added, int updated)> SyncDevicesAsync(int gatewayNodeId, List<SyncDeviceItem> devices) public async Task<(int added, int updated)> SyncDevicesAsync(int gatewayNodeId, List<SyncDeviceItem> devices)
{ {
var db = _repository.DbContext; _logger.LogInformation("[A3] 设备同步开始: NodeId={Id}, 设备数={Count}", gatewayNodeId, devices.Count);
try
var adapterCodes = devices.Select(d => d.AdapterCode).Distinct().ToList();
var existingIds = db.Queryable<base_device>()
.Where(x => x.NodeId == gatewayNodeId && adapterCodes.Contains(x.AdapterCode))
.ToList()
.ToDictionary(x => (x.AdapterCode, x.SourceId), x => x.DeviceId);
int added = 0, updated = 0;
foreach (var d in devices)
{ {
var key = (d.AdapterCode, d.SourceId); var db = _repository.DbContext;
existingIds.TryGetValue(key, out var existingId);
bool isNew = existingId == 0;
int? parentDeviceId = null; var adapterCodes = devices.Select(d => d.AdapterCode).Distinct().ToList();
if (!string.IsNullOrEmpty(d.ParentSourceId)) // 全局去重——不限定 NodeId,防止网关重启后 NodeId 变化导致重复插入
{ var existingIds = db.Queryable<base_device>()
existingIds.TryGetValue((d.AdapterCode, d.ParentSourceId), out var pid); .Where(x => adapterCodes.Contains(x.AdapterCode))
if (pid > 0) parentDeviceId = pid; .ToList()
} .ToDictionary(x => (x.AdapterCode, x.SourceId), x => x.DeviceId);
if (isNew) int added = 0, updated = 0;
foreach (var d in devices)
{ {
var entity = new base_device var key = (d.AdapterCode, d.SourceId);
existingIds.TryGetValue(key, out var existingId);
bool isNew = existingId == 0;
int? parentDeviceId = null;
if (!string.IsNullOrEmpty(d.ParentSourceId))
{ {
DeviceName = d.Name ?? $"DEV_{d.SourceId}", existingIds.TryGetValue((d.AdapterCode, d.ParentSourceId), out var pid);
AdapterCode = d.AdapterCode, if (pid > 0) parentDeviceId = pid;
SourceId = d.SourceId, }
DeviceCategory = d.Category,
DeviceGroup = d.Group, if (isNew)
NodeId = gatewayNodeId,
IsParent = d.IsParent ? "是" : "否",
ParentDeviceId = parentDeviceId,
IsOnline = d.IsOnline ? "在线" : "离线",
IpAddress = d.IpAddress,
Port = d.Port,
ExtraData = d.ExtraDataJson,
Enable = "启用",
LastSyncTime = DateTime.Now,
CreateDate = DateTime.Now
};
db.Insertable(entity).ExecuteCommand();
added++;
}
else
{
var entity = db.Queryable<base_device>().InSingle(existingId);
if (entity != null)
{ {
entity.IsOnline = d.IsOnline ? "在线" : "离线"; var entity = new base_device
entity.IsParent = d.IsParent ? "是" : "否"; {
entity.ParentDeviceId = parentDeviceId ?? entity.ParentDeviceId; DeviceName = d.Name ?? $"DEV_{d.SourceId}",
entity.IpAddress = d.IpAddress; AdapterCode = d.AdapterCode,
entity.Port = d.Port; SourceId = d.SourceId,
entity.ExtraData = d.ExtraDataJson ?? entity.ExtraData; DeviceCategory = d.Category,
entity.LastSyncTime = DateTime.Now; DeviceGroup = d.Group,
db.Updateable(entity).ExecuteCommand(); NodeId = gatewayNodeId,
updated++; IsParent = d.IsParent ? "是" : "否",
ParentDeviceId = parentDeviceId,
IsOnline = d.IsOnline ? "在线" : "离线",
IpAddress = d.IpAddress,
Port = d.Port,
ExtraData = d.ExtraDataJson,
Enable = "启用",
LastSyncTime = DateTime.Now,
CreateDate = DateTime.Now
};
var newId = db.Insertable(entity).ExecuteReturnIdentity();
// 补入去重字典,同批次子设备可查到父设备
existingIds[(d.AdapterCode, d.SourceId)] = Convert.ToInt32(newId);
added++;
}
else
{
var entity = db.Queryable<base_device>().InSingle(existingId);
if (entity != null)
{
entity.NodeId = gatewayNodeId; // 重新归属到当前网关
entity.IsOnline = d.IsOnline ? "在线" : "离线";
entity.IsParent = d.IsParent ? "是" : "否";
entity.ParentDeviceId = parentDeviceId ?? entity.ParentDeviceId;
entity.IpAddress = d.IpAddress;
entity.Port = d.Port;
entity.ExtraData = d.ExtraDataJson ?? entity.ExtraData;
entity.LastSyncTime = DateTime.Now;
db.Updateable(entity).ExecuteCommand();
updated++;
}
} }
} }
_logger.LogInformation("[A3] 设备同步完成: 新增{Added}台, 更新{Updated}台", added, updated);
return (added, updated);
}
catch (Exception ex)
{
_logger.LogError(ex, "[A3] 设备同步异常: NodeId={Id}, 设备数={Count}", gatewayNodeId, devices.Count);
throw;
} }
return (added, updated);
} }
} }
+20 -18
View File
@@ -44,7 +44,7 @@ CREATE TABLE base_device (
INDEX IX_Sync (AdapterCode, SourceId), INDEX IX_Sync (AdapterCode, SourceId),
INDEX IX_Point (PointId), INDEX IX_Point (PointId),
INDEX IX_Parent (ParentDeviceId), INDEX IX_Parent (ParentDeviceId),
INDEX IX_Gateway (GatewayNodeId), INDEX IX_Gateway (NodeId),
INDEX IX_Group (DeviceGroup) INDEX IX_Group (DeviceGroup)
) COMMENT '统一设备主表'; ) COMMENT '统一设备主表';
@@ -169,25 +169,27 @@ CREATE TABLE gateway_nodes (
-- 规则条件/动作的 ValueId 绑定到此表的 VariableId -- 规则条件/动作的 ValueId 绑定到此表的 VariableId
-- DeviceId 关联 base_device.DeviceId -- DeviceId 关联 base_device.DeviceId
-- ================================================= -- =================================================
DROP TABLE IF EXISTS warehouse_variable;
CREATE TABLE warehouse_variable ( CREATE TABLE warehouse_variable (
VariableId INT IDENTITY(1,1) PRIMARY KEY, VariableId INT AUTO_INCREMENT COMMENT '变量ID(自增主键)',
DeviceId INT NOT NULL, DeviceId INT NOT NULL COMMENT '关联设备ID(base_device.DeviceId)',
VariableName NVARCHAR(255) NOT NULL, -- 温度/湿度/人数 VariableName VARCHAR(255) NOT NULL COMMENT '变量名称(温度/湿度/人数等)',
PointIndex INT DEFAULT 0, -- MC4 pointIndex / Owl 统计量编码 PointIndex INT DEFAULT 0 COMMENT '点位索引(MC4 pointIndex / Owl统计量编码)',
Unit NVARCHAR(50) NULL, -- ℃/%/人 Unit VARCHAR(50) NULL COMMENT '单位(℃/%/人)',
SortOrder INT DEFAULT 0 SortOrder INT DEFAULT 0 COMMENT '排序顺序',
); PRIMARY KEY (VariableId),
INDEX IX_warehouse_variable_DeviceId (DeviceId)
CREATE INDEX IX_warehouse_variable_DeviceId ON warehouse_variable (DeviceId); ) COMMENT '规则引擎变量定义表';
-- F3.2 规则引擎滞后窗 (hysteresis) -- F3.2 规则引擎滞后窗 (hysteresis)
ALTER TABLE warehouse_rulecondition ADD -- 触发阈值与恢复阈值之间留缓冲区间,防止阈值附近反复抖动
RecoveryThreshold_Numeric DECIMAL(18,2) NULL, ALTER TABLE warehouse_rulecondition
RecoveryThreshold_Switch NVARCHAR(50) NULL; ADD COLUMN RecoveryThreshold_Numeric DECIMAL(18,2) NULL COMMENT '数值型恢复阈值(如>28℃触发,≤26℃恢复)',
ADD COLUMN RecoveryThreshold_Switch VARCHAR(50) NULL COMMENT '开关型恢复阈值(如开触发,关恢复)';
-- F3.3 条件级冷却 -- F3.3 条件级冷却 (cooldown)
ALTER TABLE warehouse_rulecondition ADD -- 冷却期内条件再次命中不重复执行动作,防止告警轰炸
LastTriggered DATETIME NULL, ALTER TABLE warehouse_rulecondition
LastTriggerValue DECIMAL(18,2) NULL; ADD COLUMN LastTriggered DATETIME NULL COMMENT '该条件上次触发的时间',
ADD COLUMN LastTriggerValue DECIMAL(18,2) NULL COMMENT '该条件上次触发时的实际值';
@@ -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 个 Task5 个适配器 = 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/getN=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.Abstractions;
using IntegrationGateway.Core.Infrastructure; using IntegrationGateway.Core.Infrastructure;
using IntegrationGateway.Core.Models; using IntegrationGateway.Core.Models;
using Microsoft.Extensions.Logging;
using System.Net.Http.Json; using System.Net.Http.Json;
using System.Text; using System.Text;
using System.Text.Json; using System.Text.Json;
@@ -22,6 +23,8 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
private readonly HttpClient _http; private readonly HttpClient _http;
private readonly KmsAuthHelper _auth; private readonly KmsAuthHelper _auth;
private readonly RateLimiter _limiter = new(5); private readonly RateLimiter _limiter = new(5);
private static readonly JsonSerializerOptions JsonOpts = new() { PropertyNameCaseInsensitive = true };
private readonly ILogger<KmsAdapter> _logger;
/// <summary>适配器编码,格式 "KMS:{实例名}"</summary> /// <summary>适配器编码,格式 "KMS:{实例名}"</summary>
public string AdapterCode { get; } public string AdapterCode { get; }
@@ -38,11 +41,12 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
/// <param name="baseUrl">KMS 服务地址</param> /// <param name="baseUrl">KMS 服务地址</param>
/// <param name="clientId">KMS 客户端 ID</param> /// <param name="clientId">KMS 客户端 ID</param>
/// <param name="clientSecret">KMS 客户端密钥</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; AdapterCode = adapterCode;
_http = http; _http = http;
_auth = new KmsAuthHelper(http, baseUrl, clientId, clientSecret); _logger = logger;
_auth = new KmsAuthHelper(http, baseUrl, clientId, clientSecret, logger);
} }
/// <summary>初始化适配器:获取 KMS Token</summary> /// <summary>初始化适配器:获取 KMS Token</summary>
@@ -59,9 +63,11 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
{ {
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
var resp = await client.GetAsync("/prod-api/heartBeat"); var resp = await client.GetAsync("/prod-api/heartBeat");
return resp.IsSuccessStatusCode; var ok = resp.IsSuccessStatusCode;
_logger.LogDebug("[{Code}] 健康检查完成,状态码={Status}", AdapterCode, ok ? 200 : resp.StatusCode);
return ok;
} }
catch (Exception ex) { Console.Error.WriteLine($"[{AdapterCode}] HealthCheck 失败: {ex.Message}"); return false; } catch (Exception ex) { Console.Error.WriteLine($"[{AdapterCode}] 健康检查失败: {ex.Message}"); return false; }
} }
// ═══════════════════════════════════════════ // ═══════════════════════════════════════════
@@ -76,55 +82,83 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
{ {
await _limiter.WaitAsync(); await _limiter.WaitAsync();
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
var resp = await client.PostAsync("/prod-api/getOpenerList",
new StringContent("{}", Encoding.UTF8, "application/json"));
resp.EnsureSuccessStatusCode();
var data = await resp.Content.ReadFromJsonAsync<KmsOpenerListResponse>()!;
var devices = new List<StandardDevice>(); var devices = new List<StandardDevice>();
foreach (var locker in data.Rows ?? new())
// ① 获取锁柜列表(父设备)
var lockerJson = await client.GetStringAsync(
$"/prod-api/kms/locker/list?pageNum=1&pageSize=100");
var lockerData = JsonSerializer.Deserialize<KmsLockerListResponse>(lockerJson, JsonOpts);
_logger.LogDebug("[{Code}] 获取锁柜列表,共{Count}个", AdapterCode, lockerData?.rows?.Count ?? 0);
// ② 获取锁芯列表(子设备)
var lockholeJson = await client.GetStringAsync(
$"/prod-api/kms/lockhole/list?state=1&pageNum=1&pageSize=1000");
var lockholeData = JsonSerializer.Deserialize<KmsLockholeListResponse>(lockholeJson, JsonOpts);
_logger.LogDebug("[{Code}] 获取锁芯列表,共{Count}个", AdapterCode, lockholeData?.Rows?.Count ?? 0);
// ③ 映射:锁柜 → 父设备
var lockerDict = new Dictionary<int, string>(); // lockerId → lockerName
if (lockerData?.rows != null)
{ {
devices.Add(MapLockerToDevice(locker)); foreach (var locker in lockerData.rows)
if (locker.LockholeList != null) {
devices.AddRange(locker.LockholeList.Select(h => MapLockholeToDevice(h, locker.LockerId))); var lockerSourceId = $"locker_{locker.id}";
lockerDict[locker.id] = locker.name ?? $"锁柜{locker.id}";
devices.Add(new StandardDevice
{
AdapterCode = AdapterCode,
SourceId = lockerSourceId,
Name = locker.name ?? $"锁柜{locker.id}",
Category = "智能钥匙柜",
Group = "门禁设备",
IsParent = true,
IsOnline = locker.state == 1,
Extra = new Dictionary<string, object?>
{
["lockerCode"] = locker.code,
["lockholeCount"] = locker.num ?? 0
//["inCount"] = locker.inNum ?? 0,
//["outCount"] = locker.outNum ?? 0
}
});
}
} }
// ④ 映射:锁芯 → 子设备
if (lockholeData?.Rows != null)
{
foreach (var hole in lockholeData.Rows)
{
var lockerId = hole.Locker?.Id ?? hole.LockerId;
var lockerName = hole.Locker?.Name ?? (lockerId > 0 ? lockerDict.GetValueOrDefault(lockerId) : null);
devices.Add(new StandardDevice
{
AdapterCode = AdapterCode,
SourceId = $"lockhole_{hole.Id}",
Name = hole.Opener?.CnName ?? $"锁芯{hole.Sort ?? hole.Id}",
Category = "钥匙位",
Group = "门禁设备",
IsParent = false,
IsOnline = hole.State == 1,
ParentSourceId = lockerId > 0 ? $"locker_{lockerId}" : null,
Extra = new Dictionary<string, object?>
{
["lockholeSort"] = hole.Sort,
["lockerName"] = lockerName,
["openerId"] = hole.Opener?.Id,
["lockholeState"] = hole.State,
["openerName"] = hole.Opener?.CnName,
["openerNumber"] = hole.Opener?.Number
}
});
}
}
_logger.LogDebug("[{Code}] 设备映射完成,共{Total}台(父{Parent}/子{Child})",
AdapterCode, devices.Count, devices.Count(d => d.IsParent), devices.Count(d => !d.IsParent));
return new PagedResult<StandardDevice> { Items = devices, Total = devices.Count }; return new PagedResult<StandardDevice> { Items = devices, Total = devices.Count };
} }
/// <summary>KMS 柜体 → StandardDevice(父设备)</summary>
private static StandardDevice MapLockerToDevice(KmsLocker locker) => new()
{
SourceId = $"locker_{locker.LockerId}",
Name = locker.LockerName ?? $"柜体{locker.LockerId}",
Category = "智能钥匙柜",
Group = "门禁设备",
IsParent = true,
IsOnline = true,
Extra = new Dictionary<string, object?>
{
["lockerCode"] = locker.LockerCode,
["lockholeCount"] = locker.LockholeList?.Count ?? 0
}
};
/// <summary>KMS 锁孔 → StandardDevice(子设备)</summary>
private static StandardDevice MapLockholeToDevice(KmsLockhole hole, int lockerId) => new()
{
SourceId = $"lockhole_{lockerId}_{hole.LockholeSort}",
Name = hole.OpenerName ?? $"锁孔{hole.LockholeSort}",
Category = "钥匙位",
Group = "门禁设备",
IsParent = false,
// KMS openerState: 1=在柜, 2=借出, 3=录入, 10=丢失 (数值编码或中文)
IsOnline = hole.OpenerState == "1" || hole.OpenerState == "在柜", // KMS: 1=在柜/2=借出/3=录入/10=丢失
ParentSourceId = $"locker_{lockerId}",
Extra = new Dictionary<string, object?>
{
["openerId"] = hole.OpenerId,
["openerType"] = hole.OpenerType,
["openerState"] = hole.OpenerState
}
};
// ═══════════════════════════════════════════ // ═══════════════════════════════════════════
// IHasAlarms — 告警(2.18.7 告警列表) // IHasAlarms — 告警(2.18.7 告警列表)
@@ -147,6 +181,7 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
new StringContent(body, Encoding.UTF8, "application/json")); new StringContent(body, Encoding.UTF8, "application/json"));
resp.EnsureSuccessStatusCode(); resp.EnsureSuccessStatusCode();
var data = await resp.Content.ReadFromJsonAsync<KmsWarningListResponse>()!; var data = await resp.Content.ReadFromJsonAsync<KmsWarningListResponse>()!;
_logger.LogDebug("[{Code}] 获取告警列表,共{Count}条", AdapterCode, data.Rows?.Count ?? 0);
var alarms = (data.Rows ?? new()).Select(w => new StandardAlarm var alarms = (data.Rows ?? new()).Select(w => new StandardAlarm
{ {
@@ -207,6 +242,7 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
new StringContent(body, Encoding.UTF8, "application/json")); new StringContent(body, Encoding.UTF8, "application/json"));
resp.EnsureSuccessStatusCode(); resp.EnsureSuccessStatusCode();
var data = await resp.Content.ReadFromJsonAsync<KmsRecordListResponse>()!; var data = await resp.Content.ReadFromJsonAsync<KmsRecordListResponse>()!;
_logger.LogDebug("[{Code}] 获取借还记录,共{Count}条", AdapterCode, data.Rows?.Count ?? 0);
return new PagedResult<KmsRecord> { Items = data.Rows ?? new(), Total = data.Total }; return new PagedResult<KmsRecord> { Items = data.Rows ?? new(), Total = data.Total };
} }
@@ -234,6 +270,7 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
new StringContent(body, Encoding.UTF8, "application/json")); new StringContent(body, Encoding.UTF8, "application/json"));
resp.EnsureSuccessStatusCode(); resp.EnsureSuccessStatusCode();
var data = await resp.Content.ReadFromJsonAsync<KmsPermissionListResponse>()!; var data = await resp.Content.ReadFromJsonAsync<KmsPermissionListResponse>()!;
_logger.LogDebug("[{Code}] 获取授权记录,共{Count}条", AdapterCode, data.Rows?.Count ?? 0);
return new PagedResult<KmsPermission> { Items = data.Rows ?? new(), Total = data.Total }; return new PagedResult<KmsPermission> { Items = data.Rows ?? new(), Total = data.Total };
} }
@@ -254,6 +291,7 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
var resp = await client.PostAsJsonAsync("/prod-api/batchDeleteStaff", staffUuids); var resp = await client.PostAsJsonAsync("/prod-api/batchDeleteStaff", staffUuids);
resp.EnsureSuccessStatusCode(); resp.EnsureSuccessStatusCode();
_logger.LogDebug("[{Code}] 批量删除员工({Count}人),状态={Status}", AdapterCode, staffUuids.Count, resp.StatusCode);
} }
/// <summary>2.4.3 远程授权开门</summary> /// <summary>2.4.3 远程授权开门</summary>
@@ -271,6 +309,7 @@ public class KmsAdapter : IHasFlatDevices, IHasAlarms, IAcceptsControl, IHasBusi
await _limiter.WaitAsync(); await _limiter.WaitAsync();
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
var resp = await client.PostAsync($"/thirdPlatlogin?username={Uri.EscapeDataString(username)}", null); var resp = await client.PostAsync($"/thirdPlatlogin?username={Uri.EscapeDataString(username)}", null);
_logger.LogDebug("[{Code}] 第三方登录({User}),状态={Status}", AdapterCode, username, resp.StatusCode);
if (resp.StatusCode == System.Net.HttpStatusCode.Redirect) if (resp.StatusCode == System.Net.HttpStatusCode.Redirect)
return resp.Headers.Location?.ToString(); return resp.Headers.Location?.ToString();
resp.EnsureSuccessStatusCode(); resp.EnsureSuccessStatusCode();
@@ -298,10 +337,12 @@ OpenerIds = parameters.TryGetValue("lockholeSort", out var lh) ? new List<int> {
}; };
await RemoteAuthorizeAsync(req); await RemoteAuthorizeAsync(req);
} }
_logger.LogDebug("[{Code}] 设备控制完成 {Id} cmd={Cmd}", AdapterCode, sourceDeviceId, command);
return new ControlResult { Success = true }; return new ControlResult { Success = true };
} }
catch (Exception ex) catch (Exception ex)
{ {
_logger.LogDebug("[{Code}] 设备控制失败: {Msg}", AdapterCode, ex.Message);
return new ControlResult { Success = false, Message = ex.Message }; return new ControlResult { Success = false, Message = ex.Message };
} }
} }
@@ -360,7 +401,7 @@ OpenerIds = parameters.TryGetValue("lockholeSort", out var lh) ? new List<int> {
await BatchSyncStaffAsync(staffList); await BatchSyncStaffAsync(staffList);
return new SyncResult { SuccessCount = staffList.Count }; return new SyncResult { SuccessCount = staffList.Count };
} }
catch (Exception ex) { return new SyncResult { FailCount = items.Count, Message = ex.Message }; } catch (Exception ex) { _logger.LogDebug("[{Code}] 员工数据同步失败: {Msg}", AdapterCode, ex.Message); return new SyncResult { FailCount = items.Count, Message = ex.Message }; }
} }
/// <summary>从 KMS 批量删除数据</summary> /// <summary>从 KMS 批量删除数据</summary>
@@ -372,6 +413,48 @@ OpenerIds = parameters.TryGetValue("lockholeSort", out var lh) ? new List<int> {
await BatchDeleteStaffAsync(ids); await BatchDeleteStaffAsync(ids);
return new SyncResult { SuccessCount = ids.Count }; return new SyncResult { SuccessCount = ids.Count };
} }
catch (Exception ex) { return new SyncResult { FailCount = ids.Count, Message = ex.Message }; } catch (Exception ex) { _logger.LogDebug("[{Code}] 员工数据删除失败: {Msg}", AdapterCode, ex.Message); return new SyncResult { FailCount = ids.Count, Message = ex.Message }; }
}
// ═══════════════════════════════════════════
// 标准管理接口 — 详情查询(暂不实现增删改)
// ═══════════════════════════════════════════
/// <summary>获取锁柜详情 (2.16.5)</summary>
public async Task<KmsLockerInfo?> GetLockerDetailAsync(int lockerId)
{
await _limiter.WaitAsync();
var client = await _auth.GetAuthenticatedClientAsync();
var json = await client.GetStringAsync($"/prod-api/kms/locker/{lockerId}");
return JsonSerializer.Deserialize<KmsLockerInfo>(json);
}
/// <summary>获取锁芯详情 (2.17.4)</summary>
public async Task<KmsLockholeDetail?> GetLockholeDetailAsync(int lockholeId)
{
await _limiter.WaitAsync();
var client = await _auth.GetAuthenticatedClientAsync();
var json = await client.GetStringAsync($"/prod-api/kms/lockhole/{lockholeId}");
return JsonSerializer.Deserialize<KmsLockholeDetail>(json);
}
/// <summary>获取钥匙详情 (2.14.8)</summary>
public async Task<KmsOpenerInfo?> GetOpenerDetailAsync(int openerId)
{
await _limiter.WaitAsync();
var client = await _auth.GetAuthenticatedClientAsync();
var json = await client.GetStringAsync($"/prod-api/kms/opener/{openerId}");
return JsonSerializer.Deserialize<KmsOpenerInfo>(json);
}
/// <summary>获取锁芯列表,按锁柜ID过滤 (2.17.2)</summary>
public async Task<List<KmsLockholeInfo>> GetLockholesByLockerAsync(int lockerId)
{
await _limiter.WaitAsync();
var client = await _auth.GetAuthenticatedClientAsync();
var json = await client.GetStringAsync(
$"/prod-api/kms/lockhole/list?state=1&pageNum=1&pageSize=1000");
var data = JsonSerializer.Deserialize<KmsLockholeListResponse>(json);
return data?.Rows?.Where(h => (h.Locker?.Id ?? h.LockerId) == lockerId).ToList() ?? new();
} }
} }
@@ -1,4 +1,6 @@
using Microsoft.Extensions.Logging;
using System.Net.Http.Json; using System.Net.Http.Json;
using System.Text;
using System.Text.Json; using System.Text.Json;
namespace IntegrationGateway.Adapters.Kms; namespace IntegrationGateway.Adapters.Kms;
@@ -14,6 +16,8 @@ public class KmsAuthHelper
private readonly string _baseUrl; private readonly string _baseUrl;
private readonly string _clientId; private readonly string _clientId;
private readonly string _clientSecret; private readonly string _clientSecret;
private static readonly JsonSerializerOptions JsonOpts = new() { PropertyNameCaseInsensitive = true };
private readonly ILogger _logger;
private string? _token; private string? _token;
private DateTime _tokenExpiry = DateTime.MinValue; private DateTime _tokenExpiry = DateTime.MinValue;
@@ -24,12 +28,14 @@ public class KmsAuthHelper
/// <param name="baseUrl">KMS 服务地址</param> /// <param name="baseUrl">KMS 服务地址</param>
/// <param name="clientId">KMS 客户端 ID</param> /// <param name="clientId">KMS 客户端 ID</param>
/// <param name="clientSecret">KMS 客户端密钥</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; _http = http;
_baseUrl = baseUrl.TrimEnd('/'); _baseUrl = baseUrl.TrimEnd('/');
_clientId = clientId; _clientId = clientId;
_clientSecret = clientSecret; _clientSecret = clientSecret;
_logger = logger;
_http.DefaultRequestHeaders.TryAddWithoutValidation("User-Agent", "SecMPS-Gateway/1.0");
} }
/// <summary> /// <summary>
@@ -38,22 +44,31 @@ public class KmsAuthHelper
public async Task<string> GetTokenAsync() public async Task<string> GetTokenAsync()
{ {
if (!string.IsNullOrEmpty(_token) && DateTime.UtcNow < _tokenExpiry) if (!string.IsNullOrEmpty(_token) && DateTime.UtcNow < _tokenExpiry)
{
_logger?.LogDebug("KMS Token缓存命中(过期时间={Exp})", _tokenExpiry);
return _token; return _token;
}
var url = $"{_baseUrl}/prod-api/getToken?clientId={Uri.EscapeDataString(_clientId)}&clientSecret={Uri.EscapeDataString(_clientSecret)}"; var json = JsonSerializer.Serialize(new { clientId = _clientId, clientSecret = _clientSecret });
var resp = await _http.PostAsync(url, null); var content = new ByteArrayContent(Encoding.UTF8.GetBytes(json));
content.Headers.ContentType = new System.Net.Http.Headers.MediaTypeHeaderValue("application/json");
_logger?.LogDebug("KMS 开始获取Token(clientId={Cid})", _clientId);
var resp = await _http.PostAsync($"{_baseUrl}/prod-api/getToken", content);
resp.EnsureSuccessStatusCode(); resp.EnsureSuccessStatusCode();
_logger?.LogDebug("KMS Token请求完成,状态={Status}", resp.StatusCode);
var result = await resp.Content.ReadFromJsonAsync<KmsTokenResponse>() var result = await resp.Content.ReadFromJsonAsync<KmsTokenResponse>()
?? throw new Exception("KMS Token 响应为空"); ?? throw new Exception("KMS Token 响应为空");
if (result.Code != 200) if (result.code != 200)
throw new Exception($"KMS 认证失败: code={result.Code}"); throw new Exception($"KMS 认证失败: code={result.code}, msg={result.msg}");
if (string.IsNullOrEmpty(result.data?.token))
throw new Exception("KMS Token 响应中 data.token 为空");
_token = result.Token; _token = result.data.token;
_tokenExpiry = DateTime.UtcNow.AddMinutes(25); _tokenExpiry = DateTime.UtcNow.AddMinutes(25);
_logger?.LogDebug("KMS Token获取成功,25分钟有效");
return _token; return _token;
} }
/// <summary> /// <summary>
/// 创建一个已认证的 HttpClient,自动附带 Authorization: Bearer 头。 /// 创建一个已认证的 HttpClient,自动附带 Authorization: Bearer 头。
/// </summary> /// </summary>
@@ -62,6 +77,7 @@ public class KmsAuthHelper
var token = await GetTokenAsync(); var token = await GetTokenAsync();
var client = new HttpClient { BaseAddress = new Uri(_baseUrl) }; var client = new HttpClient { BaseAddress = new Uri(_baseUrl) };
client.DefaultRequestHeaders.Add("Authorization", $"Bearer {token}"); client.DefaultRequestHeaders.Add("Authorization", $"Bearer {token}");
client.DefaultRequestHeaders.TryAddWithoutValidation("User-Agent", "SecMPS-Gateway/1.0");
return client; return client;
} }
@@ -11,9 +11,15 @@ namespace IntegrationGateway.Adapters.Kms;
/// <summary>POST /prod-api/getToken 响应</summary> /// <summary>POST /prod-api/getToken 响应</summary>
public class KmsTokenResponse public class KmsTokenResponse
{ {
public int Code { get; set; } public int code { get; set; }
public string Token { get; set; } = ""; public string? msg { get; set; }
public string? Msg { get; set; } public KmsTokenData? data { get; set; }
}
/// <summary>Token 数据嵌套对象</summary>
public class KmsTokenData
{
public string token { get; set; } = "";
} }
// ═══════════════════════════════════════════ // ═══════════════════════════════════════════
@@ -158,19 +164,29 @@ public class KmsStaff
/// <summary>KMS 柜体列表响应</summary> /// <summary>KMS 柜体列表响应</summary>
public class KmsLockerListResponse public class KmsLockerListResponse
{ {
public int Code { get; set; } public int code { get; set; }
public int Total { get; set; } public string? msg { get; set; }
public List<KmsLockerInfo>? Rows { get; set; } public int total { get; set; }
public List<KmsLockerInfo>? rows { get; set; }
} }
/// <summary>KMS 柜体详细信息</summary> /// <summary>KMS 柜体详细信息</summary>
public class KmsLockerInfo public class KmsLockerInfo
{ {
public int Id { get; set; } public int id { get; set; }
public string? Name { get; set; } public string? name { get; set; }
public string? Code { get; set; } public string? code { get; set; }
public int State { get; set; } public int state { get; set; }
public int? DeptId { get; set; } public int? deptId { get; set; }
public string? remark { get; set; }
public string? createBy { get; set; }
public string? createTime { get; set; }
public string? updateBy { get; set; }
public string? updateTime { get; set; }
public object? openers { get; set; }
public int? num { get; set; }
public int? InNum { get; set; }
public int? OutNum { get; set; }
public List<KmsLockhole>? LockholeList { get; set; } public List<KmsLockhole>? LockholeList { get; set; }
} }
@@ -178,6 +194,7 @@ public class KmsLockerInfo
public class KmsLockholeListResponse public class KmsLockholeListResponse
{ {
public int Code { get; set; } public int Code { get; set; }
public string? Msg { get; set; }
public int Total { get; set; } public int Total { get; set; }
public List<KmsLockholeInfo>? Rows { get; set; } public List<KmsLockholeInfo>? Rows { get; set; }
} }
@@ -187,9 +204,64 @@ public class KmsLockholeInfo
{ {
public int Id { get; set; } public int Id { get; set; }
public int LockerId { get; set; } public int LockerId { get; set; }
public int LockholeSort { get; set; } public string? Code { get; set; }
public int? Sort { get; set; }
public int State { get; set; }
public string? Remark { get; set; }
public string? CreateBy { get; set; }
public string? CreateTime { get; set; }
public string? UpdateBy { get; set; }
public string? UpdateTime { get; set; }
public KmsLockerBrief? Locker { get; set; }
public KmsLockholeOpener? Opener { get; set; }
}
/// <summary>锁芯列表项中嵌套的钥匙信息</summary>
public class KmsLockholeOpener
{
public long? Id { get; set; }
public int LockerId { get; set; }
public string? Code { get; set; }
public long? LockholeId { get; set; }
public int? Sort { get; set; }
public long? OpenerGroupId { get; set; }
public string? CnName { get; set; }
public string? CnNamePy { get; set; }
public string? Number { get; set; }
public int? Type { get; set; }
public int? State { get; set; }
public int? WarnInterval { get; set; }
public string? Remark { get; set; }
public string? CreateBy { get; set; }
public string? CreateTime { get; set; }
public string? UpdateBy { get; set; }
public string? UpdateTime { get; set; }
public string? DetailImgUrl { get; set; }
public int? DeptId { get; set; }
public string? LendStaffName { get; set; }
public string? BorrowTime { get; set; }
public string? BackStaffName { get; set; }
public string? BackTime { get; set; }
}
public class KmsLockerBrief
{
public int Id { get; set; }
public string? Code { get; set; }
public string? Name { get; set; }
}
public class KmsLockholeDetail
{
public int Id { get; set; }
public int LockerId { get; set; }
public string? Code { get; set; }
public int? Sort { get; set; }
public int State { get; set; } public int State { get; set; }
public int? OpenerId { get; set; } public int? OpenerId { get; set; }
public string? Remark { get; set; }
public string? CreateTime { get; set; }
public string? UpdateTime { get; set; }
} }
/// <summary>KMS 钥匙列表响应</summary> /// <summary>KMS 钥匙列表响应</summary>
@@ -1,6 +1,7 @@
using IntegrationGateway.Core.Abstractions; using IntegrationGateway.Core.Abstractions;
using IntegrationGateway.Core.Infrastructure; using IntegrationGateway.Core.Infrastructure;
using IntegrationGateway.Core.Models; using IntegrationGateway.Core.Models;
using Microsoft.Extensions.Logging;
using System.Text; using System.Text;
using System.Text.Json; using System.Text.Json;
@@ -23,6 +24,7 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
private readonly Mc4AuthHelper _auth; private readonly Mc4AuthHelper _auth;
/// <summary>令牌桶限流器(2 QPS</summary> /// <summary>令牌桶限流器(2 QPS</summary>
private readonly RateLimiter _limiter = new(2); private readonly RateLimiter _limiter = new(2);
private readonly ILogger<Mc4Adapter> _logger;
/// <summary>适配器编码,格式 "MC4:实例名"</summary> /// <summary>适配器编码,格式 "MC4:实例名"</summary>
public string AdapterCode { get; } public string AdapterCode { get; }
@@ -38,11 +40,12 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
/// <param name="adapterCode">适配器编码</param> /// <param name="adapterCode">适配器编码</param>
/// <param name="http">HttpClient 实例</param> /// <param name="http">HttpClient 实例</param>
/// <param name="baseUrl">MC4.0 服务地址</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; AdapterCode = adapterCode;
_http = http; _http = http;
_auth = new Mc4AuthHelper(http, baseUrl, account, password); _logger = logger;
_auth = new Mc4AuthHelper(http, baseUrl, account, password, logger);
} }
/// <summary>初始化适配器:获取 MC4.0 Token</summary> /// <summary>初始化适配器:获取 MC4.0 Token</summary>
@@ -55,9 +58,11 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
{ {
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
var resp = await client.PostAsync("/api/central/auth/conf/get", null); var resp = await client.PostAsync("/api/central/auth/conf/get", null);
return resp.IsSuccessStatusCode; var ok = resp.IsSuccessStatusCode;
_logger.LogDebug("[{Code}] 健康检查完成,状态码={Status}", AdapterCode, ok ? 200 : resp.StatusCode);
return ok;
} }
catch (Exception ex) { Console.Error.WriteLine($"[{AdapterCode}] HealthCheck 失败: {ex.Message}"); return false; } catch (Exception ex) { Console.Error.WriteLine($"[{AdapterCode}] 健康检查失败: {ex.Message}"); return false; }
} }
// ═══════════════════════════════════════════ // ═══════════════════════════════════════════
@@ -76,6 +81,7 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
resp.EnsureSuccessStatusCode(); resp.EnsureSuccessStatusCode();
var json = await resp.Content.ReadAsStringAsync(); var json = await resp.Content.ReadAsStringAsync();
var tree = JsonSerializer.Deserialize<List<Mc4TreeNode>>(json)!; var tree = JsonSerializer.Deserialize<List<Mc4TreeNode>>(json)!;
_logger.LogDebug("[{Code}] 获取对象树,响应{Sz}字节,{Ct}个节点", AdapterCode, json.Length, tree.Count);
return tree.Select(MapNode).ToList(); return tree.Select(MapNode).ToList();
} }
@@ -107,6 +113,7 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
resp.EnsureSuccessStatusCode(); resp.EnsureSuccessStatusCode();
var json = await resp.Content.ReadAsStringAsync(); var json = await resp.Content.ReadAsStringAsync();
var values = JsonSerializer.Deserialize<List<Mc4PointValue>>(json)!; var values = JsonSerializer.Deserialize<List<Mc4PointValue>>(json)!;
_logger.LogDebug("[{Code}] 获取实时点位({Id}),响应{Sz}字节,{Ct}个点位", AdapterCode, sourceDeviceId, json.Length, values.Count);
return values.Select(v => new PointValue return values.Select(v => new PointValue
{ {
SourceDeviceId = sourceDeviceId, SourceDeviceId = sourceDeviceId,
@@ -123,8 +130,9 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
await _limiter.WaitAsync(); await _limiter.WaitAsync();
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
var body = JsonSerializer.Serialize(new { id = int.Parse(sourceDeviceId), index = pointIndex, value }); var body = JsonSerializer.Serialize(new { id = int.Parse(sourceDeviceId), index = pointIndex, value });
await client.PostAsync("/api/central/point/value/set", var resp = await client.PostAsync("/api/central/point/value/set",
new StringContent(body, Encoding.UTF8, "application/json")); new StringContent(body, Encoding.UTF8, "application/json"));
_logger.LogDebug("[{Code}] 设备控制写入 {Id}[{Idx}]={Val},状态={Status}", AdapterCode, sourceDeviceId, pointIndex, value, resp.StatusCode);
} }
// ═══════════════════════════════════════════ // ═══════════════════════════════════════════
@@ -142,8 +150,8 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
var body = JsonSerializer.Serialize(new Mc4AlarmQuery var body = JsonSerializer.Serialize(new Mc4AlarmQuery
{ {
From = from.ToString("yyyy-MM-dd HH:mm:ss"), From = from == DateTime.MinValue ? "" : from.ToString("yyyy-MM-dd HH:mm:ss"),
To = to.ToString("yyyy-MM-dd HH:mm:ss"), To = to == DateTime.MinValue ? "" : to.ToString("yyyy-MM-dd HH:mm:ss"),
Skip = (page - 1) * size, Skip = (page - 1) * size,
Limit = size, Limit = size,
Sort = 1 // 按时间降序 Sort = 1 // 按时间降序
@@ -153,6 +161,7 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
resp.EnsureSuccessStatusCode(); resp.EnsureSuccessStatusCode();
var json = await resp.Content.ReadAsStringAsync(); var json = await resp.Content.ReadAsStringAsync();
var result = JsonSerializer.Deserialize<Mc4AlarmQueryResult>(json)!; var result = JsonSerializer.Deserialize<Mc4AlarmQueryResult>(json)!;
_logger.LogDebug("[{Code}] 获取当前告警,响应{Sz}字节,{Ct}条", AdapterCode, json.Length, result.List?.Count ?? 0);
return new PagedResult<StandardAlarm> return new PagedResult<StandardAlarm>
{ {
Items = result.List?.Select(a => new StandardAlarm Items = result.List?.Select(a => new StandardAlarm
@@ -177,8 +186,10 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
await _limiter.WaitAsync(); await _limiter.WaitAsync();
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
var body = JsonSerializer.Serialize(new { id = alarmId, option = new { } }); 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")); new StringContent(body, Encoding.UTF8, "application/json"));
_logger.LogDebug("[{Code}] 确认告警({Id}),状态={Status}", AdapterCode, alarmId, cresp.StatusCode);
cresp.EnsureSuccessStatusCode();
} }
/// <summary>结束告警(同时写回 MC4.0</summary> /// <summary>结束告警(同时写回 MC4.0</summary>
@@ -187,8 +198,10 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
await _limiter.WaitAsync(); await _limiter.WaitAsync();
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
var body = JsonSerializer.Serialize(new { id = alarmId, option = new { } }); 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")); new StringContent(body, Encoding.UTF8, "application/json"));
_logger.LogDebug("[{Code}] 结束告警({Id}),状态={Status}", AdapterCode, alarmId, eresp.StatusCode);
eresp.EnsureSuccessStatusCode();
} }
/// <summary>MC4.0 告警等级数字 → 中文映射</summary> /// <summary>MC4.0 告警等级数字 → 中文映射</summary>
@@ -231,8 +244,8 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
var body = JsonSerializer.Serialize(new Mc4HisAlarmQuery var body = JsonSerializer.Serialize(new Mc4HisAlarmQuery
{ {
From = from.ToString("yyyy-MM-dd HH:mm:ss"), From = from == DateTime.MinValue ? "" : from.ToString("yyyy-MM-dd HH:mm:ss"),
To = to.ToString("yyyy-MM-dd HH:mm:ss"), To = to == DateTime.MinValue ? "" : to.ToString("yyyy-MM-dd HH:mm:ss"),
Skip = (page - 1) * size, Skip = (page - 1) * size,
Limit = size, Limit = size,
Sort = 1 Sort = 1
@@ -242,6 +255,7 @@ public class Mc4Adapter : IHasOwnDeviceTree, IHasPoints, IHasAlarms
resp.EnsureSuccessStatusCode(); resp.EnsureSuccessStatusCode();
var json = await resp.Content.ReadAsStringAsync(); var json = await resp.Content.ReadAsStringAsync();
var result = JsonSerializer.Deserialize<Mc4AlarmQueryResult>(json)!; var result = JsonSerializer.Deserialize<Mc4AlarmQueryResult>(json)!;
_logger.LogDebug("[{Code}] 获取当前告警,响应{Sz}字节,{Ct}条", AdapterCode, json.Length, result.List?.Count ?? 0);
return new PagedResult<StandardAlarm> return new PagedResult<StandardAlarm>
{ {
Items = (result.List ?? new()).Select(MapAlarmItem).ToList(), Items = (result.List ?? new()).Select(MapAlarmItem).ToList(),
@@ -1,3 +1,4 @@
using Microsoft.Extensions.Logging;
using System.Security.Cryptography; using System.Security.Cryptography;
using System.Text; using System.Text;
using System.Text.Json; using System.Text.Json;
@@ -20,22 +21,27 @@ public class Mc4AuthHelper
private readonly string _baseUrl; private readonly string _baseUrl;
private readonly string _account; private readonly string _account;
private readonly string _password; private readonly string _password;
private readonly ILogger _logger;
private string? _token; private string? _token;
private DateTime _tokenExpiry = DateTime.MinValue; private DateTime _tokenExpiry = DateTime.MinValue;
private bool? _needMd5; private bool? _needMd5 = false;
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; _http = http;
_baseUrl = baseUrl.TrimEnd('/'); _baseUrl = baseUrl.TrimEnd('/');
_account = account; _account = account;
_password = password; _password = password;
_logger = logger;
} }
public async Task<string> GetTokenAsync() public async Task<string> GetTokenAsync()
{ {
if (!string.IsNullOrEmpty(_token) && DateTime.UtcNow < _tokenExpiry) if (!string.IsNullOrEmpty(_token) && DateTime.UtcNow < _tokenExpiry)
{
_logger?.LogDebug("MC4 Token缓存命中(过期时间={Exp})", _tokenExpiry);
return _token; return _token;
}
// 1. 获取加密配置 // 1. 获取加密配置
if (!_needMd5.HasValue) if (!_needMd5.HasValue)
@@ -48,14 +54,16 @@ public class Mc4AuthHelper
var confJson = await confResp.Content.ReadAsStringAsync(); var confJson = await confResp.Content.ReadAsStringAsync();
var conf = JsonSerializer.Deserialize<Mc4ConfResponse>(confJson); var conf = JsonSerializer.Deserialize<Mc4ConfResponse>(confJson);
_needMd5 = conf?.Encrypt ?? false; _needMd5 = conf?.Encrypt ?? false;
_logger?.LogDebug("MC4 加密配置: encrypt={Enc}", _needMd5);
} }
else { _needMd5 = false; } else { _needMd5 = false; }
} }
catch { _needMd5 = false; } catch (Exception ex) { _needMd5 = false; _logger?.LogDebug("MC4 加密配置查询失败,MD5已关闭: {Msg}", ex.Message); }
} }
// 2. 登录获取 Token // 2. 登录获取 Token
var pwd = _needMd5 == true ? ComputeMd5(_password) : _password; var pwd = _needMd5 == true ? ComputeMd5(_password) : _password;
_logger?.LogDebug("MC4 开始登录: 账号={Acct}, MD5={Md5}", _account, _needMd5);
var loginBody = JsonSerializer.Serialize(new { account = _account, password = pwd }); var loginBody = JsonSerializer.Serialize(new { account = _account, password = pwd });
var resp = await _http.PostAsync($"{_baseUrl}/api/central/auth/login", var resp = await _http.PostAsync($"{_baseUrl}/api/central/auth/login",
new StringContent(loginBody, Encoding.UTF8, "application/json")); new StringContent(loginBody, Encoding.UTF8, "application/json"));
@@ -67,6 +75,7 @@ public class Mc4AuthHelper
throw new Exception("MC4 登录失败: Token 为空"); throw new Exception("MC4 登录失败: Token 为空");
_token = result.Token; _token = result.Token;
_tokenExpiry = DateTime.UtcNow.AddHours(7); _tokenExpiry = DateTime.UtcNow.AddHours(7);
_logger?.LogDebug("MC4 登录成功(账号={Acct})", _account);
return _token; return _token;
} }
@@ -1,6 +1,7 @@
using IntegrationGateway.Core.Abstractions; using IntegrationGateway.Core.Abstractions;
using IntegrationGateway.Core.Infrastructure; using IntegrationGateway.Core.Infrastructure;
using IntegrationGateway.Core.Models; using IntegrationGateway.Core.Models;
using Microsoft.Extensions.Logging;
using System.Text.Json; using System.Text.Json;
using System.Net.Http.Json; using System.Net.Http.Json;
@@ -24,6 +25,8 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
private readonly HttpClient _http; private readonly HttpClient _http;
private readonly OwlAuthHelper _auth; private readonly OwlAuthHelper _auth;
private readonly RateLimiter _limiter = new(5); private readonly RateLimiter _limiter = new(5);
private static readonly JsonSerializerOptions JsonOpts = new() { PropertyNameCaseInsensitive = true };
private readonly ILogger<OwlAdapter> _logger;
public string AdapterCode { get; } public string AdapterCode { get; }
public string DisplayName => $"Owl ({AdapterCode})"; public string DisplayName => $"Owl ({AdapterCode})";
@@ -33,11 +36,12 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
HasRecordings = true, AcceptsMetadataPush = true, HasAlarms = true 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; AdapterCode = adapterCode;
_http = http; _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(); public async Task InitializeAsync() => await _auth.GetTokenAsync();
@@ -52,9 +56,11 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
{ {
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
var resp = await client.GetAsync("/health"); var resp = await client.GetAsync("/health");
return resp.IsSuccessStatusCode; bool ok = resp.IsSuccessStatusCode;
_logger.LogDebug("[{Code}] 健康检查完成,状态码={Status}", AdapterCode, ok ? 200 : resp.StatusCode);
return ok;
} }
catch (Exception ex) { Console.Error.WriteLine($"[{AdapterCode}] HealthCheck 失败: {ex.Message}"); return false; } catch (Exception ex) { Console.Error.WriteLine($"[{AdapterCode}] 健康检查失败: {ex.Message}"); return false; }
} }
// ═══════════════════════════════════════════ // ═══════════════════════════════════════════
@@ -65,62 +71,71 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
{ {
await _limiter.WaitAsync(); await _limiter.WaitAsync();
var client = await _auth.GetAuthenticatedClientAsync(); 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)}"; if (!string.IsNullOrEmpty(keyword)) url += $"&key={Uri.EscapeDataString(keyword)}";
var json = await client.GetStringAsync(url); var json = await client.GetStringAsync(url);
var result = JsonSerializer.Deserialize<OwlPagedResult<OwlDeviceChannel>>(json)!; _logger.LogDebug("[{Code}] 获取设备列表({Url}),响应{Size}字节", AdapterCode, url, json.Length);
var result = JsonSerializer.Deserialize<OwlDevicesResponse>(json, JsonOpts)!;
var devices = new List<StandardDevice>(); var devices = new List<StandardDevice>();
var deviceItems = result.Items.Where(x => x.Type == "DEVICE").ToList(); foreach (var d in result.Items)
var channelItems = result.Items.Where(x => x.Type == "CHANNEL").ToList();
foreach (var d in deviceItems)
{ {
var childChannels = channelItems.Where(c => c.Did == d.Id).ToList(); var dev = MapDevice(d);
devices.Add(MapDevice(d, childChannels)); dev.AdapterCode = AdapterCode;
foreach (var ch in childChannels) devices.Add(dev);
devices.Add(MapChannel(ch, d.Id));
if (d.Children != null)
{
foreach (var ch in d.Children)
{
var chDev = MapChannel(ch, d.Id);
chDev.AdapterCode = AdapterCode;
devices.Add(chDev);
}
}
} }
return new PagedResult<StandardDevice> { Items = devices, Total = devices.Count }; return new PagedResult<StandardDevice> { Items = devices, Total = devices.Count };
} }
private static StandardDevice MapDevice(OwlDeviceChannel d, List<OwlDeviceChannel> channels) => new() private StandardDevice MapDevice(OwlDeviceItem d) => new()
{ {
SourceId = d.Id ?? "", SourceId = d.Id,
Name = d.Name ?? d.Id ?? "", Name = d.Name,
Category = "硬盘录像机", Category = "硬盘录像机",
Group = "视频设备", Group = "视频设备",
IsOnline = d.IsOnline == "1", IsOnline = d.IsOnline,
IsParent = true, IsParent = true,
IpAddress = d.Address, IpAddress = d.Address,
Port = int.TryParse(d.Port, out var p) ? p : null, Port = d.Port > 0 ? d.Port : null,
Extra = new Dictionary<string, object?> Extra = new Dictionary<string, object?>
{ {
["manufacturer"] = d.Manufacturer, ["manufacturer"] = d.Ext?.Manufacturer,
["model"] = d.Model, ["model"] = d.Ext?.Model,
["firmware"] = d.Firmware, ["firmware"] = d.Ext?.Firmware,
["longitude"] = d.Longitude, ["protocol"] = "GB28181",
["latitude"] = d.Latitude,
["protocol"] = d.Protocol ?? "GB28181",
["transport"] = d.Transport, ["transport"] = d.Transport,
["channelCount"] = d.ChannelCount ?? channels.Count ["channelCount"] = d.Children?.Count ?? d.Channels,
["deviceId"] = d.DeviceId,
["registeredAt"] = d.RegisteredAt
} }
}; };
private static StandardDevice MapChannel(OwlDeviceChannel ch, string? parentDeviceId) => new() private StandardDevice MapChannel(OwlChannelItem ch, string parentDeviceId) => new()
{ {
SourceId = ch.Id ?? "", SourceId = ch.Id,
Name = ch.Name ?? $"通道{ch.Id}", Name = ch.Name,
Category = "摄像机", Category = "摄像机",
Group = "视频设备", Group = "视频设备",
IsOnline = ch.IsOnline?.ToLower() == "true" || ch.IsOnline == "1", IsOnline = ch.IsOnline,
IsParent = false, IsParent = false,
ParentSourceId = parentDeviceId, ParentSourceId = parentDeviceId,
Extra = new Dictionary<string, object?> Extra = new Dictionary<string, object?>
{ {
["hasPtz"] = (ch.Ptztype ?? 0) > 0 ? "1" : "0", ["hasPtz"] = ch.Ptztype > 0 ? "1" : "0",
["app"] = ch.App, ["app"] = ch.App,
["streamId"] = ch.StreamId ["streamId"] = ch.Stream,
["channelId"] = ch.ChannelId,
["hasRecording"] = ch.HasRecording ? "1" : "0"
} }
}; };
@@ -135,7 +150,8 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
var resp = await client.PostAsync($"/channels/{channelId}/play", null); var resp = await client.PostAsync($"/channels/{channelId}/play", null);
resp.EnsureSuccessStatusCode(); resp.EnsureSuccessStatusCode();
var json = await resp.Content.ReadAsStringAsync(); var json = await resp.Content.ReadAsStringAsync();
var play = JsonSerializer.Deserialize<OwlPlayResponse>(json)!; _logger.LogDebug("[{Code}] 获取实时流地址({Id}),响应{Sz}字节", AdapterCode, channelId, json.Length);
var play = JsonSerializer.Deserialize<OwlPlayResponse>(json, JsonOpts)!;
return MapStreamUrls(play); return MapStreamUrls(play);
} }
@@ -147,6 +163,8 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
var startMs = new DateTimeOffset(start).ToUnixTimeMilliseconds(); var startMs = new DateTimeOffset(start).ToUnixTimeMilliseconds();
var endMs = new DateTimeOffset(end).ToUnixTimeMilliseconds(); var endMs = new DateTimeOffset(end).ToUnixTimeMilliseconds();
var baseUrl = (client.BaseAddress?.ToString() ?? "").TrimEnd('/'); var baseUrl = (client.BaseAddress?.ToString() ?? "").TrimEnd('/');
var playbackUrl = $"{baseUrl}/recordings/channels/{channelId}/index.m3u8?start_ms={startMs}&end_ms={endMs}&token=***";
_logger.LogDebug("[{Code}] 获取回放地址: {Url}", AdapterCode, playbackUrl);
return new StreamUrls return new StreamUrls
{ {
Hls = $"{baseUrl}/recordings/channels/{channelId}/index.m3u8?start_ms={startMs}&end_ms={endMs}&token={token}" Hls = $"{baseUrl}/recordings/channels/{channelId}/index.m3u8?start_ms={startMs}&end_ms={endMs}&token={token}"
@@ -157,6 +175,7 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
{ {
await _limiter.WaitAsync(); await _limiter.WaitAsync();
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
_logger.LogDebug("[{Code}] 云台控制({Ch}): 方向={Dir}, 速度={Spd}", AdapterCode, channelId, direction, speed);
if (direction.StartsWith("preset_")) if (direction.StartsWith("preset_"))
{ {
var idx = int.Parse(direction.Replace("preset_", "")); var idx = int.Parse(direction.Replace("preset_", ""));
@@ -173,6 +192,7 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
{ {
await _limiter.WaitAsync(); await _limiter.WaitAsync();
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
_logger.LogDebug("[{Code}] 云台停止({Ch})", AdapterCode, channelId);
await client.PostAsJsonAsync($"/channels/{channelId}/ptz/control", new { action = "stop" }); await client.PostAsJsonAsync($"/channels/{channelId}/ptz/control", new { action = "stop" });
} }
@@ -183,6 +203,7 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
var resp = await client.PostAsync($"/channels/{channelId}/snapshot", var resp = await client.PostAsync($"/channels/{channelId}/snapshot",
new StringContent("{}", System.Text.Encoding.UTF8, "application/json")); new StringContent("{}", System.Text.Encoding.UTF8, "application/json"));
var json = await resp.Content.ReadAsStringAsync(); var json = await resp.Content.ReadAsStringAsync();
_logger.LogDebug("[{Code}] 实时截图({Ch}),响应{Sz}字节", AdapterCode, channelId, json.Length);
var snap = JsonSerializer.Deserialize<OwlSnapshotResponse>(json)!; var snap = JsonSerializer.Deserialize<OwlSnapshotResponse>(json)!;
return new StreamUrls { Hls = snap.Link }; return new StreamUrls { Hls = snap.Link };
} }
@@ -200,6 +221,7 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
var endMs = new DateTimeOffset(end).ToUnixTimeMilliseconds(); var endMs = new DateTimeOffset(end).ToUnixTimeMilliseconds();
var json = await client.GetStringAsync( var json = await client.GetStringAsync(
$"/recordings?cid={channelId}&start_ms={startMs}&end_ms={endMs}&page={page}&size={size}"); $"/recordings?cid={channelId}&start_ms={startMs}&end_ms={endMs}&page={page}&size={size}");
_logger.LogDebug("[{Code}] 录像查询({Ch}, {S}-{E}),响应{Sz}字节", AdapterCode, channelId, start, end, json.Length);
var owl = JsonSerializer.Deserialize<OwlPagedResult<OwlRecording>>(json)!; var owl = JsonSerializer.Deserialize<OwlPagedResult<OwlRecording>>(json)!;
return new PagedResult<StandardRecording> return new PagedResult<StandardRecording>
{ {
@@ -222,6 +244,7 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
var client = await _auth.GetAuthenticatedClientAsync(); var client = await _auth.GetAuthenticatedClientAsync();
var body = new Dictionary<string, object>(); var body = new Dictionary<string, object>();
if (changes.Name != null) body["name"] = changes.Name; if (changes.Name != null) body["name"] = changes.Name;
_logger.LogDebug("[{Code}] 元数据推送({Id}): 名称={Name}", AdapterCode, sourceDeviceId, changes.Name);
await client.PutAsJsonAsync($"/devices/{sourceDeviceId}", body); await client.PutAsJsonAsync($"/devices/{sourceDeviceId}", body);
return new MetadataPushResult { Success = true }; return new MetadataPushResult { Success = true };
} }
@@ -240,6 +263,7 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
var json = await client.GetStringAsync( var json = await client.GetStringAsync(
$"/events?page={page}&size={size}&start_ms={fromMs}&end_ms={toMs}"); $"/events?page={page}&size={size}&start_ms={fromMs}&end_ms={toMs}");
var result = JsonSerializer.Deserialize<OwlPagedResult<OwlAiEvent>>(json)!; var result = JsonSerializer.Deserialize<OwlPagedResult<OwlAiEvent>>(json)!;
_logger.LogDebug("[{Code}] AI事件查询(第{Pg}页,每页{Sz}条),响应{Bl}字节,共{Ct}条", AdapterCode, page, size, json.Length, result.Items?.Count ?? 0);
return new PagedResult<StandardAlarm> return new PagedResult<StandardAlarm>
{ {
Items = result.Items.Select(MapEventToAlarm).ToList(), Items = result.Items.Select(MapEventToAlarm).ToList(),
@@ -276,12 +300,15 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
switch (command) switch (command)
{ {
case "ai-enable": case "ai-enable":
_logger.LogDebug("[{Code}] 设备控制 {Id}: 启用AI检测", AdapterCode, sourceDeviceId);
await client.PostAsync($"/channels/{sourceDeviceId}/ai/enable", null); await client.PostAsync($"/channels/{sourceDeviceId}/ai/enable", null);
break; break;
case "ai-disable": case "ai-disable":
_logger.LogDebug("[{Code}] 设备控制 {Id}: 禁用AI检测", AdapterCode, sourceDeviceId);
await client.PostAsync($"/channels/{sourceDeviceId}/ai/disable", null); await client.PostAsync($"/channels/{sourceDeviceId}/ai/disable", null);
break; break;
case "zone-add": case "zone-add":
_logger.LogDebug("[{Code}] 设备控制 {Id}: 添加AI检测区域", AdapterCode, sourceDeviceId);
await client.PostAsJsonAsync($"/channels/{sourceDeviceId}/zones", parameters!); await client.PostAsJsonAsync($"/channels/{sourceDeviceId}/zones", parameters!);
break; break;
default: default:
@@ -289,7 +316,7 @@ public class OwlAdapter : IHasFlatDevices, IHasStreams, IHasRecordings, IAccepts
} }
return new ControlResult { Success = true }; return new ControlResult { Success = true };
} }
catch (Exception ex) { return new ControlResult { Success = false, Message = ex.Message }; } catch (Exception ex) { _logger.LogDebug("[{Code}] 设备控制失败: {Msg}", AdapterCode, ex.Message); return new ControlResult { Success = false, Message = ex.Message }; }
} }
// ═══════════════════════════════════════════ // ═══════════════════════════════════════════
@@ -3,6 +3,8 @@ using System.Text;
using System.Text.Json; using System.Text.Json;
using System.Net.Http.Json; using System.Net.Http.Json;
using Microsoft.Extensions.Logging;
namespace IntegrationGateway.Adapters.Owl; namespace IntegrationGateway.Adapters.Owl;
/// <summary> /// <summary>
@@ -21,6 +23,7 @@ public class OwlAuthHelper
private readonly string _baseUrl; private readonly string _baseUrl;
private readonly string _username; private readonly string _username;
private readonly string _password; private readonly string _password;
private readonly ILogger _logger;
/// <summary>缓存的 JWT Token</summary> /// <summary>缓存的 JWT Token</summary>
private string? _token; private string? _token;
/// <summary>Token 过期时间(UTC</summary> /// <summary>Token 过期时间(UTC</summary>
@@ -31,10 +34,11 @@ public class OwlAuthHelper
/// <param name="baseUrl">Owl 服务地址,如 http://localhost:15123</param> /// <param name="baseUrl">Owl 服务地址,如 http://localhost:15123</param>
/// <param name="username">Owl 登录用户名</param> /// <param name="username">Owl 登录用户名</param>
/// <param name="password">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('/'); _http = http; _baseUrl = baseUrl.TrimEnd('/');
_username = username; _password = password; _username = username; _password = password;
_logger = logger;
} }
/// <summary> /// <summary>
@@ -42,27 +46,36 @@ public class OwlAuthHelper
/// </summary> /// </summary>
public async Task<string> GetTokenAsync() public async Task<string> GetTokenAsync()
{ {
if (!string.IsNullOrEmpty(_token) && DateTime.UtcNow < _tokenExpiry) return _token; if (!string.IsNullOrEmpty(_token) && DateTime.UtcNow < _tokenExpiry)
{
_logger.LogDebug("Owl Token缓存命中(过期时间={Exp})", _tokenExpiry);
return _token;
}
// 第一步:获取 RSA 公钥 // 第一步:获取 RSA 公钥
var keyResp = await _http.GetStringAsync($"{_baseUrl}/login/key"); var keyResp = await _http.GetStringAsync($"{_baseUrl}/login/key");
_logger.LogDebug("Owl 获取登录公钥,响应{Len}字节", keyResp.Length);
var keyData = JsonSerializer.Deserialize<LoginKeyResponse>(keyResp); var keyData = JsonSerializer.Deserialize<LoginKeyResponse>(keyResp);
var publicKey = Encoding.UTF8.GetString(Convert.FromBase64String(keyData!.Key!)); var publicKey = Encoding.UTF8.GetString(Convert.FromBase64String(keyData!.key!));
_logger.LogDebug("Owl RSA公钥解码成功({Len}字符)", publicKey.Length);
// 第二步:RSA 加密凭据 // 第二步:RSA 加密凭据
using var rsa = RSA.Create(); using var rsa = RSA.Create();
rsa.ImportFromPem(publicKey); rsa.ImportFromPem(publicKey);
var plain = JsonSerializer.Serialize(new { username = _username, password = _password }); var plain = JsonSerializer.Serialize(new { username = _username, password = _password });
var encrypted = rsa.Encrypt(Encoding.UTF8.GetBytes(plain), RSAEncryptionPadding.Pkcs1); var encrypted = rsa.Encrypt(Encoding.UTF8.GetBytes(plain), RSAEncryptionPadding.OaepSHA256);
var payload = JsonSerializer.Serialize(new { data = Convert.ToBase64String(encrypted) }); var payload = JsonSerializer.Serialize(new { data = Convert.ToBase64String(encrypted) });
_logger.LogDebug("Owl 凭据RSA-OAEP加密完成({Len}字节)", encrypted.Length);
// 第三步:登录换取 Token // 第三步:登录换取 Token
var resp = await _http.PostAsync($"{_baseUrl}/login", var resp = await _http.PostAsync($"{_baseUrl}/login",
new StringContent(payload, Encoding.UTF8, "application/json")); new StringContent(payload, Encoding.UTF8, "application/json"));
_logger.LogDebug("Owl 登录请求完成,状态={Status}", resp.StatusCode);
resp.EnsureSuccessStatusCode(); resp.EnsureSuccessStatusCode();
var loginResult = await resp.Content.ReadFromJsonAsync<LoginResponse>(); var loginResult = await resp.Content.ReadFromJsonAsync<LoginResponse>();
_token = loginResult!.Token; _token = loginResult!.Token;
_tokenExpiry = DateTime.UtcNow.AddDays(2.5); // 保守设置,Owl 默认 3 天 _tokenExpiry = DateTime.UtcNow.AddDays(2.5);
_logger.LogDebug("Owl Token获取成功(用户={User}, 2.5天有效)", loginResult.User ?? _username); // 保守设置,Owl 默认 3 天
return _token; return _token;
} }
@@ -82,7 +95,7 @@ public class OwlAuthHelper
} }
/// <summary>登录密钥响应</summary> /// <summary>登录密钥响应</summary>
public class LoginKeyResponse { public string? Key { get; set; } } public class LoginKeyResponse { public string? key { get; set; } }
/// <summary>登录响应</summary> /// <summary>登录响应</summary>
public class LoginResponse { public string Token { get; set; } = ""; public string? User { get; set; } } public class LoginResponse { public string Token { get; set; } = ""; public string? User { get; set; } }
} }
@@ -8,7 +8,7 @@ namespace IntegrationGateway.Adapters.Owl;
// 通用 // 通用
// ═══════════════════════════════════════════ // ═══════════════════════════════════════════
/// <summary>Owl API 分页响应</summary> /// <summary>Owl API 通用分页响应(录像/AI事件用)</summary>
public class OwlPagedResult<T> public class OwlPagedResult<T>
{ {
public List<T> Items { get; set; } = new(); public List<T> Items { get; set; } = new();
@@ -16,50 +16,115 @@ public class OwlPagedResult<T>
} }
// ═══════════════════════════════════════════ // ═══════════════════════════════════════════
// 设备+通道联合模型 (GET /devices/channels) // 设备+通道联合接口 (GET /devices/channels) — 2026-06-07 实际返回结构
// ═══════════════════════════════════════════ // ═══════════════════════════════════════════
/// <summary>Owl 设备或通道(联合接口返回)</summary> /// <summary>/devices/channels 响应</summary>
public class OwlDeviceChannel public class OwlDevicesResponse
{
public List<OwlDeviceItem> Items { get; set; } = new();
public int Total { get; set; }
}
/// <summary>设备项(含子通道列表)</summary>
public class OwlDeviceItem
{
public string Id { get; set; } = "";
public string Type { get; set; } = "";
public string DeviceId { get; set; } = ""; // device_id
public string Name { get; set; } = "";
public string Transport { get; set; } = "";
public int StreamMode { get; set; } // stream_mode
public string? Ip { get; set; }
public int Port { get; set; }
public bool IsOnline { get; set; } // is_online — bool 非 string!
public string? RegisteredAt { get; set; }
public string? KeepaliveAt { get; set; }
public int Keepalives { get; set; }
public int Expires { get; set; }
public int Channels { get; set; }
public string? CreatedAt { get; set; }
public string? UpdatedAt { get; set; }
public string? Password { get; set; }
public string? Address { get; set; }
public OwlDeviceExt? Ext { get; set; }
public string? Username { get; set; }
public List<OwlChannelItem>? Children { get; set; }
}
/// <summary>通道项(嵌套在设备 children 数组中)</summary>
public class OwlChannelItem
{
public string Id { get; set; } = "";
public string Did { get; set; } = "";
public string DeviceId { get; set; } = ""; // device_id
public string ChannelId { get; set; } = ""; // channel_id
public string Name { get; set; } = "";
public int Ptztype { get; set; }
public bool IsOnline { get; set; }
public bool IsPlaying { get; set; } // is_playing
public OwlDeviceExt? Ext { get; set; }
public string? CreatedAt { get; set; }
public string? UpdatedAt { get; set; }
public string Type { get; set; } = "";
public string? App { get; set; }
public string? Stream { get; set; }
public OwlChannelConfig? Config { get; set; }
public bool HasRecording { get; set; } // has_recording
}
/// <summary>设备/通道扩展信息(嵌套在 ext 字段中)</summary>
public class OwlDeviceExt
{ {
public string? Id { get; set; }
public string? Type { get; set; } // "DEVICE" | "CHANNEL"
public string? Did { get; set; } // 通道所属设备 ID
public string? Name { get; set; }
public string? IsOnline { get; set; } // DEVICE: "1"/"0", CHANNEL: true/false
public string? Manufacturer { get; set; } public string? Manufacturer { get; set; }
public string? Model { get; set; } public string? Model { get; set; }
public string? Firmware { get; set; } public string? Firmware { get; set; }
public string? Longitude { get; set; } public string? Name { get; set; }
public string? Latitude { get; set; } public string? GbVersion { get; set; } // gb_version
public int? ChannelCount { get; set; } public object? Zones { get; set; }
public int? Ptztype { get; set; } // 0=无云台, 1=方向, 2=预置位 public bool EnabledAi { get; set; } // enabled_ai
public string? Protocol { get; set; } public string? RecordMode { get; set; } // record_mode
public string? Address { get; set; } }
public string? Port { get; set; }
public string? Transport { get; set; } /// <summary>通道推送/流配置</summary>
public string? App { get; set; } // 流应用名 public class OwlChannelConfig
public string? StreamId { get; set; } // 流ID {
public string? Status { get; set; } public bool IsAuthDisabled { get; set; }
public string? RegisterWay { get; set; } public object? PushedAt { get; set; }
public object? StoppedAt { get; set; }
public string? MediaServerId { get; set; }
public string? SourceUrl { get; set; }
public int Transport { get; set; }
public int TimeoutS { get; set; }
public bool EnabledAudio { get; set; }
public bool EnabledRemoveNoneReader { get; set; }
public bool EnabledDisabledNoneReader { get; set; }
public string? StreamKey { get; set; }
public bool Enabled { get; set; }
} }
// ═══════════════════════════════════════════ // ═══════════════════════════════════════════
// 播放/流 // 播放/流
// ═══════════════════════════════════════════ // ═══════════════════════════════════════════
/// <summary>Owl 播放响应</summary> /// <summary>Owl 播放响应 (POST /channels/{id}/play)</summary>
public class OwlPlayResponse public class OwlPlayResponse
{ {
public string? App { get; set; }
public string? Stream { get; set; }
public List<OwlPlayItem>? Items { get; set; } public List<OwlPlayItem>? Items { get; set; }
} }
/// <summary>Owl 播放流条目</summary> /// <summary>Owl 播放流条目 (snake_case)</summary>
public class OwlPlayItem public class OwlPlayItem
{ {
public string? Label { get; set; }
[System.Text.Json.Serialization.JsonPropertyName("ws_flv")]
public string? WsFlv { get; set; } public string? WsFlv { get; set; }
[System.Text.Json.Serialization.JsonPropertyName("http_flv")]
public string? HttpFlv { get; set; } public string? HttpFlv { get; set; }
public string? Hls { get; set; } public string? Hls { get; set; }
[System.Text.Json.Serialization.JsonPropertyName("webrtc")]
public string? WebRtc { get; set; } public string? WebRtc { get; set; }
public string? Rtmp { get; set; } public string? Rtmp { get; set; }
public string? Rtsp { get; set; } public string? Rtsp { get; set; }
@@ -46,7 +46,7 @@ public class AdapterRegistry
public IGatewayAdapter? FindByCode(string adapterCode) public IGatewayAdapter? FindByCode(string adapterCode)
=> _adapters.FirstOrDefault(a => a.AdapterCode == adapterCode); => _adapters.FirstOrDefault(a => a.AdapterCode == adapterCode);
/// <summary>获取所有在线适配器</summary> /// <summary>获取所有已注册适配器(含在线和离线)</summary>
public IReadOnlyList<IGatewayAdapter> GetOnlineAdapters() public IReadOnlyList<IGatewayAdapter> GetAllAdapters()
=> _adapters.AsReadOnly(); => _adapters.AsReadOnly();
} }
@@ -9,6 +9,7 @@ namespace IntegrationGateway.Core.Infrastructure;
/// </summary> /// </summary>
public class GatewayClientFactory public class GatewayClientFactory
{ {
private static readonly JsonSerializerOptions JsonOpts = new() { PropertyNamingPolicy = JsonNamingPolicy.CamelCase };
private readonly IHttpClientFactory _httpFactory; private readonly IHttpClientFactory _httpFactory;
private readonly string _volProBaseUrl; private readonly string _volProBaseUrl;
@@ -48,7 +49,7 @@ public class GatewayClientFactory
{ {
var http = CreateClient(); var http = CreateClient();
var resp = await http.PostAsJsonAsync($"{_volProBaseUrl}/api/gateway/sync/devices", var resp = await http.PostAsJsonAsync($"{_volProBaseUrl}/api/gateway/sync/devices",
new { nodeCode, token, devices }); new { nodeCode, token, devices }, JsonOpts);
if (!resp.IsSuccessStatusCode) return null; if (!resp.IsSuccessStatusCode) return null;
return await resp.Content.ReadFromJsonAsync<JsonDocument>(); return await resp.Content.ReadFromJsonAsync<JsonDocument>();
} }
@@ -58,7 +59,7 @@ public class GatewayClientFactory
{ {
var http = CreateClient(); var http = CreateClient();
var resp = await http.PostAsJsonAsync($"{_volProBaseUrl}/api/gateway/sync/alarms", var resp = await http.PostAsJsonAsync($"{_volProBaseUrl}/api/gateway/sync/alarms",
new { nodeCode, token, alarms }); new { nodeCode, token, alarms }, JsonOpts);
if (!resp.IsSuccessStatusCode) return null; if (!resp.IsSuccessStatusCode) return null;
return await resp.Content.ReadFromJsonAsync<JsonDocument>(); return await resp.Content.ReadFromJsonAsync<JsonDocument>();
} }
@@ -4,13 +4,12 @@ namespace IntegrationGateway.Core.Infrastructure;
/// 令牌桶限流器。控制对第三方子系统的请求频率,防止超出 API 配额。 /// 令牌桶限流器。控制对第三方子系统的请求频率,防止超出 API 配额。
/// 每个适配器实例持有独立的限流器。 /// 每个适配器实例持有独立的限流器。
/// ///
/// 算法:启动时桶内有 tokensPerSecond 个令牌,每次请求消耗一个令牌, /// 使用单个后台 PeriodicTimer 持续补充令牌,避免每次 WaitAsync 创建新 Task。
/// 令牌按 (1000/tokensPerSecond) 毫秒的速率补充。
/// </summary> /// </summary>
public class RateLimiter public class RateLimiter
{ {
private readonly SemaphoreSlim _semaphore; private readonly SemaphoreSlim _semaphore;
private readonly int _intervalMs; private readonly int _maxTokens;
/// <summary> /// <summary>
/// 创建限流器 /// 创建限流器
@@ -18,22 +17,23 @@ public class RateLimiter
/// <param name="tokensPerSecond">每秒允许的请求数(QPS</param> /// <param name="tokensPerSecond">每秒允许的请求数(QPS</param>
public RateLimiter(int tokensPerSecond) public RateLimiter(int tokensPerSecond)
{ {
_maxTokens = tokensPerSecond;
_semaphore = new SemaphoreSlim(tokensPerSecond, 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>
/// 等待获取一个令牌。如果当前没有可用令牌,阻塞直到有令牌被释放。 /// 等待获取一个令牌。如果当前没有可用令牌,阻塞直到有令牌被释放。
/// </summary> /// </summary>
/// <param name="ct">取消令牌</param>
public async Task WaitAsync(CancellationToken ct = default) public async Task WaitAsync(CancellationToken ct = default)
{ => await _semaphore.WaitAsync(ct);
await _semaphore.WaitAsync(ct);
// 在后台任务中延迟补充令牌
_ = Task.Run(async () =>
{
await Task.Delay(_intervalMs, ct);
try { _semaphore.Release(); } catch { }
}, ct);
}
} }
@@ -15,12 +15,14 @@ public class AdapterCapabilities
/// <summary>是否支持视频取流</summary> /// <summary>是否支持视频取流</summary>
public bool HasStreams { get; set; } public bool HasStreams { get; set; }
/// <summary>是否支持云台控制(PTZ</summary> /// <summary>是否支持云台控制(PTZ</summary>
[Obsolete("Use HasStreams for PTZ capability check")]
public bool HasPtz { get; set; } public bool HasPtz { get; set; }
/// <summary>是否支持录像回放</summary> /// <summary>是否支持录像回放</summary>
public bool HasRecordings { get; set; } public bool HasRecordings { get; set; }
/// <summary>是否支持告警查询与处理</summary> /// <summary>是否支持告警查询与处理</summary>
public bool HasAlarms { get; set; } public bool HasAlarms { get; set; }
/// <summary>是否接受反向控制(点位写值)</summary> /// <summary>是否接受反向控制(点位写值)</summary>
[Obsolete("Check IAcceptsControl interface directly")]
public bool AcceptsControl { get; set; } public bool AcceptsControl { get; set; }
/// <summary>是否接受元数据回写(如设备改名)</summary> /// <summary>是否接受元数据回写(如设备改名)</summary>
public bool AcceptsMetadataPush { get; set; } public bool AcceptsMetadataPush { get; set; }
@@ -8,7 +8,7 @@ namespace IntegrationGateway.Core.Models;
public class StandardDevice public class StandardDevice
{ {
/// <summary>Vol.Pro 侧主键(同步后由 Vol.Pro 回填)</summary> /// <summary>Vol.Pro 侧主键(同步后由 Vol.Pro 回填)</summary>
public int DeviceId { get; set; } public int? DeviceId { get; set; }
/// <summary>来源适配器标识,格式 "类型:实例",如 "Owl:main"</summary> /// <summary>来源适配器标识,格式 "类型:实例",如 "Owl:main"</summary>
public string AdapterCode { get; set; } = ""; public string AdapterCode { get; set; } = "";
/// <summary>子系统原始设备 IDGB28181 编码 / MC4 sid</summary> /// <summary>子系统原始设备 IDGB28181 编码 / MC4 sid</summary>
+23 -8
View File
@@ -1,6 +1,8 @@
using IntegrationGateway.Core.Abstractions; using IntegrationGateway.Core.Abstractions;
using IntegrationGateway.Core.Infrastructure; using IntegrationGateway.Core.Infrastructure;
using IntegrationGateway.Core.Models; using IntegrationGateway.Core.Models;
using Microsoft.AspNetCore.Mvc;
using System.Net;
// ═══════════════════════════════════════════════════════════════ // ═══════════════════════════════════════════════════════════════
// IntegrationGateway 宿主启动程序 // IntegrationGateway 宿主启动程序
@@ -56,7 +58,8 @@ foreach (var o in owlList)
var code = $"Owl:{o.InstanceName ?? "default"}"; var code = $"Owl:{o.InstanceName ?? "default"}";
var a = new IntegrationGateway.Adapters.Owl.OwlAdapter(code, var a = new IntegrationGateway.Adapters.Owl.OwlAdapter(code,
app.Services.GetRequiredService<IHttpClientFactory>().CreateClient("VolPro"), 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); registry.Register(a);
} }
@@ -67,7 +70,8 @@ foreach (var k in kmsList)
var code = $"KMS:{k.InstanceName ?? "default"}"; var code = $"KMS:{k.InstanceName ?? "default"}";
var a = new IntegrationGateway.Adapters.Kms.KmsAdapter(code, var a = new IntegrationGateway.Adapters.Kms.KmsAdapter(code,
app.Services.GetRequiredService<IHttpClientFactory>().CreateClient("VolPro"), 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); registry.Register(a);
} }
@@ -78,7 +82,8 @@ foreach (var m in mc4List)
var code = $"MC4:{m.InstanceName ?? "default"}"; var code = $"MC4:{m.InstanceName ?? "default"}";
var a = new IntegrationGateway.Adapters.MC4.Mc4Adapter(code, var a = new IntegrationGateway.Adapters.MC4.Mc4Adapter(code,
app.Services.GetRequiredService<IHttpClientFactory>().CreateClient("VolPro"), 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); registry.Register(a);
} }
@@ -93,7 +98,16 @@ Console.WriteLine($"[Gateway] {registry.All.Count} 个适配器已注册: {adapt
var nodeCode = gwCfg["NodeCode"] ?? "gw-default"; var nodeCode = gwCfg["NodeCode"] ?? "gw-default";
var nodeToken = Environment.GetEnvironmentVariable("SECMPS_GATEWAY_TOKEN") ?? gwCfg["NodeToken"] ?? ""; var nodeToken = Environment.GetEnvironmentVariable("SECMPS_GATEWAY_TOKEN") ?? gwCfg["NodeToken"] ?? "";
var port = app.Urls.FirstOrDefault()?.Split(':').LastOrDefault() ?? "5100"; var port = app.Urls.FirstOrDefault()?.Split(':').LastOrDefault() ?? "5100";
var selfUrl = gwCfg["SelfUrl"] ?? $"http://localhost:{port}"; var selfUrl = gwCfg["SelfUrl"];
if (string.IsNullOrEmpty(selfUrl))
{
// 自动检测本机局域网 IP,避免注册 localhost 导致前端不可达
var host = Dns.GetHostEntry(Dns.GetHostName());
var lanIp = host.AddressList
.FirstOrDefault(ip => ip.AddressFamily == System.Net.Sockets.AddressFamily.InterNetwork
&& !IPAddress.IsLoopback(ip));
selfUrl = $"http://{(lanIp != null ? lanIp.ToString() : "localhost")}:{port}";
}
try try
{ {
var registerReq = new GatewayRegisterRequest var registerReq = new GatewayRegisterRequest
@@ -168,8 +182,9 @@ async Task SyncAllDevicesAsync(string nc, string nt, string baseUrl)
FlattenTree(allDevices, nodes, adapter.AdapterCode, null); FlattenTree(allDevices, nodes, adapter.AdapterCode, null);
} }
} }
catch { } catch (Exception ex) { Console.Error.WriteLine($"[Gateway] A3: 适配器 {adapter.AdapterCode} 取设备失败: {ex.Message}"); }
} }
Console.WriteLine($"[Gateway] A3: 收集到 {allDevices.Count} 台设备,开始推送到 Vol.Pro");
if (allDevices.Any()) if (allDevices.Any())
await clientFactory.SyncDevicesAsync(nc, nt, allDevices); await clientFactory.SyncDevicesAsync(nc, nt, allDevices);
} }
@@ -194,7 +209,7 @@ app.MapGet("/api/gateway/health", async () =>
foreach (var a in registry.All) foreach (var a in registry.All)
{ {
bool healthy = false; bool healthy = false;
try { healthy = await a.HealthCheckAsync(); } catch { } try { healthy = await a.HealthCheckAsync(); } catch (Exception ex) { Console.Error.WriteLine($"[Gateway] B1: HealthCheck {a.AdapterCode} 失败: {ex.Message}"); }
results.Add(new { a.AdapterCode, a.DisplayName, Healthy = healthy, a.Capabilities }); results.Add(new { a.AdapterCode, a.DisplayName, Healthy = healthy, a.Capabilities });
} }
return Results.Ok(results); return Results.Ok(results);
@@ -276,7 +291,7 @@ app.MapPost("/api/gateway/realtime/{adapter}/batch", async (string adapter, Batc
var results = new Dictionary<string, List<PointValue>>(); var results = new Dictionary<string, List<PointValue>>();
foreach (var deviceId in req.DeviceIds ?? new()) foreach (var deviceId in req.DeviceIds ?? new())
try { results[deviceId] = await a.GetRealtimeValuesAsync(deviceId); } catch { } try { results[deviceId] = await a.GetRealtimeValuesAsync(deviceId); } catch (Exception ex) { Console.Error.WriteLine($"[Gateway] B4-batch: {adapter}/{deviceId} 取实时值失败: {ex.Message}"); }
return Results.Ok(results); return Results.Ok(results);
}); });
@@ -349,7 +364,7 @@ app.MapPost("/api/gateway/sync/{adapter}", async (string adapter, SyncRequest re
}); });
// B13: 数据删除 — 从子系统删除数据 // B13: 数据删除 — 从子系统删除数据
app.MapDelete("/api/gateway/sync/{adapter}", async (string adapter, SyncDeleteRequest req) => app.MapDelete("/api/gateway/sync/{adapter}", async (string adapter, [FromBody] SyncDeleteRequest req) =>
{ {
var a = registry.FindByCode<IAcceptsDataSync>(adapter); var a = registry.FindByCode<IAcceptsDataSync>(adapter);
if (a == null) return Results.NotFound(new { error = "CAPABILITY_NOT_SUPPORTED" }); if (a == null) return Results.NotFound(new { error = "CAPABILITY_NOT_SUPPORTED" });
@@ -13,7 +13,7 @@
"commandName": "Project", "commandName": "Project",
"dotnetRunMessages": true, "dotnetRunMessages": true,
"launchBrowser": true, "launchBrowser": true,
"applicationUrl": "http://localhost:5260", "applicationUrl": "http://0.0.0.0:5100",
"environmentVariables": { "environmentVariables": {
"ASPNETCORE_ENVIRONMENT": "Development" "ASPNETCORE_ENVIRONMENT": "Development"
} }
@@ -2,27 +2,49 @@
"Logging": { "Logging": {
"LogLevel": { "LogLevel": {
"Default": "Information", "Default": "Information",
"Microsoft.AspNetCore": "Warning" "Microsoft.AspNetCore": "Warning",
"IntegrationGateway": "Debug"
} }
}, },
"Owl": [ "Owl": [
{ {
"InstanceName": "main", "InstanceName": "main",
"BaseUrl": "http://localhost:15123", "BaseUrl": "http://192.168.3.108:15123",
"Username": "admin", "Username": "admin",
"Password": "your_owl_password" "Password": "admin"
} }
], ],
"MC4": [ "MC4": [
{ {
"InstanceName": "31ku", "InstanceName": "31ku",
"BaseUrl": "http://localhost:3000" "BaseUrl": "http://192.168.3.90:3000",
"Username": "admin",
"Password": "@ZLcx8980"
},
{
"InstanceName": "32ku",
"BaseUrl": "http://192.168.3.91:3000",
"Username": "admin",
"Password": "@ZLcx8980"
},
{
"InstanceName": "33ku",
"BaseUrl": "http://192.168.3.92:3000/api/central/auth/conf/get",
"Username": "admin",
"Password": "@ZLcx8980"
},
{
"InstanceName": "34ku",
"BaseUrl": "http://192.168.3.93:3000",
"Username": "admin",
"Password": "@ZLcx8980"
} }
], ],
"Gateway": { "Gateway": {
"VolProBaseUrl": "http://localhost:9100", "VolProBaseUrl": "http://192.168.3.108:9100",
"NodeCode": "gw-31ku", "NodeCode": "gw-test",
"NodeToken": "changeme", "NodeToken": "changeme",
"SelfUrl": "http://192.168.3.108:5100",
"HeartbeatIntervalSec": 15, "HeartbeatIntervalSec": 15,
"AdapterInitTimeoutSec": 30, "AdapterInitTimeoutSec": 30,
"GatewayKey": null, "GatewayKey": null,
@@ -31,9 +53,9 @@
"KMS": [ "KMS": [
{ {
"InstanceName": "main", "InstanceName": "main",
"BaseUrl": "http://192.168.1.50:8080", "BaseUrl": "http://192.168.3.198",
"ClientId": "your_client_id", "ClientId": "g82tt",
"ClientSecret": "your_client_secret" "ClientSecret": "2w1q821130@W!Q"
} }
] ]
} }