|
| 1 | +package kubernetes |
| 2 | + |
| 3 | +import ( |
| 4 | + "context" |
| 5 | + "fmt" |
| 6 | + "os" |
| 7 | + "path/filepath" |
| 8 | + "strings" |
| 9 | + |
| 10 | + k8serrors "k8s.io/apimachinery/pkg/api/errors" |
| 11 | + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" |
| 12 | + "k8s.io/apimachinery/pkg/runtime/serializer/yaml" |
| 13 | + "sigs.k8s.io/controller-runtime/pkg/client" |
| 14 | + |
| 15 | + "github.com/dapr/cli/pkg/print" |
| 16 | +) |
| 17 | + |
| 18 | +func getResources(resourcesFolder string) ([]client.Object, error) { |
| 19 | + // Create YAML decoder |
| 20 | + decUnstructured := yaml.NewDecodingSerializer(unstructured.UnstructuredJSONScheme) |
| 21 | + |
| 22 | + // Read files from the resources folder |
| 23 | + files, err := os.ReadDir(resourcesFolder) |
| 24 | + if err != nil { |
| 25 | + return nil, fmt.Errorf("error reading resources folder: %w", err) |
| 26 | + } |
| 27 | + |
| 28 | + var resources []client.Object |
| 29 | + for _, file := range files { |
| 30 | + if file.IsDir() || (!strings.HasSuffix(file.Name(), ".yaml") && !strings.HasSuffix(file.Name(), ".json")) { |
| 31 | + continue |
| 32 | + } |
| 33 | + |
| 34 | + // Read file content |
| 35 | + content, err := os.ReadFile(filepath.Join(resourcesFolder, file.Name())) |
| 36 | + if err != nil { |
| 37 | + return nil, fmt.Errorf("error reading file %s: %w", file.Name(), err) |
| 38 | + } |
| 39 | + |
| 40 | + // Decode YAML/JSON to unstructured |
| 41 | + obj := &unstructured.Unstructured{} |
| 42 | + _, _, err = decUnstructured.Decode(content, nil, obj) |
| 43 | + if err != nil { |
| 44 | + return nil, fmt.Errorf("error decoding file %s: %w", file.Name(), err) |
| 45 | + } |
| 46 | + |
| 47 | + resources = append(resources, obj) |
| 48 | + } |
| 49 | + |
| 50 | + return resources, nil |
| 51 | +} |
| 52 | + |
| 53 | +func createOrUpdateResources(ctx context.Context, cl client.Client, resources []client.Object, namespace string) error { |
| 54 | + // create resources in k8s |
| 55 | + for _, resource := range resources { |
| 56 | + // clone the resource to avoid modifying the original |
| 57 | + obj := resource.DeepCopyObject().(*unstructured.Unstructured) |
| 58 | + // Set namespace on the resource metadata |
| 59 | + obj.SetNamespace(namespace) |
| 60 | + |
| 61 | + print.InfoStatusEvent(os.Stdout, "Deploying resource %q kind %q to Kubernetes", obj.GetName(), obj.GetKind()) |
| 62 | + |
| 63 | + if err := cl.Create(ctx, obj); err != nil { |
| 64 | + if k8serrors.IsAlreadyExists(err) { |
| 65 | + print.InfoStatusEvent(os.Stdout, "Resource %q kind %q already exists, updating", obj.GetName(), obj.GetKind()) |
| 66 | + if err := cl.Update(ctx, obj); err != nil { |
| 67 | + return err |
| 68 | + } |
| 69 | + } else { |
| 70 | + return fmt.Errorf("error deploying resource %q kind %q to Kubernetes: %w", obj.GetName(), obj.GetKind(), err) |
| 71 | + } |
| 72 | + } |
| 73 | + } |
| 74 | + return nil |
| 75 | +} |
0 commit comments