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]