Copilot commented on code in PR #2890:
URL: 
https://github.com/apache/apisix-ingress-controller/pull/2890#discussion_r4080724230


##########
internal/adc/cache/store.go:
##########
@@ -118,11 +178,18 @@ func routeOwner(targetCache Cache, id string) 
(types.NamespacedNameKind, bool) {
 
 // setGlobalRules is Insert and SetGlobalRules' shared implementation. Callers 
must
 // already hold s.Lock.
-func (s *Store) setGlobalRules(targetCache Cache, owner 
types.NamespacedNameKind, plugins adctypes.GlobalRule) error {
+func (s *Store) setGlobalRules(name string, targetCache Cache, owner 
types.NamespacedNameKind, plugins adctypes.GlobalRule) error {

Review Comment:
   `recordChange` is deferred during mutations, so the store 
revision/lastChange can be advanced even when the operation returns an error 
partway through (e.g., a delete/create fails in `setGlobalRules`, or a failure 
occurs after `changed` is set in `Insert`). This can make 
`ChangedSince`/`OwnerChangedSince` report changes that never successfully 
happened, which can incorrectly clear/avoid exclusions and invalidate the 
“stale sync result” protection. Prefer recording changes only after the 
mutation completes successfully (e.g., compute a `shouldRecordChange` flag, 
perform all cache operations, and call `recordChange` only right before 
returning `nil`). Apply the same pattern in `Delete`, `setPluginMetadata`, and 
any other paths using deferred change recording.



##########
internal/adc/cache/store.go:
##########
@@ -118,11 +178,18 @@ func routeOwner(targetCache Cache, id string) 
(types.NamespacedNameKind, bool) {
 
 // setGlobalRules is Insert and SetGlobalRules' shared implementation. Callers 
must
 // already hold s.Lock.
-func (s *Store) setGlobalRules(targetCache Cache, owner 
types.NamespacedNameKind, plugins adctypes.GlobalRule) error {
+func (s *Store) setGlobalRules(name string, targetCache Cache, owner 
types.NamespacedNameKind, plugins adctypes.GlobalRule) error {
        rows, err := targetCache.ListGlobalRules(&OwnerSelector{Owner: owner})
        if err != nil {
                return err
        }
+       existing := make(adctypes.Plugins, len(rows))
+       for _, row := range rows {
+               existing[row.ID] = row.Config
+       }
+       if pluginsContent(existing) != 
pluginsContent(adctypes.Plugins(plugins)) {
+               defer s.recordChange(name, owner)
+       }

Review Comment:
   `recordChange` is deferred during mutations, so the store 
revision/lastChange can be advanced even when the operation returns an error 
partway through (e.g., a delete/create fails in `setGlobalRules`, or a failure 
occurs after `changed` is set in `Insert`). This can make 
`ChangedSince`/`OwnerChangedSince` report changes that never successfully 
happened, which can incorrectly clear/avoid exclusions and invalidate the 
“stale sync result” protection. Prefer recording changes only after the 
mutation completes successfully (e.g., compute a `shouldRecordChange` flag, 
perform all cache operations, and call `recordChange` only right before 
returning `nil`). Apply the same pattern in `Delete`, `setPluginMetadata`, and 
any other paths using deferred change recording.



##########
internal/provider/apisix/status.go:
##########
@@ -137,13 +173,50 @@ func (d *apisixProvider) classifySyncResult(
        return gatewayProxyMsgs, failedEndpoints
 }
 
+// dropUnit maps a rejected resource to what has to be dropped for the rest to 
apply. A
+// service is dropped by itself. A route, stream route or named upstream is 
dropped along
+// with the service it lives in: the service is what one Kubernetes rule 
translates to,
+// and a named upstream is referenced by the service's traffic-split by id, so 
dropping
+// it alone would leave that reference dangling. Nested resources are found 
through their
+// parent, which also tells them apart when several services embed upstreams 
of the same
+// id.
+func (d *apisixProvider) dropUnit(configName string, ev adctypes.StatusEvent) 
(wireKey, exclusion, bool) {
+       serviceID := ev.ResourceID
+       switch ev.ResourceType {
+       case adctypes.TypeService:
+       case adctypes.TypeRoute, adctypes.TypeStreamRoute, 
adctypes.TypeUpstream:
+               serviceID = ev.ParentID
+       default:
+               return wireKey{}, exclusion{}, false
+       }
+       service, ok := d.store.Lookup(configName, adctypes.TypeService, 
serviceID)
+       if !ok {
+               return wireKey{}, exclusion{}, false
+       }
+       owner, name := service.Owner, service.Name
+       // A route's own owner can differ from its service's: a traffic-split 
service can
+       // combine rules several ApisixRoutes each contributed. Attribute to 
the route
+       // itself when the store can tell them apart, so fixing the actual bad 
route clears
+       // the exclusion instead of leaving it stuck on whichever owner the 
shared service
+       // happens to carry.
+       if ev.ResourceType == adctypes.TypeRoute {
+               if route, ok := d.store.Lookup(configName, adctypes.TypeRoute, 
ev.ResourceID); ok && route.Owner != (types.NamespacedNameKind{}) {
+                       owner = route.Owner
+               }
+       }
+       return wireKey{adctypes.TypeService, service.ID}, exclusion{owner: 
owner, name: name}, true
+}
+
 // applyResourceFailures writes this round's newly (or still) failing 
resources, and
 // clears the recorded error from any resource that was failing last round but 
isn't in
 // newFailures now. See updateStatusFromSyncResults for why resources use this 
delta
 // instead of GatewayProxy's full recompute.
-func (d *apisixProvider) applyResourceFailures(newFailures 
map[types.NamespacedNameKind][]string) {
+func (d *apisixProvider) applyResourceFailures(ctx context.Context, 
newFailures map[types.NamespacedNameKind][]string) {
        for resourceKey, msgs := range newFailures {
                d.updateStatus(resourceKey, failureCondition(resourceKey, 
strings.Join(msgs, "; ")))
+               if resourceKey.Kind == types.KindIngress {
+                       d.recordIngressFailureEvent(ctx, resourceKey, msgs)
+               }
        }

Review Comment:
   This emits a Warning event for every failing Ingress on every sync round, 
and performs a live GET each time. In steady-state failures (or frequent syncs) 
this can generate high write volume to the Events API and unnecessary apiserver 
reads. Consider emitting the event only when the failure message changes or 
when the Ingress transitions into failing (compare against 
`d.resourceFailures`), and/or add a simple rate limit (e.g., per Ingress key + 
time window) to reduce event spam while still providing visibility.



##########
internal/provider/apisix/provider.go:
##########
@@ -287,6 +292,11 @@ func (d *apisixProvider) Delete(ctx context.Context, obj 
client.Object) error {
 // applyResourceState upserts a resource's config associations and its 
contribution to each
 // target config's cached resource snapshot, the AIC-side bookkeeping the adc 
client
 // package no longer holds itself.
+// applyResourceState upserts a resource's config associations and its 
contribution to each
+// target config's cached resource snapshot. Whatever the skip table excluded 
for a
+// resource whose content this changed gets another try; rewriting identical 
content, which
+// reconciles triggered by unrelated events do all the time, must not retry a 
known-bad
+// resource.

Review Comment:
   The `applyResourceState` doc comment is duplicated/repeated (two consecutive 
“applyResourceState upserts…” paragraphs). Please collapse into a single 
coherent comment to avoid confusion and keep documentation consistent.



##########
internal/provider/apisix/status.go:
##########
@@ -153,6 +226,21 @@ func (d *apisixProvider) applyResourceFailures(newFailures 
map[types.NamespacedN
        d.resourceFailures = newFailures
 }
 
+// recordIngressFailureEvent fires a Warning event for an Ingress whose status 
has no
+// conditions to carry a failure. It is fired every round the Ingress stays 
failing, since
+// events expire.
+func (d *apisixProvider) recordIngressFailureEvent(ctx context.Context, nnk 
types.NamespacedNameKind, msgs []string) {
+       if d.EventRecorder == nil || d.K8sClient == nil {
+               return
+       }
+       ingress := &networkingv1.Ingress{}
+       if err := d.K8sClient.Get(ctx, nnk.NamespacedName(), ingress); err != 
nil {
+               d.log.Error(err, "failed to get Ingress to record a failure 
event", "name", nnk.Name, "namespace", nnk.Namespace)
+               return
+       }
+       d.EventRecorder.Event(ingress, corev1.EventTypeWarning, 
string(apiv2.ConditionReasonSyncFailed), strings.Join(msgs, "; "))
+}

Review Comment:
   This emits a Warning event for every failing Ingress on every sync round, 
and performs a live GET each time. In steady-state failures (or frequent syncs) 
this can generate high write volume to the Events API and unnecessary apiserver 
reads. Consider emitting the event only when the failure message changes or 
when the Ingress transitions into failing (compare against 
`d.resourceFailures`), and/or add a simple rate limit (e.g., per Ingress key + 
time window) to reduce event spam while still providing visibility.



##########
internal/adc/cache/store.go:
##########
@@ -86,6 +100,52 @@ func gatewayProxyOf(name string) (types.NamespacedNameKind, 
bool) {
        return gatewayProxy, true
 }
 
+func (s *Store) recordChange(name string, owner types.NamespacedNameKind) {
+       s.revision++
+       if s.lastChange[name] == nil {
+               s.lastChange[name] = make(map[types.NamespacedNameKind]uint64)
+       }
+       s.lastChange[name][owner] = s.revision
+}
+
+// contentOf renders items, ordered by id, the way they reach ADC, so two 
writes can be
+// compared by what they would push.
+func contentOf[T any](items []T, id func(T) string) string {
+       if len(items) == 0 {
+               return ""
+       }
+       sorted := slices.Clone(items)
+       slices.SortFunc(sorted, func(a, b T) int { return 
strings.Compare(id(a), id(b)) })
+       return canonicalJSON(sorted)
+}
+
+// canonicalJSON renders v as JSON with every object's keys sorted. A stored 
object's
+// plugin configs have been through a JSON round trip while a freshly 
translated one may
+// still hold typed structs, whose fields marshal in declaration order rather 
than sorted.
+func canonicalJSON(v any) string {
+       b, err := json.Marshal(v)
+       if err != nil {
+               return ""
+       }
+       var generic any
+       if err := json.Unmarshal(b, &generic); err != nil {
+               return string(b)
+       }
+       b, _ = json.Marshal(generic)
+       return string(b)
+}

Review Comment:
   `canonicalJSON` silently returns `""` on `json.Marshal` error and ignores 
the error on the second `json.Marshal`. For change detection, this can mask 
non-JSON-marshalable plugin configs and produce false “no change” results (or 
inconsistent behavior), which directly affects retry/exclusion logic. Consider 
surfacing errors to the caller (best), or at least logging and returning a 
deterministic fallback that forces “changed” semantics when marshaling fails; 
also handle/propagate the second marshal error rather than discarding it.



##########
internal/provider/apisix/status.go:
##########
@@ -153,6 +226,21 @@ func (d *apisixProvider) applyResourceFailures(newFailures 
map[types.NamespacedN
        d.resourceFailures = newFailures
 }
 
+// recordIngressFailureEvent fires a Warning event for an Ingress whose status 
has no
+// conditions to carry a failure. It is fired every round the Ingress stays 
failing, since
+// events expire.
+func (d *apisixProvider) recordIngressFailureEvent(ctx context.Context, nnk 
types.NamespacedNameKind, msgs []string) {

Review Comment:
   This emits a Warning event for every failing Ingress on every sync round, 
and performs a live GET each time. In steady-state failures (or frequent syncs) 
this can generate high write volume to the Events API and unnecessary apiserver 
reads. Consider emitting the event only when the failure message changes or 
when the Ingress transitions into failing (compare against 
`d.resourceFailures`), and/or add a simple rate limit (e.g., per Ingress key + 
time window) to reduce event spam while still providing visibility.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to