Skip to content

Commit 1ae432a

Browse files
committed
Arbitrary context call for riverlog middleware
This one's inspired by #914. Our current `riverlog` implementation should work pretty well for most users, but it does force them to use slog. I'm inclined to provide additional options here because in all honesty, slog's API is terrible, and I wouldn't want to use it myself. Here, we add the new initializer `NewMiddlewareArbitrary`. Users will have to bring their own context key like: type ArbitraryContextKey struct{} New middleware takes a `newContext` function that should initialize whatever it wants with an input writer and put it in context: riverlog.NewMiddlewareArbitrary(func(ctx context.Context, w io.Writer) context.Context { logger := log.New(w, "", 0) return context.WithValue(ctx, ArbitraryContextKey{}, logger) }, nil), Work functions then re-extract whatever they added to context (in reality, they'd probably define their own shorthand helpers over what's shown below): func (w *ArbitraryLoggingWorker) Work(ctx context.Context, job *river.Job[ArbitraryLogginArgs]) error { logger := ctx.Value(ArbitraryContextKey{}).(*log.Logger) logger.Printf("Raw log from worker") return nil } An alternative to this might be to expose more River innards that'd let callers set their own metadata, but I'm kind of thinking that a project like that should involve very thorough consideration, and may not even be desirable. I think that even if we go that direction, the new `NewMiddlewareArbitrary` would continue to be useful as a higher level abstraction. Fixes #914.
1 parent 487814e commit 1ae432a

5 files changed

Lines changed: 212 additions & 27 deletions

File tree

CHANGELOG.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
1111

1212
- Preliminary River driver for SQLite (`riverdriver/riversqlite`). This driver seems to produce good results as judged by the test suite, but so far has minimal real world vetting. Try it and let us know how it works out. [PR #870](https://github.com/riverqueue/river/pull/870).
1313
- CLI `river migrate-get` now takes a `--schema` option to inject a custom schema into dumped migrations and schema comments are hidden if `--schema` option isn't provided. [PR #903](https://github.com/riverqueue/river/pull/903).
14+
- Added `riverlog.NewMiddlewareArbitrary` that makes the use of `riverlog` job-persisted logging possible with non-slog loggers. [PR #919](https://github.com/riverqueue/river/pull/919).
1415

1516
### Changed
1617

Lines changed: 108 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,108 @@
1+
package riverlog_test
2+
3+
import (
4+
"context"
5+
"encoding/json"
6+
"fmt"
7+
"io"
8+
"log"
9+
"log/slog"
10+
11+
"github.com/jackc/pgx/v5/pgxpool"
12+
13+
"github.com/riverqueue/river"
14+
"github.com/riverqueue/river/riverdbtest"
15+
"github.com/riverqueue/river/riverdriver/riverpgxv5"
16+
"github.com/riverqueue/river/riverlog"
17+
"github.com/riverqueue/river/rivershared/riversharedtest"
18+
"github.com/riverqueue/river/rivershared/util/slogutil"
19+
"github.com/riverqueue/river/rivershared/util/testutil"
20+
"github.com/riverqueue/river/rivertype"
21+
)
22+
23+
// Callers should define their own context key to extract their a logger back
24+
// out of work context.
25+
type ArbitraryContextKey struct{}
26+
27+
type ArbitraryLogginArgs struct{}
28+
29+
func (ArbitraryLogginArgs) Kind() string { return "logging" }
30+
31+
type ArbitraryLoggingWorker struct {
32+
river.WorkerDefaults[ArbitraryLogginArgs]
33+
}
34+
35+
func (w *ArbitraryLoggingWorker) Work(ctx context.Context, job *river.Job[ArbitraryLogginArgs]) error {
36+
logger := ctx.Value(ArbitraryContextKey{}).(*log.Logger) //nolint:forcetypeassert
37+
logger.Printf("Raw log from worker")
38+
return nil
39+
}
40+
41+
// ExampleNewMiddlewareArbitrary demonstrates the use of riverlog middleware
42+
// with an arbitrary new context function that can be used to inject any sort of
43+
// logger into context.
44+
func ExampleNewMiddlewareArbitrary() {
45+
ctx := context.Background()
46+
47+
dbPool, err := pgxpool.New(ctx, riversharedtest.TestDatabaseURL())
48+
if err != nil {
49+
panic(err)
50+
}
51+
defer dbPool.Close()
52+
53+
workers := river.NewWorkers()
54+
river.AddWorker(workers, &ArbitraryLoggingWorker{})
55+
56+
riverClient, err := river.NewClient(riverpgxv5.New(dbPool), &river.Config{
57+
Logger: slog.New(&slogutil.SlogMessageOnlyHandler{Level: slog.LevelWarn}),
58+
Queues: map[string]river.QueueConfig{
59+
river.QueueDefault: {MaxWorkers: 100},
60+
},
61+
Middleware: []rivertype.Middleware{
62+
riverlog.NewMiddlewareArbitrary(func(ctx context.Context, w io.Writer) context.Context {
63+
// For demonstration purposes we show the use of a built-in
64+
// non-slog logger, but this could be anything like Logrus or
65+
// Zap. Even the raw writer could be stored if so desired.
66+
logger := log.New(w, "", 0)
67+
return context.WithValue(ctx, ArbitraryContextKey{}, logger)
68+
}, nil),
69+
},
70+
Schema: riverdbtest.TestSchema(ctx, testutil.PanicTB(), riverpgxv5.New(dbPool), nil), // only necessary for the example test
71+
TestOnly: true, // suitable only for use in tests; remove for live environments
72+
Workers: workers,
73+
})
74+
if err != nil {
75+
panic(err)
76+
}
77+
78+
// Out of example scope, but used to wait until a job is worked.
79+
subscribeChan, subscribeCancel := riverClient.Subscribe(river.EventKindJobCompleted)
80+
defer subscribeCancel()
81+
82+
if err := riverClient.Start(ctx); err != nil {
83+
panic(err)
84+
}
85+
86+
_, err = riverClient.Insert(ctx, ArbitraryLogginArgs{}, nil)
87+
if err != nil {
88+
panic(err)
89+
}
90+
91+
// Wait for job to complete, extract log data out of metadata, and print it.
92+
for _, event := range riversharedtest.WaitOrTimeoutN(testutil.PanicTB(), subscribeChan, 1) {
93+
var metadataWithLog metadataWithLog
94+
if err := json.Unmarshal(event.Job.Metadata, &metadataWithLog); err != nil {
95+
panic(err)
96+
}
97+
for _, logAttempt := range metadataWithLog.RiverLog {
98+
fmt.Print(logAttempt.Log)
99+
}
100+
}
101+
102+
if err := riverClient.Stop(ctx); err != nil {
103+
panic(err)
104+
}
105+
106+
// Output:
107+
// Raw log from worker
108+
}
Lines changed: 8 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -33,9 +33,9 @@ func (w *LoggingWorker) Work(ctx context.Context, job *river.Job[LoggingArgs]) e
3333
return nil
3434
}
3535

36-
// Example_middleware demonstrates the use of riverlog middleware to inject a
36+
// ExampleNewMiddleware demonstrates the use of riverlog middleware to inject a
3737
// logger into context that'll persist its output onto the job record.
38-
func Example_middleware() {
38+
func ExampleNewMiddleware() {
3939
ctx := context.Background()
4040

4141
dbPool, err := pgxpool.New(ctx, riversharedtest.TestDatabaseURL())
@@ -83,12 +83,6 @@ func Example_middleware() {
8383
panic(err)
8484
}
8585

86-
type metadataWithLog struct {
87-
RiverLog []struct {
88-
Log string `json:"log"`
89-
} `json:"river:log"`
90-
}
91-
9286
// Wait for job to complete, extract log data out of metadata, and print it.
9387
for _, event := range riversharedtest.WaitOrTimeoutN(testutil.PanicTB(), subscribeChan, 1) {
9488
var metadataWithLog metadataWithLog
@@ -108,3 +102,9 @@ func Example_middleware() {
108102
// Logged from worker
109103
// Another line logged from worker
110104
}
105+
106+
type metadataWithLog struct {
107+
RiverLog []struct {
108+
Log string `json:"log"`
109+
} `json:"river:log"`
110+
}

riverlog/river_log.go

Lines changed: 51 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -42,8 +42,9 @@ func Logger(ctx context.Context) *slog.Logger {
4242
type Middleware struct {
4343
baseservice.BaseService
4444
rivertype.Middleware
45-
config *MiddlewareConfig
46-
newHandler func(w io.Writer) slog.Handler
45+
config *MiddlewareConfig
46+
newContextLogger func(ctx context.Context, w io.Writer) context.Context
47+
newSlogHandler func(w io.Writer) slog.Handler
4748
}
4849

4950
// MiddlewareConfig is configuration for Middleware.
@@ -64,8 +65,8 @@ type MiddlewareConfig struct {
6465
MaxSizeBytes int
6566
}
6667

67-
// NewMiddleware initializes a new Middleware with the given handler function
68-
// and configuration.
68+
// NewMiddleware initializes a new Middleware with the given slog handler
69+
// initialization function and configuration.
6970
//
7071
// newHandler is a function which is invoked on every Work execution to generate
7172
// a new slog.Handler for a work-specific slog.Logger. It should take an
@@ -77,20 +78,46 @@ type MiddlewareConfig struct {
7778
// riverlog.NewMiddleware(func(w io.Writer) slog.Handler {
7879
// return slog.NewJSONHandler(w, nil)
7980
// }, nil)
80-
func NewMiddleware(newHandler func(w io.Writer) slog.Handler, config *MiddlewareConfig) *Middleware {
81-
if config == nil {
82-
config = &MiddlewareConfig{}
81+
func NewMiddleware(newSlogHandler func(w io.Writer) slog.Handler, config *MiddlewareConfig) *Middleware {
82+
return &Middleware{
83+
config: defaultConfig(config),
84+
newSlogHandler: newSlogHandler,
8385
}
86+
}
8487

85-
// Assign defaults.
86-
config = &MiddlewareConfig{
87-
MaxSizeBytes: cmp.Or(config.MaxSizeBytes, maxSizeBytes),
88+
// NewMiddlewareArbitrary initializes a new Middleware with the given arbitrary
89+
// context initialization function and configuration.
90+
//
91+
// newContext is a function which is invoked on every Work execution to generate
92+
// a new context for the worker. It's generally used to initialize a logger with
93+
// the given writer and put it in context under a user-defined context key for
94+
// later use.
95+
//
96+
// This variant is meant to provide callers with a version of the middleware
97+
// that's not tied to slog. A non-slog standard library logger, Logrus, or Zap
98+
// logger could all be placed in context according to preferred convention.
99+
//
100+
// For example:
101+
//
102+
// riverlog.NewMiddlewareArbitrary(func(ctx context.Context, w io.Writer) context.Context {
103+
// logger := log.New(w, "", 0)
104+
// return context.WithValue(ctx, ArbitraryContextKey{}, logger)
105+
// }, nil),
106+
func NewMiddlewareArbitrary(newContext func(ctx context.Context, w io.Writer) context.Context, config *MiddlewareConfig) *Middleware {
107+
return &Middleware{
108+
config: defaultConfig(config),
109+
newContextLogger: newContext,
88110
}
111+
}
89112

90-
return &Middleware{
91-
config: config,
92-
newHandler: newHandler,
113+
func defaultConfig(config *MiddlewareConfig) *MiddlewareConfig {
114+
if config == nil {
115+
config = &MiddlewareConfig{}
93116
}
117+
118+
config.MaxSizeBytes = cmp.Or(config.MaxSizeBytes, maxSizeBytes)
119+
120+
return config
94121
}
95122

96123
type logAttempt struct {
@@ -106,9 +133,18 @@ func (m *Middleware) Work(ctx context.Context, job *rivertype.JobRow, doInner fu
106133
var (
107134
existingLogData metadataWithLog
108135
logBuf bytes.Buffer
109-
logger = slog.New(m.newHandler(&logBuf))
110136
)
111137

138+
switch {
139+
case m.newContextLogger != nil:
140+
ctx = m.newContextLogger(ctx, &logBuf)
141+
case m.newSlogHandler != nil:
142+
logger := slog.New(m.newSlogHandler(&logBuf))
143+
ctx = context.WithValue(ctx, contextKey{}, logger)
144+
default:
145+
return errors.New("expected either newContextLogger or newSlogHandler to be set")
146+
}
147+
112148
if err := json.Unmarshal(job.Metadata, &existingLogData); err != nil {
113149
return err
114150
}
@@ -145,5 +181,5 @@ func (m *Middleware) Work(ctx context.Context, job *rivertype.JobRow, doInner fu
145181
metadataUpdates[metadataKey] = json.RawMessage(allLogDataBytes)
146182
}()
147183

148-
return doInner(context.WithValue(ctx, contextKey{}, logger))
184+
return doInner(ctx)
149185
}

riverlog/river_log_test.go

Lines changed: 44 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import (
55
"encoding/json"
66
"errors"
77
"io"
8+
"log"
89
"log/slog"
910
"testing"
1011

@@ -54,8 +55,10 @@ func TestMiddleware(t *testing.T) {
5455
ctx := context.Background()
5556

5657
type testBundle struct {
57-
driver *riverpgxv5.Driver
58-
tx pgx.Tx
58+
clientConfig *river.Config
59+
driver *riverpgxv5.Driver
60+
middleware *Middleware
61+
tx pgx.Tx
5962
}
6063

6164
setup := func(t *testing.T, config *MiddlewareConfig) (*rivertest.Worker[loggingArgs, pgx.Tx], *testBundle) {
@@ -74,8 +77,10 @@ func TestMiddleware(t *testing.T) {
7477
)
7578

7679
return rivertest.NewWorker(t, driver, clientConfig, worker), &testBundle{
77-
driver: driver,
78-
tx: tx,
80+
clientConfig: clientConfig,
81+
driver: driver,
82+
middleware: middleware,
83+
tx: tx,
7984
}
8085
}
8186

@@ -205,4 +210,39 @@ func TestMiddleware(t *testing.T) {
205210
metadataWithLog.RiverLog,
206211
)
207212
})
213+
214+
t.Run("RawMiddleware", func(t *testing.T) {
215+
t.Parallel()
216+
217+
_, bundle := setup(t, nil)
218+
219+
type writerContextKey struct{}
220+
221+
bundle.middleware.newContextLogger = func(ctx context.Context, w io.Writer) context.Context {
222+
logger := log.New(w, "", 0)
223+
return context.WithValue(ctx, writerContextKey{}, logger)
224+
}
225+
bundle.middleware.newSlogHandler = nil
226+
227+
testWorker := rivertest.NewWorker(t, bundle.driver, bundle.clientConfig, river.WorkFunc(func(ctx context.Context, job *river.Job[loggingArgs]) error {
228+
logger := ctx.Value(writerContextKey{}).(*log.Logger) //nolint:forcetypeassert
229+
logger.Printf(job.Args.Message)
230+
return nil
231+
}))
232+
233+
workRes, err := testWorker.Work(ctx, t, bundle.tx, loggingArgs{Message: "Raw log from worker"}, nil)
234+
require.NoError(t, err)
235+
236+
var metadataWithLog metadataWithLog
237+
require.NoError(t, json.Unmarshal(workRes.Job.Metadata, &metadataWithLog))
238+
239+
require.Equal(t, []logAttempt{
240+
{
241+
Attempt: 1,
242+
Log: "Raw log from worker\n",
243+
},
244+
},
245+
metadataWithLog.RiverLog,
246+
)
247+
})
208248
}

0 commit comments

Comments
 (0)