From b7968db4e06efd6c5ba67e801e591c1d77eb924f Mon Sep 17 00:00:00 2001 From: Denozordec Date: Sun, 5 Apr 2026 17:42:07 +0700 Subject: [PATCH] refactor: update worker processes in EvoBGP to accept dependencies for shared store and job registry. Enhance scheduler, ingest, render, and deploy components to utilize a unified context and improve logging for drift detection. Update architecture documentation to reflect changes in process interactions and worker functionalities. --- cmd/evobgp-all/main.go | 14 +- cmd/evobgp-deploy/main.go | 20 ++- cmd/evobgp-ingest/main.go | 20 ++- cmd/evobgp-render/main.go | 20 ++- cmd/evobgp-scheduler/main.go | 35 ++++- deploy/compose/docker-compose.yaml | 6 + docs/architecture.md | 29 ++-- docs/quickstart.md | 4 +- internal/deploy/run.go | 58 +++++++- internal/httpapi/bootstrap.go | 52 +++++++ internal/httpapi/server.go | 36 +---- internal/ingest/run.go | 37 ++++- internal/jobs/worker.go | 26 +++- internal/pipeline/parse.go | 47 +++++++ internal/pipeline/prefetch.go | 61 ++++++++ internal/pipeline/refresh.go | 214 +++++++++++++++++++++++++++++ internal/render/publish.go | 31 +++++ internal/render/run.go | 41 +++++- internal/repository/postgres.go | 76 ++++++++++ internal/scheduler/run.go | 125 ++++++++++++++++- internal/store/backend.go | 5 + internal/store/memory.go | 56 ++++++++ 22 files changed, 936 insertions(+), 77 deletions(-) create mode 100644 internal/httpapi/bootstrap.go create mode 100644 internal/pipeline/parse.go create mode 100644 internal/pipeline/prefetch.go create mode 100644 internal/pipeline/refresh.go create mode 100644 internal/render/publish.go diff --git a/cmd/evobgp-all/main.go b/cmd/evobgp-all/main.go index b02cc02..5913045 100644 --- a/cmd/evobgp-all/main.go +++ b/cmd/evobgp-all/main.go @@ -21,7 +21,7 @@ import ( "evobgp/internal/scheduler" ) -// microVPS entrypoint: one process — HTTP API (same as evobgp-api) plus in-process stubs for scheduler/ingest/render/deploy. +// microVPS entrypoint: one process — HTTP API plus in-process scheduler, ingest, render, deploy workers (shared store + job registry). func main() { cfg := config.Load() opts := httpapi.Options{ @@ -42,10 +42,14 @@ func main() { ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() - go scheduler.Run(ctx) - go ingest.Run(ctx) - go render.Run(ctx) - go deploy.Run(ctx) + schedDeps := &scheduler.Deps{Store: srv.Store(), Jobs: srv.Jobs()} + ingestDeps := &ingest.Deps{Store: srv.Store()} + renderDeps := &render.Deps{Store: srv.Store()} + deployDeps := &deploy.Deps{Store: srv.Store()} + go scheduler.Run(ctx, schedDeps) + go ingest.Run(ctx, ingestDeps) + go render.Run(ctx, renderDeps) + go deploy.Run(ctx, deployDeps) startBirdMetricsPoller() diff --git a/cmd/evobgp-deploy/main.go b/cmd/evobgp-deploy/main.go index 7539075..2b95996 100644 --- a/cmd/evobgp-deploy/main.go +++ b/cmd/evobgp-deploy/main.go @@ -5,10 +5,13 @@ import ( "log" "os" "os/signal" + "strings" "syscall" + "time" "evobgp/internal/config" "evobgp/internal/deploy" + "evobgp/internal/httpapi" ) func main() { @@ -20,5 +23,20 @@ func main() { log.Printf("%s starting (reference profile worker)", name) ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() - deploy.Run(ctx) + + bctx, bcancel := context.WithTimeout(context.Background(), 90*time.Second) + defer bcancel() + opts := httpapi.Options{ + DatabaseURL: strings.TrimSpace(os.Getenv("EVOBGP_DATABASE_URL")), + SeedDemo: os.Getenv("EVOBGP_SEED_DEMO") != "0", + } + st, _, pool, err := httpapi.BootstrapWorkers(bctx, opts) + if err != nil { + log.Fatal(err) + } + if pool != nil { + defer pool.Close() + } + + deploy.Run(ctx, &deploy.Deps{Store: st}) } diff --git a/cmd/evobgp-ingest/main.go b/cmd/evobgp-ingest/main.go index 4d424ff..31eff7b 100644 --- a/cmd/evobgp-ingest/main.go +++ b/cmd/evobgp-ingest/main.go @@ -5,9 +5,12 @@ import ( "log" "os" "os/signal" + "strings" "syscall" + "time" "evobgp/internal/config" + "evobgp/internal/httpapi" "evobgp/internal/ingest" ) @@ -20,5 +23,20 @@ func main() { log.Printf("%s starting (reference profile worker)", name) ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() - ingest.Run(ctx) + + bctx, bcancel := context.WithTimeout(context.Background(), 90*time.Second) + defer bcancel() + opts := httpapi.Options{ + DatabaseURL: strings.TrimSpace(os.Getenv("EVOBGP_DATABASE_URL")), + SeedDemo: os.Getenv("EVOBGP_SEED_DEMO") != "0", + } + st, _, pool, err := httpapi.BootstrapWorkers(bctx, opts) + if err != nil { + log.Fatal(err) + } + if pool != nil { + defer pool.Close() + } + + ingest.Run(ctx, &ingest.Deps{Store: st}) } diff --git a/cmd/evobgp-render/main.go b/cmd/evobgp-render/main.go index 99ae92c..e745534 100644 --- a/cmd/evobgp-render/main.go +++ b/cmd/evobgp-render/main.go @@ -5,9 +5,12 @@ import ( "log" "os" "os/signal" + "strings" "syscall" + "time" "evobgp/internal/config" + "evobgp/internal/httpapi" "evobgp/internal/render" ) @@ -20,5 +23,20 @@ func main() { log.Printf("%s starting (reference profile worker)", name) ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() - render.Run(ctx) + + bctx, bcancel := context.WithTimeout(context.Background(), 90*time.Second) + defer bcancel() + opts := httpapi.Options{ + DatabaseURL: strings.TrimSpace(os.Getenv("EVOBGP_DATABASE_URL")), + SeedDemo: os.Getenv("EVOBGP_SEED_DEMO") != "0", + } + st, _, pool, err := httpapi.BootstrapWorkers(bctx, opts) + if err != nil { + log.Fatal(err) + } + if pool != nil { + defer pool.Close() + } + + render.Run(ctx, &render.Deps{Store: st}) } diff --git a/cmd/evobgp-scheduler/main.go b/cmd/evobgp-scheduler/main.go index 071dc94..ef8fdb0 100644 --- a/cmd/evobgp-scheduler/main.go +++ b/cmd/evobgp-scheduler/main.go @@ -3,11 +3,15 @@ package main import ( "context" "log" + "net/http" "os" "os/signal" + "strings" "syscall" + "time" "evobgp/internal/config" + "evobgp/internal/httpapi" "evobgp/internal/scheduler" ) @@ -20,5 +24,34 @@ func main() { log.Printf("%s starting (reference profile worker)", name) ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() - scheduler.Run(ctx) + + bctx, bcancel := context.WithTimeout(context.Background(), 90*time.Second) + defer bcancel() + opts := httpapi.Options{ + DatabaseURL: strings.TrimSpace(os.Getenv("EVOBGP_DATABASE_URL")), + SeedDemo: os.Getenv("EVOBGP_SEED_DEMO") != "0", + } + st, reg, pool, err := httpapi.BootstrapWorkers(bctx, opts) + if err != nil { + log.Fatal(err) + } + if pool != nil { + defer pool.Close() + } + + apiBase := strings.TrimSpace(os.Getenv("EVOBGP_CONTROL_PLANE_URL")) + apiTok := strings.TrimSpace(os.Getenv("EVOBGP_SCHEDULER_BEARER")) + deps := &scheduler.Deps{ + Store: st, + HTTP: &http.Client{Timeout: 45 * time.Second}, + } + if apiBase != "" && apiTok != "" { + deps.APIBase = apiBase + deps.APIToken = apiTok + log.Printf("evobgp-scheduler: control plane HTTP mode (%s)", apiBase) + } else { + deps.Jobs = reg + log.Printf("evobgp-scheduler: in-process job registry mode (use EVOBGP_CONTROL_PLANE_URL + EVOBGP_SCHEDULER_BEARER for split containers)") + } + scheduler.Run(ctx, deps) } diff --git a/deploy/compose/docker-compose.yaml b/deploy/compose/docker-compose.yaml index 8620921..322ac1e 100644 --- a/deploy/compose/docker-compose.yaml +++ b/deploy/compose/docker-compose.yaml @@ -91,6 +91,8 @@ services: <<: *env-ref EVOBGP_HTTP_ADDR: ":8080" EVOBGP_SEED_DEMO: "1" + # Local reference only: allows Bearer dev for scheduler HTTP client (EVOBGP_SCHEDULER_BEARER). + EVOBGP_DEV_INSECURE: "1" EVOBGP_BIRDC_SOCKET: /run/bird/bird.ctl EVOBGP_BIRDC_INTERVAL: 30s EVOBGP_BIRD_ACTIVE_DIR: /etc/bird @@ -120,8 +122,12 @@ services: condition: service_healthy nats: condition: service_started + evobgp-api: + condition: service_started environment: <<: *env-ref + EVOBGP_CONTROL_PLANE_URL: http://evobgp-api:8080 + EVOBGP_SCHEDULER_BEARER: dev logging: *default-logging deploy: resources: diff --git a/docs/architecture.md b/docs/architecture.md index d33c16c..4329a3b 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -13,11 +13,11 @@ | Бинарник | Роль | |----------|------| | `evobgp-api` | Только HTTP API и связанная логика в одном процессе. | -| `evobgp-all` | Режим одной VPS: тот же API + in-process запуск заглушек scheduler, ingest, render, deploy. | -| `evobgp-scheduler` | Планировщик cron/интервалов модулей (в коде сейчас **stub**). | -| `evobgp-ingest` | Воркеры загрузки внешних источников (CDN и т.д.) (**stub**). | -| `evobgp-render` | Генерация артефактов BIRD из ревизий (**stub**). | -| `evobgp-deploy` | Выкладка на спикеры / взаимодействие с BIRD на стороне деплоя (**stub**). | +| `evobgp-all` | Тот же API + in-process **scheduler** (очередь `module_refresh` в общем Registry), **ingest** (prefetch ETag CDN), **render** (опционально auto-publish), **deploy** (лог расхождений published/applied). | +| `evobgp-scheduler` | По `refresh_interval_sec` ставит refresh: в одном процессе с API — через `jobs.Registry`; в reference Compose — **HTTP** `POST /v1/modules/{id}/refresh` (`EVOBGP_CONTROL_PLANE_URL`, `EVOBGP_SCHEDULER_BEARER`). | +| `evobgp-ingest` | Периодический conditional GET по URL CDN-источников и обновление `etag` в БД. | +| `evobgp-render` | По умолчанию только heartbeat; при `EVOBGP_RENDER_AUTOPUBLISH=1` выставляет всем спикерам tenant последнюю ревизию (упрощение для демо). | +| `evobgp-deploy` | Периодически логирует **drift**: `last_applied_revision_id` vs опубликованная ревизия для ноды. | | `evobgp-node` | CLI реплики: `pull-bundle`, `verify-bundle`, `apply-bundle`. | | `evobgp-agent` | Локальный агент рядом с BIRD (например `watch` по сокету). | @@ -39,6 +39,7 @@ | `config` | Переменные окружения `EVOBGP_*`. | | `observability` | Метрики Prometheus, HTTP middleware. | | `broker` | Заготовка под NATS/Redis (логирование подключения в воркерах). | +| `pipeline` | Ingest+render в одном шаге для `module_refresh`: выборка префиксов (CDN/AS/IP/пустые DOMAINS), `CreateRenderRevision`, превью BIRD через `birdfmt`. | ## Диаграмма: эталонный Compose (reference) @@ -51,12 +52,11 @@ flowchart LR end subgraph control [Control_plane] API[evobgp_api] - Sched[evobgp_scheduler_stub] - Ingest[evobgp_ingest_stub] - Render[evobgp_render_stub] - Deploy[evobgp_deploy_stub] + Sched[evobgp_scheduler] + Ingest[evobgp_ingest] + Render[evobgp_render] + Deploy[evobgp_deploy] PG[(PostgreSQL)] - NATS[NATS_JetStream] end subgraph data [Data_plane] BIRD[BIRD2] @@ -66,10 +66,7 @@ flowchart LR Operator --> API NodeCLI --> API API --> PG - Sched --> NATS - Ingest --> NATS - Render --> NATS - Deploy --> NATS + Sched -->|HTTP_or_DB| API Sched --> PG Ingest --> PG Render --> PG @@ -77,7 +74,7 @@ flowchart LR Agent --> BIRD ``` -На практике воркеры **пока не выполняют** полноценную работу с очередью — они резервируют место в топологии и пишут в лог. API и БД уже обеспечивают основной сценарий разработки и тестов. +Очередь задач по-прежнему **in-memory в процессе API** (`jobs.Registry`); отдельный контейнер `evobgp-scheduler` не разделяет память с API и дергает refresh по HTTP. Полноценный брокер (NATS) и общая очередь `job_audit` между процессами — в следующих итерациях. ## Диаграмма: microvps (`evobgp-all`) @@ -92,7 +89,7 @@ flowchart LR All --> BIRD ``` -Внутри процесса `evobgp-all` горутины scheduler/ingest/render/deploy — те же **stub**, что и отдельные бинарники. +Внутри `evobgp-all` все воркеры используют **тот же** `store` и `jobs.Registry`, что и HTTP handlers, поэтому `module_refresh` выполняется в том же процессе без HTTP. ## Поток: ревизия и бандл для ноды diff --git a/docs/quickstart.md b/docs/quickstart.md index 5fbcfd7..b93a798 100644 --- a/docs/quickstart.md +++ b/docs/quickstart.md @@ -60,7 +60,9 @@ docker compose --profile reference up -d --build | 9090 | Prometheus (в compose) | | 179 | BGP (BIRD2) | -**Важно:** процессы `evobgp-scheduler`, `evobgp-ingest`, `evobgp-render`, `evobgp-deploy` в текущей версии кода — **заглушки** (логирование и периодический тик). Реальная очередь задач и брокер подключаются в будущих итерациях; API и БД при этом уже работают. +**Воркеры reference:** `evobgp-scheduler` ходит в API по HTTP (`EVOBGP_CONTROL_PLANE_URL`, `EVOBGP_SCHEDULER_BEARER`); в [docker-compose.yaml](../deploy/compose/docker-compose.yaml) для локального запуска включены `EVOBGP_DEV_INSECURE=1` на API и токен `dev` у планировщика. `evobgp-ingest` обновляет ETag CDN-источников; `evobgp-render` по умолчанию не трогает `published_revision` (включите `EVOBGP_RENDER_AUTOPUBLISH=1` осознанно); `evobgp-deploy` пишет в лог расхождение applied vs published. Очередь `jobs` остаётся in-process у **evobgp-api**; общий брокер — в планах. + +В **evobgp-all** (microvps) те же пакеты крутятся в одном процессе и используют общий `jobs.Registry` без HTTP. ## Вариант 3: Локально без Docker (только API) diff --git a/internal/deploy/run.go b/internal/deploy/run.go index 508bd24..be377c6 100644 --- a/internal/deploy/run.go +++ b/internal/deploy/run.go @@ -1,4 +1,3 @@ -// Package deploy delivers generated BIRD fragments to evobgp-agent / publishes signed bundles for evobgp-node. package deploy import ( @@ -10,18 +9,67 @@ import ( "evobgp/internal/broker" "evobgp/internal/config" + "evobgp/internal/store" ) -// Run blocks until ctx is cancelled. Reference deployment: worker after render, talks to agent API or shared volume. -func Run(ctx context.Context) { +// Deps enables deploy-side drift logging between published and last-applied revision per speaker. +type Deps struct { + Store store.Backend +} + +// Run blocks until ctx is cancelled. +func Run(ctx context.Context, deps *Deps) { cfg := config.Load() broker.LogConnect(ctx, cfg.BrokerURL) if d := strings.TrimSpace(os.Getenv("EVOBGP_BIRD_ACTIVE_DIR")); d != "" { - log.Printf("evobgp-deploy: EVOBGP_BIRD_ACTIVE_DIR=%q (apply handled by API jobs + birddeploy when set on API)", d) + log.Printf("evobgp-deploy: EVOBGP_BIRD_ACTIVE_DIR=%q (apply via API jobs when API has same env)", d) } + if deps == nil || deps.Store == nil { + runStub(ctx) + return + } + t := time.NewTicker(90 * time.Second) + defer t.Stop() + log.Printf("evobgp-deploy: active (speaker published vs applied drift log)") + for { + select { + case <-ctx.Done(): + log.Printf("evobgp-deploy: stopped") + return + case <-t.C: + logDrift(context.Background(), deps.Store) + } + } +} + +func logDrift(ctx context.Context, st store.Backend) { + _ = ctx + tenants, err := st.ListTenantIDs() + if err != nil { + log.Printf("evobgp-deploy: list tenants: %v", err) + return + } + for _, tid := range tenants { + for _, sp := range st.ListSpeakersForTenant(tid) { + pub, _, err := st.LatestPublishedRevision(sp.ID) + if err != nil { + continue + } + applied := "" + if sp.LastAppliedRevisionID != nil { + applied = *sp.LastAppliedRevisionID + } + if applied != "" && applied != pub { + log.Printf("evobgp-deploy: drift speaker=%s applied=%s published=%s", sp.ID, applied, pub) + } + } + } +} + +func runStub(ctx context.Context) { t := time.NewTicker(60 * time.Second) defer t.Stop() - log.Printf("evobgp-deploy: started (orchestration stub; two-phase apply runs in evobgp-api/evobgp-all job worker when EVOBGP_BIRD_ACTIVE_DIR is set)") + log.Printf("evobgp-deploy: idle stub (no store in Deps)") for { select { case <-ctx.Done(): diff --git a/internal/httpapi/bootstrap.go b/internal/httpapi/bootstrap.go new file mode 100644 index 0000000..566b2d4 --- /dev/null +++ b/internal/httpapi/bootstrap.go @@ -0,0 +1,52 @@ +package httpapi + +import ( + "context" + "net/http" + "strings" + "time" + + "evobgp/internal/db" + "evobgp/internal/jobs" + "evobgp/internal/observability" + "evobgp/internal/repository" + "evobgp/internal/store" + + "github.com/jackc/pgx/v5/pgxpool" +) + +// BootstrapWorkers opens the same store.Backend and jobs.Registry as New (without HTTP or bundle keys). +// Used by standalone worker binaries (scheduler, ingest, …) that share PostgreSQL with the API. +func BootstrapWorkers(ctx context.Context, opts Options) (store.Backend, *jobs.Registry, *pgxpool.Pool, error) { + ctx, cancel := context.WithTimeout(ctx, 60*time.Second) + defer cancel() + + var backend store.Backend + var pool *pgxpool.Pool + + if u := strings.TrimSpace(opts.DatabaseURL); u != "" { + p, err := db.OpenPostgresPool(ctx, u) + if err != nil { + return nil, nil, nil, err + } + pool = p + pgbe, err := repository.NewPostgres(ctx, p, opts.SeedDemo) + if err != nil { + pool.Close() + return nil, nil, nil, err + } + backend = pgbe + } else { + mem := store.NewMemory() + if opts.SeedDemo { + mem.SeedDemo() + } + backend = mem + } + + cdnHTTP := &http.Client{Timeout: 45 * time.Second} + wk := &jobs.Worker{Store: backend, HTTPClient: cdnHTTP} + reg := jobs.NewRegistry(wk.Process) + observability.RegisterStoreBackend(backend) + return backend, reg, pool, nil +} diff --git a/internal/httpapi/server.go b/internal/httpapi/server.go index 58c96ac..7a76b5f 100644 --- a/internal/httpapi/server.go +++ b/internal/httpapi/server.go @@ -10,10 +10,7 @@ import ( "strings" "time" - "evobgp/internal/db" "evobgp/internal/jobs" - "evobgp/internal/observability" - "evobgp/internal/repository" "evobgp/internal/store" "github.com/jackc/pgx/v5/pgxpool" @@ -44,34 +41,11 @@ type Options struct { // New constructs Server and wiring for async jobs. func New(opts Options) (*Server, error) { - var backend store.Backend - var pool *pgxpool.Pool - ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) - defer cancel() - - if u := strings.TrimSpace(opts.DatabaseURL); u != "" { - p, err := db.OpenPostgresPool(ctx, u) - if err != nil { - return nil, err - } - pool = p - pgbe, err := repository.NewPostgres(ctx, p, opts.SeedDemo) - if err != nil { - pool.Close() - return nil, err - } - backend = pgbe - } else { - mem := store.NewMemory() - if opts.SeedDemo { - mem.SeedDemo() - } - backend = mem + backend, pool, reg, err := BootstrapWorkers(context.Background(), opts) + if err != nil { + return nil, err } - wk := &jobs.Worker{Store: backend} - reg := jobs.NewRegistry(wk.Process) - var priv ed25519.PrivateKey if strings.TrimSpace(opts.BundleSeedHex) != "" { seed, err := hex.DecodeString(strings.TrimSpace(opts.BundleSeedHex)) @@ -95,7 +69,6 @@ func New(opts Options) (*Server, error) { insecureDev: opts.InsecureDev && opts.SeedDemo, corsOrigins: parseCORSOrigins(opts.CORSAllowedOrigins), } - observability.RegisterStoreBackend(backend) s.mux = http.NewServeMux() s.registerRoutes() return s, nil @@ -110,3 +83,6 @@ func (s *Server) Close() { // Store exposes the backing store (for operators / tests). func (s *Server) Store() store.Backend { return s.store } + +// Jobs exposes the in-process async job registry (for scheduler / evobgp-all). +func (s *Server) Jobs() *jobs.Registry { return s.jobs } diff --git a/internal/ingest/run.go b/internal/ingest/run.go index 2641b64..7da4470 100644 --- a/internal/ingest/run.go +++ b/internal/ingest/run.go @@ -1,22 +1,51 @@ -// Package ingest fetches CDN lists, resolves DoH domains, and writes normalized prefixes into the database. package ingest import ( "context" "log" + "net/http" "time" "evobgp/internal/broker" "evobgp/internal/config" + "evobgp/internal/pipeline" + "evobgp/internal/store" ) -// Run blocks until ctx is cancelled. Reference deployment: separate process consuming the job queue. -func Run(ctx context.Context) { +// Deps runs lightweight CDN ETag prefetch against the shared store. +type Deps struct { + Store store.Backend +} + +// Run blocks until ctx is cancelled. +func Run(ctx context.Context, deps *Deps) { cfg := config.Load() broker.LogConnect(ctx, cfg.BrokerURL) + if deps == nil || deps.Store == nil { + runStub(ctx) + return + } + hc := &http.Client{Timeout: 45 * time.Second} t := time.NewTicker(60 * time.Second) defer t.Stop() - log.Printf("evobgp-ingest: started (stub; ingest workers dequeue from broker or job_audit)") + log.Printf("evobgp-ingest: active (CDN conditional GET / ETag prefetch)") + for { + select { + case <-ctx.Done(): + log.Printf("evobgp-ingest: stopped") + return + case <-t.C: + if err := pipeline.PrefetchCDNSourceETags(context.Background(), deps.Store, hc); err != nil { + log.Printf("evobgp-ingest: prefetch: %v", err) + } + } + } +} + +func runStub(ctx context.Context) { + t := time.NewTicker(60 * time.Second) + defer t.Stop() + log.Printf("evobgp-ingest: idle stub (no store in Deps)") for { select { case <-ctx.Done(): diff --git a/internal/jobs/worker.go b/internal/jobs/worker.go index f921436..e68283b 100644 --- a/internal/jobs/worker.go +++ b/internal/jobs/worker.go @@ -2,12 +2,15 @@ package jobs import ( "context" + "net/http" "os" "strings" + "time" "evobgp/internal/birddeploy" "evobgp/internal/birdfmt" "evobgp/internal/observability" + "evobgp/internal/pipeline" "evobgp/internal/store" ) @@ -20,7 +23,17 @@ const ( // Worker executes queued jobs against store.Backend (memory or SQL). type Worker struct { - Store store.Backend + Store store.Backend + HTTPClient *http.Client // optional; CDN refresh uses this (default 45s timeout). +} + +var defaultWorkerHTTP = &http.Client{Timeout: 45 * time.Second} + +func (w *Worker) httpClient() *http.Client { + if w != nil && w.HTTPClient != nil { + return w.HTTPClient + } + return defaultWorkerHTTP } // Process is registered as Registry.workerStart. @@ -42,6 +55,17 @@ func (w *Worker) Process(j *Job) { switch j.Kind { case KindModuleRefresh: + mid, _ := j.Meta["module_id"].(string) + if strings.TrimSpace(mid) == "" { + j.Fail("missing module_id in job meta") + return + } + rev, err := pipeline.RefreshModule(context.Background(), w.Store, w.httpClient(), j.TenantID, mid) + if err != nil { + j.Fail(err.Error()) + return + } + j.mergeMeta(map[string]any{"revision_id": rev}) j.Succeed() case KindDeployApply: w.runDeployApply(j) diff --git a/internal/pipeline/parse.go b/internal/pipeline/parse.go new file mode 100644 index 0000000..d1b495c --- /dev/null +++ b/internal/pipeline/parse.go @@ -0,0 +1,47 @@ +package pipeline + +import ( + "bufio" + "net/netip" + "strings" +) + +// ParseCIDRLines extracts unique IPv4/IPv6 CIDRs from plain text (one per line, # comments, empty lines skipped). +func ParseCIDRLines(body string) []netip.Prefix { + seen := make(map[string]struct{}) + var out []netip.Prefix + sc := bufio.NewScanner(strings.NewReader(body)) + for sc.Scan() { + line := strings.TrimSpace(sc.Text()) + if line == "" || strings.HasPrefix(line, "#") { + continue + } + pfx := parseOneCIDR(line) + if !pfx.IsValid() { + continue + } + m := pfx.Masked() + s := m.String() + if _, ok := seen[s]; ok { + continue + } + seen[s] = struct{}{} + out = append(out, m) + } + return out +} + +func parseOneCIDR(s string) netip.Prefix { + if p, err := netip.ParsePrefix(s); err == nil { + return p + } + if addr, err := netip.ParseAddr(s); err == nil { + if addr.Is4() { + p, _ := addr.Prefix(32) + return p + } + p, _ := addr.Prefix(128) + return p + } + return netip.Prefix{} +} diff --git a/internal/pipeline/prefetch.go b/internal/pipeline/prefetch.go new file mode 100644 index 0000000..ad4be87 --- /dev/null +++ b/internal/pipeline/prefetch.go @@ -0,0 +1,61 @@ +package pipeline + +import ( + "context" + "io" + "net/http" + "strings" + + "evobgp/internal/store" +) + +// PrefetchCDNSourceETags performs conditional GETs for CDN module sources and updates stored ETags when the origin responds 200. +func PrefetchCDNSourceETags(ctx context.Context, st store.Backend, hc *http.Client) error { + if hc == nil { + hc = http.DefaultClient + } + tenants, err := st.ListTenantIDs() + if err != nil { + return err + } + for _, tid := range tenants { + for _, mod := range st.ListModules(tid) { + if !mod.Enabled || mod.Type != "CDN_CIDRS" { + continue + } + sources, err := st.ListCDNSources(tid, mod.ID) + if err != nil { + continue + } + for _, src := range sources { + u := strings.TrimSpace(src.URL) + if u == "" { + continue + } + req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil) + if err != nil { + continue + } + if strings.TrimSpace(src.Etag) != "" { + req.Header.Set("If-None-Match", strings.TrimSpace(src.Etag)) + } + resp, err := hc.Do(req) + if err != nil { + continue + } + _, _ = io.Copy(io.Discard, resp.Body) + _ = resp.Body.Close() + if resp.StatusCode != http.StatusOK { + continue + } + etag := strings.TrimSpace(resp.Header.Get("ETag")) + if etag == "" || etag == strings.TrimSpace(src.Etag) { + continue + } + e := etag + _, _ = st.UpdateCDNSource(tid, mod.ID, src.ID, &store.CDNSourcePatch{Etag: &e}) + } + } + } + return nil +} diff --git a/internal/pipeline/refresh.go b/internal/pipeline/refresh.go new file mode 100644 index 0000000..f8129ae --- /dev/null +++ b/internal/pipeline/refresh.go @@ -0,0 +1,214 @@ +package pipeline + +import ( + "context" + "crypto/sha256" + "fmt" + "io" + "net/http" + "net/netip" + "sort" + "strings" + + "evobgp/internal/birdfmt" + "evobgp/internal/store" + + "github.com/google/uuid" +) + +// RefreshModule runs ingest (where applicable) and creates a new rendered revision for the module. +func RefreshModule(ctx context.Context, st store.Backend, hc *http.Client, tenantID, moduleID string) (revisionID string, err error) { + if hc == nil { + hc = http.DefaultClient + } + mod, err := st.GetModule(tenantID, moduleID) + if err != nil { + return "", err + } + if !mod.Enabled { + return "", fmt.Errorf("module disabled") + } + + var rows []store.PrefixRow + switch mod.Type { + case "IP_RANGES": + list, err := st.ListIPRangeEntries(tenantID, moduleID) + if err != nil { + return "", err + } + for _, e := range list { + comm := e.CommunityID + if comm == nil && mod.DefaultCommunityID != nil { + c := *mod.DefaultCommunityID + comm = &c + } + rows = append(rows, store.PrefixRow{Prefix: e.Prefix, CommunityID: comm, Source: "ip_range"}) + } + case "AS_PREFIXES": + list, err := st.ListASEntries(tenantID, moduleID) + if err != nil { + return "", err + } + for _, e := range list { + if e.Prefix == nil || strings.TrimSpace(*e.Prefix) == "" { + continue + } + comm := e.CommunityID + if comm == nil && mod.DefaultCommunityID != nil { + c := *mod.DefaultCommunityID + comm = &c + } + rows = append(rows, store.PrefixRow{Prefix: strings.TrimSpace(*e.Prefix), CommunityID: comm, Source: "as_entry"}) + } + case "CDN_CIDRS": + sources, err := st.ListCDNSources(tenantID, moduleID) + if err != nil { + return "", err + } + for _, src := range sources { + u := strings.TrimSpace(src.URL) + if u == "" { + continue + } + req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil) + if err != nil { + return "", err + } + if strings.TrimSpace(src.Etag) != "" { + req.Header.Set("If-None-Match", strings.TrimSpace(src.Etag)) + } + resp, err := hc.Do(req) + if err != nil { + return "", fmt.Errorf("cdn fetch %s: %w", u, err) + } + if resp.StatusCode == http.StatusNotModified { + _ = resp.Body.Close() + continue + } + if resp.StatusCode != http.StatusOK { + _, _ = io.Copy(io.Discard, resp.Body) + _ = resp.Body.Close() + return "", fmt.Errorf("cdn url %s: %s", u, resp.Status) + } + body, err := io.ReadAll(io.LimitReader(resp.Body, 8<<20)) + _ = resp.Body.Close() + if err != nil { + return "", err + } + etag := strings.TrimSpace(resp.Header.Get("ETag")) + if etag != "" && etag != strings.TrimSpace(src.Etag) { + e := etag + _, _ = st.UpdateCDNSource(tenantID, moduleID, src.ID, &store.CDNSourcePatch{Etag: &e}) + } + for _, pfx := range ParseCIDRLines(string(body)) { + comm := src.CommunityID + if comm == nil && mod.DefaultCommunityID != nil { + c := *mod.DefaultCommunityID + comm = &c + } + rows = append(rows, store.PrefixRow{Prefix: pfx.String(), CommunityID: comm, Source: "cdn:" + src.ID}) + } + } + case "DOMAINS": + if _, err := st.ListDomainEntries(tenantID, moduleID); err != nil { + return "", err + } + // DNS/DoH resolution not wired yet; emit empty prefix set (valid revision). + default: + return "", fmt.Errorf("unknown module type %q", mod.Type) + } + + revisionID = uuid.NewString() + parent := parentRevision(st, tenantID, moduleID) + hash := hashMaterialization(moduleID, rows) + preview, err := buildPreviewFragments(revisionID, rows) + if err != nil { + return "", err + } + if err := st.CreateRenderRevision(revisionID, tenantID, moduleID, parent, hash, preview, rows); err != nil { + return "", err + } + return revisionID, nil +} + +func parentRevision(st store.Backend, tenantID, moduleID string) *string { + items, _, _ := st.ListRevisions(tenantID, moduleID, "", 1) + if len(items) == 0 { + return nil + } + id := items[0].ID + return &id +} + +func hashMaterialization(moduleID string, rows []store.PrefixRow) string { + type line struct{ p, c, s string } + var lines []line + for _, r := range rows { + c := "" + if r.CommunityID != nil { + c = *r.CommunityID + } + lines = append(lines, line{r.Prefix, c, r.Source}) + } + sort.Slice(lines, func(i, j int) bool { + if lines[i].p != lines[j].p { + return lines[i].p < lines[j].p + } + if lines[i].c != lines[j].c { + return lines[i].c < lines[j].c + } + return lines[i].s < lines[j].s + }) + h := sha256.New() + h.Write([]byte(moduleID)) + h.Write([]byte{0}) + for _, l := range lines { + h.Write([]byte(l.p)) + h.Write([]byte{1}) + h.Write([]byte(l.c)) + h.Write([]byte{1}) + h.Write([]byte(l.s)) + h.Write([]byte{0}) + } + return fmt.Sprintf("sha256:%x", h.Sum(nil)) +} + +func buildPreviewFragments(revisionID string, rows []store.PrefixRow) (map[string]string, error) { + var v4, v6 []netip.Prefix + for _, pr := range rows { + pfx, err := netip.ParsePrefix(strings.TrimSpace(pr.Prefix)) + if err != nil { + continue + } + if pfx.Addr().Is4() { + v4 = append(v4, pfx.Masked()) + } else { + v6 = append(v6, pfx.Masked()) + } + } + f4, err := birdfmt.RenderExportFilterIPv4("evobgp_export_v4", v4) + if err != nil { + return nil, err + } + f6, err := birdfmt.RenderExportFilterIPv6("evobgp_export_v6", v6) + if err != nil { + return nil, err + } + birdD := birdfmt.JoinFragments(birdfmt.ManagedBanner(revisionID), f4, f6) + main := `# EvoBGP generated (pipeline refresh) +router id 192.0.2.1; +include "bird.d/evobgp_generated.conf"; + +protocol device { +} + +protocol direct { + ipv4; + ipv6; +} +` + return map[string]string{ + "bird.conf": main, + "bird.d/evobgp_generated.conf": birdD, + }, nil +} diff --git a/internal/render/publish.go b/internal/render/publish.go new file mode 100644 index 0000000..2c11f6f --- /dev/null +++ b/internal/render/publish.go @@ -0,0 +1,31 @@ +package render + +import ( + "context" + "log" + + "evobgp/internal/store" +) + +// PublishLatestTenantRevision sets each speaker's published_revision_id to the newest tenant revision (by created_at). +// Simplification for single-tenant / demo stacks; multi-module tenants may prefer explicit publish policies later. +func PublishLatestTenantRevision(ctx context.Context, st store.Backend) { + _ = ctx + tenants, err := st.ListTenantIDs() + if err != nil { + log.Printf("evobgp-render: list tenants: %v", err) + return + } + for _, tid := range tenants { + revs, _, _ := st.ListRevisions(tid, "", "", 1) + if len(revs) == 0 { + continue + } + rid := revs[0].ID + for _, sp := range st.ListSpeakersForTenant(tid) { + if err := st.PublishRevisionForSpeaker(sp.ID, rid); err != nil { + log.Printf("evobgp-render: publish speaker %s: %v", sp.ID, err) + } + } + } +} diff --git a/internal/render/run.go b/internal/render/run.go index bebe30f..62b8953 100644 --- a/internal/render/run.go +++ b/internal/render/run.go @@ -1,22 +1,55 @@ -// Package render materializes prefix sets, creates config_revision rows, and builds BIRD text via internal/birdfmt. package render import ( "context" "log" + "os" + "strings" "time" "evobgp/internal/broker" "evobgp/internal/config" + "evobgp/internal/store" ) -// Run blocks until ctx is cancelled. Reference deployment: worker process after ingest completes. -func Run(ctx context.Context) { +// Deps enables auto-publish of the latest rendered revision to all speakers in each tenant. +type Deps struct { + Store store.Backend +} + +// Run blocks until ctx is cancelled. With Store set, periodically publishes the head revision for each tenant. +func Run(ctx context.Context, deps *Deps) { cfg := config.Load() broker.LogConnect(ctx, cfg.BrokerURL) + if deps == nil || deps.Store == nil { + runStub(ctx) + return + } + t := time.NewTicker(45 * time.Second) + defer t.Stop() + auto := strings.TrimSpace(os.Getenv("EVOBGP_RENDER_AUTOPUBLISH")) == "1" + if auto { + log.Printf("evobgp-render: active (EVOBGP_RENDER_AUTOPUBLISH=1: publish head revision per tenant)") + } else { + log.Printf("evobgp-render: active (idle publish; set EVOBGP_RENDER_AUTOPUBLISH=1 to auto-publish head revision)") + } + for { + select { + case <-ctx.Done(): + log.Printf("evobgp-render: stopped") + return + case <-t.C: + if auto { + PublishLatestTenantRevision(context.Background(), deps.Store) + } + } + } +} + +func runStub(ctx context.Context) { t := time.NewTicker(60 * time.Second) defer t.Stop() - log.Printf("evobgp-render: started (stub; render jobs produce revisions and artifacts)") + log.Printf("evobgp-render: idle stub (no store in Deps)") for { select { case <-ctx.Done(): diff --git a/internal/repository/postgres.go b/internal/repository/postgres.go index af3dad9..4795721 100644 --- a/internal/repository/postgres.go +++ b/internal/repository/postgres.go @@ -760,6 +760,82 @@ func (p *Postgres) LatestPublishedRevision(speakerID string) (string, time.Time, return rid, at, nil } +func (p *Postgres) ListTenantIDs() ([]string, error) { + ctx := context.Background() + rows, err := p.pool.Query(ctx, `SELECT id::text FROM tenant ORDER BY id`) + if err != nil { + return nil, err + } + defer rows.Close() + var out []string + for rows.Next() { + var id string + if err := rows.Scan(&id); err != nil { + continue + } + out = append(out, id) + } + return out, nil +} + +func (p *Postgres) CreateRenderRevision(revisionID, tenantID, moduleID string, parentRevisionID *string, contentHash string, previewFragments map[string]string, prefixes []store.PrefixRow) error { + if strings.TrimSpace(revisionID) == "" { + return store.ErrInvalidInput + } + if _, err := p.GetModule(tenantID, moduleID); err != nil { + return err + } + ctx := context.Background() + if previewFragments == nil { + previewFragments = map[string]string{} + } + meta, err := json.Marshal(map[string]any{ + "preview_fragments": previewFragments, + "materialized_prefix_count": len(prefixes), + }) + if err != nil { + return err + } + tx, err := p.pool.Begin(ctx) + if err != nil { + return err + } + defer func() { _ = tx.Rollback(ctx) }() + + var parent any + if parentRevisionID != nil && strings.TrimSpace(*parentRevisionID) != "" { + parent = strings.TrimSpace(*parentRevisionID) + } + _, err = tx.Exec(ctx, ` + INSERT INTO config_revision (id, tenant_id, module_id, content_hash, parent_revision_id, meta_json) + VALUES ($1::uuid, $2::uuid, $3::uuid, $4, $5::uuid, $6::jsonb)`, + strings.TrimSpace(revisionID), tenantID, moduleID, strings.TrimSpace(contentHash), parent, string(meta)) + if err != nil { + return err + } + for _, pr := range prefixes { + var comm any + if pr.CommunityID != nil && strings.TrimSpace(*pr.CommunityID) != "" { + comm = strings.TrimSpace(*pr.CommunityID) + } + src := pr.Source + if strings.TrimSpace(src) == "" { + src = "render" + } + _, err = tx.Exec(ctx, ` + INSERT INTO revision_materialized_prefix (revision_id, prefix, community_id, source) + VALUES ($1::uuid, $2::cidr, $3::uuid, $4)`, + strings.TrimSpace(revisionID), strings.TrimSpace(pr.Prefix), comm, src) + if err != nil { + return err + } + } + if err := tx.Commit(ctx); err != nil { + return err + } + return nil +} + func (p *Postgres) ListDohProfiles(tenantID string) ([]*store.DohProfile, error) { ctx := context.Background() rows, err := p.pool.Query(ctx, `SELECT id::text, name, url, timeout_ms, secret_ref FROM doh_profile WHERE tenant_id=$1`, tenantID) diff --git a/internal/scheduler/run.go b/internal/scheduler/run.go index 281ccbc..2f39cdf 100644 --- a/internal/scheduler/run.go +++ b/internal/scheduler/run.go @@ -1,30 +1,141 @@ -// Package scheduler drives module/CDN refresh intervals and enqueues work (broker or job_audit). -// This process entrypoint is a stub until PostgreSQL and NATS/Redis are wired. +// Package scheduler drives module refresh intervals and enqueues module_refresh jobs on the shared Registry. package scheduler import ( "context" + "fmt" + "io" "log" + "net/http" + "strings" "time" "evobgp/internal/broker" "evobgp/internal/config" + "evobgp/internal/jobs" + "evobgp/internal/store" ) -// Run blocks until ctx is cancelled. In production it connects to DB + broker and schedules refresh events. -func Run(ctx context.Context) { +// Deps wires the scheduler to the store. Use either Jobs (same process as API, e.g. evobgp-all) or +// APIBase+APIToken to call the control plane over HTTP (separate containers in reference Compose). +type Deps struct { + Store store.Backend + Jobs *jobs.Registry + APIBase string // e.g. http://evobgp-api:8080 + APIToken string // Bearer token (operator/editor) + HTTP *http.Client +} + +// Run blocks until ctx is cancelled. When deps is nil or incomplete, falls back to connect-only stub logging. +func Run(ctx context.Context, deps *Deps) { cfg := config.Load() broker.LogConnect(ctx, cfg.BrokerURL) - t := time.NewTicker(60 * time.Second) + if deps == nil || deps.Store == nil { + runStub(ctx) + return + } + if deps.Jobs == nil && (strings.TrimSpace(deps.APIBase) == "" || strings.TrimSpace(deps.APIToken) == "") { + runStub(ctx) + return + } + t := time.NewTicker(30 * time.Second) defer t.Stop() - log.Printf("evobgp-scheduler: started (stub; connect EVOBGP_DATABASE_URL / EVOBGP_BROKER_URL for full stack)") + if deps.Jobs != nil { + log.Printf("evobgp-scheduler: active (in-process enqueue module_refresh)") + } else { + log.Printf("evobgp-scheduler: active (HTTP POST .../modules/{id}/refresh → %s)", strings.TrimSpace(deps.APIBase)) + } + for { + select { + case <-ctx.Done(): + log.Printf("evobgp-scheduler: stopped") + return + case <-t.C: + tick(context.Background(), deps) + } + } +} + +func tick(ctx context.Context, deps *Deps) { + tenants, err := deps.Store.ListTenantIDs() + if err != nil { + log.Printf("evobgp-scheduler: list tenants: %v", err) + return + } + for _, tid := range tenants { + for _, mod := range deps.Store.ListModules(tid) { + if !mod.Enabled || mod.Type == "IP_RANGES" { + continue + } + interval := mod.RefreshIntervalSec + if interval <= 0 { + continue + } + win := interval + if win < 60 { + win = 60 + } + bucket := time.Now().Unix() / int64(win) + key := fmt.Sprintf("sched-%s-%d", mod.ID, bucket) + if deps.Jobs != nil { + mid := mod.ID + _, created, err := deps.Jobs.Enqueue(tid, jobs.KindModuleRefresh, &key, &mid, map[string]any{ + "module_id": mod.ID, + "trigger": "scheduler", + }) + if err != nil { + log.Printf("evobgp-scheduler: enqueue module %s: %v", mod.ID, err) + continue + } + if created { + log.Printf("evobgp-scheduler: queued refresh for module %s (%s)", mod.ID, mod.Type) + } + continue + } + if err := postModuleRefresh(ctx, deps, mod.ID, key); err != nil { + log.Printf("evobgp-scheduler: http refresh module %s: %v", mod.ID, err) + } + } + } +} + +func postModuleRefresh(ctx context.Context, deps *Deps, moduleID, idempotencyKey string) error { + base := strings.TrimRight(strings.TrimSpace(deps.APIBase), "/") + u := base + "/v1/modules/" + moduleID + "/refresh" + req, err := http.NewRequestWithContext(ctx, http.MethodPost, u, nil) + if err != nil { + return err + } + req.Header.Set("Authorization", "Bearer "+strings.TrimSpace(deps.APIToken)) + if idempotencyKey != "" { + req.Header.Set("Idempotency-Key", idempotencyKey) + } + hc := deps.HTTP + if hc == nil { + hc = http.DefaultClient + } + resp, err := hc.Do(req) + if err != nil { + return err + } + defer resp.Body.Close() + if resp.StatusCode == http.StatusNoContent || resp.StatusCode == http.StatusAccepted { + return nil + } + b, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) + return fmt.Errorf("%s: %s", resp.Status, strings.TrimSpace(string(b))) +} + +func runStub(ctx context.Context) { + t := time.NewTicker(60 * time.Second) + defer t.Stop() + log.Printf("evobgp-scheduler: idle stub (no store/jobs: pass scheduler.Deps from evobgp-all or bootstrap worker)") for { select { case <-ctx.Done(): log.Printf("evobgp-scheduler: stopped") return case <-t.C: - // Placeholder: evaluate module cron / refresh_interval and publish jobs. } } } diff --git a/internal/store/backend.go b/internal/store/backend.go index 1563505..cea9069 100644 --- a/internal/store/backend.go +++ b/internal/store/backend.go @@ -10,6 +10,9 @@ type Backend interface { // DemoIDs is non-empty only after SeedDemo (memory or seeded SQL). DemoIDs() (tenant, moduleCDN, moduleIP, revision, speaker string) + // ListTenantIDs returns distinct tenant identifiers (for background workers). + ListTenantIDs() ([]string, error) + ListModules(tenantID string) []*Module GetModule(tenantID, moduleID string) (*Module, error) CreateModule(tenantID string, in *Module) (*Module, error) @@ -64,6 +67,8 @@ type Backend interface { ListRevisions(tenantID, moduleID string, cursor string, limit int) (items []*Revision, nextCursor string, hasMore bool) ListRevisionPrefixes(tenantID, revisionID string, cursor string, limit int) (prefixes []PrefixRow, next string, more bool) CreateRollbackRevision(tenantID, sourceRevisionID string) (newID string, err error) + // CreateRenderRevision inserts a new config_revision (revID must be unique) with materialized prefixes and preview fragments. + CreateRenderRevision(revisionID, tenantID, moduleID string, parentRevisionID *string, contentHash string, previewFragments map[string]string, prefixes []PrefixRow) error RevisionDiff(tenantID, aID, bID string) (map[string]any, error) SetLastAppliedRevision(tenantID, speakerID, revisionID string) error PublishRevisionForSpeaker(speakerID, revisionID string) error diff --git a/internal/store/memory.go b/internal/store/memory.go index 1d1f66e..f2696c9 100644 --- a/internal/store/memory.go +++ b/internal/store/memory.go @@ -284,6 +284,62 @@ func (m *Memory) DemoIDs() (tenant, moduleCDN, moduleIP, revision, speaker strin return m.demoTenantID, m.demoModuleCDN, m.demoModuleIP, m.demoRevisionID, m.demoSpeakerID } +// ListTenantIDs returns tenant ids sorted lexicographically. +func (m *Memory) ListTenantIDs() ([]string, error) { + m.mu.RLock() + defer m.mu.RUnlock() + out := make([]string, 0, len(m.tenants)) + for id := range m.tenants { + out = append(out, id) + } + sort.Strings(out) + return out, nil +} + +// CreateRenderRevision stores a new revision and its materialized prefixes. +func (m *Memory) CreateRenderRevision(revisionID, tenantID, moduleID string, parentRevisionID *string, contentHash string, previewFragments map[string]string, prefixes []PrefixRow) error { + if strings.TrimSpace(revisionID) == "" || strings.TrimSpace(moduleID) == "" || strings.TrimSpace(contentHash) == "" { + return ErrInvalidInput + } + m.mu.Lock() + defer m.mu.Unlock() + if _, exists := m.revisions[revisionID]; exists { + return ErrInvalidInput + } + mod, ok := m.modules[moduleID] + if !ok || mod.DeletedAt != nil || mod.TenantID != tenantID { + return ErrNotFound + } + if parentRevisionID != nil && *parentRevisionID != "" { + if pr, ok := m.revisions[*parentRevisionID]; !ok || pr.TenantID != tenantID { + return ErrNotFound + } + } + frag := make(map[string]string, len(previewFragments)) + for k, v := range previewFragments { + frag[k] = v + } + parent := parentRevisionID + m.revisions[revisionID] = &Revision{ + ID: revisionID, + TenantID: tenantID, + ModuleID: moduleID, + ContentHash: contentHash, + ParentRevisionID: parent, + CreatedAt: time.Now().UTC(), + MaterializedPrefixCount: len(prefixes), + PreviewFragments: frag, + } + if len(prefixes) > 0 { + cp := make([]PrefixRow, len(prefixes)) + copy(cp, prefixes) + m.revPrefixes[revisionID] = cp + } else { + m.revPrefixes[revisionID] = nil + } + return nil +} + // ListModules returns modules for a tenant (sorted by priority, then name). func (m *Memory) ListModules(tenantID string) []*Module { m.mu.RLock()