From 3416ccf55ca1f521e2db0d8f972e0b1274bedc22 Mon Sep 17 00:00:00 2001 From: Jeefos Date: Mon, 20 Apr 2026 16:41:52 -0400 Subject: [PATCH] Monitoring --- cmd/kleffd/main.go | 10 + internal/adapters/out/platform/client.go | 56 +++++- .../repository/memory/server_repository.go | 10 + .../adapters/out/runtime/docker/docker.go | 82 +++++++- internal/app/config/config.go | 5 +- internal/application/ports/log_shipper.go | 23 +++ .../application/ports/server_repository.go | 2 + internal/application/ports/workload.go | 4 + .../ports/workload_status_reporter.go | 6 + internal/logs/tailer.go | 187 ++++++++++++++++++ internal/metrics/scraper.go | 148 ++++++++++++++ internal/workers/provision_worker.go | 1 + 12 files changed, 522 insertions(+), 12 deletions(-) create mode 100644 internal/application/ports/log_shipper.go create mode 100644 internal/logs/tailer.go create mode 100644 internal/metrics/scraper.go diff --git a/cmd/kleffd/main.go b/cmd/kleffd/main.go index 05edf81..8c69151 100644 --- a/cmd/kleffd/main.go +++ b/cmd/kleffd/main.go @@ -7,6 +7,7 @@ import ( "os" "os/signal" "syscall" + "time" "github.com/kleffio/kleff-daemon/internal/adapters/out/db" "github.com/kleffio/kleff-daemon/internal/adapters/out/observability/logging" @@ -17,6 +18,8 @@ import ( k8sadapter "github.com/kleffio/kleff-daemon/internal/adapters/out/runtime/kubernetes" "github.com/kleffio/kleff-daemon/internal/app/config" "github.com/kleffio/kleff-daemon/internal/application/ports" + "github.com/kleffio/kleff-daemon/internal/logs" + "github.com/kleffio/kleff-daemon/internal/metrics" "github.com/kleffio/kleff-daemon/internal/workers" "github.com/kleffio/kleff-daemon/internal/workers/jobs" "k8s.io/client-go/rest" @@ -86,6 +89,13 @@ func main() { ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() + scrapeInterval := time.Duration(cfg.MetricsScrapeInterval) * time.Second + scraper := metrics.NewScraper(runtime, repo, platformClient, scrapeInterval, cfg.NodeID, daemonLog) + go scraper.Run(ctx) + + tailer := logs.NewTailer(runtime, repo, platformClient, daemonLog) + go tailer.Run(ctx) + dispatcher.Run(ctx) daemonLog.Info("Daemon shutdown complete") } diff --git a/internal/adapters/out/platform/client.go b/internal/adapters/out/platform/client.go index b8ae77c..4ca3318 100644 --- a/internal/adapters/out/platform/client.go +++ b/internal/adapters/out/platform/client.go @@ -92,17 +92,61 @@ func (c *Client) RegisterNode(ctx context.Context) error { return nil } +func (c *Client) ShipLogs(ctx context.Context, workloadID, projectID string, lines []ports.LogEntry) error { + if len(lines) == 0 { + return nil + } + type lineDTO struct { + Ts string `json:"ts"` + Stream string `json:"stream"` + Line string `json:"line"` + } + dtos := make([]lineDTO, len(lines)) + for i, l := range lines { + dtos[i] = lineDTO{Ts: l.Ts.UTC().Format(time.RFC3339Nano), Stream: l.Stream, Line: l.Line} + } + payload := map[string]any{"project_id": projectID, "lines": dtos} + body, err := json.Marshal(payload) + if err != nil { + return fmt.Errorf("marshal log payload: %w", err) + } + url := fmt.Sprintf("%s/api/v1/internal/workloads/%s/log-lines", c.baseURL, workloadID) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(body)) + if err != nil { + return fmt.Errorf("build log ship request: %w", err) + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bearer "+c.nodeToken) + + resp, err := c.httpClient.Do(req) + if err != nil { + return fmt.Errorf("log ship request failed: %w", err) + } + defer resp.Body.Close() + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + raw, _ := io.ReadAll(io.LimitReader(resp.Body, 512)) + return fmt.Errorf("log ship failed: status=%d body=%s", resp.StatusCode, strings.TrimSpace(string(raw))) + } + return nil +} + func (c *Client) ReportStatus(ctx context.Context, update ports.WorkloadStatusUpdate) error { if c.nodeToken == "" { return fmt.Errorf("node token is not set; call RegisterNode first") } payload := map[string]any{ - "status": update.Status, - "runtime_ref": update.RuntimeRef, - "endpoint": update.Endpoint, - "node_id": update.NodeID, - "error_message": update.ErrorMessage, - "observed_at": time.Now().UTC().Format(time.RFC3339), + "status": update.Status, + "runtime_ref": update.RuntimeRef, + "endpoint": update.Endpoint, + "node_id": update.NodeID, + "error_message": update.ErrorMessage, + "observed_at": time.Now().UTC().Format(time.RFC3339), + "cpu_millicores": update.CPUMillicores, + "memory_mb": update.MemoryMB, + "network_rx_mb": update.NetworkRxMB, + "network_tx_mb": update.NetworkTxMB, + "disk_read_mb": update.DiskReadMB, + "disk_write_mb": update.DiskWriteMB, } body, err := json.Marshal(payload) if err != nil { diff --git a/internal/adapters/out/repository/memory/server_repository.go b/internal/adapters/out/repository/memory/server_repository.go index 00c3b89..3e8e34b 100644 --- a/internal/adapters/out/repository/memory/server_repository.go +++ b/internal/adapters/out/repository/memory/server_repository.go @@ -46,3 +46,13 @@ func (r *ServerRepository) UpdateStatus(ctx context.Context, id string, status s s.Status = status return nil } + +func (r *ServerRepository) ListAll(ctx context.Context) ([]*ports.ServerRecord, error) { + r.mu.Lock() + defer r.mu.Unlock() + out := make([]*ports.ServerRecord, 0, len(r.servers)) + for _, s := range r.servers { + out = append(out, s) + } + return out, nil +} diff --git a/internal/adapters/out/runtime/docker/docker.go b/internal/adapters/out/runtime/docker/docker.go index 43487aa..78c71eb 100644 --- a/internal/adapters/out/runtime/docker/docker.go +++ b/internal/adapters/out/runtime/docker/docker.go @@ -2,9 +2,11 @@ package docker import ( "context" + "encoding/json" "errors" "fmt" "io" + "sort" "strings" "github.com/docker/docker/api/types/container" @@ -280,7 +282,7 @@ func (a *Adapter) Remove(ctx context.Context, projectID, workloadID string) erro return nil } -// Status returns the current state of the container. +// Status returns the current state and resource metrics of the container. func (a *Adapter) Status(ctx context.Context, projectID, workloadID string) (*ports.WorkloadHealth, error) { containerID, err := a.findContainer(ctx, projectID, workloadID) if err != nil { @@ -290,8 +292,71 @@ func (a *Adapter) Status(ctx context.Context, projectID, workloadID string) (*po if err != nil { return nil, fmt.Errorf("failed to inspect container: %w", err) } - state := strings.ToLower(info.State.Status) - return &ports.WorkloadHealth{WorkloadID: workloadID, State: state}, nil + health := &ports.WorkloadHealth{ + WorkloadID: workloadID, + State: strings.ToLower(info.State.Status), + } + if info.State.Running { + if err := a.collectStats(ctx, containerID, health); err != nil { + // Non-fatal: state is already populated; metrics will be zero. + _ = err + } + } + return health, nil +} + +func (a *Adapter) collectStats(ctx context.Context, containerID string, h *ports.WorkloadHealth) error { + resp, err := a.client.ContainerStats(ctx, containerID, false) + if err != nil { + return fmt.Errorf("container stats: %w", err) + } + defer resp.Body.Close() + + var stats container.StatsResponse + if err := json.NewDecoder(resp.Body).Decode(&stats); err != nil { + return fmt.Errorf("decode stats: %w", err) + } + + // CPU: delta-based percentage converted to millicores. + cpuDelta := stats.CPUStats.CPUUsage.TotalUsage - stats.PreCPUStats.CPUUsage.TotalUsage + sysDelta := stats.CPUStats.SystemUsage - stats.PreCPUStats.SystemUsage + numCPUs := uint64(stats.CPUStats.OnlineCPUs) + if numCPUs == 0 { + numCPUs = uint64(len(stats.CPUStats.CPUUsage.PercpuUsage)) + } + if numCPUs == 0 { + numCPUs = 1 + } + if sysDelta > 0 && cpuDelta > 0 { + h.CPUMillicores = int64((float64(cpuDelta) / float64(sysDelta)) * float64(numCPUs) * 1000) + } + + // Memory. + h.MemoryMB = int64(stats.MemoryStats.Usage / (1024 * 1024)) + + // Network: sum all interfaces. + var rxBytes, txBytes uint64 + for _, iface := range stats.Networks { + rxBytes += iface.RxBytes + txBytes += iface.TxBytes + } + h.NetworkRxMB = float64(rxBytes) / (1024 * 1024) + h.NetworkTxMB = float64(txBytes) / (1024 * 1024) + + // Disk I/O. + var diskRead, diskWrite uint64 + for _, entry := range stats.BlkioStats.IoServiceBytesRecursive { + switch strings.ToLower(entry.Op) { + case "read": + diskRead += entry.Value + case "write": + diskWrite += entry.Value + } + } + h.DiskReadMB = float64(diskRead) / (1024 * 1024) + h.DiskWriteMB = float64(diskWrite) / (1024 * 1024) + + return nil } // Endpoint returns the first exposed host port. @@ -304,8 +369,15 @@ func (a *Adapter) Endpoint(ctx context.Context, projectID, workloadID string) (s if err != nil { return "", fmt.Errorf("failed to inspect container: %w", err) } - for _, bindings := range info.NetworkSettings.Ports { - if len(bindings) > 0 { + // Sort port keys so we always pick the lowest container port deterministically. + keys := make([]string, 0, len(info.NetworkSettings.Ports)) + for k := range info.NetworkSettings.Ports { + keys = append(keys, string(k)) + } + sort.Strings(keys) + for _, k := range keys { + bindings := info.NetworkSettings.Ports[nat.Port(k)] + if len(bindings) > 0 && bindings[0].HostPort != "" { return fmt.Sprintf("127.0.0.1:%s", bindings[0].HostPort), nil } } diff --git a/internal/app/config/config.go b/internal/app/config/config.go index 915ef7f..12635b6 100644 --- a/internal/app/config/config.go +++ b/internal/app/config/config.go @@ -25,7 +25,8 @@ type Config struct { ClusterRegion string `mapstructure:"cluster.region"` NodeID string `mapstructure:"node.id"` GRPCPort int `mapstructure:"grpc.port"` - MetricsPort int `mapstructure:"metrics.port"` + MetricsPort int `mapstructure:"metrics.port"` + MetricsScrapeInterval int `mapstructure:"metrics.scrape_interval"` QueueBackend QueueBackend `mapstructure:"queue.backend"` DatabasePath string `mapstructure:"database.path"` RedisURL string `mapstructure:"redis.url"` @@ -72,6 +73,7 @@ func Load() (*Config, error) { v.SetDefault("node.id", hostname) v.SetDefault("grpc.port", 50051) v.SetDefault("metrics.port", 9090) + v.SetDefault("metrics.scrape_interval", 30) v.SetDefault("queue.backend", string(QueueBackendMemory)) v.SetDefault("database.path", "./data/kleff.db") v.SetDefault("redis.url", "redis://localhost:6379/0") @@ -116,6 +118,7 @@ func Load() (*Config, error) { cfg.NodeID = v.GetString("node.id") cfg.GRPCPort = v.GetInt("grpc.port") cfg.MetricsPort = v.GetInt("metrics.port") + cfg.MetricsScrapeInterval = v.GetInt("metrics.scrape_interval") cfg.QueueBackend = QueueBackend(v.GetString("queue.backend")) cfg.DatabasePath = v.GetString("database.path") cfg.RedisURL = v.GetString("redis.url") diff --git a/internal/application/ports/log_shipper.go b/internal/application/ports/log_shipper.go new file mode 100644 index 0000000..347e9f5 --- /dev/null +++ b/internal/application/ports/log_shipper.go @@ -0,0 +1,23 @@ +package ports + +import ( + "context" + "time" +) + +// LogEntry is one line of container output. +type LogEntry struct { + Ts time.Time + Stream string // "stdout" or "stderr" + Line string +} + +// LogShipper ships batches of log lines to the platform. +type LogShipper interface { + ShipLogs(ctx context.Context, workloadID, projectID string, lines []LogEntry) error +} + +// NoopLogShipper discards all log lines. Used when log shipping is disabled. +type NoopLogShipper struct{} + +func (NoopLogShipper) ShipLogs(_ context.Context, _, _ string, _ []LogEntry) error { return nil } diff --git a/internal/application/ports/server_repository.go b/internal/application/ports/server_repository.go index 2b1865e..f0126fd 100644 --- a/internal/application/ports/server_repository.go +++ b/internal/application/ports/server_repository.go @@ -11,10 +11,12 @@ type ServerRecord struct { NodeID string Runtime string RuntimeRef string + ProjectID string } type ServerRepository interface { Save(ctx context.Context, server *ServerRecord) error FindByID(ctx context.Context, id string) (*ServerRecord, error) UpdateStatus(ctx context.Context, id string, status string) error + ListAll(ctx context.Context) ([]*ServerRecord, error) } diff --git a/internal/application/ports/workload.go b/internal/application/ports/workload.go index a475347..83d4f3a 100644 --- a/internal/application/ports/workload.go +++ b/internal/application/ports/workload.go @@ -67,6 +67,10 @@ type WorkloadHealth struct { State string `json:"state"` // running, stopped, failed CPUMillicores int64 `json:"cpu_millicores"` MemoryMB int64 `json:"memory_mb"` + NetworkRxMB float64 `json:"network_rx_mb"` + NetworkTxMB float64 `json:"network_tx_mb"` + DiskReadMB float64 `json:"disk_read_mb"` + DiskWriteMB float64 `json:"disk_write_mb"` // Game server extension (zero-valued for non-game workloads) ActivePlayers int `json:"active_players,omitempty"` // HTTP service extension diff --git a/internal/application/ports/workload_status_reporter.go b/internal/application/ports/workload_status_reporter.go index 876915e..036443f 100644 --- a/internal/application/ports/workload_status_reporter.go +++ b/internal/application/ports/workload_status_reporter.go @@ -10,6 +10,12 @@ type WorkloadStatusUpdate struct { Endpoint string NodeID string ErrorMessage string + CPUMillicores int64 + MemoryMB int64 + NetworkRxMB float64 + NetworkTxMB float64 + DiskReadMB float64 + DiskWriteMB float64 } type WorkloadStatusReporter interface { diff --git a/internal/logs/tailer.go b/internal/logs/tailer.go new file mode 100644 index 0000000..f8dae68 --- /dev/null +++ b/internal/logs/tailer.go @@ -0,0 +1,187 @@ +package logs + +import ( + "bufio" + "context" + "io" + "sync" + "time" + + "github.com/docker/docker/pkg/stdcopy" + "github.com/kleffio/kleff-daemon/internal/application/ports" +) + +const ( + batchSize = 100 + flushInterval = 5 * time.Second +) + +// Tailer manages one log-tailing goroutine per running workload. +// It mirrors the Scraper pattern: call Run to start, cancel the context to stop. +type Tailer struct { + runtime ports.RuntimeAdapter + repo ports.ServerRepository + shipper ports.LogShipper + logger ports.Logger + + mu sync.Mutex + running map[string]context.CancelFunc // workloadID → cancel +} + +func NewTailer( + runtime ports.RuntimeAdapter, + repo ports.ServerRepository, + shipper ports.LogShipper, + logger ports.Logger, +) *Tailer { + return &Tailer{ + runtime: runtime, + repo: repo, + shipper: shipper, + logger: logger, + running: make(map[string]context.CancelFunc), + } +} + +// Run blocks until ctx is cancelled, reconciling the set of active tailers +// every 10 seconds to match the set of running workloads. +func (t *Tailer) Run(ctx context.Context) { + ticker := time.NewTicker(10 * time.Second) + defer ticker.Stop() + // Reconcile immediately on start. + t.reconcile(ctx) + for { + select { + case <-ctx.Done(): + t.stopAll() + return + case <-ticker.C: + t.reconcile(ctx) + } + } +} + +func (t *Tailer) reconcile(ctx context.Context) { + servers, err := t.repo.ListAll(ctx) + if err != nil { + t.logger.Warn("log tailer: failed to list servers", "error", err) + return + } + + active := make(map[string]struct{}, len(servers)) + for _, srv := range servers { + if srv.Status != "Running" && srv.Status != "running" { + continue + } + active[srv.ID] = struct{}{} + + t.mu.Lock() + _, already := t.running[srv.ID] + t.mu.Unlock() + + if !already { + wctx, cancel := context.WithCancel(ctx) + t.mu.Lock() + t.running[srv.ID] = cancel + t.mu.Unlock() + go t.tail(wctx, srv.ID, srv.ProjectID) + } + } + + // Stop tailers for workloads that are no longer running. + t.mu.Lock() + for id, cancel := range t.running { + if _, ok := active[id]; !ok { + cancel() + delete(t.running, id) + } + } + t.mu.Unlock() +} + +func (t *Tailer) stopAll() { + t.mu.Lock() + defer t.mu.Unlock() + for id, cancel := range t.running { + cancel() + delete(t.running, id) + } +} + +// tail streams logs for one workload until ctx is cancelled or the container exits. +func (t *Tailer) tail(ctx context.Context, workloadID, projectID string) { + defer func() { + t.mu.Lock() + delete(t.running, workloadID) + t.mu.Unlock() + }() + + rc, err := t.runtime.Logs(ctx, projectID, workloadID, true) + if err != nil { + t.logger.Warn("log tailer: failed to open log stream", "workload_id", workloadID, "error", err) + return + } + defer rc.Close() + + // Docker multiplexes stdout/stderr in an 8-byte framed format. + // stdcopy.StdCopy demuxes them into separate writers. + stdoutR, stdoutW := io.Pipe() + stderrR, stderrW := io.Pipe() + + go func() { + _, _ = stdcopy.StdCopy(stdoutW, stderrW, rc) + stdoutW.Close() + stderrW.Close() + }() + + batch := make([]ports.LogEntry, 0, batchSize) + flush := time.NewTicker(flushInterval) + defer flush.Stop() + + lines := make(chan ports.LogEntry, 256) + + go scanStream(ctx, stdoutR, "stdout", lines) + go scanStream(ctx, stderrR, "stderr", lines) + + for { + select { + case <-ctx.Done(): + if len(batch) > 0 { + _ = t.shipper.ShipLogs(context.Background(), workloadID, projectID, batch) + } + return + case entry, ok := <-lines: + if !ok { + if len(batch) > 0 { + _ = t.shipper.ShipLogs(context.Background(), workloadID, projectID, batch) + } + return + } + batch = append(batch, entry) + if len(batch) >= batchSize { + if err := t.shipper.ShipLogs(ctx, workloadID, projectID, batch); err != nil { + t.logger.Warn("log tailer: ship failed", "workload_id", workloadID, "error", err) + } + batch = batch[:0] + } + case <-flush.C: + if len(batch) > 0 { + if err := t.shipper.ShipLogs(ctx, workloadID, projectID, batch); err != nil { + t.logger.Warn("log tailer: ship failed", "workload_id", workloadID, "error", err) + } + batch = batch[:0] + } + } + } +} + +func scanStream(ctx context.Context, r io.Reader, stream string, out chan<- ports.LogEntry) { + scanner := bufio.NewScanner(r) + for scanner.Scan() { + select { + case <-ctx.Done(): + return + case out <- ports.LogEntry{Ts: time.Now().UTC(), Stream: stream, Line: scanner.Text()}: + } + } +} diff --git a/internal/metrics/scraper.go b/internal/metrics/scraper.go new file mode 100644 index 0000000..c3b5449 --- /dev/null +++ b/internal/metrics/scraper.go @@ -0,0 +1,148 @@ +package metrics + +import ( + "context" + "sync" + "time" + + "github.com/kleffio/kleff-daemon/internal/application/ports" +) + +type ioSnapshot struct { + rxMB float64 + txMB float64 + diskRMB float64 + diskWMB float64 +} + +// Scraper periodically collects per-container metrics and reports them to the platform. +type Scraper struct { + runtime ports.RuntimeAdapter + repo ports.ServerRepository + reporter ports.WorkloadStatusReporter + interval time.Duration + nodeID string + logger ports.Logger + + mu sync.Mutex + prevIO map[string]ioSnapshot // workloadID → last cumulative network+disk totals +} + +func NewScraper( + runtime ports.RuntimeAdapter, + repo ports.ServerRepository, + reporter ports.WorkloadStatusReporter, + interval time.Duration, + nodeID string, + logger ports.Logger, +) *Scraper { + return &Scraper{ + runtime: runtime, + repo: repo, + reporter: reporter, + interval: interval, + nodeID: nodeID, + logger: logger, + prevIO: make(map[string]ioSnapshot), + } +} + +// Run blocks until ctx is cancelled, scraping metrics on each tick. +func (s *Scraper) Run(ctx context.Context) { + ticker := time.NewTicker(s.interval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + s.scrape(ctx) + } + } +} + +func (s *Scraper) scrape(ctx context.Context) { + servers, err := s.repo.ListAll(ctx) + if err != nil { + s.logger.Warn("metrics scraper: failed to list servers", "error", err) + return + } + + activeIDs := make(map[string]struct{}, len(servers)) + for _, srv := range servers { + if srv.Status != "Running" && srv.Status != "running" { + continue + } + activeIDs[srv.ID] = struct{}{} + + health, err := s.runtime.Status(ctx, srv.ProjectID, srv.ID) + if err != nil { + s.logger.Warn("metrics scraper: failed to get status", "workload_id", srv.ID, "error", err) + continue + } + + // Convert cumulative network+disk totals to per-interval deltas. + s.mu.Lock() + prev := s.prevIO[srv.ID] + rxDelta := health.NetworkRxMB - prev.rxMB + txDelta := health.NetworkTxMB - prev.txMB + diskRDelta := health.DiskReadMB - prev.diskRMB + diskWDelta := health.DiskWriteMB - prev.diskWMB + if rxDelta < 0 { + rxDelta = health.NetworkRxMB // container restarted + } + if txDelta < 0 { + txDelta = health.NetworkTxMB + } + if diskRDelta < 0 { + diskRDelta = health.DiskReadMB + } + if diskWDelta < 0 { + diskWDelta = health.DiskWriteMB + } + s.prevIO[srv.ID] = ioSnapshot{rxMB: health.NetworkRxMB, txMB: health.NetworkTxMB, diskRMB: health.DiskReadMB, diskWMB: health.DiskWriteMB} + s.mu.Unlock() + + update := ports.WorkloadStatusUpdate{ + WorkloadID: srv.ID, + ProjectID: srv.ProjectID, + Status: mapDockerState(health.State), + RuntimeRef: srv.RuntimeRef, + NodeID: s.nodeID, + CPUMillicores: health.CPUMillicores, + MemoryMB: health.MemoryMB, + NetworkRxMB: rxDelta, + NetworkTxMB: txDelta, + DiskReadMB: diskRDelta, + DiskWriteMB: diskWDelta, + } + if err := s.reporter.ReportStatus(ctx, update); err != nil { + s.logger.Warn("metrics scraper: failed to report status", "workload_id", srv.ID, "error", err) + } + } + + // Clean up snapshots for workloads that are no longer active. + s.mu.Lock() + for id := range s.prevIO { + if _, ok := activeIDs[id]; !ok { + delete(s.prevIO, id) + } + } + s.mu.Unlock() +} + +// mapDockerState converts Docker container states to platform workload states. +func mapDockerState(dockerState string) string { + switch dockerState { + case "running": + return "running" + case "exited", "dead": + return "stopped" + case "created", "paused", "restarting": + return "pending" + case "removing": + return "deleted" + default: + return "failed" + } +} diff --git a/internal/workers/provision_worker.go b/internal/workers/provision_worker.go index c5bc363..7ddcd10 100644 --- a/internal/workers/provision_worker.go +++ b/internal/workers/provision_worker.go @@ -64,6 +64,7 @@ func (w *ProvisionWorker) Handle(ctx context.Context, job *jobs.Job) error { Status: server.State, NodeID: server.Labels.NodeID, RuntimeRef: server.RuntimeRef, + ProjectID: spec.ProjectID, } if err := w.repository.Save(ctx, record); err != nil {