2018-10-02 00:25:41 +01:00
|
|
|
// Copyright (C) 2018 Storj Labs, Inc.
|
|
|
|
// See LICENSE for copying information.
|
|
|
|
|
|
|
|
package queue
|
|
|
|
|
|
|
|
import (
|
|
|
|
"github.com/golang/protobuf/proto"
|
|
|
|
"storj.io/storj/pkg/pb"
|
|
|
|
"storj.io/storj/storage"
|
|
|
|
)
|
|
|
|
|
|
|
|
// RepairQueue is the interface for the data repair queue
|
|
|
|
type RepairQueue interface {
|
|
|
|
Enqueue(qi *pb.InjuredSegment) error
|
|
|
|
Dequeue() (pb.InjuredSegment, error)
|
|
|
|
}
|
|
|
|
|
|
|
|
// Queue implements the RepairQueue interface
|
|
|
|
type Queue struct {
|
2018-11-08 13:53:27 +00:00
|
|
|
db storage.Queue
|
2018-10-02 00:25:41 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
// NewQueue returns a pointer to a new Queue instance with an initialized connection to Redis
|
2018-11-08 13:53:27 +00:00
|
|
|
func NewQueue(client storage.Queue) *Queue {
|
|
|
|
return &Queue{db: client}
|
2018-10-02 00:25:41 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
// Enqueue adds a repair segment to the queue
|
|
|
|
func (q *Queue) Enqueue(qi *pb.InjuredSegment) error {
|
|
|
|
val, err := proto.Marshal(qi)
|
|
|
|
if err != nil {
|
2018-10-09 17:09:33 +01:00
|
|
|
return Error.New("error marshalling injured seg %s", err)
|
2018-10-02 00:25:41 +01:00
|
|
|
}
|
2018-11-08 13:53:27 +00:00
|
|
|
|
|
|
|
err = q.db.Enqueue(val)
|
2018-10-02 00:25:41 +01:00
|
|
|
if err != nil {
|
2018-10-09 17:09:33 +01:00
|
|
|
return Error.New("error adding injured seg to queue %s", err)
|
2018-10-02 00:25:41 +01:00
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
// Dequeue returns the next repair segement and removes it from the queue
|
|
|
|
func (q *Queue) Dequeue() (pb.InjuredSegment, error) {
|
2018-11-08 13:53:27 +00:00
|
|
|
val, err := q.db.Dequeue()
|
2018-10-02 00:25:41 +01:00
|
|
|
if err != nil {
|
2018-11-08 13:53:27 +00:00
|
|
|
return pb.InjuredSegment{}, Error.New("error obtaining item from repair queue %s", err)
|
2018-10-02 00:25:41 +01:00
|
|
|
}
|
|
|
|
seg := &pb.InjuredSegment{}
|
|
|
|
err = proto.Unmarshal(val, seg)
|
|
|
|
if err != nil {
|
2018-10-09 17:09:33 +01:00
|
|
|
return pb.InjuredSegment{}, Error.New("error unmarshalling segment %s", err)
|
2018-10-02 00:25:41 +01:00
|
|
|
}
|
|
|
|
return *seg, nil
|
|
|
|
}
|