This is an automated email from the ASF dual-hosted git repository.

zqr10159 pushed a commit to branch 2.0.0
in repository https://gitbox.apache.org/repos/asf/hertzbeat.git

commit 6baa3638cb353bd45e5bd7877a78911613820328
Author: Logic <[email protected]>
AuthorDate: Fri Aug 28 18:18:38 2026 +0800

    Add investigation correlation model
---
 .../enricher/OtlpEntityIdentityResolver.java       |   8 +-
 .../enricher/OtlpEntityIdentityResolverTest.java   |  33 ++-
 .../model/investigation-anchor-model.ts            | 125 ++++++++++
 .../model/investigation-correlation-model.test.ts  | 274 +++++++++++++++++++++
 .../model/investigation-correlation-model.ts       | 190 ++++++++++++++
 .../model/investigation-model-validation.ts        |  52 ++++
 6 files changed, 677 insertions(+), 5 deletions(-)

diff --git 
a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/enricher/OtlpEntityIdentityResolver.java
 
b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/enricher/OtlpEntityIdentityResolver.java
index d2b4633344..9b7a9aead4 100644
--- 
a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/enricher/OtlpEntityIdentityResolver.java
+++ 
b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/enricher/OtlpEntityIdentityResolver.java
@@ -303,7 +303,7 @@ public class OtlpEntityIdentityResolver {
             if (!entities.containsKey(entry.getKey())) {
                 continue;
             }
-            int score = entry.getValue().matchedIdentityCount();
+            int score = entry.getValue().canonicalIdentityScore();
             if (score > topScore) {
                 topScore = score;
                 topEntityIds.clear();
@@ -554,8 +554,10 @@ public class OtlpEntityIdentityResolver {
             matchedIdentityKeys.add(identityKey);
         }
 
-        private int matchedIdentityCount() {
-            return matchedIdentityKeys.size();
+        private int canonicalIdentityScore() {
+            return matchedIdentityKeys.stream()
+                    .mapToInt(EntityCanonicalIdentityRegistry::defaultPriority)
+                    .sum();
         }
     }
 }
diff --git 
a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/enricher/OtlpEntityIdentityResolverTest.java
 
b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/enricher/OtlpEntityIdentityResolverTest.java
index 3caea5ec96..98f7c04dae 100644
--- 
a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/enricher/OtlpEntityIdentityResolverTest.java
+++ 
b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/enricher/OtlpEntityIdentityResolverTest.java
@@ -103,11 +103,40 @@ class OtlpEntityIdentityResolverTest {
     }
 
     @Test
-    void doesNotResolveEntityWhenCanonicalEvidenceIsSplitAcrossEntities() {
+    void 
doesNotLetThreeWeakIdentityMatchesSilentlyOverrideOneStrongIdentityMatch() {
+        
when(workspaceQueryGateway.findIdentitiesByKeysAndNormalizedValues(eq("prod-west"),
 anySet(), anySet()))
+                .thenReturn(List.of(
+                        identity(41L, "service.instance.id", "checkout-1", 
"checkout-1", 140, true),
+                        identity(42L, "service.name", "checkout", "checkout", 
90, true),
+                        identity(42L, "service.namespace", "commerce", 
"commerce", 30, false),
+                        identity(42L, "deployment.environment.name", "prod", 
"prod", 20, false)));
+        when(workspaceQueryGateway.findEntitiesByIds("prod-west", Set.of(41L, 
42L)))
+                .thenReturn(Map.of(
+                        41L, entity(41L, "prod-west", "service", "checkout-1", 
"Checkout Instance"),
+                        42L, entity(42L, "prod-west", "service", "checkout", 
"Checkout Service")));
+
+        Optional<String> resolved = resolver.resolveEntityId(Map.of(
+                "service.instance.id", "checkout-1",
+                "service.name", "checkout",
+                "service.namespace", "commerce",
+                "deployment.environment.name", "prod"), "prod-west");
+
+        assertTrue(resolved.isEmpty());
+        verify(workspaceQueryGateway).recordEntityDiscoveryGovernanceActivity(
+                eq("prod-west"),
+                eq("identity_conflict"),
+                eq("needs_governance"),
+                eq("OTLP resource identity matched multiple entities"),
+                
org.mockito.ArgumentMatchers.contains("service.instance.id=checkout-1"),
+                eq(Map.of(41L, "Checkout Instance", 42L, "Checkout Service")));
+    }
+
+    @Test
+    void recordsConflictWhenTopWeightedEvidenceIsEqualAcrossEntities() {
         
when(workspaceQueryGateway.findIdentitiesByKeysAndNormalizedValues(eq("prod-west"),
 anySet(), anySet()))
                 .thenReturn(List.of(
                         identity(41L, "service.name", "checkout", "checkout", 
90, true),
-                        identity(41L, "deployment.environment.name", "prod", 
"prod", 20, false),
+                        identity(41L, "service.namespace", "commerce", 
"commerce", 30, false),
                         identity(42L, "service.name", "checkout", "checkout", 
90, true),
                         identity(42L, "service.namespace", "commerce", 
"commerce", 30, false)));
         when(workspaceQueryGateway.findEntitiesByIds("prod-west", Set.of(41L, 
42L)))
diff --git 
a/web-app/src/features/investigation/model/investigation-anchor-model.ts 
b/web-app/src/features/investigation/model/investigation-anchor-model.ts
new file mode 100644
index 0000000000..b66bef3226
--- /dev/null
+++ b/web-app/src/features/investigation/model/investigation-anchor-model.ts
@@ -0,0 +1,125 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0.
+ */
+
+import {
+  normalizeInvestigationTimeZone,
+  parseQueryContext,
+  QUERY_CONTEXT_FIELDS,
+  writeQueryContext,
+  type InvestigationTimeWindow,
+  type QueryContext
+} from '@/shared/query-context';
+
+import {
+  isPositiveJavaLong,
+  normalizeOpaqueId,
+  normalizePositiveId,
+  requireRecord
+} from './investigation-model-validation';
+
+export type InvestigationSource = 'entity' | 'monitor' | 'alert' | 'metric' | 
'log' | 'trace' | 'topology' | 'ai';
+
+export type InvestigationAnchor = {
+  source: InvestigationSource;
+  context: QueryContext;
+  window: InvestigationTimeWindow;
+  traceId?: string | undefined;
+  spanId?: string | undefined;
+  alertId?: string | undefined;
+};
+
+export type InvestigationCapabilityState = 'available' | 'empty' | 
'unavailable' | 'unknown';
+type CapabilityKey =
+  | 'metrics'
+  | 'logs'
+  | 'traces'
+  | 'topology'
+  | 'collection'
+  | 'alerts'
+  | 'nativeMetrics'
+  | 'otelMetrics'
+  | 'redMetrics'
+  | 'traceCorrelation'
+  | 'logTraceCorrelation'
+  | 'semanticGraph';
+
+export type SignalCapabilities = Readonly<Record<CapabilityKey, 
InvestigationCapabilityState>>;
+
+const sources: readonly InvestigationSource[] = [
+  'entity',
+  'monitor',
+  'alert',
+  'metric',
+  'log',
+  'trace',
+  'topology',
+  'ai'
+];
+const capabilityKeys: readonly CapabilityKey[] = [
+  'metrics',
+  'logs',
+  'traces',
+  'topology',
+  'collection',
+  'alerts',
+  'nativeMetrics',
+  'otelMetrics',
+  'redMetrics',
+  'traceCorrelation',
+  'logTraceCorrelation',
+  'semanticGraph'
+];
+const capabilityStates: readonly InvestigationCapabilityState[] = 
['available', 'empty', 'unavailable', 'unknown'];
+const contextKeys = Object.values(QUERY_CONTEXT_FIELDS);
+
+export function createInvestigationAnchor(input: InvestigationAnchor): 
InvestigationAnchor {
+  requireRecord(input, ['source', 'context', 'window', 'traceId', 'spanId', 
'alertId']);
+  if (!sources.includes(input.source)) throw new Error('Investigation anchor 
source is invalid');
+  requireRecord(input.context, contextKeys);
+  requireRecord(input.window, ['from', 'to', 'timeZone']);
+
+  const context = parseQueryContext(writeQueryContext(new URLSearchParams(), 
input.context));
+  for (const key of ['entityId', 'monitorId', 'intakeProfileId'] as const) {
+    if (context[key] != null && !isPositiveJavaLong(context[key])) {
+      throw new Error('Investigation anchor identity is invalid');
+    }
+  }
+  if (!validWindow(input.window)) throw new Error('Investigation anchor window 
is invalid');
+  const timeZone = normalizeInvestigationTimeZone(input.window.timeZone);
+  if (!timeZone) throw new Error('Investigation anchor time zone is invalid');
+
+  const traceId = normalizeOpaqueId(input.traceId);
+  const spanId = normalizeOpaqueId(input.spanId);
+  const alertId = normalizePositiveId(input.alertId);
+  if (spanId && !traceId) throw new Error('Investigation span identity 
requires a trace identity');
+
+  return {
+    source: input.source,
+    context,
+    window: { from: input.window.from, to: input.window.to, timeZone },
+    ...(traceId ? { traceId } : {}),
+    ...(spanId ? { spanId } : {}),
+    ...(alertId ? { alertId } : {})
+  };
+}
+
+export function createSignalCapabilities(input: Partial<SignalCapabilities>): 
SignalCapabilities {
+  requireRecord(input, capabilityKeys);
+  return Object.fromEntries(
+    capabilityKeys.map(key => {
+      const state = input[key] ?? 'unknown';
+      if (!capabilityStates.includes(state)) throw new Error('Investigation 
capability state is invalid');
+      return [key, state];
+    })
+  ) as SignalCapabilities;
+}
+
+function validWindow(window: InvestigationTimeWindow) {
+  return (
+    Number.isSafeInteger(window.from) && Number.isSafeInteger(window.to) && 
window.from > 0 && window.from < window.to
+  );
+}
diff --git 
a/web-app/src/features/investigation/model/investigation-correlation-model.test.ts
 
b/web-app/src/features/investigation/model/investigation-correlation-model.test.ts
new file mode 100644
index 0000000000..181c165fd1
--- /dev/null
+++ 
b/web-app/src/features/investigation/model/investigation-correlation-model.test.ts
@@ -0,0 +1,274 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0.
+ */
+
+import { describe, expect, it } from 'vitest';
+
+import {
+  correlateInvestigationEvidence,
+  type CorrelationReason,
+  type InvestigationCorrelationCandidate,
+  type InvestigationEvidenceSummary
+} from './investigation-correlation-model';
+import {
+  createInvestigationAnchor,
+  createSignalCapabilities,
+  type InvestigationAnchor,
+  type InvestigationSource,
+  type SignalCapabilities
+} from './investigation-anchor-model';
+
+const window = { from: 1_750_000_000_000, to: 1_750_000_060_000, timeZone: 
'UTC' } as const;
+const sources: InvestigationSource[] = ['entity', 'monitor', 'alert', 
'metric', 'log', 'trace', 'topology', 'ai'];
+
+describe('investigation correlation model', () => {
+  it.each(sources)('accepts the supported %s anchor source without duplicating 
shared context fields', source => {
+    const anchor = createInvestigationAnchor({
+      source,
+      context: {
+        entityId: ' 7 ',
+        monitorId: ' 42 ',
+        serviceName: ' checkout ',
+        serviceNamespace: ' commerce ',
+        environment: ' production '
+      },
+      window: { ...window, timeZone: ' UTC ' },
+      traceId: ' trace-1 ',
+      spanId: ' span-1 ',
+      alertId: ' 91 '
+    });
+
+    expect(anchor).toEqual({
+      source,
+      context: {
+        entityId: '7',
+        monitorId: '42',
+        serviceName: 'checkout',
+        serviceNamespace: 'commerce',
+        environment: 'production'
+      },
+      window,
+      traceId: 'trace-1',
+      spanId: 'span-1',
+      alertId: '91'
+    });
+    expect(anchor).not.toHaveProperty('entityId');
+    expect(anchor).not.toHaveProperty('monitorId');
+    expect(anchor).not.toHaveProperty('serviceName');
+  });
+
+  it.each([
+    { source: 'unsupported', context: {}, window },
+    { source: 'entity', context: {}, window: { ...window, from: 0 } },
+    { source: 'entity', context: {}, window: { ...window, to: window.from } },
+    { source: 'entity', context: {}, window: { ...window, timeZone: 
'not/a-zone' } },
+    { source: 'entity', context: { entityId: '01' }, window },
+    { source: 'monitor', context: { monitorId: '9223372036854775808' }, window 
},
+    { source: 'trace', context: {}, window, traceId: ' ' },
+    { source: 'trace', context: {}, window, spanId: 'span-without-trace' },
+    { source: 'alert', context: {}, window, alertId: '0' },
+    { source: 'entity', context: { authorization: 'private' }, window }
+  ])('fails closed on invalid anchor evidence %#', candidate => {
+    expect(() => createInvestigationAnchor(candidate as 
InvestigationAnchor)).toThrow();
+  });
+
+  it('keeps every capability independently honest and defaults only missing 
evidence to unknown', () => {
+    const capabilities: SignalCapabilities = createSignalCapabilities({
+      metrics: 'available',
+      logs: 'empty',
+      traces: 'unavailable',
+      topology: 'unknown',
+      nativeMetrics: 'available',
+      traceCorrelation: 'empty'
+    });
+
+    expect(capabilities).toEqual({
+      metrics: 'available',
+      logs: 'empty',
+      traces: 'unavailable',
+      topology: 'unknown',
+      collection: 'unknown',
+      alerts: 'unknown',
+      nativeMetrics: 'available',
+      otelMetrics: 'unknown',
+      redMetrics: 'unknown',
+      traceCorrelation: 'empty',
+      logTraceCorrelation: 'unknown',
+      semanticGraph: 'unknown'
+    });
+    expect(() => createSignalCapabilities({ logs: 'ready' } as 
never)).toThrow();
+    expect(() => createSignalCapabilities({ metrics: 'available', telemetry: 
'available' } as never)).toThrow();
+  });
+
+  it('orders correlations by authoritative reason and derives confidence 
without caller input', () => {
+    const anchor = investigationAnchor();
+    const candidates = [
+      candidate('time', anchorFor({ entityId: '13', serviceName: 'inventory' 
})),
+      candidate('topology', anchorFor({ entityId: '12', serviceName: 
'inventory' }), {
+        sourceEntityId: '7',
+        targetEntityId: '12'
+      }),
+      candidate('otel', anchorFor({ entityId: '11', monitorId: '50' })),
+      candidate('monitor', anchorFor({ entityId: '10', monitorId: '42', 
serviceName: 'inventory' })),
+      candidate('entity', anchorFor({ entityId: '7', monitorId: '50', 
serviceName: 'inventory' })),
+      candidate('trace', anchorFor({ traceId: 'trace-1', spanId: 'span-2', 
entityId: '20' })),
+      candidate('span', anchorFor({ traceId: 'trace-1', spanId: 'span-1', 
entityId: '21' }))
+    ];
+
+    const evidence: InvestigationEvidenceSummary[] = 
correlateInvestigationEvidence(anchor, candidates);
+    const reasons: CorrelationReason[] = evidence.map(item => item.reason);
+
+    expect(evidence.map(item => item.candidateQuery.key)).toEqual([
+      'span',
+      'trace',
+      'entity',
+      'monitor',
+      'otel',
+      'topology',
+      'time'
+    ]);
+    expect(reasons).toEqual([
+      { kind: 'exact-span', level: 1 },
+      { kind: 'exact-trace', level: 1 },
+      { kind: 'same-entity', level: 2 },
+      { kind: 'bound-monitor', level: 3 },
+      { kind: 'canonical-otel-identity', level: 4 },
+      { kind: 'topology-related', level: 5 },
+      { kind: 'time-proximity', level: 6 }
+    ]);
+    expect(evidence.map(item => item.confidence)).toEqual([
+      'exact',
+      'exact',
+      'high',
+      'high',
+      'medium',
+      'medium',
+      'low'
+    ]);
+  });
+
+  it('uses time only as proximity evidence and never copies same-entity 
identity into the candidate', () => {
+    const [evidence] = correlateInvestigationEvidence(investigationAnchor(), [
+      candidate('nearby', anchorFor({ entityId: '99', monitorId: '100', 
serviceName: 'inventory' }))
+    ]);
+
+    expect(evidence?.reason).toEqual({ kind: 'time-proximity', level: 6 });
+    expect(evidence?.candidateQuery.anchor.context.entityId).toBe('99');
+    expect(evidence?.candidateQuery.anchor.context.entityId).not.toBe('7');
+  });
+
+  it('uses the same non-empty service instance as canonical OTel identity 
before the service tuple fallback', () => {
+    const anchor = createInvestigationAnchor({
+      ...investigationAnchor(),
+      context: { ...investigationAnchor().context, instance: 'checkout-7d9' }
+    });
+    const target = anchorFor({
+      entityId: '99',
+      monitorId: '100',
+      serviceName: 'inventory',
+      serviceNamespace: 'warehouse',
+      environment: 'staging',
+      instance: 'checkout-7d9'
+    });
+
+    expect(correlateInvestigationEvidence(anchor, [candidate('instance', 
target)])[0]?.reason).toEqual({
+      kind: 'canonical-otel-identity',
+      level: 4
+    });
+  });
+
+  it('deduplicates by candidate query key deterministically and keeps the 
strongest evidence', () => {
+    const anchor = investigationAnchor();
+    const nearby = candidate('same-query', anchorFor({ entityId: '99', 
serviceName: 'inventory' }));
+    const exact = candidate('same-query', anchorFor({ traceId: 'trace-1', 
spanId: 'span-1', entityId: '21' }));
+
+    const forward = correlateInvestigationEvidence(anchor, [nearby, exact]);
+    const reversed = correlateInvestigationEvidence(anchor, [exact, nearby]);
+
+    expect(forward).toEqual(reversed);
+    expect(forward).toHaveLength(1);
+    expect(forward[0]?.reason).toEqual({ kind: 'exact-span', level: 1 });
+  });
+
+  it('keeps evidence summary-only and rejects raw telemetry or arbitrary 
confidence fields', () => {
+    const anchor = investigationAnchor();
+    const valid = candidate('summary', anchorFor({ traceId: 'trace-1' }));
+    const [evidence] = correlateInvestigationEvidence(anchor, [valid]);
+
+    expect(evidence).toEqual({
+      summary: { state: 'available', count: 2 },
+      candidateQuery: { key: 'summary', anchor: valid.candidateQuery.anchor },
+      reason: { kind: 'exact-trace', level: 1 },
+      confidence: 'exact'
+    });
+    expect(Object.keys(evidence ?? {})).toEqual(['summary', 'candidateQuery', 
'reason', 'confidence']);
+
+    expect(() =>
+      correlateInvestigationEvidence(anchor, [{ ...valid, rawTelemetry: { 
body: 'private' } } as never])
+    ).toThrow();
+    expect(() => correlateInvestigationEvidence(anchor, [{ ...valid, 
confidence: 'exact' } as never])).toThrow();
+    expect(() =>
+      correlateInvestigationEvidence(anchor, [{ ...valid, summary: { 
...valid.summary, body: 'private' } } as never])
+    ).toThrow();
+  });
+
+  it('omits candidates outside the bounded proximity window without stronger 
correlation evidence', () => {
+    const distant = anchorFor(
+      { entityId: '99', monitorId: '100', serviceName: 'inventory' },
+      { from: window.to + 300_001, to: window.to + 360_000, timeZone: 'UTC' }
+    );
+    expect(correlateInvestigationEvidence(investigationAnchor(), 
[candidate('distant', distant)])).toEqual([]);
+  });
+});
+
+function investigationAnchor(): InvestigationAnchor {
+  return createInvestigationAnchor({
+    source: 'trace',
+    context: {
+      entityId: '7',
+      monitorId: '42',
+      serviceName: 'checkout',
+      serviceNamespace: 'commerce',
+      environment: 'production'
+    },
+    window,
+    traceId: 'trace-1',
+    spanId: 'span-1'
+  });
+}
+
+function anchorFor(
+  patch: Partial<InvestigationAnchor['context']> & { traceId?: string; 
spanId?: string },
+  candidateWindow: InvestigationAnchor['window'] = window
+): InvestigationAnchor {
+  const { traceId, spanId, ...contextPatch } = patch;
+  return createInvestigationAnchor({
+    source: traceId ? 'trace' : 'log',
+    context: {
+      entityId: '30',
+      monitorId: '60',
+      serviceName: 'checkout',
+      serviceNamespace: 'commerce',
+      environment: 'production',
+      ...contextPatch
+    },
+    window: candidateWindow,
+    ...(traceId ? { traceId } : {}),
+    ...(spanId ? { spanId } : {})
+  });
+}
+
+function candidate(
+  key: string,
+  anchor: InvestigationAnchor,
+  topologyRelation?: InvestigationCorrelationCandidate['topologyRelation']
+): InvestigationCorrelationCandidate {
+  return {
+    candidateQuery: { key, anchor },
+    summary: { state: 'available', count: 2 },
+    ...(topologyRelation ? { topologyRelation } : {})
+  };
+}
diff --git 
a/web-app/src/features/investigation/model/investigation-correlation-model.ts 
b/web-app/src/features/investigation/model/investigation-correlation-model.ts
new file mode 100644
index 0000000000..6eb14e676b
--- /dev/null
+++ 
b/web-app/src/features/investigation/model/investigation-correlation-model.ts
@@ -0,0 +1,190 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0.
+ */
+
+import type { InvestigationTimeWindow, QueryContext } from 
'@/shared/query-context';
+
+import {
+  createInvestigationAnchor,
+  type InvestigationAnchor,
+  type InvestigationCapabilityState
+} from './investigation-anchor-model';
+import { normalizeOpaqueId, normalizePositiveId, requireRecord } from 
'./investigation-model-validation';
+
+export type CorrelationReason =
+  | { kind: 'exact-span'; level: 1 }
+  | { kind: 'exact-trace'; level: 1 }
+  | { kind: 'same-entity'; level: 2 }
+  | { kind: 'bound-monitor'; level: 3 }
+  | { kind: 'canonical-otel-identity'; level: 4 }
+  | { kind: 'topology-related'; level: 5 }
+  | { kind: 'time-proximity'; level: 6 };
+
+type CorrelationConfidence = 'exact' | 'high' | 'medium' | 'low';
+type EvidenceSummary = { state: InvestigationCapabilityState; count?: number | 
undefined };
+type CandidateQueryMetadata = { key: string; anchor: InvestigationAnchor };
+type TopologyRelation = { sourceEntityId: string; targetEntityId: string };
+
+export type InvestigationCorrelationCandidate = {
+  candidateQuery: CandidateQueryMetadata;
+  summary: EvidenceSummary;
+  topologyRelation?: TopologyRelation | undefined;
+};
+
+export type InvestigationEvidenceSummary = {
+  summary: EvidenceSummary;
+  candidateQuery: CandidateQueryMetadata;
+  reason: CorrelationReason;
+  confidence: CorrelationConfidence;
+};
+
+const evidenceStates: readonly InvestigationCapabilityState[] = ['available', 
'empty', 'unavailable', 'unknown'];
+const MAX_TIME_PROXIMITY_MS = 5 * 60_000;
+
+const correlationRules = {
+  'exact-span': { level: 1, order: 0, confidence: 'exact' },
+  'exact-trace': { level: 1, order: 1, confidence: 'exact' },
+  'same-entity': { level: 2, order: 2, confidence: 'high' },
+  'bound-monitor': { level: 3, order: 3, confidence: 'high' },
+  'canonical-otel-identity': { level: 4, order: 4, confidence: 'medium' },
+  'topology-related': { level: 5, order: 5, confidence: 'medium' },
+  'time-proximity': { level: 6, order: 6, confidence: 'low' }
+} as const satisfies Record<
+  CorrelationReason['kind'],
+  { level: CorrelationReason['level']; order: number; confidence: 
CorrelationConfidence }
+>;
+
+export function correlateInvestigationEvidence(
+  inputAnchor: InvestigationAnchor,
+  inputCandidates: readonly InvestigationCorrelationCandidate[]
+): InvestigationEvidenceSummary[] {
+  const anchor = createInvestigationAnchor(inputAnchor);
+  const ranked = inputCandidates.flatMap(input => {
+    const candidate = normalizeCandidate(input);
+    const reasonKind = correlationReason(anchor, candidate);
+    if (!reasonKind) return [];
+    const rule = correlationRules[reasonKind];
+    const evidence: InvestigationEvidenceSummary = {
+      summary: candidate.summary,
+      candidateQuery: candidate.candidateQuery,
+      reason: { kind: reasonKind, level: rule.level } as CorrelationReason,
+      confidence: rule.confidence
+    };
+    return [{ evidence, order: rule.order, fingerprint: 
JSON.stringify(evidence) }];
+  });
+
+  ranked.sort(
+    (left, right) =>
+      left.order - right.order ||
+      compareText(left.evidence.candidateQuery.key, 
right.evidence.candidateQuery.key) ||
+      compareText(left.fingerprint, right.fingerprint)
+  );
+  const seen = new Set<string>();
+  return ranked.flatMap(item => {
+    const key = item.evidence.candidateQuery.key;
+    if (seen.has(key)) return [];
+    seen.add(key);
+    return [item.evidence];
+  });
+}
+
+function normalizeCandidate(input: InvestigationCorrelationCandidate): 
InvestigationCorrelationCandidate {
+  requireRecord(input, ['candidateQuery', 'summary', 'topologyRelation']);
+  requireRecord(input.candidateQuery, ['key', 'anchor']);
+  const key = normalizeRequiredMetadata(input.candidateQuery.key);
+  const topologyRelation = input.topologyRelation ? 
normalizeTopologyRelation(input.topologyRelation) : undefined;
+  return {
+    candidateQuery: { key, anchor: 
createInvestigationAnchor(input.candidateQuery.anchor) },
+    summary: normalizeEvidenceSummary(input.summary),
+    ...(topologyRelation ? { topologyRelation } : {})
+  };
+}
+
+function normalizeEvidenceSummary(input: EvidenceSummary): EvidenceSummary {
+  requireRecord(input, ['state', 'count']);
+  if (!evidenceStates.includes(input.state)) throw new Error('Investigation 
evidence state is invalid');
+  const count = input.count;
+  if (count == null) return { state: input.state };
+  const countContradictsState =
+    input.state === 'unavailable' ||
+    input.state === 'unknown' ||
+    (input.state === 'empty' && count !== 0) ||
+    (input.state === 'available' && count === 0);
+  if (!Number.isSafeInteger(count) || count < 0 || countContradictsState) {
+    throw new Error('Investigation evidence count is invalid');
+  }
+  return { state: input.state, count };
+}
+
+function normalizeTopologyRelation(input: TopologyRelation): TopologyRelation {
+  requireRecord(input, ['sourceEntityId', 'targetEntityId']);
+  const sourceEntityId = normalizePositiveId(input.sourceEntityId);
+  const targetEntityId = normalizePositiveId(input.targetEntityId);
+  if (!sourceEntityId || !targetEntityId || sourceEntityId === targetEntityId) 
{
+    throw new Error('Investigation topology relation is invalid');
+  }
+  return { sourceEntityId, targetEntityId };
+}
+
+function correlationReason(
+  anchor: InvestigationAnchor,
+  candidate: InvestigationCorrelationCandidate
+): CorrelationReason['kind'] | undefined {
+  const target = candidate.candidateQuery.anchor;
+  if (anchor.traceId && anchor.traceId === target.traceId) {
+    if (anchor.spanId && anchor.spanId === target.spanId) return 'exact-span';
+    return 'exact-trace';
+  }
+  if (sameContextId(anchor.context.entityId, target.context.entityId)) return 
'same-entity';
+  if (sameContextId(anchor.context.monitorId, target.context.monitorId)) 
return 'bound-monitor';
+  if (sameOtelIdentity(anchor.context, target.context)) return 
'canonical-otel-identity';
+  if (candidate.topologyRelation && connectsEntities(anchor, target, 
candidate.topologyRelation)) {
+    return 'topology-related';
+  }
+  if (windowGap(anchor.window, target.window) <= MAX_TIME_PROXIMITY_MS) return 
'time-proximity';
+  return undefined;
+}
+
+function sameContextId(left: string | undefined, right: string | undefined) {
+  return left != null && right != null && left === right;
+}
+
+function sameOtelIdentity(left: QueryContext, right: QueryContext) {
+  if (left.instance != null && left.instance === right.instance) return true;
+  return (
+    left.serviceName != null &&
+    left.serviceName === right.serviceName &&
+    left.serviceNamespace === right.serviceNamespace &&
+    left.environment === right.environment
+  );
+}
+
+function connectsEntities(left: InvestigationAnchor, right: 
InvestigationAnchor, relation: TopologyRelation) {
+  const leftId = left.context.entityId;
+  const rightId = right.context.entityId;
+  return (
+    leftId != null &&
+    rightId != null &&
+    ((relation.sourceEntityId === leftId && relation.targetEntityId === 
rightId) ||
+      (relation.sourceEntityId === rightId && relation.targetEntityId === 
leftId))
+  );
+}
+
+function windowGap(left: InvestigationTimeWindow, right: 
InvestigationTimeWindow) {
+  if (left.to < right.from) return right.from - left.to;
+  if (right.to < left.from) return left.from - right.to;
+  return 0;
+}
+
+function normalizeRequiredMetadata(value: string) {
+  const normalized = normalizeOpaqueId(value);
+  if (!normalized) throw new Error('Investigation candidate query key is 
invalid');
+  return normalized;
+}
+
+function compareText(left: string, right: string) {
+  return left < right ? -1 : left > right ? 1 : 0;
+}
diff --git 
a/web-app/src/features/investigation/model/investigation-model-validation.ts 
b/web-app/src/features/investigation/model/investigation-model-validation.ts
new file mode 100644
index 0000000000..2c2d144471
--- /dev/null
+++ b/web-app/src/features/investigation/model/investigation-model-validation.ts
@@ -0,0 +1,52 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0.
+ */
+
+const JAVA_LONG_MAX = '9223372036854775807';
+
+export function normalizeOpaqueId(value: string | undefined) {
+  if (value === undefined) return undefined;
+  if (typeof value !== 'string') throw new Error('Investigation identity is 
invalid');
+  const normalized = value.trim();
+  if (!normalized || normalized.length > 256 || 
hasControlCharacter(normalized)) {
+    throw new Error('Investigation identity is invalid');
+  }
+  return normalized;
+}
+
+export function normalizePositiveId(value: string | undefined) {
+  if (value === undefined) return undefined;
+  if (typeof value !== 'string') throw new Error('Investigation identity is 
invalid');
+  const normalized = value.trim();
+  if (!isPositiveJavaLong(normalized)) throw new Error('Investigation identity 
is invalid');
+  return normalized;
+}
+
+export function isPositiveJavaLong(value: string) {
+  return (
+    /^[1-9]\d{0,18}$/u.test(value) &&
+    (value.length < JAVA_LONG_MAX.length || (value.length === 
JAVA_LONG_MAX.length && value <= JAVA_LONG_MAX))
+  );
+}
+
+export function requireRecord(
+  value: unknown,
+  allowedKeys: readonly string[]
+): asserts value is Record<string, unknown> {
+  if (value == null || typeof value !== 'object' || Array.isArray(value)) {
+    throw new Error('Investigation evidence must be a record');
+  }
+  if (Object.keys(value).some(key => !allowedKeys.includes(key))) {
+    throw new Error('Investigation evidence contains unsupported fields');
+  }
+}
+
+function hasControlCharacter(value: string) {
+  return [...value].some(character => {
+    const code = character.charCodeAt(0);
+    return code <= 31 || code === 127;
+  });
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to