Skip to content

Commit a3f7625

Browse files
authored
using context log (#151)
Signed-off-by: Wei Liu <liuweixa@redhat.com>
1 parent b82d34b commit a3f7625

1 file changed

Lines changed: 18 additions & 11 deletions

File tree

pkg/cloudevents/generic/baseclient.go

Lines changed: 18 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,8 @@ type baseClient struct {
4848
}
4949

5050
func (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

110112
func (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

Comments
 (0)