diff --git a/third_party/event-subscriber/events/listener.go b/third_party/event-subscriber/events/listener.go index 57d2622..25d70bb 100644 --- a/third_party/event-subscriber/events/listener.go +++ b/third_party/event-subscriber/events/listener.go @@ -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) @@ -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 { diff --git a/third_party/event-subscriber/events/safe_log.go b/third_party/event-subscriber/events/safe_log.go new file mode 100644 index 0000000..edfb6ee --- /dev/null +++ b/third_party/event-subscriber/events/safe_log.go @@ -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", " ") +} diff --git a/third_party/event-subscriber/events/safe_log_test.go b/third_party/event-subscriber/events/safe_log_test.go new file mode 100644 index 0000000..463decc --- /dev/null +++ b/third_party/event-subscriber/events/safe_log_test.go @@ -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) + } +} diff --git a/third_party/event-subscriber/events/worker.go b/third_party/event-subscriber/events/worker.go index 8d5ad73..f4d1e53 100644 --- a/third_party/event-subscriber/events/worker.go +++ b/third_party/event-subscriber/events/worker.go @@ -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") } } @@ -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 } @@ -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{ @@ -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") } }