Skip to content

Commit 6efc16e

Browse files
committed
Set the Node UID on recorded events
Recorded events refer to the node with an ObjectReference that has no UID. The client already reads the Node object at startup to check that kube-apiserver is ready, and this change stores the UID from that read. The node status patch returns the Node object, so the client stores that UID too. This change removes the TODO in problem_client.go.
1 parent eeb802f commit 6efc16e

3 files changed

Lines changed: 187 additions & 15 deletions

File tree

‎pkg/exporters/k8sexporter/k8s_exporter.go‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -123,7 +123,9 @@ func (ke *k8sExporter) startHTTPReporting(npdo *options.NodeProblemDetectorOptio
123123
func waitForAPIServerReadyWithTimeout(ctx context.Context, c problemclient.Client, npdo *options.NodeProblemDetectorOptions) error {
124124
return wait.PollUntilContextTimeout(ctx, npdo.APIServerWaitInterval, npdo.APIServerWaitTimeout, true, func(ctx context.Context) (done bool, err error) {
125125
// If NPD can get the node object from kube-apiserver, the server is
126-
// ready and the RBAC permission is set correctly.
126+
// ready and the RBAC permission is set correctly. The call also caches
127+
// the Node UID that recorded events refer to. If this check fails, the
128+
// first node status patch caches the UID instead.
127129
if _, err := c.GetNode(ctx); err != nil {
128130
klog.Errorf("Can't get node object: %v", err)
129131
return false, err

‎pkg/exporters/k8sexporter/problemclient/problem_client.go‎

Lines changed: 50 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -23,10 +23,12 @@ import (
2323
"net/url"
2424
"os"
2525
"path/filepath"
26+
"sync/atomic"
2627

2728
v1 "k8s.io/api/core/v1"
2829
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
2930
"k8s.io/apimachinery/pkg/runtime"
31+
"k8s.io/apimachinery/pkg/types"
3032
clientset "k8s.io/client-go/kubernetes"
3133
typedcorev1 "k8s.io/client-go/kubernetes/typed/core/v1"
3234
"k8s.io/client-go/tools/record"
@@ -52,11 +54,14 @@ type Client interface {
5254
}
5355

5456
type nodeProblemClient struct {
55-
nodeName string
56-
client typedcorev1.CoreV1Interface
57-
clock clock.Clock
58-
recorders map[string]record.EventRecorder
59-
nodeRef *v1.ObjectReference
57+
nodeName string
58+
client typedcorev1.CoreV1Interface
59+
clock clock.Clock
60+
recorders map[string]record.EventRecorder
61+
// nodeRef identifies the node in recorded events.
62+
nodeRef *v1.ObjectReference
63+
// cachedNodeRef holds a copy of nodeRef with the Node UID set.
64+
cachedNodeRef atomic.Pointer[v1.ObjectReference]
6065
eventNamespace string
6166
}
6267

@@ -113,7 +118,10 @@ func (c *nodeProblemClient) SetConditions(ctx context.Context, newConditions []v
113118
return true
114119
},
115120
func() error {
116-
_, err := c.client.Nodes().PatchStatus(ctx, c.nodeName, patch)
121+
node, err := c.client.Nodes().PatchStatus(ctx, c.nodeName, patch)
122+
if err == nil {
123+
c.cacheNodeRef(node.UID)
124+
}
117125
return err
118126
},
119127
)
@@ -126,13 +134,47 @@ func (c *nodeProblemClient) Eventf(eventType, source, reason, messageFmt string,
126134
recorder = getEventRecorder(c.client, c.eventNamespace, c.nodeName, source)
127135
c.recorders[source] = recorder
128136
}
129-
recorder.Eventf(c.nodeRef, eventType, reason, messageFmt, args...)
137+
recorder.Eventf(c.nodeRefWithUID(), eventType, reason, messageFmt, args...)
130138
}
131139

132140
func (c *nodeProblemClient) GetNode(ctx context.Context) (*v1.Node, error) {
133141
// To reduce the load on APIServer & etcd, we are serving GET operations from
134142
// apiserver cache (the data might be slightly delayed).
135-
return c.client.Nodes().Get(ctx, c.nodeName, metav1.GetOptions{ResourceVersion: "0"})
143+
node, err := c.client.Nodes().Get(ctx, c.nodeName, metav1.GetOptions{ResourceVersion: "0"})
144+
if err == nil {
145+
c.cacheNodeRef(node.UID)
146+
}
147+
return node, err
148+
}
149+
150+
// nodeRefWithUID returns the node reference to record events against. The
151+
// reference carries the Node UID after the first successful read of the Node
152+
// object.
153+
//
154+
// Events with an empty involvedObject.uid are invisible to consumers that
155+
// correlate by UID. `kubectl describe node` is one of them: it queries events
156+
// with the Node UID and with the node name as the UID, and shows neither an
157+
// event with no UID. See kubernetes/kubernetes issue 138524.
158+
//
159+
// The UID is resolved once. If the Node object is deleted and created again
160+
// while node-problem-detector runs, events keep the first UID until restart.
161+
func (c *nodeProblemClient) nodeRefWithUID() *v1.ObjectReference {
162+
if ref := c.cachedNodeRef.Load(); ref != nil {
163+
return ref
164+
}
165+
return c.nodeRef
166+
}
167+
168+
// cacheNodeRef stores a copy of nodeRef with the given UID set. It does not
169+
// write to nodeRef, so that concurrent readers of the shared reference stay
170+
// safe.
171+
func (c *nodeProblemClient) cacheNodeRef(uid types.UID) {
172+
if uid == "" || c.cachedNodeRef.Load() != nil {
173+
return
174+
}
175+
ref := *c.nodeRef
176+
ref.UID = uid
177+
c.cachedNodeRef.CompareAndSwap(nil, &ref)
136178
}
137179

138180
// generatePatch generates condition patch
@@ -154,7 +196,6 @@ func getEventRecorder(c typedcorev1.CoreV1Interface, namespace, nodeName, source
154196
}
155197

156198
func getNodeRef(namespace, nodeName string) *v1.ObjectReference {
157-
// TODO(random-liu): Get node to initialize the node reference
158199
return &v1.ObjectReference{
159200
APIVersion: "v1",
160201
Kind: "Node",

‎pkg/exporters/k8sexporter/problemclient/problem_client_test.go‎

Lines changed: 134 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -17,28 +17,33 @@ limitations under the License.
1717
package problemclient
1818

1919
import (
20+
"context"
2021
"encoding/json"
2122
"fmt"
23+
"net/http"
24+
"net/http/httptest"
2225
"testing"
2326
"time"
2427

2528
"github.com/stretchr/testify/assert"
2629
v1 "k8s.io/api/core/v1"
2730
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
31+
"k8s.io/apimachinery/pkg/runtime"
32+
typedcorev1 "k8s.io/client-go/kubernetes/typed/core/v1"
33+
"k8s.io/client-go/rest"
2834
"k8s.io/client-go/tools/record"
2935
testclock "k8s.io/utils/clock/testing"
3036
)
3137

3238
const (
33-
testSource = "test"
34-
testNode = "test-node"
39+
testSource = "test"
40+
testNode = "test-node"
41+
testNodeUID = "11111111-1111-1111-1111-111111111111"
3542
)
3643

3744
func newFakeProblemClient() *nodeProblemClient {
3845
return &nodeProblemClient{
39-
nodeName: testNode,
40-
// There is no proper fake for *client.Client for now
41-
// TODO(random-liu): Add test for SetConditions when we have good fake for *client.Client
46+
nodeName: testNode,
4247
clock: testclock.NewFakeClock(time.Now()),
4348
recorders: make(map[string]record.EventRecorder),
4449
nodeRef: getNodeRef("", testNode),
@@ -93,3 +98,127 @@ func TestNodeRefHasAPIVersionV1(t *testing.T) {
9398
t.Errorf("expected nodeRef.APIVersion to be 'v1', got %q", client.nodeRef.APIVersion)
9499
}
95100
}
101+
102+
func TestNodeRefWithUID(t *testing.T) {
103+
client := newFakeProblemClient()
104+
105+
if got := client.nodeRefWithUID().UID; got != "" {
106+
t.Errorf("expected no UID before the node is read, got %q", got)
107+
}
108+
109+
client.cacheNodeRef(testNodeUID)
110+
111+
if got := client.nodeRefWithUID().UID; got != testNodeUID {
112+
t.Errorf("expected UID %q, got %q", testNodeUID, got)
113+
}
114+
if got := client.nodeRef.UID; got != "" {
115+
t.Errorf("expected the shared nodeRef to keep no UID, got %q", got)
116+
}
117+
}
118+
119+
func TestCacheNodeRefIgnoresEmptyUID(t *testing.T) {
120+
client := newFakeProblemClient()
121+
122+
client.cacheNodeRef("")
123+
124+
if got := client.nodeRefWithUID().UID; got != "" {
125+
t.Errorf("expected no UID, got %q", got)
126+
}
127+
}
128+
129+
func TestCacheNodeRefKeepsFirstUID(t *testing.T) {
130+
client := newFakeProblemClient()
131+
132+
client.cacheNodeRef(testNodeUID)
133+
client.cacheNodeRef("22222222-2222-2222-2222-222222222222")
134+
135+
if got := client.nodeRefWithUID().UID; got != testNodeUID {
136+
t.Errorf("expected the first UID %q, got %q", testNodeUID, got)
137+
}
138+
}
139+
140+
// capturingRecorder records the object that Eventf reports the event against.
141+
type capturingRecorder struct {
142+
object runtime.Object
143+
}
144+
145+
func (r *capturingRecorder) Event(object runtime.Object, eventType, reason, message string) {
146+
r.object = object
147+
}
148+
149+
func (r *capturingRecorder) Eventf(object runtime.Object, eventType, reason, messageFmt string, args ...interface{}) {
150+
r.object = object
151+
}
152+
153+
func (r *capturingRecorder) AnnotatedEventf(object runtime.Object, annotations map[string]string, eventType, reason, messageFmt string, args ...interface{}) {
154+
r.object = object
155+
}
156+
157+
func TestEventfReportsNodeUID(t *testing.T) {
158+
recorder := &capturingRecorder{}
159+
client := newFakeProblemClient()
160+
client.recorders[testSource] = recorder
161+
client.cacheNodeRef(testNodeUID)
162+
163+
client.Eventf(v1.EventTypeWarning, testSource, "test reason", "test message")
164+
165+
ref, ok := recorder.object.(*v1.ObjectReference)
166+
if !ok {
167+
t.Fatalf("expected an *v1.ObjectReference, got %T", recorder.object)
168+
}
169+
if ref.UID != testNodeUID {
170+
t.Errorf("expected the reported event to carry UID %q, got %q", testNodeUID, ref.UID)
171+
}
172+
if ref.Name != testNode {
173+
t.Errorf("expected the reported event to carry name %q, got %q", testNode, ref.Name)
174+
}
175+
}
176+
177+
// newProblemClientAgainstAPI returns a client that talks to a server which
178+
// always answers with the test node.
179+
func newProblemClientAgainstAPI(t *testing.T) *nodeProblemClient {
180+
t.Helper()
181+
182+
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
183+
node := &v1.Node{ObjectMeta: metav1.ObjectMeta{Name: testNode, UID: testNodeUID}}
184+
w.Header().Set("Content-Type", "application/json")
185+
if err := json.NewEncoder(w).Encode(node); err != nil {
186+
t.Errorf("failed to encode the node: %v", err)
187+
}
188+
}))
189+
t.Cleanup(server.Close)
190+
191+
coreClient, err := typedcorev1.NewForConfig(&rest.Config{Host: server.URL})
192+
if err != nil {
193+
t.Fatalf("failed to create the core client: %v", err)
194+
}
195+
196+
client := newFakeProblemClient()
197+
client.client = coreClient
198+
return client
199+
}
200+
201+
func TestGetNodeCachesNodeUID(t *testing.T) {
202+
client := newProblemClientAgainstAPI(t)
203+
204+
if _, err := client.GetNode(context.Background()); err != nil {
205+
t.Fatalf("GetNode returned an error: %v", err)
206+
}
207+
208+
if got := client.nodeRefWithUID().UID; got != testNodeUID {
209+
t.Errorf("expected GetNode to cache UID %q, got %q", testNodeUID, got)
210+
}
211+
}
212+
213+
func TestSetConditionsCachesNodeUID(t *testing.T) {
214+
client := newProblemClientAgainstAPI(t)
215+
216+
conditions := []v1.NodeCondition{{Type: "TestType", Status: v1.ConditionTrue}}
217+
if err := client.SetConditions(context.Background(), conditions); err != nil {
218+
t.Fatalf("SetConditions returned an error: %v", err)
219+
}
220+
221+
if got := client.nodeRefWithUID().UID; got != testNodeUID {
222+
t.Errorf("expected SetConditions to cache UID %q, got %q", testNodeUID, got)
223+
}
224+
}

0 commit comments

Comments
 (0)