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 92a5314..a1d04c4 100644 --- a/internal/bootstrap/container.go +++ b/internal/bootstrap/container.go @@ -37,7 +37,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" @@ -62,22 +65,23 @@ type Container struct { PluginManager *pluginapplication.Manager // HTTP handler groups per domain module - AuthHandler *pluginhttp.AuthHandler - SetupHandler *pluginhttp.SetupHandler - CatalogHandler *cataloghttp.Handler - OrganizationsHandler *organizationshttp.Handler - ProjectsHandler *projectshttp.Handler - WorkloadsHandler *workloadshttp.Handler - DeploymentsHandler *deploymentshttp.Handler - NodesHandler *nodeshttp.Handler - BillingHandler *billinghttp.Handler - UsageHandler *usagehttp.Handler - AuditHandler *audithttp.Handler - AdminHandler *adminhttp.Handler - PluginsHandler *pluginhttp.Handler - NotificationsHandler *notificationshttp.Handler - NotificationService *notificationsapp.Service - NotificationHub *notificationsapp.Hub + AuthHandler *pluginhttp.AuthHandler + SetupHandler *pluginhttp.SetupHandler + CatalogHandler *cataloghttp.Handler + OrganizationsHandler *organizationshttp.Handler + ProjectsHandler *projectshttp.Handler + WorkloadsHandler *workloadshttp.Handler + DeploymentsHandler *deploymentshttp.Handler + NodesHandler *nodeshttp.Handler + BillingHandler *billinghttp.Handler + UsageHandler *usagehttp.Handler + LogsHandler *logshttp.Handler + AuditHandler *audithttp.Handler + AdminHandler *adminhttp.Handler + PluginsHandler *pluginhttp.Handler + NotificationsHandler *notificationshttp.Handler + NotificationService *notificationsapp.Service + NotificationHub *notificationsapp.Hub } // NewContainer wires all dependencies and returns the composition root. @@ -175,7 +179,8 @@ func NewContainer(cfg *Config, logger *slog.Logger) (*Container, error) { WorkloadsHandler: workloadshttp.NewHandler(projectsStore, orgStore, workloadsStore, provisionHandler, workloadAction, bus, logger), 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 4c36372..ed455da 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) c.NotificationsHandler.RegisterRoutes(r) 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 a7b77a4..3d4a4c7 100644 --- a/internal/core/workloads/adapters/http/handler.go +++ b/internal/core/workloads/adapters/http/handler.go @@ -172,12 +172,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"}) @@ -220,13 +226,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) { 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);