Skip to content

Commit 227a50c

Browse files
committed
[tunnel] recover relay registration after API restore
Recreate relay liveness objects during periodic renewal when a restored API server no longer has the cached objects. Keep relay discovery available without requiring a relay process restart.
1 parent 4432875 commit 227a50c

2 files changed

Lines changed: 103 additions & 15 deletions

File tree

pkg/tunnel/controllers/relay_registrar.go

Lines changed: 33 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -139,6 +139,9 @@ func (r *RelayRegistrar) Start(ctx context.Context) error {
139139
slog.Info("Relay registrar shutting down", "relay", r.relay.Name())
140140
return ctx.Err()
141141
case <-ticker.C:
142+
if err := r.ensureRelay(ctx); err != nil {
143+
slog.Warn("Failed to restore relay registration", "relay", r.relay.Name(), "error", err)
144+
}
142145
if err := r.renewLease(ctx); err != nil {
143146
slog.Warn("Failed to renew relay lease", "relay", r.relay.Name(), "error", err)
144147
}
@@ -182,7 +185,10 @@ func (r *RelayRegistrar) ensureRelay(ctx context.Context) error {
182185
if !apierrors.IsNotFound(err) {
183186
return fmt.Errorf("failed to get relay: %w", err)
184187
}
188+
return r.createRelay(ctx)
189+
}
185190

191+
func (r *RelayRegistrar) createRelay(ctx context.Context) error {
186192
relay := &vpcv1alpha1.Relay{
187193
ObjectMeta: metav1.ObjectMeta{Name: r.relay.Name()},
188194
Spec: vpcv1alpha1.RelaySpec{
@@ -208,36 +214,48 @@ func (r *RelayRegistrar) renewLease(ctx context.Context) error {
208214
existing := &apoxycoordv1.Lease{}
209215
err := r.leaseClient.Get(ctx, key, existing)
210216
if apierrors.IsNotFound(err) {
211-
lease := &apoxycoordv1.Lease{
212-
ObjectMeta: metav1.ObjectMeta{Namespace: r.leaseNamespace, Name: key.Name},
213-
Spec: coordinationv1.LeaseSpec{
214-
HolderIdentity: ptr.To(r.relay.Name()),
215-
LeaseDurationSeconds: ptr.To(leaseDurationSeconds(r.leaseDuration)),
216-
AcquireTime: &now,
217-
RenewTime: &now,
218-
},
219-
}
220-
if err := r.leaseClient.Create(ctx, lease); err != nil {
221-
return fmt.Errorf("failed to create lease: %w", err)
222-
}
223-
return nil
217+
return r.createLease(ctx, now)
224218
}
225219
if err != nil {
226220
return fmt.Errorf("failed to get lease: %w", err)
227221
}
228222

229223
existing.Spec.HolderIdentity = ptr.To(r.relay.Name())
230-
existing.Spec.LeaseDurationSeconds = ptr.To(int32(r.leaseDuration.Seconds()))
224+
existing.Spec.LeaseDurationSeconds = ptr.To(leaseDurationSeconds(r.leaseDuration))
231225
existing.Spec.RenewTime = &now
232226
if existing.Spec.AcquireTime == nil {
233227
existing.Spec.AcquireTime = &now
234228
}
235-
if err := r.leaseClient.Update(ctx, existing); err != nil {
229+
if err := r.leaseClient.Update(ctx, existing); apierrors.IsNotFound(err) {
230+
// A restored API server can lose recent liveness objects while this
231+
// process still has them in its informer cache. Restore both objects
232+
// from the registrar's authoritative configuration.
233+
if err := r.createRelay(ctx); err != nil {
234+
return fmt.Errorf("failed to restore relay: %w", err)
235+
}
236+
return r.createLease(ctx, now)
237+
} else if err != nil {
236238
return fmt.Errorf("failed to renew lease: %w", err)
237239
}
238240
return nil
239241
}
240242

243+
func (r *RelayRegistrar) createLease(ctx context.Context, now metav1.MicroTime) error {
244+
lease := &apoxycoordv1.Lease{
245+
ObjectMeta: metav1.ObjectMeta{Namespace: r.leaseNamespace, Name: LeaseName(r.relay.Name())},
246+
Spec: coordinationv1.LeaseSpec{
247+
HolderIdentity: ptr.To(r.relay.Name()),
248+
LeaseDurationSeconds: ptr.To(leaseDurationSeconds(r.leaseDuration)),
249+
AcquireTime: &now,
250+
RenewTime: &now,
251+
},
252+
}
253+
if err := r.leaseClient.Create(ctx, lease); err != nil && !apierrors.IsAlreadyExists(err) {
254+
return fmt.Errorf("failed to create lease: %w", err)
255+
}
256+
return nil
257+
}
258+
241259
// Drain tears down the relay's control-plane presence by deleting its Lease and
242260
// Relay objects (§5). Deletion is the terminal signal consumers act on; there is
243261
// no separate ready=false write, since Drain runs synchronously at shutdown with

pkg/tunnel/controllers/relay_registrar_test.go

Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import (
1010
apierrors "k8s.io/apimachinery/pkg/api/errors"
1111
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
1212
"k8s.io/apimachinery/pkg/runtime"
13+
"k8s.io/apimachinery/pkg/runtime/schema"
1314
"sigs.k8s.io/controller-runtime/pkg/client"
1415
"sigs.k8s.io/controller-runtime/pkg/client/fake"
1516

@@ -52,6 +53,29 @@ func newRegistrar(t *testing.T, now time.Time, objs ...client.Object) (*RelayReg
5253
return r, c
5354
}
5455

56+
type deleteLeaseOnUpdateClient struct {
57+
client.Client
58+
deleteOnNextUpdate bool
59+
}
60+
61+
func (c *deleteLeaseOnUpdateClient) Update(ctx context.Context, obj client.Object, opts ...client.UpdateOption) error {
62+
if _, ok := obj.(*apoxycoordv1.Lease); ok && c.deleteOnNextUpdate {
63+
c.deleteOnNextUpdate = false
64+
lease := &apoxycoordv1.Lease{ObjectMeta: metav1.ObjectMeta{
65+
Namespace: obj.GetNamespace(),
66+
Name: obj.GetName(),
67+
}}
68+
if err := c.Client.Delete(ctx, lease); err != nil && !apierrors.IsNotFound(err) {
69+
return err
70+
}
71+
return apierrors.NewNotFound(schema.GroupResource{
72+
Group: apoxycoordv1.GroupName,
73+
Resource: "leases",
74+
}, obj.GetName())
75+
}
76+
return c.Client.Update(ctx, obj, opts...)
77+
}
78+
5579
func TestRelayRegistrarEnsureRelay(t *testing.T) {
5680
ctx := context.Background()
5781
now := time.Unix(1_700_000_000, 0)
@@ -110,6 +134,52 @@ func TestRelayRegistrarRenewLease(t *testing.T) {
110134
})
111135
}
112136

137+
func TestRelayRegistrarRecoversAfterRestoredObjectsDisappear(t *testing.T) {
138+
ctx := context.Background()
139+
now := time.Unix(1_700_000_000, 0)
140+
lease := &apoxycoordv1.Lease{
141+
ObjectMeta: metav1.ObjectMeta{Namespace: DefaultLeaseNamespace, Name: LeaseName("r0")},
142+
}
143+
_, base := newRegistrar(t, now, lease)
144+
leaseClient := &deleteLeaseOnUpdateClient{Client: base, deleteOnNextUpdate: true}
145+
r := NewRelayRegistrar(leaseClient, base, stubRelay{name: "r0"}, []string{"1.2.3.4:6081"}, nil)
146+
r.now = func() time.Time { return now }
147+
148+
require.NoError(t, r.renewLease(ctx))
149+
150+
var gotRelay vpcv1alpha1.Relay
151+
require.NoError(t, base.Get(ctx, client.ObjectKey{Name: "r0"}, &gotRelay))
152+
var gotLease apoxycoordv1.Lease
153+
require.NoError(t, base.Get(ctx, client.ObjectKey{
154+
Namespace: DefaultLeaseNamespace,
155+
Name: LeaseName("r0"),
156+
}, &gotLease))
157+
require.Equal(t, now.Unix(), gotLease.Spec.RenewTime.Unix())
158+
}
159+
160+
func TestRelayRegistrarRecreatesRelayDuringRenewal(t *testing.T) {
161+
r, c := newRegistrar(t, time.Unix(1_700_000_000, 0))
162+
r.renewInterval = 10 * time.Millisecond
163+
ctx, cancel := context.WithCancel(context.Background())
164+
done := make(chan error, 1)
165+
go func() {
166+
done <- r.Start(ctx)
167+
}()
168+
169+
require.Eventually(t, func() bool {
170+
return c.Get(context.Background(), client.ObjectKey{Name: "r0"}, &vpcv1alpha1.Relay{}) == nil
171+
}, time.Second, 10*time.Millisecond)
172+
require.NoError(t, c.Delete(context.Background(), &vpcv1alpha1.Relay{
173+
ObjectMeta: metav1.ObjectMeta{Name: "r0"},
174+
}))
175+
require.Eventually(t, func() bool {
176+
return c.Get(context.Background(), client.ObjectKey{Name: "r0"}, &vpcv1alpha1.Relay{}) == nil
177+
}, time.Second, 10*time.Millisecond)
178+
179+
cancel()
180+
require.ErrorIs(t, <-done, context.Canceled)
181+
}
182+
113183
func TestRelayRegistrarSubSecondLeaseDuration(t *testing.T) {
114184
ctx := context.Background()
115185
r, c := newRegistrar(t, time.Unix(1_700_000_000, 0))

0 commit comments

Comments
 (0)