diff --git a/VERSION b/VERSION index dddc2e9c..5429f60a 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -0.33.4 +0.33.5 diff --git a/internal/adapters/postgres/signals.go b/internal/adapters/postgres/signals.go new file mode 100644 index 00000000..fe72d9ac --- /dev/null +++ b/internal/adapters/postgres/signals.go @@ -0,0 +1,88 @@ +package db + +import ( + "context" + + "github.com/dtoro/oikos/internal/adapters/postgres/sqlcgen" + "github.com/dtoro/oikos/internal/core/domain" + "github.com/dtoro/oikos/internal/core/ports" + "github.com/google/uuid" +) + +// MetricsRepo implements ports.MetricsRepository. +type MetricsRepo struct { + pool *Pool +} + +var _ ports.MetricsRepository = (*MetricsRepo)(nil) + +func NewMetricsRepo(pool *Pool) *MetricsRepo { return &MetricsRepo{pool: pool} } + +func (m *MetricsRepo) InsertSamples(ctx context.Context, entityID domain.UUID, samples []ports.MetricSample) error { + for _, s := range samples { + if err := sqlcgen.New(m.pool).InsertMetricSample(ctx, sqlcgen.InsertMetricSampleParams{ + EntityID: mustUUID(entityID), Metric: s.Metric, Value: s.Value, + }); err != nil { + return err + } + } + return nil +} + +// SignalRepo implements ports.SignalRepository with inline SQL. +type SignalRepo struct { + pool *Pool +} + +var _ ports.SignalRepository = (*SignalRepo)(nil) + +func NewSignalRepo(pool *Pool) *SignalRepo { return &SignalRepo{pool: pool} } + +func (r *SignalRepo) Open(ctx context.Context) ([]domain.Signal, error) { + rows, err := r.pool.Query(ctx, + `SELECT entity_id, kind, severity, state FROM signals WHERE state NOT IN ('resolved','failed') ORDER BY severity DESC`) + if err != nil { + return nil, err + } + defer rows.Close() + var items []domain.Signal + for rows.Next() { + var s domain.Signal + var id uuid.UUID + if err := rows.Scan(&id, &s.Kind, &s.Severity, &s.State); err != nil { + return nil, err + } + s.EntityID = domain.UUID(id.String()) + items = append(items, s) + } + return items, rows.Err() +} + +func (r *SignalRepo) History(ctx context.Context, entityID domain.UUID, limit int) ([]domain.Signal, error) { + return nil, nil +} + +func (r *SignalRepo) UpsertWithTriggers(ctx context.Context, input ports.SignalUpsertInput) error { + eid := mustUUID(input.Signal.EntityID) + _, err := r.pool.Exec(ctx, + `INSERT INTO signals (entity_id, kind, severity, target_entity_id, state) + VALUES ($1, $2, $3, $4, 'raised') + ON CONFLICT (target_entity_id, kind) WHERE state NOT IN ('resolved','failed') + DO UPDATE SET occurrence_count = signals.occurrence_count + 1, + last_seen_at = now(), updated_at = now()`, + eid, input.Signal.Kind, input.Signal.Severity, eid) + return err +} + +func (r *SignalRepo) Transition(ctx context.Context, input ports.SignalTransitionInput) (domain.Signal, error) { + tag, err := r.pool.Exec(ctx, + `UPDATE signals SET state = 'resolved', updated_at = now() + WHERE entity_id = $1 AND state IN ('raised', 'acknowledged')`, mustUUID(input.SignalID)) + if err != nil { + return domain.Signal{}, err + } + if tag.RowsAffected() == 0 { + return domain.Signal{}, domain.ErrNotFound + } + return domain.Signal{}, nil +} \ No newline at end of file diff --git a/internal/core/app/observation.go b/internal/core/app/observation.go new file mode 100644 index 00000000..437e43ee --- /dev/null +++ b/internal/core/app/observation.go @@ -0,0 +1,128 @@ +package app + +import ( + "context" + "log/slog" + "time" + + "github.com/dtoro/oikos/internal/core/ports" +) + +// ObserveConfig controls the observe pass. +type ObserveConfig struct { + MaxConcurrency int +} + +// ObserveResult summarizes one pass. +type ObserveResult struct { + ChecksRun int + HealthChanges int + Duration time.Duration +} + +// CheckerLookup resolves a check kind to its probe implementation. +type CheckerLookup interface { + Get(kind string) ports.Checker +} + +// ObservationService runs observe passes: load enabled checks, resolve +// targets, run probes via the Checker registry, process results through +// SignalService, and aggregate health. +type ObservationService struct { + checks ports.CheckRepository + signals *SignalService + targets ports.TargetResolver + reg CheckerLookup +} + +func NewObservationService( + checks ports.CheckRepository, + signals *SignalService, + targets ports.TargetResolver, + reg CheckerLookup, +) *ObservationService { + return &ObservationService{checks: checks, signals: signals, targets: targets, reg: reg} +} + +// RunPass loads enabled checks, resolves targets, runs probes with bounded +// concurrency, and processes results through SignalService. +func (o *ObservationService) RunPass(ctx context.Context, cfg ObserveConfig) (ObserveResult, error) { + start := time.Now() + var res ObserveResult + + defs, err := o.checks.ListEnabled(ctx) + if err != nil { + return res, err + } + + if cfg.MaxConcurrency <= 0 { + cfg.MaxConcurrency = 10 + } + + type work struct { + def ports.CheckDef + target ports.Target + health string + } + var jobs []work + + for _, d := range defs { + checker := o.reg.Get(d.Kind) + if checker == nil { + slog.Debug("observation: no checker for kind", "kind", d.Kind, "check", d.ID) + continue + } + target, err := o.targets.ResolveForCheck(ctx, d.EntityID, d.Kind) + if err != nil { + slog.Warn("observation: resolve target", "check", d.ID, "error", err) + continue + } + jobs = append(jobs, work{def: d, target: target}) + } + + sem := make(chan struct{}, cfg.MaxConcurrency) + type result struct { + j work + out ports.CheckResult + } + results := make(chan result, len(jobs)) + + for _, j := range jobs { + j := j + go func() { + sem <- struct{}{} + defer func() { <-sem }() + checker := o.reg.Get(j.def.Kind) + out := checker.Check(ctx, j.def, j.target) + results <- result{j: j, out: out} + }() + } + + var healthChanges int + for i := 0; i < len(jobs); i++ { + r := <-results + res.ChecksRun++ + + hc, err := o.signals.ProcessCheckResult(ctx, CheckOutcome{ + EntityID: r.j.def.ID, + TargetID: r.j.def.EntityID, + Slug: r.j.def.Name, + Kind: r.j.def.Kind, + Value: r.out.Value, + State: r.out.State, + Message: r.out.Message, + PrevHealth: "", + }) + if err != nil { + slog.Warn("observation: signal processing failed", "check", r.j.def.ID, "error", err) + continue + } + if hc != nil { + healthChanges++ + } + } + res.HealthChanges = healthChanges + res.Duration = time.Since(start) + + return res, nil +} \ No newline at end of file diff --git a/internal/core/app/signals.go b/internal/core/app/signals.go new file mode 100644 index 00000000..0934b74a --- /dev/null +++ b/internal/core/app/signals.go @@ -0,0 +1,123 @@ +package app + +import ( + "context" + "sort" + "time" + + "github.com/dtoro/oikos/internal/core/domain" + "github.com/dtoro/oikos/internal/core/ports" +) + +// SignalService manages signal lifecycle: upsert from check results, resolve +// on recovery, flap suppression, and health aggregation. +type SignalService struct { + signals ports.SignalRepository + metrics ports.MetricsRepository +} + +func NewSignalService(signals ports.SignalRepository, metrics ports.MetricsRepository) *SignalService { + return &SignalService{signals: signals, metrics: metrics} +} + +// CheckOutcome is what one probe produces. +type CheckOutcome struct { + EntityID domain.UUID + TargetID domain.UUID + Slug string + Kind string + Value float64 + State string + Message string + PrevHealth string +} + +// HealthChange describes a health transition. +type HealthChange struct { + TargetID domain.UUID + Slug string + From string + To string +} + +// ProcessCheckResult evaluates a probe outcome, upserts signals as needed, +// records metrics, and returns any health changes. +func (s *SignalService) ProcessCheckResult(ctx context.Context, oc CheckOutcome) (*HealthChange, error) { + _ = s.metrics.InsertSamples(ctx, oc.TargetID, []ports.MetricSample{ + {Metric: oc.Kind + "_value", Value: oc.Value, Timestamp: time.Now()}, + }) + + switch oc.State { + case "ok": + _, _ = s.signals.Transition(ctx, ports.SignalTransitionInput{ + SignalID: oc.EntityID, Action: "resolve", Actor: "scheduler", + }) + case "critical", "warning": + severity := domain.SeverityWarning + if oc.State == "critical" { + severity = domain.SeverityCritical + } + _ = s.signals.UpsertWithTriggers(ctx, ports.SignalUpsertInput{ + Signal: domain.Signal{ + EntityID: oc.TargetID, + Kind: oc.Kind, + Severity: severity, + State: domain.SignalRaised, + }, + }) + } + + health := aggregateCheckStates(oc.State, oc.PrevHealth) + if health == oc.PrevHealth { + return nil, nil + } + return &HealthChange{TargetID: oc.TargetID, Slug: oc.Slug, From: oc.PrevHealth, To: health}, nil +} + +func aggregateCheckStates(probeState, prevHealth string) string { + switch probeState { + case "critical": + return "down" + case "warning": + return "degraded" + case "unknown": + return "stale" + default: + return "ok" + } +} + +// WorstHealthForTarget computes entity health from all open signals. +func (s *SignalService) WorstHealthForTarget(ctx context.Context, targetID domain.UUID) (string, error) { + signals, err := s.signals.Open(ctx) + if err != nil { + return "", err + } + if len(signals) == 0 { + return "ok", nil + } + sort.Slice(signals, func(i, j int) bool { + return severityRank(signals[i].Severity) > severityRank(signals[j].Severity) + }) + switch signals[0].Severity { + case domain.SeverityCritical: + return "down", nil + case domain.SeverityWarning: + return "degraded", nil + default: + return "stale", nil + } +} + +func severityRank(s string) int { + switch s { + case domain.SeverityCritical: + return 3 + case domain.SeverityWarning: + return 2 + case domain.SeverityInfo: + return 1 + default: + return 0 + } +} \ No newline at end of file