Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 5 additions & 5 deletions third_party/event-subscriber/events/listener.go
Original file line number Diff line number Diff line change
Expand Up @@ -168,8 +168,8 @@ func (router *EventRouter) run(wp WorkerPool, ready chan<- bool, eventSuffix str
err = json.Unmarshal(message, &event)
if err != nil {
log.WithFields(log.Fields{
"message": string(message),
}).Warnf("Error parsing message: %s", err)
"messageBytes": len(message),
}).Warnf("Error parsing message: %s", safeLogValue(err))
continue
}
wp.HandleWork(event, handlers, router.apiClient)
Expand Down Expand Up @@ -197,11 +197,11 @@ func (router *EventRouter) subscribeToEvents(subscribeURL string, accessKey stri
if err != nil {
endpoint, _ := url.Parse(target)
log.WithFields(log.Fields{
"subscribeEndpoint": endpoint.Scheme + "://" + endpoint.Host + endpoint.Path,
}).Errorf("Error subscribing to events: %s", err)
"subscribeEndpoint": safeLogValue(endpoint.Scheme + "://" + endpoint.Host + endpoint.Path),
}).Errorf("Error subscribing to events: %s", safeLogValue(err))
if resp != nil {
log.WithFields(log.Fields{
"status": resp.Status,
"status": safeLogValue(resp.Status),
"statusCode": resp.StatusCode,
}).Error("Got error response")
if resp.Body != nil {
Expand Down
12 changes: 12 additions & 0 deletions third_party/event-subscriber/events/safe_log.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
package events

import (
"fmt"
"strings"
)

func safeLogValue(value interface{}) string {
text := fmt.Sprint(value)
text = strings.ReplaceAll(text, "\r", "")
return strings.ReplaceAll(text, "\n", " ")
}
9 changes: 9 additions & 0 deletions third_party/event-subscriber/events/safe_log_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
package events

import "testing"

func TestSafeLogValueProducesSingleRecord(t *testing.T) {
if got := safeLogValue("first\r\nforged\nthird"); got != "first forged third" {
t.Fatalf("unexpected safe log value: %q", got)
}
}
25 changes: 16 additions & 9 deletions third_party/event-subscriber/events/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,12 @@ func (wp *skippingWorkerPool) HandleWork(event *Event, eventHandlers map[string]
doWork(event, eventHandlers, apiClient, wp.eventLocker(event))
}()
default:
log.Warnf("No workers available, dropping event. workerCount: %v, event: %v", cap(wp.workers), *event)
log.WithFields(log.Fields{
"workerCount": cap(wp.workers),
"eventName": safeLogValue(event.Name),
"eventId": safeLogValue(event.ID),
"resourceId": safeLogValue(event.ResourceID),
}).Warn("No workers available; dropping event")
}
}

Expand Down Expand Up @@ -77,14 +82,16 @@ func doWork(event *Event, eventHandlers map[string]EventHandler, apiClient *clie

if event.Name != "ping" {
log.WithFields(log.Fields{
"event": *event,
"eventName": safeLogValue(event.Name),
"eventId": safeLogValue(event.ID),
"resourceId": safeLogValue(event.ResourceID),
}).Debug("Processing event.")
}

unlocker := locker.Lock()
if unlocker == nil {
log.WithFields(log.Fields{
"resourceId": event.ResourceID,
"resourceId": safeLogValue(event.ResourceID),
}).Debug("Resource locked. Dropping event")
return
}
Expand All @@ -93,10 +100,10 @@ func doWork(event *Event, eventHandlers map[string]EventHandler, apiClient *clie
if fn, ok := eventHandlers[event.Name]; ok {
if err := fn(event, apiClient); err != nil {
log.WithFields(log.Fields{
"eventName": event.Name,
"eventId": event.ID,
"resourceId": event.ResourceID,
"err": err,
"eventName": safeLogValue(event.Name),
"eventId": safeLogValue(event.ID),
"resourceId": safeLogValue(event.ResourceID),
"err": safeLogValue(err),
}).Error("Error processing event")

reply := &client.Publish{
Expand All @@ -108,13 +115,13 @@ func doWork(event *Event, eventHandlers map[string]EventHandler, apiClient *clie
_, err := apiClient.Publish.Create(reply)
if err != nil {
log.WithFields(log.Fields{
"err": err,
"err": safeLogValue(err),
}).Error("Error sending error-reply")
}
}
} else {
log.WithFields(log.Fields{
"eventName": event.Name,
"eventName": safeLogValue(event.Name),
}).Warn("No event handler registered for event")
}
}
Loading