- Add jarRef and jarPullSecret fields to FlinkJob CRD (jarUri/basicAuth deprecated) - OCI pull via go-containerregistry with dual auth (K8s pull secret + env vars) - Media type validation on pulled layers - Atomic status patches with runningJarRef/runningJarDigest tracking - NeedsUpgrade/RunningRef/RunningRefPatchData domain helpers - README with usage guide, pushing JARs, and GitHub Actions CI/CD workflow - CONTEXT.md domain glossary
41 lines
981 B
Go
41 lines
981 B
Go
package managed_job
|
|
|
|
import (
|
|
"flink-kube-operator/pkg"
|
|
|
|
"go.uber.org/zap"
|
|
)
|
|
|
|
func (job *ManagedJob) upgrade() {
|
|
effectiveRef := job.def.Spec.EffectiveJarRef()
|
|
prevRef := job.def.Status.RunningRef()
|
|
|
|
pkg.Logger.Info("[managed-job] [upgrade] pausing...",
|
|
zap.String("jobName", job.def.GetName()),
|
|
zap.String("currentRef", effectiveRef),
|
|
zap.String("prevRef", prevRef),
|
|
)
|
|
job.def.Status.JarId = nil
|
|
job.crd.Patch(job.def.UID, map[string]interface{}{
|
|
"status": map[string]interface{}{
|
|
"jarId": job.def.Status.JarId,
|
|
},
|
|
})
|
|
err := job.Pause()
|
|
if err != nil {
|
|
pkg.Logger.Error("[managed-job] [upgrade] error in pausing", zap.Error(err))
|
|
return
|
|
}
|
|
pkg.Logger.Info("[managed-job] [upgrade] restoring...",
|
|
zap.String("jobName", job.def.GetName()),
|
|
zap.String("currentRef", effectiveRef),
|
|
zap.String("prevRef", prevRef),
|
|
)
|
|
|
|
err = job.Run(true)
|
|
if err != nil {
|
|
pkg.Logger.Error("[managed-job] [upgrade] error in running", zap.Error(err))
|
|
return
|
|
}
|
|
}
|