Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions Models/MonitorModels.cs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ public sealed class Iec61850MonitorDevice : ObservableObject
private string _sclIedName = string.Empty;
private string _sclAccessPointName = string.Empty;
private string _sclEndpointOrigin = "Unbound";
private Iec61850ReportContinuitySnapshot? _reportContinuityEvidence;

public string DeviceId { get; set; } = Guid.NewGuid().ToString("N");
public BulkObservableCollection<SignalDefinition> Signals { get; } = new();
Expand All @@ -53,6 +54,15 @@ public sealed class Iec61850MonitorDevice : ObservableObject
[JsonIgnore]
public StaticAcquisitionParitySnapshot StaticAcquisitionParity { get; set; } = new();

// Snapshot publication is atomic across runtime/diagnostic threads, never a
// mutable live ReportStreams dictionary. Association reset clears its scope.
[JsonIgnore]
public Iec61850ReportContinuitySnapshot? ReportContinuityEvidence
{
get => System.Threading.Volatile.Read(ref _reportContinuityEvidence);
set => System.Threading.Volatile.Write(ref _reportContinuityEvidence, value);
}

public SclIedWorkspace? SclWorkspace
{
get => _sclWorkspace;
Expand Down
43 changes: 43 additions & 0 deletions Services/DiagnosticReportBuilder.cs
Original file line number Diff line number Diff line change
Expand Up @@ -171,6 +171,7 @@ public static async Task<string> BuildAsync(
builder.AppendLine($"Detail : {device.Detail}");
builder.AppendLine($"Acquisition : {device.AcquisitionMode}");
AppendStaticIngressParity(builder, device.StaticAcquisitionParity);
AppendReportContinuity(builder, device.ReportContinuityEvidence);
if (device.IsMonitoring && Iec61850MonitoringModeRegistry.IsStaticDataSetReportOnly(device))
AppendStaticInitialImage(builder, device.Points);
builder.AppendLine($"Saved model : {device.HasDiscoveryCache} ({device.SignalCount:N0} signal(s))");
Expand Down Expand Up @@ -271,6 +272,48 @@ private static void AppendStaticInitialImage(
builder.AppendLine($" Additional pending: {image.PendingPoints.Count - 12} (showing first 12; never inferred unavailable)");
}

internal static void AppendReportContinuity(
StringBuilder builder, Iec61850ReportContinuitySnapshot? snapshot)
{
if (snapshot is null)
{
builder.AppendLine("Report continuity : NOT CAPTURED • no monitor evidence in current association");
return;
}

var verdict = snapshot.Frames == 0 ? "NO FRAME METADATA" :
snapshot.UntrackedStreamCount > 0 ? "INCOMPLETE • metadata stream limit reached" :
snapshot.Findings > 0 ? "ANOMALY OBSERVED" :
snapshot.Sequenced < snapshot.Frames ? "PARTIAL SqNum EVIDENCE" :
"OBSERVED WITHOUT SEQUENCE ALERTS";
builder.AppendLine($"Report continuity : {verdict} • frames={snapshot.Frames}, " +
$"sequenced={snapshot.Sequenced}, processUpdates={snapshot.ProcessUpdatesSeen}, " +
$"streams={snapshot.StreamCount}, findings={snapshot.Findings}, " +
$"BufOvfl={snapshot.Overflows}, untracked={snapshot.UntrackedStreamCount}");
builder.AppendLine(" Qualification : association-local decoded metadata only; no alerts " +
"does not prove SOE/event continuity, GI causality or reconnect/replay completeness.");

foreach (var (stream, index) in snapshot.Streams.Select((stream, index) => (stream, index)))
{
static string Safe(string value)
{
var text = (value ?? string.Empty).Replace("\r", @"\r", StringComparison.Ordinal)
.Replace("\n", @"\n", StringComparison.Ordinal);
return text.Length > 240 ? text[..240] + "…[TRUNCATED]" : text;
}
builder.AppendLine($" stream[{index}] : {(stream.Buffered ? "BRCB" : "URCB")}, " +
$"rcb={Safe(stream.Rcb)}, dataset={Safe(stream.DataSet)}, " +
$"rptId={Safe(stream.ReportId)}, frames={stream.Frames}, " +
$"SqNum={stream.Sequenced}, missingSqNum={stream.SequenceMissing}, " +
$"segmented={stream.Segmented}, findings={stream.Findings}, " +
$"BufOvfl={stream.Overflow}, ConfRevChanges={stream.ConfRevChanges}, " +
$"EntryIDPresent={stream.EntryIdPresent}, firstSqNum={stream.FirstSqNum?.ToString() ?? "-"}, " +
$"lastSqNum={stream.LastSqNum?.ToString() ?? "-"}");
}
if (snapshot.StreamCount > snapshot.Streams.Count)
builder.AppendLine($" stream list : TRUNCATED • shown={snapshot.Streams.Count}, total={snapshot.StreamCount}");
}

internal static void AppendStaticIngressParity(
StringBuilder builder,
StaticAcquisitionParitySnapshot parity)
Expand Down
49 changes: 48 additions & 1 deletion Services/Iec61850MonitorRuntime.cs
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,8 @@ private sealed class DeviceSession
public Dictionary<string, Iec61850MonitorPoint> ReportReferenceIndex { get; } = new(StringComparer.OrdinalIgnoreCase);
public PriorityQueue<string, long> PollQueue { get; } = new();
public Dictionary<string, Iec61850ReportContinuityState> ReportStreams { get; } = new(StringComparer.OrdinalIgnoreCase);
public long ContinuityProcessUpdatesSeen { get; set; }
public int ContinuityUntrackedStreams { get; set; }
public int LastUnroutedReportCount { get; set; }
public int ReportPlanCursor { get; set; }
public IReadOnlyList<ReportControlPlan> PendingReportPlans { get; set; } = Array.Empty<ReportControlPlan>();
Expand Down Expand Up @@ -489,6 +491,9 @@ public async Task<IReadOnlyList<Iec61850MonitorPoint>> StartMonitoringAsync(
session.ReportReferenceIndex.Clear();
session.PollQueue.Clear();
session.ReportStreams.Clear();
session.ContinuityProcessUpdatesSeen = 0;
session.ContinuityUntrackedStreams = 0;
session.Device.ReportContinuityEvidence = null;
session.LastUnroutedReportCount = 0;
session.ReportPlanCursor = 0;
session.PendingReportPlans = Array.Empty<ReportControlPlan>();
Expand Down Expand Up @@ -1314,12 +1319,39 @@ private void ProcessReportHealth(
$"{delta} IEC 61850 InformationReport frame(s) were not routed because RptID/DataSet identity was ambiguous. The engine intentionally refused unsafe DataSet projection.");
}

if (slice.Updates.Count > 0)
session.ContinuityProcessUpdatesSeen =
Math.Min(long.MaxValue - slice.Updates.Count, session.ContinuityProcessUpdatesSeen) +
slice.Updates.Count;

// Bound metadata cardinality even if an IED emits pathological RptID
// variations. Report acceptance and value delivery remain untouched.
const int maxTrackedStreams = 64;
foreach (var frame in slice.ReportFrames)
{
var streamKey = BuildReportStreamKey(plan, frame);
if (!session.ReportStreams.TryGetValue(streamKey, out var state))
{
state = new Iec61850ReportContinuityState();
if (session.ReportStreams.Count >= maxTrackedStreams)
{
if (session.ContinuityUntrackedStreams < int.MaxValue)
session.ContinuityUntrackedStreams++;
if (session.ContinuityUntrackedStreams == 1)
Log("WARN", session.Device.Name,
"Report continuity evidence reached the 64-stream diagnostic limit; " +
"additional metadata identities remain untracked; report processing continues.");
continue;
}

state = new Iec61850ReportContinuityState
{
ReportControlReference = string.IsNullOrWhiteSpace(frame.ReportControlReference)
? plan.ReportControlReference : frame.ReportControlReference,
DataSetReference = string.IsNullOrWhiteSpace(frame.DataSetReference)
? plan.DataSetReference : frame.DataSetReference,
ReportId = frame.ReportId,
Buffered = plan.Buffered
};
session.ReportStreams[streamKey] = state;
}

Expand All @@ -1333,6 +1365,15 @@ private void ProcessReportHealth(
$"{ReportName(frame.ReportControlReference, plan.ReportControlReference)}: {finding}");
}
}

// One immutable publication per non-empty slice, never per value/frame.
// The GUI may concurrently Copy Diagnostic without locking the MMS loop.
if (slice.Updates.Count > 0 || slice.ReportFrames.Count > 0)
session.Device.ReportContinuityEvidence =
Iec61850ReportContinuityInspector.Snapshot(
session.ReportStreams,
session.ContinuityProcessUpdatesSeen,
session.ContinuityUntrackedStreams);
}

private static string BuildReportStreamKey(ReportControlPlan plan, NativeReportFrameMetadata frame)
Expand Down Expand Up @@ -1760,6 +1801,9 @@ await DisposeClientForReconnectAsync(
session.ActiveReportPlanOrder.Clear();
session.PointPlanIds.Clear();
session.ReportStreams.Clear();
session.ContinuityProcessUpdatesSeen = 0;
session.ContinuityUntrackedStreams = 0;
session.Device.ReportContinuityEvidence = null;
session.LastUnroutedReportCount = 0;
ResetAssociationReportEvidence(session);

Expand Down Expand Up @@ -2064,6 +2108,9 @@ private async Task StopMonitoringSessionAsync(DeviceSession session, bool preser
session.StaticReportProjection.Reset();
session.PollQueue.Clear();
session.ReportStreams.Clear();
session.ContinuityProcessUpdatesSeen = 0;
session.ContinuityUntrackedStreams = 0;
session.Device.ReportContinuityEvidence = null;
session.LastUnroutedReportCount = 0;
session.ReportPlanCursor = 0;
session.PendingReportPlans = Array.Empty<ReportControlPlan>();
Expand Down
116 changes: 115 additions & 1 deletion Services/Iec61850ReportContinuityInspector.cs
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,60 @@ internal sealed class Iec61850ReportContinuityState
public bool AwaitingMoreSegments { get; set; }
public ulong? ConfigurationRevision { get; set; }
public string LastEntryIdHex { get; set; } = string.Empty;

// All counts belong to this exact report stream within the current
// association. No static state, no process/control-plane side effects.
public string ReportControlReference { get; set; } = string.Empty;
public string DataSetReference { get; set; } = string.Empty;
public string ReportId { get; set; } = string.Empty;
public bool Buffered { get; set; }
public long FramesSeen { get; set; }
public long SequencedFrames { get; set; }
public long MissingSequenceFrames { get; set; }
public long SegmentedFrames { get; set; }
public long AnomalyFindings { get; set; }
public long BufferOverflows { get; set; }
public long ConfRevChanges { get; set; }
public long EntryIdPresentFrames { get; set; }
public ulong? FirstSequenceNumber { get; set; }
}

/// <summary>
/// Immutable copyable diagnostic evidence. Counts describe decoded metadata,
/// NOT completeness of SOE delivery or absence of lost events.
/// </summary>
public sealed record Iec61850ReportContinuityStreamSnapshot
{
public string Rcb { get; init; } = string.Empty;
public string DataSet { get; init; } = string.Empty;
public string ReportId { get; init; } = string.Empty;
public bool Buffered { get; init; }
public long Frames { get; init; }
public long Sequenced { get; init; }
public long SequenceMissing { get; init; }
public long Segmented { get; init; }
public long Findings { get; init; }
public long Overflow { get; init; }
public long ConfRevChanges { get; init; }
public long EntryIdPresent { get; init; }
public ulong? FirstSqNum { get; init; }
public ulong? LastSqNum { get; init; }
}

public sealed record Iec61850ReportContinuitySnapshot
{
public long ProcessUpdatesSeen { get; init; }
public int StreamCount { get; init; }
public int UntrackedStreamCount { get; init; }
public long Frames { get; init; }
public long Sequenced { get; init; }
public long Findings { get; init; }
public long Overflows { get; init; }
public IReadOnlyList<Iec61850ReportContinuityStreamSnapshot> Streams { get; init; } =
Array.Empty<Iec61850ReportContinuityStreamSnapshot>();
}


internal static class Iec61850ReportContinuityInspector
{
private const ulong MaxSequenceNumber = ushort.MaxValue;
Expand All @@ -35,8 +87,25 @@ internal static IReadOnlyList<string> Observe(
ArgumentNullException.ThrowIfNull(state);
ArgumentNullException.ThrowIfNull(frame);

if (state.FramesSeen < long.MaxValue) state.FramesSeen++;
if (frame.SequenceNumber.HasValue)
{
if (state.SequencedFrames < long.MaxValue) state.SequencedFrames++;
if (frame.SequenceNumber.Value <= MaxSequenceNumber)
state.FirstSequenceNumber ??= frame.SequenceNumber;
}
else if (state.MissingSequenceFrames < long.MaxValue) state.MissingSequenceFrames++;
if (frame.SubSequenceNumber.HasValue && state.SegmentedFrames < long.MaxValue)
state.SegmentedFrames++;
if (!string.IsNullOrEmpty(frame.EntryIdHex) && state.EntryIdPresentFrames < long.MaxValue)
state.EntryIdPresentFrames++;

List<string>? findings = null;
void Warn(string message) => (findings ??= new List<string>(2)).Add(message);
void Warn(string message)
{
(findings ??= new List<string>(2)).Add(message);
if (state.AnomalyFindings < long.MaxValue) state.AnomalyFindings++;
}

var priorEntryIdPresent = !string.IsNullOrEmpty(state.LastEntryIdHex);
var entryIdPresent = !string.IsNullOrEmpty(frame.EntryIdHex);
Expand All @@ -50,6 +119,7 @@ string EntryContext() => buffered

if (frame.BufferOverflow == true)
{
if (state.BufferOverflows < long.MaxValue) state.BufferOverflows++;
Warn(buffered
? "BRCB BufOvfl=true: possible loss of buffered entries; continuity cannot be certified from SqNum alone"
: "URCB carried unexpected BufOvfl=true; inspect report OptionFields/decoder attribution");
Expand All @@ -60,6 +130,7 @@ string EntryContext() => buffered
if (state.ConfigurationRevision.HasValue &&
state.ConfigurationRevision.Value != frame.ConfRev.Value)
{
if (state.ConfRevChanges < long.MaxValue) state.ConfRevChanges++;
Warn($"Report ConfRev changed: previous={state.ConfigurationRevision.Value}, " +
$"current={frame.ConfRev.Value}; DataSet schema and member bindings require revalidation");
}
Expand Down Expand Up @@ -139,6 +210,49 @@ string EntryContext() => buffered
return findings is null ? Array.Empty<string>() : findings;
}

internal static Iec61850ReportContinuitySnapshot Snapshot(
IReadOnlyDictionary<string, Iec61850ReportContinuityState> streams,
long processUpdatesSeen,
int untrackedStreamCount,
int maxShown = 12)
{
ArgumentNullException.ThrowIfNull(streams);
var ordered = streams.Values
.OrderBy(state => state.ReportControlReference, StringComparer.OrdinalIgnoreCase)
.ThenBy(state => state.ReportId, StringComparer.OrdinalIgnoreCase)
.ThenBy(state => state.DataSetReference, StringComparer.OrdinalIgnoreCase)
.ToArray();
var shown = ordered.Take(Math.Clamp(maxShown, 0, 32)).Select(state =>
new Iec61850ReportContinuityStreamSnapshot
{
Rcb = state.ReportControlReference,
DataSet = state.DataSetReference,
ReportId = state.ReportId,
Buffered = state.Buffered,
Frames = state.FramesSeen,
Sequenced = state.SequencedFrames,
SequenceMissing = state.MissingSequenceFrames,
Segmented = state.SegmentedFrames,
Findings = state.AnomalyFindings,
Overflow = state.BufferOverflows,
ConfRevChanges = state.ConfRevChanges,
EntryIdPresent = state.EntryIdPresentFrames,
FirstSqNum = state.FirstSequenceNumber,
LastSqNum = state.LastSequenceNumber
}).ToArray();
return new Iec61850ReportContinuitySnapshot
{
ProcessUpdatesSeen = processUpdatesSeen,
StreamCount = ordered.Length,
UntrackedStreamCount = untrackedStreamCount,
Frames = ordered.Sum(state => state.FramesSeen),
Sequenced = ordered.Sum(state => state.SequencedFrames),
Findings = ordered.Sum(state => state.AnomalyFindings),
Overflows = ordered.Sum(state => state.BufferOverflows),
Streams = shown
};
}

private static string? DescribeSequenceAnomaly(ulong? previous, ulong current)
{
if (!previous.HasValue)
Expand Down
Loading