From 52ca801fbd1c4147b3a0b7dd03ebcb82e920b0b4 Mon Sep 17 00:00:00 2001 From: Joel Wetzell Date: Wed, 27 May 2026 18:30:23 -0500 Subject: [PATCH] add context to keyvaluemodule get/set --- internal/common/module.go | 4 ++-- internal/module/redis-client.go | 8 ++++---- internal/processor/kv-get.go | 2 +- internal/processor/kv-set.go | 2 +- internal/test/module.go | 4 ++-- 5 files changed, 10 insertions(+), 10 deletions(-) diff --git a/internal/common/module.go b/internal/common/module.go index e9616ef..f25c526 100644 --- a/internal/common/module.go +++ b/internal/common/module.go @@ -17,8 +17,8 @@ type OutputModule interface { } type KeyValueModule interface { - Get(key string) (any, error) - Set(key string, value any) error + Get(ctx context.Context, key string) (any, error) + Set(ctx context.Context, key string, value any) error } type DatabaseModule interface { diff --git a/internal/module/redis-client.go b/internal/module/redis-client.go index 556b1d0..be980c8 100644 --- a/internal/module/redis-client.go +++ b/internal/module/redis-client.go @@ -110,9 +110,9 @@ func (rc *RedisClient) Stop() { rc.logger.Debug("done") } -func (rc *RedisClient) Get(key string) (any, error) { +func (rc *RedisClient) Get(ctx context.Context, key string) (any, error) { if rc.client != nil { - val, err := rc.client.Get(rc.ctx, key).Result() + val, err := rc.client.Get(ctx, key).Result() if err != nil { return nil, err } @@ -121,9 +121,9 @@ func (rc *RedisClient) Get(key string) (any, error) { return nil, errors.New("redis.client not setup") } -func (rc *RedisClient) Set(key string, value any) error { +func (rc *RedisClient) Set(ctx context.Context, key string, value any) error { if rc.client != nil { - status := rc.client.Set(rc.ctx, key, value, 0) + status := rc.client.Set(ctx, key, value, 0) return status.Err() } return errors.New("redis.client not setup") diff --git a/internal/processor/kv-get.go b/internal/processor/kv-get.go index bc270fd..2102e67 100644 --- a/internal/processor/kv-get.go +++ b/internal/processor/kv-get.go @@ -77,7 +77,7 @@ func (kvg *KVGet) Process(ctx context.Context, wrappedPayload common.WrappedPayl kvg.module = kvModule } - value, err := kvg.module.Get(kvg.Key) + value, err := kvg.module.Get(ctx, kvg.Key) if err != nil { wrappedPayload.End = true return wrappedPayload, fmt.Errorf("kv.get error getting key: %w", err) diff --git a/internal/processor/kv-set.go b/internal/processor/kv-set.go index ce14979..d67b2c5 100644 --- a/internal/processor/kv-set.go +++ b/internal/processor/kv-set.go @@ -78,7 +78,7 @@ func (kvs *KVSet) Process(ctx context.Context, wrappedPayload common.WrappedPayl kvs.module = kvModule } - err := kvs.module.Set(kvs.Key, wrappedPayload.Payload) + err := kvs.module.Set(ctx, kvs.Key, wrappedPayload.Payload) if err != nil { wrappedPayload.End = true return wrappedPayload, fmt.Errorf("kv.set error setting key: %w", err) diff --git a/internal/test/module.go b/internal/test/module.go index 21a25fc..00f3427 100644 --- a/internal/test/module.go +++ b/internal/test/module.go @@ -89,7 +89,7 @@ func (m *TestKVModule) Id() string { return m.id } -func (m *TestKVModule) Get(key string) (any, error) { +func (m *TestKVModule) Get(ctx context.Context, key string) (any, error) { if m.kvData == nil { return nil, nil } @@ -100,7 +100,7 @@ func (m *TestKVModule) Get(key string) (any, error) { return value, nil } -func (m *TestKVModule) Set(key string, value any) error { +func (m *TestKVModule) Set(ctx context.Context, key string, value any) error { if m.kvData == nil { m.kvData = make(map[string]any) }