新增原料入场记录功能,包含免密接口和数据同步,更新相关控制器、实体和服务,支持条码/批次号生成及管理,优化用户体验和系统实时数据处理能力。
This commit is contained in:
@@ -15,6 +15,7 @@ public class CategorySyncCoordinator : ISingletonDependency
|
||||
public CategorySyncCoordinator(
|
||||
IEventAggregator eventAggregator,
|
||||
IJeecgCategorySyncService categorySyncService,
|
||||
SyncPollManager pollManager,
|
||||
ILoggerService logger)
|
||||
{
|
||||
_eventAggregator = eventAggregator;
|
||||
@@ -26,6 +27,8 @@ public class CategorySyncCoordinator : ISingletonDependency
|
||||
_eventAggregator.GetEvent<NetworkStatusChangedEvent>()
|
||||
.Subscribe(OnNetworkStatusChanged, ThreadOption.BackgroundThread);
|
||||
|
||||
pollManager.Register("分类字典", () => SyncAndPublishAsync("poll", null));
|
||||
|
||||
_logger.Information("[分类字典] CategorySyncCoordinator 已启动");
|
||||
_ = Task.Run(() => SyncAndPublishAsync("startup", null));
|
||||
}
|
||||
|
||||
@@ -11,7 +11,10 @@ public class CustomerSyncCoordinator : ISingletonDependency
|
||||
private readonly IEventAggregator _eventAggregator;
|
||||
private readonly ILoggerService _logger;
|
||||
|
||||
public CustomerSyncCoordinator(IEventAggregator eventAggregator, ILoggerService logger)
|
||||
public CustomerSyncCoordinator(
|
||||
IEventAggregator eventAggregator,
|
||||
SyncPollManager pollManager,
|
||||
ILoggerService logger)
|
||||
{
|
||||
_eventAggregator = eventAggregator;
|
||||
_logger = logger;
|
||||
@@ -19,6 +22,14 @@ public class CustomerSyncCoordinator : ISingletonDependency
|
||||
.Subscribe(OnRemoteCommand, ThreadOption.BackgroundThread);
|
||||
_eventAggregator.GetEvent<NetworkStatusChangedEvent>()
|
||||
.Subscribe(OnNetworkStatusChanged, ThreadOption.BackgroundThread);
|
||||
|
||||
pollManager.Register("客户", () =>
|
||||
{
|
||||
_eventAggregator.GetEvent<CustomerChangedEvent>()
|
||||
.Publish(new CustomerChangedPayload { Action = "poll" });
|
||||
return Task.CompletedTask;
|
||||
});
|
||||
|
||||
_logger.Information("[客户推送] CustomerSyncCoordinator 已启动");
|
||||
}
|
||||
|
||||
|
||||
@@ -15,6 +15,7 @@ public class DictSyncCoordinator : ISingletonDependency
|
||||
public DictSyncCoordinator(
|
||||
IEventAggregator eventAggregator,
|
||||
IJeecgDictSyncService dictSyncService,
|
||||
SyncPollManager pollManager,
|
||||
ILoggerService logger)
|
||||
{
|
||||
_eventAggregator = eventAggregator;
|
||||
@@ -26,6 +27,8 @@ public class DictSyncCoordinator : ISingletonDependency
|
||||
_eventAggregator.GetEvent<NetworkStatusChangedEvent>()
|
||||
.Subscribe(OnNetworkStatusChanged, ThreadOption.BackgroundThread);
|
||||
|
||||
pollManager.Register("数据字典", () => SyncAndPublishAsync("poll", null));
|
||||
|
||||
_logger.Information("[数据字典] DictSyncCoordinator 已启动");
|
||||
_ = Task.Run(() => SyncAndPublishAsync("startup", null));
|
||||
}
|
||||
|
||||
@@ -82,6 +82,7 @@ public class JeecgDictSyncService : IJeecgDictSyncService, ISingletonDependency
|
||||
const int pageSize = 500;
|
||||
var pageNo = 1;
|
||||
var synced = 0;
|
||||
var seenIds = new HashSet<string>(StringComparer.OrdinalIgnoreCase);
|
||||
|
||||
while (true)
|
||||
{
|
||||
@@ -111,6 +112,8 @@ public class JeecgDictSyncService : IJeecgDictSyncService, ISingletonDependency
|
||||
continue;
|
||||
}
|
||||
|
||||
seenIds.Add(id);
|
||||
|
||||
var existing = await _dbContext.Queryable<JeecgSysDictItem>()
|
||||
.ClearFilter()
|
||||
.Where(x => x.Id == id)
|
||||
@@ -161,6 +164,15 @@ public class JeecgDictSyncService : IJeecgDictSyncService, ISingletonDependency
|
||||
pageNo++;
|
||||
}
|
||||
|
||||
// 删除本地存在但后端已移除的字典项(如后端删除重建导致 ID 变化)
|
||||
if (seenIds.Count > 0)
|
||||
{
|
||||
var seenList = seenIds.ToList();
|
||||
await _dbContext.Deleteable<JeecgSysDictItem>()
|
||||
.Where(x => !seenList.Contains(x.Id))
|
||||
.ExecuteCommandAsync();
|
||||
}
|
||||
|
||||
return synced;
|
||||
}
|
||||
|
||||
|
||||
@@ -51,6 +51,7 @@ public class MixerMaterialService : IMixerMaterialService, ISingletonDependency
|
||||
}
|
||||
|
||||
private string BaseUrl => (_configuration.GetValue<string>("JeecgIntegration:BaseUrl") ?? "http://localhost:8080/jeecg-boot").TrimEnd('/');
|
||||
private int DefaultTenantId => (int?)_configuration.GetValue<long?>("JeecgIntegration:DefaultTenantId") ?? 1002;
|
||||
private HttpClient CreateClient() => _httpClientFactory.CreateClient("JeecgApi");
|
||||
|
||||
private async Task LoadCacheAsync()
|
||||
@@ -84,7 +85,8 @@ public class MixerMaterialService : IMixerMaterialService, ISingletonDependency
|
||||
|
||||
public async Task SyncFromRemoteAsync(CancellationToken ct = default)
|
||||
{
|
||||
if (!await _syncLock.WaitAsync(0, ct).ConfigureAwait(false)) return;
|
||||
// 等待当前正在进行的同步结束,再执行本次同步(WaitAsync(0) 会静默跳过,导致 STOMP 通知丢失)
|
||||
await _syncLock.WaitAsync(ct).ConfigureAwait(false);
|
||||
try
|
||||
{
|
||||
var all = new List<MesMixerMaterial>();
|
||||
@@ -96,10 +98,15 @@ public class MixerMaterialService : IMixerMaterialService, ISingletonDependency
|
||||
var qs = HttpUtility.ParseQueryString(string.Empty);
|
||||
qs["pageNo"] = pageNo.ToString();
|
||||
qs["pageSize"] = pageSize.ToString();
|
||||
// mes_mixer_material 不在多租户隔离表中,不传 tenantId 可拉取全部记录
|
||||
var url = $"{BaseUrl}/mes/material/mixerMaterial/anon/list?{qs}";
|
||||
using var client = CreateClient();
|
||||
var resp = await client.GetAsync(url, ct).ConfigureAwait(false);
|
||||
if (!resp.IsSuccessStatusCode) break;
|
||||
if (!resp.IsSuccessStatusCode)
|
||||
{
|
||||
_logger.Warning($"[密炼物料] 同步请求失败,状态码:{(int)resp.StatusCode}");
|
||||
break;
|
||||
}
|
||||
|
||||
var json = await resp.Content.ReadAsStringAsync(ct).ConfigureAwait(false);
|
||||
using var doc = JsonDocument.Parse(json);
|
||||
@@ -119,6 +126,10 @@ public class MixerMaterialService : IMixerMaterialService, ISingletonDependency
|
||||
await SaveCacheAsync(all).ConfigureAwait(false);
|
||||
_logger.Information($"[密炼物料] 同步完成,共 {all.Count} 条");
|
||||
}
|
||||
else
|
||||
{
|
||||
_logger.Warning("[密炼物料] 同步返回 0 条,可能是 tenantId 配置有误或网络异常,保留原缓存");
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
|
||||
@@ -3,6 +3,7 @@ using System.Text.Json;
|
||||
using YY.Admin.Core;
|
||||
using YY.Admin.Core.Events;
|
||||
using YY.Admin.Core.Services;
|
||||
using YY.Admin.Services.Service;
|
||||
|
||||
namespace YY.Admin.Services.Service.MixerMaterial;
|
||||
|
||||
@@ -15,6 +16,7 @@ public class MixerMaterialSyncCoordinator : ISingletonDependency
|
||||
public MixerMaterialSyncCoordinator(
|
||||
IEventAggregator eventAggregator,
|
||||
IMixerMaterialService mixerMaterialService,
|
||||
SyncPollManager pollManager,
|
||||
ILoggerService logger)
|
||||
{
|
||||
_eventAggregator = eventAggregator;
|
||||
@@ -26,6 +28,13 @@ public class MixerMaterialSyncCoordinator : ISingletonDependency
|
||||
_eventAggregator.GetEvent<NetworkStatusChangedEvent>()
|
||||
.Subscribe(OnNetworkStatusChanged, ThreadOption.BackgroundThread);
|
||||
|
||||
pollManager.Register("密炼物料", async () =>
|
||||
{
|
||||
await _mixerMaterialService.SyncFromRemoteAsync().ConfigureAwait(false);
|
||||
_eventAggregator.GetEvent<MixerMaterialChangedEvent>()
|
||||
.Publish(new MixerMaterialChangedPayload { Action = "poll" });
|
||||
});
|
||||
|
||||
_logger.Information("[密炼物料推送] MixerMaterialSyncCoordinator 已启动");
|
||||
_ = _mixerMaterialService.SyncFromRemoteAsync();
|
||||
}
|
||||
|
||||
@@ -0,0 +1,626 @@
|
||||
using Microsoft.Extensions.Configuration;
|
||||
using System.Net.Http;
|
||||
using System.IO;
|
||||
using System.Text;
|
||||
using System.Text.Json;
|
||||
using System.Text.Json.Serialization;
|
||||
using System.Web;
|
||||
using Prism.Events;
|
||||
using YY.Admin.Core;
|
||||
using YY.Admin.Core.Entity;
|
||||
using YY.Admin.Core.Events;
|
||||
using YY.Admin.Core.Services;
|
||||
|
||||
namespace YY.Admin.Services.Service.RawMaterialEntry;
|
||||
|
||||
public class RawMaterialEntryService : IRawMaterialEntryService, ISingletonDependency
|
||||
{
|
||||
private readonly IHttpClientFactory _httpClientFactory;
|
||||
private readonly IConfiguration _configuration;
|
||||
private readonly INetworkMonitor _networkMonitor;
|
||||
private readonly IEventAggregator _eventAggregator;
|
||||
private readonly ILoggerService _logger;
|
||||
private readonly SemaphoreSlim _syncLock = new(1, 1);
|
||||
private readonly object _cacheLock = new();
|
||||
private readonly string _pendingOpsFilePath;
|
||||
private readonly string _cacheFilePath;
|
||||
private List<EntryPendingOperation> _pendingOps = new();
|
||||
private List<MesXslRawMaterialEntry> _localCache = new();
|
||||
|
||||
private static readonly JsonSerializerOptions _jsonOpts = new()
|
||||
{
|
||||
PropertyNameCaseInsensitive = true,
|
||||
PropertyNamingPolicy = JsonNamingPolicy.CamelCase,
|
||||
Converters = { new NullableDateTimeJsonConverter() }
|
||||
};
|
||||
|
||||
public RawMaterialEntryService(
|
||||
IHttpClientFactory httpClientFactory,
|
||||
IConfiguration configuration,
|
||||
INetworkMonitor networkMonitor,
|
||||
IEventAggregator eventAggregator,
|
||||
ILoggerService logger)
|
||||
{
|
||||
_httpClientFactory = httpClientFactory;
|
||||
_configuration = configuration;
|
||||
_networkMonitor = networkMonitor;
|
||||
_eventAggregator = eventAggregator;
|
||||
_logger = logger;
|
||||
|
||||
var appDataDir = Path.Combine(
|
||||
Environment.GetFolderPath(Environment.SpecialFolder.LocalApplicationData),
|
||||
"YY.Admin", "sync-cache");
|
||||
Directory.CreateDirectory(appDataDir);
|
||||
_pendingOpsFilePath = Path.Combine(appDataDir, "mes-xsl-raw-material-entry-pending-ops.json");
|
||||
_cacheFilePath = Path.Combine(appDataDir, "mes-xsl-raw-material-entry-cache.json");
|
||||
|
||||
LoadPendingOpsFromDisk();
|
||||
LoadCacheFromDisk();
|
||||
_logger.Information($"[原料入场] 服务初始化完成,缓存={_localCache.Count},待上传={_pendingOps.Count},在线={_networkMonitor.IsOnline}");
|
||||
|
||||
_networkMonitor.StatusChanged += OnNetworkStatusChanged;
|
||||
if (_networkMonitor.IsOnline)
|
||||
{
|
||||
_ = Task.Run(() => SyncAfterReconnectAsync(CancellationToken.None));
|
||||
}
|
||||
}
|
||||
|
||||
private const int MaxPendingRetries = 5;
|
||||
private string BaseUrl => (_configuration.GetValue<string>("JeecgIntegration:BaseUrl") ?? "http://localhost:8080/jeecg-boot").TrimEnd('/');
|
||||
private int DefaultTenantId => (int?)_configuration.GetValue<long?>("JeecgIntegration:DefaultTenantId") ?? 1002;
|
||||
|
||||
private HttpClient CreateClient() => _httpClientFactory.CreateClient("JeecgApi");
|
||||
|
||||
public async Task<RawMaterialEntryPageResult> PageAsync(
|
||||
int pageNo, int pageSize,
|
||||
string? barcode = null, string? batchNo = null, string? billNo = null,
|
||||
string? materialName = null, string? supplierName = null,
|
||||
CancellationToken ct = default)
|
||||
{
|
||||
List<MesXslRawMaterialEntry>? source = null;
|
||||
if (_networkMonitor.IsOnline)
|
||||
{
|
||||
try
|
||||
{
|
||||
source = await FetchRemoteListAsync(ct).ConfigureAwait(false);
|
||||
lock (_cacheLock)
|
||||
{
|
||||
_localCache = source.Select(Clone).ToList();
|
||||
SaveCacheToDiskUnsafe();
|
||||
}
|
||||
_logger.Information($"[原料入场列表] 远端拉取成功 count={source.Count}");
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
source = null;
|
||||
_logger.Warning($"[原料入场列表] 远端拉取失败,回退缓存:{ex.Message}");
|
||||
}
|
||||
}
|
||||
|
||||
lock (_cacheLock)
|
||||
{
|
||||
source ??= _localCache.Select(Clone).ToList();
|
||||
source = ApplyPendingOpsSnapshotUnsafe(source);
|
||||
}
|
||||
|
||||
var filtered = ApplyFilters(source, barcode, batchNo, billNo, materialName, supplierName);
|
||||
var total = filtered.Count;
|
||||
var pageRecords = filtered.Skip(Math.Max(0, (pageNo - 1) * pageSize)).Take(pageSize).ToList();
|
||||
return new RawMaterialEntryPageResult(pageRecords, total, pageNo, pageSize);
|
||||
}
|
||||
|
||||
public async Task<MesXslRawMaterialEntry?> GetByIdAsync(string id, CancellationToken ct = default)
|
||||
{
|
||||
if (_networkMonitor.IsOnline)
|
||||
{
|
||||
try
|
||||
{
|
||||
var url = $"{BaseUrl}/xslmes/mesXslRawMaterialEntry/anon/queryById?id={Uri.EscapeDataString(id)}&tenantId={DefaultTenantId}";
|
||||
using var client = CreateClient();
|
||||
var resp = await client.GetAsync(url, ct).ConfigureAwait(false);
|
||||
if (!resp.IsSuccessStatusCode) return null;
|
||||
var json = await resp.Content.ReadAsStringAsync(ct).ConfigureAwait(false);
|
||||
using var doc = JsonDocument.Parse(json);
|
||||
if (!doc.RootElement.TryGetProperty("result", out var resultEl)) return null;
|
||||
return resultEl.Deserialize<MesXslRawMaterialEntry>(_jsonOpts);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.Warning($"[原料入场详情] 远端查询失败,回退缓存 id={id}:{ex.Message}");
|
||||
}
|
||||
}
|
||||
|
||||
lock (_cacheLock)
|
||||
{
|
||||
return _localCache.FirstOrDefault(e => string.Equals(e.Id, id, StringComparison.OrdinalIgnoreCase)) is { } found
|
||||
? Clone(found) : null;
|
||||
}
|
||||
}
|
||||
|
||||
public async Task<bool> AddAsync(MesXslRawMaterialEntry entry, CancellationToken ct = default)
|
||||
{
|
||||
if (!entry.TenantId.HasValue || entry.TenantId.Value <= 0) entry.TenantId = DefaultTenantId;
|
||||
var local = Clone(entry);
|
||||
if (string.IsNullOrWhiteSpace(local.Id)) local.Id = $"local-{Guid.NewGuid():N}";
|
||||
|
||||
if (_networkMonitor.IsOnline)
|
||||
{
|
||||
try
|
||||
{
|
||||
var ok = await RemoteAddAsync(local, ct).ConfigureAwait(false);
|
||||
if (ok) { UpsertLocalCache(local); return true; }
|
||||
return false;
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.Warning($"[原料入场新增] 远端失败,转离线 id={local.Id}:{ex.Message}");
|
||||
}
|
||||
}
|
||||
|
||||
EnqueuePendingOperation(new EntryPendingOperation { OpType = EntryOpType.Add, EntryId = local.Id, Entry = local, CreatedAt = DateTime.UtcNow });
|
||||
UpsertLocalCache(local);
|
||||
return true;
|
||||
}
|
||||
|
||||
public async Task<bool> EditAsync(MesXslRawMaterialEntry entry, CancellationToken ct = default)
|
||||
{
|
||||
if (!entry.TenantId.HasValue || entry.TenantId.Value <= 0) entry.TenantId = DefaultTenantId;
|
||||
var local = Clone(entry);
|
||||
|
||||
if (_networkMonitor.IsOnline)
|
||||
{
|
||||
try
|
||||
{
|
||||
var (ok, _) = await RemoteEditAsync(local, ct).ConfigureAwait(false);
|
||||
if (ok) { UpsertLocalCache(local); return true; }
|
||||
return false;
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.Warning($"[原料入场修改] 远端失败,转离线 id={local.Id}:{ex.Message}");
|
||||
}
|
||||
}
|
||||
|
||||
EnqueuePendingOperation(new EntryPendingOperation { OpType = EntryOpType.Edit, EntryId = local.Id, Entry = local, AnchorUpdateTime = local.UpdateTime, CreatedAt = DateTime.UtcNow });
|
||||
UpsertLocalCache(local);
|
||||
return true;
|
||||
}
|
||||
|
||||
public async Task<bool> DeleteAsync(string id, CancellationToken ct = default)
|
||||
{
|
||||
if (_networkMonitor.IsOnline)
|
||||
{
|
||||
try
|
||||
{
|
||||
var ok = await RemoteDeleteAsync(id, ct).ConfigureAwait(false);
|
||||
if (ok) { RemoveFromLocalCache(id); return true; }
|
||||
return false;
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.Warning($"[原料入场删除] 远端失败,转离线 id={id}:{ex.Message}");
|
||||
}
|
||||
}
|
||||
|
||||
DateTime? anchor;
|
||||
lock (_cacheLock) { anchor = _localCache.FirstOrDefault(e => string.Equals(e.Id, id, StringComparison.OrdinalIgnoreCase))?.UpdateTime; }
|
||||
EnqueuePendingOperation(new EntryPendingOperation { OpType = EntryOpType.Delete, EntryId = id, AnchorUpdateTime = anchor, CreatedAt = DateTime.UtcNow });
|
||||
RemoveFromLocalCache(id);
|
||||
return true;
|
||||
}
|
||||
|
||||
public async Task<bool> DeleteBatchAsync(string ids, CancellationToken ct = default)
|
||||
{
|
||||
var idList = ids.Split(',', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries);
|
||||
var allSuccess = true;
|
||||
foreach (var id in idList) allSuccess &= await DeleteAsync(id, ct).ConfigureAwait(false);
|
||||
return allSuccess;
|
||||
}
|
||||
|
||||
public async Task<string?> GenerateBarcodeAsync(string materialCode, CancellationToken ct = default)
|
||||
{
|
||||
var url = $"{BaseUrl}/xslmes/mesXslRawMaterialEntry/anon/generateBarcode?materialCode={Uri.EscapeDataString(materialCode ?? "")}";
|
||||
try
|
||||
{
|
||||
using var client = CreateClient();
|
||||
var resp = await client.GetAsync(url, ct).ConfigureAwait(false);
|
||||
resp.EnsureSuccessStatusCode();
|
||||
var json = await resp.Content.ReadAsStringAsync(ct).ConfigureAwait(false);
|
||||
using var doc = JsonDocument.Parse(json);
|
||||
if (doc.RootElement.TryGetProperty("result", out var el))
|
||||
return el.GetString();
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.Warning($"[原料入场] 生成条码失败: {ex.Message}");
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
// ─── Remote ────────────────────────────────────────────────────────────────
|
||||
|
||||
private async Task<List<MesXslRawMaterialEntry>> FetchRemoteListAsync(CancellationToken ct)
|
||||
{
|
||||
var query = HttpUtility.ParseQueryString(string.Empty);
|
||||
query["pageNo"] = "1";
|
||||
query["pageSize"] = "10000";
|
||||
query["tenantId"] = DefaultTenantId.ToString();
|
||||
var url = $"{BaseUrl}/xslmes/mesXslRawMaterialEntry/anon/list?{query}";
|
||||
using var client = CreateClient();
|
||||
var resp = await client.GetAsync(url, ct).ConfigureAwait(false);
|
||||
resp.EnsureSuccessStatusCode();
|
||||
var json = await resp.Content.ReadAsStringAsync(ct).ConfigureAwait(false);
|
||||
using var doc = JsonDocument.Parse(json);
|
||||
var result = doc.RootElement.GetProperty("result");
|
||||
return result.GetProperty("records").Deserialize<List<MesXslRawMaterialEntry>>(_jsonOpts) ?? new();
|
||||
}
|
||||
|
||||
private async Task<bool> RemoteAddAsync(MesXslRawMaterialEntry entry, CancellationToken ct)
|
||||
{
|
||||
var url = $"{BaseUrl}/xslmes/mesXslRawMaterialEntry/anon/add?tenantId={DefaultTenantId}";
|
||||
var payload = Clone(entry);
|
||||
if (IsLocalTempId(payload.Id)) payload.Id = null;
|
||||
return await PostJsonAsync(url, payload, ct).ConfigureAwait(false);
|
||||
}
|
||||
|
||||
private async Task<(bool Ok, bool IsVersionConflict)> RemoteEditAsync(MesXslRawMaterialEntry entry, CancellationToken ct)
|
||||
{
|
||||
var url = $"{BaseUrl}/xslmes/mesXslRawMaterialEntry/anon/edit?tenantId={DefaultTenantId}";
|
||||
return await PostJsonCheckVersionAsync(url, entry, ct).ConfigureAwait(false);
|
||||
}
|
||||
|
||||
private async Task<bool> RemoteDeleteAsync(string id, CancellationToken ct)
|
||||
{
|
||||
var url = $"{BaseUrl}/xslmes/mesXslRawMaterialEntry/anon/delete?id={Uri.EscapeDataString(id)}&tenantId={DefaultTenantId}";
|
||||
using var client = CreateClient();
|
||||
var resp = await client.DeleteAsync(url, ct).ConfigureAwait(false);
|
||||
return resp.IsSuccessStatusCode && await IsSuccessResultAsync(resp, ct).ConfigureAwait(false);
|
||||
}
|
||||
|
||||
private async Task<bool> PostJsonAsync(string url, object body, CancellationToken ct)
|
||||
{
|
||||
var content = new StringContent(JsonSerializer.Serialize(body, _jsonOpts), Encoding.UTF8, "application/json");
|
||||
using var client = CreateClient();
|
||||
var resp = await client.PostAsync(url, content, ct).ConfigureAwait(false);
|
||||
return resp.IsSuccessStatusCode && await IsSuccessResultAsync(resp, ct).ConfigureAwait(false);
|
||||
}
|
||||
|
||||
private async Task<(bool Ok, bool IsVersionConflict)> PostJsonCheckVersionAsync(string url, object body, CancellationToken ct)
|
||||
{
|
||||
var content = new StringContent(JsonSerializer.Serialize(body, _jsonOpts), Encoding.UTF8, "application/json");
|
||||
using var client = CreateClient();
|
||||
var resp = await client.PostAsync(url, content, ct).ConfigureAwait(false);
|
||||
if (!resp.IsSuccessStatusCode) return (false, false);
|
||||
try
|
||||
{
|
||||
var json = await resp.Content.ReadAsStringAsync(ct).ConfigureAwait(false);
|
||||
using var doc = JsonDocument.Parse(json);
|
||||
int code = 200;
|
||||
if (doc.RootElement.TryGetProperty("code", out var codeEl)) code = codeEl.GetInt32();
|
||||
if (code == 200) return (true, false);
|
||||
if (doc.RootElement.TryGetProperty("message", out var msgEl))
|
||||
if ((msgEl.GetString() ?? "").Contains("已被他人修改")) return (false, true);
|
||||
return (false, false);
|
||||
}
|
||||
catch { return (true, false); }
|
||||
}
|
||||
|
||||
private static async Task<bool> IsSuccessResultAsync(HttpResponseMessage resp, CancellationToken ct)
|
||||
{
|
||||
try
|
||||
{
|
||||
var json = await resp.Content.ReadAsStringAsync(ct).ConfigureAwait(false);
|
||||
using var doc = JsonDocument.Parse(json);
|
||||
if (doc.RootElement.TryGetProperty("code", out var code)) return code.GetInt32() == 200;
|
||||
if (doc.RootElement.TryGetProperty("success", out var success)) return success.GetBoolean();
|
||||
return true;
|
||||
}
|
||||
catch { return true; }
|
||||
}
|
||||
|
||||
// ─── Network / Reconnect ───────────────────────────────────────────────────
|
||||
|
||||
private void OnNetworkStatusChanged(bool isOnline)
|
||||
{
|
||||
if (!isOnline) return;
|
||||
_ = Task.Run(() => SyncAfterReconnectAsync(CancellationToken.None));
|
||||
}
|
||||
|
||||
private async Task SyncAfterReconnectAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
var pushResult = await PushPendingOnReconnectAsync(cancellationToken).ConfigureAwait(false);
|
||||
if (!_networkMonitor.IsOnline) return;
|
||||
try
|
||||
{
|
||||
var remote = await FetchRemoteListAsync(cancellationToken).ConfigureAwait(false);
|
||||
lock (_cacheLock) { _localCache = remote.Select(Clone).ToList(); SaveCacheToDiskUnsafe(); }
|
||||
_eventAggregator.GetEvent<RawMaterialEntryChangedEvent>().Publish(new RawMaterialEntryChangedPayload { Action = "pull" });
|
||||
_logger.Information($"[原料入场重连] 全量回拉成功 count={remote.Count}");
|
||||
}
|
||||
catch (Exception ex) { _logger.Warning($"[原料入场重连] 全量回拉失败:{ex.Message}"); }
|
||||
|
||||
if (pushResult.PushedCount > 0 || pushResult.ConflictCount > 0 || pushResult.NewRecordsPushed > 0)
|
||||
{
|
||||
_eventAggregator.GetEvent<SyncConflictEvent>().Publish(new SyncConflictPayload
|
||||
{
|
||||
EntityName = "原料入场",
|
||||
PushedCount = pushResult.PushedCount,
|
||||
ConflictCount = pushResult.ConflictCount,
|
||||
NewRecordsPushed = pushResult.NewRecordsPushed
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
private sealed record PendingReplayResult(bool Ok, bool IsConflict, string? EntityId);
|
||||
|
||||
private async Task<PushPendingResult> PushPendingOnReconnectAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
if (!await _syncLock.WaitAsync(0, cancellationToken).ConfigureAwait(false))
|
||||
return new PushPendingResult(0, 0, 0);
|
||||
try
|
||||
{
|
||||
List<EntryPendingOperation> snapshot;
|
||||
lock (_cacheLock) { snapshot = _pendingOps.OrderBy(x => x.CreatedAt).ToList(); }
|
||||
|
||||
int pushed = 0, conflicts = 0, newPushed = 0;
|
||||
foreach (var op in snapshot)
|
||||
{
|
||||
if (!_networkMonitor.IsOnline) break;
|
||||
lock (_cacheLock) { if (!_pendingOps.Any(x => x.Id == op.Id)) continue; }
|
||||
|
||||
var result = await ExecutePendingOperationAsync(op, cancellationToken).ConfigureAwait(false);
|
||||
if (!result.Ok)
|
||||
{
|
||||
lock (_cacheLock)
|
||||
{
|
||||
op.RetryCount++;
|
||||
if (op.RetryCount >= MaxPendingRetries) _pendingOps.RemoveAll(x => x.Id == op.Id);
|
||||
SavePendingOpsToDiskUnsafe();
|
||||
}
|
||||
break;
|
||||
}
|
||||
if (result.IsConflict)
|
||||
{
|
||||
conflicts++;
|
||||
if (!string.IsNullOrWhiteSpace(result.EntityId)) RemovePendingOpsByEntryId(result.EntityId!);
|
||||
continue;
|
||||
}
|
||||
lock (_cacheLock)
|
||||
{
|
||||
if (op.OpType == EntryOpType.Add) newPushed++; else pushed++;
|
||||
_pendingOps.RemoveAll(x => x.Id == op.Id);
|
||||
SavePendingOpsToDiskUnsafe();
|
||||
}
|
||||
}
|
||||
return new PushPendingResult(pushed, conflicts, newPushed);
|
||||
}
|
||||
finally { _syncLock.Release(); }
|
||||
}
|
||||
|
||||
private async Task<PendingReplayResult> ExecutePendingOperationAsync(EntryPendingOperation op, CancellationToken ct)
|
||||
{
|
||||
try
|
||||
{
|
||||
switch (op.OpType)
|
||||
{
|
||||
case EntryOpType.Add:
|
||||
var ok = op.Entry != null && await RemoteAddAsync(op.Entry, ct).ConfigureAwait(false);
|
||||
return ok ? new PendingReplayResult(true, false, op.EntryId) : new PendingReplayResult(false, false, null);
|
||||
|
||||
case EntryOpType.Edit:
|
||||
if (op.Entry == null || string.IsNullOrWhiteSpace(op.Entry.Id))
|
||||
return new PendingReplayResult(false, false, null);
|
||||
var id = op.Entry.Id;
|
||||
var remote = await FetchRemoteSingleAsync(id, ct).ConfigureAwait(false);
|
||||
if (remote != null && op.AnchorUpdateTime != null && remote.UpdateTime != op.AnchorUpdateTime)
|
||||
{
|
||||
UpsertLocalCache(remote);
|
||||
return new PendingReplayResult(true, true, id);
|
||||
}
|
||||
var (editOk, isConflict) = await RemoteEditAsync(op.Entry, ct).ConfigureAwait(false);
|
||||
if (isConflict) { var fresh = await FetchRemoteSingleAsync(id, ct).ConfigureAwait(false); if (fresh != null) UpsertLocalCache(fresh); return new PendingReplayResult(true, true, id); }
|
||||
return editOk ? new PendingReplayResult(true, false, id) : new PendingReplayResult(false, false, null);
|
||||
|
||||
case EntryOpType.Delete:
|
||||
if (string.IsNullOrWhiteSpace(op.EntryId)) return new PendingReplayResult(false, false, null);
|
||||
var delId = op.EntryId!;
|
||||
var delRemote = await FetchRemoteSingleAsync(delId, ct).ConfigureAwait(false);
|
||||
if (delRemote == null) return new PendingReplayResult(true, false, delId);
|
||||
if (op.AnchorUpdateTime != null && delRemote.UpdateTime != op.AnchorUpdateTime) { UpsertLocalCache(delRemote); return new PendingReplayResult(true, true, delId); }
|
||||
var delOk = await RemoteDeleteAsync(delId, ct).ConfigureAwait(false);
|
||||
return delOk ? new PendingReplayResult(true, false, delId) : new PendingReplayResult(false, false, null);
|
||||
|
||||
default:
|
||||
return new PendingReplayResult(true, false, null);
|
||||
}
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.Warning($"[原料入场推送] 执行异常 op={op.OpType},entryId={op.EntryId}:{ex.Message}");
|
||||
return new PendingReplayResult(false, false, null);
|
||||
}
|
||||
}
|
||||
|
||||
private async Task<MesXslRawMaterialEntry?> FetchRemoteSingleAsync(string id, CancellationToken ct)
|
||||
{
|
||||
try
|
||||
{
|
||||
var url = $"{BaseUrl}/xslmes/mesXslRawMaterialEntry/anon/queryById?id={Uri.EscapeDataString(id)}&tenantId={DefaultTenantId}";
|
||||
using var client = CreateClient();
|
||||
var resp = await client.GetAsync(url, ct).ConfigureAwait(false);
|
||||
if (!resp.IsSuccessStatusCode) return null;
|
||||
var json = await resp.Content.ReadAsStringAsync(ct).ConfigureAwait(false);
|
||||
using var doc = JsonDocument.Parse(json);
|
||||
if (doc.RootElement.TryGetProperty("result", out var resultEl))
|
||||
return resultEl.Deserialize<MesXslRawMaterialEntry>(_jsonOpts);
|
||||
return null;
|
||||
}
|
||||
catch { return null; }
|
||||
}
|
||||
|
||||
// ─── Local cache helpers ───────────────────────────────────────────────────
|
||||
|
||||
private static List<MesXslRawMaterialEntry> ApplyFilters(
|
||||
List<MesXslRawMaterialEntry> source,
|
||||
string? barcode, string? batchNo, string? billNo,
|
||||
string? materialName, string? supplierName)
|
||||
{
|
||||
IEnumerable<MesXslRawMaterialEntry> q = source;
|
||||
if (!string.IsNullOrWhiteSpace(barcode))
|
||||
q = q.Where(e => (e.Barcode ?? "").Contains(barcode, StringComparison.OrdinalIgnoreCase));
|
||||
if (!string.IsNullOrWhiteSpace(batchNo))
|
||||
q = q.Where(e => (e.BatchNo ?? "").Contains(batchNo, StringComparison.OrdinalIgnoreCase));
|
||||
if (!string.IsNullOrWhiteSpace(billNo))
|
||||
q = q.Where(e => (e.BillNo ?? "").Contains(billNo, StringComparison.OrdinalIgnoreCase));
|
||||
if (!string.IsNullOrWhiteSpace(materialName))
|
||||
q = q.Where(e => (e.MaterialName ?? "").Contains(materialName, StringComparison.OrdinalIgnoreCase));
|
||||
if (!string.IsNullOrWhiteSpace(supplierName))
|
||||
q = q.Where(e => (e.SupplierName ?? "").Contains(supplierName, StringComparison.OrdinalIgnoreCase));
|
||||
return q.OrderByDescending(e => e.EntryTime ?? e.CreateTime ?? DateTime.MinValue).ToList();
|
||||
}
|
||||
|
||||
private List<MesXslRawMaterialEntry> ApplyPendingOpsSnapshotUnsafe(List<MesXslRawMaterialEntry> source)
|
||||
{
|
||||
var map = source.Where(e => !string.IsNullOrWhiteSpace(e.Id))
|
||||
.ToDictionary(e => e.Id!, Clone, StringComparer.OrdinalIgnoreCase);
|
||||
|
||||
foreach (var op in _pendingOps.OrderBy(x => x.CreatedAt))
|
||||
{
|
||||
switch (op.OpType)
|
||||
{
|
||||
case EntryOpType.Add:
|
||||
case EntryOpType.Edit:
|
||||
if (op.Entry != null && !string.IsNullOrWhiteSpace(op.Entry.Id))
|
||||
map[op.Entry.Id] = Clone(op.Entry);
|
||||
break;
|
||||
case EntryOpType.Delete:
|
||||
if (!string.IsNullOrWhiteSpace(op.EntryId))
|
||||
map.Remove(op.EntryId);
|
||||
break;
|
||||
}
|
||||
}
|
||||
return map.Values.ToList();
|
||||
}
|
||||
|
||||
private void EnqueuePendingOperation(EntryPendingOperation op)
|
||||
{
|
||||
lock (_cacheLock) { _pendingOps.Add(op); SavePendingOpsToDiskUnsafe(); }
|
||||
}
|
||||
|
||||
private void UpsertLocalCache(MesXslRawMaterialEntry entry)
|
||||
{
|
||||
lock (_cacheLock)
|
||||
{
|
||||
var idx = _localCache.FindIndex(e => string.Equals(e.Id, entry.Id, StringComparison.OrdinalIgnoreCase));
|
||||
if (idx >= 0) _localCache[idx] = Clone(entry);
|
||||
else _localCache.Insert(0, Clone(entry));
|
||||
SaveCacheToDiskUnsafe();
|
||||
}
|
||||
}
|
||||
|
||||
private void RemoveFromLocalCache(string id)
|
||||
{
|
||||
lock (_cacheLock) { _localCache.RemoveAll(e => string.Equals(e.Id, id, StringComparison.OrdinalIgnoreCase)); SaveCacheToDiskUnsafe(); }
|
||||
}
|
||||
|
||||
private void RemovePendingOpsByEntryId(string entryId)
|
||||
{
|
||||
lock (_cacheLock)
|
||||
{
|
||||
_pendingOps.RemoveAll(x =>
|
||||
(!string.IsNullOrWhiteSpace(x.EntryId) && string.Equals(x.EntryId, entryId, StringComparison.OrdinalIgnoreCase)) ||
|
||||
(x.Entry?.Id != null && string.Equals(x.Entry.Id, entryId, StringComparison.OrdinalIgnoreCase)));
|
||||
SavePendingOpsToDiskUnsafe();
|
||||
}
|
||||
}
|
||||
|
||||
private void LoadPendingOpsFromDisk()
|
||||
{
|
||||
try
|
||||
{
|
||||
if (!File.Exists(_pendingOpsFilePath)) return;
|
||||
var data = JsonSerializer.Deserialize<List<EntryPendingOperation>>(File.ReadAllText(_pendingOpsFilePath), _jsonOpts);
|
||||
_pendingOps = data ?? new();
|
||||
}
|
||||
catch { _pendingOps = new(); }
|
||||
}
|
||||
|
||||
private void LoadCacheFromDisk()
|
||||
{
|
||||
try
|
||||
{
|
||||
if (!File.Exists(_cacheFilePath)) return;
|
||||
var data = JsonSerializer.Deserialize<List<MesXslRawMaterialEntry>>(File.ReadAllText(_cacheFilePath), _jsonOpts);
|
||||
_localCache = data ?? new();
|
||||
}
|
||||
catch { _localCache = new(); }
|
||||
}
|
||||
|
||||
private void SavePendingOpsToDiskUnsafe() =>
|
||||
File.WriteAllText(_pendingOpsFilePath, JsonSerializer.Serialize(_pendingOps, _jsonOpts));
|
||||
|
||||
private void SaveCacheToDiskUnsafe() =>
|
||||
File.WriteAllText(_cacheFilePath, JsonSerializer.Serialize(_localCache, _jsonOpts));
|
||||
|
||||
private static MesXslRawMaterialEntry Clone(MesXslRawMaterialEntry e) => new()
|
||||
{
|
||||
Id = e.Id, Barcode = e.Barcode, BatchNo = e.BatchNo, EntryTime = e.EntryTime,
|
||||
WeightRecordId = e.WeightRecordId, BillNo = e.BillNo, MaterialId = e.MaterialId,
|
||||
MaterialName = e.MaterialName, SupplyCustomer = e.SupplyCustomer, SupplierId = e.SupplierId,
|
||||
SupplierName = e.SupplierName, ManufacturerMaterialName = e.ManufacturerMaterialName,
|
||||
ShelfLife = e.ShelfLife, TotalWeight = e.TotalWeight, TotalPortions = e.TotalPortions,
|
||||
PortionWeight = e.PortionWeight, PortionPackages = e.PortionPackages,
|
||||
TestResult = e.TestResult, TestStatus = e.TestStatus, PrintFlag = e.PrintFlag,
|
||||
StockBalance = e.StockBalance, WarehouseLocation = e.WarehouseLocation,
|
||||
UnloadOperator = e.UnloadOperator, IsSpecialAdoption = e.IsSpecialAdoption,
|
||||
SpecialAdoptionOperator = e.SpecialAdoptionOperator, SpecialAdoptionTime = e.SpecialAdoptionTime,
|
||||
SpecialAdoptionReason = e.SpecialAdoptionReason, Status = e.Status, Remark = e.Remark,
|
||||
CreateBy = e.CreateBy, CreateTime = e.CreateTime, UpdateBy = e.UpdateBy,
|
||||
UpdateTime = e.UpdateTime, TenantId = e.TenantId
|
||||
};
|
||||
|
||||
private static bool IsLocalTempId(string? id) =>
|
||||
!string.IsNullOrWhiteSpace(id) && id.StartsWith("local-", StringComparison.OrdinalIgnoreCase);
|
||||
|
||||
private sealed class EntryPendingOperation
|
||||
{
|
||||
public string Id { get; set; } = Guid.NewGuid().ToString("N");
|
||||
public EntryOpType OpType { get; set; }
|
||||
public string? EntryId { get; set; }
|
||||
public MesXslRawMaterialEntry? Entry { get; set; }
|
||||
public DateTime? AnchorUpdateTime { get; set; }
|
||||
public DateTime CreatedAt { get; set; } = DateTime.UtcNow;
|
||||
public int RetryCount { get; set; } = 0;
|
||||
}
|
||||
|
||||
private enum EntryOpType { Add = 1, Edit = 2, Delete = 3 }
|
||||
|
||||
private sealed class NullableDateTimeJsonConverter : JsonConverter<DateTime?>
|
||||
{
|
||||
private static readonly string[] SupportedFormats =
|
||||
[
|
||||
"yyyy-MM-dd HH:mm:ss", "yyyy-MM-dd HH:mm:ss.fff",
|
||||
"yyyy-MM-ddTHH:mm:ss", "yyyy-MM-ddTHH:mm:ss.fff",
|
||||
"yyyy-MM-ddTHH:mm:ssZ", "yyyy-MM-ddTHH:mm:ss.fffZ"
|
||||
];
|
||||
|
||||
public override DateTime? Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options)
|
||||
{
|
||||
if (reader.TokenType == JsonTokenType.Null) return null;
|
||||
if (reader.TokenType == JsonTokenType.String)
|
||||
{
|
||||
var raw = reader.GetString();
|
||||
if (string.IsNullOrWhiteSpace(raw)) return null;
|
||||
if (DateTime.TryParseExact(raw, SupportedFormats, System.Globalization.CultureInfo.InvariantCulture, System.Globalization.DateTimeStyles.AssumeLocal, out var exact)) return exact;
|
||||
if (DateTime.TryParse(raw, out var fallback)) return fallback;
|
||||
}
|
||||
throw new JsonException($"无法将 JSON 值转换为 DateTime?,token={reader.TokenType}");
|
||||
}
|
||||
|
||||
public override void Write(Utf8JsonWriter writer, DateTime? value, JsonSerializerOptions options)
|
||||
{
|
||||
if (value.HasValue) writer.WriteStringValue(value.Value.ToString("yyyy-MM-dd HH:mm:ss"));
|
||||
else writer.WriteNullValue();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,83 @@
|
||||
using Prism.Events;
|
||||
using System.Text.Json;
|
||||
using YY.Admin.Core;
|
||||
using YY.Admin.Core.Events;
|
||||
using YY.Admin.Core.Services;
|
||||
|
||||
namespace YY.Admin.Services.Service.RawMaterialEntry;
|
||||
|
||||
/// <summary>
|
||||
/// 监听 STOMP 收到的原料入场记录变更信号,转发为桌面端 Prism 事件,触发列表刷新。
|
||||
/// </summary>
|
||||
public class RawMaterialEntrySyncCoordinator : ISingletonDependency
|
||||
{
|
||||
private readonly IEventAggregator _eventAggregator;
|
||||
private readonly ILoggerService _logger;
|
||||
private SubscriptionToken? _remoteCommandToken;
|
||||
private SubscriptionToken? _networkStatusToken;
|
||||
|
||||
public RawMaterialEntrySyncCoordinator(
|
||||
IEventAggregator eventAggregator,
|
||||
SyncPollManager pollManager,
|
||||
ILoggerService logger)
|
||||
{
|
||||
_eventAggregator = eventAggregator;
|
||||
_logger = logger;
|
||||
_remoteCommandToken = _eventAggregator
|
||||
.GetEvent<RemoteCommandReceivedEvent>()
|
||||
.Subscribe(OnRemoteCommand, ThreadOption.BackgroundThread);
|
||||
_networkStatusToken = _eventAggregator
|
||||
.GetEvent<NetworkStatusChangedEvent>()
|
||||
.Subscribe(OnNetworkStatusChanged, ThreadOption.BackgroundThread);
|
||||
|
||||
pollManager.Register("原料入场", () =>
|
||||
{
|
||||
_eventAggregator.GetEvent<RawMaterialEntryChangedEvent>()
|
||||
.Publish(new RawMaterialEntryChangedPayload { Action = "poll" });
|
||||
return Task.CompletedTask;
|
||||
});
|
||||
|
||||
_logger.Information("[原料入场] RawMaterialEntrySyncCoordinator 已启动");
|
||||
}
|
||||
|
||||
private void OnNetworkStatusChanged(NetworkStatusChangedPayload payload)
|
||||
{
|
||||
if (!payload.IsOnline) return;
|
||||
_logger.Information("[原料入场] 网络恢复,触发补偿刷新");
|
||||
_eventAggregator.GetEvent<RawMaterialEntryChangedEvent>().Publish(new RawMaterialEntryChangedPayload { Action = "reconnect" });
|
||||
}
|
||||
|
||||
private void OnRemoteCommand(RemoteCommandPayload payload)
|
||||
{
|
||||
try
|
||||
{
|
||||
var json = payload.CommandJson ?? string.Empty;
|
||||
if (string.IsNullOrWhiteSpace(json)) return;
|
||||
|
||||
using var doc = JsonDocument.Parse(json);
|
||||
if (!doc.RootElement.TryGetProperty("cmd", out var cmdEl)) return;
|
||||
|
||||
var cmd = cmdEl.GetString() ?? string.Empty;
|
||||
if (!cmd.Equals("MES_RAW_MATERIAL_ENTRY_CHANGED", StringComparison.OrdinalIgnoreCase))
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
doc.RootElement.TryGetProperty("action", out var actionEl);
|
||||
doc.RootElement.TryGetProperty("entryId", out var idEl);
|
||||
|
||||
var changedPayload = new RawMaterialEntryChangedPayload
|
||||
{
|
||||
Action = actionEl.GetString() ?? string.Empty,
|
||||
EntryId = idEl.ValueKind == JsonValueKind.String ? idEl.GetString() : null,
|
||||
};
|
||||
|
||||
_logger.Information($"收到原料入场变更信号: action={changedPayload.Action}, entryId={changedPayload.EntryId}");
|
||||
_eventAggregator.GetEvent<RawMaterialEntryChangedEvent>().Publish(changedPayload);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.Warning($"处理 STOMP 原料入场变更信号失败: {ex.Message}");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -15,6 +15,7 @@ public class SupplierSyncCoordinator : ISingletonDependency
|
||||
public SupplierSyncCoordinator(
|
||||
IEventAggregator eventAggregator,
|
||||
ISupplierService supplierService,
|
||||
SyncPollManager pollManager,
|
||||
ILoggerService logger)
|
||||
{
|
||||
_eventAggregator = eventAggregator;
|
||||
@@ -24,6 +25,13 @@ public class SupplierSyncCoordinator : ISingletonDependency
|
||||
.Subscribe(OnRemoteCommand, ThreadOption.BackgroundThread);
|
||||
_eventAggregator.GetEvent<NetworkStatusChangedEvent>()
|
||||
.Subscribe(OnNetworkStatusChanged, ThreadOption.BackgroundThread);
|
||||
|
||||
pollManager.Register("供应商", () =>
|
||||
{
|
||||
_eventAggregator.GetEvent<SupplierChangedEvent>()
|
||||
.Publish(new SupplierChangedPayload { Action = "poll" });
|
||||
return Task.CompletedTask;
|
||||
});
|
||||
}
|
||||
|
||||
private async void OnNetworkStatusChanged(NetworkStatusChangedPayload payload)
|
||||
|
||||
58
yy-admin-master/YY.Admin.Services/Service/SyncPollManager.cs
Normal file
58
yy-admin-master/YY.Admin.Services/Service/SyncPollManager.cs
Normal file
@@ -0,0 +1,58 @@
|
||||
using YY.Admin.Core;
|
||||
|
||||
namespace YY.Admin.Services.Service;
|
||||
|
||||
/// <summary>
|
||||
/// 统一后台轮询管理器。
|
||||
/// 各模块通过 <see cref="Register"/> 注册轮询任务,定时器统一触发。
|
||||
/// ★ 只需修改 <see cref="PollInterval"/> 即可调整所有模块的轮询间隔。
|
||||
/// </summary>
|
||||
public class SyncPollManager : ISingletonDependency
|
||||
{
|
||||
// ★ 唯一轮询间隔配置 — 修改这里即可调整所有模块的轮询频率
|
||||
public static readonly TimeSpan PollInterval = TimeSpan.FromMinutes(5);
|
||||
|
||||
private readonly List<(string Name, Func<Task> Task)> _tasks = [];
|
||||
private readonly Timer _timer;
|
||||
private readonly ILoggerService _logger;
|
||||
|
||||
public SyncPollManager(ILoggerService logger)
|
||||
{
|
||||
_logger = logger;
|
||||
// 首次延迟一个完整周期,确保所有协调器启动后都已完成注册
|
||||
_timer = new Timer(OnTick, null, PollInterval, PollInterval);
|
||||
_logger.Information($"[轮询管理器] 已启动,间隔 {PollInterval.TotalMinutes} 分钟");
|
||||
}
|
||||
|
||||
/// <summary>注册一个需要定时执行的轮询任务</summary>
|
||||
/// <param name="name">任务名称,用于日志</param>
|
||||
/// <param name="pollTask">异步轮询委托</param>
|
||||
public void Register(string name, Func<Task> pollTask)
|
||||
{
|
||||
lock (_tasks)
|
||||
_tasks.Add((name, pollTask));
|
||||
_logger.Information($"[轮询管理器] 注册任务: {name}");
|
||||
}
|
||||
|
||||
private void OnTick(object? _)
|
||||
{
|
||||
_ = Task.Run(async () =>
|
||||
{
|
||||
List<(string Name, Func<Task> Task)> snapshot;
|
||||
lock (_tasks) snapshot = [.. _tasks];
|
||||
|
||||
_logger.Debug($"[轮询管理器] 触发,执行 {snapshot.Count} 个任务");
|
||||
foreach (var (name, task) in snapshot)
|
||||
{
|
||||
try
|
||||
{
|
||||
await task().ConfigureAwait(false);
|
||||
}
|
||||
catch (Exception ex)
|
||||
{
|
||||
_logger.Warning($"[轮询管理器] 任务 [{name}] 异常: {ex.Message}");
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -16,18 +16,28 @@ public class VehicleSyncCoordinator : ISingletonDependency
|
||||
private SubscriptionToken? _remoteCommandToken;
|
||||
private SubscriptionToken? _networkStatusToken;
|
||||
|
||||
public VehicleSyncCoordinator(IEventAggregator eventAggregator, ILoggerService logger)
|
||||
public VehicleSyncCoordinator(
|
||||
IEventAggregator eventAggregator,
|
||||
SyncPollManager pollManager,
|
||||
ILoggerService logger)
|
||||
{
|
||||
_eventAggregator = eventAggregator;
|
||||
_logger = logger;
|
||||
_remoteCommandToken = _eventAggregator
|
||||
.GetEvent<RemoteCommandReceivedEvent>()
|
||||
.Subscribe(OnRemoteCommand, ThreadOption.BackgroundThread);
|
||||
// 断线重连后补拉一次,覆盖离线期间漏掉的 STOMP 事件
|
||||
_networkStatusToken = _eventAggregator
|
||||
.GetEvent<NetworkStatusChangedEvent>()
|
||||
.Subscribe(OnNetworkStatusChanged, ThreadOption.BackgroundThread);
|
||||
_logger.Information("[车辆推送] VehicleSyncCoordinator 已启动,开始监听 RemoteCommandReceivedEvent");
|
||||
|
||||
pollManager.Register("车辆", () =>
|
||||
{
|
||||
_eventAggregator.GetEvent<VehicleChangedEvent>()
|
||||
.Publish(new VehicleChangedPayload { Action = "poll" });
|
||||
return Task.CompletedTask;
|
||||
});
|
||||
|
||||
_logger.Information("[车辆推送] VehicleSyncCoordinator 已启动");
|
||||
}
|
||||
|
||||
private void OnNetworkStatusChanged(NetworkStatusChangedPayload payload)
|
||||
|
||||
@@ -16,7 +16,10 @@ public class WeightRecordSyncCoordinator : ISingletonDependency
|
||||
private SubscriptionToken? _remoteCommandToken;
|
||||
private SubscriptionToken? _networkStatusToken;
|
||||
|
||||
public WeightRecordSyncCoordinator(IEventAggregator eventAggregator, ILoggerService logger)
|
||||
public WeightRecordSyncCoordinator(
|
||||
IEventAggregator eventAggregator,
|
||||
SyncPollManager pollManager,
|
||||
ILoggerService logger)
|
||||
{
|
||||
_eventAggregator = eventAggregator;
|
||||
_logger = logger;
|
||||
@@ -26,6 +29,14 @@ public class WeightRecordSyncCoordinator : ISingletonDependency
|
||||
_networkStatusToken = _eventAggregator
|
||||
.GetEvent<NetworkStatusChangedEvent>()
|
||||
.Subscribe(OnNetworkStatusChanged, ThreadOption.BackgroundThread);
|
||||
|
||||
pollManager.Register("磅单", () =>
|
||||
{
|
||||
_eventAggregator.GetEvent<MesXslWeightRecordChangedEvent>()
|
||||
.Publish(new MesXslWeightRecordChangedPayload { Action = "poll" });
|
||||
return Task.CompletedTask;
|
||||
});
|
||||
|
||||
_logger.Information("[磅单推送] WeightRecordSyncCoordinator 已启动");
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user