Public source validation / validate (push) Failing after 3m8s
283 lines
13 KiB
Go
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
|
|
}
|