Three insertion points: - CreateEntity (POST /api/v1/entities) - EnrollClient (POST /api/v1/clients/enroll) - seed.go (seed ingest at deploy time) Shared logic in internal/checkdefaults — resolves host IP from lan_ip > mesh.netbird.ip > mesh_ip, SSH user/port from attributes. Default checks per entity type: - proxmox-host/standalone-server: ping + cpu + memory + load + disk + updates - workstation: ping + cpu + memory + load - lxc: cpu + memory + load + disk - vm: ping - service: process_check.sh All idempotent (ON CONFLICT DO NOTHING). New machines now get monitoring automatically — no manual curl calls needed.
411 lines
13 KiB
Go
411 lines
13 KiB
Go
package db
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
|
|
"github.com/dtoro/oikos/internal/checkdefaults"
|
|
"github.com/google/uuid"
|
|
"github.com/jackc/pgx/v5"
|
|
)
|
|
|
|
// SeedResult holds counts from a seed ingest operation.
|
|
type SeedResult struct {
|
|
Lifecycles int
|
|
EntityTypes int
|
|
RelationshipTypes int
|
|
Entities int
|
|
Relationships int
|
|
RiskClasses int
|
|
ApprovalRules int
|
|
AutonomySettings int
|
|
}
|
|
|
|
// IngestOntologySeed ingests seeds/ontology.yaml into the DB.
|
|
func IngestOntologySeed(ctx context.Context, tx pgx.Tx, data map[string]any) (*SeedResult, error) {
|
|
r := &SeedResult{}
|
|
|
|
// Lifecycles
|
|
lifecycles, _ := data["lifecycles"].(map[string]any)
|
|
for id, raw := range lifecycles {
|
|
lcMap, _ := raw.(map[string]any)
|
|
states := toStringSlice(lcMap["states"])
|
|
defaultState, _ := lcMap["default_state"].(string)
|
|
terminalStates := toStringSlice(lcMap["terminal_states"])
|
|
if len(terminalStates) == 0 {
|
|
terminalStates = []string{}
|
|
}
|
|
transitionsBytes, _ := json.Marshal(lcMap["transitions"])
|
|
_, err := tx.Exec(ctx,
|
|
`INSERT INTO lifecycle_defs (id, states, default_state, terminal_states, transitions)
|
|
VALUES ($1, $2, $3, $4, $5)
|
|
ON CONFLICT (id) DO UPDATE SET states = $2, default_state = $3,
|
|
terminal_states = $4, transitions = $5`,
|
|
id, states, defaultState, terminalStates, string(transitionsBytes))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("lifecycle %s: %w", id, err)
|
|
}
|
|
r.Lifecycles++
|
|
}
|
|
|
|
// Entity types — need to handle parent_type FK, so insert in dependency order
|
|
// (types with no parent first, then their children)
|
|
types, _ := data["entity_types"].(map[string]any)
|
|
if err := insertEntityTypes(ctx, tx, types, r); err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// Relationship types
|
|
relTypes, _ := data["relationship_types"].(map[string]any)
|
|
for name, raw := range relTypes {
|
|
rtMap, _ := raw.(map[string]any)
|
|
inverse, _ := rtMap["inverse"].(string)
|
|
sourceType, _ := rtMap["source"].(string)
|
|
targetType, _ := rtMap["target"].(string)
|
|
cardinality, _ := rtMap["cardinality"].(string)
|
|
desc, _ := rtMap["description"].(string)
|
|
_, err := tx.Exec(ctx,
|
|
`INSERT INTO relationship_types (name, inverse, source_type, target_type, cardinality, description)
|
|
VALUES ($1, $2, $3, $4, $5, $6)
|
|
ON CONFLICT (name) DO UPDATE SET inverse = $2, source_type = $3,
|
|
target_type = $4, cardinality = $5, description = $6`,
|
|
name, nullableStr(inverse), sourceType, targetType, cardinality, desc)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("relationship_type %s: %w", name, err)
|
|
}
|
|
r.RelationshipTypes++
|
|
}
|
|
|
|
return r, nil
|
|
}
|
|
|
|
// IngestInventorySeed ingests seeds/inventory.yaml into the DB.
|
|
// Every entity and edge is validated against the ontology (abstract types
|
|
// rejected, lifecycle states checked, relationship endpoints hierarchy-
|
|
// validated, cardinality enforced) — a violating seed rolls back atomically.
|
|
func IngestInventorySeed(ctx context.Context, tx pgx.Tx, data map[string]any) (*SeedResult, error) {
|
|
r := &SeedResult{}
|
|
|
|
tree, err := LoadTypeTree(ctx, tx)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("load type tree: %w", err)
|
|
}
|
|
|
|
// Entities
|
|
entities, _ := data["entities"].([]any)
|
|
entityTypes := make(map[string]string) // slug -> type, for edge validation
|
|
for _, raw := range entities {
|
|
eMap, ok := raw.(map[string]any)
|
|
if !ok {
|
|
continue
|
|
}
|
|
slug, _ := eMap["slug"].(string)
|
|
typeName, _ := eMap["type"].(string)
|
|
name, _ := eMap["name"].(string)
|
|
state, _ := eMap["state"].(string)
|
|
attrs := eMap["attributes"]
|
|
|
|
if err := tree.ValidateEntity(typeName, state); err != nil {
|
|
return nil, fmt.Errorf("entity %s: %w", slug, err)
|
|
}
|
|
if state == "" {
|
|
state = tree.DefaultState(typeName)
|
|
}
|
|
entityTypes[slug] = typeName
|
|
|
|
entityID, err := getOrCreateEntityID(ctx, tx, slug)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("entity %s: %w", slug, err)
|
|
}
|
|
|
|
attrsBytes, _ := json.Marshal(attrs)
|
|
_, err = tx.Exec(ctx,
|
|
`INSERT INTO entities (id, slug, type, name, state, attributes, version, created_at, updated_at)
|
|
VALUES ($1, $2, $3, $4, $5, $6, 1, now(), now())
|
|
ON CONFLICT (slug) DO UPDATE SET
|
|
type = EXCLUDED.type, name = EXCLUDED.name,
|
|
attributes = entities.attributes || EXCLUDED.attributes,
|
|
updated_at = now()`,
|
|
entityID, slug, typeName, name, nullableStr(state), string(attrsBytes))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("entity %s: %w", slug, err)
|
|
}
|
|
|
|
// Seed initial entity_status row so health queries return
|
|
// results even before the scheduler populates check results.
|
|
_, err = tx.Exec(ctx,
|
|
`INSERT INTO entity_status (entity_id, health, updated_at)
|
|
VALUES ($1, 'unknown', now())
|
|
ON CONFLICT (entity_id) DO NOTHING`,
|
|
entityID)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("entity_status %s: %w", slug, err)
|
|
}
|
|
|
|
checkdefaults.Ensure(ctx, tx, entityID, slug, typeName, attrsBytes)
|
|
|
|
r.Entities++
|
|
}
|
|
|
|
// Relationships — upsert against the current-edge partial unique index
|
|
// (migration 007) so re-ingest never duplicates edges.
|
|
rels, _ := data["relationships"].([]any)
|
|
for _, raw := range rels {
|
|
relMap, ok := raw.(map[string]any)
|
|
if !ok {
|
|
continue
|
|
}
|
|
source, _ := relMap["source"].(string)
|
|
target, _ := relMap["target"].(string)
|
|
relType, _ := relMap["type"].(string)
|
|
attrs := relMap["attributes"]
|
|
|
|
sourceID, err := getEntityIDBySlug(ctx, tx, source)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("rel from %s: %w", source, err)
|
|
}
|
|
targetID, err := getEntityIDBySlug(ctx, tx, target)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("rel to %s: %w", target, err)
|
|
}
|
|
|
|
srcType := entityTypes[source]
|
|
tgtType := entityTypes[target]
|
|
if srcType == "" || tgtType == "" { // entity pre-existing in DB, not in this seed
|
|
if srcType == "" {
|
|
srcType, err = getEntityTypeBySlug(ctx, tx, source)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
if tgtType == "" {
|
|
tgtType, err = getEntityTypeBySlug(ctx, tx, target)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
}
|
|
if err := tree.ValidateEdge(relType, srcType, tgtType); err != nil {
|
|
return nil, fmt.Errorf("rel %s→%s: %w", source, target, err)
|
|
}
|
|
|
|
attrsBytes, _ := json.Marshal(attrs)
|
|
_, err = tx.Exec(ctx,
|
|
`INSERT INTO relationships (source_id, target_id, type, attributes, valid_from, valid_to)
|
|
VALUES ($1, $2, $3, $4, now(), NULL)
|
|
ON CONFLICT (source_id, target_id, type) WHERE valid_to IS NULL
|
|
DO UPDATE SET attributes = EXCLUDED.attributes`,
|
|
sourceID, targetID, relType, string(attrsBytes))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("rel %s→%s %s: %w", source, target, relType, err)
|
|
}
|
|
r.Relationships++
|
|
}
|
|
|
|
if err := ValidateCardinality(ctx, tx); err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return r, nil
|
|
}
|
|
|
|
// IngestPolicySeed ingests seeds/policy.yaml into the DB.
|
|
func IngestPolicySeed(ctx context.Context, tx pgx.Tx, data map[string]any) (*SeedResult, error) {
|
|
r := &SeedResult{}
|
|
|
|
// Risk classes
|
|
riskClasses, _ := data["risk_classes"].(map[string]any)
|
|
for name, raw := range riskClasses {
|
|
rcMap, _ := raw.(map[string]any)
|
|
desc, _ := rcMap["description"].(string)
|
|
approval, _ := rcMap["approval_required"].(string)
|
|
autonomy, _ := rcMap["autonomy_allowed"].(bool)
|
|
_, err := tx.Exec(ctx,
|
|
`INSERT INTO risk_classes (name, description, approval_required, autonomy_allowed)
|
|
VALUES ($1, $2, $3, $4)
|
|
ON CONFLICT (name) DO UPDATE SET description = $2, approval_required = $3, autonomy_allowed = $4`,
|
|
name, desc, approval, autonomy)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("risk_class %s: %w", name, err)
|
|
}
|
|
r.RiskClasses++
|
|
}
|
|
|
|
// Approval rules
|
|
rules, _ := data["approval_rules"].([]any)
|
|
for _, raw := range rules {
|
|
ruleMap, ok := raw.(map[string]any)
|
|
if !ok {
|
|
continue
|
|
}
|
|
entityType, _ := ruleMap["entity_type"].(string)
|
|
action, _ := ruleMap["action"].(string)
|
|
riskClass, _ := ruleMap["risk_class"].(string)
|
|
autonomy, _ := ruleMap["autonomy_level"].(string)
|
|
scopeEntity, _ := ruleMap["scope_entity"].(string)
|
|
|
|
var scopeID any
|
|
if scopeEntity != "" {
|
|
id, err := getEntityIDBySlug(ctx, tx, scopeEntity)
|
|
if err == nil {
|
|
scopeID = id
|
|
}
|
|
}
|
|
|
|
ruleID := uuid.New()
|
|
_, err := tx.Exec(ctx,
|
|
`INSERT INTO approval_rules (id, entity_type, action, risk_class, autonomy_level, scope_entity, version, updated_at)
|
|
VALUES ($1, $2, $3, $4, $5, $6, 1, now())
|
|
ON CONFLICT (entity_type, action, scope_entity)
|
|
DO UPDATE SET risk_class = $4, autonomy_level = $5, scope_entity = $6, updated_at = now()`,
|
|
ruleID, nullableStr(entityType), action, riskClass, autonomy, scopeID)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("approval_rule %s/%s: %w", entityType, action, err)
|
|
}
|
|
r.ApprovalRules++
|
|
}
|
|
|
|
// Autonomy settings
|
|
settings, _ := data["autonomy_settings"].(map[string]any)
|
|
for key, raw := range settings {
|
|
val, _ := raw.(string)
|
|
_, err := tx.Exec(ctx,
|
|
`INSERT INTO autonomy_settings (key, value, version, updated_at)
|
|
VALUES ($1, $2, 1, now())
|
|
ON CONFLICT (key) DO UPDATE SET value = $2, updated_at = now()`,
|
|
key, val)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("autonomy_setting %s: %w", key, err)
|
|
}
|
|
r.AutonomySettings++
|
|
}
|
|
|
|
return r, nil
|
|
}
|
|
|
|
// --- Helpers ---
|
|
|
|
// getOrCreateEntityID returns the UUID for a slug, generating a new
|
|
// time-ordered UUIDv7 if the slug doesn't exist yet (ADR-0005).
|
|
func getOrCreateEntityID(ctx context.Context, tx pgx.Tx, slug string) (uuid.UUID, error) {
|
|
var id uuid.UUID
|
|
err := tx.QueryRow(ctx, "SELECT id FROM entities WHERE slug = $1", slug).Scan(&id)
|
|
switch {
|
|
case err == nil:
|
|
return id, nil
|
|
case errors.Is(err, pgx.ErrNoRows):
|
|
return uuid.NewV7()
|
|
default:
|
|
return uuid.Nil, fmt.Errorf("lookup slug %s: %w", slug, err)
|
|
}
|
|
}
|
|
|
|
// getEntityTypeBySlug resolves a slug to its entity type name.
|
|
func getEntityTypeBySlug(ctx context.Context, tx pgx.Tx, slug string) (string, error) {
|
|
var t string
|
|
err := tx.QueryRow(ctx, "SELECT type FROM entities WHERE slug = $1", slug).Scan(&t)
|
|
if err != nil {
|
|
return "", fmt.Errorf("resolve type of %s: %w", slug, err)
|
|
}
|
|
return t, nil
|
|
}
|
|
|
|
// getEntityIDBySlug resolves a slug to its UUID.
|
|
func getEntityIDBySlug(ctx context.Context, tx pgx.Tx, slug string) (uuid.UUID, error) {
|
|
var id uuid.UUID
|
|
err := tx.QueryRow(ctx, "SELECT id FROM entities WHERE slug = $1", slug).Scan(&id)
|
|
if err != nil {
|
|
return uuid.Nil, fmt.Errorf("resolve slug %s: %w", slug, err)
|
|
}
|
|
return id, nil
|
|
}
|
|
|
|
// insertEntityTypes inserts entity types in dependency order (parents before children).
|
|
func insertEntityTypes(ctx context.Context, tx pgx.Tx, types map[string]any, r *SeedResult) error {
|
|
// Build a dependency graph and insert in topological order
|
|
// Simple approach: insert types with no parent first, then iterate
|
|
inserted := make(map[string]bool)
|
|
remaining := make(map[string]map[string]any)
|
|
for name, raw := range types {
|
|
tMap, _ := raw.(map[string]any)
|
|
remaining[name] = tMap
|
|
}
|
|
|
|
maxPasses := 10
|
|
for pass := 0; pass < maxPasses && len(remaining) > 0; pass++ {
|
|
for name, tMap := range remaining {
|
|
parent, _ := tMap["parent"].(string)
|
|
if parent == "" || inserted[parent] {
|
|
if err := insertOneEntityType(ctx, tx, name, tMap); err != nil {
|
|
return err
|
|
}
|
|
inserted[name] = true
|
|
delete(remaining, name)
|
|
r.EntityTypes++
|
|
}
|
|
}
|
|
}
|
|
|
|
if len(remaining) > 0 {
|
|
return fmt.Errorf("circular or missing parent in entity types: %v", keysOf(remaining))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func insertOneEntityType(ctx context.Context, tx pgx.Tx, name string, tMap map[string]any) error {
|
|
parent, _ := tMap["parent"].(string)
|
|
isAbstract, _ := tMap["abstract"].(bool)
|
|
domain, _ := tMap["domain"].(string)
|
|
layer, _ := tMap["layer"].(string)
|
|
desc, _ := tMap["description"].(string)
|
|
lifecycleID, _ := tMap["lifecycle"].(string)
|
|
attrSchema := tMap["attribute_schema"]
|
|
|
|
schemaBytes, _ := json.Marshal(attrSchema)
|
|
_, err := tx.Exec(ctx,
|
|
`INSERT INTO entity_types (name, parent_type, is_abstract, domain, layer, description,
|
|
lifecycle_id, attribute_schema, schema_version, status, created_at, updated_at)
|
|
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, 1, 'active', now(), now())
|
|
ON CONFLICT (name) DO UPDATE SET parent_type = $2, is_abstract = $3, domain = $4,
|
|
layer = $5, description = $6, lifecycle_id = $7, attribute_schema = $8, updated_at = now()`,
|
|
name, nullableStr(parent), isAbstract, domain, layer, desc, nullableStr(lifecycleID), nullableStr(string(schemaBytes)))
|
|
return err
|
|
}
|
|
|
|
func toStringSlice(v any) []string {
|
|
if v == nil {
|
|
return nil
|
|
}
|
|
switch s := v.(type) {
|
|
case []string:
|
|
return s
|
|
case []any:
|
|
out := make([]string, 0, len(s))
|
|
for _, item := range s {
|
|
if str, ok := item.(string); ok {
|
|
out = append(out, str)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func nullableStr(s string) any {
|
|
if s == "" {
|
|
return nil
|
|
}
|
|
return s
|
|
}
|
|
|
|
func keysOf(m map[string]map[string]any) []string {
|
|
keys := make([]string, 0, len(m))
|
|
for k := range m {
|
|
keys = append(keys, k)
|
|
}
|
|
return keys
|
|
}
|
|
|