using System.Text.Json; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Options; using MiGu.Server.Auth; using MiGu.Server.Launcher; using MiGu.Server.Persistence; namespace MiGu.Server.Fleet; /// /// 单例:轮询 Simple3 车辆 + 每车状态,读取「车体_AlarmInfo/车体_AlarmLevel」并对帐进 platform.db(vehicle_alarms)。 /// 出现→开 active;文案变→更新;消失→置 cleared 并记录恢复时间/时长。永不删=完整历史;Simple3 离线仍可查最近记录。 /// public sealed class AlarmCollector { private readonly IServiceScopeFactory _scopeFactory; private readonly IHttpClientFactory _httpFactory; private readonly Simple3Options _sl; private readonly InternalTokenStore _token; private readonly ILogger _log; private readonly SemaphoreSlim _gate = new(1, 1); private static readonly JsonSerializerOptions JsonOpt = new() { PropertyNameCaseInsensitive = true }; public volatile bool Online; public DateTimeOffset? LastSyncAt { get; private set; } public AlarmCollector( IServiceScopeFactory scopeFactory, IHttpClientFactory httpFactory, IOptions sl, InternalTokenStore token, ILogger log) { _scopeFactory = scopeFactory; _httpFactory = httpFactory; _sl = sl.Value; _token = token; _log = log; } private sealed class CurrentAlarm { public int CarId; public string CarName = ""; public string Info = ""; public int Level; } public async Task SyncOnceAsync(CancellationToken ct) { if (!await _gate.WaitAsync(0, ct)) return Online; try { var cars = await FetchCarsAsync(ct); if (cars == null) { Online = false; return false; } var current = new Dictionary(); await Parallel.ForEachAsync( cars, new ParallelOptions { MaxDegreeOfParallelism = 8, CancellationToken = ct }, async (car, token) => { var (info, level) = await FetchCarAlarmAsync(car.Id, token); if (string.IsNullOrWhiteSpace(info)) return; lock (current) { current[car.Id] = new CurrentAlarm { CarId = car.Id, CarName = car.Name, Info = info, Level = level }; } }); try { using var scope = _scopeFactory.CreateScope(); var db = scope.ServiceProvider.GetRequiredService(); await ReconcileAsync(db, current, ct); } catch (Exception ex) { _log.LogDebug(ex, "alarm reconcile failed"); } Online = true; LastSyncAt = DateTimeOffset.UtcNow; return true; } finally { _gate.Release(); } } private sealed class CarRow { public int Id; public string Name = ""; } private async Task?> FetchCarsAsync(CancellationToken ct) { var port = _sl.ProjectionPort > 0 ? _sl.ProjectionPort : 8222; try { using var client = CreateClient(); using var resp = await client.SendAsync(Req($"http://127.0.0.1:{port}/projection/cars"), ct); if (!resp.IsSuccessStatusCode) return null; var text = await resp.Content.ReadAsStringAsync(ct); using var doc = JsonDocument.Parse(text); if (doc.RootElement.ValueKind != JsonValueKind.Array) return new(); var list = new List(); foreach (var el in doc.RootElement.EnumerateArray()) { var id = el.TryGetProperty("rawId", out var rid) && rid.TryGetInt32(out var n) ? n : 0; if (id <= 0) continue; var name = el.TryGetProperty("name", out var nm) ? nm.GetString() ?? "" : ""; list.Add(new CarRow { Id = id, Name = name }); } return list; } catch (Exception ex) { _log.LogDebug(ex, "alarm fetch cars failed"); return null; } } private async Task<(string info, int level)> FetchCarAlarmAsync(int carId, CancellationToken ct) { var port = _sl.ProjectionPort > 0 ? _sl.ProjectionPort : 8222; try { using var client = CreateClient(); using var resp = await client.SendAsync(Req($"http://127.0.0.1:{port}/projection/reflection/status/car/{carId}"), ct); if (!resp.IsSuccessStatusCode) return ("", 0); var text = await resp.Content.ReadAsStringAsync(ct); using var doc = JsonDocument.Parse(text); if (!doc.RootElement.TryGetProperty("data", out var data) || data.ValueKind != JsonValueKind.Array) return ("", 0); string info = ""; var level = 0; foreach (var kv in data.EnumerateArray()) { var key = kv.TryGetProperty("key", out var k) ? k.GetString() : null; var val = kv.TryGetProperty("value", out var v) ? v.GetString() : null; if (key == "车体_AlarmInfo" || key == "AlarmInfo") info = val ?? ""; else if (key == "车体_AlarmLevel" || key == "AlarmLevel") int.TryParse(val, out level); } info = info.Trim(); if (info is "0" or "/" or "-") info = ""; return (info, level); } catch { return ("", 0); } } private static async Task ReconcileAsync(MiGuDbContext db, Dictionary current, CancellationToken ct) { var now = DateTimeOffset.UtcNow; var active = await db.VehicleAlarms.Where(a => a.Status == "active").ToListAsync(ct); var activeByCar = new Dictionary(); foreach (var a in active) activeByCar[a.CarId] = a; // 每车取一条 active // 出现 / 更新:文案或级别变化时先 clear 旧记录再开新 active,保留分段历史 foreach (var cur in current.Values) { if (activeByCar.TryGetValue(cur.CarId, out var rec)) { var same = string.Equals(rec.Info, cur.Info, StringComparison.Ordinal) && rec.Level == cur.Level; if (same) { rec.CarName = cur.CarName; rec.LastAt = now; continue; } rec.Status = "cleared"; rec.ResolvedAt = now; rec.DurationSecs = (long)Math.Max(0, (now - rec.FirstAt).TotalSeconds); } db.VehicleAlarms.Add(new VehicleAlarmRecord { CarId = cur.CarId, CarName = cur.CarName, Info = cur.Info, Level = cur.Level, Status = "active", FirstAt = now, LastAt = now }); } // 消失 → 恢复 foreach (var rec in active) { if (current.ContainsKey(rec.CarId)) continue; rec.Status = "cleared"; rec.ResolvedAt = now; rec.DurationSecs = (long)Math.Max(0, (now - rec.FirstAt).TotalSeconds); } await db.SaveChangesAsync(ct); } private HttpClient CreateClient() { var c = _httpFactory.CreateClient(); c.Timeout = TimeSpan.FromSeconds(10); return c; } private HttpRequestMessage Req(string url) { var req = new HttpRequestMessage(HttpMethod.Get, url); var token = _token.Token; if (!string.IsNullOrEmpty(token)) req.Headers.TryAddWithoutValidation("X-Platform-Internal-Token", token); return req; } } /// 后台循环:定时采集车辆报警到 platform.db。 public sealed class AlarmCollectorService : BackgroundService { private readonly AlarmCollector _collector; private readonly ILogger _log; public AlarmCollectorService(AlarmCollector collector, ILogger log) { _collector = collector; _log = log; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { try { await Task.Delay(TimeSpan.FromSeconds(4), stoppingToken); } catch { return; } while (!stoppingToken.IsCancellationRequested) { try { await _collector.SyncOnceAsync(stoppingToken); } catch (Exception ex) { _log.LogDebug(ex, "alarm collector loop error"); } try { await Task.Delay(TimeSpan.FromSeconds(8), stoppingToken); } catch { break; } } } }