Skip to content

Commit 4019981

Browse files
authored
Merge pull request #86 from FabianKramm/main
fix(syncer): fix issues with syncing to another namespace
2 parents f00d085 + ce30239 commit 4019981

13 files changed

Lines changed: 178 additions & 109 deletions

File tree

‎cmd/vcluster/context/context.go‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ type VirtualClusterOptions struct {
3030
DisableSyncResources string
3131
TargetNamespace string
3232
ServiceName string
33+
ServiceNamespace string
3334
OwningStatefulSet string
3435

3536
SyncAllNodes bool
@@ -85,7 +86,7 @@ func NewControllerContext(localManager ctrl.Manager, virtualManager ctrl.Manager
8586
Context: ctx,
8687
LocalManager: localManager,
8788
VirtualManager: virtualManager,
88-
NodeServiceProvider: nodeservice.NewNodeServiceProvider(localManager.GetClient(), virtualManager.GetClient(), uncachedVirtualClient),
89+
NodeServiceProvider: nodeservice.NewNodeServiceProvider(localManager.GetClient(), virtualManager.GetClient(), uncachedVirtualClient, options.TargetNamespace),
8990
LockFactory: locks.NewDefaultLockFactory(),
9091
CacheSynced: func() {
9192
cacheSynced.Do(func() {

‎cmd/vcluster/main.go‎

Lines changed: 31 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -97,6 +97,7 @@ func NewCommand() *cobra.Command {
9797

9898
cmd.Flags().StringVar(&options.TargetNamespace, "target-namespace", "", "The namespace to run the virtual cluster in (defaults to current namespace)")
9999
cmd.Flags().StringVar(&options.ServiceName, "service-name", "vcluster", "The service name where the vcluster proxy will be available")
100+
cmd.Flags().StringVar(&options.ServiceNamespace, "service-namespace", "", "The service namespace where the vcluster proxy will be available. If empty defaults to the current namespace")
100101
cmd.Flags().StringVar(&options.OwningStatefulSet, "owning-statefulset", "", "If configured, all synced resources will have this statefulset as owner reference")
101102

102103
cmd.Flags().StringVar(&options.Suffix, "suffix", "suffix", "The suffix to append to the synced resources in the namespace")
@@ -193,16 +194,22 @@ func Execute(cobraCmd *cobra.Command, args []string, options *context.VirtualClu
193194
// set kubelet port
194195
nodeservice.KubeletTargetPort = options.Port
195196

196-
// retrieve current namespace
197-
if options.TargetNamespace == "" {
198-
currentNamespace, err := clienthelper.CurrentNamespace()
199-
if err != nil {
200-
return err
201-
}
197+
// get current namespace
198+
currentNamespace, err := clienthelper.CurrentNamespace()
199+
if err != nil {
200+
return err
201+
}
202202

203+
// ensure target namespace
204+
if options.TargetNamespace == "" {
203205
options.TargetNamespace = currentNamespace
204206
}
205207

208+
// set service namespace
209+
if options.ServiceNamespace == "" {
210+
options.ServiceNamespace = currentNamespace
211+
}
212+
206213
rawConfig, err := clientConfig.RawConfig()
207214
if err != nil {
208215
return err
@@ -346,12 +353,12 @@ func syncKubernetesService(ctx *context.ControllerContext) error {
346353
return err
347354
}
348355

349-
err = services.SyncKubernetesService(ctx.Context, localClient, virtualClient, ctx.Options.TargetNamespace, ctx.Options.ServiceName)
356+
err = services.SyncKubernetesService(ctx.Context, localClient, virtualClient, ctx.Options.ServiceNamespace, ctx.Options.ServiceName)
350357
if err != nil {
351358
return errors.Wrap(err, "sync kubernetes service")
352359
}
353360

354-
err = endpoints.SyncKubernetesServiceEndpoints(ctx.Context, localClient, virtualClient, ctx.Options.TargetNamespace, ctx.Options.ServiceName)
361+
err = endpoints.SyncKubernetesServiceEndpoints(ctx.Context, localClient, virtualClient, ctx.Options.ServiceNamespace, ctx.Options.ServiceName)
355362
if err != nil {
356363
return errors.Wrap(err, "sync kubernetes service endpoints")
357364
}
@@ -363,7 +370,9 @@ func syncKubernetesService(ctx *context.ControllerContext) error {
363370
return errors.Wrap(err, "get owning statefulset")
364371
}
365372

366-
translate.OwningStatefulSet = statefulSet
373+
if statefulSet.Namespace == ctx.Options.TargetNamespace {
374+
translate.OwningStatefulSet = statefulSet
375+
}
367376
}
368377

369378
return nil
@@ -416,14 +425,20 @@ func writeKubeConfigToSecret(ctx *context.ControllerContext, config *api.Config)
416425
}
417426

418427
// which namespace should we create the secret in?
419-
secretNamespace, err := clienthelper.CurrentNamespace()
420-
if err != nil {
421-
return err
422-
} else if ctx.Options.KubeConfigSecretNamespace != "" {
423-
secretNamespace = ctx.Options.KubeConfigSecretNamespace
424-
} else if ctx.Options.TargetNamespace != "" {
428+
secretNamespace := ctx.Options.KubeConfigSecretNamespace
429+
if secretNamespace == "" {
425430
secretNamespace = ctx.Options.TargetNamespace
426431
}
427432

428-
return kubeconfig.WriteKubeConfig(ctx.Context, ctx.LocalManager.GetClient(), ctx.Options.KubeConfigSecret, secretNamespace, config, translate.OwningStatefulSet)
433+
// we have to create a new client here, because the cached version will always say
434+
// the secret does not exist in another namespace
435+
localClient, err := client.New(ctx.LocalManager.GetConfig(), client.Options{
436+
Scheme: ctx.LocalManager.GetScheme(),
437+
Mapper: ctx.LocalManager.GetRESTMapper(),
438+
})
439+
if err != nil {
440+
return errors.Wrap(err, "create uncached client")
441+
}
442+
443+
return kubeconfig.WriteKubeConfig(ctx.Context, localClient, ctx.Options.KubeConfigSecret, secretNamespace, config, translate.OwningStatefulSet)
429444
}

‎pkg/controllers/resources/endpoints/syncer.go‎

Lines changed: 32 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -26,23 +26,42 @@ func RegisterIndices(ctx *context2.ControllerContext) error {
2626
}
2727

2828
func Register(ctx *context2.ControllerContext) error {
29+
var err error
2930
eventBroadcaster := record.NewBroadcaster()
3031
eventBroadcaster.StartRecordingToSink(&v1core.EventSinkImpl{Interface: kubernetes.NewForConfigOrDie(ctx.VirtualManager.GetConfig()).CoreV1().Events("")})
32+
33+
serviceClient := ctx.LocalManager.GetClient()
34+
if ctx.Options.ServiceNamespace != ctx.Options.TargetNamespace {
35+
serviceClient, err = client.New(ctx.LocalManager.GetConfig(), client.Options{
36+
Scheme: ctx.LocalManager.GetScheme(),
37+
Mapper: ctx.LocalManager.GetRESTMapper(),
38+
})
39+
if err != nil {
40+
return errors.Wrap(err, "create uncached client")
41+
}
42+
}
43+
3144
return generic.RegisterSyncer(ctx, &syncer{
32-
eventRecoder: eventBroadcaster.NewRecorder(ctx.VirtualManager.GetScheme(), corev1.EventSource{Component: "endpoints-syncer"}),
33-
targetNamespace: ctx.Options.TargetNamespace,
34-
serviceName: ctx.Options.ServiceName,
35-
localClient: ctx.LocalManager.GetClient(),
36-
virtualClient: ctx.VirtualManager.GetClient(),
45+
eventRecoder: eventBroadcaster.NewRecorder(ctx.VirtualManager.GetScheme(), corev1.EventSource{Component: "endpoints-syncer"}),
46+
targetNamespace: ctx.Options.TargetNamespace,
47+
serviceName: ctx.Options.ServiceName,
48+
serviceNamespace: ctx.Options.ServiceNamespace,
49+
serviceClient: serviceClient,
50+
localClient: ctx.LocalManager.GetClient(),
51+
virtualClient: ctx.VirtualManager.GetClient(),
3752
}, "endpoints", generic.RegisterSyncerOptions{})
3853
}
3954

4055
type syncer struct {
4156
eventRecoder record.EventRecorder
4257
targetNamespace string
43-
serviceName string
44-
localClient client.Client
45-
virtualClient client.Client
58+
59+
serviceName string
60+
serviceNamespace string
61+
serviceClient client.Client
62+
63+
localClient client.Client
64+
virtualClient client.Client
4665
}
4766

4867
func (s *syncer) New() client.Object {
@@ -191,8 +210,8 @@ func (s *syncer) BackwardUpdateNeeded(pObj client.Object, vObj client.Object) (b
191210

192211
func (s *syncer) BackwardStart(ctx context.Context, req ctrl.Request) (bool, error) {
193212
// sync the kubernetes service
194-
if req.Name == s.serviceName && req.Namespace == s.targetNamespace {
195-
return true, SyncKubernetesServiceEndpoints(ctx, s.virtualClient, s.localClient, s.targetNamespace, s.serviceName)
213+
if req.Name == s.serviceName && req.Namespace == s.serviceNamespace {
214+
return true, SyncKubernetesServiceEndpoints(ctx, s.virtualClient, s.serviceClient, s.serviceNamespace, s.serviceName)
196215
}
197216

198217
return false, nil
@@ -205,7 +224,7 @@ func (s *syncer) BackwardEnd() {
205224
func (s *syncer) ForwardStart(ctx context.Context, req ctrl.Request) (bool, error) {
206225
// dont do anything for the kubernetes service
207226
if req.Name == "kubernetes" && req.Namespace == "default" {
208-
return true, SyncKubernetesServiceEndpoints(ctx, s.virtualClient, s.localClient, s.targetNamespace, s.serviceName)
227+
return true, SyncKubernetesServiceEndpoints(ctx, s.virtualClient, s.serviceClient, s.serviceNamespace, s.serviceName)
209228
}
210229

211230
return false, nil
@@ -215,11 +234,11 @@ func (s *syncer) ForwardEnd() {
215234

216235
}
217236

218-
func SyncKubernetesServiceEndpoints(ctx context.Context, virtualClient client.Client, localClient client.Client, targetNamespace, serviceName string) error {
237+
func SyncKubernetesServiceEndpoints(ctx context.Context, virtualClient client.Client, localClient client.Client, serviceNamespace, serviceName string) error {
219238
// get physical service endpoints
220239
pObj := &corev1.Endpoints{}
221240
err := localClient.Get(ctx, types.NamespacedName{
222-
Namespace: targetNamespace,
241+
Namespace: serviceNamespace,
223242
Name: serviceName,
224243
}, pObj)
225244
if err != nil {

‎pkg/controllers/resources/endpoints/syncer_test.go‎

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -15,10 +15,12 @@ import (
1515

1616
func newFakeSyncer(pClient *testingutil.FakeIndexClient, vClient *testingutil.FakeIndexClient) *syncer {
1717
return &syncer{
18-
eventRecoder: &testingutil.FakeEventRecorder{},
19-
targetNamespace: "test",
20-
virtualClient: vClient,
21-
localClient: pClient,
18+
eventRecoder: &testingutil.FakeEventRecorder{},
19+
targetNamespace: "test",
20+
serviceNamespace: "test",
21+
serviceClient: pClient,
22+
virtualClient: vClient,
23+
localClient: pClient,
2224
}
2325
}
2426

‎pkg/controllers/resources/nodes/nodeservice/node_service.go‎

Lines changed: 8 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -39,11 +39,12 @@ type NodeServiceProvider interface {
3939
GetNodeIP(ctx context.Context, name types.NamespacedName) (string, error)
4040
}
4141

42-
func NewNodeServiceProvider(localClient client.Client, virtualClient client.Client, uncachedVirtualClient client.Client) NodeServiceProvider {
42+
func NewNodeServiceProvider(localClient client.Client, virtualClient client.Client, uncachedVirtualClient client.Client, targetNamespace string) NodeServiceProvider {
4343
return &nodeServiceProvider{
4444
localClient: localClient,
4545
virtualClient: virtualClient,
4646
uncachedVirtualClient: uncachedVirtualClient,
47+
targetNamespace: targetNamespace,
4748
}
4849
}
4950

@@ -52,7 +53,8 @@ type nodeServiceProvider struct {
5253
virtualClient client.Client
5354
uncachedVirtualClient client.Client
5455

55-
serviceMutex sync.Mutex
56+
targetNamespace string
57+
serviceMutex sync.Mutex
5658
}
5759

5860
func (n *nodeServiceProvider) Start(ctx context.Context) {
@@ -68,13 +70,8 @@ func (n *nodeServiceProvider) cleanupNodeServices(ctx context.Context) error {
6870
n.serviceMutex.Lock()
6971
defer n.serviceMutex.Unlock()
7072

71-
namespace, err := clienthelper.CurrentNamespace()
72-
if err != nil {
73-
return errors.Wrap(err, "get current namespace")
74-
}
75-
7673
serviceList := &corev1.ServiceList{}
77-
err = n.localClient.List(ctx, serviceList, client.InNamespace(namespace), client.MatchingLabels{
74+
err := n.localClient.List(ctx, serviceList, client.InNamespace(n.targetNamespace), client.MatchingLabels{
7875
ServiceClusterLabel: translate.Suffix,
7976
})
8077
if err != nil {
@@ -124,13 +121,8 @@ func (n *nodeServiceProvider) Unlock() {
124121
}
125122

126123
func (n *nodeServiceProvider) GetNodeIP(ctx context.Context, name types.NamespacedName) (string, error) {
127-
namespace, err := clienthelper.CurrentNamespace()
128-
if err != nil {
129-
return "", errors.Wrap(err, "get current namespace")
130-
}
131-
132124
serviceList := &corev1.ServiceList{}
133-
err = n.localClient.List(ctx, serviceList, client.InNamespace(namespace), client.MatchingLabels{
125+
err := n.localClient.List(ctx, serviceList, client.InNamespace(n.targetNamespace), client.MatchingLabels{
134126
ServiceClusterLabel: translate.Suffix,
135127
ServiceNodeLabel: name.Name,
136128
})
@@ -148,7 +140,7 @@ func (n *nodeServiceProvider) GetNodeIP(ctx context.Context, name types.Namespac
148140

149141
// find out the labels to select ourself
150142
pod := &corev1.Pod{}
151-
err = n.localClient.Get(ctx, types.NamespacedName{Name: podName, Namespace: namespace}, pod)
143+
err = n.localClient.Get(ctx, types.NamespacedName{Name: podName, Namespace: n.targetNamespace}, pod)
152144
if err != nil {
153145
return "", errors.Wrap(err, "get pod")
154146
} else if len(pod.Labels) == 0 {
@@ -168,7 +160,7 @@ func (n *nodeServiceProvider) GetNodeIP(ctx context.Context, name types.Namespac
168160
// create the new service
169161
nodeService := &corev1.Service{
170162
ObjectMeta: metav1.ObjectMeta{
171-
Namespace: namespace,
163+
Namespace: n.targetNamespace,
172164
GenerateName: translate.SafeConcatGenerateName(translate.Suffix, "node") + "-",
173165
Labels: map[string]string{
174166
ServiceClusterLabel: translate.Suffix,

‎pkg/controllers/resources/pods/syncer.go‎

Lines changed: 18 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -69,11 +69,25 @@ func Register(ctx *context2.ControllerContext) error {
6969
return errors.Wrap(err, "create pod translator")
7070
}
7171

72+
// service client
73+
serviceClient := ctx.LocalManager.GetClient()
74+
if ctx.Options.ServiceNamespace != ctx.Options.TargetNamespace {
75+
serviceClient, err = client.New(ctx.LocalManager.GetConfig(), client.Options{
76+
Scheme: ctx.LocalManager.GetScheme(),
77+
Mapper: ctx.LocalManager.GetRESTMapper(),
78+
})
79+
if err != nil {
80+
return errors.Wrap(err, "create uncached client")
81+
}
82+
}
83+
7284
return generic.RegisterSyncer(ctx, &syncer{
7385
sharedNodesMutex: ctx.LockFactory.GetLock("nodes-controller"),
7486
eventRecoder: eventBroadcaster.NewRecorder(ctx.VirtualManager.GetScheme(), corev1.EventSource{Component: "pod-syncer"}),
7587
targetNamespace: ctx.Options.TargetNamespace,
7688
serviceName: ctx.Options.ServiceName,
89+
serviceNamespace: ctx.Options.ServiceNamespace,
90+
serviceClient: serviceClient,
7791
localClient: ctx.LocalManager.GetClient(),
7892
virtualClient: ctx.VirtualManager.GetClient(),
7993
virtualClusterClient: virtualClusterClient,
@@ -93,6 +107,8 @@ type syncer struct {
93107
eventRecoder record.EventRecorder
94108
targetNamespace string
95109
serviceName string
110+
serviceNamespace string
111+
serviceClient client.Client
96112
podTranslator translatepods.Translator
97113
localClient client.Client
98114
virtualClient client.Client
@@ -266,9 +282,9 @@ func (s *syncer) translatePod(vPod *corev1.Pod) (*corev1.Pod, error) {
266282

267283
func (s *syncer) findKubernetesIP() (string, error) {
268284
pService := &corev1.Service{}
269-
err := s.localClient.Get(context.TODO(), types.NamespacedName{
285+
err := s.serviceClient.Get(context.TODO(), types.NamespacedName{
270286
Name: s.serviceName,
271-
Namespace: s.targetNamespace,
287+
Namespace: s.serviceNamespace,
272288
}, pService)
273289
if err != nil {
274290
return "", err

0 commit comments

Comments
 (0)