Skip to content

Commit acb0b1e

Browse files
committed
polish
1 parent 4e62678 commit acb0b1e

2 files changed

Lines changed: 51 additions & 72 deletions

File tree

pkg/controller/launcher-populator/digest-updater.go

Lines changed: 35 additions & 51 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import (
2222

2323
fmav1alpha1 "github.com/llm-d-incubation/llm-d-fast-model-actuation/api/fma/v1alpha1"
2424
"github.com/llm-d-incubation/llm-d-fast-model-actuation/pkg/controller/utils"
25+
corev1 "k8s.io/api/core/v1"
2526
apierrors "k8s.io/apimachinery/pkg/api/errors"
2627
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
2728
"k8s.io/apimachinery/pkg/labels"
@@ -50,23 +51,19 @@ func (ctl *controller) updateDigestForLC(ctx context.Context, name string) error
5051
if !apierrors.IsNotFound(err) {
5152
return fmt.Errorf("failed to get LC %s: %w", name, err)
5253
}
53-
// LC deleted.
54+
if !prevExists {
55+
return nil
56+
}
5457
logger.Info("LC deleted, marking referencing keys as handsOff", "config", name)
5558
delete(ctl.policy.lcs, name)
56-
for _, key := range ctl.policy.keysForLC(name) {
57-
if entry := ctl.policy.getEntry(key.NodeName, key.LauncherConfigName); entry != nil {
58-
entry.handsOff = true
59-
entry.spec = nil
60-
entry.count = 0
61-
ctl.keyQueue.Queue.Add(keyItem{key})
62-
}
59+
for key, entry := range ctl.policy.entriesForLC(name) {
60+
entry.handsOff = true
61+
entry.spec = nil
62+
entry.count = 0
63+
ctl.keyQueue.Queue.Add(keyItem{key})
6364
}
64-
// Existence flip: re-enqueue LPPs that reference this LC so they can
65-
// recompute missingLCs.
66-
if prevExists {
67-
for _, lppName := range ctl.policy.lppNamesRefByLC(name) {
68-
ctl.digestQueue.Queue.Add(funcItem{kind: kindLPP, name: lppName})
69-
}
65+
for _, lppName := range ctl.policy.lppNamesRefByLC(name) {
66+
ctl.digestQueue.Queue.Add(funcItem{kind: kindLPP, name: lppName})
7067
}
7168
return nil
7269
}
@@ -81,9 +78,11 @@ func (ctl *controller) updateDigestForLC(ctx context.Context, name string) error
8178
if templateErr == "" {
8279
h, hashErr := utils.ComputeLauncherTemplateHash(lc.Spec.PodTemplate)
8380
if hashErr != nil {
84-
return fmt.Errorf("failed to compute template hash for LC %s: %w", name, hashErr)
81+
templateErr = hashErr.Error()
82+
logger.Error(hashErr, "Failed to hash PodTemplate, reporting in Status", "config", name)
83+
} else {
84+
templateHash = h
8585
}
86-
templateHash = h
8786
}
8887
if statusErr := ctl.setLCStatusErrors(ctx, lc, nonNilSlice(templateErr)); statusErr != nil {
8988
return fmt.Errorf("failed to set Status for LC %s: %w", name, statusErr)
@@ -110,18 +109,15 @@ func (ctl *controller) updateDigestForLC(ctx context.Context, name string) error
110109
}
111110

112111
// Refresh per-key entries that reference this LC.
113-
for _, key := range ctl.policy.keysForLC(name) {
114-
entry := ctl.policy.getEntry(key.NodeName, key.LauncherConfigName)
115-
if entry == nil {
116-
continue
117-
}
112+
lcd := ctl.policy.lcs[name]
113+
for key, entry := range ctl.policy.entriesForLC(name) {
118114
if templateErr != "" {
119115
entry.handsOff = true
120116
entry.spec = nil
121117
} else {
122118
entry.handsOff = false
123-
entry.spec = &lc.Spec
124-
entry.ownerRef = makeLCOwnerRef(lc)
119+
entry.spec = &lcd.object.Spec
120+
entry.ownerRef = lcd.ownerRef
125121
}
126122
ctl.keyQueue.Queue.Add(keyItem{key})
127123
}
@@ -216,8 +212,9 @@ func (ctl *controller) updateDigestForLPP(ctx context.Context, name string) erro
216212
}
217213

218214
// Apply LPP to each matched node and enqueue affected keys.
219-
for _, node := range currentMatchedNodes {
220-
ctl.applyLPPToDigestForNode(lpp, node.Name)
215+
for i := range currentMatchedNodes {
216+
node := &currentMatchedNodes[i]
217+
ctl.applyLPPToDigestForNode(lpp, node)
221218
for _, cr := range lpp.Spec.CountForLauncher {
222219
ctl.keyQueue.Queue.Add(keyItem{NodeLauncherKey{NodeName: node.Name, LauncherConfigName: cr.LauncherConfigName}})
223220
}
@@ -233,52 +230,41 @@ func (ctl *controller) updateDigestForLPP(ctx context.Context, name string) erro
233230
func (ctl *controller) updateDigestForNode(ctx context.Context, nodeName string) error {
234231
logger := klog.FromContext(ctx)
235232

236-
_, err := ctl.nodeLister.Get(nodeName)
233+
node, err := ctl.nodeLister.Get(nodeName)
237234
if err != nil {
238235
if apierrors.IsNotFound(err) {
239236
logger.Info("Node deleted, removing from digest", "node", nodeName)
240-
removedKeys := ctl.policy.removeNode(nodeName)
241-
for _, key := range removedKeys {
242-
ctl.keyQueue.Queue.Add(keyItem{key})
243-
}
237+
delete(ctl.policy.digest, nodeName)
244238
return nil
245239
}
246240
return fmt.Errorf("failed to get node %s: %w", nodeName, err)
247241
}
248-
return ctl.recomputeDigestForNode(nodeName)
242+
ctl.recomputeDigestForNode(node)
243+
return nil
249244
}
250245

251246
// recomputeDigestForNode rebuilds digest entries for a single node by replaying
252247
// every cached LPP. Pure read of ctl.policy.lpps + ctl.policy.lcs.
253-
func (ctl *controller) recomputeDigestForNode(nodeName string) error {
254-
ctl.policy.digest[nodeName] = make(map[string]*digestEntry)
248+
func (ctl *controller) recomputeDigestForNode(node *corev1.Node) {
249+
ctl.policy.digest[node.Name] = make(map[string]*digestEntry)
255250

256251
for _, lppd := range ctl.policy.lpps {
257-
if lppd == nil || lppd.object == nil {
258-
continue
259-
}
260-
ctl.applyLPPToDigestForNode(lppd.object, nodeName)
252+
ctl.applyLPPToDigestForNode(lppd.object, node)
261253
}
262254

263-
nodeMap := ctl.policy.digest[nodeName]
255+
nodeMap := ctl.policy.digest[node.Name]
264256
for lcName := range nodeMap {
265-
ctl.keyQueue.Queue.Add(keyItem{NodeLauncherKey{NodeName: nodeName, LauncherConfigName: lcName}})
257+
ctl.keyQueue.Queue.Add(keyItem{NodeLauncherKey{NodeName: node.Name, LauncherConfigName: lcName}})
266258
}
267259
if len(nodeMap) == 0 {
268-
delete(ctl.policy.digest, nodeName)
260+
delete(ctl.policy.digest, node.Name)
269261
}
270-
return nil
271262
}
272263

273264
// applyLPPToDigestForNode evaluates one LPP for one node and updates the digest
274265
// using ONLY ctl.policy.lcs as the source of truth for LC status. It never
275266
// validates templates, writes Status, or fetches from listers.
276-
func (ctl *controller) applyLPPToDigestForNode(lpp *fmav1alpha1.LauncherPopulationPolicy, nodeName string) {
277-
node, err := ctl.nodeLister.Get(nodeName)
278-
if err != nil {
279-
// Node missing or transient lister error: nothing to apply.
280-
return
281-
}
267+
func (ctl *controller) applyLPPToDigestForNode(lpp *fmav1alpha1.LauncherPopulationPolicy, node *corev1.Node) {
282268
labelSelector, selectorErr := metav1.LabelSelectorAsSelector(&lpp.Spec.EnhancedNodeSelector.LabelSelector)
283269
if selectorErr != nil {
284270
return
@@ -293,10 +279,10 @@ func (ctl *controller) applyLPPToDigestForNode(lpp *fmav1alpha1.LauncherPopulati
293279
for _, cr := range lpp.Spec.CountForLauncher {
294280
lcName := cr.LauncherConfigName
295281

296-
entry := ctl.policy.getEntry(nodeName, lcName)
282+
entry := ctl.policy.getEntry(node.Name, lcName)
297283
if entry == nil {
298284
entry = &digestEntry{lpps: make(map[string]*fmav1alpha1.LauncherPopulationPolicy)}
299-
ctl.policy.setEntry(nodeName, lcName, entry)
285+
ctl.policy.setEntry(node.Name, lcName, entry)
300286
}
301287

302288
lcd := ctl.policy.lcDigestFor(lcName)
@@ -388,5 +374,3 @@ func (ctl *controller) recomputeEntryFromLPPs(entry *digestEntry, lcName string)
388374
entry.ownerRef = lcd.ownerRef
389375
}
390376
}
391-
392-

pkg/controller/launcher-populator/digested-policy.go

Lines changed: 16 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ limitations under the License.
1717
package launcherpopulator
1818

1919
import (
20+
"iter"
2021
"sync"
2122

2223
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -137,29 +138,23 @@ func (dp *digestedPolicy) setEntry(nodeName, lcName string, entry *digestEntry)
137138
nodeMap[lcName] = entry
138139
}
139140

140-
// removeNode removes all digest entries for a node and returns the removed keys.
141-
func (dp *digestedPolicy) removeNode(nodeName string) []NodeLauncherKey {
142-
nodeMap, ok := dp.digest[nodeName]
143-
if !ok {
144-
return nil
145-
}
146-
keys := make([]NodeLauncherKey, 0, len(nodeMap))
147-
for lcName := range nodeMap {
148-
keys = append(keys, NodeLauncherKey{NodeName: nodeName, LauncherConfigName: lcName})
149-
}
150-
delete(dp.digest, nodeName)
151-
return keys
152-
}
153-
154-
// keysForLC returns all NodeLauncherKeys that reference the given LC name.
155-
func (dp *digestedPolicy) keysForLC(lcName string) []NodeLauncherKey {
156-
var keys []NodeLauncherKey
157-
for nodeName, nodeMap := range dp.digest {
158-
if _, ok := nodeMap[lcName]; ok {
159-
keys = append(keys, NodeLauncherKey{NodeName: nodeName, LauncherConfigName: lcName})
141+
// entriesForLC yields every (key, entry) pair that references lcName. Each
142+
// yielded entry is non-nil. The caller must not insert into or delete from
143+
// dp.digest while iterating; mutating fields on the yielded *digestEntry is
144+
// allowed.
145+
func (dp *digestedPolicy) entriesForLC(lcName string) iter.Seq2[NodeLauncherKey, *digestEntry] {
146+
return func(yield func(NodeLauncherKey, *digestEntry) bool) {
147+
for nodeName, nodeMap := range dp.digest {
148+
entry, ok := nodeMap[lcName]
149+
if !ok {
150+
continue
151+
}
152+
key := NodeLauncherKey{NodeName: nodeName, LauncherConfigName: lcName}
153+
if !yield(key, entry) {
154+
return
155+
}
160156
}
161157
}
162-
return keys
163158
}
164159

165160
// allKeys returns all NodeLauncherKeys in the digest.

0 commit comments

Comments
 (0)