Files
ITWorx Pulse release export bd774932d5
Public source validation / validate (push) Failing after 3m8s
Publish ITWorx Pulse source
2026-09-03 02:09:19 +02:00

283 lines
13 KiB
Go

package inventory
import (
"context"
"encoding/json"
"errors"
"fmt"
"sort"
"strings"
"time"
)
type EntityFilter struct {
Limit int
AfterName, AfterID string
Search string
EntityType string
Status string
Direction string
}
type EntitySummary struct {
ID string `json:"id"`
EntityType string `json:"entityType"`
CanonicalName string `json:"canonicalName"`
DisplayName string `json:"displayName"`
Status string `json:"status"`
FirstSeenAt *time.Time `json:"firstSeenAt"`
LastSeenAt *time.Time `json:"lastSeenAt,omitempty"`
TombstonedAt *time.Time `json:"tombstonedAt,omitempty"`
FactCount int `json:"factCount"`
OverrideCount int `json:"overrideCount"`
RelationCount int `json:"relationCount"`
SourceCount int `json:"sourceCount"`
StaleFactCount int `json:"staleFactCount"`
}
type AliasView struct {
SourceID string `json:"sourceId"`
SourceName string `json:"sourceName"`
ExternalType string `json:"externalType"`
ExternalID string `json:"externalId"`
}
type FactView struct {
FieldName string `json:"fieldName"`
SourceID string `json:"sourceId"`
SourceName string `json:"sourceName"`
Value json.RawMessage `json:"value"`
ObservedAt time.Time `json:"observedAt"`
Confidence float64 `json:"confidence"`
ValidUntil *time.Time `json:"validUntil,omitempty"`
Stale bool `json:"stale"`
}
type OverrideView struct {
FieldName string `json:"fieldName"`
Value json.RawMessage `json:"value"`
UserID string `json:"userId,omitempty"`
UpdatedAt time.Time `json:"updatedAt"`
}
type EffectiveValue struct {
FieldName string `json:"fieldName"`
Value json.RawMessage `json:"value"`
Origin string `json:"origin"`
SourceID string `json:"sourceId,omitempty"`
SourceName string `json:"sourceName,omitempty"`
ObservedAt *time.Time `json:"observedAt,omitempty"`
Confidence *float64 `json:"confidence,omitempty"`
Stale bool `json:"stale"`
OverriddenAt *time.Time `json:"overriddenAt,omitempty"`
}
type RelationView struct {
ID string `json:"id"`
Direction string `json:"direction"`
RelationType string `json:"relationType"`
PeerID string `json:"peerId"`
PeerType string `json:"peerType"`
PeerName string `json:"peerName"`
PeerStatus string `json:"peerStatus"`
PeerTombstonedAt *time.Time `json:"peerTombstonedAt,omitempty"`
SourceID string `json:"sourceId"`
SourceName string `json:"sourceName"`
Confidence float64 `json:"confidence"`
Confirmed bool `json:"confirmed"`
FirstSeenAt *time.Time `json:"firstSeenAt"`
LastSeenAt *time.Time `json:"lastSeenAt,omitempty"`
TombstonedAt *time.Time `json:"tombstonedAt,omitempty"`
}
type EntityDetail struct {
Entity EntitySummary `json:"entity"`
Aliases []AliasView `json:"aliases"`
Facts []FactView `json:"facts"`
Overrides []OverrideView `json:"overrides"`
Effective []EffectiveValue `json:"effectiveValues"`
Relations []RelationView `json:"relations"`
}
func (f EntityFilter) Validate() error {
if f.Limit < 1 || f.Limit > 100 {
return errors.New("entity page limit must be between 1 and 100")
}
if len(f.Search) > 120 || len(f.EntityType) > 80 || len(f.Status) > 40 {
return errors.New("entity filters exceed bounds")
}
if f.Direction != "asc" && f.Direction != "desc" {
return errors.New("entity sort direction must be asc or desc")
}
return nil
}
func (r *Repository) SearchEntities(ctx context.Context, filter EntityFilter) ([]EntitySummary, error) {
if err := filter.Validate(); err != nil {
return nil, err
}
rows, err := r.pool.Query(ctx, `WITH projected AS (
SELECT e.id::text,e.entity_type,e.canonical_name,
COALESCE(NULLIF(od.value #>> '{}',''),e.display_name) AS effective_name,
COALESCE(NULLIF(os.value #>> '{}',''),e.status) AS effective_status,
e.first_seen_at,e.last_seen_at,e.tombstoned_at,
(SELECT count(*) FROM entity_facts f WHERE f.entity_id=e.id)::int AS fact_count,
(SELECT count(*) FROM entity_overrides o WHERE o.entity_id=e.id)::int AS override_count,
(SELECT count(*) FROM entity_relations rel WHERE rel.source_entity_id=e.id OR rel.target_entity_id=e.id)::int AS relation_count,
(SELECT count(DISTINCT source_id) FROM entity_facts f WHERE f.entity_id=e.id)::int AS source_count,
(SELECT count(*) FROM entity_facts f WHERE f.entity_id=e.id AND f.valid_until IS NOT NULL AND f.valid_until < now())::int AS stale_fact_count
FROM entities e
LEFT JOIN entity_overrides od ON od.entity_id=e.id AND od.field_name='displayName'
LEFT JOIN entity_overrides os ON os.entity_id=e.id AND os.field_name='status'
WHERE e.tombstoned_at IS NULL)
SELECT id,entity_type,canonical_name,effective_name,effective_status,first_seen_at,last_seen_at,tombstoned_at,fact_count,override_count,relation_count,source_count,stale_fact_count
FROM projected
WHERE ($1='' OR lower(effective_name) LIKE '%'||lower($1)||'%' OR lower(canonical_name) LIKE '%'||lower($1)||'%')
AND ($2='' OR entity_type=$2) AND ($3='' OR effective_status=$3)
AND ($4='' OR ($6='asc' AND (lower(effective_name),id) > (lower($4),$5)) OR ($6='desc' AND (lower(effective_name),id) < (lower($4),$5)))
ORDER BY CASE WHEN $6='asc' THEN lower(effective_name) END ASC, CASE WHEN $6='desc' THEN lower(effective_name) END DESC,
CASE WHEN $6='asc' THEN id END ASC, CASE WHEN $6='desc' THEN id END DESC LIMIT $7`, strings.TrimSpace(filter.Search), filter.EntityType, filter.Status, filter.AfterName, filter.AfterID, filter.Direction, filter.Limit+1)
if err != nil {
return nil, fmt.Errorf("search entities: %w", err)
}
defer rows.Close()
result := make([]EntitySummary, 0, filter.Limit+1)
for rows.Next() {
var item EntitySummary
if err := rows.Scan(&item.ID, &item.EntityType, &item.CanonicalName, &item.DisplayName, &item.Status, &item.FirstSeenAt, &item.LastSeenAt, &item.TombstonedAt, &item.FactCount, &item.OverrideCount, &item.RelationCount, &item.SourceCount, &item.StaleFactCount); err != nil {
return nil, fmt.Errorf("scan entity search: %w", err)
}
result = append(result, item)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("search entity rows: %w", err)
}
return result, nil
}
func (r *Repository) GetEntityDetail(ctx context.Context, id string) (EntityDetail, error) {
rows, err := r.SearchEntitiesByID(ctx, id)
if err != nil {
return EntityDetail{}, err
}
if len(rows) == 0 {
return EntityDetail{}, fmt.Errorf("get entity detail: %w", ErrNotFound)
}
detail := EntityDetail{Entity: rows[0], Aliases: []AliasView{}, Facts: []FactView{}, Overrides: []OverrideView{}, Effective: []EffectiveValue{}, Relations: []RelationView{}}
aliasRows, err := r.pool.Query(ctx, `SELECT a.source_id::text,s.name,a.external_type,a.external_id FROM entity_aliases a JOIN data_sources s ON s.id=a.source_id WHERE a.entity_id=$1 ORDER BY s.name,a.external_type,a.external_id`, id)
if err != nil {
return EntityDetail{}, fmt.Errorf("list entity aliases: %w", err)
}
for aliasRows.Next() {
var item AliasView
if err := aliasRows.Scan(&item.SourceID, &item.SourceName, &item.ExternalType, &item.ExternalID); err != nil {
aliasRows.Close()
return EntityDetail{}, err
}
detail.Aliases = append(detail.Aliases, item)
}
aliasRows.Close()
if err := aliasRows.Err(); err != nil {
return EntityDetail{}, err
}
factRows, err := r.pool.Query(ctx, `SELECT f.field_name,f.source_id::text,s.name,f.value,f.observed_at,f.confidence,f.valid_until,(f.valid_until IS NOT NULL AND f.valid_until<now()) FROM entity_facts f JOIN data_sources s ON s.id=f.source_id WHERE f.entity_id=$1 ORDER BY f.field_name,(f.valid_until IS NOT NULL AND f.valid_until<now()),f.confidence DESC,f.observed_at DESC,f.source_id`, id)
if err != nil {
return EntityDetail{}, fmt.Errorf("list entity fact views: %w", err)
}
for factRows.Next() {
var item FactView
if err := factRows.Scan(&item.FieldName, &item.SourceID, &item.SourceName, &item.Value, &item.ObservedAt, &item.Confidence, &item.ValidUntil, &item.Stale); err != nil {
factRows.Close()
return EntityDetail{}, err
}
detail.Facts = append(detail.Facts, item)
}
factRows.Close()
if err := factRows.Err(); err != nil {
return EntityDetail{}, err
}
overrideRows, err := r.pool.Query(ctx, `SELECT field_name,value,COALESCE(user_id::text,''),updated_at FROM entity_overrides WHERE entity_id=$1 ORDER BY field_name`, id)
if err != nil {
return EntityDetail{}, fmt.Errorf("list entity overrides: %w", err)
}
for overrideRows.Next() {
var item OverrideView
if err := overrideRows.Scan(&item.FieldName, &item.Value, &item.UserID, &item.UpdatedAt); err != nil {
overrideRows.Close()
return EntityDetail{}, err
}
detail.Overrides = append(detail.Overrides, item)
}
overrideRows.Close()
if err := overrideRows.Err(); err != nil {
return EntityDetail{}, err
}
detail.Effective = effectiveValues(detail.Facts, detail.Overrides)
relationRows, err := r.pool.Query(ctx, `SELECT rel.id::text,CASE WHEN rel.source_entity_id=$1 THEN 'outgoing' ELSE 'incoming' END,rel.relation_type,peer.id::text,peer.entity_type,peer.display_name,CASE WHEN peer.tombstoned_at IS NULL THEN peer.status ELSE 'Ontbrekende entiteit' END,peer.tombstoned_at,rel.source_id::text,s.name,rel.confidence,rel.confirmed,rel.first_seen_at,rel.last_seen_at,rel.tombstoned_at FROM entity_relations rel JOIN entities peer ON peer.id=CASE WHEN rel.source_entity_id=$1 THEN rel.target_entity_id ELSE rel.source_entity_id END JOIN data_sources s ON s.id=rel.source_id WHERE rel.source_entity_id=$1 OR rel.target_entity_id=$1 ORDER BY rel.relation_type,peer.display_name,rel.id`, id)
if err != nil {
return EntityDetail{}, fmt.Errorf("list entity relation views: %w", err)
}
for relationRows.Next() {
var item RelationView
if err := relationRows.Scan(&item.ID, &item.Direction, &item.RelationType, &item.PeerID, &item.PeerType, &item.PeerName, &item.PeerStatus, &item.PeerTombstonedAt, &item.SourceID, &item.SourceName, &item.Confidence, &item.Confirmed, &item.FirstSeenAt, &item.LastSeenAt, &item.TombstonedAt); err != nil {
relationRows.Close()
return EntityDetail{}, err
}
detail.Relations = append(detail.Relations, item)
}
relationRows.Close()
if err := relationRows.Err(); err != nil {
return EntityDetail{}, err
}
return detail, nil
}
var ErrNotFound = errors.New("inventory entity not found")
func (r *Repository) SearchEntitiesByID(ctx context.Context, id string) ([]EntitySummary, error) {
rows, err := r.pool.Query(ctx, `SELECT e.id::text,e.entity_type,e.canonical_name,COALESCE(NULLIF(od.value #>> '{}',''),e.display_name),COALESCE(NULLIF(os.value #>> '{}',''),e.status),e.first_seen_at,e.last_seen_at,e.tombstoned_at,(SELECT count(*) FROM entity_facts f WHERE f.entity_id=e.id)::int,(SELECT count(*) FROM entity_overrides o WHERE o.entity_id=e.id)::int,(SELECT count(*) FROM entity_relations rel WHERE rel.source_entity_id=e.id OR rel.target_entity_id=e.id)::int,(SELECT count(DISTINCT source_id) FROM entity_facts f WHERE f.entity_id=e.id)::int,(SELECT count(*) FROM entity_facts f WHERE f.entity_id=e.id AND f.valid_until IS NOT NULL AND f.valid_until<now())::int FROM entities e LEFT JOIN entity_overrides od ON od.entity_id=e.id AND od.field_name='displayName' LEFT JOIN entity_overrides os ON os.entity_id=e.id AND os.field_name='status' WHERE e.id=$1`, id)
if err != nil {
return nil, err
}
defer rows.Close()
var result []EntitySummary
for rows.Next() {
var item EntitySummary
if err := rows.Scan(&item.ID, &item.EntityType, &item.CanonicalName, &item.DisplayName, &item.Status, &item.FirstSeenAt, &item.LastSeenAt, &item.TombstonedAt, &item.FactCount, &item.OverrideCount, &item.RelationCount, &item.SourceCount, &item.StaleFactCount); err != nil {
return nil, err
}
result = append(result, item)
}
return result, rows.Err()
}
func effectiveValues(facts []FactView, overrides []OverrideView) []EffectiveValue {
overrideByField := make(map[string]OverrideView, len(overrides))
for _, item := range overrides {
overrideByField[item.FieldName] = item
}
fields := make(map[string]struct{}, len(facts)+len(overrides))
for _, item := range facts {
fields[item.FieldName] = struct{}{}
}
for _, item := range overrides {
fields[item.FieldName] = struct{}{}
}
names := make([]string, 0, len(fields))
for name := range fields {
names = append(names, name)
}
sort.Strings(names)
result := make([]EffectiveValue, 0, len(names))
for _, name := range names {
if item, ok := overrideByField[name]; ok {
at := item.UpdatedAt
result = append(result, EffectiveValue{FieldName: name, Value: item.Value, Origin: "override", OverriddenAt: &at})
continue
}
for _, fact := range facts {
if fact.FieldName == name {
observed, confidence := fact.ObservedAt, fact.Confidence
result = append(result, EffectiveValue{FieldName: name, Value: fact.Value, Origin: "discovered", SourceID: fact.SourceID, SourceName: fact.SourceName, ObservedAt: &observed, Confidence: &confidence, Stale: fact.Stale})
break
}
}
}
return result
}