@@ -41,7 +41,7 @@ func NewRedisQueue(opts Options, logger *zap.SugaredLogger) (queue.Queue, error)
4141 return q , nil
4242}
4343
44- func (q * RedisQueue ) WriteMessage (ctx context.Context , message * queue.Message ) error {
44+ func (q * RedisQueue ) Enqueue (ctx context.Context , message * queue.Message ) error {
4545 ctx , span := tracing .Start (ctx , "redis.queue.enqueue" , trace .WithSpanKind (trace .SpanKindServer ))
4646 defer span .End ()
4747
@@ -117,15 +117,15 @@ func (q *RedisQueue) delete(ctx context.Context, xmessages []redis.XMessage) err
117117 return err
118118}
119119
120- func (q * RedisQueue ) StartListen (ctx context.Context , handle queue.HandleFunc ) {
120+ func (q * RedisQueue ) StartListen (ctx context.Context , handler queue.HandlerFunc ) {
121121 q .log .Infof ("starting %d listeners" , q .opts .Listeners )
122122 for i := 0 ; i < q .opts .Listeners ; i ++ {
123- go q .listen (ctx , handle )
123+ go q .listen (ctx , handler )
124124 }
125125 go q .process (ctx )
126126}
127127
128- func (q * RedisQueue ) listen (ctx context.Context , handle queue.HandleFunc ) {
128+ func (q * RedisQueue ) listen (ctx context.Context , handler queue.HandlerFunc ) {
129129 for {
130130 select {
131131 case <- ctx .Done ():
@@ -146,7 +146,7 @@ func (q *RedisQueue) listen(ctx context.Context, handle queue.HandleFunc) {
146146 messages = append (messages , toMessage (msg .Values ))
147147 }
148148
149- err = handle (ctx , messages )
149+ err = handler (ctx , messages )
150150 if err != nil {
151151 q .log .Warnf ("failed to handle message: %v" , err )
152152 continue
0 commit comments