diff --git a/apis/core.oam.dev/v1beta1/resourcetracker_types.go b/apis/core.oam.dev/v1beta1/resourcetracker_types.go index 4fa1ef4ae..e3fb132f1 100644 --- a/apis/core.oam.dev/v1beta1/resourcetracker_types.go +++ b/apis/core.oam.dev/v1beta1/resourcetracker_types.go @@ -21,8 +21,8 @@ import ( "reflect" "strings" - errors2 "github.com/pkg/errors" - v1 "k8s.io/api/core/v1" + "github.com/pkg/errors" + corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime" @@ -33,7 +33,8 @@ import ( "github.com/oam-dev/kubevela/apis/interfaces" velatypes "github.com/oam-dev/kubevela/apis/types" "github.com/oam-dev/kubevela/pkg/oam" - "github.com/oam-dev/kubevela/pkg/utils/errors" + "github.com/oam-dev/kubevela/pkg/utils/compression" + velaerr "github.com/oam-dev/kubevela/pkg/utils/errors" ) // +kubebuilder:object:root=true @@ -69,9 +70,65 @@ const ( // ResourceTrackerSpec define the spec of resourceTracker type ResourceTrackerSpec struct { - Type ResourceTrackerType `json:"type,omitempty"` - ApplicationGeneration int64 `json:"applicationGeneration"` - ManagedResources []ManagedResource `json:"managedResources,omitempty"` + Type ResourceTrackerType `json:"type,omitempty"` + ApplicationGeneration int64 `json:"applicationGeneration"` + ManagedResources []ManagedResource `json:"managedResources,omitempty"` + Compression ResourceTrackerCompression `json:"compression,omitempty"` +} + +// ResourceTrackerCompression the compression for ResourceTracker ManagedResources +type ResourceTrackerCompression struct { + Type compression.Type `json:"type,omitempty"` + Data string `json:"data,omitempty"` +} + +// MarshalJSON will encode ResourceTrackerSpec according to the compression type. If type specified, +// it will encode data to compression data. +// Note: this is not the standard json Marshal process but re-use the framework function. +func (in *ResourceTrackerSpec) MarshalJSON() ([]byte, error) { + type Alias ResourceTrackerSpec + tmp := &struct{ *Alias }{} + switch in.Compression.Type { + case compression.Uncompressed: + tmp.Alias = (*Alias)(in) + case compression.Gzip: + cpy := in.DeepCopy() + data, err := compression.GzipObjectToString(in.ManagedResources) + if err != nil { + return nil, err + } + cpy.ManagedResources = nil + cpy.Compression.Data = data + tmp.Alias = (*Alias)(cpy) + default: + return nil, compression.NewUnsupportedCompressionTypeError(string(in.Compression.Type)) + } + return json.Marshal(tmp.Alias) +} + +// UnmarshalJSON will decode ResourceTrackerSpec according to the compression type. If type specified, +// it will decode data from compression data. +// Note: this is not the standard json Unmarshal process but re-use the framework function. +func (in *ResourceTrackerSpec) UnmarshalJSON(src []byte) error { + type Alias ResourceTrackerSpec + tmp := &struct{ *Alias }{} + if err := json.Unmarshal(src, tmp); err != nil { + return err + } + switch tmp.Compression.Type { + case compression.Uncompressed: + break + case compression.Gzip: + tmp.ManagedResources = []ManagedResource{} + if err := compression.GunzipStringToObject(tmp.Compression.Data, &tmp.ManagedResources); err != nil { + return err + } + tmp.Compression.Data = "" + default: + return compression.NewUnsupportedCompressionTypeError(string(in.Compression.Type)) + } + (*ResourceTrackerSpec)(tmp.Alias).DeepCopyInto(in) + return nil } // ManagedResource define the resource to be managed by ResourceTracker @@ -140,7 +197,7 @@ func (in ManagedResource) ComponentKey() string { // UnmarshalTo unmarshal ManagedResource into target object func (in ManagedResource) UnmarshalTo(obj interface{}) error { if in.Data == nil || in.Data.Raw == nil { - return errors.ManagedResourceHasNoDataError{} + return velaerr.ManagedResourceHasNoDataError{} } return json.Unmarshal(in.Data.Raw, obj) } @@ -161,7 +218,7 @@ func (in ManagedResource) ToUnstructured() *unstructured.Unstructured { func (in ManagedResource) ToUnstructuredWithData() (*unstructured.Unstructured, error) { obj := in.ToUnstructured() if err := in.UnmarshalTo(obj); err != nil { - if errors2.Is(err, errors.ManagedResourceHasNoDataError{}) { + if errors.Is(err, velaerr.ManagedResourceHasNoDataError{}) { return nil, err } } @@ -198,7 +255,7 @@ func newManagedResourceFromResource(rsc client.Object) ManagedResource { gvk := rsc.GetObjectKind().GroupVersionKind() return ManagedResource{ ClusterObjectReference: common.ClusterObjectReference{ - ObjectReference: v1.ObjectReference{ + ObjectReference: corev1.ObjectReference{ APIVersion: gvk.GroupVersion().String(), Kind: gvk.Kind, Name: rsc.GetName(), @@ -246,7 +303,7 @@ func (in *ResourceTracker) DeleteManagedResource(rsc client.Object, remove bool) gvk := rsc.GetObjectKind().GroupVersionKind() mr := ManagedResource{ ClusterObjectReference: common.ClusterObjectReference{ - ObjectReference: v1.ObjectReference{ + ObjectReference: corev1.ObjectReference{ APIVersion: gvk.GroupVersion().String(), Kind: gvk.Kind, Name: rsc.GetName(), @@ -289,7 +346,7 @@ func (in *ResourceTracker) addClusterObjectReference(ref common.ClusterObjectRef // Deprecated func (in *ResourceTracker) AddTrackedResource(rsc interfaces.TrackableResource) bool { return in.addClusterObjectReference(common.ClusterObjectReference{ - ObjectReference: v1.ObjectReference{ + ObjectReference: corev1.ObjectReference{ APIVersion: rsc.GetAPIVersion(), Kind: rsc.GetKind(), Name: rsc.GetName(), diff --git a/apis/core.oam.dev/v1beta1/resourcetracker_types_test.go b/apis/core.oam.dev/v1beta1/resourcetracker_types_test.go index a75ae2145..c4c55c716 100644 --- a/apis/core.oam.dev/v1beta1/resourcetracker_types_test.go +++ b/apis/core.oam.dev/v1beta1/resourcetracker_types_test.go @@ -18,18 +18,22 @@ package v1beta1 import ( "encoding/json" + "fmt" + "strings" "testing" + "time" "github.com/stretchr/testify/require" - v12 "k8s.io/api/apps/v1" - v1 "k8s.io/api/core/v1" - v13 "k8s.io/apimachinery/pkg/apis/meta/v1" + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime" "k8s.io/utils/pointer" "github.com/oam-dev/kubevela/apis/core.oam.dev/common" "github.com/oam-dev/kubevela/pkg/oam" + "github.com/oam-dev/kubevela/pkg/utils/compression" "github.com/oam-dev/kubevela/pkg/utils/errors" ) @@ -111,10 +115,10 @@ func TestManagedResourceKeys(t *testing.T) { input := ManagedResource{ ClusterObjectReference: common.ClusterObjectReference{ Cluster: "cluster", - ObjectReference: v1.ObjectReference{ + ObjectReference: corev1.ObjectReference{ Namespace: "namespace", Name: "name", - APIVersion: v12.SchemeGroupVersion.String(), + APIVersion: appsv1.SchemeGroupVersion.String(), Kind: "Deployment", }, }, @@ -128,7 +132,7 @@ func TestManagedResourceKeys(t *testing.T) { r.Equal("apps/Deployment/cluster/namespace/name", input.ResourceKey()) r.Equal("env/component", input.ComponentKey()) r.Equal("Deployment name (Cluster: cluster, Namespace: namespace)", input.DisplayName()) - var deploy1, deploy2 v12.Deployment + var deploy1, deploy2 appsv1.Deployment deploy1.Spec.Replicas = pointer.Int32(5) bs, err := json.Marshal(deploy1) r.NoError(err) @@ -155,13 +159,13 @@ func TestManagedResourceKeys(t *testing.T) { func TestResourceTracker_ManagedResource(t *testing.T) { r := require.New(t) input := &ResourceTracker{} - deploy1 := v12.Deployment{ObjectMeta: v13.ObjectMeta{Name: "deploy1"}} + deploy1 := appsv1.Deployment{ObjectMeta: metav1.ObjectMeta{Name: "deploy1"}} input.AddManagedResource(&deploy1, true, false, "") r.Equal(1, len(input.Spec.ManagedResources)) - cm2 := v1.ConfigMap{ObjectMeta: v13.ObjectMeta{Name: "cm2"}} + cm2 := corev1.ConfigMap{ObjectMeta: metav1.ObjectMeta{Name: "cm2"}} input.AddManagedResource(&cm2, false, false, "") r.Equal(2, len(input.Spec.ManagedResources)) - pod3 := v1.Pod{ObjectMeta: v13.ObjectMeta{Name: "pod3"}} + pod3 := corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "pod3"}} input.AddManagedResource(&pod3, false, false, "") r.Equal(3, len(input.Spec.ManagedResources)) deploy1.Spec.Replicas = pointer.Int32(5) @@ -176,9 +180,55 @@ func TestResourceTracker_ManagedResource(t *testing.T) { r.Equal(1, len(input.Spec.ManagedResources)) input.DeleteManagedResource(&pod3, true) r.Equal(0, len(input.Spec.ManagedResources)) - secret4 := v1.Secret{ObjectMeta: v13.ObjectMeta{Name: "secret4"}} + secret4 := corev1.Secret{ObjectMeta: metav1.ObjectMeta{Name: "secret4"}} input.DeleteManagedResource(&secret4, true) r.Equal(0, len(input.Spec.ManagedResources)) input.DeleteManagedResource(&secret4, false) r.Equal(1, len(input.Spec.ManagedResources)) } + +func TestResourceTrackerCompression(t *testing.T) { + size := 1000 + r := require.New(t) + rt := &ResourceTracker{} + for i := 0; i < size; i++ { + rt.AddManagedResource(&corev1.ConfigMap{ObjectMeta: metav1.ObjectMeta{Name: fmt.Sprintf("cm%d", i)}}, false, false, "") + rt.AddManagedResource(&corev1.Secret{ObjectMeta: metav1.ObjectMeta{Name: fmt.Sprintf("secret%d", i)}}, true, false, "") + } + rt.Spec.Compression.Type = compression.Gzip + t0 := time.Now() + bs, err := json.Marshal(rt) + r.NoError(err) + afterElapsed := time.Since(t0).Nanoseconds() + r.Contains(string(bs), `"type":"gzip","data":`) + _rt := &ResourceTracker{} + r.NoError(json.Unmarshal(bs, _rt)) + r.Equal(size*2, len(_rt.Spec.ManagedResources)) + for i, rsc := range _rt.Spec.ManagedResources { + r.Equal(i%2 == 1, rsc.Data == nil) + } + + _rt.Spec.Compression.Type = compression.Uncompressed + t0 = time.Now() + _bs, err := json.Marshal(_rt) + beforeElapsed := time.Since(t0) + r.NoError(err) + before, after := len(_bs), len(bs) + r.Less(after, before) + fmt.Printf("Compression Size:\n before: %d\n after: %d\n rate: %.2f%%\n", + before, after, float64(after)*100.0/float64(before)) + fmt.Printf("Compression Time:\n before: %d ns\n after: %d ns\n rate: %.2f%%\n", + beforeElapsed, afterElapsed, float64(afterElapsed)*100.0/float64(beforeElapsed)) +} + +func TestResourceTrackerInvalidMarshal(t *testing.T) { + r := require.New(t) + rt := &ResourceTracker{} + rt.Spec.Compression.Type = "invalid" + _, err := json.Marshal(rt) + r.ErrorIs(err, compression.NewUnsupportedCompressionTypeError("invalid")) + r.True(strings.Contains(err.Error(), "invalid")) + r.ErrorIs(json.Unmarshal([]byte(`{"spec":{"compression":{"type":"invalid"}}}`), rt), compression.NewUnsupportedCompressionTypeError("invalid")) + r.NotNil(json.Unmarshal([]byte(`{"spec":{"compression":{"type":"gzip","data":"xxx"}}}`), rt)) + r.NotNil(json.Unmarshal([]byte(`{"spec":["invalid"]}`), rt)) +} diff --git a/apis/core.oam.dev/v1beta1/zz_generated.deepcopy.go b/apis/core.oam.dev/v1beta1/zz_generated.deepcopy.go index 838b67634..eb3e558ad 100644 --- a/apis/core.oam.dev/v1beta1/zz_generated.deepcopy.go +++ b/apis/core.oam.dev/v1beta1/zz_generated.deepcopy.go @@ -650,6 +650,21 @@ func (in *ResourceTracker) DeepCopyObject() runtime.Object { return nil } +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *ResourceTrackerCompression) DeepCopyInto(out *ResourceTrackerCompression) { + *out = *in +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new ResourceTrackerCompression. +func (in *ResourceTrackerCompression) DeepCopy() *ResourceTrackerCompression { + if in == nil { + return nil + } + out := new(ResourceTrackerCompression) + in.DeepCopyInto(out) + return out +} + // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. func (in *ResourceTrackerList) DeepCopyInto(out *ResourceTrackerList) { *out = *in @@ -692,6 +707,7 @@ func (in *ResourceTrackerSpec) DeepCopyInto(out *ResourceTrackerSpec) { (*in)[i].DeepCopyInto(&(*out)[i]) } } + out.Compression = in.Compression } // DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new ResourceTrackerSpec. diff --git a/charts/vela-core/README.md b/charts/vela-core/README.md index 235ef8537..40a299ff5 100644 --- a/charts/vela-core/README.md +++ b/charts/vela-core/README.md @@ -94,6 +94,7 @@ helm install --create-namespace -n vela-system kubevela kubevela/vela-core --wai | `optimize.disableResourceApplyDoubleCheck` | Optimize workflow by ignoring resource double check after apply. | `false` | | `optimize.enableResourceTrackerDeleteOnlyTrigger` | Optimize resourcetracker by only trigger reconcile when resourcetracker is deleted. | `true` | | `featureGates.enableLegacyComponentRevision` | if disabled, only component with rollout trait will create component revisions | `false` | +| `featureGates.gzipResourceTracker` | if enabled, resourceTracker will be compressed before stored | `false` | ### MultiCluster parameters diff --git a/charts/vela-core/crds/core.oam.dev_resourcetrackers.yaml b/charts/vela-core/crds/core.oam.dev_resourcetrackers.yaml index 246d0dc13..5a94a3db2 100644 --- a/charts/vela-core/crds/core.oam.dev_resourcetrackers.yaml +++ b/charts/vela-core/crds/core.oam.dev_resourcetrackers.yaml @@ -56,6 +56,16 @@ spec: applicationGeneration: format: int64 type: integer + compression: + description: ResourceTrackerCompression the compression for ResourceTracker + ManagedResources + properties: + data: + type: string + type: + description: Type the compression type + type: string + type: object managedResources: items: description: ManagedResource define the resource to be managed by diff --git a/charts/vela-core/templates/kubevela-controller.yaml b/charts/vela-core/templates/kubevela-controller.yaml index 263d4b0fb..be57db998 100644 --- a/charts/vela-core/templates/kubevela-controller.yaml +++ b/charts/vela-core/templates/kubevela-controller.yaml @@ -218,6 +218,7 @@ spec: - "--feature-gates=EnableSuspendOnFailure={{- .Values.workflow.enableSuspendOnFailure | toString -}}" - "--feature-gates=AuthenticateApplication={{- .Values.authentication.enabled | toString -}}" - "--feature-gates=LegacyComponentRevision={{- .Values.featureGates.enableLegacyComponentRevision | toString -}}" + - "--feature-gates=GzipResourceTracker={{- .Values.featureGates.gzipResourceTracker | toString -}}" {{ if .Values.authentication.enabled }} {{ if .Values.authentication.withUser }} - "--authentication-with-user" diff --git a/charts/vela-core/values.yaml b/charts/vela-core/values.yaml index c88a6edaa..9b2913044 100644 --- a/charts/vela-core/values.yaml +++ b/charts/vela-core/values.yaml @@ -110,8 +110,10 @@ optimize: enableResourceTrackerDeleteOnlyTrigger: true ##@param featureGates.enableLegacyComponentRevision if disabled, only component with rollout trait will create component revisions +##@param featureGates.gzipResourceTracker if enabled, resourceTracker will be compressed before stored featureGates: enableLegacyComponentRevision: false + gzipResourceTracker: false ## @section MultiCluster parameters diff --git a/charts/vela-minimal/crds/core.oam.dev_resourcetrackers.yaml b/charts/vela-minimal/crds/core.oam.dev_resourcetrackers.yaml index 246d0dc13..5a94a3db2 100644 --- a/charts/vela-minimal/crds/core.oam.dev_resourcetrackers.yaml +++ b/charts/vela-minimal/crds/core.oam.dev_resourcetrackers.yaml @@ -56,6 +56,16 @@ spec: applicationGeneration: format: int64 type: integer + compression: + description: ResourceTrackerCompression the compression for ResourceTracker + ManagedResources + properties: + data: + type: string + type: + description: Type the compression type + type: string + type: object managedResources: items: description: ManagedResource define the resource to be managed by diff --git a/legacy/charts/vela-core-legacy/crds/core.oam.dev_resourcetrackers.yaml b/legacy/charts/vela-core-legacy/crds/core.oam.dev_resourcetrackers.yaml index 29f12678e..0ab2a80e0 100644 --- a/legacy/charts/vela-core-legacy/crds/core.oam.dev_resourcetrackers.yaml +++ b/legacy/charts/vela-core-legacy/crds/core.oam.dev_resourcetrackers.yaml @@ -56,6 +56,16 @@ spec: applicationGeneration: format: int64 type: integer + compression: + description: ResourceTrackerCompression the compression for ResourceTracker + ManagedResources + properties: + data: + type: string + type: + description: Type the compression type + type: string + type: object managedResources: items: description: ManagedResource define the resource to be managed by diff --git a/makefiles/e2e.mk b/makefiles/e2e.mk index dda4ef40a..2850d3dbf 100644 --- a/makefiles/e2e.mk +++ b/makefiles/e2e.mk @@ -17,7 +17,7 @@ e2e-setup-core-wo-auth: .PHONY: e2e-setup-core-w-auth e2e-setup-core-w-auth: - helm upgrade --install --create-namespace --namespace vela-system --set image.pullPolicy=IfNotPresent --set image.repository=vela-core-test --set applicationRevisionLimit=5 --set dependCheckWait=10s --set image.tag=$(GIT_COMMIT) --wait kubevela ./charts/vela-core --set authentication.enabled=true --set authentication.withUser=true --set authentication.groupPattern=* + helm upgrade --install --create-namespace --namespace vela-system --set image.pullPolicy=IfNotPresent --set image.repository=vela-core-test --set applicationRevisionLimit=5 --set dependCheckWait=10s --set image.tag=$(GIT_COMMIT) --wait kubevela ./charts/vela-core --set authentication.enabled=true --set authentication.withUser=true --set authentication.groupPattern=* --set featureGates.gzipResourceTracker=true .PHONY: e2e-setup-core e2e-setup-core: e2e-setup-core-pre-hook e2e-setup-core-wo-auth e2e-setup-core-post-hook diff --git a/pkg/features/controller_features.go b/pkg/features/controller_features.go index 5ffb0c6c8..e86882ab2 100644 --- a/pkg/features/controller_features.go +++ b/pkg/features/controller_features.go @@ -58,6 +58,10 @@ const ( // AuthenticateApplication enable the authentication for application AuthenticateApplication featuregate.Feature = "AuthenticateApplication" + // GzipResourceTracker enable the gzip compression for ResourceTracker. It can be useful if you have large + // application that needs to dispatch lots of resources or large resources (like CRD or huge ConfigMap), + // which at the cost of slower processing speed due to the extra overhead for compression and decompression. + GzipResourceTracker featuregate.Feature = "GzipResourceTracker" ) var defaultFeatureGates = map[featuregate.Feature]featuregate.FeatureSpec{ @@ -71,6 +75,7 @@ var defaultFeatureGates = map[featuregate.Feature]featuregate.FeatureSpec{ DisableReferObjectsFromURL: {Default: false, PreRelease: featuregate.Alpha}, ApplyResourceByUpdate: {Default: false, PreRelease: featuregate.Alpha}, AuthenticateApplication: {Default: false, PreRelease: featuregate.Alpha}, + GzipResourceTracker: {Default: false, PreRelease: featuregate.Alpha}, } func init() { diff --git a/pkg/resourcetracker/app.go b/pkg/resourcetracker/app.go index daac4ddc7..da8ee298d 100644 --- a/pkg/resourcetracker/app.go +++ b/pkg/resourcetracker/app.go @@ -27,12 +27,15 @@ import ( "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" + utilfeature "k8s.io/apiserver/pkg/util/feature" "sigs.k8s.io/controller-runtime/pkg/client" "github.com/oam-dev/kubevela/apis/core.oam.dev/common" "github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1" + "github.com/oam-dev/kubevela/pkg/features" "github.com/oam-dev/kubevela/pkg/monitor/metrics" "github.com/oam-dev/kubevela/pkg/oam" + "github.com/oam-dev/kubevela/pkg/utils/compression" velaerrors "github.com/oam-dev/kubevela/pkg/utils/errors" ) @@ -89,6 +92,9 @@ func createResourceTracker(ctx context.Context, cli client.Client, app *v1beta1. } else { rt.Spec.ApplicationGeneration = 0 } + if utilfeature.DefaultMutableFeatureGate.Enabled(features.GzipResourceTracker) { + rt.Spec.Compression.Type = compression.Gzip + } if err := cli.Create(ctx, rt); err != nil { return nil, err } diff --git a/pkg/utils/compression/error.go b/pkg/utils/compression/error.go new file mode 100644 index 000000000..36d45443b --- /dev/null +++ b/pkg/utils/compression/error.go @@ -0,0 +1,30 @@ +/* +Copyright 2022 The KubeVela Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package compression + +import "fmt" + +type unsupportedCompressionTypeError string + +func (e unsupportedCompressionTypeError) Error() string { + return fmt.Sprintf("unsupported compression type: %s", string(e)) +} + +// NewUnsupportedCompressionTypeError create a new unsupported compression type error +func NewUnsupportedCompressionTypeError(t string) error { + return unsupportedCompressionTypeError(t) +} diff --git a/pkg/utils/compression/error_test.go b/pkg/utils/compression/error_test.go new file mode 100644 index 000000000..cfc8b8177 --- /dev/null +++ b/pkg/utils/compression/error_test.go @@ -0,0 +1,27 @@ +/* +Copyright 2022 The KubeVela Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package compression + +import ( + "testing" + + "github.com/stretchr/testify/require" +) + +func TestErrors(t *testing.T) { + require.Equal(t, NewUnsupportedCompressionTypeError("x").Error(), "unsupported compression type: x") +} diff --git a/pkg/utils/compression/gzip.go b/pkg/utils/compression/gzip.go new file mode 100644 index 000000000..86bd08333 --- /dev/null +++ b/pkg/utils/compression/gzip.go @@ -0,0 +1,61 @@ +/* +Copyright 2022 The KubeVela Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package compression + +import ( + "bytes" + "compress/gzip" + "encoding/base64" + "encoding/json" + "io/ioutil" +) + +// GzipObjectToString marshal object into json, compress it with gzip, encode the result with base64 +func GzipObjectToString(obj interface{}) (string, error) { + bs, err := json.Marshal(obj) + if err != nil { + return "", err + } + var b bytes.Buffer + gz := gzip.NewWriter(&b) + if _, err = gz.Write(bs); err != nil { + return "", err + } + if err = gz.Flush(); err != nil { + return "", err + } + if err = gz.Close(); err != nil { + return "", err + } + return base64.StdEncoding.EncodeToString(b.Bytes()), nil +} + +// GunzipStringToObject decode the compressed string with base64, decompress it with gzip, unmarshal it into obj +func GunzipStringToObject(compressed string, obj interface{}) error { + bs, err := base64.StdEncoding.DecodeString(compressed) + if err != nil { + return err + } + reader, err := gzip.NewReader(bytes.NewReader(bs)) + if err != nil { + return err + } + if bs, err = ioutil.ReadAll(reader); err != nil { + return err + } + return json.Unmarshal(bs, obj) +} diff --git a/pkg/utils/compression/types.go b/pkg/utils/compression/types.go new file mode 100644 index 000000000..f49640f77 --- /dev/null +++ b/pkg/utils/compression/types.go @@ -0,0 +1,27 @@ +/* +Copyright 2022 The KubeVela Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package compression + +// Type the compression type +type Type string + +const ( + // Uncompressed . + Uncompressed Type = "" + // Gzip . + Gzip Type = "gzip" +)