Documentation
¶
Overview ¶
Package graphrag provides a layered in-memory graph for real-time observability retrieval — error chains, root cause analysis, impact analysis. It replaces the simpler internal/graph package with typed stores.
Index ¶
- Constants
- func AutoMigrateGraphRAG(db *gorm.DB) error
- func Preprocess(line string) string
- func SaveDrainTemplates(db *gorm.DB, tenant string, templates []Template) error
- func SetPanicMetrics(m *telemetry.Metrics)
- type AffectedEntry
- type AggregateSource
- type AnomalyNode
- type AnomalySeverity
- type AnomalyStore
- func (as *AnomalyStore) AddAnomaly(anomaly AnomalyNode)
- func (as *AnomalyStore) AddPrecededByEdge(anomalyID, precedingID string, ts time.Time)
- func (as *AnomalyStore) AnomaliesForService(service string, since time.Time) []*AnomalyNode
- func (as *AnomalyStore) AnomaliesSince(since time.Time) []*AnomalyNode
- func (as *AnomalyStore) AnomaliesSinceLimit(since time.Time, n int) []*AnomalyNode
- type AnomalyType
- type Config
- type CorrelatedSignalsResult
- type Coverage
- type Drain
- type DrainOption
- type DrainTemplateRow
- type Edge
- type EdgeType
- type ErrorChainResult
- type GraphRAG
- func (g *GraphRAG) AggregateMode() bool
- func (g *GraphRAG) AllServiceEdges(ctx context.Context) []*Edge
- func (g *GraphRAG) AnomaliesForService(ctx context.Context, service string, since time.Time) []*AnomalyNode
- func (g *GraphRAG) AnomalyTimeline(ctx context.Context, since time.Time) []*AnomalyNode
- func (g *GraphRAG) CorrelatedSignals(ctx context.Context, service string, since time.Time) *CorrelatedSignalsResult
- func (g *GraphRAG) DependencyChain(ctx context.Context, traceID string) []SpanNode
- func (g *GraphRAG) DrainTemplateCount() int
- func (g *GraphRAG) DroppedLogsCount() int64
- func (g *GraphRAG) DroppedMetricsCount() int64
- func (g *GraphRAG) DroppedSpansCount() int64
- func (g *GraphRAG) ErrorChain(ctx context.Context, service string, since time.Time, limit int) []ErrorChainResult
- func (g *GraphRAG) EventBufferDepth() int
- func (g *GraphRAG) GetInvestigation(ctx context.Context, id string) (*Investigation, error)
- func (g *GraphRAG) GetInvestigations(ctx context.Context, service, severity, status string, limit int) ([]Investigation, error)
- func (g *GraphRAG) ImpactAnalysis(ctx context.Context, service string, maxDepth int) *ImpactResult
- func (g *GraphRAG) InvestigationInsertCount() int64
- func (g *GraphRAG) IsRunning() bool
- func (g *GraphRAG) ObserveSpanTopology(tenant, traceID, spanID, parentSpanID, service string)
- func (g *GraphRAG) OnLogIngested(log storage.Log)
- func (g *GraphRAG) OnMetricIngested(metric tsdb.RawMetric)
- func (g *GraphRAG) OnSpanIngested(span storage.Span)
- func (g *GraphRAG) OnTemplateFact(fact aggregate.TemplateFact)
- func (g *GraphRAG) PersistInvestigation(tenant, triggerService string, chains []ErrorChainResult, ...)
- func (g *GraphRAG) RegisterAnomaly(tenant string, anomaly AnomalyNode)
- func (g *GraphRAG) RootCauseAnalysis(ctx context.Context, service string, since time.Time) []RankedCause
- func (g *GraphRAG) ServiceMap(ctx context.Context, depth int) []ServiceMapEntry
- func (g *GraphRAG) ServiceMapAround(ctx context.Context, seed string, depth int) []ServiceMapEntry
- func (g *GraphRAG) ServiceNames(ctx context.Context) []string
- func (g *GraphRAG) SetAggregateSource(src AggregateSource)
- func (g *GraphRAG) SetMetrics(m *telemetry.Metrics)
- func (g *GraphRAG) ShortestPath(ctx context.Context, from, to string) []string
- func (g *GraphRAG) Shutdown(ctx context.Context) error
- func (g *GraphRAG) SpanCapacityDropsCount() int64
- func (g *GraphRAG) Start(ctx context.Context)
- func (g *GraphRAG) Stop()
- func (g *GraphRAG) StoreCounts() StoreCounts
- func (g *GraphRAG) TenantsEvictedCount() int64
- func (g *GraphRAG) TraceGraph(ctx context.Context, traceID string) TraceGraphResult
- type ImpactResult
- type Investigation
- type LogClusterNode
- type MetricNode
- type NodeType
- type OperationNode
- type RankedCause
- type RootCauseInfo
- type ServiceMapEntry
- type ServiceNode
- type ServiceStore
- func (s *ServiceStore) AllEdges() []*Edge
- func (s *ServiceStore) AllServices() []*ServiceNode
- func (s *ServiceStore) CallEdgesFrom(service string) []*Edge
- func (s *ServiceStore) CallEdgesTo(service string) []*Edge
- func (s *ServiceStore) EnsureCallEdge(source, target string, ts time.Time) bool
- func (s *ServiceStore) EnsureService(name string, ts time.Time) bool
- func (s *ServiceStore) GetService(name string) (*ServiceNode, bool)
- func (s *ServiceStore) OperationsForService(service string) []*OperationNode
- func (s *ServiceStore) ReplaceTopology(services []*ServiceNode, operations []*OperationNode, edges []*Edge)
- func (s *ServiceStore) UpsertCallEdge(source, target string, durationMs float64, isError bool, ts time.Time)
- func (s *ServiceStore) UpsertOperation(service, operation string, durationMs float64, isError bool, ts time.Time)
- func (s *ServiceStore) UpsertService(name string, durationMs float64, isError bool, ts time.Time)
- type SignalStore
- func (ss *SignalStore) AddLoggedDuringEdge(clusterID, spanID string, ts time.Time)
- func (ss *SignalStore) LogClustersForService(service string) []*LogClusterNode
- func (ss *SignalStore) MetricsForService(service string) []*MetricNode
- func (ss *SignalStore) Prune(cutoff time.Time, maxMetrics, maxLogClusters int) int
- func (ss *SignalStore) ReplaceMetrics(metrics []*MetricNode)
- func (ss *SignalStore) UpsertLogCluster(id, template, severity, service string, ts time.Time)
- func (ss *SignalStore) UpsertLogClusterWithTemplate(id, template, severity, service string, templateID uint64, tokens []string, ...)
- func (ss *SignalStore) UpsertMetric(metricName, service string, value float64, ts time.Time)
- type SpanNode
- type StoreCounts
- type Template
- type TraceGraphResult
- type TraceNode
- type TraceStore
- func (ts *TraceStore) ErrorSpans(service string, since time.Time) []*SpanNode
- func (ts *TraceStore) GetSpan(spanID string) (*SpanNode, bool)
- func (ts *TraceStore) GetTrace(traceID string) (*TraceNode, bool)
- func (ts *TraceStore) Prune() int
- func (ts *TraceStore) SpansForTrace(traceID string) []*SpanNode
- func (ts *TraceStore) UpsertSpan(span SpanNode) bool
- func (ts *TraceStore) UpsertTrace(traceID, rootService, status string, durationMs float64, timestamp time.Time)
Constants ¶
const ( // CoverageRetained — the complete parent/child chain was retained. CoverageRetained = "retained_exemplar" // CoveragePartial — the exemplar exists but its chain is incomplete: // truncated by the per-trace span/byte bound, or still arriving, or aged // out of the in-memory trace TTL. CoveragePartial = "partial_exemplar" // CoverageNone — no exemplar. The honest answer is "not retained or not // found", never a chain assembled from whatever happened to be around. CoverageNone = "not_retained_or_not_found" )
Coverage source values. They are part of the MCP response contract: a client reads Source to know whether an answer is a fact about a complete retained exemplar, a fact about a deliberately truncated one, or an admission that nothing was retained.
const Wildcard = "<*>"
Wildcard is the placeholder used for variable positions within a template.
Variables ¶
This section is empty.
Functions ¶
func AutoMigrateGraphRAG ¶
AutoMigrateGraphRAG runs GORM auto-migration for GraphRAG models and applies tenant backfill + drain_templates composite-PK promotion. Safe to call repeatedly.
func Preprocess ¶
Preprocess masks common variable patterns in a raw log line prior to tokenization. This implements Drain's configurable regex-replacement stage.
func SaveDrainTemplates ¶
SaveDrainTemplates upserts the given templates into the drain_templates table under the supplied tenant. Tokens are JSON-encoded. Uses GORM's clause.OnConflict upsert which is dialect-aware (ON CONFLICT for SQLite/ PostgreSQL, ON DUPLICATE KEY UPDATE for MySQL).
The conflict target is the composite (tenant_id, id) primary key so the same Drain template hash can coexist across tenants — a future per-tenant Drain miner can rely on this to keep cluster IDs stable per tenant.
func SetPanicMetrics ¶
SetPanicMetrics wires the telemetry metrics so GraphRAG worker recovery closures can increment OtelContext_panics_recovered_total{subsystem="graphrag"}. Safe to leave unset in tests.
Types ¶
type AffectedEntry ¶
type AffectedEntry struct {
Service string `json:"service"`
Depth int `json:"depth"`
CallCount int64 `json:"call_count"`
ImpactScore float64 `json:"impact_score"`
}
AffectedEntry is a service affected by an upstream failure.
type AggregateSource ¶ added in v0.5.0
type AggregateSource interface {
// TopologyEpoch identifies the engine instance. A change means the
// revision counter restarted.
TopologyEpoch() uint64
// TopologyTenants lists the tenants the projection holds topology for.
TopologyTenants() []string
// TopologyRevision reports a tenant's revision without rendering a
// snapshot.
TopologyRevision(tenant string) uint64
// TopologySnapshot renders one tenant's topology.
TopologySnapshot(tenant string) aggregate.TopologySnapshot
// PruneTopology drops projection windows past the retention horizon.
PruneTopology()
}
AggregateSource is the read side of the aggregate engine that GraphRAG consumes. It is an interface so the coordinator depends on the four calls it makes rather than on the engine's whole surface, and so tests can drive reconciliation without building a store.
type AnomalyNode ¶
type AnomalyNode struct {
ID string `json:"id"`
Type AnomalyType `json:"type"`
Severity AnomalySeverity `json:"severity"`
Service string `json:"service"`
Evidence string `json:"evidence"`
Timestamp time.Time `json:"timestamp"`
}
AnomalyNode represents a detected anomaly.
type AnomalySeverity ¶
type AnomalySeverity string
AnomalySeverity indicates the severity of an anomaly.
const ( SeverityCritical AnomalySeverity = "critical" SeverityWarning AnomalySeverity = "warning" SeverityInfo AnomalySeverity = "info" )
type AnomalyStore ¶
type AnomalyStore struct {
Anomalies map[string]*AnomalyNode // key: anomaly ID
Edges map[string]*Edge // key: type|from|to
// contains filtered or unexported fields
}
AnomalyStore holds detected anomalies and their temporal correlations.
func (*AnomalyStore) AddAnomaly ¶
func (as *AnomalyStore) AddAnomaly(anomaly AnomalyNode)
func (*AnomalyStore) AddPrecededByEdge ¶
func (as *AnomalyStore) AddPrecededByEdge(anomalyID, precedingID string, ts time.Time)
func (*AnomalyStore) AnomaliesForService ¶
func (as *AnomalyStore) AnomaliesForService(service string, since time.Time) []*AnomalyNode
func (*AnomalyStore) AnomaliesSince ¶
func (as *AnomalyStore) AnomaliesSince(since time.Time) []*AnomalyNode
func (*AnomalyStore) AnomaliesSinceLimit ¶
func (as *AnomalyStore) AnomaliesSinceLimit(since time.Time, n int) []*AnomalyNode
AnomaliesSinceLimit is AnomaliesSince with a result cap (n <= 0 means unlimited). correlateWithRecent walks this on every detection tick, so the cap keeps a pathological anomaly backlog from turning each tick into an O(N) scan plus O(N) edge fan-out. Selection past the cap follows map iteration order — correlation is best-effort by design.
type AnomalyType ¶
type AnomalyType string
AnomalyType indicates the kind of anomaly detected.
const ( AnomalyErrorSpike AnomalyType = "error_spike" AnomalyLatencySpike AnomalyType = "latency_spike" AnomalyMetricZScore AnomalyType = "metric_zscore" )
type Config ¶
type Config struct {
TraceTTL time.Duration
RefreshEvery time.Duration
SnapshotEvery time.Duration
AnomalyEvery time.Duration
WorkerCount int
ChannelSize int
// MaxSpansPerTenant caps each tenant's in-memory TraceStore span map.
// 0 = defaultMaxSpansPerTenant; negative disables the cap.
MaxSpansPerTenant int
// TenantIdleTTL evicts a tenant's store slice after this much time
// without any ingest event or query. 0 = defaultTenantIdleTTL;
// negative disables eviction.
TenantIdleTTL time.Duration
// Mode is the aggregate mode (AGGREGATE_MODE): aggregate.ModeLegacy,
// ModeShadow or ModeAggregate. Empty means legacy. In shadow mode the only
// change is that log templates come from the ingest-owned miner; in
// aggregate mode the raw-span rebuild and the per-span topology upserts
// are retired in favour of engine snapshots (#163, #174).
Mode string
}
Config holds GraphRAG configuration.
type CorrelatedSignalsResult ¶
type CorrelatedSignalsResult struct {
}
CorrelatedSignals gathers all related signals for a service within a time range.
type Coverage ¶ added in v0.5.0
type Coverage struct {
// Complete is true only for a chain that reaches a root span with every
// parent link resolved.
Complete bool `json:"complete"`
// Truncated marks a chain cut short by the per-trace span or byte bound.
Truncated bool `json:"truncated"`
// RetainedSpans and ObservedSpans are the retained/observed counts when
// they are known.
RetainedSpans int `json:"retained_spans,omitempty"`
ObservedSpans int `json:"observed_spans,omitempty"`
// Source is one of the Coverage* constants.
Source string `json:"source"`
// Note explains a non-complete answer in plain words.
Note string `json:"note,omitempty"`
}
Coverage states what a causal-analysis answer is actually backed by.
Errors are always ELIGIBLE for retention; they are never all PERSISTED (#161, #163). ErrorChain, RootCauseAnalysis, DependencyChain and the trace_graph tool are guaranteed only for complete retained exemplars. A truncated exemplar reports explicit partial coverage and an absent one says so — fabricated certainty is the one answer that is never acceptable.
type Drain ¶
type Drain struct {
// contains filtered or unexported fields
}
Drain is a thread-safe log template miner.
func NewDrain ¶
func NewDrain(opts ...DrainOption) *Drain
NewDrain constructs a Drain miner with the supplied options.
func (*Drain) LoadTemplates ¶
LoadTemplates restores templates from a previous snapshot. Existing state is cleared. Intended for startup recovery. The LRU list is rebuilt by inserting templates in LastSeen order (oldest first, each PushFront'd), so the most recently-seen template ends up at the front.
func (*Drain) Match ¶
Match preprocesses the log line, descends the prefix tree, and either merges into an existing template or creates a new one. Returns the resulting template. Returns nil only for empty input.
func (*Drain) TemplateCount ¶
TemplateCount reports the number of live templates under the read lock.
type DrainOption ¶
type DrainOption func(*Drain)
DrainOption configures a Drain instance.
func WithDepth ¶
func WithDepth(depth int) DrainOption
WithDepth sets the prefix tree depth (default 4). Minimum 2.
func WithMaxChildren ¶
func WithMaxChildren(n int) DrainOption
WithMaxChildren sets the max distinct children per internal node (default 100). Beyond this, tokens collapse to Wildcard.
func WithMaxTemplates ¶
func WithMaxTemplates(n int) DrainOption
WithMaxTemplates sets the total template cap (default 50000). When exceeded, the least-recently-seen template is evicted.
func WithSimilarityThreshold ¶
func WithSimilarityThreshold(st float64) DrainOption
WithSimilarityThreshold sets st (default 0.4). Clamped to (0, 1].
type DrainTemplateRow ¶
type DrainTemplateRow struct {
TenantID string `gorm:"primaryKey;size:64;default:'default';not null" json:"tenant_id"`
ID int64 `gorm:"primaryKey;autoIncrement:false" json:"id"` // int64(Template.ID)
Tokens string `gorm:"type:text;not null" json:"tokens"` // JSON-encoded []string
Count int `json:"count"`
FirstSeen time.Time `gorm:"index" json:"first_seen"`
LastSeen time.Time `gorm:"index" json:"last_seen"`
Sample string `gorm:"type:text" json:"sample"`
}
DrainTemplateRow is the persisted GORM representation of a Drain log template. Tokens are JSON-encoded to stay schema-simple across SQLite/ MySQL/PostgreSQL/MSSQL.
ID is stored as int64 (bit-reinterpretation of the uint64 FNV-64 hash): the standard SQL drivers reject uint64 values with the high bit set, and signed int64 carries the same 64 bits without loss. Conversion happens in the persistence helpers.
The primary key is composite (tenant_id, id): the same template tokens can legitimately recur across tenants, and we want the cluster ID to stay stable per tenant once the in-memory Drain miner is partitioned per-tenant. TenantID is declared first so it leads the PK index.
func (DrainTemplateRow) TableName ¶
func (DrainTemplateRow) TableName() string
TableName overrides GORM's default table name.
type Edge ¶
type Edge struct {
Type EdgeType `json:"type"`
FromID string `json:"from_id"`
ToID string `json:"to_id"`
Weight float64 `json:"weight,omitempty"`
CallCount int64 `json:"call_count,omitempty"`
ErrorRate float64 `json:"error_rate,omitempty"`
AvgMs float64 `json:"avg_latency_ms,omitempty"`
TotalMs float64 `json:"-"`
ErrorCount int64 `json:"-"`
UpdatedAt time.Time `json:"updated_at"`
}
Edge represents a directed relationship between two nodes.
type EdgeType ¶
type EdgeType string
EdgeType distinguishes different relationship categories.
const ( EdgeCalls EdgeType = "CALLS" EdgeExposes EdgeType = "EXPOSES" EdgeContains EdgeType = "CONTAINS" EdgeChildOf EdgeType = "CHILD_OF" EdgeEmittedBy EdgeType = "EMITTED_BY" EdgeLoggedDuring EdgeType = "LOGGED_DURING" EdgeMeasuredBy EdgeType = "MEASURED_BY" EdgePrecededBy EdgeType = "PRECEDED_BY" EdgeTriggeredBy EdgeType = "TRIGGERED_BY" )
type ErrorChainResult ¶
type ErrorChainResult struct {
RootCause *RootCauseInfo `json:"root_cause"`
SpanChain []SpanNode `json:"span_chain"`
AnomalousMetrics []MetricNode `json:"anomalous_metrics,omitempty"`
TraceID string `json:"trace_id"`
// Coverage states whether this chain is backed by a complete retained
// exemplar. A chain whose upstream walk stopped at a missing parent is
// reported as partial, not presented as a root cause.
Coverage Coverage `json:"coverage"`
}
ErrorChainResult is the output of an error chain query.
type GraphRAG ¶
type GraphRAG struct {
// contains filtered or unexported fields
}
GraphRAG is the main coordinator for the layered graph system.
Every in-memory store is partitioned by tenant. The coordinator holds a map of tenant ID → *tenantStores and a reader/writer mutex that protects only the outer map; per-tenant stores keep their own RWMutexes for fine-grained concurrent access. All event ingestion and queries route through storesFor(ctx) / storesForTenant(tenant) — there is no "global" slice.
func New ¶
func New(repo *storage.Repository, tsdbAgg *tsdb.Aggregator, ringBuf *tsdb.RingBuffer, cfg Config) *GraphRAG
New creates a new GraphRAG coordinator.
The vectordb-backed semantic similarity path was removed on 2026-05-24 along with the find_similar_logs MCP tool — log clustering now relies solely on the Drain template miner (see drain.go).
func (*GraphRAG) AggregateMode ¶ added in v0.5.0
AggregateMode reports whether this coordinator consumes aggregate snapshots instead of rebuilding from raw spans.
func (*GraphRAG) AllServiceEdges ¶
AllServiceEdges returns every edge in the caller's tenant's ServiceStore. Kept as a narrow helper so API handlers do not need to traverse the tenantStores composite themselves.
func (*GraphRAG) AnomaliesForService ¶
func (g *GraphRAG) AnomaliesForService(ctx context.Context, service string, since time.Time) []*AnomalyNode
AnomaliesForService is a tenant-aware read-through to the per-tenant AnomalyStore, exported so handlers outside this package never need to reach into the store maps directly.
func (*GraphRAG) AnomalyTimeline ¶
AnomalyTimeline returns recent anomalies sorted by time.
func (*GraphRAG) CorrelatedSignals ¶
func (*GraphRAG) DependencyChain ¶
DependencyChain returns the full span tree for a trace.
func (*GraphRAG) DrainTemplateCount ¶
DrainTemplateCount reports the number of live Drain templates (bounded by maxTemplates, but the gauge proves it in production).
func (*GraphRAG) DroppedLogsCount ¶
DroppedLogsCount reports the number of log events dropped because the ingestion channel was full.
func (*GraphRAG) DroppedMetricsCount ¶
DroppedMetricsCount reports the number of metric events dropped because the ingestion channel was full.
func (*GraphRAG) DroppedSpansCount ¶
DroppedSpansCount reports the number of span events dropped because the ingestion channel was full. Exported for tests and readiness probes; atomic, safe from any goroutine.
func (*GraphRAG) ErrorChain ¶
func (g *GraphRAG) ErrorChain(ctx context.Context, service string, since time.Time, limit int) []ErrorChainResult
ErrorChain traces error spans upstream to find the root cause service. The tenant slice is selected via ctx — callers without a tenant ctx collapse to storage.DefaultTenantID at the coordinator boundary.
func (*GraphRAG) EventBufferDepth ¶
EventBufferDepth returns the current number of events queued in the ingestion channel. Exported for telemetry polling; never blocks.
func (*GraphRAG) GetInvestigation ¶
GetInvestigation retrieves a single investigation by ID, scoped to the tenant carried by ctx. Returning ErrRecordNotFound for cross-tenant lookups prevents id-guessing from leaking another tenant's row.
func (*GraphRAG) GetInvestigations ¶
func (g *GraphRAG) GetInvestigations(ctx context.Context, service, severity, status string, limit int) ([]Investigation, error)
GetInvestigations queries persisted investigations scoped to the tenant carried by ctx. The composite (tenant_id, created_at) index supports the recency-ordered scan.
func (*GraphRAG) ImpactAnalysis ¶
ImpactAnalysis performs BFS downstream from a service to find affected services.
func (*GraphRAG) InvestigationInsertCount ¶
InvestigationInsertCount reports cooldown-allowed PersistInvestigation calls. Semantics: this counter increments when the cooldown check passes, BEFORE the DB write — so a subsequent DB failure still increments this. It is NOT a strict DB insert count. Intended for tests to assert cooldown behavior without requiring a live repo.
func (*GraphRAG) IsRunning ¶
IsRunning reports whether the coordinator's stop channel has not been closed. Used by readiness probes to confirm the background workers are still live.
func (*GraphRAG) ObserveSpanTopology ¶ added in v0.3.1
ObserveSpanTopology records one received span's cross-service call topology INDEPENDENT of the sampler. The OTLP trace receiver calls this for every span BEFORE the sampler's keep/drop decision, so the service map retains flow direction even when sampling drops the spans that would have formed the edge.
tenant is the resolved tenant (empty collapses to DefaultTenantID, matching the rest of GraphRAG). traceID is accepted for symmetry with the event path and possible future per-trace scoping; the current implementation keys solely on span IDs, which are globally unique within a tenant. Safe to call from many goroutines.
func (*GraphRAG) OnLogIngested ¶
OnLogIngested is the callback wired into the log ingestion pipeline.
func (*GraphRAG) OnMetricIngested ¶
OnMetricIngested is the callback wired into the metric ingestion pipeline. tsdb.RawMetric already carries a resolved TenantID (set in ingest/otlp.go Export), so we read it here instead of adding a second argument — keeping the metric callback signature identical across TSDB and GraphRAG.
func (*GraphRAG) OnSpanIngested ¶
OnSpanIngested is the callback wired into the trace ingestion pipeline. Tenant is taken straight from the persisted Span (already resolved upstream by the OTLP Export handlers) and carried on the event — the callback signature is intentionally unchanged so external wiring stays trivial.
func (*GraphRAG) OnTemplateFact ¶ added in v0.5.0
func (g *GraphRAG) OnTemplateFact(fact aggregate.TemplateFact)
OnTemplateFact consumes one ingest-mined log template fact. It is the aggregate/shadow replacement for GraphRAG's own Drain miner: the template ID is minted once, on the ingest path, so a template never has two identities (#163).
The miner calls this SYNCHRONOUSLY on the OTLP goroutine, once per log line, so it must not take a store lock. It enqueues onto the same best-effort channel every other ingest callback uses and returns; a full channel drops the fact, exactly as a full channel drops a span.
func (*GraphRAG) PersistInvestigation ¶
func (g *GraphRAG) PersistInvestigation(tenant, triggerService string, chains []ErrorChainResult, anomalies []*AnomalyNode)
PersistInvestigation saves an investigation record from an error chain analysis. Tenant is accepted explicitly so the caller (the per-tenant anomaly loop) can re-enter ImpactAnalysis on the correct tenant slice and so the persisted row carries its originating tenant_id.
func (*GraphRAG) RegisterAnomaly ¶
func (g *GraphRAG) RegisterAnomaly(tenant string, anomaly AnomalyNode)
RegisterAnomaly inserts an anomaly into the AnomalyStore for tenant. Mirrors PersistInvestigation's "tenant accepted explicitly" shape so out-of-band anomaly producers (synthetic detectors, integration tests, future external anomaly feeds) can land directly on the right tenant slice without going through the metric/error detection loops. Empty tenant collapses to storage.DefaultTenantID.
func (*GraphRAG) RootCauseAnalysis ¶
func (g *GraphRAG) RootCauseAnalysis(ctx context.Context, service string, since time.Time) []RankedCause
RootCauseAnalysis combines ErrorChain with anomaly correlation to rank probable causes.
func (*GraphRAG) ServiceMap ¶
func (g *GraphRAG) ServiceMap(ctx context.Context, depth int) []ServiceMapEntry
func (*GraphRAG) ServiceMapAround ¶ added in v0.5.0
ServiceMapAround returns the subgraph reachable from seed within depth hops in either direction (downstream via CALLS-from, upstream via CALLS-to). It bounds both the compute and the payload for focused "map around service X" queries at 100–200 services; callers wanting the full map use ServiceMap. depth<=0 falls back to a single hop so a focus query always returns the seed + neighbours.
func (*GraphRAG) ServiceNames ¶
ServiceNames returns every service the caller's tenant has emitted any span for, sorted ascending. Reads from the in-memory ServiceStore — no DB scan. Used by /api/metadata/services so the dropdown matches the system map exactly (both are sourced from the same store).
func (*GraphRAG) SetAggregateSource ¶ added in v0.5.0
func (g *GraphRAG) SetAggregateSource(src AggregateSource)
SetAggregateSource wires the aggregate engine's topology projection. Main calls it before Start; the lock also keeps tests and alternate construction paths from racing the refresh loop.
func (*GraphRAG) SetMetrics ¶
SetMetrics wires the Prometheus registry so GraphRAG event drops are observable via otelcontext_graphrag_events_dropped_total. Safe to call before Start; pass nil to disable Prometheus recording (atomic counters still tick).
func (*GraphRAG) ShortestPath ¶
ShortestPath finds the shortest path between two services using Dijkstra.
func (*GraphRAG) Shutdown ¶ added in v0.5.0
Shutdown closes event admission, drains accepted events, joins every worker, and only then persists the final Drain template state.
func (*GraphRAG) SpanCapacityDropsCount ¶
SpanCapacityDropsCount reports the number of spans skipped because the tenant's TraceStore was at its MaxSpans cap. Atomic, safe from any goroutine; exported for tests and readiness probes.
func (*GraphRAG) Start ¶
Start begins background goroutines: workers, refresh, snapshot, anomaly detection. Each goroutine is wrapped in a panic recovery so one misbehaving event can't take down the whole subsystem.
func (*GraphRAG) Stop ¶
func (g *GraphRAG) Stop()
Stop preserves the prior internal lifecycle surface. Production shutdown uses Shutdown so deadline and persistence failures reach the process result.
func (*GraphRAG) StoreCounts ¶
func (g *GraphRAG) StoreCounts() StoreCounts
StoreCounts walks a tenant snapshot and takes len() under each store's read lock — no slice building, cheap enough for a periodic sampler.
func (*GraphRAG) TenantsEvictedCount ¶
TenantsEvictedCount reports the number of tenant store slices evicted for exceeding the idle TTL since startup.
func (*GraphRAG) TraceGraph ¶ added in v0.5.0
func (g *GraphRAG) TraceGraph(ctx context.Context, traceID string) TraceGraphResult
TraceGraph returns the retained span tree for a trace together with an explicit statement of what it covers.
The tree is assembled solely from retained exemplars. A trace nobody retained returns CoverageNone and an empty span list — "not retained or not found" is the honest answer, and it is a successful response, not an error.
type ImpactResult ¶
type ImpactResult struct {
Service string `json:"service"`
AffectedServices []AffectedEntry `json:"affected_services"`
TotalDownstream int `json:"total_downstream"`
}
ImpactResult describes the blast radius of a service failure.
type Investigation ¶
type Investigation struct {
TenantID string `gorm:"size:64;default:'default';not null;index:idx_investigations_tenant_created,priority:1" json:"tenant_id"`
ID string `gorm:"primaryKey;size:64" json:"id"`
CreatedAt time.Time `gorm:"index:idx_investigations_tenant_created,priority:2" json:"created_at"`
Status string `gorm:"size:20" json:"status"` // detected, triaged, resolved
Severity string `gorm:"size:20" json:"severity"` // critical, warning, info
TriggerService string `gorm:"size:255;index" json:"trigger_service"`
TriggerOperation string `gorm:"size:255" json:"trigger_operation"`
ErrorMessage string `gorm:"type:text" json:"error_message"`
RootService string `gorm:"size:255" json:"root_service"`
RootOperation string `gorm:"size:255" json:"root_operation"`
CausalChain json.RawMessage `gorm:"type:text" json:"causal_chain"`
TraceIDs json.RawMessage `gorm:"type:text" json:"trace_ids"`
ErrorLogs json.RawMessage `gorm:"type:text" json:"error_logs"`
AnomalousMetrics json.RawMessage `gorm:"type:text" json:"anomalous_metrics"`
AffectedServices json.RawMessage `gorm:"type:text" json:"affected_services"`
SpanChain json.RawMessage `gorm:"type:text" json:"span_chain"`
}
Investigation is a persisted record of an automated error investigation.
TenantID scopes the row to its originating tenant. The composite (tenant_id, created_at) index supports the recency-ordered "investigations for tenant X" query that GetInvestigations runs on every read.
func (Investigation) TableName ¶
func (Investigation) TableName() string
TableName overrides GORM's default table name.
type LogClusterNode ¶
type LogClusterNode struct {
ID string `json:"id"` // service-scoped cluster id (stable)
Template string `json:"template"`
TemplateID uint64 `json:"template_id,omitempty"`
TemplateTokens []string `json:"template_tokens,omitempty"`
SampleLog string `json:"sample_log,omitempty"`
Count int64 `json:"count"`
FirstSeen time.Time `json:"first_seen"`
LastSeen time.Time `json:"last_seen"`
SeverityDist map[string]int64 `json:"severity_distribution"`
}
LogClusterNode groups similar log messages.
Clustering is performed by the Drain template miner. TemplateID is the stable FNV-64 hash of TemplateTokens; ID is the user-facing cluster identifier (service-scoped) that remains stable across Drain re-merges.
type MetricNode ¶
type MetricNode struct {
ID string `json:"id"` // metric_name + "|" + service
MetricName string `json:"metric_name"`
Service string `json:"service"`
RollingMin float64 `json:"rolling_min"`
RollingMax float64 `json:"rolling_max"`
RollingAvg float64 `json:"rolling_avg"`
SampleCount int64 `json:"sample_count"`
LastSeen time.Time `json:"last_seen"`
}
MetricNode represents a metric series for a service.
type OperationNode ¶
type OperationNode struct {
ID string `json:"id"` // service + "|" + operation
Service string `json:"service"`
Operation string `json:"operation"`
FirstSeen time.Time `json:"first_seen"`
LastSeen time.Time `json:"last_seen"`
HealthScore float64 `json:"health_score"`
CallCount int64 `json:"call_count"`
ErrorCount int64 `json:"error_count"`
ErrorRate float64 `json:"error_rate"`
AvgLatency float64 `json:"avg_latency_ms"`
P50Latency float64 `json:"p50_latency_ms"`
P95Latency float64 `json:"p95_latency_ms"`
P99Latency float64 `json:"p99_latency_ms"`
LatencyProvenance *latency.Provenance `json:"latency_provenance,omitempty"`
TotalMs float64 `json:"-"`
}
OperationNode represents an endpoint/RPC within a service.
type RankedCause ¶
type RankedCause struct {
Service string `json:"service"`
Operation string `json:"operation"`
Score float64 `json:"score"`
Evidence []string `json:"evidence"`
ErrorChain []SpanNode `json:"error_chain,omitempty"`
Anomalies []AnomalyNode `json:"anomalies,omitempty"`
// Coverage states what the ranking rests on: a complete retained
// exemplar, a partial one, or aggregate anomaly evidence with no chain.
Coverage Coverage `json:"coverage"`
}
RankedCause is a probable root cause with evidence.
type RootCauseInfo ¶
type RootCauseInfo struct {
Service string `json:"service"`
Operation string `json:"operation"`
ErrorMessage string `json:"error_message"`
SpanID string `json:"span_id"`
TraceID string `json:"trace_id"`
}
RootCauseInfo identifies the responsible service and operation.
type ServiceMapEntry ¶
type ServiceMapEntry struct {
Service *ServiceNode `json:"service"`
Operations []*OperationNode `json:"operations,omitempty"`
CallsTo []*Edge `json:"calls_to,omitempty"`
CalledBy []*Edge `json:"called_by,omitempty"`
}
ServiceMap returns the service topology with health scores for the API.
type ServiceNode ¶
type ServiceNode struct {
ID string `json:"id"`
Name string `json:"name"`
FirstSeen time.Time `json:"first_seen"`
LastSeen time.Time `json:"last_seen"`
HealthScore float64 `json:"health_score"` // 0.0–1.0
CallCount int64 `json:"call_count"`
ErrorCount int64 `json:"error_count"`
ErrorRate float64 `json:"error_rate"`
AvgLatency float64 `json:"avg_latency_ms"`
P95Latency float64 `json:"p95_latency_ms,omitempty"`
P99Latency float64 `json:"p99_latency_ms,omitempty"`
LatencyProvenance *latency.Provenance `json:"latency_provenance,omitempty"`
TotalMs float64 `json:"-"` // for computing avg
}
ServiceNode represents a microservice with aggregated health stats.
type ServiceStore ¶
type ServiceStore struct {
Services map[string]*ServiceNode // key: service name
Operations map[string]*OperationNode // key: service|operation
Edges map[string]*Edge // key: type|from|to
// contains filtered or unexported fields
}
ServiceStore holds permanent service topology data.
func (*ServiceStore) AllEdges ¶
func (s *ServiceStore) AllEdges() []*Edge
func (*ServiceStore) AllServices ¶
func (s *ServiceStore) AllServices() []*ServiceNode
func (*ServiceStore) CallEdgesFrom ¶
func (s *ServiceStore) CallEdgesFrom(service string) []*Edge
func (*ServiceStore) CallEdgesTo ¶
func (s *ServiceStore) CallEdgesTo(service string) []*Edge
func (*ServiceStore) EnsureCallEdge ¶ added in v0.3.1
func (s *ServiceStore) EnsureCallEdge(source, target string, ts time.Time) bool
EnsureCallEdge guarantees a CALLS edge source→target EXISTS without touching its aggregates. If the edge is absent it is created with zeroed CallCount/latency/error stats (UpdatedAt = ts); if it already exists this is a no-op. The pre-sample topology observer uses this so the service map shows flow direction even when sampling dropped every span that would have formed the edge — while the sampled path (UpsertCallEdge) remains the sole source of CallCount/latency/error-rate aggregates. Returns true if a new edge was created.
func (*ServiceStore) EnsureService ¶ added in v0.3.1
func (s *ServiceStore) EnsureService(name string, ts time.Time) bool
EnsureService guarantees a ServiceNode for name EXISTS without touching its aggregates. Absent → created with zeroed call/error stats (FirstSeen/LastSeen = ts); present → no-op. The pre-sample topology observer uses this so the service map has a row to hang flow-direction edges on even when sampling dropped every span for the service; the sampled path (UpsertService) and the 60s DB rebuild remain the sole source of call/error/latency aggregates. Returns true if a new node was created.
func (*ServiceStore) GetService ¶
func (s *ServiceStore) GetService(name string) (*ServiceNode, bool)
func (*ServiceStore) OperationsForService ¶ added in v0.5.0
func (s *ServiceStore) OperationsForService(service string) []*OperationNode
OperationsForService returns the operations exposed by a service via the per-service index, avoiding a full Operations-map scan per service.
func (*ServiceStore) ReplaceTopology ¶ added in v0.5.0
func (s *ServiceStore) ReplaceTopology(services []*ServiceNode, operations []*OperationNode, edges []*Edge)
ReplaceTopology swaps the entire service topology for a new one under a single write lock.
Aggregate mode consumes a per-revision snapshot from the aggregate engine rather than accumulating per-span upserts, so the only correct way to apply it is replacement: re-applying a cumulative snapshot through UpsertService / UpsertCallEdge would multiply every counter by the number of ticks it survived. Rebuilding the adjacency indexes here (rather than mutating them) is what keeps their append-only invariant intact — after the swap they describe exactly the maps they were built from.
func (*ServiceStore) UpsertCallEdge ¶
func (*ServiceStore) UpsertOperation ¶
func (*ServiceStore) UpsertService ¶
type SignalStore ¶
type SignalStore struct {
LogClusters map[string]*LogClusterNode // key: cluster ID
Metrics map[string]*MetricNode // key: metric|service
Edges map[string]*Edge // key: type|from|to
// contains filtered or unexported fields
}
SignalStore holds log cluster and metric correlation data.
func (*SignalStore) AddLoggedDuringEdge ¶
func (ss *SignalStore) AddLoggedDuringEdge(clusterID, spanID string, ts time.Time)
func (*SignalStore) LogClustersForService ¶
func (ss *SignalStore) LogClustersForService(service string) []*LogClusterNode
func (*SignalStore) MetricsForService ¶
func (ss *SignalStore) MetricsForService(service string) []*MetricNode
func (*SignalStore) Prune ¶
func (ss *SignalStore) Prune(cutoff time.Time, maxMetrics, maxLogClusters int) int
Prune bounds the SignalStore (modeled on TraceStore.Prune): MetricNodes and LogClusterNodes whose LastSeen predates cutoff are removed; if either map still exceeds its cap (maxMetrics / maxLogClusters; <=0 = uncapped) the oldest-LastSeen overflow is evicted. Each removed node takes its edges with it — metric MEASURED_BY edges inline, evicted-cluster EMITTED_BY / LOGGED_DURING edges in the final sweep — and any edge whose UpdatedAt predates cutoff is swept too (upsert paths refresh edge timestamps, so live correlations survive). Returns the number of nodes (metrics + clusters) removed.
func (*SignalStore) ReplaceMetrics ¶ added in v0.5.0
func (ss *SignalStore) ReplaceMetrics(metrics []*MetricNode)
ReplaceMetrics swaps the metric nodes (and their MEASURED_BY edges) for a new set. Same reasoning as ReplaceTopology: aggregate-mode metric state is a projection of a revision, not an accumulation. Log-cluster nodes and their edges are untouched — they arrive as ingest-owned template facts on their own schedule.
func (*SignalStore) UpsertLogCluster ¶
func (ss *SignalStore) UpsertLogCluster(id, template, severity, service string, ts time.Time)
func (*SignalStore) UpsertLogClusterWithTemplate ¶
func (ss *SignalStore) UpsertLogClusterWithTemplate(id, template, severity, service string, templateID uint64, tokens []string, sample string, ts time.Time)
UpsertLogClusterWithTemplate is the Drain-aware upsert. It stores the mined template tokens, the stable template ID, and a sample raw log. The older UpsertLogCluster is preserved for backward compatibility.
func (*SignalStore) UpsertMetric ¶
func (ss *SignalStore) UpsertMetric(metricName, service string, value float64, ts time.Time)
type SpanNode ¶
type SpanNode struct {
ID string `json:"id"` // span_id
TraceID string `json:"trace_id"`
ParentSpanID string `json:"parent_span_id"`
Service string `json:"service"`
Operation string `json:"operation"`
Duration float64 `json:"duration_ms"`
StatusCode string `json:"status_code"`
IsError bool `json:"is_error"`
Timestamp time.Time `json:"timestamp"`
}
SpanNode represents a single span within a trace.
type StoreCounts ¶
type StoreCounts struct {
Tenants int
Services int
Operations int
// LatencySketches counts the per-service duration sketches (#291), each
// a fixed-size aggregate.Sketch value.
LatencySketches int
Traces int
Spans int
LogClusters int
Metrics int
Anomalies int
ServiceEdges int
TraceEdges int
SignalEdges int
AnomalyEdges int
}
StoreCounts is a point-in-time census of every long-lived in-memory structure the coordinator owns, aggregated across tenants. It feeds the otelcontext_graphrag_* gauges so operators can attribute RSS growth to a specific store before reaching for a heap profile.
type Template ¶
type Template struct {
ID uint64 `json:"id"`
Tokens []string `json:"tokens"`
Count int `json:"count"`
FirstSeen time.Time `json:"first_seen"`
LastSeen time.Time `json:"last_seen"`
Sample string `json:"sample,omitempty"`
// contains filtered or unexported fields
}
Template represents a mined log group.
func LoadDrainTemplates ¶
LoadDrainTemplates reads persisted Drain templates for the supplied tenant and returns them in a format ready to pass to Drain.LoadTemplates. Returns an empty slice (and nil error) if no rows match.
func (*Template) TemplateString ¶
TemplateString returns the template rendered as a single string.
type TraceGraphResult ¶ added in v0.5.0
type TraceGraphResult struct {
TraceID string `json:"trace_id"`
Spans []SpanNode `json:"spans"`
Coverage Coverage `json:"coverage"`
}
TraceGraphResult is the trace_graph tool's response shape: the retained span tree plus an explicit statement of what it covers.
type TraceNode ¶
type TraceNode struct {
ID string `json:"id"` // trace_id
RootService string `json:"root_service"`
Duration float64 `json:"duration_ms"`
Status string `json:"status"`
Timestamp time.Time `json:"timestamp"`
SpanCount int `json:"span_count"`
}
TraceNode represents a distributed trace.
type TraceStore ¶
type TraceStore struct {
Traces map[string]*TraceNode // key: trace_id
Spans map[string]*SpanNode // key: span_id
Edges map[string]*Edge // key: type|from|to
TTL time.Duration
// MaxSpans hard-caps the Spans map: at the cap, NEW span IDs are
// skipped (UpsertSpan returns false) while updates to resident IDs
// still apply. <=0 disables the cap.
MaxSpans int
// contains filtered or unexported fields
}
TraceStore holds trace/span detail with TTL-based pruning.
func (*TraceStore) ErrorSpans ¶
func (ts *TraceStore) ErrorSpans(service string, since time.Time) []*SpanNode
func (*TraceStore) Prune ¶
func (ts *TraceStore) Prune() int
Prune removes spans and traces older than TTL.
func (*TraceStore) SpansForTrace ¶
func (ts *TraceStore) SpansForTrace(traceID string) []*SpanNode
func (*TraceStore) UpsertSpan ¶
func (ts *TraceStore) UpsertSpan(span SpanNode) bool
UpsertSpan inserts or updates a span node and its CONTAINS/CHILD_OF edges. Returns false when the span is NEW and MaxSpans is already reached — the span is skipped entirely (the graph is best-effort; the DB is the source of truth, same doctrine as the event-channel overflow in builder.go).
func (*TraceStore) UpsertTrace ¶
func (ts *TraceStore) UpsertTrace(traceID, rootService, status string, durationMs float64, timestamp time.Time)