Repository navigation
[HYPERSHELL-177] Implement GatewayNetwork reconciliation #246
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,219 @@ | ||
| package reconciler | ||
|
|
||
| import ( | ||
| "context" | ||
| "strings" | ||
| "testing" | ||
|
|
||
| pb "github.com/openshift-online/hypershell/components/api-server/pkg/api/grpc/hypershell/v1" | ||
| "github.com/openshift-online/hypershell/components/control-plane/internal/watcher" | ||
| "google.golang.org/grpc" | ||
| "google.golang.org/grpc/codes" | ||
| "google.golang.org/grpc/status" | ||
| ) | ||
|
|
||
| // fakeNetworkClient records UpdateGatewayNetwork calls and can inject an error. | ||
| type fakeNetworkClient struct { | ||
| pb.GatewayNetworkServiceClient | ||
| updates []*pb.UpdateGatewayNetworkRequest | ||
| updateErr error | ||
| } | ||
|
|
||
| func (f *fakeNetworkClient) UpdateGatewayNetwork(ctx context.Context, in *pb.UpdateGatewayNetworkRequest, opts ...grpc.CallOption) (*pb.UpdateGatewayNetworkResponse, error) { | ||
| f.updates = append(f.updates, in) | ||
| if f.updateErr != nil { | ||
| return nil, f.updateErr | ||
| } | ||
| return &pb.UpdateGatewayNetworkResponse{}, nil | ||
| } | ||
|
|
||
| // fakeNetworkGatewayClient resolves GetGateway against a fixed hub inventory and | ||
| // can inject a lookup error to simulate not-found or transient failures. | ||
| type fakeNetworkGatewayClient struct { | ||
| pb.GatewayServiceClient | ||
| existing map[string]bool | ||
| getErr error | ||
| } | ||
|
|
||
| func (f *fakeNetworkGatewayClient) GetGateway(ctx context.Context, in *pb.GetGatewayRequest, opts ...grpc.CallOption) (*pb.GetGatewayResponse, error) { | ||
| if f.getErr != nil { | ||
| return nil, f.getErr | ||
| } | ||
| if f.existing[in.Id] { | ||
| return &pb.GetGatewayResponse{Gateway: &pb.Gateway{Metadata: &pb.ObjectReference{Id: in.Id}}}, nil | ||
| } | ||
| return nil, status.Error(codes.NotFound, "gateway not found") | ||
| } | ||
|
|
||
| func newTestNetworkReconciler(gw pb.GatewayServiceClient, net pb.GatewayNetworkServiceClient) *GatewayNetworkReconciler { | ||
| return &GatewayNetworkReconciler{ | ||
| active: make(map[string]struct{}), | ||
| gateways: gw, | ||
| networks: net, | ||
| } | ||
| } | ||
|
|
||
| func networkEvent(t watcher.EventType, id, topology, hubID, status string) watcher.Event[*pb.GatewayNetwork] { | ||
| net := &pb.GatewayNetwork{ | ||
| Metadata: &pb.ObjectReference{Id: id}, | ||
| Name: "net-" + id, | ||
| } | ||
| if topology != "" { | ||
| net.Topology = &topology | ||
| } | ||
| if hubID != "" { | ||
| net.HubGatewayId = &hubID | ||
| } | ||
| if status != "" { | ||
| net.Status = &status | ||
| } | ||
| return watcher.Event[*pb.GatewayNetwork]{Type: t, ResourceID: id, Resource: net} | ||
| } | ||
|
|
||
| func TestGatewayNetwork_ValidHubSpokeSetsValid(t *testing.T) { | ||
| net := &fakeNetworkClient{} | ||
| gw := &fakeNetworkGatewayClient{existing: map[string]bool{"hub1": true}} | ||
| r := newTestNetworkReconciler(gw, net) | ||
|
|
||
| if err := r.Handle(context.Background(), networkEvent(watcher.EventCreated, "n1", "hub-spoke", "hub1", "")); err != nil { | ||
| t.Fatalf("unexpected error: %v", err) | ||
| } | ||
| if len(net.updates) != 1 || net.updates[0].GetStatus() != networkStatusValid { | ||
| t.Fatalf("expected status Valid, got updates=%v", net.updates) | ||
| } | ||
| } | ||
|
|
||
| func TestGatewayNetwork_ValidMeshSetsValid(t *testing.T) { | ||
| net := &fakeNetworkClient{} | ||
| r := newTestNetworkReconciler(&fakeNetworkGatewayClient{}, net) | ||
|
|
||
| if err := r.Handle(context.Background(), networkEvent(watcher.EventCreated, "n1", "mesh", "", "")); err != nil { | ||
| t.Fatalf("unexpected error: %v", err) | ||
| } | ||
| if len(net.updates) != 1 || net.updates[0].GetStatus() != networkStatusValid { | ||
| t.Fatalf("expected status Valid, got updates=%v", net.updates) | ||
| } | ||
| } | ||
|
|
||
| func TestGatewayNetwork_UnrecognizedTopologyIsInvalid(t *testing.T) { | ||
| net := &fakeNetworkClient{} | ||
| r := newTestNetworkReconciler(&fakeNetworkGatewayClient{}, net) | ||
|
|
||
| if err := r.Handle(context.Background(), networkEvent(watcher.EventCreated, "n1", "ring", "", "")); err != nil { | ||
| t.Fatalf("unexpected error: %v", err) | ||
| } | ||
| if len(net.updates) != 1 || !strings.HasPrefix(net.updates[0].GetStatus(), networkStatusInvalid) { | ||
| t.Fatalf("expected Invalid status, got updates=%v", net.updates) | ||
| } | ||
| if !strings.Contains(net.updates[0].GetStatus(), "ring") { | ||
| t.Fatalf("expected reason to mention the unrecognized topology, got %q", net.updates[0].GetStatus()) | ||
| } | ||
| } | ||
|
|
||
| func TestGatewayNetwork_HubSpokeWithoutHubIsInvalid(t *testing.T) { | ||
| net := &fakeNetworkClient{} | ||
| r := newTestNetworkReconciler(&fakeNetworkGatewayClient{}, net) | ||
|
|
||
| if err := r.Handle(context.Background(), networkEvent(watcher.EventCreated, "n1", "hub-spoke", "", "")); err != nil { | ||
| t.Fatalf("unexpected error: %v", err) | ||
| } | ||
| if len(net.updates) != 1 || !strings.HasPrefix(net.updates[0].GetStatus(), networkStatusInvalid) { | ||
| t.Fatalf("expected Invalid status, got updates=%v", net.updates) | ||
| } | ||
| } | ||
|
|
||
| func TestGatewayNetwork_DanglingHubIsInvalid(t *testing.T) { | ||
| net := &fakeNetworkClient{} | ||
| // hub inventory is empty, so GetGateway returns NotFound for hub1. | ||
| gw := &fakeNetworkGatewayClient{existing: map[string]bool{}} | ||
| r := newTestNetworkReconciler(gw, net) | ||
|
|
||
| if err := r.Handle(context.Background(), networkEvent(watcher.EventCreated, "n1", "hub-spoke", "hub1", "")); err != nil { | ||
| t.Fatalf("unexpected error: %v", err) | ||
| } | ||
| if len(net.updates) != 1 || !strings.HasPrefix(net.updates[0].GetStatus(), networkStatusInvalid) { | ||
| t.Fatalf("expected Invalid status, got updates=%v", net.updates) | ||
| } | ||
| if !strings.Contains(net.updates[0].GetStatus(), "hub1") { | ||
| t.Fatalf("expected reason to mention the missing hub gateway, got %q", net.updates[0].GetStatus()) | ||
| } | ||
| } | ||
|
|
||
| func TestGatewayNetwork_NoRedundantStatusWrite(t *testing.T) { | ||
| net := &fakeNetworkClient{} | ||
| gw := &fakeNetworkGatewayClient{existing: map[string]bool{"hub1": true}} | ||
| r := newTestNetworkReconciler(gw, net) | ||
|
|
||
| // Persisted status already equals the reconciled outcome. | ||
| if err := r.Handle(context.Background(), networkEvent(watcher.EventUpdated, "n1", "hub-spoke", "hub1", networkStatusValid)); err != nil { | ||
| t.Fatalf("unexpected error: %v", err) | ||
| } | ||
| if len(net.updates) != 0 { | ||
| t.Fatalf("expected no status write when unchanged, got %v", net.updates) | ||
| } | ||
| } | ||
|
|
||
| func TestGatewayNetwork_NoRedundantInvalidStatusWrite(t *testing.T) { | ||
| net := &fakeNetworkClient{} | ||
| r := newTestNetworkReconciler(&fakeNetworkGatewayClient{}, net) | ||
|
|
||
| // Persisted status already equals the recomputed Invalid outcome (same reason). | ||
| persisted := networkStatusInvalid + ": unrecognized topology \"ring\"" | ||
| if err := r.Handle(context.Background(), networkEvent(watcher.EventUpdated, "n1", "ring", "", persisted)); err != nil { | ||
| t.Fatalf("unexpected error: %v", err) | ||
| } | ||
| if len(net.updates) != 0 { | ||
| t.Fatalf("expected no status write when Invalid status unchanged, got %v", net.updates) | ||
| } | ||
| } | ||
|
|
||
| func TestGatewayNetwork_DeleteIsNoOp(t *testing.T) { | ||
| net := &fakeNetworkClient{} | ||
| r := newTestNetworkReconciler(&fakeNetworkGatewayClient{}, net) | ||
|
|
||
| if err := r.Handle(context.Background(), networkEvent(watcher.EventDeleted, "n1", "hub-spoke", "hub1", "")); err != nil { | ||
| t.Fatalf("unexpected error: %v", err) | ||
| } | ||
| if len(net.updates) != 0 { | ||
| t.Fatalf("expected no status write on delete, got %v", net.updates) | ||
| } | ||
| } | ||
|
|
||
| func TestGatewayNetwork_NilResourceIsNoOp(t *testing.T) { | ||
| net := &fakeNetworkClient{} | ||
| r := newTestNetworkReconciler(&fakeNetworkGatewayClient{}, net) | ||
|
|
||
| ev := watcher.Event[*pb.GatewayNetwork]{Type: watcher.EventCreated, ResourceID: "n1", Resource: nil} | ||
| if err := r.Handle(context.Background(), ev); err != nil { | ||
| t.Fatalf("unexpected error: %v", err) | ||
| } | ||
| if len(net.updates) != 0 { | ||
| t.Fatalf("expected no status write for nil resource, got %v", net.updates) | ||
| } | ||
| } | ||
|
|
||
| func TestGatewayNetwork_TransientHubLookupSurfacesAsError(t *testing.T) { | ||
| net := &fakeNetworkClient{} | ||
| gw := &fakeNetworkGatewayClient{getErr: status.Error(codes.Unavailable, "hub lookup down")} | ||
| r := newTestNetworkReconciler(gw, net) | ||
|
|
||
| err := r.Handle(context.Background(), networkEvent(watcher.EventCreated, "n1", "hub-spoke", "hub1", "")) | ||
| if err == nil { | ||
| t.Fatalf("expected transient hub lookup failure to return an error") | ||
| } | ||
| // The network must not be settled to Invalid on account of a transient failure. | ||
| if len(net.updates) != 0 { | ||
| t.Fatalf("expected no status write on transient failure, got %v", net.updates) | ||
| } | ||
| } | ||
|
|
||
| func TestGatewayNetwork_StatusWriteFailureSurfacesAsError(t *testing.T) { | ||
| net := &fakeNetworkClient{updateErr: status.Error(codes.Unavailable, "api down")} | ||
| gw := &fakeNetworkGatewayClient{existing: map[string]bool{"hub1": true}} | ||
| r := newTestNetworkReconciler(gw, net) | ||
|
|
||
| err := r.Handle(context.Background(), networkEvent(watcher.EventCreated, "n1", "hub-spoke", "hub1", "")) | ||
| if err == nil { | ||
| t.Fatalf("expected status write failure to return an error") | ||
| } | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -2417,13 +2417,49 @@ func (r *StubGatewayReconciler) Handle(ctx context.Context, event watcher.Event[ | |
| return nil | ||
| } | ||
|
|
||
| // GatewayNetwork topology vocabulary. See | ||
| // specs/platform/gateway-network-reconciliation.spec.md. | ||
| const ( | ||
| networkTopologyMesh = "mesh" | ||
| networkTopologyHubSpoke = "hub-spoke" | ||
| ) | ||
|
|
||
| // GatewayNetwork control-plane-owned status values. A network owns no Kubernetes | ||
| // resources in this scope, so status reflects configuration validity only, not | ||
| // provisioned connectivity. | ||
| const ( | ||
| networkStatusValid = "Valid" | ||
| networkStatusInvalid = "Invalid" | ||
| ) | ||
|
|
||
| // GatewayNetworkReconciler reconciles GatewayNetwork resources. A network owns no | ||
| // Kubernetes resources in this scope, so reconciliation means: validate the | ||
| // network's topology vocabulary and topology/hub coherence, validate that a | ||
| // designated hub_gateway_id references an existing Gateway, and write a | ||
| // deterministic status back to the API server. Applying real gateway-to-gateway | ||
| // connectivity (mesh/tunnel provisioning) is future work owned by a sibling spec | ||
| // once product defines the network membership model and connectivity technology; | ||
| // this reconciler only records whether the declared configuration is well-formed. | ||
| type GatewayNetworkReconciler struct { | ||
| mu sync.Mutex | ||
| active map[string]struct{} | ||
|
|
||
| gateways pb.GatewayServiceClient | ||
| networks pb.GatewayNetworkServiceClient | ||
| } | ||
|
|
||
| func NewGatewayNetworkReconciler() *GatewayNetworkReconciler { | ||
| return &GatewayNetworkReconciler{active: make(map[string]struct{})} | ||
| // NewGatewayNetworkReconciler builds the network reconciler. conn is the API | ||
| // server gRPC connection used to look up the designated hub gateway and to write | ||
| // network status back. conn may be nil (e.g. in unit tests), in which case the | ||
| // hub existence check and status write-back are skipped but the rest of | ||
| // validation still runs. | ||
| func NewGatewayNetworkReconciler(conn *grpc.ClientConn) *GatewayNetworkReconciler { | ||
| r := &GatewayNetworkReconciler{active: make(map[string]struct{})} | ||
| if conn != nil { | ||
| r.gateways = pb.NewGatewayServiceClient(conn) | ||
| r.networks = pb.NewGatewayNetworkServiceClient(conn) | ||
| } | ||
| return r | ||
| } | ||
|
|
||
| func (r *GatewayNetworkReconciler) Handle(ctx context.Context, event watcher.Event[*pb.GatewayNetwork]) error { | ||
|
|
@@ -2441,8 +2477,102 @@ func (r *GatewayNetworkReconciler) Handle(ctx context.Context, event watcher.Eve | |
| }() | ||
|
|
||
| _, endSpan := cpotel.StartReconcileSpan(ctx, "GatewayNetwork", event.Type.String(), event.Resource.GetMetadata().GetTraceparent()) | ||
| defer func() { endSpan(nil) }() | ||
| var reconcileErr error | ||
| defer func() { endSpan(reconcileErr) }() | ||
|
|
||
| log.Printf("INFO reconciling GatewayNetwork %s (event=%d)", event.ResourceID, event.Type) | ||
| // A network owns no cluster resources, so a delete is a terminal, idempotent | ||
| // no-op with respect to Kubernetes: gateways designated by the network are | ||
| // left untouched. | ||
| if event.Type == watcher.EventDeleted { | ||
| log.Printf("INFO gateway network %s deleted; no cluster resources to remove", event.ResourceID) | ||
| return nil | ||
| } | ||
|
|
||
| net := event.Resource | ||
| if net == nil { | ||
| log.Printf("WARN gateway network event %s has nil resource, skipping", event.ResourceID) | ||
| return nil | ||
| } | ||
|
|
||
| // Validate the declared configuration. A transient dependency failure (e.g. a | ||
| // transient hub lookup error) is returned so the failure is surfaced (logged | ||
| // by the watch loop) rather than silently swallowed or settled to a misleading | ||
| // Invalid. The network watch is inline log-only (no reconcile queue) and does | ||
| // not replay state on reconnect, so a surfaced error re-converges only when the | ||
| // network is next mutated, not automatically. | ||
| desiredStatus, retryErr := r.validate(ctx, net) | ||
| if retryErr != nil { | ||
| reconcileErr = fmt.Errorf("validate gateway network %s: %w", event.ResourceID, retryErr) | ||
| return reconcileErr | ||
| } | ||
|
|
||
| // Deterministic, idempotent status write-back: only update when the persisted | ||
| // status differs from the reconciled outcome. | ||
| if net.GetStatus() != desiredStatus { | ||
| if err := r.updateStatus(ctx, event.ResourceID, desiredStatus); err != nil { | ||
| reconcileErr = fmt.Errorf("update gateway network %s status: %w", event.ResourceID, err) | ||
| return reconcileErr | ||
| } | ||
| } | ||
| return nil | ||
| } | ||
|
|
||
| // validate applies the network's structural and referential coherence rules and | ||
| // returns the deterministic desired status (networkStatusValid, or | ||
| // "networkStatusInvalid: reason"). It returns a non-nil error only for a | ||
| // transient dependency failure that should be surfaced rather than swallowed; a | ||
| // definitive not-found for the hub gateway is a deterministic Invalid, not an | ||
| // error. | ||
| func (r *GatewayNetworkReconciler) validate(ctx context.Context, net *pb.GatewayNetwork) (string, error) { | ||
| invalid := func(reason string) string { | ||
| return fmt.Sprintf("%s: %s", networkStatusInvalid, reason) | ||
| } | ||
|
|
||
| topology := net.GetTopology() | ||
| switch topology { | ||
| case "": | ||
| return invalid("topology is required"), nil | ||
| case networkTopologyMesh, networkTopologyHubSpoke: | ||
| // recognized | ||
| default: | ||
| return invalid(fmt.Sprintf("unrecognized topology %q", topology)), nil | ||
| } | ||
|
|
||
| hubID := net.GetHubGatewayId() | ||
| if topology == networkTopologyHubSpoke && hubID == "" { | ||
| return invalid("hub-spoke network requires a hub_gateway_id"), nil | ||
| } | ||
|
|
||
| if hubID != "" { | ||
| // A configured hub must reference an existing Gateway. Skip the lookup when | ||
| // no gateway client is configured (started without an API-server gRPC | ||
| // connection, e.g. in unit tests). | ||
| if r.gateways == nil { | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Minor: when |
||
| return networkStatusValid, nil | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. When
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. [Info] When |
||
| } | ||
| _, err := r.gateways.GetGateway(ctx, &pb.GetGatewayRequest{Id: hubID}) | ||
| if err != nil { | ||
| if status.Code(err) == codes.NotFound { | ||
| return invalid(fmt.Sprintf("hub gateway %q does not exist", hubID)), nil | ||
| } | ||
| // Transient failure: surface as an error rather than settle to a | ||
| // misleading Invalid. | ||
| return "", err | ||
| } | ||
| } | ||
|
|
||
| return networkStatusValid, nil | ||
| } | ||
|
|
||
| // updateStatus writes the network's reconciled status back to the API server. It | ||
| // is a no-op when the network client is not configured. | ||
| func (r *GatewayNetworkReconciler) updateStatus(ctx context.Context, id, desired string) error { | ||
| if r.networks == nil { | ||
| return nil | ||
| } | ||
| _, err := r.networks.UpdateGatewayNetwork(ctx, &pb.UpdateGatewayNetworkRequest{ | ||
| Id: id, | ||
| Status: &desired, | ||
| }) | ||
| return err | ||
| } | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
[Minor] Transient failures are correctly surfaced here instead of being swallowed or settled to a misleading
Invalid- good. The caveat:WatchGatewayNetworksonly logs this returned error (watcher.go:840), and there is no initial list or periodic resync for networks, so a network stuck on a transient hub-lookup or status-write failure stays stale until it is next mutated. The spec documents this as intentional, but it is the exact assumption a separate open PR's reconciliation contract would override (see the Cross-PR coordination section).