From 5f9199f2646502905348781023c64eb33fb54561 Mon Sep 17 00:00:00 2001 From: Jeefos Date: Mon, 20 Apr 2026 16:42:50 -0400 Subject: [PATCH] Monitoring --- go.work.sum | 1 + internal/bootstrap/container.go | 10 +- internal/bootstrap/http.go | 2 + internal/core/logs/adapters/http/handler.go | 113 ++++++++++++++++++ .../core/logs/adapters/persistence/store.go | 75 ++++++++++++ internal/core/logs/domain/log.go | 13 ++ internal/core/logs/ports/repository.go | 20 ++++ internal/core/usage/adapters/http/handler.go | 45 +++++-- .../core/usage/adapters/persistence/store.go | 101 ++++++++++++++++ internal/core/usage/domain/usage.go | 39 +++++- internal/core/usage/ports/repository.go | 12 ++ .../core/workloads/adapters/http/handler.go | 87 +++++++++++--- .../workloads/adapters/persistence/store.go | 14 ++- .../commands/provision_workload.go | 2 + internal/core/workloads/domain/workload.go | 22 ++-- .../database/migrations/009_usage_records.sql | 16 +++ .../migrations/010_usage_display_columns.sql | 11 ++ .../migrations/011_workload_resources.sql | 3 + .../migrations/012_workload_log_lines.sql | 12 ++ 19 files changed, 548 insertions(+), 50 deletions(-) create mode 100644 internal/core/logs/adapters/http/handler.go create mode 100644 internal/core/logs/adapters/persistence/store.go create mode 100644 internal/core/logs/domain/log.go create mode 100644 internal/core/logs/ports/repository.go create mode 100644 internal/core/usage/adapters/persistence/store.go create mode 100644 internal/core/usage/ports/repository.go create mode 100644 internal/database/migrations/009_usage_records.sql create mode 100644 internal/database/migrations/010_usage_display_columns.sql create mode 100644 internal/database/migrations/011_workload_resources.sql create mode 100644 internal/database/migrations/012_workload_log_lines.sql diff --git a/go.work.sum b/go.work.sum index 64f9e80..bc29da1 100644 --- a/go.work.sum +++ b/go.work.sum @@ -11,6 +11,7 @@ github.com/envoyproxy/go-control-plane v0.12.0/go.mod h1:ZBTaoJ23lqITozF0M6G4/Ir github.com/envoyproxy/protoc-gen-validate v1.0.4/go.mod h1:qys6tmnRsYrQqIhm2bvKZH4Blx/1gTIZ2UKVY1M+Yew= github.com/golang/glog v1.2.0/go.mod h1:6AhwSGph0fcJtXVM/PEHPqZlFeoLxhs7/t5UDAwmO+w= github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= +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/klauspost/cpuid/v2 v2.0.9/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg= github.com/kleffio/plugin-sdk v0.0.0-20260327000334-ef46a02d9d74 h1:55pnh083fiuFK82jDb7JHR1IEcQLvd+yQJuSucXUC9s= diff --git a/internal/bootstrap/container.go b/internal/bootstrap/container.go index 5547c92..82e5a90 100644 --- a/internal/bootstrap/container.go +++ b/internal/bootstrap/container.go @@ -33,7 +33,10 @@ import ( pluginapplication "github.com/kleffio/platform/internal/core/plugins/application" projectshttp "github.com/kleffio/platform/internal/core/projects/adapters/http" projectspersistence "github.com/kleffio/platform/internal/core/projects/adapters/persistence" + logshttp "github.com/kleffio/platform/internal/core/logs/adapters/http" + logspersistence "github.com/kleffio/platform/internal/core/logs/adapters/persistence" usagehttp "github.com/kleffio/platform/internal/core/usage/adapters/http" + usagepersistence "github.com/kleffio/platform/internal/core/usage/adapters/persistence" workloadshttp "github.com/kleffio/platform/internal/core/workloads/adapters/http" workloadspersistence "github.com/kleffio/platform/internal/core/workloads/adapters/persistence" workloadcmd "github.com/kleffio/platform/internal/core/workloads/application/commands" @@ -68,6 +71,7 @@ type Container struct { NodesHandler *nodeshttp.Handler BillingHandler *billinghttp.Handler UsageHandler *usagehttp.Handler + LogsHandler *logshttp.Handler AuditHandler *audithttp.Handler AdminHandler *adminhttp.Handler PluginsHandler *pluginhttp.Handler @@ -159,10 +163,12 @@ func NewContainer(cfg *Config, logger *slog.Logger) (*Container, error) { OrganizationsHandler: organizationshttp.NewHandler(logger), DeploymentsHandler: deploymentshttp.NewHandler(createDeployment, serverAction, deploymentStore, cfg.SecretKey, logger), ProjectsHandler: projectshttp.NewHandler(projectsStore, logger), - WorkloadsHandler: workloadshttp.NewHandler(projectsStore, workloadsStore, provisionHandler, workloadAction, bus, logger), + WorkloadsHandler: workloadshttp.NewHandler(projectsStore, workloadsStore, provisionHandler, workloadAction, bus, logger). + WithUsageRepository(usagepersistence.NewPostgresUsageStore(db)), NodesHandler: nodeshttp.NewHandler(nodeStore, logger), BillingHandler: billinghttp.NewHandler(logger), - UsageHandler: usagehttp.NewHandler(logger), + UsageHandler: usagehttp.NewHandler(usagepersistence.NewPostgresUsageStore(db), logger), + LogsHandler: logshttp.NewHandler(logspersistence.NewPostgresLogStore(db), logger), AuditHandler: audithttp.NewHandler(logger), AdminHandler: adminhttp.NewHandler(logger), PluginsHandler: pluginhttp.NewHandler(pluginMgr, catalogRegistry, logger), diff --git a/internal/bootstrap/http.go b/internal/bootstrap/http.go index 0c88b7c..817ca64 100644 --- a/internal/bootstrap/http.go +++ b/internal/bootstrap/http.go @@ -52,6 +52,7 @@ func buildRouter(c *Container) http.Handler { r.Group(func(r chi.Router) { r.Use(middleware.RequireNodeAuth(c.NodeVerifier)) c.WorkloadsHandler.RegisterInternalRoutes(r) + c.LogsHandler.RegisterInternalRoutes(r) }) // Authenticated routes @@ -68,6 +69,7 @@ func buildRouter(c *Container) http.Handler { c.NodesHandler.RegisterRoutes(r) c.BillingHandler.RegisterRoutes(r) c.UsageHandler.RegisterRoutes(r) + c.LogsHandler.RegisterRoutes(r) c.AuditHandler.RegisterRoutes(r) // Admin routes — additionally require the "admin" role. diff --git a/internal/core/logs/adapters/http/handler.go b/internal/core/logs/adapters/http/handler.go new file mode 100644 index 0000000..1bcc287 --- /dev/null +++ b/internal/core/logs/adapters/http/handler.go @@ -0,0 +1,113 @@ +package http + +import ( + "encoding/json" + "log/slog" + "net/http" + "strconv" + "time" + + "github.com/go-chi/chi/v5" + "github.com/kleffio/platform/internal/core/logs/domain" + "github.com/kleffio/platform/internal/core/logs/ports" +) + +type Handler struct { + repo ports.LogRepository + logger *slog.Logger +} + +func NewHandler(repo ports.LogRepository, logger *slog.Logger) *Handler { + return &Handler{repo: repo, logger: logger} +} + +// RegisterInternalRoutes wires the daemon-facing ingest endpoint. +// Requires node token auth (called within the RequireNodeAuth middleware group). +func (h *Handler) RegisterInternalRoutes(r chi.Router) { + r.Post("/api/v1/internal/workloads/{workloadID}/log-lines", h.ingestLines) +} + +// RegisterRoutes wires the panel-facing read endpoint. +// Requires user JWT auth. +func (h *Handler) RegisterRoutes(r chi.Router) { + r.Get("/api/v1/projects/{projectID}/workloads/{workloadID}/logs", h.getLogs) +} + +// ingestLines accepts a batch of log lines from the daemon. +func (h *Handler) ingestLines(w http.ResponseWriter, r *http.Request) { + workloadID := chi.URLParam(r, "workloadID") + + var body struct { + ProjectID string `json:"project_id"` + Lines []struct { + Ts string `json:"ts"` + Stream string `json:"stream"` + Line string `json:"line"` + } `json:"lines"` + } + if err := json.NewDecoder(r.Body).Decode(&body); err != nil { + writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid body"}) + return + } + + lines := make([]*domain.LogLine, 0, len(body.Lines)) + for _, l := range body.Lines { + ts, err := parseRFC3339(l.Ts) + if err != nil { + continue + } + stream := l.Stream + if stream == "" { + stream = "stdout" + } + lines = append(lines, &domain.LogLine{ + WorkloadID: workloadID, + ProjectID: body.ProjectID, + Ts: ts, + Stream: stream, + Line: l.Line, + }) + } + + if err := h.repo.SaveBatch(r.Context(), lines); err != nil { + h.logger.Error("save log batch", "error", err, "workload_id", workloadID) + writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "failed to save logs"}) + return + } + w.WriteHeader(http.StatusNoContent) +} + +// getLogs returns recent log lines for a workload. +func (h *Handler) getLogs(w http.ResponseWriter, r *http.Request) { + projectID := chi.URLParam(r, "projectID") + workloadID := chi.URLParam(r, "workloadID") + _ = projectID // used for future access-control checks + + limit := 200 + if s := r.URL.Query().Get("limit"); s != "" { + if n, err := strconv.Atoi(s); err == nil && n > 0 && n <= 2000 { + limit = n + } + } + + lines, err := h.repo.ListByWorkload(r.Context(), workloadID, limit) + if err != nil { + h.logger.Error("list log lines", "error", err, "workload_id", workloadID) + writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "failed to fetch logs"}) + return + } + if lines == nil { + lines = []*domain.LogLine{} + } + writeJSON(w, http.StatusOK, map[string]any{"lines": lines}) +} + +func parseRFC3339(s string) (time.Time, error) { + return time.Parse(time.RFC3339Nano, s) +} + +func writeJSON(w http.ResponseWriter, status int, body any) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(body) +} diff --git a/internal/core/logs/adapters/persistence/store.go b/internal/core/logs/adapters/persistence/store.go new file mode 100644 index 0000000..05438e4 --- /dev/null +++ b/internal/core/logs/adapters/persistence/store.go @@ -0,0 +1,75 @@ +package persistence + +import ( + "context" + "database/sql" + "fmt" + + "github.com/kleffio/platform/internal/core/logs/domain" +) + +type PostgresLogStore struct { + db *sql.DB +} + +func NewPostgresLogStore(db *sql.DB) *PostgresLogStore { + return &PostgresLogStore{db: db} +} + +func (s *PostgresLogStore) SaveBatch(ctx context.Context, lines []*domain.LogLine) error { + if len(lines) == 0 { + return nil + } + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("begin tx: %w", err) + } + defer tx.Rollback() //nolint:errcheck + + stmt, err := tx.PrepareContext(ctx, ` + INSERT INTO workload_log_lines (workload_id, project_id, ts, stream, line) + VALUES ($1, $2, $3, $4, $5) + `) + if err != nil { + return fmt.Errorf("prepare log insert: %w", err) + } + defer stmt.Close() + + for _, l := range lines { + if _, err := stmt.ExecContext(ctx, l.WorkloadID, l.ProjectID, l.Ts, l.Stream, l.Line); err != nil { + return fmt.Errorf("insert log line: %w", err) + } + } + return tx.Commit() +} + +func (s *PostgresLogStore) ListByWorkload(ctx context.Context, workloadID string, limit int) ([]*domain.LogLine, error) { + if limit <= 0 { + limit = 200 + } + rows, err := s.db.QueryContext(ctx, ` + SELECT id, workload_id, project_id, ts, stream, line + FROM workload_log_lines + WHERE workload_id = $1 + ORDER BY ts DESC + LIMIT $2 + `, workloadID, limit) + if err != nil { + return nil, fmt.Errorf("list log lines: %w", err) + } + defer rows.Close() + + var results []*domain.LogLine + for rows.Next() { + l := &domain.LogLine{} + if err := rows.Scan(&l.ID, &l.WorkloadID, &l.ProjectID, &l.Ts, &l.Stream, &l.Line); err != nil { + return nil, fmt.Errorf("scan log line: %w", err) + } + results = append(results, l) + } + // Reverse so results are in chronological order (oldest first). + for i, j := 0, len(results)-1; i < j; i, j = i+1, j-1 { + results[i], results[j] = results[j], results[i] + } + return results, rows.Err() +} diff --git a/internal/core/logs/domain/log.go b/internal/core/logs/domain/log.go new file mode 100644 index 0000000..572dbc3 --- /dev/null +++ b/internal/core/logs/domain/log.go @@ -0,0 +1,13 @@ +package domain + +import "time" + +// LogLine is one line of output from a running workload container. +type LogLine struct { + ID int64 `json:"id"` + WorkloadID string `json:"workload_id"` + ProjectID string `json:"project_id"` + Ts time.Time `json:"ts"` + Stream string `json:"stream"` + Line string `json:"line"` +} diff --git a/internal/core/logs/ports/repository.go b/internal/core/logs/ports/repository.go new file mode 100644 index 0000000..d5c9473 --- /dev/null +++ b/internal/core/logs/ports/repository.go @@ -0,0 +1,20 @@ +package ports + +import ( + "context" + + "github.com/kleffio/platform/internal/core/logs/domain" +) + +// LogRepository persists and retrieves workload log lines. +// Implementations of this interface are the extension point for log-sink +// plugins (Loki, Elasticsearch, etc.) — a plugin can wrap or replace the +// default Postgres store by implementing this interface. +type LogRepository interface { + // SaveBatch persists a batch of log lines. Implementations should be + // idempotent or at least tolerant of duplicate lines on retry. + SaveBatch(ctx context.Context, lines []*domain.LogLine) error + + // ListByWorkload returns up to limit lines for workloadID ordered by ts DESC. + ListByWorkload(ctx context.Context, workloadID string, limit int) ([]*domain.LogLine, error) +} diff --git a/internal/core/usage/adapters/http/handler.go b/internal/core/usage/adapters/http/handler.go index 65c6008..52814ca 100644 --- a/internal/core/usage/adapters/http/handler.go +++ b/internal/core/usage/adapters/http/handler.go @@ -1,34 +1,55 @@ package http import ( + "encoding/json" "log/slog" "net/http" "github.com/go-chi/chi/v5" + usagedomain "github.com/kleffio/platform/internal/core/usage/domain" + usageports "github.com/kleffio/platform/internal/core/usage/ports" ) const basePath = "/api/v1/usage" -// Handler groups all HTTP endpoints for the usage module. type Handler struct { + repo usageports.UsageRepository logger *slog.Logger } -func NewHandler(logger *slog.Logger) *Handler { - return &Handler{logger: logger} +func NewHandler(repo usageports.UsageRepository, logger *slog.Logger) *Handler { + return &Handler{repo: repo, logger: logger} } -// RegisterRoutes attaches all usage routes to the provided router. func (h *Handler) RegisterRoutes(r chi.Router) { - r.Get(basePath+"/summary", h.getSummary) - r.Get(basePath+"/records", h.listRecords) + r.Get(basePath+"/metrics", h.getMetrics) } -func notImplemented(w http.ResponseWriter) { - w.Header().Set("Content-Type", "application/json") - w.WriteHeader(http.StatusNotImplemented) - _, _ = w.Write([]byte(`{"error":"not implemented"}`)) +// getMetrics returns the latest per-workload metrics snapshot for a project. +// Query param: project_id (required) +func (h *Handler) getMetrics(w http.ResponseWriter, r *http.Request) { + projectID := r.URL.Query().Get("project_id") + if projectID == "" { + writeJSON(w, http.StatusBadRequest, map[string]string{"error": "project_id is required"}) + return + } + + metrics, err := h.repo.ListLatestByProject(r.Context(), projectID) + if err != nil { + h.logger.Error("list metrics by project", "error", err, "project_id", projectID) + writeJSON(w, http.StatusInternalServerError, map[string]string{"error": "failed to fetch metrics"}) + return + } + + if metrics == nil { + metrics = []*usagedomain.WorkloadMetrics{} + } + + writeJSON(w, http.StatusOK, map[string]any{"workloads": metrics}) } -func (h *Handler) getSummary(w http.ResponseWriter, _ *http.Request) { notImplemented(w) } -func (h *Handler) listRecords(w http.ResponseWriter, _ *http.Request) { notImplemented(w) } +func writeJSON(w http.ResponseWriter, status int, body any) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(body) +} diff --git a/internal/core/usage/adapters/persistence/store.go b/internal/core/usage/adapters/persistence/store.go new file mode 100644 index 0000000..d80cfdf --- /dev/null +++ b/internal/core/usage/adapters/persistence/store.go @@ -0,0 +1,101 @@ +package persistence + +import ( + "context" + "database/sql" + "fmt" + + "github.com/kleffio/platform/internal/core/usage/domain" +) + +type PostgresUsageStore struct { + db *sql.DB +} + +func NewPostgresUsageStore(db *sql.DB) *PostgresUsageStore { + return &PostgresUsageStore{db: db} +} + +func (s *PostgresUsageStore) Save(ctx context.Context, r *domain.UsageRecord) error { + _, err := s.db.ExecContext(ctx, ` + INSERT INTO usage_records ( + id, organization_id, project_id, workload_id, node_id, recorded_at, + cpu_seconds, memory_gb_hours, network_in_mb, network_out_mb, + disk_read_mb, disk_write_mb, + cpu_millicores, memory_mb, + network_in_kbps, network_out_kbps, + disk_read_kbps, disk_write_kbps + ) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,$18) + `, + r.ID, + r.OrganizationID, + r.ProjectID, + r.GameServerID, + r.NodeID, + r.RecordedAt, + r.CPUSeconds, + r.MemoryGBHours, + r.NetworkInMB, + r.NetworkOutMB, + r.DiskReadMB, + r.DiskWriteMB, + r.CPUMillicores, + r.MemoryMB, + r.NetworkInKbps, + r.NetworkOutKbps, + r.DiskReadKbps, + r.DiskWriteKbps, + ) + if err != nil { + return fmt.Errorf("save usage record: %w", err) + } + return nil +} + +// ListLatestByProject returns the most recent metrics snapshot per workload for a given project, +// joined with the workload's allocated CPU/memory limits. +func (s *PostgresUsageStore) ListLatestByProject(ctx context.Context, projectID string) ([]*domain.WorkloadMetrics, error) { + rows, err := s.db.QueryContext(ctx, ` + SELECT + u.workload_id, u.project_id, + u.cpu_millicores, u.memory_mb, + u.network_in_kbps, u.network_out_kbps, + u.disk_read_kbps, u.disk_write_kbps, + u.recorded_at, + COALESCE(w.cpu_millicores, 0), + COALESCE(w.memory_bytes, 0) + FROM ( + SELECT DISTINCT ON (workload_id) + workload_id, project_id, + cpu_millicores, memory_mb, + network_in_kbps, network_out_kbps, + disk_read_kbps, disk_write_kbps, + recorded_at + FROM usage_records + WHERE project_id = $1 + ORDER BY workload_id, recorded_at DESC + ) u + JOIN workloads w ON w.id = u.workload_id AND w.state != 'deleted' + `, projectID) + if err != nil { + return nil, fmt.Errorf("list latest usage by project: %w", err) + } + defer rows.Close() + + var results []*domain.WorkloadMetrics + for rows.Next() { + m := &domain.WorkloadMetrics{} + if err := rows.Scan( + &m.WorkloadID, &m.ProjectID, + &m.CPUMillicores, &m.MemoryMB, + &m.NetworkInKbps, &m.NetworkOutKbps, + &m.DiskReadKbps, &m.DiskWriteKbps, + &m.RecordedAt, + &m.CPULimitMillicores, &m.MemoryLimitBytes, + ); err != nil { + return nil, fmt.Errorf("scan usage row: %w", err) + } + results = append(results, m) + } + return results, rows.Err() +} diff --git a/internal/core/usage/domain/usage.go b/internal/core/usage/domain/usage.go index 135f56f..e37f6ec 100644 --- a/internal/core/usage/domain/usage.go +++ b/internal/core/usage/domain/usage.go @@ -17,14 +17,41 @@ type UsageSummary struct { type UsageRecord struct { ID string OrganizationID string + ProjectID string GameServerID string NodeID string RecordedAt time.Time - CPUSeconds float64 - MemoryGBHours float64 - NetworkInMB float64 - NetworkOutMB float64 - DiskReadMB float64 - DiskWriteMB float64 + // Billing units + CPUSeconds float64 + MemoryGBHours float64 + NetworkInMB float64 + NetworkOutMB float64 + DiskReadMB float64 + DiskWriteMB float64 + + // Display units (human-readable monitoring) + CPUMillicores int64 + MemoryMB int64 + NetworkInKbps float64 + NetworkOutKbps float64 + DiskReadKbps float64 + DiskWriteKbps float64 +} + +// WorkloadMetrics is the latest snapshot for a single workload, used by the monitoring page. +type WorkloadMetrics struct { + WorkloadID string `json:"workload_id"` + ProjectID string `json:"project_id"` + CPUMillicores int64 `json:"cpu_millicores"` + MemoryMB int64 `json:"memory_mb"` + NetworkInKbps float64 `json:"network_in_kbps"` + NetworkOutKbps float64 `json:"network_out_kbps"` + DiskReadKbps float64 `json:"disk_read_kbps"` + DiskWriteKbps float64 `json:"disk_write_kbps"` + RecordedAt time.Time `json:"recorded_at"` + + // Allocation limits from the workload provisioning request. + CPULimitMillicores int64 `json:"cpu_limit_millicores"` + MemoryLimitBytes int64 `json:"memory_limit_bytes"` } diff --git a/internal/core/usage/ports/repository.go b/internal/core/usage/ports/repository.go new file mode 100644 index 0000000..880bab6 --- /dev/null +++ b/internal/core/usage/ports/repository.go @@ -0,0 +1,12 @@ +package ports + +import ( + "context" + + "github.com/kleffio/platform/internal/core/usage/domain" +) + +type UsageRepository interface { + Save(ctx context.Context, record *domain.UsageRecord) error + ListLatestByProject(ctx context.Context, projectID string) ([]*domain.WorkloadMetrics, error) +} diff --git a/internal/core/workloads/adapters/http/handler.go b/internal/core/workloads/adapters/http/handler.go index 0c567ea..1c9437d 100644 --- a/internal/core/workloads/adapters/http/handler.go +++ b/internal/core/workloads/adapters/http/handler.go @@ -14,6 +14,9 @@ import ( "github.com/go-chi/chi/v5" projectports "github.com/kleffio/platform/internal/core/projects/ports" + "github.com/kleffio/platform/internal/shared/ids" + usagedomain "github.com/kleffio/platform/internal/core/usage/domain" + usageports "github.com/kleffio/platform/internal/core/usage/ports" "github.com/kleffio/platform/internal/core/workloads/application/commands" "github.com/kleffio/platform/internal/core/workloads/domain" "github.com/kleffio/platform/internal/core/workloads/ports" @@ -29,12 +32,13 @@ const ( ) type Handler struct { - projects projectports.ProjectRepository - repo ports.Repository - provision *commands.ProvisionWorkloadHandler - action *commands.WorkloadActionHandler - bus *events.Bus - logger *slog.Logger + projects projectports.ProjectRepository + repo ports.Repository + usageRepo usageports.UsageRepository + provision *commands.ProvisionWorkloadHandler + action *commands.WorkloadActionHandler + bus *events.Bus + logger *slog.Logger } var orgSlugCleaner = regexp.MustCompile(`[^a-z0-9-]+`) @@ -43,6 +47,11 @@ func NewHandler(projects projectports.ProjectRepository, repo ports.Repository, return &Handler{projects: projects, repo: repo, provision: provision, action: action, bus: bus, logger: logger} } +func (h *Handler) WithUsageRepository(ur usageports.UsageRepository) *Handler { + h.usageRepo = ur + return h +} + func (h *Handler) RegisterRoutes(r chi.Router) { r.Get(projectBasePath, h.list) r.Post(projectBasePath, h.provisionWorkload) @@ -170,12 +179,18 @@ func (h *Handler) get(w http.ResponseWriter, r *http.Request) { func (h *Handler) updateStatus(w http.ResponseWriter, r *http.Request) { workloadID := chi.URLParam(r, "id") var req struct { - Status string `json:"status"` - RuntimeRef string `json:"runtime_ref"` - Endpoint string `json:"endpoint"` - NodeID string `json:"node_id"` - ErrorMessage string `json:"error_message"` - ObservedAt string `json:"observed_at"` + Status string `json:"status"` + RuntimeRef string `json:"runtime_ref"` + Endpoint string `json:"endpoint"` + NodeID string `json:"node_id"` + ErrorMessage string `json:"error_message"` + ObservedAt string `json:"observed_at"` + 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"` } if err := json.NewDecoder(r.Body).Decode(&req); err != nil { writeJSON(w, http.StatusBadRequest, map[string]string{"error": "invalid json body"}) @@ -218,13 +233,19 @@ func (h *Handler) updateStatus(w http.ResponseWriter, r *http.Request) { } update := domain.DaemonStatusUpdate{ - WorkloadID: workloadID, - Status: status, - RuntimeRef: req.RuntimeRef, - Endpoint: req.Endpoint, - NodeID: nodeID, - ErrorMessage: req.ErrorMessage, - ObservedAt: observedAt, + WorkloadID: workloadID, + Status: status, + RuntimeRef: req.RuntimeRef, + Endpoint: req.Endpoint, + NodeID: nodeID, + ErrorMessage: req.ErrorMessage, + ObservedAt: observedAt, + CPUMillicores: req.CPUMillicores, + MemoryMB: req.MemoryMB, + NetworkRxMB: req.NetworkRxMB, + NetworkTxMB: req.NetworkTxMB, + DiskReadMB: req.DiskReadMB, + DiskWriteMB: req.DiskWriteMB, } if err := h.repo.UpdateFromDaemon(r.Context(), update); err != nil { if errors.Is(err, sql.ErrNoRows) { @@ -236,6 +257,34 @@ func (h *Handler) updateStatus(w http.ResponseWriter, r *http.Request) { return } + if h.usageRepo != nil && (req.CPUMillicores > 0 || req.MemoryMB > 0) { + const scrapeIntervalSeconds = 30.0 + usageRecord := &usagedomain.UsageRecord{ + ID: ids.New(), + OrganizationID: existing.OrganizationID, + ProjectID: existing.ProjectID, + GameServerID: workloadID, + NodeID: nodeID, + RecordedAt: observedAt, + CPUSeconds: float64(req.CPUMillicores) / 1000.0 * scrapeIntervalSeconds, + MemoryGBHours: float64(req.MemoryMB) / 1024.0 * (scrapeIntervalSeconds / 3600.0), + NetworkInMB: req.NetworkRxMB, + NetworkOutMB: req.NetworkTxMB, + DiskReadMB: req.DiskReadMB, + DiskWriteMB: req.DiskWriteMB, + // Display units + CPUMillicores: req.CPUMillicores, + MemoryMB: req.MemoryMB, + NetworkInKbps: req.NetworkRxMB * 1024.0 / scrapeIntervalSeconds, + NetworkOutKbps: req.NetworkTxMB * 1024.0 / scrapeIntervalSeconds, + DiskReadKbps: req.DiskReadMB * 1024.0 / scrapeIntervalSeconds, + DiskWriteKbps: req.DiskWriteMB * 1024.0 / scrapeIntervalSeconds, + } + if err := h.usageRepo.Save(r.Context(), usageRecord); err != nil { + h.logger.Warn("save usage record", "error", err, "workload_id", workloadID) + } + } + if h.bus != nil { _ = h.bus.Publish(r.Context(), domain.WorkloadStatusChanged{ WorkloadID: workloadID, diff --git a/internal/core/workloads/adapters/persistence/store.go b/internal/core/workloads/adapters/persistence/store.go index 330aab7..b9f9d8b 100644 --- a/internal/core/workloads/adapters/persistence/store.go +++ b/internal/core/workloads/adapters/persistence/store.go @@ -26,11 +26,13 @@ func (s *PostgresStore) CreateWorkload(ctx context.Context, workload *domain.Wor INSERT INTO workloads ( id, name, organization_id, project_id, owner_id, blueprint_id, image, runtime_ref, endpoint, node_id, state, error_message, + cpu_millicores, memory_bytes, created_at, updated_at ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, - $13, $14 + $13, $14, + $15, $16 )`, workload.ID, workload.Name, @@ -44,6 +46,8 @@ func (s *PostgresStore) CreateWorkload(ctx context.Context, workload *domain.Wor nullIfEmpty(workload.NodeID), workload.State, workload.ErrorMessage, + workload.CPUMillicores, + workload.MemoryBytes, workload.CreatedAt, workload.UpdatedAt, ) @@ -57,7 +61,7 @@ func (s *PostgresStore) FindByProjectAndName(ctx context.Context, projectID, nam row := s.db.QueryRowContext(ctx, ` SELECT id, name, organization_id, project_id, owner_id, blueprint_id, image, runtime_ref, endpoint, COALESCE(node_id, ''), state, - error_message, created_at, updated_at + error_message, cpu_millicores, memory_bytes, created_at, updated_at FROM workloads WHERE project_id = $1 AND name = $2 ORDER BY updated_at DESC @@ -69,7 +73,7 @@ func (s *PostgresStore) FindByID(ctx context.Context, workloadID string) (*domai row := s.db.QueryRowContext(ctx, ` SELECT id, name, organization_id, project_id, owner_id, blueprint_id, image, runtime_ref, endpoint, COALESCE(node_id, ''), state, - error_message, created_at, updated_at + error_message, cpu_millicores, memory_bytes, created_at, updated_at FROM workloads WHERE id = $1`, workloadID) return scanWorkload(row) } @@ -78,7 +82,7 @@ func (s *PostgresStore) ListByProject(ctx context.Context, projectID string) ([] rows, err := s.db.QueryContext(ctx, ` SELECT id, name, organization_id, project_id, owner_id, blueprint_id, image, runtime_ref, endpoint, COALESCE(node_id, ''), state, - error_message, created_at, updated_at + error_message, cpu_millicores, memory_bytes, created_at, updated_at FROM workloads WHERE project_id = $1 ORDER BY created_at DESC`, projectID) @@ -262,6 +266,8 @@ func scanWorkload(s scanner) (*domain.Workload, error) { &w.NodeID, &w.State, &w.ErrorMessage, + &w.CPUMillicores, + &w.MemoryBytes, &w.CreatedAt, &w.UpdatedAt, ); err != nil { diff --git a/internal/core/workloads/application/commands/provision_workload.go b/internal/core/workloads/application/commands/provision_workload.go index aaeb677..2a7d815 100644 --- a/internal/core/workloads/application/commands/provision_workload.go +++ b/internal/core/workloads/application/commands/provision_workload.go @@ -167,6 +167,8 @@ func (h *ProvisionWorkloadHandler) Handle(ctx context.Context, cmd ProvisionWork BlueprintID: cmd.BlueprintID, Image: image, State: domain.WorkloadPending, + CPUMillicores: cpuMillicores, + MemoryBytes: memoryBytes, CreatedAt: now, UpdatedAt: now, } diff --git a/internal/core/workloads/domain/workload.go b/internal/core/workloads/domain/workload.go index dc7391a..410db44 100644 --- a/internal/core/workloads/domain/workload.go +++ b/internal/core/workloads/domain/workload.go @@ -25,18 +25,26 @@ type Workload struct { NodeID string `json:"node_id"` State WorkloadState `json:"state"` ErrorMessage string `json:"error_message"` + CPUMillicores int64 `json:"cpu_millicores"` + MemoryBytes int64 `json:"memory_bytes"` CreatedAt time.Time `json:"created_at"` UpdatedAt time.Time `json:"updated_at"` } type DaemonStatusUpdate struct { - WorkloadID string - Status WorkloadState - RuntimeRef string - Endpoint string - NodeID string - ErrorMessage string - ObservedAt time.Time + WorkloadID string + Status WorkloadState + RuntimeRef string + Endpoint string + NodeID string + ErrorMessage string + ObservedAt time.Time + CPUMillicores int64 + MemoryMB int64 + NetworkRxMB float64 + NetworkTxMB float64 + DiskReadMB float64 + DiskWriteMB float64 } // WorkloadStatusChanged is emitted after daemon callbacks are persisted. diff --git a/internal/database/migrations/009_usage_records.sql b/internal/database/migrations/009_usage_records.sql new file mode 100644 index 0000000..a202f30 --- /dev/null +++ b/internal/database/migrations/009_usage_records.sql @@ -0,0 +1,16 @@ +CREATE TABLE IF NOT EXISTS usage_records ( + id TEXT PRIMARY KEY, + organization_id TEXT NOT NULL, + workload_id TEXT NOT NULL, + node_id TEXT NOT NULL, + recorded_at TIMESTAMPTZ NOT NULL, + cpu_seconds DOUBLE PRECISION NOT NULL DEFAULT 0, + memory_gb_hours DOUBLE PRECISION NOT NULL DEFAULT 0, + network_in_mb DOUBLE PRECISION NOT NULL DEFAULT 0, + network_out_mb DOUBLE PRECISION NOT NULL DEFAULT 0, + disk_read_mb DOUBLE PRECISION NOT NULL DEFAULT 0, + disk_write_mb DOUBLE PRECISION NOT NULL DEFAULT 0 +); + +CREATE INDEX IF NOT EXISTS idx_usage_records_workload ON usage_records (workload_id, recorded_at DESC); +CREATE INDEX IF NOT EXISTS idx_usage_records_org ON usage_records (organization_id, recorded_at DESC); diff --git a/internal/database/migrations/010_usage_display_columns.sql b/internal/database/migrations/010_usage_display_columns.sql new file mode 100644 index 0000000..1f67418 --- /dev/null +++ b/internal/database/migrations/010_usage_display_columns.sql @@ -0,0 +1,11 @@ +ALTER TABLE usage_records + ADD COLUMN IF NOT EXISTS project_id TEXT NOT NULL DEFAULT '', + ADD COLUMN IF NOT EXISTS cpu_millicores BIGINT NOT NULL DEFAULT 0, + ADD COLUMN IF NOT EXISTS memory_mb BIGINT NOT NULL DEFAULT 0, + ADD COLUMN IF NOT EXISTS network_in_kbps DOUBLE PRECISION NOT NULL DEFAULT 0, + ADD COLUMN IF NOT EXISTS network_out_kbps DOUBLE PRECISION NOT NULL DEFAULT 0, + ADD COLUMN IF NOT EXISTS disk_read_kbps DOUBLE PRECISION NOT NULL DEFAULT 0, + ADD COLUMN IF NOT EXISTS disk_write_kbps DOUBLE PRECISION NOT NULL DEFAULT 0; + +CREATE INDEX IF NOT EXISTS idx_usage_records_project_recorded + ON usage_records (project_id, recorded_at DESC); diff --git a/internal/database/migrations/011_workload_resources.sql b/internal/database/migrations/011_workload_resources.sql new file mode 100644 index 0000000..487e8c4 --- /dev/null +++ b/internal/database/migrations/011_workload_resources.sql @@ -0,0 +1,3 @@ +ALTER TABLE workloads + ADD COLUMN IF NOT EXISTS cpu_millicores BIGINT NOT NULL DEFAULT 0, + ADD COLUMN IF NOT EXISTS memory_bytes BIGINT NOT NULL DEFAULT 0; diff --git a/internal/database/migrations/012_workload_log_lines.sql b/internal/database/migrations/012_workload_log_lines.sql new file mode 100644 index 0000000..08a7c1e --- /dev/null +++ b/internal/database/migrations/012_workload_log_lines.sql @@ -0,0 +1,12 @@ +-- Stores raw log lines shipped from daemon containers. +CREATE TABLE IF NOT EXISTS workload_log_lines ( + id BIGSERIAL PRIMARY KEY, + workload_id TEXT NOT NULL, + project_id TEXT NOT NULL, + ts TIMESTAMPTZ NOT NULL, + stream TEXT NOT NULL DEFAULT 'stdout', -- 'stdout' or 'stderr' + line TEXT NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_wll_workload_ts ON workload_log_lines (workload_id, ts DESC); +CREATE INDEX IF NOT EXISTS idx_wll_project_ts ON workload_log_lines (project_id, ts DESC);