From bf31a8db37a1ca8231210446c5435709671fbf9a Mon Sep 17 00:00:00 2001 From: Philipp Date: Thu, 11 Jun 2026 10:24:36 +0200 Subject: [PATCH] feat: share proxmox cluster client with worker --- CHANGELOG.md | 1 + README.md | 4 +- TODO.md | 8 +- backend/cmd/api/main.go | 4 +- backend/internal/clusteradmin/handler.go | 2 +- backend/internal/clusteradmin/handler_test.go | 2 +- .../cluster/repository.go | 2 +- .../cluster/repository_test.go | 2 +- .../encryption/cipher.go | 0 .../encryption/cipher_test.go | 0 .../internal => platform}/proxmox/client.go | 40 +++++++++- .../proxmox/client_test.go | 25 +++++- worker/cmd/worker/main.go | 15 +++- worker/internal/tasks/proxmox_poll.go | 3 + worker/internal/tasks/proxmox_resolver.go | 69 +++++++++++++++++ .../internal/tasks/proxmox_resolver_test.go | 77 +++++++++++++++++++ 16 files changed, 240 insertions(+), 14 deletions(-) rename {backend/internal => platform}/cluster/repository.go (98%) rename {backend/internal => platform}/cluster/repository_test.go (98%) rename {backend/internal => platform}/encryption/cipher.go (100%) rename {backend/internal => platform}/encryption/cipher_test.go (100%) rename {backend/internal => platform}/proxmox/client.go (84%) rename {backend/internal => platform}/proxmox/client_test.go (80%) create mode 100644 worker/internal/tasks/proxmox_resolver.go create mode 100644 worker/internal/tasks/proxmox_resolver_test.go diff --git a/CHANGELOG.md b/CHANGELOG.md index f1bfcb4..8bb7062 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,7 @@ ## Unreleased +- Cluster-, Krypto- und Proxmox-Client-Bausteine nach `platform/` verschoben und im Worker fuer echte UPID-Polls verdrahtet. - UPID-Polling-Job fuer Proxmox-Tasks mit VM-Status-Update und Audit-Writer angelegt. - Worker-Grundgeruest mit asynq, Redis-Anbindung, DB-Ping und Dummy-Job-Handler angelegt. - Interne Operator-Endpunkte fuer Cluster-Anlage, Token-Rotation und Status-Updates ohne Token-Leak angelegt. diff --git a/README.md b/README.md index 6c9dd84..929f44d 100644 --- a/README.md +++ b/README.md @@ -40,7 +40,7 @@ Der Proxmox-Client nutzt HTTPS mit normaler Zertifikatsvalidierung plus SHA-256- Interne Cluster-Admin-Routen unter `/internal/clusters` sind nicht tenant-gebunden. Sie erfordern den Header `X-ProxUI-Operator-Token` mit `OPERATOR_TOKEN` und die isolierte Rolle `operator`; Responses enthalten nie `token_secret`. -Der Worker nutzt `asynq` gegen `REDIS_ADDR`, pingt beim Start dieselbe Supabase-DB ueber `DATABASE_URL` und registriert aktuell den Dummy-Job `proxui.dummy` als Smoke-Test fuer die Queue-Verarbeitung. +Der Worker nutzt `asynq` gegen `REDIS_ADDR`, pingt beim Start dieselbe Supabase-DB ueber `DATABASE_URL`, entschluesselt Cluster-Credentials mit `MASTER_KEY_BASE64` und nutzt den gemeinsamen Proxmox-Client aus `platform/proxmox` fuer UPID-Polls. Lokale Dienste: @@ -65,7 +65,7 @@ Aktuelle Targets: Worker-Jobs: - `proxui.dummy`: Dummy-Job mit Payload `{"message":"..."}` fuer Smoke-Tests der asynq-Verarbeitung -- `proxmox.task.poll`: Pollt einen Proxmox-UPID, setzt VM-Status auf Erfolgs- oder Fehlerstatus und schreibt Audit-Metadaten. Die echte Cluster-Client-Aufloesung wird im naechsten Integrationsschritt angeschlossen. +- `proxmox.task.poll`: Pollt einen Proxmox-UPID ueber den gespeicherten Cluster, setzt VM-Status auf Erfolgs- oder Fehlerstatus und schreibt Audit-Metadaten. Backend-Endpunkte: diff --git a/TODO.md b/TODO.md index 3a127f0..f20f3a6 100644 --- a/TODO.md +++ b/TODO.md @@ -132,7 +132,12 @@ Arbeitsliste auf Basis von `proxmox-console-entwicklungsplan.md`. Die Entwurfsda - [x] Polling-Handler mit Proxmox-Status-Client-Interface angelegt - [x] VM-Status-Update und Audit-Writer fuer finale Task-Ergebnisse angelegt - [x] Mock-Proxmox-Test fuer `running` -> `stopped/OK` und Fehlerstatus angelegt - - [x] Echte Cluster-Client-Aufloesung als separater Integrationspunkt gekapselt + - [x] Echte Cluster-Client-Aufloesung ueber gemeinsames Cluster-Repository und Proxmox-Client verdrahtet +- [x] Shared-Plattform-Follow-up fuer Worker-Integration + - [x] Cluster-Repository von `backend/internal` nach `platform/cluster` verschoben + - [x] Krypto-Layer nach `platform/encryption` verschoben + - [x] Proxmox-Client nach `platform/proxmox` verschoben + - [x] Backend und Worker nutzen dieselben Implementierungen ## MVP-Backlog @@ -192,3 +197,4 @@ Arbeitsliste auf Basis von `proxmox-console-entwicklungsplan.md`. Die Entwurfsda - 2026-06-11: Interne Operator-Endpunkte fuer Cluster-Anlage, Token-Rotation und Status-Updates angelegt. - 2026-06-11: Worker-Grundgeruest mit asynq, Redis-Anbindung, DB-Ping und Dummy-Job-Tests angelegt. - 2026-06-11: UPID-Polling-Job mit VM-Status-Update, Audit-Writer und Mock-Proxmox-Tests angelegt. +- 2026-06-11: Shared Cluster-/Krypto-/Proxmox-Packages nach `platform/` verschoben und Worker-UPID-Polling an echte Cluster-Aufloesung angeschlossen. diff --git a/backend/cmd/api/main.go b/backend/cmd/api/main.go index cf283b0..b059202 100644 --- a/backend/cmd/api/main.go +++ b/backend/cmd/api/main.go @@ -12,15 +12,15 @@ import ( "syscall" "time" + "forgejo.digital-droplets.de/philschlo/proxui/platform/cluster" "forgejo.digital-droplets.de/philschlo/proxui/platform/config" + "forgejo.digital-droplets.de/philschlo/proxui/platform/encryption" "forgejo.digital-droplets.de/philschlo/proxui/platform/logging" _ "github.com/jackc/pgx/v5/stdlib" "proxui/backend/internal/auth" "proxui/backend/internal/authorization" - "proxui/backend/internal/cluster" "proxui/backend/internal/clusteradmin" - "proxui/backend/internal/encryption" "proxui/backend/internal/membership" "proxui/backend/internal/operator" "proxui/backend/internal/profile" diff --git a/backend/internal/clusteradmin/handler.go b/backend/internal/clusteradmin/handler.go index 349201f..429a051 100644 --- a/backend/internal/clusteradmin/handler.go +++ b/backend/internal/clusteradmin/handler.go @@ -7,7 +7,7 @@ import ( "strings" "time" - "proxui/backend/internal/cluster" + "forgejo.digital-droplets.de/philschlo/proxui/platform/cluster" ) type Repository interface { diff --git a/backend/internal/clusteradmin/handler_test.go b/backend/internal/clusteradmin/handler_test.go index a2e63cb..be3cc70 100644 --- a/backend/internal/clusteradmin/handler_test.go +++ b/backend/internal/clusteradmin/handler_test.go @@ -11,7 +11,7 @@ import ( "net/http" "net/http/httptest" - "proxui/backend/internal/cluster" + "forgejo.digital-droplets.de/philschlo/proxui/platform/cluster" ) func TestCreateClusterDoesNotReturnTokenSecret(t *testing.T) { diff --git a/backend/internal/cluster/repository.go b/platform/cluster/repository.go similarity index 98% rename from backend/internal/cluster/repository.go rename to platform/cluster/repository.go index 1411739..5e6f6cb 100644 --- a/backend/internal/cluster/repository.go +++ b/platform/cluster/repository.go @@ -8,7 +8,7 @@ import ( "strings" "time" - "proxui/backend/internal/encryption" + "forgejo.digital-droplets.de/philschlo/proxui/platform/encryption" ) type Cluster struct { diff --git a/backend/internal/cluster/repository_test.go b/platform/cluster/repository_test.go similarity index 98% rename from backend/internal/cluster/repository_test.go rename to platform/cluster/repository_test.go index 0b8f2fb..5b95cd2 100644 --- a/backend/internal/cluster/repository_test.go +++ b/platform/cluster/repository_test.go @@ -7,7 +7,7 @@ import ( "testing" "time" - "proxui/backend/internal/encryption" + "forgejo.digital-droplets.de/philschlo/proxui/platform/encryption" ) func TestRepositoryEncryptsStoredTokenAndDecryptsOnLoad(t *testing.T) { diff --git a/backend/internal/encryption/cipher.go b/platform/encryption/cipher.go similarity index 100% rename from backend/internal/encryption/cipher.go rename to platform/encryption/cipher.go diff --git a/backend/internal/encryption/cipher_test.go b/platform/encryption/cipher_test.go similarity index 100% rename from backend/internal/encryption/cipher_test.go rename to platform/encryption/cipher_test.go diff --git a/backend/internal/proxmox/client.go b/platform/proxmox/client.go similarity index 84% rename from backend/internal/proxmox/client.go rename to platform/proxmox/client.go index 5cc19e8..0927d8a 100644 --- a/backend/internal/proxmox/client.go +++ b/platform/proxmox/client.go @@ -8,6 +8,7 @@ import ( "crypto/tls" "crypto/x509" "encoding/hex" + "encoding/json" "fmt" "io" "net/http" @@ -15,7 +16,7 @@ import ( "strings" "time" - "proxui/backend/internal/cluster" + "forgejo.digital-droplets.de/philschlo/proxui/platform/cluster" ) const defaultTimeout = 15 * time.Second @@ -125,6 +126,43 @@ func (c *Client) Delete(ctx context.Context, path string) (*http.Response, error return c.Do(ctx, http.MethodDelete, path, nil) } +type TaskStatus struct { + Status string + ExitStatus string +} + +func (c *Client) GetTaskStatus(ctx context.Context, node string, upid string) (TaskStatus, error) { + response, err := c.Get(ctx, fmt.Sprintf( + "/nodes/%s/tasks/%s/status", + url.PathEscape(node), + url.PathEscape(upid), + )) + if err != nil { + return TaskStatus{}, err + } + defer response.Body.Close() + + if response.StatusCode >= http.StatusBadRequest { + _, _ = io.Copy(io.Discard, response.Body) + return TaskStatus{}, fmt.Errorf("proxmox returned %s", response.Status) + } + + var body struct { + Data struct { + Status string `json:"status"` + ExitStatus string `json:"exitstatus"` + } `json:"data"` + } + if err := json.NewDecoder(response.Body).Decode(&body); err != nil { + return TaskStatus{}, err + } + + return TaskStatus{ + Status: body.Data.Status, + ExitStatus: body.Data.ExitStatus, + }, nil +} + func (c *Client) Do(ctx context.Context, method string, path string, body []byte) (*http.Response, error) { var lastErr error attempts := c.retries + 1 diff --git a/backend/internal/proxmox/client_test.go b/platform/proxmox/client_test.go similarity index 80% rename from backend/internal/proxmox/client_test.go rename to platform/proxmox/client_test.go index 8827733..5b30437 100644 --- a/backend/internal/proxmox/client_test.go +++ b/platform/proxmox/client_test.go @@ -12,7 +12,7 @@ import ( "testing" "time" - "proxui/backend/internal/cluster" + "forgejo.digital-droplets.de/philschlo/proxui/platform/cluster" ) func TestClientAcceptsMatchingFingerprintAndSendsTokenHeader(t *testing.T) { @@ -85,6 +85,29 @@ func TestClientRetriesServerErrors(t *testing.T) { } } +func TestClientGetsTaskStatus(t *testing.T) { + server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/api2/json/nodes/pve/tasks/UPID:pve:1/status" { + t.Fatalf("path = %q", r.URL.Path) + } + _, _ = w.Write([]byte(`{"data":{"status":"stopped","exitstatus":"OK"}}`)) + })) + defer server.Close() + + client := newTestClient(t, server, fingerprintForServer(server)) + status, err := client.GetTaskStatus(context.Background(), "pve", "UPID:pve:1") + if err != nil { + t.Fatalf("GetTaskStatus() error = %v", err) + } + + if status.Status != "stopped" { + t.Fatalf("Status = %q, want stopped", status.Status) + } + if status.ExitStatus != "OK" { + t.Fatalf("ExitStatus = %q, want OK", status.ExitStatus) + } +} + func TestNewClientRejectsInvalidFingerprint(t *testing.T) { _, err := NewClient(cluster.Cluster{ APIEndpoint: "https://pve.example.test:8006/api2/json", diff --git a/worker/cmd/worker/main.go b/worker/cmd/worker/main.go index cb04b65..e5cafab 100644 --- a/worker/cmd/worker/main.go +++ b/worker/cmd/worker/main.go @@ -11,7 +11,9 @@ import ( "syscall" "time" + "forgejo.digital-droplets.de/philschlo/proxui/platform/cluster" "forgejo.digital-droplets.de/philschlo/proxui/platform/config" + "forgejo.digital-droplets.de/philschlo/proxui/platform/encryption" "forgejo.digital-droplets.de/philschlo/proxui/platform/logging" "github.com/hibiken/asynq" _ "github.com/jackc/pgx/v5/stdlib" @@ -28,6 +30,12 @@ func main() { } logger := logging.New("worker", cfg.AppEnv, cfg.LogLevel) + tokenCipher, err := encryption.NewFromBase64(cfg.MasterKeyBase64) + if err != nil { + logger.Error("failed to initialize cluster token encryption", "error", err) + os.Exit(1) + } + db, err := openDatabase(cfg.DatabaseURL) if err != nil { logger.Error("failed to connect database", "error", err) @@ -43,7 +51,8 @@ func main() { Concurrency: cfg.WorkerConcurrency, }) mux := asynq.NewServeMux() - registerHandlers(mux, db, logger) + clusterRepository := cluster.NewRepository(db, tokenCipher) + registerHandlers(mux, db, clusterRepository, logger) go func() { logger.Info("worker started", "concurrency", cfg.WorkerConcurrency, "redis_addr", cfg.RedisAddr) @@ -58,12 +67,12 @@ func main() { logger.Info("worker stopped") } -func registerHandlers(mux *asynq.ServeMux, db *sql.DB, logger *slog.Logger) { +func registerHandlers(mux *asynq.ServeMux, db *sql.DB, clusters cluster.Repository, logger *slog.Logger) { mux.Handle(tasks.TypeDummy, tasks.NewDummyHandler(logger)) mux.Handle( tasks.TypeProxmoxTaskPoll, tasks.NewProxmoxTaskPollHandler( - tasks.UnconfiguredProxmoxTaskClientResolver{}, + tasks.NewDefaultClusterTaskClientResolver(clusters), tasks.NewSQLVMStatusStore(db), tasks.NewSQLAuditWriter(db), logger, diff --git a/worker/internal/tasks/proxmox_poll.go b/worker/internal/tasks/proxmox_poll.go index 1e044cb..1e00487 100644 --- a/worker/internal/tasks/proxmox_poll.go +++ b/worker/internal/tasks/proxmox_poll.go @@ -4,6 +4,7 @@ import ( "context" "database/sql" "encoding/json" + "errors" "fmt" "log/slog" "time" @@ -13,6 +14,8 @@ import ( const TypeProxmoxTaskPoll = "proxmox.task.poll" +var ErrClusterNotFound = errors.New("cluster not found") + type ProxmoxTaskPollPayload struct { ClusterID string `json:"cluster_id"` Node string `json:"node"` diff --git a/worker/internal/tasks/proxmox_resolver.go b/worker/internal/tasks/proxmox_resolver.go new file mode 100644 index 0000000..5e9b782 --- /dev/null +++ b/worker/internal/tasks/proxmox_resolver.go @@ -0,0 +1,69 @@ +package tasks + +import ( + "context" + + "forgejo.digital-droplets.de/philschlo/proxui/platform/cluster" + "forgejo.digital-droplets.de/philschlo/proxui/platform/proxmox" +) + +type ClusterRepository interface { + GetCluster(ctx context.Context, id string) (cluster.Cluster, bool, error) +} + +type ProxmoxTaskClientFactory func(cluster cluster.Cluster) (PlatformProxmoxTaskClient, error) + +type PlatformProxmoxTaskClient interface { + GetTaskStatus(ctx context.Context, node string, upid string) (proxmox.TaskStatus, error) +} + +type ClusterTaskClientResolver struct { + clusters ClusterRepository + factory ProxmoxTaskClientFactory +} + +func NewClusterTaskClientResolver(clusters ClusterRepository, factory ProxmoxTaskClientFactory) ClusterTaskClientResolver { + return ClusterTaskClientResolver{ + clusters: clusters, + factory: factory, + } +} + +func NewDefaultClusterTaskClientResolver(clusters ClusterRepository) ClusterTaskClientResolver { + return NewClusterTaskClientResolver(clusters, func(cluster cluster.Cluster) (PlatformProxmoxTaskClient, error) { + return proxmox.NewClient(cluster) + }) +} + +func (r ClusterTaskClientResolver) ResolveTaskClient(ctx context.Context, clusterID string) (ProxmoxTaskClient, error) { + cluster, found, err := r.clusters.GetCluster(ctx, clusterID) + if err != nil { + return nil, err + } + if !found { + return nil, ErrClusterNotFound + } + + client, err := r.factory(cluster) + if err != nil { + return nil, err + } + + return platformTaskClientAdapter{client: client}, nil +} + +type platformTaskClientAdapter struct { + client PlatformProxmoxTaskClient +} + +func (a platformTaskClientAdapter) GetTaskStatus(ctx context.Context, node string, upid string) (ProxmoxTaskStatus, error) { + status, err := a.client.GetTaskStatus(ctx, node, upid) + if err != nil { + return ProxmoxTaskStatus{}, err + } + + return ProxmoxTaskStatus{ + Status: status.Status, + ExitStatus: status.ExitStatus, + }, nil +} diff --git a/worker/internal/tasks/proxmox_resolver_test.go b/worker/internal/tasks/proxmox_resolver_test.go new file mode 100644 index 0000000..2c8012a --- /dev/null +++ b/worker/internal/tasks/proxmox_resolver_test.go @@ -0,0 +1,77 @@ +package tasks + +import ( + "context" + "errors" + "testing" + + "forgejo.digital-droplets.de/philschlo/proxui/platform/cluster" + "forgejo.digital-droplets.de/philschlo/proxui/platform/proxmox" +) + +func TestClusterTaskClientResolverLoadsClusterAndAdaptsStatus(t *testing.T) { + repository := &stubClusterRepository{ + cluster: cluster.Cluster{ + ID: "cluster-1", + APIEndpoint: "https://pve.example.test:8006", + TLSFingerprint: "fingerprint", + TokenID: "root@pam!proxui", + TokenSecret: "secret-token", + }, + found: true, + } + factoryClient := &stubPlatformTaskClient{ + status: proxmox.TaskStatus{Status: "stopped", ExitStatus: "OK"}, + } + resolver := NewClusterTaskClientResolver(repository, func(cluster cluster.Cluster) (PlatformProxmoxTaskClient, error) { + if cluster.ID != "cluster-1" { + t.Fatalf("cluster ID = %q, want cluster-1", cluster.ID) + } + return factoryClient, nil + }) + + client, err := resolver.ResolveTaskClient(context.Background(), "cluster-1") + if err != nil { + t.Fatalf("ResolveTaskClient() error = %v", err) + } + + status, err := client.GetTaskStatus(context.Background(), "pve", "UPID:pve:1") + if err != nil { + t.Fatalf("GetTaskStatus() error = %v", err) + } + if status.Status != "stopped" || status.ExitStatus != "OK" { + t.Fatalf("status = %+v, want stopped/OK", status) + } + if repository.id != "cluster-1" { + t.Fatalf("repository id = %q, want cluster-1", repository.id) + } +} + +func TestClusterTaskClientResolverReturnsClusterNotFound(t *testing.T) { + resolver := NewClusterTaskClientResolver(&stubClusterRepository{}, nil) + + if _, err := resolver.ResolveTaskClient(context.Background(), "missing"); !errors.Is(err, ErrClusterNotFound) { + t.Fatalf("ResolveTaskClient() error = %v, want ErrClusterNotFound", err) + } +} + +type stubClusterRepository struct { + id string + cluster cluster.Cluster + found bool + err error +} + +func (s *stubClusterRepository) GetCluster(_ context.Context, id string) (cluster.Cluster, bool, error) { + s.id = id + return s.cluster, s.found, s.err +} + +type stubPlatformTaskClient struct { + status proxmox.TaskStatus + err error +} + +func (s *stubPlatformTaskClient) GetTaskStatus(context.Context, string, string) (proxmox.TaskStatus, error) { + return s.status, s.err +}