From d2c5be9d12d8c8342b799041fbea4669a68a97d0 Mon Sep 17 00:00:00 2001 From: Blank Date: Tue, 22 Sep 2026 15:37:38 +0800 Subject: [PATCH] =?UTF-8?q?fix(discovery):=20=E6=8B=86=E5=88=86=20ApplySta?= =?UTF-8?q?teTransitions=20=E9=99=8D=E4=BD=8E=20S3776=20=E8=AE=A4=E7=9F=A5?= =?UTF-8?q?=E5=A4=8D=E6=9D=82=E5=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 重构 MongoEndpointWatcher.ApplyStateTransitions,把首观测 / 已知观测五分支派发 + 路由表准入 + Evicted 兜底抽为 RecordTransition / RecordFirstObservation / EmitIncarnationChange / RecordStatusChange / EvictMissingInstances / CollectObservedInstanceIds / IncarnationChanged / RecoveredFromStale / StatusChanged / BecameStale / IsRoutable / IsFirstObservationSkipped / IsActiveOrDraining 共 13 个 private/private static 辅助方法;保持对外签名 / 行为 / 事件序列 / _knownInstances 与 _lastSeenIncarnations 更新顺序完全等价。 Linear: GFX-861 --- .../Discovery/MongoEndpointWatcher.cs | 256 +++++++++++++----- 1 file changed, 194 insertions(+), 62 deletions(-) diff --git a/GameFrameX.NetWork.RemoteMessaging/Discovery/MongoEndpointWatcher.cs b/GameFrameX.NetWork.RemoteMessaging/Discovery/MongoEndpointWatcher.cs index 6552a0da..3f43b3c2 100644 --- a/GameFrameX.NetWork.RemoteMessaging/Discovery/MongoEndpointWatcher.cs +++ b/GameFrameX.NetWork.RemoteMessaging/Discovery/MongoEndpointWatcher.cs @@ -338,7 +338,6 @@ private async Task PollOnceAsync(CancellationToken cancellationToken) private List ApplyStateTransitions(List documents, DateTime nowUtc, List> pendingEvents) { var liveInstances = new List(documents.Count); - var observedInstanceIds = new HashSet(StringComparer.Ordinal); foreach (var document in documents) { var descriptor = TryToDescriptor(document); @@ -347,82 +346,139 @@ private List ApplyStateTransitions(List _stalenessThreshold; _lastSeenIncarnations.TryGetValue(descriptor.InstanceId, out var lastIncarnation); + _knownInstances.TryGetValue(descriptor.InstanceId, out var known); - if (!_knownInstances.TryGetValue(descriptor.InstanceId, out var known)) - { - // 新出现的 instanceId:仅对具备路由资格(非 stale 且 Active/Draining,与路由表准入一致)的首次观测发事件; - // stale 或 Stopped/Booting 首次观测不进路由表,也不发 Online,避免订阅者看到表中不存在的实例。 - if (!isStale && (descriptor.Status == InstanceStatus.Active || descriptor.Status == InstanceStatus.Draining)) - { - if (lastIncarnation != default && lastIncarnation != descriptor.Incarnation) - { - // 曾在 graveyard 里见过且 incarnation 变化 → 重启语义(Offline+Online,D15 规则)。 - pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Offline, descriptor)); - pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Online, descriptor)); - } - else if (descriptor.Status == InstanceStatus.Draining) - { - // 首次观测即为 Draining:发 Draining(不接新流量、保留在途投递)。 - pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Draining, descriptor)); - } - else - { - pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Online, descriptor)); - } - } - } - else if (known.Descriptor.Incarnation != descriptor.Incarnation) - { - // 已知实例的 incarnation 变化:旧身份下线 + 新身份上线(D15 incarnation 规则)。 - pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Offline, known.Descriptor)); - pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Online, descriptor)); - } - else if (known.IsStale && !isStale) - { - // 同 incarnation 从陈旧恢复新鲜 → Recovered。 - pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Recovered, descriptor)); - } - else if (!known.IsStale && known.Descriptor.Status != descriptor.Status && !isStale) - { - // 状态跃迁:→ Draining 发 Draining;→ Active(自 Draining 恢复接流)发 Online;→ Stopped 发 Offline。 - if (descriptor.Status == InstanceStatus.Draining) - { - pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Draining, descriptor)); - } - else if (descriptor.Status == InstanceStatus.Active) - { - pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Online, descriptor)); - } - else if (descriptor.Status == InstanceStatus.Stopped) - { - pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Offline, descriptor)); - } - } - else if (!known.IsStale && isStale) - { - // 新鲜 → 陈旧:三周期阈值判死,摘出路由表。 - pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Offline, descriptor)); - } - + RecordTransition(known, lastIncarnation, descriptor, isStale, pendingEvents); _knownInstances[descriptor.InstanceId] = new KnownInstance(descriptor, isStale); _lastSeenIncarnations[descriptor.InstanceId] = descriptor.Incarnation; - if (!isStale && descriptor.Status != InstanceStatus.Stopped) + + if (IsRoutable(descriptor, isStale)) { liveInstances.Add(descriptor); } } - // 曾知实例本轮文档消失(TTL 已清除)→ Evicted,彻底移出观测。 + EvictMissingInstances(documents, pendingEvents); + return liveInstances; + } + + /// + /// 状态机派发:按优先级匹配首观测 / incarnation 变化 / 陈旧恢复 / 状态跃迁 / 新鲜→陈旧 五条事件路径。 + /// + /// + /// Dispatches the state machine. Each branch is a single-responsibility predicate + /// so the method stays under the S3776 threshold. + /// + private void RecordTransition(KnownInstance known, long lastIncarnation, InstanceDescriptor descriptor, bool isStale, List> pendingEvents) + { + if (known == null) + { + RecordFirstObservation(isStale, lastIncarnation, descriptor, pendingEvents); + } + else if (IncarnationChanged(known, descriptor)) + { + EmitIncarnationChange(known.Descriptor, descriptor, pendingEvents); + } + else if (RecoveredFromStale(known, isStale)) + { + // 同 incarnation 从陈旧恢复新鲜 → Recovered。 + pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Recovered, descriptor)); + } + else if (StatusChanged(known, descriptor, isStale)) + { + RecordStatusChange(descriptor, pendingEvents); + } + else if (BecameStale(known, isStale)) + { + // 新鲜 → 陈旧:三周期阈值判死,摘出路由表。 + pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Offline, descriptor)); + } + } + + /// + /// 首观测分支:仅对具备路由资格(非 stale 且 Active/Draining)的首次观测发事件。 + /// + /// + /// The first-observation branch: only routable first sightings emit events; stale or + /// non-routable first sightings stay silent so subscribers never see instances that + /// the route table would refuse. + /// + private static void RecordFirstObservation(bool isStale, long lastIncarnation, InstanceDescriptor descriptor, List> pendingEvents) + { + if (IsFirstObservationSkipped(isStale, descriptor.Status)) + { + return; + } + + if (lastIncarnation != default && lastIncarnation != descriptor.Incarnation) + { + // 曾在 graveyard 里见过且 incarnation 变化 → 重启语义(Offline+Online,D15 规则)。 + pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Offline, descriptor)); + pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Online, descriptor)); + } + else if (descriptor.Status == InstanceStatus.Draining) + { + // 首次观测即为 Draining:发 Draining(不接新流量、保留在途投递)。 + pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Draining, descriptor)); + } + else + { + pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Online, descriptor)); + } + } + + /// + /// 已知实例的 incarnation 变化:旧身份下线 + 新身份上线(D15 incarnation 规则)。 + /// + /// + /// Emits the incarnation-change pair (Offline the old identity, Online the new one). + /// + private static void EmitIncarnationChange(InstanceDescriptor oldDescriptor, InstanceDescriptor newDescriptor, List> pendingEvents) + { + pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Offline, oldDescriptor)); + pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Online, newDescriptor)); + } + + /// + /// 状态跃迁:→ Draining 发 Draining;→ Active(自 Draining 恢复接流)发 Online;→ Stopped 发 Offline。 + /// + /// + /// Maps the new status to the matching event kind; unknown / non-transitioning states + /// stay silent. + /// + private static void RecordStatusChange(InstanceDescriptor descriptor, List> pendingEvents) + { + RoleInstanceChangeKind? kind = descriptor.Status switch + { + InstanceStatus.Draining => RoleInstanceChangeKind.Draining, + InstanceStatus.Active => RoleInstanceChangeKind.Online, + InstanceStatus.Stopped => RoleInstanceChangeKind.Offline, + _ => null, + }; + if (kind.HasValue) + { + pendingEvents.Add(new KeyValuePair(kind.Value, descriptor)); + } + } + + /// + /// 曾知实例本轮文档消失(TTL 已清除)→ Evicted,彻底移出观测。 + /// + /// + /// Evicts known instances whose heartbeat documents disappeared this round (TTL cleanup). + /// + private void EvictMissingInstances(List documents, List> pendingEvents) + { + var observedInstanceIds = CollectObservedInstanceIds(documents); var evictedIds = new List(); foreach (var pair in _knownInstances) { if (!observedInstanceIds.Contains(pair.Key)) { - evictedIds.Add(pair.Key); pendingEvents.Add(new KeyValuePair(RoleInstanceChangeKind.Evicted, pair.Value.Descriptor)); + evictedIds.Add(pair.Key); } } @@ -430,8 +486,84 @@ private List ApplyStateTransitions(List + /// 收集本轮文档中可解析的实例 id(用于 Evicted 判定)。 + /// + /// + /// Collects the parseable instance ids from this round's documents for the missing-instance + /// check (matches the per-document parse policy of the main loop). + /// + private static HashSet CollectObservedInstanceIds(List documents) + { + var observedInstanceIds = new HashSet(StringComparer.Ordinal); + foreach (var document in documents) + { + var descriptor = TryToDescriptor(document); + if (descriptor == null) + { + continue; + } + observedInstanceIds.Add(descriptor.InstanceId); + } + return observedInstanceIds; + } + + /// + /// 是否 incarnation 变化(D15:旧身份被换新身份)。 + /// + private static bool IncarnationChanged(KnownInstance known, InstanceDescriptor descriptor) + { + return known.Descriptor.Incarnation != descriptor.Incarnation; + } + + /// + /// 是否从陈旧恢复新鲜(同 incarnation)。 + /// + private static bool RecoveredFromStale(KnownInstance known, bool isStale) + { + return known.IsStale && !isStale; + } + + /// + /// 是否发生状态跃迁(fresh → 另一 fresh 状态)。 + /// + private static bool StatusChanged(KnownInstance known, InstanceDescriptor descriptor, bool isStale) + { + return !known.IsStale && known.Descriptor.Status != descriptor.Status && !isStale; + } + + /// + /// 是否从新鲜变为陈旧(三周期阈值判死)。 + /// + private static bool BecameStale(KnownInstance known, bool isStale) + { + return !known.IsStale && isStale; + } + + /// + /// 是否进入路由表(非 stale 且非 Stopped)。 + /// + private static bool IsRoutable(InstanceDescriptor descriptor, bool isStale) + { + return !isStale && descriptor.Status != InstanceStatus.Stopped; + } + + /// + /// 首观测是否被跳过滤(stale 或非 Active/Draining 时不发事件)。 + /// + private static bool IsFirstObservationSkipped(bool isStale, InstanceStatus status) + { + return isStale || !IsActiveOrDraining(status); + } + + /// + /// 是否为 Active 或 Draining 状态。 + /// + private static bool IsActiveOrDraining(InstanceStatus status) + { + return status == InstanceStatus.Active || status == InstanceStatus.Draining; } ///