@@ -48,6 +48,8 @@ type baseClient struct {
4848}
4949
5050func (c * baseClient ) connect (ctx context.Context ) error {
51+ logger := klog .FromContext (ctx )
52+
5153 var err error
5254 c .cloudEventsClient , err = c .newCloudEventsClient (ctx )
5355 if err != nil {
@@ -58,7 +60,7 @@ func (c *baseClient) connect(ctx context.Context) error {
5860 go func () {
5961 for {
6062 if ! c .isClientReady () {
61- klog .V (4 ). Infof ("reconnecting the cloudevents client" )
63+ logger .V (2 ). Info ("reconnecting the cloudevents client" )
6264
6365 c .cloudEventsClient , err = c .newCloudEventsClient (ctx )
6466 // TODO enhance the cloudevents SKD to avoid wrapping the error type to distinguish the net connection
@@ -70,7 +72,7 @@ func (c *baseClient) connect(ctx context.Context) error {
7072 continue
7173 }
7274 // the cloudevents network connection is back, mark the client ready and send the receiver restart signal
73- klog .V (4 ). Infof ("the cloudevents client is reconnected" )
75+ logger .V (2 ). Info ("the cloudevents client is reconnected" )
7476 increaseClientReconnectedCounter (c .clientID )
7577 c .setClientReady (true )
7678 c .sendReceiverSignal (restartReceiverSignal )
@@ -108,6 +110,7 @@ func (c *baseClient) connect(ctx context.Context) error {
108110}
109111
110112func (c * baseClient ) publish (ctx context.Context , evt cloudevents.Event ) error {
113+ logger := klog .FromContext (ctx )
111114 now := time .Now ()
112115
113116 if err := c .cloudEventsRateLimiter .Wait (ctx ); err != nil {
@@ -116,8 +119,11 @@ func (c *baseClient) publish(ctx context.Context, evt cloudevents.Event) error {
116119
117120 latency := time .Since (now )
118121 if latency > longThrottleLatency {
119- klog .Warningf ("Waited for %v due to client-side throttling, not priority and fairness, request: %s" ,
120- latency , evt .Context )
122+ logger .V (3 ).Info (
123+ "Client-side throttling delay (not priority and fairness)" ,
124+ "latency" , latency ,
125+ "request" , evt .Context .GetID (),
126+ )
121127 }
122128
123129 sendingCtx , err := c .cloudEventsOptions .WithContext (ctx , evt .Context )
@@ -129,8 +135,8 @@ func (c *baseClient) publish(ctx context.Context, evt cloudevents.Event) error {
129135 return fmt .Errorf ("the cloudevents client is not ready" )
130136 }
131137
132- klog .V (4 ). Infof ("Sending event: %v \n %s" , sendingCtx , evt .Context )
133- klog .V (5 ).Infof ("Sending event: evt=%s" , evt )
138+ logger .V (2 ). Info ("Sending event" , "context" , sendingCtx , "event" , evt .Context )
139+ logger .V (5 ).Info ("Sending event" , "event" , func () any { return evt . String () } )
134140 if err := c .cloudEventsClient .Send (sendingCtx , evt ); cloudevents .IsUndelivered (err ) {
135141 return err
136142 }
@@ -142,9 +148,10 @@ func (c *baseClient) subscribe(ctx context.Context, receive receiveFn) {
142148 c .Lock ()
143149 defer c .Unlock ()
144150
151+ logger := klog .FromContext (ctx )
145152 // make sure there is only one subscription go routine starting for one client.
146153 if c .receiverChan != nil {
147- klog . Warningf ("the subscription has already started" )
154+ logger . V ( 2 ). Info ("the subscription has already started" )
148155 return
149156 }
150157
@@ -159,8 +166,8 @@ func (c *baseClient) subscribe(ctx context.Context, receive receiveFn) {
159166 if startReceiving {
160167 go func () {
161168 if err := c .cloudEventsClient .StartReceiver (receiverCtx , func (evt cloudevents.Event ) {
162- klog .V (4 ). Infof ("Received event: %s " , evt .Context )
163- klog .V (5 ).Infof ("Received event: evt=%s" , evt )
169+ logger .V (2 ). Info ("Received event" , "event " , evt .Context )
170+ logger .V (5 ).Info ("Received event" , "event" , func () any { return evt . String () } )
164171
165172 receive (receiverCtx , evt )
166173 }); err != nil {
@@ -183,12 +190,12 @@ func (c *baseClient) subscribe(ctx context.Context, receive receiveFn) {
183190
184191 switch signal {
185192 case restartReceiverSignal :
186- klog .V (4 ). Infof ("restart the cloudevents receiver" )
193+ logger .V (2 ). Info ("restart the cloudevents receiver" )
187194 // rebuild the receiver context and restart receiving
188195 receiverCtx , receiverCancel = context .WithCancel (context .TODO ())
189196 startReceiving = true
190197 case stopReceiverSignal :
191- klog .V (4 ). Infof ("stop the cloudevents receiver" )
198+ logger .V (2 ). Info ("stop the cloudevents receiver" )
192199 receiverCancel ()
193200 default :
194201 runtime .HandleError (fmt .Errorf ("unknown receiver signal %d" , signal ))
0 commit comments