Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 23 additions & 7 deletions pkg/cloudevents/generic/options/v2/pubsub/transport.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package pubsub
import (
"context"
"fmt"
"sync"

"cloud.google.com/go/pubsub/v2"
cloudevents "github.com/cloudevents/sdk-go/v2"
Expand Down Expand Up @@ -135,36 +136,51 @@ func (o *pubsubTransport) Send(ctx context.Context, evt cloudevents.Event) error
func (o *pubsubTransport) Receive(ctx context.Context, fn options.ReceiveHandlerFn) error {
errChan := make(chan error)

// Use a mutex to ensure sequential processing across both subscribers.
// This prevents race conditions when concurrent events for the same
// resource arrive on different subscriptions.
var mu sync.Mutex

// start the subscriber for spec/status updates
go o.receiveFromSubscriber(ctx, o.subscriber, fn, errChan)
go o.receiveFromSubscriber(ctx, o.subscriber, fn, &mu, errChan)

// start the resync subscriber for resync events
go o.receiveFromSubscriber(ctx, o.resyncSubscriber, fn, errChan)
go o.receiveFromSubscriber(ctx, o.resyncSubscriber, fn, &mu, errChan)

// Return the error from either subscriber (including context cancellation).
// Return the first error from either subscriber (including context cancellation).
// We return errors directly instead of writing to the transport errorChan because
// Pub/Sub client has internal retry logic for transient errors. Only non-retryable
// errors or context cancellation will be returned here.
return <-errChan
}

// receiveFromSubscriber handles receiving messages from a subscriber.
// It uses a mutex to ensure sequential processing across all subscribers.
func (o *pubsubTransport) receiveFromSubscriber(
ctx context.Context,
subscriber *pubsub.Subscriber,
fn options.ReceiveHandlerFn,
mu *sync.Mutex,
errChan chan<- error,
) {
logger := klog.FromContext(ctx)
err := subscriber.Receive(ctx, func(ctx context.Context, msg *pubsub.Message) {
evt, err := Decode(msg)
if err != nil {
// also send ACK on decode error since redelivery won't fix it.
// ACK decode errors immediately since redelivery won't fix them.
logger.Error(err, "failed to decode pubsub message")
} else {
fn(ctx, evt)
msg.Ack()
return
}
// send ACK after all receiver handlers complete.

// Lock to ensure sequential processing across both subscribers.
// This prevents race conditions when concurrent events for the same
// resource arrive on different subscriptions.
mu.Lock()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I am not sure having lock is good here. Can we consider put this into a threadsafe queue, and start another gorouting to read from queue and call handler?

@morvencao morvencao Mar 19, 2026

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

code updated to use a separate goroutine to process messages/events using a queue, but ACK message is still called in Receive function to respect Pub/Sub flow control by only Ack'ing after the handler completes.

defer mu.Unlock()
fn(ctx, evt)

// ACK after successful processing.
msg.Ack()
})

Expand Down
168 changes: 168 additions & 0 deletions test/integration/cloudevents/manifestworkclients_resync_pubsub_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,168 @@
package cloudevents

import (
"context"
"fmt"
"time"

"github.com/onsi/ginkgo"
"github.com/onsi/gomega"

metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/util/rand"

workv1 "open-cluster-management.io/api/work/v1"

"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/common"
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/work"
agentcodec "open-cluster-management.io/sdk-go/pkg/cloudevents/clients/work/agent/codec"
"open-cluster-management.io/sdk-go/pkg/cloudevents/clients/work/payload"
sourcecodec "open-cluster-management.io/sdk-go/pkg/cloudevents/clients/work/source/codec"
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic"
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/clients"
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/options/v2/pubsub"
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
"open-cluster-management.io/sdk-go/test/integration/cloudevents/agent"
"open-cluster-management.io/sdk-go/test/integration/cloudevents/util"
)

// emptyManifestWorkLister is a simple lister that returns empty list
// Used for testing where we only need to publish events, not list resources
type emptyManifestWorkLister struct{}

func (e *emptyManifestWorkLister) List(ctx context.Context, options types.ListOptions) ([]*workv1.ManifestWork, error) {
return []*workv1.ManifestWork{}, nil
}

// This test simulates the PubSub race condition by sending create_request and resync_response events concurrently
// to a running work agent, to ensure pubsub transport can handle messages arriving on the different subscriptions simultaneously.
var _ = ginkgo.Describe("ManifestWork Clients Test - Resync PubSub", func() {
ginkgo.Context("Concurrent message delivery on pubsub", func() {
var ctx context.Context
var cancel context.CancelFunc

var sourceID string
var clusterName string
var workName string

var agentClientHolder *work.ClientHolder
var sourceCloudEventsClient generic.CloudEventsClient[*workv1.ManifestWork]

ginkgo.BeforeEach(func() {
ctx, cancel = context.WithCancel(context.Background())
sourceID = fmt.Sprintf("mw-pubsub-race-%s", rand.String(5))
clusterName = fmt.Sprintf("cluster-race-%s", rand.String(5))
workName = "race-test-work"

// Setup PubSub topics and subscriptions
gomega.Expect(setupTopicsAndSubscriptions(ctx, clusterName, sourceID)).ToNot(gomega.HaveOccurred())

// Start the agent and keep it running
ginkgo.By("starting the agent")
pubsubAgentOptions := util.NewPubSubAgentOptions(pubsubServer.Addr, pubsubProjectID, clusterName, true)
var err error
agentClientHolder, _, err = agent.StartWorkAgent(ctx, clusterName, pubsubAgentOptions, agentcodec.NewManifestBundleCodec())
gomega.Expect(err).ToNot(gomega.HaveOccurred())

// wait for agent ready
<-time.After(time.Second)

// Create a pure CloudEvents client (source) to send events directly
ginkgo.By("creating pure cloudevents client to send events")
pubsubSourceOptions := util.NewPubSubSourceOptions(pubsubServer.Addr, pubsubProjectID, sourceID, true)
sourceOptions := pubsub.NewSourceOptions(pubsubSourceOptions, sourceID)

// Use a simple lister that returns empty list (we're not listing, just publishing)
lister := &emptyManifestWorkLister{}
hashGetter := func(obj *workv1.ManifestWork) (string, error) {
return "", nil // not used for this test
}

sourceCloudEventsClient, err = clients.NewCloudEventSourceClient(
ctx,
sourceOptions,
lister,
hashGetter,
sourcecodec.NewManifestBundleCodec(),
)
gomega.Expect(err).ToNot(gomega.HaveOccurred())

// wait for source client ready
<-time.After(time.Second)
})

ginkgo.AfterEach(func() {
cancel()
})

ginkgo.It("should handle concurrent create_request and resync_response without race", func() {
// Create the manifestwork that we'll send events for
work := util.NewManifestWork(clusterName, workName, true)
work.UID = "test-uid-123"
work.Labels = map[string]string{
common.CloudEventsOriginalSourceLabelKey: sourceID,
}

// Define the event types we'll send concurrently
createRequest := types.CloudEventsType{
CloudEventsDataType: payload.ManifestBundleEventDataType,
SubResource: types.SubResourceSpec,
Action: types.CreateRequestAction,
}

resyncResponse := types.CloudEventsType{
CloudEventsDataType: payload.ManifestBundleEventDataType,
SubResource: types.SubResourceSpec,
Action: types.ResyncResponseAction,
}

// THE RACE CONDITION TEST:
// Send both create_request and resync_response events concurrently.
// Both arrive on the agent's sourceevents and sourcebroadcast subscriptions from the source.
// Without the sequential channel fix: two goroutines in receiveFromSubscriber
// could process these concurrently, causing race conditions in the agent store
// and potential resource version conflicts when controllers patch the work.
// With the fix: events are funneled through a single channel for sequential processing.
ginkgo.By("sending create_request and resync_response concurrently")

start := make(chan struct{})
errChan := make(chan error, 2)

// Send create_request
go func() {
<-start
err := sourceCloudEventsClient.Publish(ctx, createRequest, work)
errChan <- err
}()

// Send resync_response
go func() {
<-start
err := sourceCloudEventsClient.Publish(ctx, resyncResponse, work.DeepCopy())
errChan <- err
}()
Comment thread
coderabbitai[bot] marked this conversation as resolved.

close(start)

// Wait for both publishes to complete
for range 2 {
err := <-errChan
gomega.Expect(err).ToNot(gomega.HaveOccurred())
}

// Verify the work is correctly applied in the agent despite concurrent message delivery
ginkgo.By("verifying manifestwork is applied correctly")
gomega.Eventually(func() error {
retrievedWork, err := agentClientHolder.ManifestWorks(clusterName).Get(
ctx, workName, metav1.GetOptions{})
if err != nil {
return err
}
if len(retrievedWork.Spec.Workload.Manifests) == 0 {
return fmt.Errorf("work has no manifests")
}
return nil
}, 10*time.Second, 500*time.Millisecond).Should(gomega.Succeed())
})
})
})
Loading