From 1933b1036010d33e051f9c92f55fc1b29c8226bb Mon Sep 17 00:00:00 2001 From: Ashraf Fouda Date: Wed, 23 Sep 2026 17:10:17 +0300 Subject: [PATCH] fixes provisiond hanging specially with zoslight nodes Signed-off-by: Ashraf Fouda --- pkg/kernel/kernel.go | 8 +- pkg/primitives/provisioner.go | 39 ++++- pkg/provision/engine.go | 293 ++++++++++++++++++++++++++-------- 3 files changed, 260 insertions(+), 80 deletions(-) diff --git a/pkg/kernel/kernel.go b/pkg/kernel/kernel.go index 1e1131fa..dd897f18 100644 --- a/pkg/kernel/kernel.go +++ b/pkg/kernel/kernel.go @@ -72,25 +72,25 @@ func (k Params) IsLight() bool { // ISV3Light check if node is running in v3light mode func (k Params) IsV3Light() bool { - version, found := GetParams().GetOne("version") + version, found := k.GetOne("version") return found && version == "v3light" } // getVersion checks if node is running version 4 func (k Params) GetVersion() string { - version, _ := GetParams().GetOne("version") + version, _ := k.GetOne("version") return version } // IsV4 checks if node is running version 4 func (k Params) IsV4() bool { - version, found := GetParams().GetOne("version") + version, found := k.GetOne("version") return found && version == "v4" } // IsV3 checks if node is running normal version 3 func (k Params) IsV3() bool { - version, found := GetParams().GetOne("version") + version, found := k.GetOne("version") return found && version == "v3" } diff --git a/pkg/primitives/provisioner.go b/pkg/primitives/provisioner.go index a549f276..e4b6923b 100644 --- a/pkg/primitives/provisioner.go +++ b/pkg/primitives/provisioner.go @@ -4,6 +4,7 @@ import ( "github.com/threefoldtech/zbus" "github.com/threefoldtech/zos_base/pkg/gridtypes" "github.com/threefoldtech/zos_base/pkg/gridtypes/zos" + "github.com/threefoldtech/zos_base/pkg/kernel" "github.com/threefoldtech/zos_base/pkg/primitives/gateway" "github.com/threefoldtech/zos_base/pkg/primitives/network" netlight "github.com/threefoldtech/zos_base/pkg/primitives/network-light" @@ -19,22 +20,46 @@ import ( ) // NewPrimitivesProvisioner creates a new 0-OS provisioner +// +// Only the workload types this node variant can actually serve are registered. +// A manager reaches its worker over zbus, and a light node runs netlightd +// (module "netlight") in place of networkd (module "network"), so registering +// the full network backed types there would leave a request queued for a module +// that never answers. The provisioner rejects an unregistered type outright, +// which surfaces to the owner as a failed workload instead of a deployment that +// never resolves. Keep this in step with systemMonitor.GetNodeFeatures, which +// advertises the same split to the grid. func NewPrimitivesProvisioner(zbus zbus.Client) provision.Provisioner { + return provision.NewMapProvisioner(managersFor(zbus, kernel.GetParams())) +} + +// managersFor builds the set of workload managers the given node variant can +// actually serve. Split out from NewPrimitivesProvisioner so the variant gating +// is testable without a matching /proc/cmdline. +func managersFor(zbus zbus.Client, params kernel.Params) map[gridtypes.WorkloadType]provision.Manager { + // types every variant serves. their managers talk to storaged, flistd, + // contd and vmd, which run everywhere. managers := map[gridtypes.WorkloadType]provision.Manager{ zos.ZMountType: zmount.NewManager(zbus), zos.ZLogsType: zlogs.NewManager(zbus), zos.QuantumSafeFSType: qsfs.NewManager(zbus), zos.ZDBType: zdb.NewManager(zbus), - zos.NetworkType: network.NewManager(zbus), - zos.PublicIPType: pubip.NewManager(zbus), - zos.PublicIPv4Type: pubip.NewManager(zbus), // backward compatibility - zos.ZMachineType: vm.NewManager(zbus), - zos.NetworkLightType: netlight.NewManager(zbus), - zos.ZMachineLightType: vmlight.NewManager(zbus), zos.VolumeType: volume.NewManager(zbus), zos.GatewayNameProxyType: gateway.NewNameManager(zbus), zos.GatewayFQDNProxyType: gateway.NewFQDNManager(zbus), + zos.PublicIPType: pubip.NewManager(zbus), + zos.PublicIPv4Type: pubip.NewManager(zbus), // backward compatibility + } + + if params.IsLight() { + // served by netlightd, module "netlight". networkd does not run here. + managers[zos.NetworkLightType] = netlight.NewManager(zbus) + managers[zos.ZMachineLightType] = vmlight.NewManager(zbus) + } else { + // served by networkd, module "network". netlightd does not run here. + managers[zos.NetworkType] = network.NewManager(zbus) + managers[zos.ZMachineType] = vm.NewManager(zbus) } - return provision.NewMapProvisioner(managers) + return managers } diff --git a/pkg/provision/engine.go b/pkg/provision/engine.go index e0eeaf3f..42e9e27d 100644 --- a/pkg/provision/engine.go +++ b/pkg/provision/engine.go @@ -10,6 +10,7 @@ import ( "os" "path/filepath" "sort" + "sync" "time" "github.com/cenkalti/backoff/v3" @@ -102,6 +103,16 @@ const ( opStart // servers default timeout defaultHttpTimeout = 10 * time.Second + + // jobTimeout is the maximum time a single job is allowed to occupy the + // engine loop. The loop is strictly serial, so a job that never returns + // stalls provisioning for the entire node. Most of the work a job does + // is zbus calls, which poll for a reply and honor context cancellation, + // so a deadline on the job context is enough to break out of a reply + // that is never coming (e.g. a workload type whose module does not run + // on this node variant). It is deliberately generous: it is a stall + // breaker, not an SLA on deployment time. + jobTimeout = 30 * time.Minute ) // engineJob is a persisted job instance that is @@ -116,6 +127,45 @@ type engineJob struct { Message string } +func (o jobOperation) String() string { + switch o { + case opProvision: + return "provision" + case opDeprovision: + return "deprovision" + case opUpdate: + return "update" + case opProvisionNoValidation: + return "provision-no-validation" + case opPause: + return "pause" + case opResume: + return "resume" + case opPrepare: + return "prepare" + case opStart: + return "start" + default: + return fmt.Sprintf("unknown(%d)", int(o)) + } +} + +// dedupable reports whether a repeated enqueue of this operation for the same +// deployment is redundant. These ops carry no payload beyond the identity of +// the deployment, so a second copy queued while the first is still waiting has +// nothing new to do. Payload carrying ops (provision, update, prepare, start) +// are never deduplicated. +func (o jobOperation) dedupable() bool { + return o == opDeprovision || o == opPause || o == opResume +} + +// jobKey identifies a queued job for deduplication purposes. +type jobKey struct { + twin uint32 + contract uint64 + op jobOperation +} + // NativeEngine is the core of this package // The engine is responsible to manage provision and decomission of workloads on the system type NativeEngine struct { @@ -124,6 +174,13 @@ type NativeEngine struct { queue *dque.DQue + // pending tracks dedupable jobs that are currently sitting in the queue, + // so the same work is not queued twice. It is in memory only: the queue + // survives restarts, this does not, which at worst allows the duplicates + // that were possible before. + pending map[jobKey]struct{} + pendingMu sync.Mutex + // options // janitor Janitor twins Twins @@ -335,6 +392,7 @@ func New(storage Storage, provisioner Provisioner, root string, opts ...EngineOp } e.queue = queue + e.pending = make(map[jobKey]struct{}) return e, nil } @@ -373,7 +431,7 @@ func (e *NativeEngine) Provision(ctx context.Context, deployment gridtypes.Deplo Op: opProvision, } - return e.queue.Enqueue(&job) + return e.enqueue(&job) } // Prepare stages a deployment on this node: it persists the whole deployment @@ -400,7 +458,7 @@ func (e *NativeEngine) Prepare(ctx context.Context, deployment gridtypes.Deploym Op: opPrepare, } - return e.queue.Enqueue(&job) + return e.enqueue(&job) } // Pause deployment @@ -420,7 +478,7 @@ func (e *NativeEngine) Pause(ctx context.Context, twin uint32, id uint64) error Op: opPause, } - return e.queue.Enqueue(&job) + return e.enqueue(&job) } // Resume deployment @@ -440,7 +498,7 @@ func (e *NativeEngine) Resume(ctx context.Context, twin uint32, id uint64) error Op: opResume, } - return e.queue.Enqueue(&job) + return e.enqueue(&job) } // Start boots the zmachine(s) of a previously prepared deployment (see opPrepare). @@ -462,7 +520,7 @@ func (e *NativeEngine) Start(ctx context.Context, twin uint32, id uint64) error Op: opStart, } - return e.queue.Enqueue(&job) + return e.enqueue(&job) } // Deprovision workload @@ -484,7 +542,7 @@ func (e *NativeEngine) Deprovision(ctx context.Context, twin uint32, id uint64, Message: reason, } - return e.queue.Enqueue(&job) + return e.enqueue(&job) } // Update workloads @@ -536,7 +594,61 @@ func (e *NativeEngine) Update(ctx context.Context, update gridtypes.Deployment) Source: &deployment, } - return e.queue.Enqueue(&job) + return e.enqueue(&job) +} + +// enqueue adds a job to the persisted queue, dropping it when an equivalent +// job is already waiting to be picked up. +// +// The periodic chain reconciler re-issues deprovision for the same contracts on +// every pass. Without this, any slowdown at the head of the queue lets those +// repeats accumulate without bound, and the backlog outlives the original +// problem. +func (e *NativeEngine) enqueue(job *engineJob) error { + if !job.Op.dedupable() { + return e.queue.Enqueue(job) + } + + key := jobKey{ + twin: job.Target.TwinID, + contract: job.Target.ContractID, + op: job.Op, + } + + e.pendingMu.Lock() + if _, ok := e.pending[key]; ok { + e.pendingMu.Unlock() + log.Debug(). + Uint32("twin", key.twin). + Uint64("contract", key.contract). + Str("op", job.Op.String()). + Msg("an equivalent job is already queued, skipping duplicate") + return nil + } + e.pending[key] = struct{}{} + e.pendingMu.Unlock() + + if err := e.queue.Enqueue(job); err != nil { + e.forgetPending(job) + return err + } + + return nil +} + +// forgetPending drops the deduplication marker for a job that has left the queue. +func (e *NativeEngine) forgetPending(job *engineJob) { + if !job.Op.dedupable() { + return + } + + e.pendingMu.Lock() + defer e.pendingMu.Unlock() + delete(e.pending, jobKey{ + twin: job.Target.TwinID, + contract: job.Target.ContractID, + op: job.Op, + }) } // Run starts reader reservation from the Source and handle them @@ -560,77 +672,120 @@ func (e *NativeEngine) Run(root context.Context) error { } job := obj.(*engineJob) - ctx := withDeployment(root, job.Target.TwinID, job.Target.ContractID) - l := log.With(). - Uint32("twin", job.Target.TwinID). - Uint64("contract", job.Target.ContractID). - Logger() - - // contract validation - // this should ONLY be done on provosion and update operation - if job.Op == opProvision || - job.Op == opUpdate || - job.Op == opPrepare || - job.Op == opProvisionNoValidation { - // otherwise, contract validation is needed - ctx, err = e.validate(ctx, &job.Target, job.Op == opProvisionNoValidation) - if err != nil { - l.Error().Err(err).Msg("contact validation fails") - // job.Target.SetError(err) - if err := e.storage.Error(job.Target.TwinID, job.Target.ContractID, err); err != nil { - l.Error().Err(err).Msg("failed to set deployment global error") - } - _, _ = e.queue.Dequeue() + callback := e.runJob(root, job) + + // the job is always dequeued, even if it failed or timed out. leaving + // it at the head would stall every other deployment on the node. + if _, err := e.queue.Dequeue(); err != nil { + log.Error().Err(err). + Uint32("twin", job.Target.TwinID). + Uint64("contract", job.Target.ContractID). + Msg("failed to dequeue job") + } - continue - } + e.forgetPending(job) - l.Debug().Msg("contact validation pass") + if callback { + e.safeCallback(&job.Target, job.Op == opDeprovision) } + } +} - switch job.Op { - case opProvisionNoValidation: - fallthrough - case opProvision: - e.installDeployment(ctx, &job.Target) - case opPrepare: - e.prepareDeployment(ctx, &job.Target) - case opStart: - e.startDeployment(ctx, &job.Target) - case opDeprovision: - e.uninstallDeployment(ctx, &job.Target, job.Message) - case opPause: - e.lockDeployment(ctx, &job.Target) - case opResume: - e.unlockDeployment(ctx, &job.Target) - case opUpdate: - // update is tricky because we need to work against - // 2 versions of the object. Once that reflects the current state - // and the new one that is the target state but it does not know - // the current state of already deployed workloads - // so (1st) we need to get the difference - // this call will return 3 lists - // - things to remove - // - things to add - // - things to update (not supported atm) - // - things that is not in any of the 3 lists are basically stay as is - // the call will also make sure the Result of those workload in both the (did not change) - // and update to reflect the current result on those workloads. - update, err := job.Source.Upgrade(&job.Target) - if err != nil { - l.Error().Err(err).Msg("failed to get update procedure") - break +// runJob executes a single queued job and reports whether the deployment +// callback should run for it. +// +// The engine loop is strictly serial, so this must always return: a job that +// blocks forever stalls provisioning for the whole node. The job context +// therefore carries jobTimeout. Nearly all the work below is zbus calls, which +// poll for a reply and honor context cancellation, so the deadline is enough +// to break out of a reply that is never coming, e.g. a workload whose manager +// talks to a module that does not run on this node variant. +func (e *NativeEngine) runJob(root context.Context, job *engineJob) (callback bool) { + ctx, cancel := context.WithTimeout( + withDeployment(root, job.Target.TwinID, job.Target.ContractID), + jobTimeout, + ) + defer cancel() + + l := log.With(). + Uint32("twin", job.Target.TwinID). + Uint64("contract", job.Target.ContractID). + Logger() + + // contract validation + // this should ONLY be done on provosion and update operation + if job.Op == opProvision || + job.Op == opUpdate || + job.Op == opPrepare || + job.Op == opProvisionNoValidation { + // otherwise, contract validation is needed + validated, err := e.validate(ctx, &job.Target, job.Op == opProvisionNoValidation) + if err != nil { + l.Error().Err(err).Msg("contact validation fails") + // job.Target.SetError(err) + if err := e.storage.Error(job.Target.TwinID, job.Target.ContractID, err); err != nil { + l.Error().Err(err).Msg("failed to set deployment global error") } - e.updateDeployment(ctx, update) - } - _, err = e.queue.Dequeue() + return false + } + ctx = validated + + l.Debug().Msg("contact validation pass") + } + + switch job.Op { + case opProvisionNoValidation: + fallthrough + case opProvision: + e.installDeployment(ctx, &job.Target) + case opPrepare: + e.prepareDeployment(ctx, &job.Target) + case opStart: + e.startDeployment(ctx, &job.Target) + case opDeprovision: + e.uninstallDeployment(ctx, &job.Target, job.Message) + case opPause: + e.lockDeployment(ctx, &job.Target) + case opResume: + e.unlockDeployment(ctx, &job.Target) + case opUpdate: + // update is tricky because we need to work against + // 2 versions of the object. Once that reflects the current state + // and the new one that is the target state but it does not know + // the current state of already deployed workloads + // so (1st) we need to get the difference + // this call will return 3 lists + // - things to remove + // - things to add + // - things to update (not supported atm) + // - things that is not in any of the 3 lists are basically stay as is + // the call will also make sure the Result of those workload in both the (did not change) + // and update to reflect the current result on those workloads. + update, err := job.Source.Upgrade(&job.Target) if err != nil { - l.Error().Err(err).Msg("failed to dequeue job") + l.Error().Err(err).Msg("failed to get update procedure") + break } + e.updateDeployment(ctx, update) + } - e.safeCallback(&job.Target, job.Op == opDeprovision) + // a job that ran out of time left the deployment in whatever state it + // reached. record that on the deployment so the owner can see why, instead + // of the node going quiet. + if err := ctx.Err(); errors.Is(err, context.DeadlineExceeded) { + l.Error(). + Str("op", job.Op.String()). + Dur("timeout", jobTimeout). + Msg("job did not finish in time, dropping it to keep the queue draining") + + err = fmt.Errorf("provisioning did not finish within %s", jobTimeout) + if err := e.storage.Error(job.Target.TwinID, job.Target.ContractID, err); err != nil { + l.Error().Err(err).Msg("failed to set deployment global error") + } } + + return true } func (e *NativeEngine) safeCallback(d *gridtypes.Deployment, delete bool) { @@ -724,7 +879,7 @@ func (e *NativeEngine) boot(root context.Context) error { Op: opProvisionNoValidation, } - if err := e.queue.Enqueue(&job); err != nil { + if err := e.enqueue(&job); err != nil { log.Error(). Err(err). Uint32("twin", dl.TwinID).