Skip to content
This repository was archived by the owner on Apr 20, 2026. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions go.work.sum
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down
39 changes: 22 additions & 17 deletions internal/bootstrap/container.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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.
Expand Down Expand Up @@ -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),
Expand Down
2 changes: 2 additions & 0 deletions internal/bootstrap/http.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)

Expand Down
113 changes: 113 additions & 0 deletions internal/core/logs/adapters/http/handler.go
Original file line number Diff line number Diff line change
@@ -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)
}
75 changes: 75 additions & 0 deletions internal/core/logs/adapters/persistence/store.go
Original file line number Diff line number Diff line change
@@ -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()
}
13 changes: 13 additions & 0 deletions internal/core/logs/domain/log.go
Original file line number Diff line number Diff line change
@@ -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"`
}
20 changes: 20 additions & 0 deletions internal/core/logs/ports/repository.go
Original file line number Diff line number Diff line change
@@ -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)
}
45 changes: 33 additions & 12 deletions internal/core/usage/adapters/http/handler.go
Original file line number Diff line number Diff line change
@@ -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)
}
Loading
Loading