tigerquoll commented on code in PR #1060:
URL: https://github.com/apache/yunikorn-k8shim/pull/1060#discussion_r3889398975


##########
pkg/admission/webhook_manager.go:
##########
@@ -86,17 +86,11 @@ type webhookManagerImpl struct {
 
 // NewWebhookManager is used to create a new webhook manager
 func NewWebhookManager(conf *conf.AdmissionControllerConf) (WebhookManager, 
error) {
-       kubeconfig, err := client.CreateRestConfig(conf.GetKubeConfig())
-       if err != nil {
-               log.Log(log.AdmissionWebhook).Error("Unable to create 
kubernetes config", zap.Error(err))
-               return nil, err
-       }
-       clientset, err := kubernetes.NewForConfig(kubeconfig)
-       if err != nil {
-               log.Log(log.AdmissionWebhook).Error("Unable to create 
kubernetes clientset", zap.Error(err))
-               return nil, err
-       }
-       return newWebhookManagerImpl(conf, clientset), nil
+       // the webhook and secret writes must be attributed to the admission 
controller; the
+       // admission controller never populates the scheduler configuration, so 
the client runs

Review Comment:
   You're right on the outcome, and the comment was misleading either way, so 
it's gone.
   
   For the record the values didn't come from the configmap: nothing in the 
admission controller calls `conf.UpdateConfigMaps`, so `NewKubeClient` read the 
*defaults* of the scheduler configuration — 1000/1000 for the informer client 
in `main.go` — while the webhook client, built straight from 
`CreateRestConfig`, ran on client-go's 5/10. Both now run on the client-go 
defaults, see the next thread.
   



##########
pkg/admission/webhook_manager.go:
##########
@@ -86,17 +86,11 @@ type webhookManagerImpl struct {
 
 // NewWebhookManager is used to create a new webhook manager
 func NewWebhookManager(conf *conf.AdmissionControllerConf) (WebhookManager, 
error) {
-       kubeconfig, err := client.CreateRestConfig(conf.GetKubeConfig())
-       if err != nil {
-               log.Log(log.AdmissionWebhook).Error("Unable to create 
kubernetes config", zap.Error(err))
-               return nil, err
-       }
-       clientset, err := kubernetes.NewForConfig(kubeconfig)
-       if err != nil {
-               log.Log(log.AdmissionWebhook).Error("Unable to create 
kubernetes clientset", zap.Error(err))
-               return nil, err
-       }
-       return newWebhookManagerImpl(conf, clientset), nil
+       // the webhook and secret writes must be attributed to the admission 
controller; the
+       // admission controller never populates the scheduler configuration, so 
the client runs
+       // on the defaults: no client side rate limiting

Review Comment:
   Done. Both admission controller clients — the informers in `main.go` and the 
webhook/secret writes in `NewWebhookManager` — now leave QPS and burst at 0, 
i.e. client-go's 5/10, through `client.NewAdmissionControllerKubeClient`. The 
bootstrap client is on the same defaults, which is what it had on master.
   



##########
pkg/client/apifactory.go:
##########
@@ -89,9 +89,18 @@ type APIFactory struct {
        lock     *locking.RWMutex
 }
 
-func NewAPIFactory(scheduler api.SchedulerAPI, informerFactory 
informers.SharedInformerFactory, configs *conf.SchedulerConf, testMode bool) 
(*APIFactory, error) {
+// NewAPIFactory creates the clients shared by the shim.
+func NewAPIFactory(scheduler api.SchedulerAPI, configs *conf.SchedulerConf, 
testMode bool) (*APIFactory, error) {
        kubeClient := NewKubeClient(configs.KubeConfig)
-       namespaceInformerFactory := 
informers.NewSharedInformerFactoryWithOptions(kubeClient.GetClientSet(), 0, 
informers.WithNamespace(configs.Namespace))
+
+       // all informers, cluster wide and namespaced, share one clientset: 
both only run
+       // informers so they need the same unlimited client and the same 
attribution.

Review Comment:
   Agreed. With the events off the shared bucket, the relist-versus-event-storm 
coupling that the informer client was meant to break is gone, and a separate 
clientset only bought a distinct user agent in the audit log. The informer 
clientset is removed; both informer factories are back on the scheduler client.
   



##########
pkg/client/apifactory.go:
##########
@@ -89,9 +90,12 @@ type APIFactory struct {
        lock     *locking.RWMutex
 }
 
-func NewAPIFactory(scheduler api.SchedulerAPI, informerFactory 
informers.SharedInformerFactory, configs *conf.SchedulerConf, testMode bool) 
(*APIFactory, error) {
+// NewAPIFactory creates the clients shared by the shim. The clientset backing 
informerFactory
+// is passed in so that the namespaced factory created here shares it: both 
only run informers
+// so they need the same unlimited client and the same attribution.
+func NewAPIFactory(scheduler api.SchedulerAPI, informerClientSet 
kubernetes.Interface, informerFactory informers.SharedInformerFactory, configs 
*conf.SchedulerConf, testMode bool) (*APIFactory, error) {
        kubeClient := NewKubeClient(configs.KubeConfig)

Review Comment:
   Done — the layout is now the two clients you describe: `yunikorn-scheduler` 
for the scheduling actions and the informers, and `yunikorn-scheduler/events` 
for the events, plus the pre-existing bootstrap and admission controller 
clients, both on client-go defaults.
   



##########
pkg/client/eventsink.go:
##########


Review Comment:
   Agreed that not all events are the same; the bucket was type-blind and that 
was the wrong shape. What's changed:
   
   - **Warning events are never shed by the bucket.** `FailedScheduling`, 
`TaskRejected`, `NodeRejected`, `ApplicationFailed`, the gang scheduling 
failures — everything that explains why a pod is not running — always goes out 
and does not consume a token.
   - **`PodUnschedulable` becomes a Warning.** It was typed Normal, which would 
have put it in the sheddable class; kube-scheduler's equivalent, 
`FailedScheduling`, is a Warning too.
   - **Only Normal events are shed above `eventQPS`.** Counting what the shim 
emits, the volume is three Normal events per *successfully* scheduled pod 
(`Scheduling`, `Scheduled`, `PodBindSuccessful`) plus the core's 
`Informational` records. At the ~1,900 binds/s measured on the KWOK rig that is 
~5,700 events/s saying "this worked", and that is what a storm sheds.
   - **`kubernetes.eventLevel`** (`normal` | `warning` | `none`, default 
`normal`) is the verbosity knob you suggested. It is applied at the recorder, 
so a suppressed event is never copied, cached or written, and it is 
hot-reloadable like the log level.
   
   The server-side mute still applies to every type. Once APF is rejecting 
events, client-go drops them itself after its Retry-After retries 
(`recordEvent` treats every `StatusError` as "Server rejected event (will not 
retry!)"), so muting only saves the round trips and the parked goroutines; it 
does not lose an event that would otherwise have landed.
   



##########
pkg/client/kubeclient.go:
##########
@@ -30,44 +30,127 @@ import (
        "k8s.io/client-go/rest"
        "k8s.io/client-go/tools/clientcmd"
        "k8s.io/client-go/util/retry"
+       "k8s.io/utils/clock"
 
        "github.com/apache/yunikorn-k8shim/pkg/conf"
        "github.com/apache/yunikorn-k8shim/pkg/log"
 )
 
 type SchedulerKubeClient struct {
        clientSet *kubernetes.Clientset
-       configs   *rest.Config
 }
 
-func newBootstrapSchedulerKubeClient(kc string) SchedulerKubeClient {
-       config := CreateRestConfigOrDie(kc)
-       configuredClient, err := kubernetes.NewForConfig(config)
-       if err != nil {
-               log.Log(log.ShimClient).Fatal("failed to get Clientset", 
zap.Error(err))
-       }
-       return SchedulerKubeClient{
-               clientSet: configuredClient,
-               configs:   config,
+// Every client identifies the concern it serves in its user agent so that 
traffic can be
+// attributed on the API server side. The concern comes first to allow prefix 
matching, the
+// build version is appended by UserAgent().
+const (
+       UserAgentAdmissionController = "yunikorn-admission-controller"
+
+       userAgentBootstrap = "yunikorn-bootstrap"
+       userAgentWrites    = "yunikorn-scheduler/writes"
+       userAgentInformers = "yunikorn-scheduler/informers"
+       userAgentEvents    = "yunikorn-scheduler/events"
+)
+
+// UserAgent appends the build version to the concern, e.g. 
"yunikorn-scheduler/writes
+// (1.7.0)". The version is only set in release builds, it is left out when 
empty.
+func UserAgent(concern string) string {
+       return userAgentWithVersion(concern, 
conf.GetBuildInfoMap()["buildVersion"])

Review Comment:
   Removed; the user agents are the bare concern strings.
   



##########
pkg/client/kubeclient.go:
##########
@@ -30,44 +30,127 @@ import (
        "k8s.io/client-go/rest"
        "k8s.io/client-go/tools/clientcmd"
        "k8s.io/client-go/util/retry"
+       "k8s.io/utils/clock"
 
        "github.com/apache/yunikorn-k8shim/pkg/conf"
        "github.com/apache/yunikorn-k8shim/pkg/log"
 )
 
 type SchedulerKubeClient struct {
        clientSet *kubernetes.Clientset
-       configs   *rest.Config
 }
 
-func newBootstrapSchedulerKubeClient(kc string) SchedulerKubeClient {
-       config := CreateRestConfigOrDie(kc)
-       configuredClient, err := kubernetes.NewForConfig(config)
-       if err != nil {
-               log.Log(log.ShimClient).Fatal("failed to get Clientset", 
zap.Error(err))
-       }
-       return SchedulerKubeClient{
-               clientSet: configuredClient,
-               configs:   config,
+// Every client identifies the concern it serves in its user agent so that 
traffic can be
+// attributed on the API server side. The concern comes first to allow prefix 
matching, the
+// build version is appended by UserAgent().
+const (
+       UserAgentAdmissionController = "yunikorn-admission-controller"
+
+       userAgentBootstrap = "yunikorn-bootstrap"
+       userAgentWrites    = "yunikorn-scheduler/writes"
+       userAgentInformers = "yunikorn-scheduler/informers"
+       userAgentEvents    = "yunikorn-scheduler/events"
+)
+
+// UserAgent appends the build version to the concern, e.g. 
"yunikorn-scheduler/writes
+// (1.7.0)". The version is only set in release builds, it is left out when 
empty.
+func UserAgent(concern string) string {
+       return userAgentWithVersion(concern, 
conf.GetBuildInfoMap()["buildVersion"])
+}
+
+func userAgentWithVersion(concern, version string) string {
+       if version == "" {
+               return concern
        }
+       return fmt.Sprintf("%s (%s)", concern, version)
 }
 
-func newSchedulerKubeClient(kc string) SchedulerKubeClient {
-       schedulerConf := conf.GetSchedulerConf()
+// rateLimitPolicy turns the configured qps and burst into the values to set 
on a REST config
+// and a description of the resulting limit for logging.
+// A qps <= 0 disables client side rate limiting: this is expressed as a 
negative QPS, which
+// client-go treats as "create no rate limiter", while a QPS of 0 silently 
falls back to its
+// own defaults of 5 QPS / 10 burst (see RESTClientFor and 
NewForConfigAndClient in
+// client-go). A QPS set without a burst is rejected by client-go, so the 
burst defaults to
+// the configured qps.

Review Comment:
   Done. `rateLimitPolicy` and its warnings are gone in favour of the 
normalisation you suggested, with one difference: burst is the configured burst 
when set, otherwise the qps — `max(burst, qps)` would silently raise an 
explicit `burst < qps`. `0/0` is left to client-go's 5/10 as you asked, so 
`qps: "0"` now behaves exactly as it did before this PR; the defaults stay at 
`-1`. The compatibility table in the description is updated to match.
   
   One thing I'd like to confirm with you: the scheduler client itself still 
defaults to no client-side limiter (`-1`), on the measurements in the 
description — the 1000/1000 bucket clipped the bind burst, and at `qps: 50` the 
limiter was the throughput ceiling at 52 binds/s. If you'd rather keep a 
default cap on the scheduler client, say so and I'll set one.
   



-- 
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