Skip to content

Support for client middleware? #159

Description

@CGA1123

I'm integration river into a project and would love to be able to have some kind of a middleware pattern whereby I can inject metadata such as trace IDs, correlation IDs, etc...

I'm looking for something akin to the ClientMiddleware available in Sidekiq, which allows for access to the job object before persistence. The equivalent for ServerMiddleware is more straightforward, as wrapping the Worker interface is easily done!

I think for the time being I'll create a smaller Client interface for river that I propagate through my application which I can then use to wrap the Insert operations with custom logic that have access to the context in a uniform manner.

It might be nice for that to be something that can be configured directly on a *river.Client as a middleware stack?

Would be great to hear thoughts about how others are approaching this problem, and views around this from the maintainers as well! Thank you 😄

Activity

  1. CGA1123 commented on Jan 18, 2024

    @CGA1123
    ContributorAuthor

    I've ended up with the following interface for the time being to enable instrumenting enqueues and propagating tracing information

    package riverutil
    
    import (
    	"context"
    
    	"github.com/jackc/pgx/v5"
    	"github.com/riverqueue/river"
    	"github.com/riverqueue/river/rivertype"
    )
    
    // Ensure *river.Client satisfies this interface at compile time.
    var _ EnqueueClient[pgx.Tx] = (*river.Client[pgx.Tx])(nil)
    
    // EnqueueClient defines the subset of the *river.Client interface which is
    // available when initialised without a `*pgxpool.Pool`.
    type EnqueueClient[TTx any] interface {
    	InsertTx(context.Context, TTx, river.JobArgs, *river.InsertOpts) (*rivertype.JobRow, error)
    	InsertManyTx(context.Context, TTx, []river.InsertManyParams) (int64, error)
    }

    And the current implementation for instrumented enqueues (doesn't create any new spans, but propagates the SpanContext through, so the work can know where the job came from and correlate through):

    package olly
    
    import (
    	"context"
    	"fmt"
    
    	"github.com/jackc/pgx/v5"
    	"github.com/CGA1123/riverplayground/riverutil"
    	"github.com/riverqueue/river"
    	"github.com/riverqueue/river/riverdriver/riverpgxv5"
    	"github.com/riverqueue/river/rivertype"
    	"go.opentelemetry.io/otel"
    	"go.opentelemetry.io/otel/propagation"
    )
    
    type enqueueClient struct {
    	river *river.Client[pgx.Tx]
    }
    
    func (ec *enqueueClient) InsertTx(ctx context.Context, tx pgx.Tx, j river.JobArgs, opts *river.InsertOpts) (*rivertype.JobRow, error) {
    	opts = propagateRiverTrace(ctx, j, opts)
    
    	row, err := ec.river.InsertTx(ctx, tx, j, opts)
    	if err != nil {
    		return row, err
    	}
    
    	return row, err
    }
    
    func (ec *enqueueClient) InsertManyTx(ctx context.Context, tx pgx.Tx, jobs []river.InsertManyParams) (int64, error) {
    	count, err := ec.river.InsertManyTx(ctx, tx, jobs)
    	if err != nil {
    		return count, err
    	}
    
    	for _, j := range jobs {
    		j.InsertOpts = propagateRiverTrace(ctx, j.Args, j.InsertOpts)
    	}
    
    	return count, err
    }
    
    func propagateRiverTrace(ctx context.Context, j river.JobArgs, opts *river.InsertOpts) *river.InsertOpts {
    	if opts == nil {
    		opts = &river.InsertOpts{}
    	}
    
    	var tags []string
    	if argsWithOpts, ok := j.(river.JobArgsWithInsertOpts); ok {
    		tags = argsWithOpts.InsertOpts().Tags
    	}
    
    	c := propagation.MapCarrier(map[string]string{})
    	otel.GetTextMapPropagator().Inject(ctx, c)
    	ollyTags := make([]string, 0, len(c))
    	for k, v := range c {
    		ollyTags = append(ollyTags, fmt.Sprintf("%s%s:%s", riverTagPrefix, k, v))
    	}
    
    	// This replicates the behaviour of `river.insertParamsFromArgsAndOptions`
    	if opts.Tags == nil {
    		opts.Tags = append(tags, ollyTags...)
    	} else {
    		opts.Tags = append(opts.Tags, ollyTags...)
    	}
    
    	return opts
    }
    
    // EnqueueClient returns a wrapped `*river.Client` with a reduced set of
    // methods exposed.
    //
    // It includes additional observability and propagates relevant metadata at
    // enqueue-time.
    //
    // This function will panic if `river.NewClient` returns an error.
    func EnqueueClient(ctx context.Context) riverutil.EnqueueClient[pgx.Tx] {
    	c, err := river.NewClient(riverpgxv5.New(nil), &river.Config{})
    	if err != nil {
    		panic(err)
    	}
    
    	return &enqueueClient{river: c}
    }
  2. locked and limited conversation to collaborators on Jan 20, 2024
  3. converted this issue into a discussion #167 on Jan 20, 2024
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions