From bef43e2fca434f36bead1c93f9fd09edb5aa81a8 Mon Sep 17 00:00:00 2001 From: Lukasz Zajaczkowski Date: Mon, 21 Sep 2026 12:55:40 +0200 Subject: [PATCH 1/4] add k8sObjMeta for lua and python --- .go/pkg/sumdb/sum.golang.org/latest | 5 + go/.gopath/pkg/sumdb/sum.golang.org/latest | 5 + .../pkg/streamline/common/component.go | 1 + .../pkg/streamline/global.go | 40 +++++ .../pkg/streamline/global_test.go | 93 +++++++++++ .../pkg/streamline/store/db_queries.go | 15 +- .../pkg/streamline/store/db_store.go | 64 +++++++- .../streamline/store/db_store_labels_test.go | 147 ++++++++++++++++++ .../pkg/streamline/store/db_store_test.go | 6 + .../store/db_store_unsynced_test.go | 4 +- .../pkg/streamline/store/labels.go | 42 +++++ .../pkg/streamline/store/labels_test.go | 57 +++++++ 12 files changed, 467 insertions(+), 12 deletions(-) create mode 100644 .go/pkg/sumdb/sum.golang.org/latest create mode 100644 go/.gopath/pkg/sumdb/sum.golang.org/latest create mode 100644 go/deployment-operator/pkg/streamline/global_test.go create mode 100644 go/deployment-operator/pkg/streamline/store/db_store_labels_test.go create mode 100644 go/deployment-operator/pkg/streamline/store/labels.go create mode 100644 go/deployment-operator/pkg/streamline/store/labels_test.go diff --git a/.go/pkg/sumdb/sum.golang.org/latest b/.go/pkg/sumdb/sum.golang.org/latest new file mode 100644 index 0000000000..ac6cee953e --- /dev/null +++ b/.go/pkg/sumdb/sum.golang.org/latest @@ -0,0 +1,5 @@ +go.sum database tree +64314384 +ZC3yZK2ozw+N7OOkZe/OHCVDokg96k9FUuWqOfuvEjs= + +— sum.golang.org Az3grlIFbSC/aQzBkvQCZ1PfCy51WV7p0Sj5LxueXWR7niRSTwLj9LCRkrufc8uCmoJOVHiJh/M8BS3/3WaSr1nJsAA= diff --git a/go/.gopath/pkg/sumdb/sum.golang.org/latest b/go/.gopath/pkg/sumdb/sum.golang.org/latest new file mode 100644 index 0000000000..a0ba796456 --- /dev/null +++ b/go/.gopath/pkg/sumdb/sum.golang.org/latest @@ -0,0 +1,5 @@ +go.sum database tree +63306336 +0KR7qhJAWpQCpbenlA7wrcrE07onVG8SmBiAmnoaTqI= + +— sum.golang.org Az3grk7fXu8pjGzqGpSAy9RRIInKkws8E7BbSqU6Gmm6LiAjbJAStefITvO7F4LO9oor6zkJqvpHSbrUpqNgNM9yaA4= diff --git a/go/deployment-operator/pkg/streamline/common/component.go b/go/deployment-operator/pkg/streamline/common/component.go index 2484bc49d0..e89bf4a588 100644 --- a/go/deployment-operator/pkg/streamline/common/component.go +++ b/go/deployment-operator/pkg/streamline/common/component.go @@ -25,6 +25,7 @@ type Component struct { TransientManifestSHA string ApplySHA string ServerSHA string + Labels map[string]string } func (in *Component) GroupVersionKind() schema.GroupVersionKind { diff --git a/go/deployment-operator/pkg/streamline/global.go b/go/deployment-operator/pkg/streamline/global.go index 14d2c03b7d..3af29a816e 100644 --- a/go/deployment-operator/pkg/streamline/global.go +++ b/go/deployment-operator/pkg/streamline/global.go @@ -4,6 +4,7 @@ import ( "sync" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" smcommon "github.com/pluralsh/console/go/deployment-operator/pkg/streamline/common" "github.com/pluralsh/console/go/deployment-operator/pkg/streamline/store" @@ -46,6 +47,45 @@ func (in *GlobalStore) GetComponent(obj unstructured.Unstructured) (result *smco return in.store.GetAppliedComponent(obj) } +// LookupObjectMeta returns cached Kubernetes object metadata from the process-wide store. +// It returns nil when the store is not initialized or the object is not cached. +func LookupObjectMeta(group, version, kind, namespace, name string) (map[string]any, error) { + return GetGlobalStore().ObjectMeta(group, version, kind, namespace, name) +} + +// ObjectMeta returns uid, name, namespace, and labels for a cached applied object. +// Cluster-scoped objects use an empty namespace. A cache miss returns nil, nil. +func (in *GlobalStore) ObjectMeta(group, version, kind, namespace, name string) (map[string]any, error) { + if in == nil || in.store == nil { + return nil, nil + } + + obj := unstructured.Unstructured{} + obj.SetGroupVersionKind(schema.GroupVersionKind{Group: group, Version: version, Kind: kind}) + obj.SetNamespace(namespace) + obj.SetName(name) + + component, err := in.GetComponent(obj) + if err != nil { + return nil, err + } + if component == nil { + return nil, nil + } + + labels := component.Labels + if labels == nil { + labels = map[string]string{} + } + + return map[string]any{ + "uid": component.UID, + "name": component.Name, + "namespace": component.Namespace, + "labels": labels, + }, nil +} + func (in *GlobalStore) UpdateComponentSHA(obj unstructured.Unstructured, shaType store.SHAType) error { return in.store.UpdateComponentSHA(obj, shaType) } diff --git a/go/deployment-operator/pkg/streamline/global_test.go b/go/deployment-operator/pkg/streamline/global_test.go new file mode 100644 index 0000000000..3a63eb6ea8 --- /dev/null +++ b/go/deployment-operator/pkg/streamline/global_test.go @@ -0,0 +1,93 @@ +package streamline_test + +import ( + "context" + "testing" + + "github.com/stretchr/testify/require" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/types" + + "github.com/pluralsh/console/go/deployment-operator/pkg/streamline" + "github.com/pluralsh/console/go/deployment-operator/pkg/streamline/store" +) + +func TestLookupObjectMeta(t *testing.T) { + t.Run("returns nil when the global store is not initialized", func(t *testing.T) { + streamline.ResetGlobalStore() + + got, err := streamline.LookupObjectMeta("", "v1", "Namespace", "", "kube-system") + require.NoError(t, err) + require.Nil(t, got) + }) + + t.Run("returns kube-system metadata from the cache", func(t *testing.T) { + streamline.ResetGlobalStore() + storeInstance, err := store.NewDatabaseStore(context.Background()) + require.NoError(t, err) + t.Cleanup(func() { + streamline.ResetGlobalStore() + require.NoError(t, storeInstance.Shutdown()) + }) + streamline.InitGlobalStore(storeInstance) + + ns := unstructured.Unstructured{} + ns.SetGroupVersionKind(schema.GroupVersionKind{Version: "v1", Kind: "Namespace"}) + ns.SetName("kube-system") + ns.SetUID(types.UID("cfb1383b-37cc-4d91-b943-aab5119e4cb1")) + ns.SetLabels(map[string]string{"kubernetes.io/metadata.name": "kube-system"}) + require.NoError(t, storeInstance.SaveComponent(ns)) + + got, err := streamline.LookupObjectMeta("", "v1", "Namespace", "", "kube-system") + require.NoError(t, err) + require.Equal(t, map[string]any{ + "uid": "cfb1383b-37cc-4d91-b943-aab5119e4cb1", + "name": "kube-system", + "namespace": "", + "labels": map[string]string{"kubernetes.io/metadata.name": "kube-system"}, + }, got) + }) + + t.Run("returns nil for a cache miss", func(t *testing.T) { + streamline.ResetGlobalStore() + storeInstance, err := store.NewDatabaseStore(context.Background()) + require.NoError(t, err) + t.Cleanup(func() { + streamline.ResetGlobalStore() + require.NoError(t, storeInstance.Shutdown()) + }) + streamline.InitGlobalStore(storeInstance) + + got, err := streamline.LookupObjectMeta("apps", "v1", "Deployment", "default", "missing") + require.NoError(t, err) + require.Nil(t, got) + }) + + t.Run("returns empty labels when the object has none", func(t *testing.T) { + streamline.ResetGlobalStore() + storeInstance, err := store.NewDatabaseStore(context.Background()) + require.NoError(t, err) + t.Cleanup(func() { + streamline.ResetGlobalStore() + require.NoError(t, storeInstance.Shutdown()) + }) + streamline.InitGlobalStore(storeInstance) + + cm := unstructured.Unstructured{} + cm.SetGroupVersionKind(schema.GroupVersionKind{Version: "v1", Kind: "ConfigMap"}) + cm.SetNamespace("default") + cm.SetName("plain") + cm.SetUID(types.UID("plain-uid")) + require.NoError(t, storeInstance.SaveComponent(cm)) + + got, err := streamline.GetGlobalStore().ObjectMeta("", "v1", "ConfigMap", "default", "plain") + require.NoError(t, err) + require.Equal(t, map[string]any{ + "uid": "plain-uid", + "name": "plain", + "namespace": "default", + "labels": map[string]string{}, + }, got) + }) +} diff --git a/go/deployment-operator/pkg/streamline/store/db_queries.go b/go/deployment-operator/pkg/streamline/store/db_queries.go index f4b4a3e58b..149d4c7ddd 100644 --- a/go/deployment-operator/pkg/streamline/store/db_queries.go +++ b/go/deployment-operator/pkg/streamline/store/db_queries.go @@ -22,7 +22,8 @@ const ( apply_sha TEXT, server_sha TEXT, manifest BOOLEAN DEFAULT 0, -- Indicates if the component was created from an original manifest set of a service - applied BOOLEAN DEFAULT 0 -- Indicates if the component was already applied to the cluster + applied BOOLEAN DEFAULT 0, -- Indicates if the component was already applied to the cluster + labels TEXT -- JSON object of Kubernetes object labels ); CREATE UNIQUE INDEX IF NOT EXISTS idx_unique_component ON component("group", version, kind, namespace, name); CREATE INDEX IF NOT EXISTS idx_parent ON component(parent_uid); @@ -68,7 +69,7 @@ const ( ` getAppliedComponent = ` - SELECT uid, "group", version, kind, namespace, name, health, parent_uid, manifest_sha, transient_manifest_sha, apply_sha, server_sha, service_id, manifest + SELECT uid, "group", version, kind, namespace, name, health, parent_uid, manifest_sha, transient_manifest_sha, apply_sha, server_sha, service_id, manifest, labels FROM component WHERE name = ? AND namespace = ? AND "group" = ? AND version = ? AND kind = ? AND applied = 1 ` @@ -100,7 +101,8 @@ const ( service_id, delete_phase, server_sha, - applied + applied, + labels ) VALUES ( ?, ?, @@ -115,6 +117,7 @@ const ( ?, ?, ?, + ?, ? ) ON CONFLICT("group", version, kind, namespace, name) DO UPDATE SET uid = excluded.uid, @@ -125,7 +128,8 @@ const ( service_id = excluded.service_id, delete_phase = excluded.delete_phase, server_sha = excluded.server_sha, - applied = excluded.applied + applied = excluded.applied, + labels = excluded.labels ` setComponentUnsynced = ` @@ -137,7 +141,8 @@ const ( server_sha = '', manifest_sha = '', transient_manifest_sha = '', - apply_sha = '' + apply_sha = '', + labels = '' WHERE "group" = ? AND version = ? AND kind = ? AND namespace = ? AND name = ? ` diff --git a/go/deployment-operator/pkg/streamline/store/db_store.go b/go/deployment-operator/pkg/streamline/store/db_store.go index 65d0285817..4dfcd8cc41 100644 --- a/go/deployment-operator/pkg/streamline/store/db_store.go +++ b/go/deployment-operator/pkg/streamline/store/db_store.go @@ -131,7 +131,29 @@ func (in *DatabaseStore) init() error { cancelFunc() }() - return sqlitex.ExecuteScript(conn, createTables, nil) + if err := sqlitex.ExecuteScript(conn, createTables, nil); err != nil { + return err + } + + return in.ensureLabelsColumn(conn) +} + +func (in *DatabaseStore) ensureLabelsColumn(conn *sqlite.Conn) error { + hasLabels := false + err := sqlitex.ExecuteTransient(conn, `SELECT 1 FROM pragma_table_info('component') WHERE name = 'labels'`, &sqlitex.ExecOptions{ + ResultFunc: func(_ *sqlite.Stmt) error { + hasLabels = true + return nil + }, + }) + if err != nil { + return err + } + if hasLabels { + return nil + } + + return sqlitex.Execute(conn, `ALTER TABLE component ADD COLUMN labels TEXT`, nil) } func (in *DatabaseStore) take() (*sqlite.Conn, context.CancelFunc, error) { @@ -260,10 +282,20 @@ func (in *DatabaseStore) SaveComponent(obj unstructured.Unstructured) error { in.maybeSaveHookComponent(conn, obj, lo.FromPtr(status), serviceID) + return in.upsertAppliedComponent(conn, obj, lo.FromPtr(ownerRef), nodeName, serviceID, serverSHA) +} + +func (in *DatabaseStore) upsertAppliedComponent(conn *sqlite.Conn, obj unstructured.Unstructured, ownerRef, nodeName, serviceID, serverSHA string) error { + labels, err := encodeComponentLabels(obj) + if err != nil { + return err + } + + gvk := obj.GroupVersionKind() return sqlitex.ExecuteTransient(conn, setComponentWithSHA, &sqlitex.ExecOptions{ Args: []interface{}{ obj.GetUID(), - lo.FromPtr(ownerRef), + ownerRef, gvk.Group, gvk.Version, gvk.Kind, @@ -276,6 +308,7 @@ func (in *DatabaseStore) SaveComponent(obj unstructured.Unstructured) error { smcommon.GetDeletePhase(obj), serverSHA, true, + labels, }, }) } @@ -315,7 +348,8 @@ func (in *DatabaseStore) SaveComponents(objects []unstructured.Unstructured) err service_id, delete_phase, server_sha, - applied + applied, + labels ) VALUES `) valueStrings := make([]string, 0, len(objects)) @@ -357,7 +391,13 @@ func (in *DatabaseStore) SaveComponents(objects []unstructured.Unstructured) err continue } - valueStrings = append(valueStrings, fmt.Sprintf("('%s','%s','%s','%s','%s','%s','%s',%d,'%s',%d,'%s','%s','%s', 1)", + labels, err := encodeComponentLabels(obj) + if err != nil { + klog.V(log.LogLevelDefault).ErrorS(err, "failed to encode resource labels", "name", obj.GetName(), "namespace", obj.GetNamespace(), "gvk", gvk.String()) + continue + } + + valueStrings = append(valueStrings, fmt.Sprintf("('%s','%s','%s','%s','%s','%s','%s',%d,'%s',%d,'%s','%s','%s', 1, '%s')", obj.GetUID(), lo.FromPtr(ownerRef), gvk.Group, @@ -371,6 +411,7 @@ func (in *DatabaseStore) SaveComponents(objects []unstructured.Unstructured) err serviceID, smcommon.GetDeletePhase(obj), serverSHA, + escapeSQLString(labels), )) } @@ -388,7 +429,8 @@ func (in *DatabaseStore) SaveComponents(objects []unstructured.Unstructured) err service_id = excluded.service_id, delete_phase = excluded.delete_phase, server_sha = excluded.server_sha, - applied = excluded.applied + applied = excluded.applied, + labels = excluded.labels `) in.maybeSaveHookComponents(conn, objects) @@ -698,6 +740,10 @@ func (in *DatabaseStore) GetAppliedComponent(obj unstructured.Unstructured) (res err = sqlitex.ExecuteTransient(conn, getAppliedComponent, &sqlitex.ExecOptions{ Args: []interface{}{obj.GetName(), obj.GetNamespace(), gvk.Group, gvk.Version, gvk.Kind}, ResultFunc: func(stmt *sqlite.Stmt) error { + labels, err := decodeComponentLabels(stmt.ColumnText(14)) + if err != nil { + return err + } result = &smcommon.Component{ UID: stmt.ColumnText(0), Group: stmt.ColumnText(1), @@ -713,6 +759,7 @@ func (in *DatabaseStore) GetAppliedComponent(obj unstructured.Unstructured) (res ServerSHA: stmt.ColumnText(11), ServiceID: stmt.ColumnText(12), Manifest: stmt.ColumnBool(13), + Labels: labels, } return nil }, @@ -827,7 +874,7 @@ func (in *DatabaseStore) GetServiceComponents(serviceID string, onlyApplied bool }() var sb strings.Builder - sb.WriteString(`SELECT uid, parent_uid, "group", version, kind, name, namespace, health, delete_phase, manifest, applied + sb.WriteString(`SELECT uid, parent_uid, "group", version, kind, name, namespace, health, delete_phase, manifest, applied, labels FROM component WHERE service_id = ?`) if onlyApplied { sb.WriteString(" AND applied = 1") @@ -843,6 +890,10 @@ func (in *DatabaseStore) GetServiceComponents(serviceID string, onlyApplied bool err = sqlitex.ExecuteTransient(conn, sb.String(), &sqlitex.ExecOptions{ Args: []interface{}{serviceID}, ResultFunc: func(stmt *sqlite.Stmt) error { + labels, err := decodeComponentLabels(stmt.ColumnText(11)) + if err != nil { + return err + } result = append(result, smcommon.Component{ UID: stmt.ColumnText(0), ParentUID: stmt.ColumnText(1), @@ -855,6 +906,7 @@ func (in *DatabaseStore) GetServiceComponents(serviceID string, onlyApplied bool ServiceID: serviceID, DeletePhase: stmt.ColumnText(8), Manifest: stmt.ColumnBool(9), + Labels: labels, }) return nil }, diff --git a/go/deployment-operator/pkg/streamline/store/db_store_labels_test.go b/go/deployment-operator/pkg/streamline/store/db_store_labels_test.go new file mode 100644 index 0000000000..54c0d05d62 --- /dev/null +++ b/go/deployment-operator/pkg/streamline/store/db_store_labels_test.go @@ -0,0 +1,147 @@ +package store_test + +import ( + "context" + "path/filepath" + "testing" + + "github.com/stretchr/testify/require" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/types" + "zombiezen.com/go/sqlite" + "zombiezen.com/go/sqlite/sqlitex" + + "github.com/pluralsh/console/go/deployment-operator/pkg/streamline/api" + "github.com/pluralsh/console/go/deployment-operator/pkg/streamline/store" +) + +func TestComponentCache_Labels(t *testing.T) { + ctx := context.Background() + storeInstance, err := store.NewDatabaseStore(ctx) + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, storeInstance.Shutdown()) + }) + + labels := map[string]string{ + "app": "observe", + "kubernetes.io/metadata.name": "kube-system", + "quoted": `foo's "bar"`, + } + + t.Run("SaveComponent round-trips labels", func(t *testing.T) { + ns := kubeSystemNamespace("cfb1383b-37cc-4d91-b943-aab5119e4cb1", labels) + require.NoError(t, storeInstance.SaveComponent(ns)) + + got, err := storeInstance.GetAppliedComponent(ns) + require.NoError(t, err) + require.NotNil(t, got) + require.Equal(t, string(ns.GetUID()), got.UID) + require.Equal(t, "kube-system", got.Name) + require.Empty(t, got.Namespace) + require.Equal(t, labels, got.Labels) + }) + + t.Run("SaveComponent stores empty labels as an empty map", func(t *testing.T) { + obj := createComponent("empty-labels-uid", WithName("no-labels")) + require.NoError(t, storeInstance.SaveComponent(obj)) + + got, err := storeInstance.GetAppliedComponent(obj) + require.NoError(t, err) + require.NotNil(t, got) + require.Empty(t, got.Labels) + }) + + t.Run("SaveComponents round-trips labels", func(t *testing.T) { + first := createComponent("batch-uid-1", WithName("batch-1"), WithLabels(map[string]string{"env": "prod"})) + second := createComponent("batch-uid-2", WithName("batch-2"), WithLabels(map[string]string{"quoted": `foo's "bar"`})) + require.NoError(t, storeInstance.SaveComponents([]unstructured.Unstructured{first, second})) + + gotFirst, err := storeInstance.GetAppliedComponent(first) + require.NoError(t, err) + require.Equal(t, map[string]string{"env": "prod"}, gotFirst.Labels) + + gotSecond, err := storeInstance.GetAppliedComponent(second) + require.NoError(t, err) + require.Equal(t, map[string]string{"quoted": `foo's "bar"`}, gotSecond.Labels) + }) + + t.Run("SaveComponent updates labels on conflict", func(t *testing.T) { + obj := createComponent("update-labels-uid", WithName("update-labels"), WithLabels(map[string]string{"v": "1"})) + require.NoError(t, storeInstance.SaveComponent(obj)) + + obj.SetLabels(map[string]string{"v": "2"}) + require.NoError(t, storeInstance.SaveComponent(obj)) + + got, err := storeInstance.GetAppliedComponent(obj) + require.NoError(t, err) + require.Equal(t, map[string]string{"v": "2"}, got.Labels) + }) +} + +func TestComponentCache_LabelsColumnMigratesExistingFile(t *testing.T) { + ctx := context.Background() + dbPath := filepath.Join(t.TempDir(), "store.db") + require.NoError(t, createLegacyComponentDB(dbPath)) + + storeInstance, err := store.NewDatabaseStore(ctx, store.WithStorage(api.StorageFile), store.WithFilePath(dbPath)) + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, storeInstance.Shutdown()) + }) + + ns := kubeSystemNamespace("migrated-uid", map[string]string{"app": "observe"}) + require.NoError(t, storeInstance.SaveComponent(ns)) + + got, err := storeInstance.GetAppliedComponent(ns) + require.NoError(t, err) + require.NotNil(t, got) + require.Equal(t, "migrated-uid", got.UID) + require.Equal(t, map[string]string{"app": "observe"}, got.Labels) +} + +func kubeSystemNamespace(uid string, labels map[string]string) unstructured.Unstructured { + ns := unstructured.Unstructured{} + ns.SetGroupVersionKind(schema.GroupVersionKind{Version: "v1", Kind: "Namespace"}) + ns.SetName("kube-system") + ns.SetUID(types.UID(uid)) + ns.SetLabels(labels) + return ns +} + +func createLegacyComponentDB(path string) error { + conn, err := sqlite.OpenConn(path, sqlite.OpenCreate, sqlite.OpenReadWrite) + if err != nil { + return err + } + defer func() { + _ = conn.Close() + }() + + return sqlitex.ExecuteScript(conn, ` + CREATE TABLE component ( + id INTEGER PRIMARY KEY, + parent_uid TEXT, + uid TEXT, + "group" TEXT, + version TEXT, + kind TEXT, + namespace TEXT, + name TEXT, + health INT, + node TEXT, + created_at TIMESTAMP, + updated_at TIMESTAMP, + service_id TEXT, + delete_phase TEXT, + manifest_sha TEXT, + transient_manifest_sha TEXT, + apply_sha TEXT, + server_sha TEXT, + manifest BOOLEAN DEFAULT 0, + applied BOOLEAN DEFAULT 0 + ); + CREATE UNIQUE INDEX idx_unique_component ON component("group", version, kind, namespace, name); + `, nil) +} diff --git a/go/deployment-operator/pkg/streamline/store/db_store_test.go b/go/deployment-operator/pkg/streamline/store/db_store_test.go index 7eae0af5be..3a66501dec 100644 --- a/go/deployment-operator/pkg/streamline/store/db_store_test.go +++ b/go/deployment-operator/pkg/streamline/store/db_store_test.go @@ -86,6 +86,12 @@ func WithName(name string) CreateComponentOption { } } +func WithLabels(labels map[string]string) CreateComponentOption { + return func(u *unstructured.Unstructured) { + u.SetLabels(labels) + } +} + type CreateStoreKeyOption func(entry *common.Component) func WithStoreKeyName(name string) CreateStoreKeyOption { diff --git a/go/deployment-operator/pkg/streamline/store/db_store_unsynced_test.go b/go/deployment-operator/pkg/streamline/store/db_store_unsynced_test.go index cf17f0e264..af8e8a35f4 100644 --- a/go/deployment-operator/pkg/streamline/store/db_store_unsynced_test.go +++ b/go/deployment-operator/pkg/streamline/store/db_store_unsynced_test.go @@ -16,7 +16,7 @@ func TestSetComponentUnsynced(t *testing.T) { require.NoError(t, err) serviceID := "test-service" - component := createComponent(testUID, WithService(serviceID)) + component := createComponent(testUID, WithService(serviceID), WithLabels(map[string]string{"app": "console"})) err = storeInstance.SaveComponent(component) require.NoError(t, err) @@ -25,6 +25,7 @@ func TestSetComponentUnsynced(t *testing.T) { require.NoError(t, err) assert.NotNil(t, savedComponent) assert.Equal(t, testUID, savedComponent.UID) + assert.Equal(t, map[string]string{"app": "console"}, savedComponent.Labels) // Set component unsynced err = storeInstance.SetComponentUnsynced(component) @@ -48,4 +49,5 @@ func TestSetComponentUnsynced(t *testing.T) { assert.Equal(t, testKind, unsyncedComponent.Kind) assert.Equal(t, testName, unsyncedComponent.Name) assert.Equal(t, testNamespace, unsyncedComponent.Namespace) + assert.Empty(t, unsyncedComponent.Labels) } diff --git a/go/deployment-operator/pkg/streamline/store/labels.go b/go/deployment-operator/pkg/streamline/store/labels.go new file mode 100644 index 0000000000..7ca8d95e31 --- /dev/null +++ b/go/deployment-operator/pkg/streamline/store/labels.go @@ -0,0 +1,42 @@ +package store + +import ( + "encoding/json" + "strings" + + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" +) + +func encodeComponentLabels(obj unstructured.Unstructured) (string, error) { + labels := obj.GetLabels() + if labels == nil { + labels = map[string]string{} + } + + encoded, err := json.Marshal(labels) + if err != nil { + return "", err + } + + return string(encoded), nil +} + +func decodeComponentLabels(raw string) (map[string]string, error) { + if raw == "" || raw == "null" { + return map[string]string{}, nil + } + + labels := map[string]string{} + if err := json.Unmarshal([]byte(raw), &labels); err != nil { + return nil, err + } + if labels == nil { + return map[string]string{}, nil + } + + return labels, nil +} + +func escapeSQLString(value string) string { + return strings.ReplaceAll(value, "'", "''") +} diff --git a/go/deployment-operator/pkg/streamline/store/labels_test.go b/go/deployment-operator/pkg/streamline/store/labels_test.go new file mode 100644 index 0000000000..ab64ed39c8 --- /dev/null +++ b/go/deployment-operator/pkg/streamline/store/labels_test.go @@ -0,0 +1,57 @@ +package store + +import ( + "testing" + + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" +) + +func TestEncodeDecodeComponentLabels(t *testing.T) { + t.Run("nil labels round-trip to empty map", func(t *testing.T) { + encoded, err := encodeComponentLabels(unstructured.Unstructured{}) + if err != nil { + t.Fatalf("encode: %v", err) + } + if encoded != "{}" { + t.Fatalf("expected {}, got %q", encoded) + } + + decoded, err := decodeComponentLabels(encoded) + if err != nil { + t.Fatalf("decode: %v", err) + } + if len(decoded) != 0 { + t.Fatalf("expected empty map, got %#v", decoded) + } + }) + + t.Run("empty and null stored values decode to empty map", func(t *testing.T) { + for _, raw := range []string{"", "null"} { + decoded, err := decodeComponentLabels(raw) + if err != nil { + t.Fatalf("decode %q: %v", raw, err) + } + if len(decoded) != 0 { + t.Fatalf("decode %q: expected empty map, got %#v", raw, decoded) + } + } + }) + + t.Run("labels with quotes round-trip", func(t *testing.T) { + obj := unstructured.Unstructured{} + obj.SetLabels(map[string]string{"quoted": `foo's "bar"`}) + + encoded, err := encodeComponentLabels(obj) + if err != nil { + t.Fatalf("encode: %v", err) + } + + decoded, err := decodeComponentLabels(encoded) + if err != nil { + t.Fatalf("decode: %v", err) + } + if decoded["quoted"] != `foo's "bar"` { + t.Fatalf("unexpected labels: %#v", decoded) + } + }) +} From 92cdf31e1bc4515db96fc12d08d4f7ddcc8c872d Mon Sep 17 00:00:00 2001 From: Lukasz Zajaczkowski Date: Mon, 21 Sep 2026 14:20:32 +0200 Subject: [PATCH 2/4] add k8s_object_meta to lua --- .../pkg/manifests/template/helm.go | 2 + .../pkg/manifests/template/helm_lua.go | 30 ++++++ .../pkg/manifests/template/helm_lua_test.go | 95 +++++++++++++++++++ 3 files changed, 127 insertions(+) create mode 100644 go/deployment-operator/pkg/manifests/template/helm_lua.go create mode 100644 go/deployment-operator/pkg/manifests/template/helm_lua_test.go diff --git a/go/deployment-operator/pkg/manifests/template/helm.go b/go/deployment-operator/pkg/manifests/template/helm.go index b30a57b904..79a5cf4fc7 100644 --- a/go/deployment-operator/pkg/manifests/template/helm.go +++ b/go/deployment-operator/pkg/manifests/template/helm.go @@ -228,6 +228,8 @@ func (h *helm) luaValues(svc *console.ServiceDeploymentForAgent) (map[string]any L := luautils.NewLuaState(h.dir) defer L.Close() + registerLuaFunctions(L) + // Register global values and valuesFiles in Lua valuesTable := L.NewTable() L.SetGlobal("values", valuesTable) diff --git a/go/deployment-operator/pkg/manifests/template/helm_lua.go b/go/deployment-operator/pkg/manifests/template/helm_lua.go new file mode 100644 index 0000000000..ce225274c1 --- /dev/null +++ b/go/deployment-operator/pkg/manifests/template/helm_lua.go @@ -0,0 +1,30 @@ +package template + +import ( + "github.com/pluralsh/console/go/deployment-operator/pkg/streamline" + "github.com/pluralsh/console/go/polly/luautils" + lua "github.com/yuin/gopher-lua" +) + +func registerLuaFunctions(l *lua.LState) { + l.SetFuncs(l.G.Global, map[string]lua.LGFunction{ + "k8s_object_meta": luaK8sObjectMeta, + }) +} + +func luaK8sObjectMeta(l *lua.LState) int { + group := l.CheckString(1) + version := l.CheckString(2) + kind := l.CheckString(3) + namespace := l.CheckString(4) + name := l.CheckString(5) + + meta, err := streamline.LookupObjectMeta(group, version, kind, namespace, name) + if err != nil { + l.RaiseError("k8s_object_meta: %s", err.Error()) + return 0 + } + + l.Push(luautils.GoValueToLuaValue(l, meta)) + return 1 +} diff --git a/go/deployment-operator/pkg/manifests/template/helm_lua_test.go b/go/deployment-operator/pkg/manifests/template/helm_lua_test.go new file mode 100644 index 0000000000..efba854125 --- /dev/null +++ b/go/deployment-operator/pkg/manifests/template/helm_lua_test.go @@ -0,0 +1,95 @@ +package template + +import ( + "context" + "testing" + + console "github.com/pluralsh/console/go/client" + "github.com/samber/lo" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/types" + + "github.com/pluralsh/console/go/deployment-operator/pkg/streamline" + "github.com/pluralsh/console/go/deployment-operator/pkg/streamline/store" +) + +func TestLuaValuesK8sObjectMeta(t *testing.T) { + t.Run("reads kube-system metadata from the cache", func(t *testing.T) { + streamline.ResetGlobalStore() + storeInstance, err := store.NewDatabaseStore(context.Background()) + if err != nil { + t.Fatalf("NewDatabaseStore: %v", err) + } + t.Cleanup(func() { + streamline.ResetGlobalStore() + _ = storeInstance.Shutdown() + }) + streamline.InitGlobalStore(storeInstance) + + ns := unstructured.Unstructured{} + ns.SetGroupVersionKind(schema.GroupVersionKind{Version: "v1", Kind: "Namespace"}) + ns.SetName("kube-system") + ns.SetUID(types.UID("cfb1383b-37cc-4d91-b943-aab5119e4cb1")) + ns.SetLabels(map[string]string{"kubernetes.io/metadata.name": "kube-system"}) + if err := storeInstance.SaveComponent(ns); err != nil { + t.Fatalf("SaveComponent: %v", err) + } + + svc := &console.ServiceDeploymentForAgent{ + Helm: &console.ServiceDeploymentForAgent_Helm{ + LuaScript: lo.ToPtr(` +local ns = k8s_object_meta("", "v1", "Namespace", "", "kube-system") +values["observeClusterId"] = ns.uid +values["name"] = ns.name +values["namespace"] = ns.namespace +values["label"] = ns.labels["kubernetes.io/metadata.name"] +`), + }, + } + + result, _, err := (&helm{dir: t.TempDir()}).luaValues(svc) + if err != nil { + t.Fatalf("luaValues: %v", err) + } + if result["observeClusterId"] != "cfb1383b-37cc-4d91-b943-aab5119e4cb1" { + t.Fatalf("unexpected uid: %#v", result) + } + if result["name"] != "kube-system" || result["namespace"] != "" { + t.Fatalf("unexpected identity: %#v", result) + } + if result["label"] != "kube-system" { + t.Fatalf("unexpected label: %#v", result) + } + }) + + t.Run("returns nil on a cache miss", func(t *testing.T) { + streamline.ResetGlobalStore() + storeInstance, err := store.NewDatabaseStore(context.Background()) + if err != nil { + t.Fatalf("NewDatabaseStore: %v", err) + } + t.Cleanup(func() { + streamline.ResetGlobalStore() + _ = storeInstance.Shutdown() + }) + streamline.InitGlobalStore(storeInstance) + + svc := &console.ServiceDeploymentForAgent{ + Helm: &console.ServiceDeploymentForAgent_Helm{ + LuaScript: lo.ToPtr(` +local missing = k8s_object_meta("apps", "v1", "Deployment", "default", "missing") +values["missing"] = missing == nil +`), + }, + } + + result, _, err := (&helm{dir: t.TempDir()}).luaValues(svc) + if err != nil { + t.Fatalf("luaValues: %v", err) + } + if result["missing"] != true { + t.Fatalf("expected nil on cache miss: %#v", result) + } + }) +} From 502aea2361af9283905949fb3ee04744d040d304 Mon Sep 17 00:00:00 2001 From: Lukasz Zajaczkowski Date: Mon, 21 Sep 2026 14:39:42 +0200 Subject: [PATCH 3/4] address comments --- .gitignore | 2 ++ .go/pkg/sumdb/sum.golang.org/latest | 5 ----- go/.gopath/pkg/sumdb/sum.golang.org/latest | 5 ----- .../pkg/streamline/store/db_store.go | 13 ++++++++++--- .../streamline/store/db_store_labels_test.go | 18 ++++++++++++++++++ 5 files changed, 30 insertions(+), 13 deletions(-) delete mode 100644 .go/pkg/sumdb/sum.golang.org/latest delete mode 100644 go/.gopath/pkg/sumdb/sum.golang.org/latest diff --git a/.gitignore b/.gitignore index 1d237ae7e7..0259f60fc1 100644 --- a/.gitignore +++ b/.gitignore @@ -25,6 +25,8 @@ node_modules/ # Go workspace go.work.sum .cache/ +.go/ +go/.gopath/ # Test binary, built with `go test -c` *.test diff --git a/.go/pkg/sumdb/sum.golang.org/latest b/.go/pkg/sumdb/sum.golang.org/latest deleted file mode 100644 index ac6cee953e..0000000000 --- a/.go/pkg/sumdb/sum.golang.org/latest +++ /dev/null @@ -1,5 +0,0 @@ -go.sum database tree -64314384 -ZC3yZK2ozw+N7OOkZe/OHCVDokg96k9FUuWqOfuvEjs= - -— sum.golang.org Az3grlIFbSC/aQzBkvQCZ1PfCy51WV7p0Sj5LxueXWR7niRSTwLj9LCRkrufc8uCmoJOVHiJh/M8BS3/3WaSr1nJsAA= diff --git a/go/.gopath/pkg/sumdb/sum.golang.org/latest b/go/.gopath/pkg/sumdb/sum.golang.org/latest deleted file mode 100644 index a0ba796456..0000000000 --- a/go/.gopath/pkg/sumdb/sum.golang.org/latest +++ /dev/null @@ -1,5 +0,0 @@ -go.sum database tree -63306336 -0KR7qhJAWpQCpbenlA7wrcrE07onVG8SmBiAmnoaTqI= - -— sum.golang.org Az3grk7fXu8pjGzqGpSAy9RRIInKkws8E7BbSqU6Gmm6LiAjbJAStefITvO7F4LO9oor6zkJqvpHSbrUpqNgNM9yaA4= diff --git a/go/deployment-operator/pkg/streamline/store/db_store.go b/go/deployment-operator/pkg/streamline/store/db_store.go index 4dfcd8cc41..30d621d06c 100644 --- a/go/deployment-operator/pkg/streamline/store/db_store.go +++ b/go/deployment-operator/pkg/streamline/store/db_store.go @@ -1248,6 +1248,11 @@ func (in *DatabaseStore) SyncAppliedResource(obj unstructured.Unstructured) erro return err } + labels, err := encodeComponentLabels(obj) + if err != nil { + return err + } + conn, cancelFunc, err := in.take() if err != nil { return err @@ -1269,7 +1274,8 @@ func (in *DatabaseStore) SyncAppliedResource(obj unstructured.Unstructured) erro END, transient_manifest_sha = NULL, manifest = 1, - applied = 1 + applied = 1, + labels = ? WHERE "group" = ? AND version = ? AND kind = ? @@ -1277,8 +1283,9 @@ func (in *DatabaseStore) SyncAppliedResource(obj unstructured.Unstructured) erro AND name = ? `, &sqlitex.ExecOptions{ Args: []interface{}{ - sha, // Apply SHA. - sha, // Server SHA. + sha, // Apply SHA. + sha, // Server SHA. + labels, gvk.Group, gvk.Version, gvk.Kind, obj.GetNamespace(), obj.GetName(), // WHERE clause parameters. }, }) diff --git a/go/deployment-operator/pkg/streamline/store/db_store_labels_test.go b/go/deployment-operator/pkg/streamline/store/db_store_labels_test.go index 54c0d05d62..65ce99b3f5 100644 --- a/go/deployment-operator/pkg/streamline/store/db_store_labels_test.go +++ b/go/deployment-operator/pkg/streamline/store/db_store_labels_test.go @@ -78,6 +78,24 @@ func TestComponentCache_Labels(t *testing.T) { require.NoError(t, err) require.Equal(t, map[string]string{"v": "2"}, got.Labels) }) + + t.Run("SyncAppliedResource persists labels from the live object", func(t *testing.T) { + manifest := createComponent("applied-uid", WithName("applied-labels"), WithGVK("apps", "v1", "Deployment")) + require.NoError(t, storeInstance.SaveUnsyncedComponents([]unstructured.Unstructured{manifest})) + + before, err := storeInstance.GetAppliedComponent(manifest) + require.NoError(t, err) + require.Nil(t, before) + + applied := createComponent("applied-uid", WithName("applied-labels"), WithGVK("apps", "v1", "Deployment"), WithLabels(map[string]string{"app": "observe"})) + require.NoError(t, storeInstance.SyncAppliedResource(applied)) + + got, err := storeInstance.GetAppliedComponent(applied) + require.NoError(t, err) + require.NotNil(t, got) + require.Equal(t, "applied-uid", got.UID) + require.Equal(t, map[string]string{"app": "observe"}, got.Labels) + }) } func TestComponentCache_LabelsColumnMigratesExistingFile(t *testing.T) { From 204d725249eea21744636e94f14572e8b5ce0a76 Mon Sep 17 00:00:00 2001 From: Lukasz Zajaczkowski Date: Mon, 21 Sep 2026 16:36:34 +0200 Subject: [PATCH 4/4] add python support --- go/deployment-operator/cmd/agent/main.go | 2 + .../manifests/template/helm_python_test.go | 120 ++++++++++++++++-- .../pkg/python/object_meta.go | 117 +++++++++++++++++ go/deployment-operator/pkg/python/pool.go | 12 +- .../pkg/python/pool_test.go | 54 ++++++++ .../continuous-deployment/helm-service.md | 27 ++++ .../continuous-deployment/lua.md | 29 +++++ 7 files changed, 345 insertions(+), 16 deletions(-) create mode 100644 go/deployment-operator/pkg/python/object_meta.go diff --git a/go/deployment-operator/cmd/agent/main.go b/go/deployment-operator/cmd/agent/main.go index 2267658cc1..e5e5283384 100644 --- a/go/deployment-operator/cmd/agent/main.go +++ b/go/deployment-operator/cmd/agent/main.go @@ -255,8 +255,10 @@ func initPythonRuntimeOrDie() func() { os.Exit(1) } + pythonruntime.SetObjectMetaLookup(streamline.LookupObjectMeta) pythonruntime.SetDefaultPool(pythonPool) return func() { + pythonruntime.SetObjectMetaLookup(nil) pythonruntime.SetDefaultPool(nil) if err := pythonPool.Close(); err != nil { setupLog.Error(err, "unable to shutdown Python runtime") diff --git a/go/deployment-operator/pkg/manifests/template/helm_python_test.go b/go/deployment-operator/pkg/manifests/template/helm_python_test.go index 1c3b22b205..08e9ba7bd7 100644 --- a/go/deployment-operator/pkg/manifests/template/helm_python_test.go +++ b/go/deployment-operator/pkg/manifests/template/helm_python_test.go @@ -9,7 +9,14 @@ import ( "testing" console "github.com/pluralsh/console/go/client" + "github.com/samber/lo" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/types" + "github.com/pluralsh/console/go/deployment-operator/pkg/python" + "github.com/pluralsh/console/go/deployment-operator/pkg/streamline" + "github.com/pluralsh/console/go/deployment-operator/pkg/streamline/store" ) func TestPythonValuesUsesBindingsAndFolderOrder(t *testing.T) { @@ -27,11 +34,11 @@ func TestPythonValuesUsesBindingsAndFolderOrder(t *testing.T) { Name: "demo", Namespace: "default", Cluster: &console.ServiceDeploymentForAgent_Cluster{ - Version: new("1.2.3"), + Version: lo.ToPtr("1.2.3"), }, Helm: &console.ServiceDeploymentForAgent_Helm{ - PythonFolder: new("python"), - PythonScript: new(` + PythonFolder: lo.ToPtr("python"), + PythonScript: lo.ToPtr(` values["order"] += "-main" values["version"] = cluster["version"] valuesFiles.append("generated.yaml") @@ -68,8 +75,8 @@ func TestPythonValuesInlineScriptWinsOverFile(t *testing.T) { svc := &console.ServiceDeploymentForAgent{ Helm: &console.ServiceDeploymentForAgent_Helm{ - PythonFile: new("values.py"), - PythonScript: new(`values["source"] = "inline"`), + PythonFile: lo.ToPtr("values.py"), + PythonScript: lo.ToPtr(`values["source"] = "inline"`), }, } result, _, err := (&helm{dir: dir, pythonPool: p}).pythonValues(context.Background(), svc) @@ -117,7 +124,7 @@ func TestPythonValuesCannotReadPythonFileOutsideDirectory(t *testing.T) { for _, test := range paths { t.Run(test.name, func(t *testing.T) { svc := &console.ServiceDeploymentForAgent{ - Helm: &console.ServiceDeploymentForAgent_Helm{PythonFile: new(test.path)}, + Helm: &console.ServiceDeploymentForAgent_Helm{PythonFile: lo.ToPtr(test.path)}, } _, _, err := (&helm{dir: dir, pythonPool: p}).pythonValues(context.Background(), svc) if err == nil { @@ -177,7 +184,7 @@ func TestTemplateValuesCannotReadValuesFileOutsideDirectory(t *testing.T) { svc := &console.ServiceDeploymentForAgent{ Helm: &console.ServiceDeploymentForAgent_Helm{ - PythonScript: new(fmt.Sprintf("valuesFiles.append(%q)", outsidePath)), + PythonScript: lo.ToPtr(fmt.Sprintf("valuesFiles.append(%q)", outsidePath)), }, } if _, err := (&helm{dir: dir, pythonPool: p}).templateValues(svc); err == nil { @@ -194,7 +201,7 @@ func TestPythonValuesErrorIsContextualized(t *testing.T) { svc := &console.ServiceDeploymentForAgent{ Helm: &console.ServiceDeploymentForAgent_Helm{ - PythonScript: new(`raise ValueError("not safe to render")`), + PythonScript: lo.ToPtr(`raise ValueError("not safe to render")`), }, } _, _, err = (&helm{dir: t.TempDir(), pythonPool: p}).pythonValues(context.Background(), svc) @@ -216,9 +223,9 @@ func TestTemplateValuesRunsLuaBeforePython(t *testing.T) { svc := &console.ServiceDeploymentForAgent{ Helm: &console.ServiceDeploymentForAgent_Helm{ - LuaScript: new(`values["collision"] = "lua" + LuaScript: lo.ToPtr(`values["collision"] = "lua" valuesFiles[1] = "lua.yaml"`), - PythonScript: new(`values["collision"] = "python" + PythonScript: lo.ToPtr(`values["collision"] = "python" valuesFiles.append("python.yaml")`), }, } @@ -231,6 +238,99 @@ valuesFiles.append("python.yaml")`), } } +func TestPythonValuesK8sObjectMeta(t *testing.T) { + python.SetObjectMetaLookup(streamline.LookupObjectMeta) + t.Cleanup(func() { python.SetObjectMetaLookup(nil) }) + + t.Run("reads kube-system metadata from the cache", func(t *testing.T) { + streamline.ResetGlobalStore() + storeInstance, err := store.NewDatabaseStore(context.Background()) + if err != nil { + t.Fatalf("NewDatabaseStore: %v", err) + } + t.Cleanup(func() { + streamline.ResetGlobalStore() + _ = storeInstance.Shutdown() + }) + streamline.InitGlobalStore(storeInstance) + + ns := unstructured.Unstructured{} + ns.SetGroupVersionKind(schema.GroupVersionKind{Version: "v1", Kind: "Namespace"}) + ns.SetName("kube-system") + ns.SetUID(types.UID("cfb1383b-37cc-4d91-b943-aab5119e4cb1")) + ns.SetLabels(map[string]string{"kubernetes.io/metadata.name": "kube-system"}) + if err := storeInstance.SaveComponent(ns); err != nil { + t.Fatalf("SaveComponent: %v", err) + } + + p, err := python.NewPoolWithConfig(python.Config{WorkerCount: 1, QueueSize: 1}) + if err != nil { + t.Fatalf("NewPoolWithConfig: %v", err) + } + t.Cleanup(func() { _ = p.Close() }) + + svc := &console.ServiceDeploymentForAgent{ + Helm: &console.ServiceDeploymentForAgent_Helm{ + PythonScript: lo.ToPtr(` +ns = k8s_object_meta("", "v1", "Namespace", "", "kube-system") +values["observeClusterId"] = ns["uid"] +values["name"] = ns["name"] +values["namespace"] = ns["namespace"] +values["label"] = ns["labels"]["kubernetes.io/metadata.name"] +`), + }, + } + result, _, err := (&helm{dir: t.TempDir(), pythonPool: p}).pythonValues(context.Background(), svc) + if err != nil { + t.Fatalf("pythonValues: %v", err) + } + if result["observeClusterId"] != "cfb1383b-37cc-4d91-b943-aab5119e4cb1" { + t.Fatalf("unexpected uid: %#v", result) + } + if result["name"] != "kube-system" || result["namespace"] != "" { + t.Fatalf("unexpected identity: %#v", result) + } + if result["label"] != "kube-system" { + t.Fatalf("unexpected label: %#v", result) + } + }) + + t.Run("returns None on a cache miss", func(t *testing.T) { + streamline.ResetGlobalStore() + storeInstance, err := store.NewDatabaseStore(context.Background()) + if err != nil { + t.Fatalf("NewDatabaseStore: %v", err) + } + t.Cleanup(func() { + streamline.ResetGlobalStore() + _ = storeInstance.Shutdown() + }) + streamline.InitGlobalStore(storeInstance) + + p, err := python.NewPoolWithConfig(python.Config{WorkerCount: 1, QueueSize: 1}) + if err != nil { + t.Fatalf("NewPoolWithConfig: %v", err) + } + t.Cleanup(func() { _ = p.Close() }) + + svc := &console.ServiceDeploymentForAgent{ + Helm: &console.ServiceDeploymentForAgent_Helm{ + PythonScript: lo.ToPtr(` +missing = k8s_object_meta("apps", "v1", "Deployment", "default", "missing") +values["missing"] = missing is None +`), + }, + } + result, _, err := (&helm{dir: t.TempDir(), pythonPool: p}).pythonValues(context.Background(), svc) + if err != nil { + t.Fatalf("pythonValues: %v", err) + } + if result["missing"] != true { + t.Fatalf("expected None on cache miss: %#v", result) + } + }) +} + func writePythonFile(t *testing.T, dir, name, contents string) { t.Helper() path := dir + "/" + name diff --git a/go/deployment-operator/pkg/python/object_meta.go b/go/deployment-operator/pkg/python/object_meta.go new file mode 100644 index 0000000000..f06a8687d3 --- /dev/null +++ b/go/deployment-operator/pkg/python/object_meta.go @@ -0,0 +1,117 @@ +package python + +import ( + "context" + "fmt" + "sync" + + monty "github.com/ewhauser/gomonty" +) + +// ObjectMetaLookup returns cached Kubernetes object metadata for Helm Python +// scripts. A nil map means the object is not in the cache. +type ObjectMetaLookup func(group, version, kind, namespace, name string) (map[string]any, error) + +var objectMetaLookup struct { + sync.RWMutex + fn ObjectMetaLookup +} + +// SetObjectMetaLookup installs the process-wide cache lookup used by +// k8s_object_meta. Pass nil to clear it. +func SetObjectMetaLookup(fn ObjectMetaLookup) { + objectMetaLookup.Lock() + objectMetaLookup.fn = fn + objectMetaLookup.Unlock() +} + +func (p *Pool) feedOptions() monty.FeedOptions { + return monty.FeedOptions{ + Functions: map[string]monty.ExternalFunction{ + "k8s_object_meta": k8sObjectMeta, + }, + } +} + +func k8sObjectMeta(_ context.Context, call monty.Call) (monty.Result, error) { + if len(call.Args) != 5 { + return raise("TypeError", "k8s_object_meta() takes 5 string arguments") + } + + group, err := stringArg(call.Args, 0) + if err != nil { + return raise("TypeError", err.Error()) + } + version, err := stringArg(call.Args, 1) + if err != nil { + return raise("TypeError", err.Error()) + } + kind, err := stringArg(call.Args, 2) + if err != nil { + return raise("TypeError", err.Error()) + } + namespace, err := stringArg(call.Args, 3) + if err != nil { + return raise("TypeError", err.Error()) + } + name, err := stringArg(call.Args, 4) + if err != nil { + return raise("TypeError", err.Error()) + } + + objectMetaLookup.RLock() + fn := objectMetaLookup.fn + objectMetaLookup.RUnlock() + if fn == nil { + return monty.Return(monty.None()), nil + } + + meta, err := fn(group, version, kind, namespace, name) + if err != nil { + return raise("RuntimeError", err.Error()) + } + if meta == nil { + return monty.Return(monty.None()), nil + } + + return monty.Return(goMapToValue(meta)), nil +} + +func stringArg(args []monty.Value, index int) (string, error) { + raw := args[index].Raw() + if raw == nil { + return "", nil + } + value, ok := raw.(string) + if !ok { + return "", fmt.Errorf("k8s_object_meta() argument %d must be a string", index+1) + } + return value, nil +} + +func raise(kind, message string) (monty.Result, error) { + return monty.Raise(monty.Exception{Type: kind, Arg: &message}), nil +} + +func goMapToValue(value any) monty.Value { + switch typed := value.(type) { + case nil: + return monty.None() + case string: + return monty.String(typed) + case map[string]string: + items := make(monty.Dict, 0, len(typed)) + for key, item := range typed { + items = append(items, monty.Pair{Key: monty.String(key), Value: monty.String(item)}) + } + return monty.DictValue(items) + case map[string]any: + items := make(monty.Dict, 0, len(typed)) + for key, item := range typed { + items = append(items, monty.Pair{Key: monty.String(key), Value: goMapToValue(item)}) + } + return monty.DictValue(items) + default: + return monty.None() + } +} diff --git a/go/deployment-operator/pkg/python/pool.go b/go/deployment-operator/pkg/python/pool.go index c5d84dd560..f8a6e3e31d 100644 --- a/go/deployment-operator/pkg/python/pool.go +++ b/go/deployment-operator/pkg/python/pool.go @@ -1,9 +1,9 @@ // Package python executes the user-provided Python used by Helm value // templating. // -// Scripts run in a fresh Gomonty REPL for every job. The package deliberately -// does not provide any Gomonty host callbacks, so scripts cannot access the -// operator process, filesystem, network, or environment. +// Scripts run in a fresh Gomonty REPL for every job. The sandbox does not +// expose OS, filesystem, network, or environment access. The only host +// callback is the read-only k8s_object_meta lookup against the agent cache. package python import ( @@ -292,14 +292,14 @@ func (p *Pool) execute(parentCtx context.Context, script string, bindings map[st "service = __helm_bindings.get('service')\n" + "values = {}\n" + "valuesFiles = []\n" - if _, err := repl.FeedRun(ctx, initialization, monty.FeedOptions{}); err != nil { + if _, err := repl.FeedRun(ctx, initialization, p.feedOptions()); err != nil { return Result{}, p.mapExecutionError(ctx, err) } - if _, err := repl.FeedRun(ctx, script, monty.FeedOptions{}); err != nil { + if _, err := repl.FeedRun(ctx, script, p.feedOptions()); err != nil { return Result{}, p.mapExecutionError(ctx, err) } - encoded, err := repl.FeedRun(ctx, "__helm_json.dumps({'values': values, 'valuesFiles': valuesFiles})", monty.FeedOptions{}) + encoded, err := repl.FeedRun(ctx, "__helm_json.dumps({'values': values, 'valuesFiles': valuesFiles})", p.feedOptions()) if err != nil { return Result{}, p.mapExecutionError(ctx, err) } diff --git a/go/deployment-operator/pkg/python/pool_test.go b/go/deployment-operator/pkg/python/pool_test.go index 87f4f9a9d4..52b02ab71f 100644 --- a/go/deployment-operator/pkg/python/pool_test.go +++ b/go/deployment-operator/pkg/python/pool_test.go @@ -258,6 +258,60 @@ func TestCloseIsIdempotentDuringRuns(t *testing.T) { runsWG.Wait() } +func TestRunK8sObjectMeta(t *testing.T) { + p := testPool(t, Config{WorkerCount: 1, QueueSize: 1}) + SetObjectMetaLookup(func(group, version, kind, namespace, name string) (map[string]any, error) { + if group != "" || version != "v1" || kind != "Namespace" || namespace != "" || name != "kube-system" { + return nil, nil + } + return map[string]any{ + "uid": "cfb1383b-37cc-4d91-b943-aab5119e4cb1", + "name": name, + "namespace": namespace, + "labels": map[string]string{"kubernetes.io/metadata.name": "kube-system"}, + }, nil + }) + t.Cleanup(func() { SetObjectMetaLookup(nil) }) + + result, err := p.Run(context.Background(), ` +ns = k8s_object_meta("", "v1", "Namespace", "", "kube-system") +values["observeClusterId"] = ns["uid"] +values["name"] = ns["name"] +values["namespace"] = ns["namespace"] +values["label"] = ns["labels"]["kubernetes.io/metadata.name"] +missing = k8s_object_meta("apps", "v1", "Deployment", "default", "missing") +values["missing"] = missing is None +`, nil) + if err != nil { + t.Fatalf("Run: %v", err) + } + if result.Values["observeClusterId"] != "cfb1383b-37cc-4d91-b943-aab5119e4cb1" { + t.Fatalf("unexpected uid: %#v", result.Values) + } + if result.Values["name"] != "kube-system" || result.Values["namespace"] != "" { + t.Fatalf("unexpected identity: %#v", result.Values) + } + if result.Values["label"] != "kube-system" { + t.Fatalf("unexpected label: %#v", result.Values) + } + if result.Values["missing"] != true { + t.Fatalf("expected None on cache miss: %#v", result.Values) + } +} + +func TestRunK8sObjectMetaRaisesStoreErrors(t *testing.T) { + p := testPool(t, Config{WorkerCount: 1, QueueSize: 1}) + SetObjectMetaLookup(func(string, string, string, string, string) (map[string]any, error) { + return nil, errors.New("store unavailable") + }) + t.Cleanup(func() { SetObjectMetaLookup(nil) }) + + _, err := p.Run(context.Background(), `k8s_object_meta("", "v1", "Namespace", "", "kube-system")`, nil) + if err == nil || !strings.Contains(err.Error(), "python RuntimeError") { + t.Fatalf("expected RuntimeError, got %v", err) + } +} + func TestRunSupportsConcurrentJobs(t *testing.T) { p := testPool(t, Config{WorkerCount: 2, QueueSize: 8}) diff --git a/js/documentation/pages/plural-features/continuous-deployment/helm-service.md b/js/documentation/pages/plural-features/continuous-deployment/helm-service.md index 87b3346f14..c9c41480cd 100644 --- a/js/documentation/pages/plural-features/continuous-deployment/helm-service.md +++ b/js/documentation/pages/plural-features/continuous-deployment/helm-service.md @@ -63,6 +63,33 @@ spec: ``` For more information, see [Dynamic Helm Configuration with Lua Scripts](lua.md). +## Dynamic Helm Configuration via pythonScript + +The same `values` / `valuesFiles` overlay is available from a sandboxed Python script. The sandbox does not expose OS, filesystem, or network access. The only host callback is `k8s_object_meta`, which reads cached Kubernetes object metadata (uid, name, namespace, and labels) from the agent. Cluster-scoped objects use an empty namespace. A cache miss returns `None`. + +```yaml +apiVersion: deployments.plural.sh/v1alpha1 +kind: ServiceDeployment +metadata: + name: observe + namespace: infra +spec: + namespace: observe + name: observe + cluster: k3s + helm: + version: 1.x.x + chart: observe + url: https://example.invalid/charts + pythonScript: | + ns = k8s_object_meta("", "v1", "Namespace", "", "kube-system") + if ns: + values["observeClusterId"] = ns["uid"] + values["label"] = ns["labels"]["kubernetes.io/metadata.name"] +``` + +The Lua equivalent of `k8s_object_meta` is documented in [Dynamic Helm Configuration with Lua Scripts](lua.md#kubernetes-object-metadata). + ## Multi-Source Helm Say you want to source the helm templates from an upstream helm repository, but the values files from a Git repository. In that case, you can define a multi-sourced service, which has both a git and helm repository defined. It would look like so: diff --git a/js/documentation/pages/plural-features/continuous-deployment/lua.md b/js/documentation/pages/plural-features/continuous-deployment/lua.md index 94d3c7e6c0..5e3c3aeb13 100644 --- a/js/documentation/pages/plural-features/continuous-deployment/lua.md +++ b/js/documentation/pages/plural-features/continuous-deployment/lua.md @@ -270,6 +270,35 @@ values["templatePath"] = templatePath valuesFiles = {templatePath} ``` +## Kubernetes Object Metadata + +Helm Lua scripts can look up cached Kubernetes object metadata from the deployment agent. This is a read against the agent's applied-object cache, not a live API request. Annotations are not included. + +### `k8s_object_meta(group, version, kind, namespace, name) -> table | nil` + +**Parameters:** +- `group` (string): API group. Use `""` for core objects such as Namespace +- `version` (string): API version, for example `"v1"` +- `kind` (string): Kind, for example `"Namespace"` +- `namespace` (string): Namespace, or `""` for cluster-scoped objects +- `name` (string): Object name + +**Returns:** +- `table` with `uid`, `name`, `namespace`, and `labels` when the object is cached +- `nil` when the object is not in the cache + +Store errors raise a Lua error. + +```lua +local ns = k8s_object_meta("", "v1", "Namespace", "", "kube-system") +if ns then + values["observeClusterId"] = ns.uid + values["name"] = ns.name + values["namespace"] = ns.namespace + values["label"] = ns.labels["kubernetes.io/metadata.name"] +end +``` + ## Error Handling All functions return errors in a consistent format: