Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 5 additions & 5 deletions pkg/microservice/aslan/core/common/repository/models/build.go
Original file line number Diff line number Diff line change
Expand Up @@ -99,8 +99,8 @@ type PreBuild struct {
ClusterSource string `bson:"cluster_source" json:"cluster_source"`
StrategyID string `bson:"strategy_id" json:"strategy_id"`
// UseHostDockerDaemon determines is dockerDaemon on host node is used in pod
UseHostDockerDaemon bool `bson:"use_host_docker_daemon" json:"use_host_docker_daemon"`
TemporaryStorage *TemporaryStorage `bson:"temporary_storage" json:"temporary_storage"`
UseHostDockerDaemon bool `bson:"use_host_docker_daemon" json:"use_host_docker_daemon"`
Storages *Storages `bson:"storages" json:"storages"`

CustomAnnotations []*util.KeyValue `bson:"custom_annotations" json:"custom_annotations" yaml:"custom_annotations"`
CustomLabels []*util.KeyValue `bson:"custom_labels" json:"custom_labels" yaml:"custom_labels"`
Expand All @@ -109,9 +109,9 @@ type PreBuild struct {
Namespace string `bson:"namespace" json:"namespace"`
}

type TemporaryStorage struct {
Enabled bool `bson:"enabled" json:"enabled"`
*types.NFSProperties `bson:",inline" json:",inline"`
type Storages struct {
Enabled bool `bson:"enabled" json:"enabled" yaml:"enabled"`
StoragesProperties []*types.NFSProperties `bson:"storages_properties" json:"storages_properties" yaml:"storages_properties"`
}

type PreDeploy struct {
Expand Down
27 changes: 14 additions & 13 deletions pkg/microservice/aslan/core/common/repository/models/workflow_v4.go
Original file line number Diff line number Diff line change
Expand Up @@ -1416,6 +1416,7 @@ type JobAdvancedSettings struct {
CustomLabels []*util.KeyValue `bson:"custom_labels" json:"custom_labels" yaml:"custom_labels"`
// 共享存储配置
ShareStorageInfo *ShareStorageInfo `bson:"share_storage_info" json:"share_storage_info" yaml:"share_storage_info"`
Storages *Storages `bson:"storages" json:"storages" yaml:"storages"`
}

type JobProperties struct {
Expand All @@ -1433,19 +1434,19 @@ type JobProperties struct {
Namespace string `bson:"namespace" json:"namespace" yaml:"namespace"`
Envs KeyValList `bson:"envs" json:"envs" yaml:"envs"`
// log user-defined variables, shows in workflow task detail.
CustomEnvs []*KeyVal `bson:"custom_envs" json:"custom_envs" yaml:"custom_envs,omitempty"`
Params []*Param `bson:"params" json:"params" yaml:"params"`
LogFileName string `bson:"log_file_name" json:"log_file_name" yaml:"log_file_name"`
DockerHost string `bson:"-" json:"docker_host,omitempty" yaml:"docker_host,omitempty"`
Registries []*RegistryNamespace `bson:"registries" json:"registries" yaml:"registries"`
Cache types.Cache `bson:"cache" json:"cache" yaml:"cache"`
CacheEnable bool `bson:"cache_enable" json:"cache_enable" yaml:"cache_enable"`
CacheDirType types.CacheDirType `bson:"cache_dir_type" json:"cache_dir_type" yaml:"cache_dir_type"`
CacheUserDir string `bson:"cache_user_dir" json:"cache_user_dir" yaml:"cache_user_dir"`
ShareStorageDetails []*StorageDetail `bson:"share_storage_details" json:"share_storage_details" yaml:"-"`
EnablePrivileged bool `bson:"enable_privileged,omitempty" json:"enable_privileged,omitempty" yaml:"enable_privileged,omitempty"`
UseHostDockerDaemon bool `bson:"use_host_docker_daemon,omitempty" json:"use_host_docker_daemon,omitempty" yaml:"use_host_docker_daemon"`
TemporaryStorage *types.NFSProperties `bson:"temporary_storage" json:"temporary_storage" yaml:"temporary_storage"`
CustomEnvs []*KeyVal `bson:"custom_envs" json:"custom_envs" yaml:"custom_envs,omitempty"`
Params []*Param `bson:"params" json:"params" yaml:"params"`
LogFileName string `bson:"log_file_name" json:"log_file_name" yaml:"log_file_name"`
DockerHost string `bson:"-" json:"docker_host,omitempty" yaml:"docker_host,omitempty"`
Registries []*RegistryNamespace `bson:"registries" json:"registries" yaml:"registries"`
Cache types.Cache `bson:"cache" json:"cache" yaml:"cache"`
CacheEnable bool `bson:"cache_enable" json:"cache_enable" yaml:"cache_enable"`
CacheDirType types.CacheDirType `bson:"cache_dir_type" json:"cache_dir_type" yaml:"cache_dir_type"`
CacheUserDir string `bson:"cache_user_dir" json:"cache_user_dir" yaml:"cache_user_dir"`
ShareStorageDetails []*StorageDetail `bson:"share_storage_details" json:"share_storage_details" yaml:"-"`
EnablePrivileged bool `bson:"enable_privileged,omitempty" json:"enable_privileged,omitempty" yaml:"enable_privileged,omitempty"`
UseHostDockerDaemon bool `bson:"use_host_docker_daemon,omitempty" json:"use_host_docker_daemon,omitempty" yaml:"use_host_docker_daemon"`
Storages []*types.NFSProperties `bson:"storages" json:"storages" yaml:"storages"`
// for VM deploy to get service name to save
ServiceName string `bson:"service_name" json:"service_name" yaml:"service_name"`

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ import (
"time"

"github.com/koderover/zadig/v2/pkg/tool/clientmanager"
"github.com/koderover/zadig/v2/pkg/types"
"github.com/pkg/errors"
"go.uber.org/zap"
"gopkg.in/yaml.v2"
Expand Down Expand Up @@ -200,16 +201,20 @@ func (c *FreestyleJobCtl) run(ctx context.Context) error {

c.logger.Infof("succeed to create cm for job %s", c.job.K8sJobName)

if c.jobTaskSpec.Properties.TemporaryStorage != nil {
err = service.CreateDynamicPVC(c.jobTaskSpec.Properties.ClusterID, getTemporaryStoragePVCName(c.job.K8sJobName), c.jobTaskSpec.Properties.TemporaryStorage, c.logger)
if err != nil {
msg := fmt.Sprintf("create dynamic PVC error: %v", err)
logError(c.job, msg, c.logger)
return errors.New(msg)
if len(c.jobTaskSpec.Properties.Storages) > 0 {
for i, storage := range c.jobTaskSpec.Properties.Storages {
if storage.ProvisionType == types.DynamicProvision {
err = service.CreateDynamicPVC(c.jobTaskSpec.Properties.ClusterID, getStoragePVCName(c.job.K8sJobName, i), storage, c.logger)
if err != nil {
msg := fmt.Sprintf("create dynamic PVC error: %v", err)
logError(c.job, msg, c.logger)
return errors.New(msg)

}
}

c.logger.Infof("succeed to create dynamic PVC for job %s", c.job.K8sJobName)
c.logger.Infof("succeed to create dynamic PVC for job %s", c.job.K8sJobName)
}
}
}

jobImage := getBaseImage(c.jobTaskSpec.Properties.BuildOS, c.jobTaskSpec.Properties.ImageFrom)
Expand Down Expand Up @@ -357,9 +362,13 @@ func (c *FreestyleJobCtl) complete(ctx context.Context) {
// 清理用户取消和超时的任务
defer func() {
go func() {
if c.jobTaskSpec.Properties.TemporaryStorage != nil {
if err := ensureDeletePVC(c.job.K8sJobName, c.jobTaskSpec.Properties.Namespace, c.jobTaskSpec.Properties.TemporaryStorage, c.kubeclient); err != nil {
c.logger.Error(err)
if len(c.jobTaskSpec.Properties.Storages) > 0 {
for _, storage := range c.jobTaskSpec.Properties.Storages {
if storage.IsTemporary {
if err := ensureDeletePVC(storage.PVC, c.jobTaskSpec.Properties.Namespace, storage, c.kubeclient); err != nil {
c.logger.Error(err)
}
}
}
}
if err := ensureDeleteJob(c.jobTaskSpec.Properties.Namespace, jobLabel, c.kubeclient); err != nil {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -112,12 +112,11 @@ type JobLabel struct {
JobType string
}

func getTemporaryStoragePVCName(k8sJobName string) string {
return fmt.Sprintf("%s-temporary", k8sJobName)
func getStoragePVCName(k8sJobName string, index int) string {
return fmt.Sprintf("%s-%d", k8sJobName, index)
}

func ensureDeletePVC(jobName, namespace string, storage *types.NFSProperties, kubeClient crClient.Client) error {
pvcName := service.GetPVCName(getTemporaryStoragePVCName(jobName), storage)
func ensureDeletePVC(pvcName, namespace string, storage *types.NFSProperties, kubeClient crClient.Client) error {
return kubeClient.Delete(context.TODO(), &corev1.PersistentVolumeClaim{
ObjectMeta: metav1.ObjectMeta{
Name: pvcName,
Expand Down Expand Up @@ -481,7 +480,7 @@ EOF`,
},
}

setJobTemporaryStorages(job, workflowCtx, jobTaskSpec.Properties.TemporaryStorage, targetCluster)
setJobStorages(job, workflowCtx, jobTaskSpec.Properties.Storages, targetCluster)
setJobShareStorages(job, workflowCtx, jobTaskSpec.Properties.ShareStorageDetails, targetCluster)

if jobTaskSpec.Properties.CacheEnable && jobTaskSpec.Properties.Cache.MediumType == commontypes.NFSMedium {
Expand Down Expand Up @@ -569,29 +568,32 @@ func BuildCleanJob(jobName, clusterID, workflowName string, taskID int64) (*batc
return job, nil
}

func setJobTemporaryStorages(job *batchv1.Job, workflowCtx *commonmodels.WorkflowTaskCtx, temporaryStorage *types.NFSProperties, cluster *commonmodels.K8SCluster) {
if temporaryStorage == nil {
func setJobStorages(job *batchv1.Job, workflowCtx *commonmodels.WorkflowTaskCtx, storages []*types.NFSProperties, cluster *commonmodels.K8SCluster) {
if len(storages) <= 0 {
return
}

// save cluster id so we can clean up share storage later
workflowCtx.ClusterIDAdd(cluster.ID.Hex())

volumeName := "temporary-storage"
job.Spec.Template.Spec.Volumes = append(job.Spec.Template.Spec.Volumes, corev1.Volume{
Name: volumeName,
VolumeSource: corev1.VolumeSource{
PersistentVolumeClaim: &corev1.PersistentVolumeClaimVolumeSource{
ClaimName: temporaryStorage.PVC,
for _, storage := range storages {
volumeName := fmt.Sprintf("storage-%s", storage.PVC)
job.Spec.Template.Spec.Volumes = append(job.Spec.Template.Spec.Volumes, corev1.Volume{
Name: volumeName,
VolumeSource: corev1.VolumeSource{
PersistentVolumeClaim: &corev1.PersistentVolumeClaimVolumeSource{
ClaimName: storage.PVC,
},
},
},
})
})

job.Spec.Template.Spec.Containers[0].VolumeMounts = append(job.Spec.Template.Spec.Containers[0].VolumeMounts, corev1.VolumeMount{
Name: volumeName,
MountPath: storage.MountPath,
SubPath: storage.Subpath,
})
}

job.Spec.Template.Spec.Containers[0].VolumeMounts = append(job.Spec.Template.Spec.Containers[0].VolumeMounts, corev1.VolumeMount{
Name: volumeName,
MountPath: "/workspace",
SubPath: fmt.Sprintf("%s/%s", job.Name, "temporary-storage"),
})
}

func setJobShareStorages(job *batchv1.Job, workflowCtx *commonmodels.WorkflowTaskCtx, storageDetails []*commonmodels.StorageDetail, cluster *commonmodels.K8SCluster) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1402,7 +1402,7 @@ func CreateDynamicPVC(clusterID, prefix string, nfsProperties *types.NFSProperti

accessMode := []corev1.PersistentVolumeAccessMode{corev1.ReadWriteMany}
if nfsProperties.AccessMode != "" {
accessMode = []corev1.PersistentVolumeAccessMode{corev1.PersistentVolumeAccessMode(nfsProperties.AccessMode)}
accessMode = []corev1.PersistentVolumeAccessMode{nfsProperties.AccessMode}
}

pvc = &corev1.PersistentVolumeClaim{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,6 @@ import (
"strings"

"go.uber.org/zap"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/util/sets"

configbase "github.com/koderover/zadig/v2/pkg/config"
Expand All @@ -41,6 +40,7 @@ import (
"github.com/koderover/zadig/v2/pkg/types"
"github.com/koderover/zadig/v2/pkg/types/job"
"github.com/koderover/zadig/v2/pkg/types/step"
pkgutil "github.com/koderover/zadig/v2/pkg/util"
)

// TODO: Change note: ServiceAndBuilds field use to be the option field for the configuration, it has been
Expand Down Expand Up @@ -400,11 +400,6 @@ func (j BuildJobController) ToTask(taskID int64) ([]*commonmodels.JobTask, error
EnablePrivileged: buildInfo.EnablePrivilegedMode,
}

if buildInfo.PreBuild != nil && buildInfo.PreBuild.TemporaryStorage != nil && buildInfo.PreBuild.TemporaryStorage.Enabled {
jobTaskSpec.Properties.TemporaryStorage = buildInfo.PreBuild.TemporaryStorage.NFSProperties
jobTaskSpec.Properties.TemporaryStorage.AccessMode = string(corev1.ReadWriteOnce)
}

paramEnvs := generateKeyValsFromWorkflowParam(j.workflow.Params)
envs := mergeKeyVals(jobTaskSpec.Properties.CustomEnvs, paramEnvs)
renderedEnv, err := replaceServiceAndModules(envs, build.ServiceName, build.ServiceModule)
Expand All @@ -415,6 +410,25 @@ func (j BuildJobController) ToTask(taskID int64) ([]*commonmodels.JobTask, error
jobTaskSpec.Properties.Envs = append(renderedEnv, getBuildJobVariables(build, taskID, j.workflow.Project, j.workflow.Name, j.workflow.DisplayName, image, pkgFile, jobTask.Infrastructure, registry, logger)...)
jobTaskSpec.Properties.UseHostDockerDaemon = buildInfo.PreBuild.UseHostDockerDaemon

if buildInfo.PreBuild != nil && buildInfo.PreBuild.Storages != nil && buildInfo.PreBuild.Storages.Enabled {
if len(buildInfo.PreBuild.Storages.StoragesProperties) > 0 {
newStorages := make([]*types.NFSProperties, 0)
for _, storage := range buildInfo.PreBuild.Storages.StoragesProperties {
newStorage := &types.NFSProperties{}
err = pkgutil.DeepCopy(newStorage, storage)
if err != nil {
return nil, fmt.Errorf("failed to deep copy storage: %v", err)
}

newStorage.MountPath = commonutil.RenderEnv(storage.MountPath, jobTaskSpec.Properties.Envs)
newStorage.Subpath = commonutil.RenderEnv(storage.Subpath, jobTaskSpec.Properties.Envs)
newStorages = append(newStorages, newStorage)
}

jobTaskSpec.Properties.Storages = newStorages
}
}

cacheS3 := &commonmodels.S3Storage{}
if jobTask.Infrastructure == setting.JobVMInfrastructure {
jobTaskSpec.Properties.CacheEnable = buildInfo.CacheEnable
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ import (
"github.com/koderover/zadig/v2/pkg/tool/log"
"github.com/koderover/zadig/v2/pkg/types"
steptypes "github.com/koderover/zadig/v2/pkg/types/step"
pkgutil "github.com/koderover/zadig/v2/pkg/util"
util2 "github.com/koderover/zadig/v2/pkg/util"
)

Expand Down Expand Up @@ -571,6 +572,25 @@ func (j FreestyleJobController) generateSubTask(taskID int64, jobSubTaskID int,
taskRunProperties.Envs = renderedEnvs
}

if j.jobSpec.AdvancedSetting.Storages != nil && j.jobSpec.AdvancedSetting.Storages.Enabled {
if len(j.jobSpec.AdvancedSetting.Storages.StoragesProperties) > 0 {
newStorages := make([]*types.NFSProperties, 0)
for _, storage := range j.jobSpec.AdvancedSetting.Storages.StoragesProperties {
newStorage := &types.NFSProperties{}
err = pkgutil.DeepCopy(newStorage, storage)
if err != nil {
return nil, fmt.Errorf("failed to deep copy storage: %v", err)
}

newStorage.MountPath = commonutil.RenderEnv(storage.MountPath, taskRunProperties.Envs)
newStorage.Subpath = commonutil.RenderEnv(storage.Subpath, taskRunProperties.Envs)
newStorages = append(newStorages, newStorage)
}

taskRunProperties.Storages = newStorages
}
}

repos := j.jobSpec.Repos
if service != nil {
repos = service.Repos
Expand Down
18 changes: 12 additions & 6 deletions pkg/types/cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,10 @@ limitations under the License.

package types

import (
corev1 "k8s.io/api/core/v1"
)

type MediumType string

const (
Expand All @@ -40,12 +44,14 @@ type ObjectProperties struct {
}

type NFSProperties struct {
ProvisionType ProvisionType `json:"provision_type" bson:"provision_type" yaml:"provision_type"`
StorageClass string `json:"storage_class" bson:"storage_class" yaml:"storage_class"`
StorageSizeInGiB int64 `json:"storage_size_in_gib" bson:"storage_size_in_gib" yaml:"storage_size_in_gib"`
PVC string `json:"pvc" bson:"pvc" yaml:"pvc"`
Subpath string `json:"subpath" bson:"subpath" yaml:"subpath"`
AccessMode string `json:"access_mode" bson:"access_mode" yaml:"access_mode"`
ProvisionType ProvisionType `json:"provision_type" bson:"provision_type" yaml:"provision_type"`
StorageClass string `json:"storage_class" bson:"storage_class" yaml:"storage_class"`
StorageSizeInGiB int64 `json:"storage_size_in_gib" bson:"storage_size_in_gib" yaml:"storage_size_in_gib"`
PVC string `json:"pvc" bson:"pvc" yaml:"pvc"`
AccessMode corev1.PersistentVolumeAccessMode `json:"access_mode" bson:"access_mode" yaml:"access_mode"`
Subpath string `json:"sub_path" bson:"sub_path" yaml:"sub_path"`
MountPath string `json:"mount_path" bson:"mount_path" yaml:"mount_path"`
IsTemporary bool `json:"is_temporary" bson:"is_temporary" yaml:"is_temporary"`
}

type Cache struct {
Expand Down