Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
83 changes: 73 additions & 10 deletions pkg/cloudevents/generic/options/v2/pubsub/transport.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,13 @@ import (

var _ options.CloudEventTransport = &pubsubTransport{}

// messageWork represents a message to be processed sequentially by the worker.
type messageWork struct {
ctx context.Context
evt cloudevents.Event
done chan struct{} // Signals when processing is complete
}

// pubsubTransport is a CloudEventTransport implementation for Pub/Sub.
type pubsubTransport struct {
PubSubOptions
Expand Down Expand Up @@ -135,37 +142,93 @@ 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 buffered channel to queue incoming messages from both subscriptions.
workChan := make(chan messageWork, 10)

// start a single worker goroutine to process messages sequentially.
// to ensure that events for the same resource are processed in order,
// preventing race conditions when concurrent events arrive on different subscriptions.
go o.processMessages(ctx, fn, workChan)

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

// start the resync subscriber for resync events
go o.receiveFromSubscriber(ctx, o.resyncSubscriber, fn, errChan)
go o.receiveFromSubscriber(ctx, o.resyncSubscriber, workChan, 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
}

// processMessages reads from workChan and processes messages one at a time to ensure sequential event processing.
func (o *pubsubTransport) processMessages(
ctx context.Context,
fn options.ReceiveHandlerFn,
workChan <-chan messageWork,
) {
for {
select {
case <-ctx.Done():
// context canceled - stop processing
return
case work, ok := <-workChan:
if !ok {
// channel closed - stop processing
return
}

// process the event
fn(work.ctx, work.evt)

// signal completion so the Receive callback can Ack the message
close(work.done)
}
}
}

// receiveFromSubscriber handles receiving messages from a subscriber.
func (o *pubsubTransport) receiveFromSubscriber(
ctx context.Context,
subscriber *pubsub.Subscriber,
fn options.ReceiveHandlerFn,
workChan chan<- messageWork,
errChan chan<- error,
) {
logger := klog.FromContext(ctx)
err := subscriber.Receive(ctx, func(ctx context.Context, msg *pubsub.Message) {
err := subscriber.Receive(ctx, func(msgCtx context.Context, msg *pubsub.Message) {
// decode the message first
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
}

// create a work item with a completion signal
work := messageWork{
ctx: msgCtx,
evt: evt,
done: make(chan struct{}),
}

// queue the work for sequential processing
select {
case workChan <- work:
// block until the worker processes the message,
// to ensures we respect Pub/Sub flow control by only Ack'ing
// after the handler completes, while still using a queue-based
// approach for sequential processing instead of a mutex.
<-work.done

// now that processing is complete, Ack the message.
msg.Ack()

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.

do we need to block here. What if we directly put msg into the queue, decode msg in processMessages and ACK after handler?

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.

we can't ACK message in processMessages, because it's in another go routine; Pub/Sub requires Ack() or Nack() to be called within the Receive handler, not from another goroutine, see:

Do not call Ack or Nack from a different goroutine. The Receive method handles concurrency for you, and all acknowledgment should happen synchronously in the callback.
Pub/Sub Go Client: Subscription.Receive

case <-ctx.Done():
// context canceled - nack the message so it can be redelivered
msg.Nack()
}
// send ACK after all receiver handlers complete.
msg.Ack()
})

// The Pub/Sub client's Receive call automatically retries on retryable errors.
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