lotus/extern/sector-storage/worker_tracked.go

139 lines
4.4 KiB
Go
Raw Normal View History

package sectorstorage
2020-07-21 18:01:25 +00:00
import (
"context"
"io"
"sync"
"time"
"github.com/ipfs/go-cid"
2020-09-07 03:49:10 +00:00
"github.com/filecoin-project/go-state-types/abi"
2020-07-21 18:01:25 +00:00
"github.com/filecoin-project/specs-storage/storage"
"github.com/filecoin-project/lotus/extern/sector-storage/sealtasks"
"github.com/filecoin-project/lotus/extern/sector-storage/storiface"
2020-07-21 18:01:25 +00:00
)
2020-09-23 12:56:37 +00:00
type trackedWork struct {
job storiface.WorkerJob
worker WorkerID
}
2020-07-21 18:01:25 +00:00
type workTracker struct {
lk sync.Mutex
2020-09-07 14:12:46 +00:00
done map[storiface.CallID]struct{}
2020-09-23 12:56:37 +00:00
running map[storiface.CallID]trackedWork
2020-07-21 18:01:25 +00:00
// TODO: done, aggregate stats, queue stats, scheduler feedback
}
2020-09-07 14:12:46 +00:00
func (wt *workTracker) onDone(callID storiface.CallID) {
2020-07-21 18:01:25 +00:00
wt.lk.Lock()
defer wt.lk.Unlock()
2020-09-07 14:12:46 +00:00
_, ok := wt.running[callID]
if !ok {
wt.done[callID] = struct{}{}
return
2020-07-21 18:01:25 +00:00
}
2020-09-07 14:12:46 +00:00
delete(wt.running, callID)
}
func (wt *workTracker) track(wid WorkerID, sid storage.SectorRef, task sealtasks.TaskType) func(storiface.CallID, error) (storiface.CallID, error) {
2020-09-07 14:12:46 +00:00
return func(callID storiface.CallID, err error) (storiface.CallID, error) {
if err != nil {
return callID, err
}
2020-07-21 18:01:25 +00:00
wt.lk.Lock()
defer wt.lk.Unlock()
2020-09-07 14:12:46 +00:00
_, done := wt.done[callID]
if done {
delete(wt.done, callID)
return callID, err
}
2020-09-23 12:56:37 +00:00
wt.running[callID] = trackedWork{
job: storiface.WorkerJob{
ID: callID,
Sector: sid.ID,
2020-09-23 12:56:37 +00:00
Task: task,
Start: time.Now(),
},
worker: wid,
2020-09-07 14:12:46 +00:00
}
return callID, err
2020-07-21 18:01:25 +00:00
}
}
2020-09-23 12:56:37 +00:00
func (wt *workTracker) worker(wid WorkerID, w Worker) Worker {
2020-07-21 18:01:25 +00:00
return &trackedWorker{
2020-09-23 12:56:37 +00:00
Worker: w,
wid: wid,
2020-07-21 18:01:25 +00:00
tracker: wt,
}
}
2020-09-23 12:56:37 +00:00
func (wt *workTracker) Running() []trackedWork {
2020-07-21 18:01:25 +00:00
wt.lk.Lock()
defer wt.lk.Unlock()
2020-09-23 12:56:37 +00:00
out := make([]trackedWork, 0, len(wt.running))
2020-07-21 18:01:25 +00:00
for _, job := range wt.running {
out = append(out, job)
}
return out
}
type trackedWorker struct {
Worker
2020-09-23 12:56:37 +00:00
wid WorkerID
2020-07-21 18:01:25 +00:00
tracker *workTracker
}
func (t *trackedWorker) SealPreCommit1(ctx context.Context, sector storage.SectorRef, ticket abi.SealRandomness, pieces []abi.PieceInfo) (storiface.CallID, error) {
2020-09-23 12:56:37 +00:00
return t.tracker.track(t.wid, sector, sealtasks.TTPreCommit1)(t.Worker.SealPreCommit1(ctx, sector, ticket, pieces))
2020-07-21 18:01:25 +00:00
}
func (t *trackedWorker) SealPreCommit2(ctx context.Context, sector storage.SectorRef, pc1o storage.PreCommit1Out) (storiface.CallID, error) {
2020-09-23 12:56:37 +00:00
return t.tracker.track(t.wid, sector, sealtasks.TTPreCommit2)(t.Worker.SealPreCommit2(ctx, sector, pc1o))
2020-07-21 18:01:25 +00:00
}
func (t *trackedWorker) SealCommit1(ctx context.Context, sector storage.SectorRef, ticket abi.SealRandomness, seed abi.InteractiveSealRandomness, pieces []abi.PieceInfo, cids storage.SectorCids) (storiface.CallID, error) {
2020-09-23 12:56:37 +00:00
return t.tracker.track(t.wid, sector, sealtasks.TTCommit1)(t.Worker.SealCommit1(ctx, sector, ticket, seed, pieces, cids))
2020-07-21 18:01:25 +00:00
}
func (t *trackedWorker) SealCommit2(ctx context.Context, sector storage.SectorRef, c1o storage.Commit1Out) (storiface.CallID, error) {
2020-09-23 12:56:37 +00:00
return t.tracker.track(t.wid, sector, sealtasks.TTCommit2)(t.Worker.SealCommit2(ctx, sector, c1o))
2020-07-21 18:01:25 +00:00
}
func (t *trackedWorker) FinalizeSector(ctx context.Context, sector storage.SectorRef, keepUnsealed []storage.Range) (storiface.CallID, error) {
2020-09-23 12:56:37 +00:00
return t.tracker.track(t.wid, sector, sealtasks.TTFinalize)(t.Worker.FinalizeSector(ctx, sector, keepUnsealed))
2020-07-21 18:01:25 +00:00
}
func (t *trackedWorker) AddPiece(ctx context.Context, sector storage.SectorRef, pieceSizes []abi.UnpaddedPieceSize, newPieceSize abi.UnpaddedPieceSize, pieceData storage.Data) (storiface.CallID, error) {
2020-09-23 12:56:37 +00:00
return t.tracker.track(t.wid, sector, sealtasks.TTAddPiece)(t.Worker.AddPiece(ctx, sector, pieceSizes, newPieceSize, pieceData))
2020-07-21 18:01:25 +00:00
}
func (t *trackedWorker) Fetch(ctx context.Context, s storage.SectorRef, ft storiface.SectorFileType, ptype storiface.PathType, am storiface.AcquireMode) (storiface.CallID, error) {
2020-09-23 12:56:37 +00:00
return t.tracker.track(t.wid, s, sealtasks.TTFetch)(t.Worker.Fetch(ctx, s, ft, ptype, am))
2020-07-21 18:01:25 +00:00
}
func (t *trackedWorker) UnsealPiece(ctx context.Context, id storage.SectorRef, index storiface.UnpaddedByteIndex, size abi.UnpaddedPieceSize, randomness abi.SealRandomness, cid cid.Cid) (storiface.CallID, error) {
2020-09-23 12:56:37 +00:00
return t.tracker.track(t.wid, id, sealtasks.TTUnseal)(t.Worker.UnsealPiece(ctx, id, index, size, randomness, cid))
2020-07-21 18:01:25 +00:00
}
func (t *trackedWorker) ReadPiece(ctx context.Context, writer io.Writer, id storage.SectorRef, index storiface.UnpaddedByteIndex, size abi.UnpaddedPieceSize) (storiface.CallID, error) {
2020-09-23 12:56:37 +00:00
return t.tracker.track(t.wid, id, sealtasks.TTReadUnsealed)(t.Worker.ReadPiece(ctx, writer, id, index, size))
2020-07-21 18:01:25 +00:00
}
var _ Worker = &trackedWorker{}