Skip to content

Commit edc0ed9

Browse files
authored
make the spec update and status update independently (#389)
Signed-off-by: Wei Liu <liuweixa@redhat.com>
1 parent aee02f9 commit edc0ed9

5 files changed

Lines changed: 145 additions & 6 deletions

File tree

Makefile

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -245,7 +245,7 @@ test-integration-mqtt:
245245
.PHONY: test-integration-mqtt
246246

247247
test-integration-grpc:
248-
BROKER=grpc MAESTRO_ENV=testing gotestsum --jsonfile-timing-events=$(grpc_integration_test_json_output) --format $(TEST_SUMMARY_FORMAT) -- -p 1 -ldflags -s -v -timeout 1h $(TESTFLAGS) \
248+
BROKER=grpc MAESTRO_ENV=testing gotestsum --jsonfile-timing-events=$(grpc_integration_test_json_output) --format $(TEST_SUMMARY_FORMAT) -- -count=1 -p 1 -ldflags -s -v -timeout 1h $(TESTFLAGS) \
249249
./test/integration
250250
.PHONY: test-integration-grpc
251251

@@ -428,7 +428,7 @@ db/teardown:
428428

429429
.PHONY: mqtt/prepare
430430
mqtt/prepare:
431-
@echo $(shell LC_CTYPE=C tr -dc 'a-zA-Z0-9' < /dev/urandom | head -c 13) > $(mqtt_password_file)
431+
@openssl rand -base64 13 | tr -dc 'a-zA-Z0-9' | head -c 13 > $(mqtt_password_file)
432432

433433
.PHONY: mqtt/setup
434434
mqtt/setup: mqtt/prepare

pkg/dao/mocks/resource.go

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,10 @@ func (d *resourceDaoMock) Update(ctx context.Context, resource *api.Resource) (*
3939
return nil, errors.NotImplemented("Resource").AsError()
4040
}
4141

42+
func (d *resourceDaoMock) UpdateStatus(ctx context.Context, resource *api.Resource) (*api.Resource, error) {
43+
return nil, errors.NotImplemented("Resource").AsError()
44+
}
45+
4246
func (d *resourceDaoMock) Delete(ctx context.Context, id string, unscoped bool) error {
4347
return errors.NotImplemented("Resource").AsError()
4448
}

pkg/dao/resource.go

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ type ResourceDao interface {
1313
Get(ctx context.Context, id string) (*api.Resource, error)
1414
Create(ctx context.Context, resource *api.Resource) (*api.Resource, error)
1515
Update(ctx context.Context, resource *api.Resource) (*api.Resource, error)
16+
UpdateStatus(ctx context.Context, resource *api.Resource) (*api.Resource, error)
1617
Delete(ctx context.Context, id string, unscoped bool) error
1718
FindByIDs(ctx context.Context, ids []string) (api.ResourceList, error)
1819
FindBySource(ctx context.Context, source string) (api.ResourceList, error)
@@ -51,7 +52,27 @@ func (d *sqlResourceDao) Create(ctx context.Context, resource *api.Resource) (*a
5152

5253
func (d *sqlResourceDao) Update(ctx context.Context, resource *api.Resource) (*api.Resource, error) {
5354
g2 := (*d.sessionFactory).New(ctx)
54-
if err := g2.Unscoped().Omit(clause.Associations).Updates(resource).Error; err != nil {
55+
if err := g2.Unscoped().Omit(clause.Associations).
56+
Where("id = ?", resource.ID).
57+
Select("version", "payload").
58+
Updates(api.Resource{
59+
Version: resource.Version,
60+
Payload: resource.Payload,
61+
}).Error; err != nil {
62+
db.MarkForRollback(ctx, err)
63+
return nil, err
64+
}
65+
return resource, nil
66+
}
67+
68+
func (d *sqlResourceDao) UpdateStatus(ctx context.Context, resource *api.Resource) (*api.Resource, error) {
69+
g2 := (*d.sessionFactory).New(ctx)
70+
if err := g2.Unscoped().Omit(clause.Associations).
71+
Where("id = ?", resource.ID).
72+
Select("status").
73+
Updates(api.Resource{
74+
Status: resource.Status,
75+
}).Error; err != nil {
5576
db.MarkForRollback(ctx, err)
5677
return nil, err
5778
}

pkg/services/resource.go

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -177,6 +177,9 @@ func (s *sqlResourceService) UpdateStatus(ctx context.Context, resource *api.Res
177177
}
178178

179179
// Make sure the requested resource version is consistent with its database version.
180+
// If they do not match, the event is stale and can be safely ignored.
181+
// The manifestwork status reflects observed generations, ensuring eventual consistency
182+
// through subsequent up-to-date events.
180183
if found.Version != resource.Version {
181184
log.Warnf("Updating status for stale resource; disregard it: id=%s, foundVersion=%d, wantedVersion=%d",
182185
resource.ID, found.Version, resource.Version)
@@ -224,8 +227,9 @@ func (s *sqlResourceService) UpdateStatus(ctx context.Context, resource *api.Res
224227
return found, false, nil
225228
}
226229

230+
// Only update resource status
227231
found.Status = resource.Status
228-
updated, err := s.resourceDao.Update(ctx, found)
232+
updated, err := s.resourceDao.UpdateStatus(ctx, found)
229233
if err != nil {
230234
return nil, false, handleUpdateError("Resource", err)
231235
}

test/integration/resource_test.go

Lines changed: 112 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -490,7 +490,7 @@ func TestMarkAsDeletingThenUpdate(t *testing.T) {
490490
ID: resource.ID,
491491
},
492492
Version: resource.Version,
493-
Status: createStatusWithSequenceID(t, h, resource.ID, fmt.Sprintf("%d", 1)),
493+
Status: createStatusWithSequenceID(t, resource.ID, fmt.Sprintf("%d", 1)),
494494
}
495495
_, updated, svcErr := resourceService.UpdateStatus(ctx, statusRes)
496496
Expect(svcErr).NotTo(HaveOccurred())
@@ -529,8 +529,118 @@ func updateWorkStatus(ctx context.Context, workClient workv1client.ManifestWorkI
529529
return nil
530530
}
531531

532+
// TestUpdateAndUpdateStatusIsolation ensures that Update and UpdateStatus operations
533+
// work independently without affecting each other. Specifically:
534+
// 1. Create a resource
535+
// 2. Update the status
536+
// 3. Update the resource payload
537+
// 4. Verify the status from step 2 is preserved and not affected by the payload update
538+
func TestUpdateAndUpdateStatusIsolation(t *testing.T) {
539+
h, _ := test.RegisterIntegration(t)
540+
541+
account := h.NewRandAccount()
542+
ctx := h.NewAuthenticatedContext(account)
543+
544+
// Step 1: Create a resource
545+
consumer, err := h.CreateConsumer("cluster-" + rand.String(5))
546+
Expect(err).NotTo(HaveOccurred())
547+
deployName := fmt.Sprintf("nginx-%s", rand.String(5))
548+
resource, err := h.CreateResource(uuid.NewString(), consumer.Name, deployName, "default", 1)
549+
Expect(err).NotTo(HaveOccurred())
550+
Expect(resource.Version).To(Equal(int32(1)))
551+
Expect(len(resource.Status)).To(Equal(0), "Initial resource should have no status")
552+
553+
resourceService := h.Env().Services.Resources()
554+
555+
// Step 2: Update the status with sequence ID "1"
556+
statusRes1 := &api.Resource{
557+
Meta: api.Meta{
558+
ID: resource.ID,
559+
},
560+
Version: resource.Version,
561+
Status: createStatusWithSequenceID(t, resource.ID, "1"),
562+
}
563+
_, updated, svcErr := resourceService.UpdateStatus(ctx, statusRes1)
564+
Expect(svcErr).NotTo(HaveOccurred())
565+
Expect(updated).Should(BeTrue(), "Status should be updated")
566+
567+
// Verify status was set correctly
568+
updatedRes, svcErr := resourceService.Get(ctx, resource.ID)
569+
Expect(svcErr).NotTo(HaveOccurred())
570+
Expect(len(updatedRes.Status)).ShouldNot(Equal(0), "Status should not be empty after update")
571+
statusEvt, err := api.JSONMAPToCloudEvent(updatedRes.Status)
572+
Expect(err).NotTo(HaveOccurred())
573+
statusPayload := &workpayload.ManifestBundleStatus{}
574+
Expect(statusEvt.DataAs(statusPayload)).NotTo(HaveOccurred())
575+
Expect(statusPayload.Conditions).To(HaveLen(1))
576+
Expect(statusPayload.Conditions[0].Type).To(Equal("Applied"))
577+
Expect(statusPayload.Conditions[0].Status).To(Equal(metav1.ConditionStatus("True")))
578+
initialStatusMessage := statusPayload.Conditions[0].Message
579+
580+
// Version should still be 1 (UpdateStatus doesn't increment version)
581+
Expect(updatedRes.Version).To(Equal(int32(1)))
582+
583+
// Step 3: Update the resource payload (change replicas from 1 to 3)
584+
newResource, err := h.NewResource(resource.ID, consumer.Name, deployName, "default", 3, updatedRes.Version)
585+
Expect(err).NotTo(HaveOccurred())
586+
newResource.ID = updatedRes.ID
587+
newResource.Status = createStatusWithSequenceID(t, resource.ID, "2")
588+
589+
updatedPayloadRes, svcErr := resourceService.Update(ctx, newResource)
590+
Expect(svcErr).NotTo(HaveOccurred())
591+
Expect(updatedPayloadRes.Version).To(Equal(int32(2)), "Version should be incremented after Update")
592+
593+
// Step 4: Verify that the status is preserved and unchanged
594+
finalRes, svcErr := resourceService.Get(ctx, resource.ID)
595+
Expect(svcErr).NotTo(HaveOccurred())
596+
597+
// Status should still be present and unchanged
598+
Expect(len(finalRes.Status)).ShouldNot(Equal(0), "Status should be preserved after Update")
599+
600+
finalStatusEvt, err := api.JSONMAPToCloudEvent(finalRes.Status)
601+
Expect(err).NotTo(HaveOccurred())
602+
finalStatusPayload := &workpayload.ManifestBundleStatus{}
603+
Expect(finalStatusEvt.DataAs(finalStatusPayload)).NotTo(HaveOccurred())
604+
Expect(finalStatusPayload.Conditions).To(HaveLen(1))
605+
Expect(finalStatusPayload.Conditions[0].Type).To(Equal("Applied"))
606+
Expect(finalStatusPayload.Conditions[0].Status).To(Equal(metav1.ConditionStatus("True")))
607+
Expect(finalStatusPayload.Conditions[0].Message).To(Equal(initialStatusMessage), "Status message should be unchanged")
608+
609+
// Verify the payload was actually updated (replicas changed from 1 to 3)
610+
payloadEvt, err := api.JSONMAPToCloudEvent(finalRes.Payload)
611+
Expect(err).NotTo(HaveOccurred())
612+
payloadData := &workpayload.ManifestBundle{}
613+
Expect(payloadEvt.DataAs(payloadData)).NotTo(HaveOccurred())
614+
Expect(payloadData.Manifests).To(HaveLen(1))
615+
616+
var manifest map[string]interface{}
617+
Expect(json.Unmarshal(payloadData.Manifests[0].Raw, &manifest)).NotTo(HaveOccurred())
618+
spec := manifest["spec"].(map[string]interface{})
619+
Expect(spec["replicas"]).To(Equal(float64(3)), "Replicas should be updated to 3")
620+
621+
// Additional test: Update status again to verify it still works after a payload update
622+
statusRes2 := &api.Resource{
623+
Meta: api.Meta{
624+
ID: resource.ID,
625+
},
626+
Version: finalRes.Version, // Now version 2
627+
Status: createStatusWithSequenceID(t, resource.ID, "2"),
628+
}
629+
updatedRes2, updated2, svcErr := resourceService.UpdateStatus(ctx, statusRes2)
630+
Expect(svcErr).NotTo(HaveOccurred())
631+
Expect(updated2).Should(BeTrue(), "Status should be updated again")
632+
Expect(updatedRes2.Version).To(Equal(int32(2)), "Version should remain 2 after UpdateStatus")
633+
634+
// Verify the new status
635+
newStatusEvt, err := api.JSONMAPToCloudEvent(updatedRes2.Status)
636+
Expect(err).NotTo(HaveOccurred())
637+
newStatusPayload := &workpayload.ManifestBundleStatus{}
638+
Expect(newStatusEvt.DataAs(newStatusPayload)).NotTo(HaveOccurred())
639+
Expect(newStatusPayload.Conditions[0].Message).To(ContainSubstring("sequence 2"), "Status should reflect new update")
640+
}
641+
532642
// createStatusWithSequenceID creates a resource status CloudEvent with the given sequence ID
533-
func createStatusWithSequenceID(t *testing.T, h *test.Helper, resourceID, sequenceID string) map[string]interface{} {
643+
func createStatusWithSequenceID(t *testing.T, resourceID, sequenceID string) map[string]interface{} {
534644
source := "test-agent"
535645
eventType := cetypes.CloudEventsType{
536646
CloudEventsDataType: workpayload.ManifestBundleEventDataType,

0 commit comments

Comments
 (0)