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 }