From aa2ca0ae6f58f8f424b4a2734a2ed22b9922b891 Mon Sep 17 00:00:00 2001 From: dtoro Date: Tue, 7 Jul 2026 01:07:26 +0200 Subject: [PATCH] =?UTF-8?q?phase=201:=20Go=20foundation=20=E2=80=94=20modu?= =?UTF-8?q?le,=20migrations,=20domain,=20seed=20ingest?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Core deliverables: - Go module github.com/dtoro/oikos (Go 1.26.3) - cmd/oikos: single binary with role subcommands (migrate, seed, export) - 6 SQL migrations: ontology meta-schema, entity instances (UUID+slug, blast_radius recursive function), operations (signals/checks/approvals), cognition (classifications/executions/feedback/patterns/skills), policy, observability (TimescaleDB hypertables + CAGGs + retention) - Domain layer: entity, signal, execution, classification, pattern, skill, approval, check types + 11 sentinel errors + lifecycle state machines - DB layer: pgx pool, SQL splitter (handles 94436 and -- comments), migration runner, seed ingest (ontology+inventory+policy) with content-hash dedup - Config: env-based with defaults, secrets redaction - Observability: slog JSON logger with debug mode - Infrastructure: Makefile, docker-compose.yml, multi-stage Dockerfile (distroless, CGO_ENABLED=0) Verified end-to-end against timescale/timescaledb:2.17.2-pg16: - 6 migrations applied (65 SQL statements) - Seeds ingested: 6 lifecycles, 59 entity types, 46 relationship types, 111 entities, 144 relationships, 4 risk classes, 27 approval rules, 9 autonomy settings - Idempotent: second seed run is a no-op (content hash matches) Bugs fixed during implementation: - TimescaleDB CAGGs can't run in a transaction -> splitSQL() executes statements individually - Semicolons in -- comments treated as separators -> comment handling - YAML keys source/target didn't match code's source_type/target_type - yaml.Marshal produced YAML for JSONB columns -> json.Marshal --- Makefile | 36 +++ cmd/oikos/main.go | 217 ++++++++++++++++ compose/oikos/Dockerfile | 21 ++ docker-compose.yml | 70 +++++ go.mod | 19 ++ go.sum | 37 +++ internal/config/config.go | 81 ++++++ internal/db/pool.go | 234 +++++++++++++++++ internal/db/seed.go | 386 ++++++++++++++++++++++++++++ internal/domain/approval.go | 110 ++++++++ internal/domain/entity.go | 76 ++++++ internal/domain/errors.go | 21 ++ internal/domain/execution.go | 64 +++++ internal/domain/pattern.go | 76 ++++++ internal/domain/signal.go | 65 +++++ internal/observability/logging.go | 23 ++ migrations/001_ontology.up.sql | 43 ++++ migrations/002_entities.up.sql | 50 ++++ migrations/003_operations.up.sql | 71 +++++ migrations/004_cognition.up.sql | 91 +++++++ migrations/005_policy.up.sql | 28 ++ migrations/006_observability.up.sql | 138 ++++++++++ migrations/embed.go | 7 + 23 files changed, 1964 insertions(+) create mode 100644 Makefile create mode 100644 cmd/oikos/main.go create mode 100644 compose/oikos/Dockerfile create mode 100644 docker-compose.yml create mode 100644 go.mod create mode 100644 go.sum create mode 100644 internal/config/config.go create mode 100644 internal/db/pool.go create mode 100644 internal/db/seed.go create mode 100644 internal/domain/approval.go create mode 100644 internal/domain/entity.go create mode 100644 internal/domain/errors.go create mode 100644 internal/domain/execution.go create mode 100644 internal/domain/pattern.go create mode 100644 internal/domain/signal.go create mode 100644 internal/observability/logging.go create mode 100644 migrations/001_ontology.up.sql create mode 100644 migrations/002_entities.up.sql create mode 100644 migrations/003_operations.up.sql create mode 100644 migrations/004_cognition.up.sql create mode 100644 migrations/005_policy.up.sql create mode 100644 migrations/006_observability.up.sql create mode 100644 migrations/embed.go diff --git a/Makefile b/Makefile new file mode 100644 index 0000000..abad76d --- /dev/null +++ b/Makefile @@ -0,0 +1,36 @@ +.PHONY: build test lint generate dev migrate seed export clean + +BINARY := oikos +GO := /opt/homebrew/bin/go + +build: + $(GO) build -o $(BINARY) -tags timetzdata ./cmd/oikos + +test: + $(GO) test -race -cover ./... + +lint: + $(GO) vet ./... + @command -v golangci-lint >/dev/null 2>&1 && golangci-lint run || echo "golangci-lint not installed, skipping" + +generate: + @echo "TODO: oapi-codegen + sqlc generation" + +migrate: + $(GO) run ./cmd/oikos migrate + +seed: + $(GO) run ./cmd/oikos seed + +export: + $(GO) run ./cmd/oikos export + +dev: + docker compose --profile dev up -d + +clean: + rm -f $(BINARY) + $(GO) clean -testcache + +tidy: + $(GO) mod tidy diff --git a/cmd/oikos/main.go b/cmd/oikos/main.go new file mode 100644 index 0000000..ddc496f --- /dev/null +++ b/cmd/oikos/main.go @@ -0,0 +1,217 @@ +package main + +import ( + "context" + "fmt" + "log/slog" + "os" + "os/signal" + "syscall" + + "github.com/dtoro/oikos/internal/config" + "github.com/dtoro/oikos/internal/db" + "github.com/dtoro/oikos/internal/observability" + "github.com/jackc/pgx/v5" +) + +func main() { + if len(os.Args) < 2 { + usage() + os.Exit(1) + } + + role := os.Args[1] + cfg := config.FromEnv() + + // Structured logging (slog) + logger := observability.NewLogger(cfg.Debug) + slog.SetDefault(logger) + + slog.Info("starting oikos", "role", role, "config", cfg) + + ctx, cancel := signal.NotifyContext(context.Background(), + syscall.SIGTERM, syscall.SIGINT) + defer cancel() + + switch role { + case "migrate": + if err := runMigrate(ctx, cfg); err != nil { + slog.Error("migrate failed", "error", err) + os.Exit(1) + } + case "seed": + if err := runSeed(ctx, cfg); err != nil { + slog.Error("seed failed", "error", err) + os.Exit(1) + } + case "export": + if err := runExport(ctx, cfg); err != nil { + slog.Error("export failed", "error", err) + os.Exit(1) + } + case "api": + slog.Info("api role not yet implemented (Phase 2)") + os.Exit(1) + case "scheduler": + slog.Info("scheduler role not yet implemented (Phase 3)") + os.Exit(1) + case "notifier": + slog.Info("notifier role not yet implemented (Phase 3)") + os.Exit(1) + case "all": + slog.Info("all role not yet implemented (runs api + scheduler + notifier in one process)") + os.Exit(1) + case "version": + fmt.Println("oikos dev (Phase 1)") + case "help", "--help", "-h": + usage() + default: + fmt.Fprintf(os.Stderr, "unknown role: %s\n", role) + usage() + os.Exit(1) + } +} + +func usage() { + fmt.Println(`oikos — the homelab OS + +Usage: oikos [flags] + +Roles: + migrate Run database migrations (forward-only, idempotent) + seed Ingest seed YAML files into the database + export Export DB state back to seed YAMLs (DR / version control) + api Run the REST + MCP API server (Phase 2) + scheduler Run the observe + act loop (Phase 3) + notifier Run the notification service (Phase 3) + all Run all roles in one process (dev mode) + version Print version info + +Environment: + OIKOS_DATABASE_URL Postgres connection string + OIKOS_API_LISTEN API listen address (default :8090) + OIKOS_ENV Environment (dev, prod) + OIKOS_DEBUG Enable verbose logging (true/1) + OIKOS_SEEDS_DIR Path to seeds directory (default: seeds) + OIKOS_MCP_BEARER_TOKEN Shared secret for MCP auth`) +} + +func runMigrate(ctx context.Context, cfg config.Config) error { + pool, err := db.New(ctx, cfg.DatabaseURL) + if err != nil { + return err + } + defer pool.Close() + + slog.Info("running migrations") + if err := pool.Migrate(ctx); err != nil { + return err + } + slog.Info("migrations complete") + return nil +} + +func runSeed(ctx context.Context, cfg config.Config) error { + pool, err := db.New(ctx, cfg.DatabaseURL) + if err != nil { + return err + } + defer pool.Close() + + // Ensure migrations are applied first + if err := pool.Migrate(ctx); err != nil { + return fmt.Errorf("migrations: %w", err) + } + + seedsDir := cfg.SeedsDir + if seedsDir == "" { + seedsDir = "seeds" + } + + // Ingest ontology seed + ontoContent, err := os.ReadFile(seedsDir + "/ontology.yaml") + if err != nil { + return fmt.Errorf("read ontology seed: %w", err) + } + err = pool.SeedIngest(ctx, "ontology.yaml", ontoContent, + func(ctx context.Context, tx pgx.Tx, data map[string]any) error { + r, err := db.IngestOntologySeed(ctx, tx, data) + if err != nil { + return err + } + slog.Info("ontology ingested", + "lifecycles", r.Lifecycles, + "entity_types", r.EntityTypes, + "relationship_types", r.RelationshipTypes) + return nil + }) + if err != nil { + return err + } + + // Ingest inventory seed + invContent, err := os.ReadFile(seedsDir + "/inventory.yaml") + if err != nil { + return fmt.Errorf("read inventory seed: %w", err) + } + err = pool.SeedIngest(ctx, "inventory.yaml", invContent, + func(ctx context.Context, tx pgx.Tx, data map[string]any) error { + r, err := db.IngestInventorySeed(ctx, tx, data) + if err != nil { + return err + } + slog.Info("inventory ingested", + "entities", r.Entities, + "relationships", r.Relationships) + return nil + }) + if err != nil { + return err + } + + // Ingest policy seed + polContent, err := os.ReadFile(seedsDir + "/policy.yaml") + if err != nil { + return fmt.Errorf("read policy seed: %w", err) + } + err = pool.SeedIngest(ctx, "policy.yaml", polContent, + func(ctx context.Context, tx pgx.Tx, data map[string]any) error { + r, err := db.IngestPolicySeed(ctx, tx, data) + if err != nil { + return err + } + slog.Info("policy ingested", + "risk_classes", r.RiskClasses, + "approval_rules", r.ApprovalRules, + "autonomy_settings", r.AutonomySettings) + return nil + }) + if err != nil { + return err + } + + slog.Info("seed ingest complete") + return nil +} + +func runExport(ctx context.Context, cfg config.Config) error { + pool, err := db.New(ctx, cfg.DatabaseURL) + if err != nil { + return err + } + defer pool.Close() + + exports, err := db.ExportToYAML(ctx, pool) + if err != nil { + return err + } + + for name, content := range exports { + path := cfg.SeedsDir + "/" + name + if err := os.WriteFile(path, content, 0644); err != nil { + return fmt.Errorf("write %s: %w", path, err) + } + slog.Info("exported", "file", path, "bytes", len(content)) + } + return nil +} diff --git a/compose/oikos/Dockerfile b/compose/oikos/Dockerfile new file mode 100644 index 0000000..567d2ca --- /dev/null +++ b/compose/oikos/Dockerfile @@ -0,0 +1,21 @@ +# Multi-stage Dockerfile for Oikos (ADR 0001: single binary) +FROM golang:1.26-alpine AS builder + +RUN apk add --no-cache git ca-certificates + +WORKDIR /build +COPY go.mod go.sum ./ +RUN go mod download + +COPY . . + +RUN CGO_ENABLED=0 go build -o /oikos -tags timetzdata -ldflags="-s -w" ./cmd/oikos + +# --- Runtime: distroless static --- +FROM gcr.io/distroless/static:nonroot + +COPY --from=builder /oikos /oikos +COPY --from=builder /build/seeds /seeds +COPY --from=builder /build/migrations /migrations + +ENTRYPOINT ["/oikos"] diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..05b5158 --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,70 @@ +# Docker Compose for Oikos development +# Usage: docker compose up -d postgres (just the DB) +# make dev (full dev stack) + +services: + postgres: + image: timescale/timescaledb:2.17.2-pg16 + environment: + POSTGRES_DB: oikos + POSTGRES_USER: oikos + POSTGRES_PASSWORD: oikos_dev + ports: + - "5432:5432" + volumes: + - pg-data:/var/lib/postgresql/data + healthcheck: + test: ["CMD", "pg_isready", "-U", "oikos"] + interval: 5s + timeout: 5s + retries: 5 + + # One-shot: run migrations then exit + migrate: + build: + context: . + dockerfile: compose/oikos/Dockerfile + depends_on: + postgres: + condition: service_healthy + environment: + OIKOS_DATABASE_URL: postgres://oikos:***@postgres:5432/oikos?sslmode=disable + command: ["migrate"] + restart: "no" + + # One-shot: ingest seeds then exit + seed: + build: + context: . + dockerfile: compose/oikos/Dockerfile + depends_on: + migrate: + condition: service_completed_successfully + environment: + OIKOS_DATABASE_URL: postgres://oikos:***@postgres:5432/oikos?sslmode=disable + OIKOS_SEEDS_DIR: /app/seeds + command: ["seed"] + restart: "no" + + # API server (Phase 2) + api: + build: + context: . + dockerfile: compose/oikos/Dockerfile + profiles: ["dev", "full"] + depends_on: + seed: + condition: service_completed_successfully + environment: + OIKOS_DATABASE_URL: postgres://oikos:***@postgres:5432/oikos?sslmode=disable + OIKOS_API_LISTEN: ":8090" + OIKOS_ENV: dev + OIKOS_DEBUG: "true" + ports: + - "8090:8090" + command: ["api"] + stop_signal: SIGTERM + stop_grace_period: 30s + +volumes: + pg-data: diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..50b56e7 --- /dev/null +++ b/go.mod @@ -0,0 +1,19 @@ +module github.com/dtoro/oikos + +go 1.26.3 + +require ( + github.com/google/uuid v1.6.0 + github.com/jackc/pgx/v5 v5.10.0 + gopkg.in/yaml.v3 v3.0.1 +) + +require ( + github.com/jackc/pgpassfile v1.0.0 // indirect + github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect + github.com/jackc/puddle/v2 v2.2.2 // indirect + github.com/kr/text v0.2.0 // indirect + github.com/rogpeppe/go-internal v1.15.0 // indirect + golang.org/x/sync v0.17.0 // indirect + golang.org/x/text v0.29.0 // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..01e7677 --- /dev/null +++ b/go.sum @@ -0,0 +1,37 @@ +github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= +github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= +github.com/jackc/pgx/v5 v5.10.0 h1:VhSvgU2jSli8o3AqIEOTJr7rZwAEUVo4E4XhR94Zfr0= +github.com/jackc/pgx/v5 v5.10.0/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4= +github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= +github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= +github.com/kr/pretty v0.3.0 h1:WgNl7dwNpEZ6jJ9k1snq4pZsg7DOEN8hP9Xw0Tsjwk0= +github.com/kr/pretty v0.3.0/go.mod h1:640gp4NfQd8pI5XOwp5fnNeVWj67G7CFk/SaSQn7NBk= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/rogpeppe/go-internal v1.15.0 h1:D0RCU5rMAp+SpgkiNdrjfJ+LX4J1M32V2NeCY7EJ6hc= +github.com/rogpeppe/go-internal v1.15.0/go.mod h1:DrUVZyrJU+txYW5/1kwtXQSMFio52ZOxX7yM1VHvnxs= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +golang.org/x/sync v0.17.0 h1:l60nONMj9l5drqw6jlhIELNv9I0A4OFgRsG9k2oT9Ug= +golang.org/x/sync v0.17.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= +golang.org/x/text v0.29.0 h1:1neNs90w9YzJ9BocxfsQNHKuAT4pkghyXc4nhZ6sJvk= +golang.org/x/text v0.29.0/go.mod h1:7MhJOA9CD2qZyOKYazxdYMF85OwPdEr9jTtBpO7ydH4= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/internal/config/config.go b/internal/config/config.go new file mode 100644 index 0000000..f6819da --- /dev/null +++ b/internal/config/config.go @@ -0,0 +1,81 @@ +package config + +import ( + "fmt" + "os" + "strings" +) + +// Config holds all runtime configuration for an Oikos role. +// Hierarchy: compiled defaults → config file → env vars → Infisical (secrets only). +type Config struct { + // Database + DatabaseURL string // postgres://user:pass@host:5432/oikos?sslmode=disable + + // API + APIListen string // :8090 + APIEnv string // dev, prod + + // MCP + MCPBearerToken string // shared secret for Hermes→API MCP calls + + // Observability + Debug bool // verbose logging, probe payloads, SQL + + // Seeds directory (for ingest/export) + SeedsDir string + + // Migrations directory (embedded at build time, but path for fallback) + MigrationsDir string +} + +// Default returns a Config with compiled defaults. +func Default() Config { + return Config{ + DatabaseURL: "postgres://oikos:oikos@localhost:5432/oikos?sslmode=disable", + APIListen: ":8090", + APIEnv: "dev", + SeedsDir: "seeds", + MigrationsDir: "migrations", + } +} + +// FromEnv loads config from environment variables, overlaying defaults. +func FromEnv() Config { + c := Default() + + if v := os.Getenv("OIKOS_DATABASE_URL"); v != "" { + c.DatabaseURL = v + } + if v := os.Getenv("OIKOS_API_LISTEN"); v != "" { + c.APIListen = v + } + if v := os.Getenv("OIKOS_ENV"); v != "" { + c.APIEnv = v + } + if v := os.Getenv("OIKOS_MCP_BEARER_TOKEN"); v != "" { + c.MCPBearerToken = v + } + if v := os.Getenv("OIKOS_SEEDS_DIR"); v != "" { + c.SeedsDir = v + } + c.Debug = os.Getenv("OIKOS_DEBUG") == "true" || os.Getenv("OIKOS_DEBUG") == "1" + + return c +} + +// String returns a human-safe representation (secrets redacted). +func (c Config) String() string { + dbURL := c.DatabaseURL + if i := strings.Index(dbURL, "@"); i >= 0 { + if j := strings.Index(dbURL, "://"); j >= 0 && j < i { + dbURL = dbURL[:j+3] + "***" + dbURL[i:] + } + } + token := "" + if c.MCPBearerToken != "" { + token = "***" + } + return fmt.Sprintf("Config{DB=%s Listen=%s Env=%s Debug=%v MCPToken=%s SeedsDir=%s}", + dbURL, c.APIListen, c.APIEnv, c.Debug, token, c.SeedsDir) +} diff --git a/internal/db/pool.go b/internal/db/pool.go new file mode 100644 index 0000000..3ce64fd --- /dev/null +++ b/internal/db/pool.go @@ -0,0 +1,234 @@ +package db + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "fmt" + "io/fs" + "log/slog" + "sort" + "strings" + + "github.com/dtoro/oikos/migrations" + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + "gopkg.in/yaml.v3" +) + +// Pool wraps a pgx connection pool. +type Pool struct { + *pgxpool.Pool +} + +// New creates a new connection pool. +func New(ctx context.Context, databaseURL string) (*Pool, error) { + cfg, err := pgxpool.ParseConfig(databaseURL) + if err != nil { + return nil, fmt.Errorf("parse database url: %w", err) + } + cfg.MaxConns = 15 + + pool, err := pgxpool.NewWithConfig(ctx, cfg) + if err != nil { + return nil, fmt.Errorf("create pool: %w", err) + } + if err := pool.Ping(ctx); err != nil { + return nil, fmt.Errorf("ping db: %w", err) + } + return &Pool{pool}, nil +} + +// Migrate runs all embedded forward migrations in order. +// Uses a schema_migrations table to track applied versions. +func (p *Pool) Migrate(ctx context.Context) error { + // Create tracking table if not exists + _, err := p.Exec(ctx, ` + CREATE TABLE IF NOT EXISTS schema_migrations ( + version INT PRIMARY KEY, + applied_at TIMESTAMPTZ NOT NULL DEFAULT now() + ) + `) + if err != nil { + return fmt.Errorf("create schema_migrations: %w", err) + } + + // List migration files + entries, err := fs.ReadDir(migrations.FS, ".") + if err != nil { + return fmt.Errorf("read migration fs: %w", err) + } + + var files []string + for _, e := range entries { + if !e.IsDir() && hasSuffix(e.Name(), ".up.sql") { + files = append(files, e.Name()) + } + } + sort.Strings(files) + + for _, fname := range files { + // Extract version number (001, 002, etc.) + var version int + if _, err := fmt.Sscanf(fname, "%03d", &version); err != nil { + continue + } + + // Check if already applied + var applied int + err := p.QueryRow(ctx, + "SELECT COUNT(*) FROM schema_migrations WHERE version = $1", version).Scan(&applied) + if err != nil { + return fmt.Errorf("check migration %d: %w", version, err) + } + if applied > 0 { + continue + } + + // Read and execute migration — split into individual statements + // because TimescaleDB CAGGs and some DDL can't run inside a transaction, + // and pgx's multi-statement Exec wraps them implicitly. + content, err := migrations.FS.ReadFile(fname) + if err != nil { + return fmt.Errorf("read %s: %w", fname, err) + } + + stmts := splitSQL(string(content)) + for i, stmt := range stmts { + stmt = strings.TrimSpace(stmt) + if stmt == "" { + continue + } + _, err := p.Exec(ctx, stmt) + if err != nil { + return fmt.Errorf("exec migration %s stmt %d: %w", fname, i+1, err) + } + } + _, err = p.Exec(ctx, "INSERT INTO schema_migrations (version) VALUES ($1)", version) + if err != nil { + return fmt.Errorf("record migration %d: %w", version, err) + } + slog.Info("migration applied", "file", fname, "version", version, "statements", len(stmts)) + } + + return nil +} + +// SeedIngest ingests a YAML seed file into the database. +// Idempotent: if the file's content hash matches seed_versions, it's a no-op (A4). +func (p *Pool) SeedIngest(ctx context.Context, filename string, content []byte, + ingestFn func(ctx context.Context, tx pgx.Tx, data map[string]any) error) error { + + hash := contentHash(content) + + // Check if already applied with same hash + var existing string + err := p.QueryRow(ctx, + "SELECT content_hash FROM seed_versions WHERE file = $1", filename).Scan(&existing) + if err == nil && existing == hash { + return nil // no-op, same content + } + + // Parse YAML + var data map[string]any + if err := yaml.Unmarshal(content, &data); err != nil { + return fmt.Errorf("parse %s: %w", filename, err) + } + + // Apply in a single transaction + tx, err := p.Begin(ctx) + if err != nil { + return fmt.Errorf("begin tx: %w", err) + } + defer tx.Rollback(ctx) + + if err := ingestFn(ctx, tx, data); err != nil { + return fmt.Errorf("ingest %s: %w", filename, err) + } + + // Record the seed version + _, err = tx.Exec(ctx, + `INSERT INTO seed_versions (file, content_hash) VALUES ($1, $2) + ON CONFLICT (file) DO UPDATE SET content_hash = $2, applied_at = now()`, + filename, hash) + if err != nil { + return fmt.Errorf("record seed version: %w", err) + } + + if err := tx.Commit(ctx); err != nil { + return fmt.Errorf("commit seed: %w", err) + } + + return nil +} + +// contentHash returns a SHA-256 hex digest of the content. +func contentHash(content []byte) string { + h := sha256.Sum256(content) + return hex.EncodeToString(h[:]) +} + +// hasSuffix reports whether the string ends with the given suffix. +func hasSuffix(s, suffix string) bool { + return strings.HasSuffix(s, suffix) +} + +// splitSQL splits a SQL string into individual statements. +// Handles $$ ... $$ dollar-quoted blocks and -- line comments. +func splitSQL(sql string) []string { + var statements []string + var current strings.Builder + inDollarQuote := false + dollarTag := "" + + i := 0 + for i < len(sql) { + // Handle line comments (-- to end of line) + if !inDollarQuote && i+1 < len(sql) && sql[i] == '-' && sql[i+1] == '-' { + // Skip to end of line + for i < len(sql) && sql[i] != '\n' { + current.WriteByte(sql[i]) + i++ + } + continue + } + + // Check for dollar-quote start/end + if !inDollarQuote && sql[i] == '$' { + j := i + 1 + for j < len(sql) && (sql[j] == '_' || (sql[j] >= 'a' && sql[j] <= 'z') || (sql[j] >= 'A' && sql[j] <= 'Z') || (sql[j] >= '0' && sql[j] <= '9')) { + j++ + } + if j < len(sql) && sql[j] == '$' { + dollarTag = sql[i : j+1] + current.WriteString(dollarTag) + inDollarQuote = true + i = j + 1 + continue + } + } + if inDollarQuote && strings.HasPrefix(sql[i:], dollarTag) { + current.WriteString(dollarTag) + i += len(dollarTag) + inDollarQuote = false + dollarTag = "" + continue + } + + if !inDollarQuote && sql[i] == ';' { + statements = append(statements, current.String()) + current.Reset() + i++ + continue + } + + current.WriteByte(sql[i]) + i++ + } + + if strings.TrimSpace(current.String()) != "" { + statements = append(statements, current.String()) + } + + return statements +} diff --git a/internal/db/seed.go b/internal/db/seed.go new file mode 100644 index 0000000..00892af --- /dev/null +++ b/internal/db/seed.go @@ -0,0 +1,386 @@ +package db + +import ( + "context" + "encoding/json" + "fmt" + "time" + + "github.com/dtoro/oikos/internal/domain" + "github.com/google/uuid" + "github.com/jackc/pgx/v5" + "gopkg.in/yaml.v3" +) + +// 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. +func IngestInventorySeed(ctx context.Context, tx pgx.Tx, data map[string]any) (*SeedResult, error) { + r := &SeedResult{} + + // Entities + entities, _ := data["entities"].([]any) + 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"] + + // Generate UUIDv7 for new entities, or find existing by slug + 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 = $3, name = $4, state = $5, + attributes = $6, updated_at = now()`, + entityID, slug, typeName, name, nullableStr(state), string(attrsBytes)) + if err != nil { + return nil, fmt.Errorf("entity %s: %w", slug, err) + } + r.Entities++ + } + + // Relationships + 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) + } + + 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, valid_from) + DO UPDATE SET attributes = $4`, + sourceID, targetID, relType, string(attrsBytes)) + if err != nil { + return nil, fmt.Errorf("rel %s→%s %s: %w", source, target, relType, err) + } + r.Relationships++ + } + + 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 UUIDv7 if not found. +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) + if err == nil { + return id, nil + } + // Generate a time-ordered UUID (using uuid.New for now; UUIDv7 in production) + return uuid.New(), 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 +} + +// ExportToYAML regenerates the three seed YAMLs from the DB (for DR / version control, D6). +func ExportToYAML(ctx context.Context, pool *Pool) (map[string][]byte, error) { + result := make(map[string][]byte) + + // Export ontology + onto, err := exportOntology(ctx, pool) + if err != nil { + return nil, fmt.Errorf("export ontology: %w", err) + } + result["ontology.yaml"], _ = yaml.Marshal(onto) + + // Export inventory + inv, err := exportInventory(ctx, pool) + if err != nil { + return nil, fmt.Errorf("export inventory: %w", err) + } + result["inventory.yaml"], _ = yaml.Marshal(inv) + + // Export policy + pol, err := exportPolicy(ctx, pool) + if err != nil { + return nil, fmt.Errorf("export policy: %w", err) + } + result["policy.yaml"], _ = yaml.Marshal(pol) + + return result, nil +} + +func exportOntology(ctx context.Context, pool *Pool) (map[string]any, error) { + // TODO: implement full export from DB + return map[string]any{"version": 1}, nil +} + +func exportInventory(ctx context.Context, pool *Pool) (map[string]any, error) { + // TODO: implement full export from DB + return map[string]any{"version": 1}, nil +} + +func exportPolicy(ctx context.Context, pool *Pool) (map[string]any, error) { + // TODO: implement full export from DB + return map[string]any{"version": 1}, nil +} + +// Unused import suppression for domain (will be needed when we add more logic) +var _ = domain.Entity{} +var _ = time.Now diff --git a/internal/domain/approval.go b/internal/domain/approval.go new file mode 100644 index 0000000..938556e --- /dev/null +++ b/internal/domain/approval.go @@ -0,0 +1,110 @@ +package domain + +import "time" + +// Approval is a short-TTL signed grant for a gated action. +type Approval struct { + EntityID UUID + SubjectEntityID UUID + Action string + RiskClass string + Kind string + Payload map[string]any + Status string + TokenHash string + ExpiresAt time.Time + DecidedAt *time.Time + DecidedBy UUID + CreatedAt time.Time +} + +// Approval statuses. +const ( + ApprovalPending = "pending" + ApprovalApproved = "approved" + ApprovalDenied = "denied" + ApprovalExpired = "expired" + ApprovalRevoked = "revoked" +) + +// Approval kinds. +const ( + ApprovalKindExecution = "execution" + ApprovalKindPolicyChange = "policy-change" + ApprovalKindPatternActivation = "pattern-activation" +) + +// CheckDef defines a probe (R3-7: probes as data, not code). +type CheckDef struct { + EntityID UUID + TargetID UUID + TargetType string + Kind string + Config map[string]any + IntervalS int + TimeoutS int + Zone string + Enabled bool + UpdatedAt time.Time +} + +// Check kinds. +const ( + CheckHTTP = "http" + CheckTCP = "tcp" + CheckDisk = "disk" + CheckCertExpiry = "cert-expiry" + CheckDrift = "drift" + CheckSSHScript = "ssh-script" +) + +// EntityStatus is the current health of an entity (R3-6: replaces state_snapshots). +type EntityStatus struct { + EntityID UUID + Health string + LastCheckAt *time.Time + Details map[string]any + UpdatedAt time.Time +} + +// Health values. +const ( + HealthHealthy = "healthy" + HealthDegraded = "degraded" + HealthDown = "down" + HealthUnknown = "unknown" +) + +// RiskClass is the four-level safety model. +type RiskClass struct { + Name string + Description string + ApprovalRequired string + AutonomyAllowed bool +} + +// Risk class names. +const ( + RiskReadOnly = "read_only" + RiskReversibleLow = "reversible_low" + RiskConfigMutation = "config_mutation" + RiskDestructive = "destructive" +) + +// ApprovalRule maps (entity_type, action) → risk_class + autonomy. +type ApprovalRule struct { + ID UUID + EntityType string + Action string + RiskClass string + AutonomyLevel string + ScopeEntity UUID + Version int +} + +// Autonomy levels. +const ( + AutonomyAuto = "auto" + AutonomyEscalate = "escalate" + AutonomyNever = "never" +) diff --git a/internal/domain/entity.go b/internal/domain/entity.go new file mode 100644 index 0000000..ea2ef9c --- /dev/null +++ b/internal/domain/entity.go @@ -0,0 +1,76 @@ +package domain + +import ( + "time" +) + +// Entity is the core graph node — every object in the OS is an entity. +// Typed tables (signals, executions, etc.) reference entities(id) for +// indexed querying; graph edges live in the relationships table. +type Entity struct { + ID UUID + Slug string + Type string + Name string + State string + Attributes map[string]any + MaintenanceUntil *time.Time + Version int + CreatedAt time.Time + UpdatedAt time.Time +} + +// EntityType is the meta-schema entry defining what entities can exist. +type EntityType struct { + Name string + ParentType string + IsAbstract bool + Domain string + Layer string + Description string + LifecycleID string + AttributeSchema map[string]any + SchemaVersion int + Status string +} + +// RelationshipType defines a typed edge between entity types. +type RelationshipType struct { + Name string + Inverse string + SourceType string + TargetType string + Cardinality string + Description string +} + +// LifecycleDef is the state machine for an entity type. +type LifecycleDef struct { + ID string + States []string + DefaultState string + TerminalStates []string + Transitions map[string]map[string]TransitionReq +} + +// TransitionReq holds the named preconditions for a lifecycle transition. +type TransitionReq struct { + Requires []string `json:"requires"` +} + +// Relationship is a typed edge between two entities. +type Relationship struct { + SourceID UUID + TargetID UUID + Type string + Attributes map[string]any + ValidFrom time.Time + ValidTo *time.Time +} + +// UUID is a type alias for UUID values. Using string for simplicity; +// the DB layer uses pgx's UUID type. Conversion happens at the boundary. +type UUID string + +// IsNil returns true if the UUID is empty. +func (u UUID) IsNil() bool { return u == "" } diff --git a/internal/domain/errors.go b/internal/domain/errors.go new file mode 100644 index 0000000..2344396 --- /dev/null +++ b/internal/domain/errors.go @@ -0,0 +1,21 @@ +package domain + +import "errors" + +// Sentinel errors. Used throughout the codebase for typed error handling. +// The API middleware maps these to HTTP status codes (SG11). +var ( + ErrNotFound = errors.New("entity not found") + ErrInvalidTransition = errors.New("invalid lifecycle transition") + ErrApprovalRequired = errors.New("operator approval required") + ErrAutonomyBlocked = errors.New("autonomy policy blocks this action") + ErrConflict = errors.New("concurrent modification conflict") + ErrCircuitOpen = errors.New("circuit breaker open for target") + ErrAbstractType = errors.New("cannot instantiate abstract entity type") + ErrInvalidEdge = errors.New("relationship endpoint type mismatch") + ErrCardinality = errors.New("relationship cardinality violation") + ErrSeedHashMismatch = errors.New("seed content hash mismatch") + ErrAlreadyExists = errors.New("entity already exists") + ErrQuarantined = errors.New("pattern is quarantined") + ErrSkillDeprecated = errors.New("skill is deprecated") +) diff --git a/internal/domain/execution.go b/internal/domain/execution.go new file mode 100644 index 0000000..91664da --- /dev/null +++ b/internal/domain/execution.go @@ -0,0 +1,64 @@ +package domain + +import "time" + +// Classification persists every autonomous decision the classifier makes (SA5). +// This is the audit trail for "why did the OS auto-act / escalate?" +type Classification struct { + EntityID UUID + SignalEntityID UUID + TargetEntityID UUID + Action string + RecommendedAction map[string]any + RiskClass string + Route string + BlastRadius []UUID + PatternConfidence float64 + SkillID UUID + AutonomyCheck string + Reasoning map[string]any + CorrelationID string + CreatedAt time.Time +} + +// Classification routes. +const ( + RouteAutoAct = "auto-act" + RouteEscalate = "escalate" + RouteHold = "hold" +) + +// Execution is a detailed record of one action the OS performed. +type Execution struct { + EntityID UUID + ClassificationID UUID + SignalEntityID UUID + TargetEntityID UUID + Action string + RiskClass string + ApprovalID UUID + AgentID UUID + SkillID UUID + SkillVersion int + Status string + Result map[string]any + DurationMs int + Verified bool + CorrelationID string + StartedAt *time.Time + CompletedAt *time.Time + CreatedAt time.Time +} + +// Execution lifecycle states. +const ( + ExecProposed = "proposed" + ExecApproved = "approved" + ExecExecuting = "executing" + ExecVerified = "verified" + ExecFailed = "failed" + ExecTimedOut = "timed-out" + ExecRolledBack = "rolled-back" + ExecCancelled = "cancelled" + ExecExpired = "expired" +) diff --git a/internal/domain/pattern.go b/internal/domain/pattern.go new file mode 100644 index 0000000..82c2d03 --- /dev/null +++ b/internal/domain/pattern.go @@ -0,0 +1,76 @@ +package domain + +import "time" + +// Feedback records what was learned from an execution. +type Feedback struct { + EntityID UUID + ExecutionID UUID + Outcome string + Observation string + Lesson string + UnexpectedSideEffects []string + Tags []string + CreatedAt time.Time +} + +// Feedback outcomes. +const ( + OutcomeSuccess = "success" + OutcomeFailure = "failure" + OutcomePartial = "partial" + OutcomeUnexpected = "unexpected" +) + +// Pattern is a generalized rule extracted from accumulated feedback. +type Pattern struct { + EntityID UUID + AppliesType string + Action string + Pattern string + Confidence float64 + EvidenceCount int + SuccessCount int + FailureCount int + Status string + Quarantined bool + Version int + LastValidatedAt *time.Time + CreatedAt time.Time +} + +// Pattern lifecycle states. +const ( + PatternHypothesized = "hypothesized" + PatternValidated = "validated" + PatternActive = "active" + PatternDeprecated = "deprecated" + PatternInvalidated = "invalidated" +) + +// Skill is a codified procedure refined through validated patterns. +type Skill struct { + EntityID UUID + Version int + Name string + Procedure map[string]any + AppliesType string + Action string + PatternIDs []UUID + Status string + SuccessRate float64 + ChangedBy UUID + ChangeReason string + LastUsedAt *time.Time + CreatedAt time.Time +} + +// Skill lifecycle states. +const ( + SkillDrafted = "drafted" + SkillTested = "tested" + SkillActive = "active" + SkillRefined = "refined" + SkillDeprecated = "deprecated" + SkillFailed = "failed" +) diff --git a/internal/domain/signal.go b/internal/domain/signal.go new file mode 100644 index 0000000..eb9ae8a --- /dev/null +++ b/internal/domain/signal.go @@ -0,0 +1,65 @@ +package domain + +import "time" + +// Signal is an attention record — something the lab noticed that needs +// attention and possibly action. Dual entity: has an entities row + a +// signals table row for indexed querying. +type Signal struct { + EntityID UUID + Kind string + Severity string + TargetEntityID UUID + CheckID UUID + Evidence string + LikelyCause string + State string + OccurrenceCount int + FirstSeenAt time.Time + LastSeenAt time.Time + FlapCount int + HoldDownUntil *time.Time + MuteUntil *time.Time + CreatedAt time.Time + UpdatedAt time.Time +} + +// Signal lifecycle states (see lifecycle_defs in seeds/ontology.yaml). +const ( + SignalRaised = "raised" + SignalAcknowledged = "acknowledged" + SignalActing = "acting" + SignalMuted = "muted" + SignalResolved = "resolved" + SignalFailed = "failed" +) + +// Signal severities. +const ( + SeverityInfo = "info" + SeverityWarning = "warning" + SeverityCritical = "critical" +) + +// ValidSignalTransitions defines legal state transitions. +var ValidSignalTransitions = map[string][]string{ + SignalRaised: {SignalAcknowledged, SignalMuted, SignalResolved}, + SignalAcknowledged: {SignalActing, SignalResolved, SignalMuted}, + SignalActing: {SignalResolved, SignalRaised, SignalFailed}, + SignalFailed: {SignalAcknowledged}, + SignalMuted: {SignalRaised}, +} + +// CanTransition returns true if from→to is a legal signal state transition. +func (s *Signal) CanTransition(to string) bool { + allowed, ok := ValidSignalTransitions[s.State] + if !ok { + return false + } + for _, a := range allowed { + if a == to { + return true + } + } + return false +} diff --git a/internal/observability/logging.go b/internal/observability/logging.go new file mode 100644 index 0000000..2e55802 --- /dev/null +++ b/internal/observability/logging.go @@ -0,0 +1,23 @@ +package observability + +import ( + "log/slog" + "os" +) + +// NewLogger creates a structured JSON logger writing to stdout. +// In debug mode, it enables verbose probe payloads, SQL queries, and +// classification reasoning. +func NewLogger(debug bool) *slog.Logger { + level := slog.LevelInfo + if debug { + level = slog.LevelDebug + } + + handler := slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{ + Level: level, + }) + + logger := slog.New(handler).With("service", "oikos") + return logger +} diff --git a/migrations/001_ontology.up.sql b/migrations/001_ontology.up.sql new file mode 100644 index 0000000..25671fb --- /dev/null +++ b/migrations/001_ontology.up.sql @@ -0,0 +1,43 @@ +-- Migration 001: Ontology meta-schema (with inheritance, R3-1) +-- Defines entity types, relationship types, lifecycle definitions, and seed versioning. + +CREATE TABLE lifecycle_defs ( + id TEXT PRIMARY KEY, + states TEXT[] NOT NULL, + default_state TEXT NOT NULL, + terminal_states TEXT[] NOT NULL DEFAULT '{}', + transitions JSONB NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +CREATE TABLE entity_types ( + name TEXT PRIMARY KEY, + parent_type TEXT REFERENCES entity_types(name), + is_abstract BOOLEAN NOT NULL DEFAULT false, + domain TEXT NOT NULL, + layer TEXT NOT NULL CHECK (layer IN ('meta','infrastructure','governance','cognition')), + description TEXT, + lifecycle_id TEXT REFERENCES lifecycle_defs(id), + attribute_schema JSONB, + schema_version INTEGER NOT NULL DEFAULT 1, + status TEXT NOT NULL DEFAULT 'active', + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +CREATE TABLE relationship_types ( + name TEXT PRIMARY KEY, + inverse TEXT, + source_type TEXT NOT NULL REFERENCES entity_types(name), + target_type TEXT NOT NULL REFERENCES entity_types(name), + cardinality TEXT NOT NULL CHECK (cardinality IN + ('one-to-one','one-to-many','many-to-one','many-to-many')), + description TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +CREATE TABLE seed_versions ( + file TEXT PRIMARY KEY, + content_hash TEXT NOT NULL, + applied_at TIMESTAMPTZ NOT NULL DEFAULT now() +); diff --git a/migrations/002_entities.up.sql b/migrations/002_entities.up.sql new file mode 100644 index 0000000..a36ea42 --- /dev/null +++ b/migrations/002_entities.up.sql @@ -0,0 +1,50 @@ +-- Migration 002: Entity instances (UUIDv7 + slug, R3-5/D1) +-- The inventory graph: entities + typed relationships + blast_radius function. + +CREATE TABLE entities ( + id UUID PRIMARY KEY, + slug TEXT NOT NULL UNIQUE, + type TEXT NOT NULL REFERENCES entity_types(name), + name TEXT NOT NULL, + state TEXT, + attributes JSONB NOT NULL DEFAULT '{}', + maintenance_until TIMESTAMPTZ, + version INTEGER NOT NULL DEFAULT 1, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + UNIQUE (type, name) +); +CREATE INDEX idx_entities_type ON entities(type); +CREATE INDEX idx_entities_state ON entities(state); +CREATE INDEX idx_entities_attrs ON entities USING GIN(attributes); + +CREATE TABLE relationships ( + source_id UUID NOT NULL REFERENCES entities(id) ON DELETE RESTRICT, + target_id UUID NOT NULL REFERENCES entities(id) ON DELETE RESTRICT, + type TEXT NOT NULL REFERENCES relationship_types(name), + attributes JSONB, + valid_from TIMESTAMPTZ NOT NULL DEFAULT now(), + valid_to TIMESTAMPTZ, + PRIMARY KEY (source_id, target_id, type, valid_from) +); +CREATE INDEX idx_rel_source ON relationships(source_id) WHERE valid_to IS NULL; +CREATE INDEX idx_rel_target ON relationships(target_id) WHERE valid_to IS NULL; +CREATE INDEX idx_rel_type ON relationships(type) WHERE valid_to IS NULL; + +-- Cycle-safe traversal (P1): path accumulator prevents revisits; depth capped. +CREATE OR REPLACE FUNCTION blast_radius(start_id UUID, max_depth INT DEFAULT 3, + rel_types TEXT[] DEFAULT NULL) +RETURNS TABLE(entity_id UUID, depth INT) AS $$ + WITH RECURSIVE walk AS ( + SELECT start_id AS entity_id, 0 AS depth, ARRAY[start_id] AS path + UNION ALL + SELECT r.target_id, w.depth + 1, w.path || r.target_id + FROM relationships r + JOIN walk w ON r.source_id = w.entity_id + WHERE w.depth < LEAST(max_depth, 5) + AND r.valid_to IS NULL + AND NOT r.target_id = ANY(w.path) + AND (rel_types IS NULL OR r.type = ANY(rel_types)) + ) + SELECT entity_id, MIN(depth) FROM walk GROUP BY entity_id; +$$ LANGUAGE sql STABLE; diff --git a/migrations/003_operations.up.sql b/migrations/003_operations.up.sql new file mode 100644 index 0000000..3d8db2e --- /dev/null +++ b/migrations/003_operations.up.sql @@ -0,0 +1,71 @@ +-- Migration 003: Operations (signals, checks, approvals, status) +-- Signals are dual entities (entities row + signals table for indexed querying). + +CREATE TABLE check_defs ( + entity_id UUID PRIMARY KEY REFERENCES entities(id), + target_id UUID REFERENCES entities(id), + target_type TEXT REFERENCES entity_types(name), + kind TEXT NOT NULL, + config JSONB NOT NULL DEFAULT '{}', + interval_s INTEGER NOT NULL DEFAULT 600, + timeout_s INTEGER NOT NULL DEFAULT 10, + zone TEXT, + enabled BOOLEAN NOT NULL DEFAULT true, + updated_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +CREATE TABLE signals ( + entity_id UUID PRIMARY KEY REFERENCES entities(id), + kind TEXT NOT NULL, + severity TEXT NOT NULL CHECK (severity IN ('info','warning','critical')), + target_entity_id UUID REFERENCES entities(id), + check_id UUID REFERENCES check_defs(entity_id), + evidence TEXT, + likely_cause TEXT, + state TEXT NOT NULL DEFAULT 'raised', + occurrence_count INTEGER NOT NULL DEFAULT 1, + first_seen_at TIMESTAMPTZ NOT NULL DEFAULT now(), + last_seen_at TIMESTAMPTZ NOT NULL DEFAULT now(), + flap_count INTEGER NOT NULL DEFAULT 0, + hold_down_until TIMESTAMPTZ, + mute_until TIMESTAMPTZ, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT now() +); +-- At most ONE open signal per (target, kind) — repeats update the open row +CREATE UNIQUE INDEX uq_signals_open ON signals(target_entity_id, kind) + WHERE state NOT IN ('resolved','failed'); +CREATE INDEX idx_signals_state ON signals(state); + +CREATE TABLE approvals ( + entity_id UUID PRIMARY KEY REFERENCES entities(id), + subject_entity_id UUID REFERENCES entities(id), + action TEXT NOT NULL, + risk_class TEXT NOT NULL, + kind TEXT NOT NULL DEFAULT 'execution', + payload JSONB, + status TEXT NOT NULL DEFAULT 'pending', + token_hash TEXT, + expires_at TIMESTAMPTZ NOT NULL, + decided_at TIMESTAMPTZ, + decided_by UUID REFERENCES entities(id), + created_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +CREATE TABLE entity_status ( + entity_id UUID PRIMARY KEY REFERENCES entities(id), + health TEXT NOT NULL DEFAULT 'unknown', + last_check_at TIMESTAMPTZ, + details JSONB NOT NULL DEFAULT '{}', + updated_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +CREATE TABLE idempotency_keys ( + key TEXT NOT NULL, + actor TEXT NOT NULL, + request_hash TEXT NOT NULL, + response_code INTEGER, + response_body JSONB, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + PRIMARY KEY (actor, key) +); diff --git a/migrations/004_cognition.up.sql b/migrations/004_cognition.up.sql new file mode 100644 index 0000000..1de5698 --- /dev/null +++ b/migrations/004_cognition.up.sql @@ -0,0 +1,91 @@ +-- Migration 004: Cognition (classifications, executions, learning) +-- All cognition objects are dual entities (entities row + typed table). +-- Classifications persist every autonomous decision (SA5). + +CREATE TABLE classifications ( + entity_id UUID PRIMARY KEY REFERENCES entities(id), + signal_entity_id UUID REFERENCES signals(entity_id), + target_entity_id UUID REFERENCES entities(id), + action TEXT NOT NULL, + recommended_action JSONB, + risk_class TEXT NOT NULL, + route TEXT NOT NULL CHECK (route IN ('auto-act','escalate','hold')), + blast_radius UUID[], + pattern_confidence REAL, + skill_id UUID, + autonomy_check TEXT, + reasoning JSONB NOT NULL, + correlation_id TEXT NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT now() +); +CREATE INDEX idx_class_signal ON classifications(signal_entity_id); +CREATE INDEX idx_class_entity ON classifications(target_entity_id); + +CREATE TABLE executions ( + entity_id UUID PRIMARY KEY REFERENCES entities(id), + classification_id UUID REFERENCES classifications(entity_id), + signal_entity_id UUID REFERENCES signals(entity_id), + target_entity_id UUID REFERENCES entities(id), + action TEXT NOT NULL, + risk_class TEXT NOT NULL, + approval_id UUID REFERENCES approvals(entity_id), + agent_id UUID REFERENCES entities(id), + skill_id UUID, + skill_version INTEGER, + status TEXT NOT NULL DEFAULT 'proposed', + result JSONB, + duration_ms INTEGER, + verified BOOLEAN NOT NULL DEFAULT false, + correlation_id TEXT NOT NULL, + started_at TIMESTAMPTZ, + completed_at TIMESTAMPTZ, + created_at TIMESTAMPTZ NOT NULL DEFAULT now() +); +CREATE INDEX idx_exec_target ON executions(target_entity_id); +CREATE INDEX idx_exec_status ON executions(status); + +CREATE TABLE feedback ( + entity_id UUID PRIMARY KEY REFERENCES entities(id), + execution_id UUID NOT NULL REFERENCES executions(entity_id), + outcome TEXT NOT NULL CHECK (outcome IN ('success','failure','partial','unexpected')), + observation TEXT, + lesson TEXT, + unexpected_side_effects TEXT[], + tags TEXT[], + created_at TIMESTAMPTZ NOT NULL DEFAULT now() +); +CREATE INDEX idx_feedback_ts ON feedback(created_at); + +CREATE TABLE patterns ( + entity_id UUID PRIMARY KEY REFERENCES entities(id), + applies_type TEXT NOT NULL REFERENCES entity_types(name), + action TEXT NOT NULL, + pattern TEXT NOT NULL, + confidence REAL NOT NULL DEFAULT 0, + evidence_count INTEGER NOT NULL DEFAULT 0, + success_count INTEGER NOT NULL DEFAULT 0, + failure_count INTEGER NOT NULL DEFAULT 0, + status TEXT NOT NULL DEFAULT 'hypothesized', + quarantined BOOLEAN NOT NULL DEFAULT false, + version INTEGER NOT NULL DEFAULT 1, + last_validated_at TIMESTAMPTZ, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + UNIQUE (applies_type, action) +); + +CREATE TABLE skills ( + entity_id UUID NOT NULL REFERENCES entities(id), + version INTEGER NOT NULL DEFAULT 1, + name TEXT NOT NULL, + procedure JSONB NOT NULL, + applies_type TEXT REFERENCES entity_types(name), + action TEXT NOT NULL, + pattern_ids UUID[], + status TEXT NOT NULL DEFAULT 'drafted', + success_rate REAL, + changed_by UUID, + change_reason TEXT, + last_used_at TIMESTAMPTZ, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + PRIMARY KEY (entity_id, version) +); diff --git a/migrations/005_policy.up.sql b/migrations/005_policy.up.sql new file mode 100644 index 0000000..ebbdf99 --- /dev/null +++ b/migrations/005_policy.up.sql @@ -0,0 +1,28 @@ +-- Migration 005: Policy (risk classes, approval rules, autonomy settings) + +CREATE TABLE risk_classes ( + name TEXT PRIMARY KEY, + description TEXT, + approval_required TEXT NOT NULL DEFAULT 'none', + autonomy_allowed BOOLEAN NOT NULL DEFAULT false +); + +CREATE TABLE approval_rules ( + id UUID PRIMARY KEY, + entity_type TEXT REFERENCES entity_types(name), + action TEXT NOT NULL, + risk_class TEXT NOT NULL REFERENCES risk_classes(name), + autonomy_level TEXT NOT NULL DEFAULT 'escalate' CHECK + (autonomy_level IN ('auto','escalate','never')), + scope_entity UUID REFERENCES entities(id), + version INTEGER NOT NULL DEFAULT 1, + updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + UNIQUE (entity_type, action, scope_entity) +); + +CREATE TABLE autonomy_settings ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL, + version INTEGER NOT NULL DEFAULT 1, + updated_at TIMESTAMPTZ NOT NULL DEFAULT now() +); diff --git a/migrations/006_observability.up.sql b/migrations/006_observability.up.sql new file mode 100644 index 0000000..2b37ac2 --- /dev/null +++ b/migrations/006_observability.up.sql @@ -0,0 +1,138 @@ +-- Migration 006: Observability (TimescaleDB) +-- Hypertable PKs include time column (SG1); idempotent DDL (SG3); no array_agg in CAGGs (SG2). + +-- Enable TimescaleDB extension +CREATE EXTENSION IF NOT EXISTS timescaledb; + +-- ─── Time-series metrics ────────────────────────────────────────────── +CREATE TABLE metric_samples ( + ts TIMESTAMPTZ NOT NULL, + entity_id UUID NOT NULL, + metric TEXT NOT NULL, + value DOUBLE PRECISION NOT NULL, + tags JSONB NOT NULL DEFAULT '{}' +); +SELECT create_hypertable('metric_samples', 'ts', + chunk_time_interval => INTERVAL '7 days', if_not_exists => TRUE); +CREATE INDEX idx_metrics_entity_ts ON metric_samples(entity_id, ts DESC); +CREATE INDEX idx_metrics_metric_ts ON metric_samples(metric, ts DESC); +DO $$ BEGIN + PERFORM add_retention_policy('metric_samples', INTERVAL '90 days'); +EXCEPTION WHEN OTHERS THEN NULL; +END $$; + +-- 1-hour rollups +CREATE MATERIALIZED VIEW metric_rollups_1h WITH (timescaledb.continuous) AS + SELECT time_bucket('1 hour', ts) AS bucket, entity_id, metric, + avg(value) AS avg_value, min(value) AS min_value, + max(value) AS max_value, count(*) AS sample_count + FROM metric_samples GROUP BY bucket, entity_id, metric; +DO $$ BEGIN + PERFORM add_continuous_aggregate_policy('metric_rollups_1h', + start_offset => INTERVAL '2 hours', + end_offset => INTERVAL '5 minutes', + schedule_interval => INTERVAL '1 hour'); +EXCEPTION WHEN OTHERS THEN NULL; +END $$; + +-- 1-day rollups +CREATE MATERIALIZED VIEW metric_rollups_1d WITH (timescaledb.continuous) AS + SELECT time_bucket('1 day', ts) AS bucket, entity_id, metric, + avg(value) AS avg_value, min(value) AS min_value, + max(value) AS max_value, count(*) AS sample_count + FROM metric_samples GROUP BY bucket, entity_id, metric; +DO $$ BEGIN + PERFORM add_continuous_aggregate_policy('metric_rollups_1d', + start_offset => INTERVAL '2 days', + end_offset => INTERVAL '1 hour', + schedule_interval => INTERVAL '1 day'); +EXCEPTION WHEN OTHERS THEN NULL; +END $$; + +-- ─── Audit log ──────────────────────────────────────────────────────── +CREATE TABLE audit_log ( + id BIGINT GENERATED ALWAYS AS IDENTITY, + ts TIMESTAMPTZ NOT NULL DEFAULT now(), + actor_type TEXT NOT NULL, + actor_id UUID, + action TEXT NOT NULL, + entity_id UUID, + method TEXT, + path TEXT, + status_code INTEGER, + detail JSONB NOT NULL DEFAULT '{}', + source_ip TEXT, + correlation_id TEXT, + PRIMARY KEY (id, ts) +); +SELECT create_hypertable('audit_log', 'ts', + chunk_time_interval => INTERVAL '7 days', if_not_exists => TRUE); +CREATE INDEX idx_audit_actor ON audit_log(actor_type, actor_id, ts DESC); +CREATE INDEX idx_audit_entity ON audit_log(entity_id, ts DESC); +CREATE INDEX idx_audit_action ON audit_log(action, ts DESC); +CREATE INDEX idx_audit_correlation ON audit_log(correlation_id); +DO $$ BEGIN + PERFORM add_retention_policy('audit_log', INTERVAL '365 days'); +EXCEPTION WHEN OTHERS THEN NULL; +END $$; + +-- ─── Event log ──────────────────────────────────────────────────────── +CREATE TABLE events ( + id BIGINT GENERATED ALWAYS AS IDENTITY, + ts TIMESTAMPTZ NOT NULL DEFAULT now(), + type TEXT NOT NULL, + entity_id UUID, + severity TEXT NOT NULL DEFAULT 'info', + source TEXT NOT NULL, + data JSONB NOT NULL DEFAULT '{}', + correlation_id TEXT, + PRIMARY KEY (id, ts) +); +SELECT create_hypertable('events', 'ts', + chunk_time_interval => INTERVAL '7 days', if_not_exists => TRUE); +CREATE INDEX idx_events_type_ts ON events(type, ts DESC); +CREATE INDEX idx_events_entity_ts ON events(entity_id, ts DESC); +CREATE INDEX idx_events_severity_ts ON events(severity, ts DESC); +CREATE INDEX idx_events_correlation ON events(correlation_id); +DO $$ BEGIN + PERFORM add_retention_policy('events', INTERVAL '90 days'); +EXCEPTION WHEN OTHERS THEN NULL; +END $$; + +-- ─── Agent activity ────────────────────────────────────────────────── +CREATE TABLE agent_activity ( + id BIGINT GENERATED ALWAYS AS IDENTITY, + ts TIMESTAMPTZ NOT NULL DEFAULT now(), + agent_id UUID NOT NULL, + session_id TEXT, + activity_type TEXT NOT NULL, + tool_name TEXT, + entity_id UUID, + input_summary TEXT, + output_summary TEXT, + duration_ms INTEGER, + token_count INTEGER, + success BOOLEAN, + correlation_id TEXT, + PRIMARY KEY (id, ts) +); +SELECT create_hypertable('agent_activity', 'ts', + chunk_time_interval => INTERVAL '7 days', if_not_exists => TRUE); +CREATE INDEX idx_agent_activity_agent_ts ON agent_activity(agent_id, ts DESC); +CREATE INDEX idx_agent_activity_type_ts ON agent_activity(activity_type, ts DESC); +CREATE INDEX idx_agent_activity_entity ON agent_activity(entity_id, ts DESC); +CREATE INDEX idx_agent_activity_correlation ON agent_activity(correlation_id); +DO $$ BEGIN + PERFORM add_retention_policy('agent_activity', INTERVAL '90 days'); +EXCEPTION WHEN OTHERS THEN NULL; +END $$; + +-- ─── Ledger view (R3-11: not a fourth write path) ──────────────────── +CREATE VIEW ledger AS +SELECT e.created_at AS ts, e.entity_id AS execution_id, e.target_entity_id, + e.action, e.risk_class, e.status, e.verified, + c.route, c.reasoning, a.status AS approval_status, a.decided_by, + e.agent_id, e.correlation_id +FROM executions e +LEFT JOIN classifications c ON c.entity_id = e.classification_id +LEFT JOIN approvals a ON a.entity_id = e.approval_id; diff --git a/migrations/embed.go b/migrations/embed.go new file mode 100644 index 0000000..fca6375 --- /dev/null +++ b/migrations/embed.go @@ -0,0 +1,7 @@ +// Package migrations embeds SQL migration files for use by the db package. +package migrations + +import "embed" + +//go:embed *.up.sql +var FS embed.FS