pubsub

package
v0.0.0-...-5edbee2 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 6, 2026 License: Apache-2.0 Imports: 13 Imported by: 0

README

Adding a PubSub Provider

This guide outlines the steps required to integrate a new PubSub provider into the router.

Modify the Router Proto

Update the router.proto file by adding your provider’s configuration. Follow these steps:

  • Define a new configuration message similar to KafkaEventConfiguration.
  • Add this configuration as a repeated field within the DataSourceCustomEvents message.
  • Field naming should reflect the provider's message grouping mechanism. For example, use channels if the provider groups messages by channel, or topics if it uses topics.

After making these changes, compile the updated proto definitions by running the following command from the root directory:

make generate-go

This will generate the new proto files in the gen folder.

Implement the PubSub Provider

To implement a new PubSub provider, the following components must be created:

  • SubscriptionEventConfiguration and PublishEventConfiguration: Define the data structures used for communication between the adapter and the engine.
  • ProviderAdapter: Implements the logic that interfaces with the provider’s client or SDK.
  • SubscriptionDataSource and PublishDataSource: Engine components that leverage the configurations to subscribe and publish data.
  • EngineDataSourceFactory: Bridges the engine and the provider.
  • ProviderBuilder: Used by the router to instantiate the provider.
SubscriptionEventConfiguration and PublishEventConfiguration

These structures should be placed at the top of the engine_datasource.go file. Their design is specific to each provider.

Refer to the kafka implementation for a working example.

ProviderAdapter

This component encapsulates the provider-specific logic. Although not required, it’s best practice to implement the following interface to facilitate testing via mocks:

type Adapter interface {
	Subscribe(ctx context.Context, event SubscriptionEventConfiguration, updater resolve.SubscriptionUpdater) error
	Publish(ctx context.Context, event PublishEventConfiguration) error
	Startup(ctx context.Context) error
	Shutdown(ctx context.Context) error
}

Refer to the kafka implementation for a working example.

SubscriptionDataSource and PublishDataSource

These are the core engine interfaces:

The engine expect two kind of structures:

  • SubscriptionDataSource: Implements resolve.SubscriptionDataSource
  • PublishDataSource: Implements resolve.DataSource

The implementation of SubscriptionDataSource and PublishDataSource should be in the engine_datasource.go file.

They are going to use the SubscriptionEventConfiguration and PublishEventConfiguration that you have implemented in the first step.

Implement these in the engine_datasource.go file, referencing the kafka implementation for a working example.

EngineDataSourceFactory

This structure connects the engine (resolve.DataSource and resolve.SubscriptionDataSource) with the provider implementation. It must implement the EngineDataSourceFactory interface defined in datasource.go.

Refer to the kafka implementation for a working example.

ProviderBuilder

The builder is responsible for instantiating the provider within the router. It must implement the ProviderBuilder interface.

The interface has two generic types:

  • P, the generic type of the options that the provider builder will need, as defined in the config.go (NatsEventSource, KafkaEventSource, ...)
  • E, the generic type of the event configuration that the provider builder will receive, as defined in the proto/wg/cosmo/node/v1/node.proto (KafkaEventConfiguration, NatsEventConfiguration, ...)

Key methods:

  • BuildProvider: Initializes the provider with its configuration and receive the provider options (defined by the P type)
  • BuildEngineDataSourceFactory: Creates the data source and receive the event configuration (defined by the E type)

Refer to the kafka implementation for a working example.

Add tests

You should also add tests to your provider.

Generate mocks

As a first step, you can use the mockery tool to generate the mocks for the ProviderAdapter interface you have implemented. To do this, add the following to the .mockery.yml file:

packages:
  github.com/wundergraph/cosmo/router/pkg/pubsub/{your-provider-name}:
    interfaces:
      Adapter:

Then run the following command from the router directory:

make generate-mocks

This will generate the mocks in the {your-provider-name}/mocks.go file.

You can then use the mocks in your tests.

Tests

You should add tests as specified in the table below.

Implementation File Test File Reference File
engine_datasource.go engine_datasource_test.go kafka implementation
engine_datasource_factory.go engine_datasource_factory_test.go kafka implementation
provider_builder.go provider_builder_test.go kafka implementation
pubsub.go pubsub_test.go TestBuildProvidersAndDataSources_Kafka_OK

Add the provider to the router

Update the BuildProvidersAndDataSources function in the pubsub.go file to include your new provider.

How to use the new PubSub Provider

After you have implemented all the above, you can use your PubSub Provider by adding the following to your router config:

pubsub:
  providers:
    - name: provider-name
      type: new-provider

But to use it in the GraphQL schema, you will have to work in the composition package.

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func BuildProvidersAndDataSources

func BuildProvidersAndDataSources(
	ctx context.Context,
	config config.EventsConfiguration,
	store metric.StreamMetricStore,
	logger *zap.Logger,
	dsConfs []DataSourceConfigurationWithMetadata,
	hostName string,
	routerListenAddr string,
	hooks pubsub_datasource.Hooks,
) ([]pubsub_datasource.Provider, []plan.DataSource, error)

BuildProvidersAndDataSources is a generic function that builds providers and data sources for the given EventsConfiguration and DataSourceConfigurationWithMetadata

Types

type DataSourceConfigurationWithMetadata

type DataSourceConfigurationWithMetadata struct {
	Configuration *nodev1.DataSourceConfiguration
	Metadata      *plan.DataSourceMetadata
}

type EngineEventConfiguration

type EngineEventConfiguration interface {
	GetTypeName() string
	GetFieldName() string
	GetProviderId() string
}

type GetEngineEventConfiguration

type GetEngineEventConfiguration interface {
	GetEngineEventConfiguration() *nodev1.EngineEventConfiguration
}

type GetID

type GetID interface {
	GetID() string
}

type ProviderNotDefinedError

type ProviderNotDefinedError struct {
	ProviderID     string
	ProviderTypeID string
}

func (*ProviderNotDefinedError) Error

func (e *ProviderNotDefinedError) Error() string

Directories

Path Synopsis

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL