CI / changes (push) Successful in 4s
CI / openapi (push) Successful in 24s
CI / go (push) Successful in 22s
CI / bird2 (push) Successful in 13s
CI / docker-images (deploy/docker/bird2/Dockerfile, evobgp-bird2) (push) Successful in 38s
CI / docker-images (deploy/docker/evobgp-agent/Dockerfile, evobgp-agent) (push) Successful in 1m8s
CI / docker-images (deploy/docker/evobgp-web/Dockerfile, evobgp-web) (push) Successful in 49s
CI / docker-images (evobgp-all, 1, deploy/docker/gobinary/Dockerfile, evobgp-all) (push) Successful in 1m30s
CI / docker-images (evobgp-api, 1, deploy/docker/gobinary/Dockerfile, evobgp-api) (push) Successful in 1m31s
CI / docker-images (evobgp-ingest, 0, deploy/docker/gobinary/Dockerfile, evobgp-ingest) (push) Has been cancelled
CI / docker-images (evobgp-deploy, 0, deploy/docker/gobinary/Dockerfile, evobgp-deploy) (push) Has been cancelled
CI / docker-images (evobgp-node, 0, deploy/docker/gobinary/Dockerfile, evobgp-node) (push) Has been cancelled
CI / docker-images (evobgp-render, 0, deploy/docker/gobinary/Dockerfile, evobgp-render) (push) Has been cancelled
CI / docker-images (evobgp-scheduler, 0, deploy/docker/gobinary/Dockerfile, evobgp-scheduler) (push) Has been cancelled
126 lines
3.7 KiB
Go
126 lines
3.7 KiB
Go
// 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"
|
|
)
|
|
|
|
// 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. Misconfigured deps terminate the process (no idle fallback).
|
|
func Run(ctx context.Context, deps *Deps) {
|
|
cfg := config.Load()
|
|
broker.LogConnect(ctx, cfg.BrokerURL)
|
|
if deps == nil || deps.Store == nil {
|
|
log.Fatalf("evobgp-scheduler: missing store (pass scheduler.Deps with Store from BootstrapWorkers or evobgp-all)")
|
|
}
|
|
if deps.Jobs == nil && (strings.TrimSpace(deps.APIBase) == "" || strings.TrimSpace(deps.APIToken) == "") {
|
|
log.Fatalf("evobgp-scheduler: need either in-process Jobs (evobgp-all) or both EVOBGP_CONTROL_PLANE_URL and EVOBGP_SCHEDULER_BEARER")
|
|
}
|
|
t := time.NewTicker(30 * time.Second)
|
|
defer t.Stop()
|
|
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)))
|
|
}
|