mirror of
https://github.com/kubernetes/node-problem-detector.git
synced 2026-08-28 01:47:20 +00:00
add log-counter go plugin
This commit is contained in:
@@ -108,7 +108,7 @@ func (p *Plugin) run(rule cpmtypes.CustomRule) (exitStatus cpmtypes.Status, outp
|
||||
}
|
||||
defer cancel()
|
||||
|
||||
cmd := exec.CommandContext(ctx, rule.Path)
|
||||
cmd := exec.CommandContext(ctx, rule.Path, rule.Args...)
|
||||
stdout, err := cmd.Output()
|
||||
if err != nil {
|
||||
if _, ok := err.(*exec.ExitError); !ok {
|
||||
|
||||
@@ -48,6 +48,8 @@ type CustomRule struct {
|
||||
Reason string `json:"reason"`
|
||||
// Path is the path to the custom plugin.
|
||||
Path string `json:"path"`
|
||||
// Args is the args passed to the custom plugin.
|
||||
Args []string `json:"args"`
|
||||
// Timeout is the timeout string for the custom plugin to execute.
|
||||
TimeoutString *string `json:"timeout"`
|
||||
// Timeout is the timeout for the custom plugin to execute.
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
/*
|
||||
Copyright 2018 The Kubernetes Authors All rights reserved.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
package logcounter
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"k8s.io/kubernetes/pkg/util/clock"
|
||||
|
||||
"k8s.io/node-problem-detector/cmd/logcounter/options"
|
||||
"k8s.io/node-problem-detector/pkg/logcounter/types"
|
||||
"k8s.io/node-problem-detector/pkg/systemlogmonitor"
|
||||
"k8s.io/node-problem-detector/pkg/systemlogmonitor/logwatchers/kmsg"
|
||||
watchertypes "k8s.io/node-problem-detector/pkg/systemlogmonitor/logwatchers/types"
|
||||
systemtypes "k8s.io/node-problem-detector/pkg/systemlogmonitor/types"
|
||||
)
|
||||
|
||||
const (
|
||||
bufferSize = 1000
|
||||
timeout = 1 * time.Second
|
||||
)
|
||||
|
||||
type logCounter struct {
|
||||
logCh <-chan *systemtypes.Log
|
||||
buffer systemlogmonitor.LogBuffer
|
||||
pattern string
|
||||
clock clock.Clock
|
||||
}
|
||||
|
||||
func NewKmsgLogCounter(options *options.LogCounterOptions) (types.LogCounter, error) {
|
||||
watcher := kmsg.NewKmsgWatcher(watchertypes.WatcherConfig{Lookback: options.Lookback})
|
||||
logCh, err := watcher.Watch()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error watching kmsg: %v", err)
|
||||
}
|
||||
return &logCounter{
|
||||
logCh: logCh,
|
||||
buffer: systemlogmonitor.NewLogBuffer(bufferSize),
|
||||
pattern: options.Pattern,
|
||||
clock: clock.RealClock{},
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (e *logCounter) Count() (count int) {
|
||||
start := e.clock.Now()
|
||||
for {
|
||||
select {
|
||||
case log := <-e.logCh:
|
||||
// We only want to count events up until the time at which we started.
|
||||
// Otherwise we would run forever
|
||||
if start.Before(log.Timestamp) {
|
||||
return
|
||||
}
|
||||
e.buffer.Push(log)
|
||||
if len(e.buffer.Match(e.pattern)) != 0 {
|
||||
count++
|
||||
}
|
||||
case <-e.clock.After(timeout):
|
||||
// Don't block forever if we do not get any new messages
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,129 @@
|
||||
/*
|
||||
Copyright 2018 The Kubernetes Authors All rights reserved.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
package logcounter
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"k8s.io/kubernetes/pkg/util/clock"
|
||||
|
||||
"k8s.io/node-problem-detector/pkg/logcounter/types"
|
||||
"k8s.io/node-problem-detector/pkg/systemlogmonitor"
|
||||
systemtypes "k8s.io/node-problem-detector/pkg/systemlogmonitor/types"
|
||||
)
|
||||
|
||||
func NewTestLogCounter(pattern string, startTime time.Time) (types.LogCounter, *clock.FakeClock, chan *systemtypes.Log) {
|
||||
logCh := make(chan *systemtypes.Log)
|
||||
clock := clock.NewFakeClock(startTime)
|
||||
return &logCounter{
|
||||
logCh: logCh,
|
||||
buffer: systemlogmonitor.NewLogBuffer(bufferSize),
|
||||
pattern: pattern,
|
||||
clock: clock,
|
||||
}, clock, logCh
|
||||
}
|
||||
|
||||
func TestCount(t *testing.T) {
|
||||
startTime := time.Now()
|
||||
for _, tc := range []struct {
|
||||
description string
|
||||
logs []*systemtypes.Log
|
||||
pattern string
|
||||
expectedCount int
|
||||
}{
|
||||
{
|
||||
description: "no logs",
|
||||
logs: []*systemtypes.Log{},
|
||||
pattern: "",
|
||||
expectedCount: 0,
|
||||
},
|
||||
{
|
||||
description: "one matching log",
|
||||
logs: []*systemtypes.Log{
|
||||
{
|
||||
Timestamp: startTime.Add(-time.Second),
|
||||
Message: "0",
|
||||
},
|
||||
},
|
||||
pattern: "0",
|
||||
expectedCount: 1,
|
||||
},
|
||||
{
|
||||
description: "one non-matching log",
|
||||
logs: []*systemtypes.Log{
|
||||
{
|
||||
Timestamp: startTime.Add(-time.Second),
|
||||
Message: "1",
|
||||
},
|
||||
},
|
||||
pattern: "0",
|
||||
expectedCount: 0,
|
||||
},
|
||||
{
|
||||
description: "log too new",
|
||||
logs: []*systemtypes.Log{
|
||||
{
|
||||
Timestamp: startTime.Add(time.Second),
|
||||
Message: "0",
|
||||
},
|
||||
},
|
||||
pattern: "0",
|
||||
expectedCount: 0,
|
||||
},
|
||||
{
|
||||
description: "many logs",
|
||||
logs: []*systemtypes.Log{
|
||||
{
|
||||
Timestamp: startTime.Add(-time.Second),
|
||||
Message: "0",
|
||||
},
|
||||
{
|
||||
Timestamp: startTime.Add(-time.Second),
|
||||
Message: "0",
|
||||
},
|
||||
{
|
||||
Timestamp: startTime.Add(-time.Second),
|
||||
Message: "1",
|
||||
},
|
||||
{
|
||||
Timestamp: startTime.Add(time.Second),
|
||||
Message: "0",
|
||||
},
|
||||
},
|
||||
pattern: "0",
|
||||
expectedCount: 2,
|
||||
},
|
||||
} {
|
||||
t.Run(tc.description, func(t *testing.T) {
|
||||
counter, fakeClock, logCh := NewTestLogCounter(tc.pattern, startTime)
|
||||
go func(logs []*systemtypes.Log, ch chan<- *systemtypes.Log) {
|
||||
for _, log := range logs {
|
||||
ch <- log
|
||||
}
|
||||
// trigger the timeout to ensure the test doesn't block permenantly
|
||||
for {
|
||||
fakeClock.Step(2 * timeout)
|
||||
}
|
||||
}(tc.logs, logCh)
|
||||
actualCount := counter.Count()
|
||||
if actualCount != tc.expectedCount {
|
||||
t.Errorf("got %d; expected %d", actualCount, tc.expectedCount)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
/*
|
||||
Copyright 2018 The Kubernetes Authors All rights reserved.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
package types
|
||||
|
||||
type LogCounter interface {
|
||||
Count() int
|
||||
}
|
||||
Reference in New Issue
Block a user