Add first version of node-problem-detector

This commit is contained in:
Lantao Liu
2016-05-17 15:55:33 -07:00
parent 802acee7e3
commit f0312655bd
31 changed files with 2370 additions and 0 deletions
+100
View File
@@ -0,0 +1,100 @@
/*
Copyright 2016 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 problemclient
import (
"fmt"
"reflect"
"sync"
"time"
"k8s.io/kubernetes/pkg/api"
)
// FakeProblemClient is a fake problem client for debug.
type FakeProblemClient struct {
sync.Mutex
conditions map[api.NodeConditionType]api.NodeCondition
errors map[string]error
}
// NewFakeProblemClient creates a new fake problem client.
func NewFakeProblemClient() *FakeProblemClient {
return &FakeProblemClient{
conditions: make(map[api.NodeConditionType]api.NodeCondition),
errors: make(map[string]error),
}
}
// InjectError injects error to specific function.
func (f *FakeProblemClient) InjectError(fun string, err error) {
f.Lock()
defer f.Unlock()
f.errors[fun] = err
}
// AssertConditions asserts that the internal conditions in fake problem client should match
// the expected conditions.
func (f *FakeProblemClient) AssertConditions(expected []api.NodeCondition) error {
conditions := map[api.NodeConditionType]api.NodeCondition{}
for _, condition := range expected {
conditions[condition.Type] = condition
}
if !reflect.DeepEqual(conditions, f.conditions) {
return fmt.Errorf("expected %+v, got %+v", conditions, f.conditions)
}
return nil
}
// SetConditions is a fake mimic of SetConditions, it only update the internal condition cache.
func (f *FakeProblemClient) SetConditions(conditions []api.NodeCondition, timeout time.Duration) error {
f.Lock()
defer f.Unlock()
if err, ok := f.errors["SetConditions"]; ok {
return err
}
for _, condition := range conditions {
t := condition.Type
if condition.Status == api.ConditionFalse {
delete(f.conditions, t)
} else {
f.conditions[t] = condition
}
}
return nil
}
// GetConditions is a fake mimic of GetConditions, it returns the conditions cached internally.
func (f *FakeProblemClient) GetConditions(types []api.NodeConditionType) ([]*api.NodeCondition, error) {
f.Lock()
defer f.Unlock()
if err, ok := f.errors["GetConditions"]; ok {
return nil, err
}
conditions := []*api.NodeCondition{}
for _, t := range types {
condition, ok := f.conditions[t]
if ok {
conditions = append(conditions, &condition)
}
}
return conditions, nil
}
// Eventf does nothing now.
func (f *FakeProblemClient) Eventf(eventType string, source, reason, messageFmt string, args ...interface{}) {
}
+202
View File
@@ -0,0 +1,202 @@
/*
Copyright 2016 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 problemclient
import (
"fmt"
"os"
"time"
"k8s.io/kubernetes/pkg/api"
"k8s.io/kubernetes/pkg/api/errors"
"k8s.io/kubernetes/pkg/api/unversioned"
clientset "k8s.io/kubernetes/pkg/client/clientset_generated/internalclientset"
unversionedcore "k8s.io/kubernetes/pkg/client/clientset_generated/internalclientset/typed/core/unversioned"
"k8s.io/kubernetes/pkg/client/record"
"k8s.io/kubernetes/pkg/client/restclient"
"k8s.io/kubernetes/pkg/types"
"k8s.io/kubernetes/pkg/util"
"github.com/golang/glog"
)
// Client is the interface of problem client
type Client interface {
// GetConditions get all specifiec conditions of current node.
GetConditions(conditionTypes []api.NodeConditionType) ([]*api.NodeCondition, error)
// SetConditions set or update conditions of current node.
// Notice that conditions with status api.ConditionFalse will be removed from the condition list, so that
// we'll only have useful conditions in the condition list.
SetConditions(conditions []api.NodeCondition, timeout time.Duration) error
// Eventf reports the event.
Eventf(eventType string, source, reason, messageFmt string, args ...interface{})
}
type nodeProblemClient struct {
nodeName string
client clientset.Interface
clock util.Clock
recorders map[string]record.EventRecorder
nodeRef *api.ObjectReference
}
// NewClientOrDie creates a new problem client, panics if error occurs.
func NewClientOrDie() Client {
c := &nodeProblemClient{clock: util.RealClock{}}
cfg, err := restclient.InClusterConfig()
if err != nil {
panic(err)
}
// TODO(random-liu): Set QPS Limit
c.client, err = clientset.NewForConfig(cfg)
if err != nil {
panic(err)
}
// TODO(random-liu): Get node name from cloud provider
c.nodeName, err = os.Hostname()
if err != nil {
panic(err)
}
c.nodeRef = getNodeRef(c.nodeName)
c.recorders = make(map[string]record.EventRecorder)
return c
}
func (c *nodeProblemClient) GetConditions(conditionTypes []api.NodeConditionType) ([]*api.NodeCondition, error) {
node, err := c.client.Core().Nodes().Get(c.nodeName)
if err != nil {
return nil, err
}
conditions := []*api.NodeCondition{}
for _, conditionType := range conditionTypes {
for _, condition := range node.Status.Conditions {
if condition.Type == conditionType {
conditions = append(conditions, &condition)
}
}
}
return conditions, nil
}
func (c *nodeProblemClient) SetConditions(newConditions []api.NodeCondition, timeout time.Duration) error {
for i := range newConditions {
// Each time we update the conditions, we update the heart beat time
newConditions[i].LastHeartbeatTime = unversioned.NewTime(c.clock.Now())
}
return c.updateNodeCondition(func(conditions []api.NodeCondition) []api.NodeCondition {
for _, condition := range newConditions {
if condition.Status == api.ConditionFalse {
conditions = unsetCondition(condition.Type, conditions)
} else {
conditions = setCondition(condition, conditions)
}
}
return conditions
}, timeout)
}
func (c *nodeProblemClient) Eventf(eventType, source, reason, messageFmt string, args ...interface{}) {
recorder, found := c.recorders[source]
if !found {
// TODO(random-liu): If needed use separate client and QPS limit for event.
recorder = getEventRecorder(c.client, c.nodeName, source)
c.recorders[source] = recorder
}
recorder.Eventf(c.nodeRef, eventType, reason, messageFmt, args...)
}
func unsetCondition(conditionType api.NodeConditionType, conditions []api.NodeCondition) []api.NodeCondition {
result := []api.NodeCondition{}
for _, condition := range conditions {
if condition.Type != conditionType {
result = append(result, condition)
}
}
return result
}
func setCondition(condition api.NodeCondition, conditions []api.NodeCondition) []api.NodeCondition {
found := false
for i := range conditions {
if conditions[i].Type == condition.Type {
target := &conditions[i]
*target = condition
found = true
break
}
}
if !found {
conditions = append(conditions, condition)
}
return conditions
}
func (c *nodeProblemClient) updateNodeCondition(updateFunc func([]api.NodeCondition) []api.NodeCondition, timeout time.Duration) error {
updateTime := c.clock.Now()
for {
node, err := c.client.Core().Nodes().Get(c.nodeName)
if err != nil {
return err
}
node.Status.Conditions = updateFunc(node.Status.Conditions)
_, err = c.client.Core().Nodes().UpdateStatus(node)
if err != nil {
if errors.IsConflict(err) {
glog.Warningf("Conflicting update node status for node %q, will retry soon: %v", c.nodeName, err)
if c.clock.Now().Sub(updateTime) >= timeout {
return timeoutError{node: c.nodeName, timeout: timeout}
}
continue
}
return err
}
return nil
}
}
// getEventRecorder generates a recorder for specific node name and source.
func getEventRecorder(c clientset.Interface, nodeName, source string) record.EventRecorder {
eventBroadcaster := record.NewBroadcaster()
recorder := eventBroadcaster.NewRecorder(api.EventSource{Component: source, Host: nodeName})
eventBroadcaster.StartRecordingToSink(&unversionedcore.EventSinkImpl{Interface: c.Core().Events("")})
return recorder
}
func getNodeRef(nodeName string) *api.ObjectReference {
return &api.ObjectReference{
Kind: "Node",
Name: nodeName,
UID: types.UID(nodeName),
Namespace: "",
}
}
// timeoutError is the error returned by problem client when condition update timeout.
type timeoutError struct {
node string
timeout time.Duration
}
func (e timeoutError) Error() string {
return fmt.Sprintf("update condition for node %q timeout %s", e.node, e.timeout)
}
// IsErrTimeout checks whether a given error is timeout error.
func IsErrTimeout(err error) bool {
_, ok := err.(timeoutError)
return ok
}
+259
View File
@@ -0,0 +1,259 @@
/*
Copyright 2016 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 problemclient
import (
"fmt"
"reflect"
"testing"
"time"
"k8s.io/kubernetes/pkg/api"
"k8s.io/kubernetes/pkg/api/errors"
"k8s.io/kubernetes/pkg/api/unversioned"
"k8s.io/kubernetes/pkg/client/clientset_generated/internalclientset/fake"
"k8s.io/kubernetes/pkg/client/record"
"k8s.io/kubernetes/pkg/client/testing/core"
"k8s.io/kubernetes/pkg/runtime"
"k8s.io/kubernetes/pkg/util"
)
const (
testSource = "test"
testNode = "test-node"
)
func newFakeProblemClient(fakeClient *fake.Clientset) *nodeProblemClient {
return &nodeProblemClient{
nodeName: testNode,
client: fakeClient,
clock: &util.FakeClock{},
recorders: make(map[string]record.EventRecorder),
nodeRef: getNodeRef(testNode),
}
}
func newFakeNode(conditions []api.NodeCondition) *api.Node {
node := &api.Node{}
node.Name = testNode
node.Status = api.NodeStatus{Conditions: conditions}
return node
}
type action struct {
verb string
resource string
subresource string
}
func TestSetConditions(t *testing.T) {
now := time.Now()
expectedActions := []action{
{
verb: "get",
resource: "nodes",
},
{
verb: "update",
resource: "nodes",
subresource: "status",
},
}
for _, test := range []struct {
init []api.NodeCondition
update []api.NodeCondition
expected []api.NodeCondition
}{
// Init condition with the same type should be override
{
init: []api.NodeCondition{
{
Type: "TestType",
Status: api.ConditionTrue,
},
},
update: []api.NodeCondition{
{
Type: "TestType",
Status: api.ConditionTrue,
LastTransitionTime: unversioned.NewTime(now),
Reason: "TestReason",
Message: "TestMessage",
},
},
expected: []api.NodeCondition{
{
// LastHeartbeatTime should be updated in SetConditions
Type: "TestType",
Status: api.ConditionTrue,
LastHeartbeatTime: unversioned.NewTime(now),
LastTransitionTime: unversioned.NewTime(now),
Reason: "TestReason",
Message: "TestMessage",
},
},
},
// Init condition with different type should be kept
{
init: []api.NodeCondition{
{
Type: "InitType",
Status: api.ConditionTrue,
LastTransitionTime: unversioned.NewTime(now),
Reason: "InitReason",
Message: "InitMessage",
},
},
update: []api.NodeCondition{
{
Type: "TestType",
Status: api.ConditionTrue,
LastTransitionTime: unversioned.NewTime(now),
Reason: "TestReason",
Message: "TestMessage",
},
},
expected: []api.NodeCondition{
{
Type: "InitType",
Status: api.ConditionTrue,
LastTransitionTime: unversioned.NewTime(now),
Reason: "InitReason",
Message: "InitMessage",
},
{
// LastHeartbeatTime should be updated in SetConditions
Type: "TestType",
Status: api.ConditionTrue,
LastHeartbeatTime: unversioned.NewTime(now),
LastTransitionTime: unversioned.NewTime(now),
Reason: "TestReason",
Message: "TestMessage",
},
},
},
// Condition with false status should be removed
{
init: []api.NodeCondition{
{
Type: "TestType",
Status: api.ConditionTrue,
LastHeartbeatTime: unversioned.NewTime(now),
LastTransitionTime: unversioned.NewTime(now),
Reason: "TestReason",
Message: "TestMessage",
},
},
update: []api.NodeCondition{
{
Type: "TestType",
Status: api.ConditionFalse,
},
},
expected: []api.NodeCondition{},
},
} {
fakeClient := fake.NewSimpleClientset(newFakeNode(test.init))
client := newFakeProblemClient(fakeClient)
clock := client.clock.(*util.FakeClock)
clock.SetTime(now)
client.SetConditions(test.update, 10*time.Second)
// The actions should match the expected actions
actions := fakeClient.Actions()
if len(expectedActions) != len(actions) {
t.Errorf("expected actions %+v, got %+v", expectedActions, fakeClient.Actions())
continue
}
for i, a := range actions {
if !a.Matches(expectedActions[i].verb, expectedActions[i].resource) || a.GetSubresource() != expectedActions[i].subresource {
t.Errorf("expected action %+v, got %+v", expectedActions[i], a)
}
}
// The last action should be an update
a, ok := actions[len(actions)-1].(core.UpdateAction)
if !ok {
t.Errorf("expected the last action to be update, got %+v", actions[len(actions)-1])
}
// The updated node conditions should match the expected conditions
node, ok := a.GetObject().(*api.Node)
if !ok {
t.Errorf("expected the update object to be node, got %+v", a.GetObject())
}
if !api.Semantic.DeepEqual(test.expected, node.Status.Conditions) {
t.Errorf("expected conditions %+v, got %+v", test.expected, node.Status.Conditions)
}
}
}
func TestSetConditionsError(t *testing.T) {
timeout := time.Duration(0)
node := newFakeNode([]api.NodeCondition{})
for c, test := range []struct {
errMap map[string]error
expectedErr error
}{
{
// Get error
errMap: map[string]error{"get": fmt.Errorf("get error")},
expectedErr: fmt.Errorf("get error"),
},
{
// Update error
errMap: map[string]error{"update": fmt.Errorf("update error")},
expectedErr: fmt.Errorf("update error"),
},
{
// Timeout error
errMap: map[string]error{
"update": &errors.StatusError{ErrStatus: unversioned.Status{Reason: unversioned.StatusReasonConflict}},
},
expectedErr: timeoutError{node: testNode, timeout: timeout},
},
{
// No error
errMap: map[string]error{},
expectedErr: nil,
},
} {
fakeClient := &fake.Clientset{}
client := newFakeProblemClient(fakeClient)
fakeClient.AddReactor("get", "nodes", func(action core.Action) (bool, runtime.Object, error) {
return true, node, test.errMap["get"]
})
fakeClient.AddReactor("update", "nodes", func(action core.Action) (bool, runtime.Object, error) {
return true, node, test.errMap["update"]
})
err := client.SetConditions([]api.NodeCondition{}, timeout)
if !reflect.DeepEqual(err, test.expectedErr) {
t.Errorf("case %d: expected error %v, got %v", c+1, test.expectedErr, err)
}
}
}
func TestEvent(t *testing.T) {
fakeRecorder := record.NewFakeRecorder(1)
client := newFakeProblemClient(&fake.Clientset{})
client.recorders[testSource] = fakeRecorder
client.Eventf(api.EventTypeWarning, testSource, "test reason", "test message")
expected := fmt.Sprintf("%s %s %s", api.EventTypeWarning, "test reason", "test message")
got := <-fakeRecorder.Events
if expected != got {
t.Errorf("expected event %q, got %q", expected, got)
}
}