Skip to content

Commit

Permalink
feat: modify k8 events to be more understandable [DET-3172] (#1215)
Browse files Browse the repository at this point in the history
  • Loading branch information
eecsliu authored Sep 1, 2020
1 parent 34a744f commit 8533bbc
Show file tree
Hide file tree
Showing 2 changed files with 31 additions and 0 deletions.
26 changes: 26 additions & 0 deletions master/internal/kubernetes/events.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
package kubernetes

import (
"regexp"

"github.com/pkg/errors"

k8sV1 "k8s.io/api/core/v1"
Expand All @@ -10,6 +12,8 @@ import (
"github.com/determined-ai/determined/master/pkg/actor"
)

const gpuTextReplacement = "Waiting for resources. "

// Messages that are sent to the event listener.
type (
startEventListener struct{}
Expand Down Expand Up @@ -69,6 +73,7 @@ func (e *eventListener) startEventListener(ctx *actor.Context) error {
ctx.Log().Info("event listener is starting")
for event := range watch.ResultChan() {
newEvent := event.Object.(*k8sV1.Event)
e.modMessage(newEvent)
ctx.Tell(e.podsHandler, podEventUpdate{event: newEvent})
}

Expand All @@ -77,3 +82,24 @@ func (e *eventListener) startEventListener(ctx *actor.Context) error {

return nil
}

func (e *eventListener) modMessage(msg *k8sV1.Event) {
replacements := map[string]string{
"nodes are available": gpuTextReplacement,
"pod triggered scale-up": "Job requires additional resources, scaling up cluster.",
"Successfully assigned": "Pod resources allocated.",
"skip schedule deleting pod": "Deleting unscheduled pod.",
}

for k, v := range replacements {
matched, error := regexp.MatchString(k, msg.Message)
if error != nil {
break
} else if matched {
if v == gpuTextReplacement {
v += string(msg.Message[0]) + " GPUs available, "
}
msg.Message = v
}
}
}
5 changes: 5 additions & 0 deletions master/internal/kubernetes/pod.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package kubernetes

import (
"fmt"
"strconv"
"time"

"github.com/pkg/errors"
Expand Down Expand Up @@ -404,6 +405,10 @@ func (p *pod) receivePodEventUpdate(ctx *actor.Context, msg podEventUpdate) {
return
}

if msg.event.Message[0:23] == gpuTextReplacement {
msg.event.Message += strconv.Itoa(p.gpus) + " GPUs required."
}

message := fmt.Sprintf("Pod %s: %s", msg.event.InvolvedObject.Name, msg.event.Message)
ctx.Tell(p.taskHandler, sproto.ContainerLog{
Container: p.container,
Expand Down

0 comments on commit 8533bbc

Please sign in to comment.