mirror of
https://github.com/samuelncui/yatm.git
synced 2026-09-06 16:16:53 +00:00
fix: change tape error
This commit is contained in:
@@ -106,7 +106,7 @@ func (e *Executor) Submit(ctx context.Context, job *Job, param *entity.JobNextPa
|
||||
return err
|
||||
}
|
||||
|
||||
exe.submit(param.GetArchive())
|
||||
exe.submit(ctx, param.GetArchive())
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
+5
-5
@@ -109,19 +109,19 @@ func (e *Executor) GetJob(ctx context.Context, id int64) (*Job, error) {
|
||||
func (e *Executor) ListJob(ctx context.Context, filter *entity.JobFilter) ([]*Job, error) {
|
||||
db := e.db.WithContext(ctx)
|
||||
if filter.Status != nil {
|
||||
db.Where("status = ?", *filter.Status)
|
||||
db = db.Where("status = ?", *filter.Status)
|
||||
}
|
||||
|
||||
if filter.Limit != nil {
|
||||
db.Limit(int(*filter.Limit))
|
||||
db = db.Limit(int(*filter.Limit))
|
||||
} else {
|
||||
db.Limit(20)
|
||||
db = db.Limit(20)
|
||||
}
|
||||
if filter.Offset != nil {
|
||||
db.Offset(int(*filter.Offset))
|
||||
db = db.Offset(int(*filter.Offset))
|
||||
}
|
||||
|
||||
db.Order("create_time DESC")
|
||||
db = db.Order("create_time DESC")
|
||||
|
||||
jobs := make([]*Job, 0, 20)
|
||||
if r := db.Find(&jobs); r.Error != nil {
|
||||
|
||||
@@ -15,7 +15,10 @@ func (e *Executor) getArchiveDisplay(ctx context.Context, job *Job) (*entity.Job
|
||||
display.CopyedFiles = atomic.LoadInt64(&exe.progress.files)
|
||||
display.TotalBytes = atomic.LoadInt64(&exe.progress.totalBytes)
|
||||
display.TotalFiles = atomic.LoadInt64(&exe.progress.totalFiles)
|
||||
display.Speed = atomic.LoadInt64(&exe.progress.speed)
|
||||
display.StartTime = exe.progress.startTime.Unix()
|
||||
|
||||
speed := atomic.LoadInt64(&exe.progress.speed)
|
||||
display.Speed = &speed
|
||||
}
|
||||
|
||||
return display, nil
|
||||
|
||||
+121
-101
@@ -46,15 +46,13 @@ func (e *Executor) newArchiveExecutor(ctx context.Context, job *Job) (*jobArchiv
|
||||
logger.SetOutput(io.MultiWriter(os.Stderr, logFile))
|
||||
|
||||
exe := &jobArchiveExecutor{
|
||||
ctx: context.Background(),
|
||||
exe: e,
|
||||
job: job,
|
||||
|
||||
state: job.State.GetArchive(),
|
||||
|
||||
progress: new(progress),
|
||||
logFile: logFile,
|
||||
logger: logger,
|
||||
logFile: logFile,
|
||||
logger: logger,
|
||||
}
|
||||
|
||||
runningArchives.Store(job.ID, exe)
|
||||
@@ -62,7 +60,6 @@ func (e *Executor) newArchiveExecutor(ctx context.Context, job *Job) (*jobArchiv
|
||||
}
|
||||
|
||||
type jobArchiveExecutor struct {
|
||||
ctx context.Context
|
||||
exe *Executor
|
||||
job *Job
|
||||
|
||||
@@ -74,22 +71,23 @@ type jobArchiveExecutor struct {
|
||||
logger *logrus.Logger
|
||||
}
|
||||
|
||||
func (a *jobArchiveExecutor) submit(param *entity.JobArchiveNextParam) {
|
||||
if err := a.handle(param); err != nil {
|
||||
a.logger.WithContext(a.ctx).Infof("handler param fail, err= %w", err)
|
||||
func (a *jobArchiveExecutor) submit(ctx context.Context, param *entity.JobArchiveNextParam) {
|
||||
if err := a.handle(ctx, param); err != nil {
|
||||
a.logger.WithContext(ctx).Infof("handler param fail, err= %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
func (a *jobArchiveExecutor) handle(param *entity.JobArchiveNextParam) error {
|
||||
func (a *jobArchiveExecutor) handle(ctx context.Context, param *entity.JobArchiveNextParam) error {
|
||||
if p := param.GetCopying(); p != nil {
|
||||
if err := a.switchStep(entity.JobArchiveStep_Copying, entity.JobStatus_Processing, mapset.NewThreadUnsafeSet(entity.JobArchiveStep_WaitForTape)); err != nil {
|
||||
if err := a.switchStep(ctx, entity.JobArchiveStep_Copying, entity.JobStatus_Processing, mapset.NewThreadUnsafeSet(entity.JobArchiveStep_WaitForTape)); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
go tools.Wrap(a.ctx, func() {
|
||||
_, err := a.makeTape(p.Device, p.Barcode, p.Name)
|
||||
if err != nil {
|
||||
a.logger.WithContext(a.ctx).WithError(err).Errorf("make type has error, barcode= '%s' name= '%s'", p.Barcode, p.Name)
|
||||
tools.Working()
|
||||
go tools.Wrap(ctx, func() {
|
||||
defer tools.Done()
|
||||
if err := a.makeTape(tools.ShutdownContext, p.Device, p.Barcode, p.Name); err != nil {
|
||||
a.logger.WithContext(ctx).WithError(err).Errorf("make type has error, barcode= '%s' name= '%s'", p.Barcode, p.Name)
|
||||
}
|
||||
})
|
||||
|
||||
@@ -97,11 +95,11 @@ func (a *jobArchiveExecutor) handle(param *entity.JobArchiveNextParam) error {
|
||||
}
|
||||
|
||||
if p := param.GetWaitForTape(); p != nil {
|
||||
return a.switchStep(entity.JobArchiveStep_WaitForTape, entity.JobStatus_Processing, mapset.NewThreadUnsafeSet(entity.JobArchiveStep_Pending, entity.JobArchiveStep_Copying))
|
||||
return a.switchStep(ctx, entity.JobArchiveStep_WaitForTape, entity.JobStatus_Processing, mapset.NewThreadUnsafeSet(entity.JobArchiveStep_Pending, entity.JobArchiveStep_Copying))
|
||||
}
|
||||
|
||||
if p := param.GetFinished(); p != nil {
|
||||
if err := a.switchStep(entity.JobArchiveStep_Finished, entity.JobStatus_Completed, mapset.NewThreadUnsafeSet(entity.JobArchiveStep_Copying)); err != nil {
|
||||
if err := a.switchStep(ctx, entity.JobArchiveStep_Finished, entity.JobStatus_Completed, mapset.NewThreadUnsafeSet(entity.JobArchiveStep_Copying)); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -113,51 +111,48 @@ func (a *jobArchiveExecutor) handle(param *entity.JobArchiveNextParam) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (a *jobArchiveExecutor) makeTape(device, barcode, name string) (*library.Tape, error) {
|
||||
func (a *jobArchiveExecutor) makeTape(ctx context.Context, device, barcode, name string) (rerr error) {
|
||||
if !a.exe.occupyDevice(device) {
|
||||
return nil, fmt.Errorf("device is using, device= %s", device)
|
||||
return fmt.Errorf("device is using, device= %s", device)
|
||||
}
|
||||
defer a.exe.releaseDevice(device)
|
||||
defer a.makeTapeFinished()
|
||||
defer a.makeTapeFinished(tools.WithoutTimeout(ctx))
|
||||
|
||||
encryption, keyPath, keyRecycle, err := a.exe.newKey()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return err
|
||||
}
|
||||
defer func() {
|
||||
time.Sleep(time.Second)
|
||||
keyRecycle()
|
||||
}()
|
||||
defer keyRecycle()
|
||||
|
||||
if err := runCmd(a.logger, a.exe.makeEncryptCmd(a.ctx, device, keyPath, barcode, name)); err != nil {
|
||||
return nil, fmt.Errorf("run encrypt script fail, %w", err)
|
||||
if err := runCmd(a.logger, a.exe.makeEncryptCmd(ctx, device, keyPath, barcode, name)); err != nil {
|
||||
return fmt.Errorf("run encrypt script fail, %w", err)
|
||||
}
|
||||
|
||||
mkfsCmd := exec.CommandContext(a.ctx, a.exe.mkfsScript)
|
||||
mkfsCmd := exec.CommandContext(ctx, a.exe.mkfsScript)
|
||||
mkfsCmd.Env = append(mkfsCmd.Env, fmt.Sprintf("DEVICE=%s", device), fmt.Sprintf("TAPE_BARCODE=%s", barcode), fmt.Sprintf("TAPE_NAME=%s", name))
|
||||
if err := runCmd(a.logger, mkfsCmd); err != nil {
|
||||
return nil, fmt.Errorf("run mkfs script fail, %w", err)
|
||||
return fmt.Errorf("run mkfs script fail, %w", err)
|
||||
}
|
||||
|
||||
mountPoint, err := os.MkdirTemp("", "*.ltfs")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create temp mountpoint, %w", err)
|
||||
return fmt.Errorf("create temp mountpoint, %w", err)
|
||||
}
|
||||
|
||||
mountCmd := exec.CommandContext(a.ctx, a.exe.mountScript)
|
||||
mountCmd := exec.CommandContext(ctx, a.exe.mountScript)
|
||||
mountCmd.Env = append(mountCmd.Env, fmt.Sprintf("DEVICE=%s", device), fmt.Sprintf("MOUNT_POINT=%s", mountPoint))
|
||||
if err := runCmd(a.logger, mountCmd); err != nil {
|
||||
return nil, fmt.Errorf("run mount script fail, %w", err)
|
||||
return fmt.Errorf("run mount script fail, %w", err)
|
||||
}
|
||||
defer func() {
|
||||
umountCmd := exec.CommandContext(a.ctx, a.exe.umountScript)
|
||||
umountCmd := exec.CommandContext(tools.WithoutTimeout(ctx), a.exe.umountScript)
|
||||
umountCmd.Env = append(umountCmd.Env, fmt.Sprintf("MOUNT_POINT=%s", mountPoint))
|
||||
if err := runCmd(a.logger, umountCmd); err != nil {
|
||||
a.logger.WithContext(a.ctx).WithError(err).Errorf("run umount script fail, %s", mountPoint)
|
||||
a.logger.WithContext(ctx).WithError(err).Errorf("run umount script fail, %s", mountPoint)
|
||||
return
|
||||
}
|
||||
if err := os.Remove(mountPoint); err != nil {
|
||||
a.logger.WithContext(a.ctx).WithError(err).Errorf("remove mount point fail, %s", mountPoint)
|
||||
a.logger.WithContext(ctx).WithError(err).Errorf("remove mount point fail, %s", mountPoint)
|
||||
return
|
||||
}
|
||||
}()
|
||||
@@ -177,6 +172,10 @@ func (a *jobArchiveExecutor) makeTape(device, barcode, name string) (*library.Ta
|
||||
|
||||
reportHander, reportGetter := acp.NewReportGetter()
|
||||
opts = append(opts, acp.WithEventHandler(reportHander))
|
||||
|
||||
a.progress = newProgress()
|
||||
defer func() { a.progress = nil }()
|
||||
|
||||
opts = append(opts, acp.WithEventHandler(func(ev acp.Event) {
|
||||
switch e := ev.(type) {
|
||||
case *acp.EventUpdateCount:
|
||||
@@ -184,9 +183,12 @@ func (a *jobArchiveExecutor) makeTape(device, barcode, name string) (*library.Ta
|
||||
atomic.StoreInt64(&a.progress.totalFiles, e.Files)
|
||||
return
|
||||
case *acp.EventUpdateProgress:
|
||||
atomic.StoreInt64(&a.progress.bytes, e.Bytes)
|
||||
a.progress.setBytes(e.Bytes)
|
||||
atomic.StoreInt64(&a.progress.files, e.Files)
|
||||
return
|
||||
case *acp.EventReportError:
|
||||
a.logger.WithContext(ctx).Errorf("acp report error, src= '%s' dst= '%s' err= '%s'", e.Error.Src, e.Error.Dst, e.Error.Err)
|
||||
return
|
||||
case *acp.EventUpdateJob:
|
||||
job := e.Job
|
||||
src := entity.NewSourceFromACPJob(job)
|
||||
@@ -196,11 +198,17 @@ func (a *jobArchiveExecutor) makeTape(device, barcode, name string) (*library.Ta
|
||||
case "pending":
|
||||
targetStatus = entity.CopyStatus_Pending
|
||||
case "preparing":
|
||||
a.logger.Infof("file '%s' starts to prepare for copy, size= %d", src.RealPath(), job.Size)
|
||||
targetStatus = entity.CopyStatus_Running
|
||||
case "finished":
|
||||
a.logger.Infof("file '%s' copy finished, size= %d", src.RealPath(), job.Size)
|
||||
a.logger.WithContext(ctx).Infof("file '%s' copy finished, size= %d", src.RealPath(), job.Size)
|
||||
targetStatus = entity.CopyStatus_Staged
|
||||
|
||||
for dst, err := range job.FailTargets {
|
||||
if err == nil {
|
||||
continue
|
||||
}
|
||||
a.logger.WithContext(ctx).WithError(err).Errorf("file '%s' copy fail, dst= '%s'", src.RealPath(), dst)
|
||||
}
|
||||
default:
|
||||
return
|
||||
}
|
||||
@@ -218,72 +226,89 @@ func (a *jobArchiveExecutor) makeTape(device, barcode, name string) (*library.Ta
|
||||
}
|
||||
target.Status = targetStatus
|
||||
|
||||
if _, err := a.exe.SaveJob(a.ctx, a.job); err != nil {
|
||||
logrus.WithContext(a.ctx).Infof("save job for update file fail, name= %s", job.Base+path.Join(job.Path...))
|
||||
if _, err := a.exe.SaveJob(ctx, a.job); err != nil {
|
||||
logrus.WithContext(ctx).Infof("save job for update file fail, name= %s", job.Base+path.Join(job.Path...))
|
||||
}
|
||||
return
|
||||
}
|
||||
}))
|
||||
|
||||
copyer, err := acp.New(a.ctx, opts...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("start copy fail, %w", err)
|
||||
}
|
||||
copyer.Wait()
|
||||
defer func() {
|
||||
ctx := tools.WithoutTimeout(ctx)
|
||||
|
||||
report := reportGetter()
|
||||
sort.Slice(report.Jobs, func(i, j int) bool {
|
||||
return entity.NewSourceFromACPJob(report.Jobs[i]).Compare(entity.NewSourceFromACPJob(report.Jobs[j])) < 0
|
||||
})
|
||||
|
||||
filteredJobs := make([]*acp.Job, 0, len(report.Jobs))
|
||||
files := make([]*library.TapeFile, 0, len(report.Jobs))
|
||||
for _, job := range report.Jobs {
|
||||
if len(job.SuccessTargets) == 0 {
|
||||
continue
|
||||
}
|
||||
if !job.Mode.IsRegular() {
|
||||
continue
|
||||
}
|
||||
|
||||
hash, err := hex.DecodeString(job.SHA256)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("decode sha256 fail, err= %w", err)
|
||||
}
|
||||
|
||||
files = append(files, &library.TapeFile{
|
||||
Path: path.Join(job.Path...),
|
||||
Size: job.Size,
|
||||
Mode: job.Mode,
|
||||
ModTime: job.ModTime,
|
||||
WriteTime: job.WriteTime,
|
||||
Hash: hash,
|
||||
report := reportGetter()
|
||||
sort.Slice(report.Jobs, func(i, j int) bool {
|
||||
return entity.NewSourceFromACPJob(report.Jobs[i]).Compare(entity.NewSourceFromACPJob(report.Jobs[j])) < 0
|
||||
})
|
||||
filteredJobs = append(filteredJobs, job)
|
||||
}
|
||||
|
||||
tape, err := a.exe.lib.CreateTape(a.ctx, &library.Tape{
|
||||
Barcode: barcode,
|
||||
Name: name,
|
||||
Encryption: encryption,
|
||||
CreateTime: time.Now(),
|
||||
}, files)
|
||||
reportFile, err := a.exe.newReportWriter(barcode)
|
||||
if err != nil {
|
||||
a.logger.WithContext(ctx).WithError(err).Warnf("open report file fail, barcode= '%s'", barcode)
|
||||
} else {
|
||||
defer reportFile.Close()
|
||||
reportFile.Write([]byte(report.ToJSONString(false)))
|
||||
}
|
||||
|
||||
filteredJobs := make([]*acp.Job, 0, len(report.Jobs))
|
||||
files := make([]*library.TapeFile, 0, len(report.Jobs))
|
||||
for _, job := range report.Jobs {
|
||||
if len(job.SuccessTargets) == 0 {
|
||||
continue
|
||||
}
|
||||
if !job.Mode.IsRegular() {
|
||||
continue
|
||||
}
|
||||
|
||||
hash, err := hex.DecodeString(job.SHA256)
|
||||
if err != nil {
|
||||
a.logger.WithContext(ctx).WithError(err).Warnf("decode sha256 fail, path= '%s'", entity.NewSourceFromACPJob(job).RealPath())
|
||||
continue
|
||||
}
|
||||
|
||||
files = append(files, &library.TapeFile{
|
||||
Path: path.Join(job.Path...),
|
||||
Size: job.Size,
|
||||
Mode: job.Mode,
|
||||
ModTime: job.ModTime,
|
||||
WriteTime: job.WriteTime,
|
||||
Hash: hash,
|
||||
})
|
||||
filteredJobs = append(filteredJobs, job)
|
||||
}
|
||||
|
||||
tape, err := a.exe.lib.CreateTape(ctx, &library.Tape{
|
||||
Barcode: barcode,
|
||||
Name: name,
|
||||
Encryption: encryption,
|
||||
CreateTime: time.Now(),
|
||||
}, files)
|
||||
if err != nil {
|
||||
rerr = tools.AppendError(rerr, fmt.Errorf("create tape fail, barcode= '%s' name= '%s', %w", barcode, name, err))
|
||||
return
|
||||
}
|
||||
a.logger.Infof("create tape success, tape_id= %d", tape.ID)
|
||||
|
||||
if err := a.exe.lib.TrimFiles(ctx); err != nil {
|
||||
a.logger.WithError(err).Warnf("trim library files fail")
|
||||
}
|
||||
|
||||
if err := a.markSourcesAsSubmited(ctx, filteredJobs); err != nil {
|
||||
rerr = tools.AppendError(rerr, fmt.Errorf("mark source as submited fail, %w", err))
|
||||
return
|
||||
}
|
||||
}()
|
||||
|
||||
copyer, err := acp.New(ctx, opts...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create tape fail, barcode= '%s' name= '%s', %w", barcode, name, err)
|
||||
}
|
||||
if err := a.exe.lib.TrimFiles(a.ctx); err != nil {
|
||||
a.logger.WithError(err).Warnf("trim library files fail")
|
||||
rerr = fmt.Errorf("start copy fail, %w", err)
|
||||
return
|
||||
}
|
||||
|
||||
if err := a.markSourcesAsSubmited(filteredJobs); err != nil {
|
||||
a.submit(&entity.JobArchiveNextParam{Param: &entity.JobArchiveNextParam_WaitForTape{WaitForTape: &entity.JobArchiveWaitForTapeParam{}}})
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return tape, nil
|
||||
copyer.Wait()
|
||||
return
|
||||
}
|
||||
|
||||
func (a *jobArchiveExecutor) switchStep(target entity.JobArchiveStep, status entity.JobStatus, expect mapset.Set[entity.JobArchiveStep]) error {
|
||||
func (a *jobArchiveExecutor) switchStep(ctx context.Context, target entity.JobArchiveStep, status entity.JobStatus, expect mapset.Set[entity.JobArchiveStep]) error {
|
||||
a.stateLock.Lock()
|
||||
defer a.stateLock.Unlock()
|
||||
|
||||
@@ -293,14 +318,14 @@ func (a *jobArchiveExecutor) switchStep(target entity.JobArchiveStep, status ent
|
||||
|
||||
a.state.Step = target
|
||||
a.job.Status = status
|
||||
if _, err := a.exe.SaveJob(a.ctx, a.job); err != nil {
|
||||
if _, err := a.exe.SaveJob(ctx, a.job); err != nil {
|
||||
return fmt.Errorf("switch to step copying, save job fail, %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (a *jobArchiveExecutor) markSourcesAsSubmited(jobs []*acp.Job) error {
|
||||
func (a *jobArchiveExecutor) markSourcesAsSubmited(ctx context.Context, jobs []*acp.Job) error {
|
||||
a.stateLock.Lock()
|
||||
defer a.stateLock.Unlock()
|
||||
|
||||
@@ -322,14 +347,9 @@ func (a *jobArchiveExecutor) markSourcesAsSubmited(jobs []*acp.Job) error {
|
||||
target.Status = entity.CopyStatus_Submited
|
||||
}
|
||||
|
||||
if _, err := a.exe.SaveJob(a.ctx, a.job); err != nil {
|
||||
if _, err := a.exe.SaveJob(ctx, a.job); err != nil {
|
||||
return fmt.Errorf("mark sources as submited, save job, %w", err)
|
||||
}
|
||||
|
||||
atomic.StoreInt64(&a.progress.bytes, 0)
|
||||
atomic.StoreInt64(&a.progress.files, 0)
|
||||
atomic.StoreInt64(&a.progress.totalBytes, 0)
|
||||
atomic.StoreInt64(&a.progress.totalFiles, 0)
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -348,10 +368,10 @@ func (a *jobArchiveExecutor) getTodoSources() int {
|
||||
return todo
|
||||
}
|
||||
|
||||
func (a *jobArchiveExecutor) makeTapeFinished() {
|
||||
func (a *jobArchiveExecutor) makeTapeFinished(ctx context.Context) {
|
||||
if a.getTodoSources() > 0 {
|
||||
a.submit(&entity.JobArchiveNextParam{Param: &entity.JobArchiveNextParam_WaitForTape{WaitForTape: &entity.JobArchiveWaitForTapeParam{}}})
|
||||
a.submit(ctx, &entity.JobArchiveNextParam{Param: &entity.JobArchiveNextParam_WaitForTape{WaitForTape: &entity.JobArchiveWaitForTapeParam{}}})
|
||||
} else {
|
||||
a.submit(&entity.JobArchiveNextParam{Param: &entity.JobArchiveNextParam_Finished{Finished: &entity.JobArchiveFinishedParam{}}})
|
||||
a.submit(ctx, &entity.JobArchiveNextParam{Param: &entity.JobArchiveNextParam_Finished{Finished: &entity.JobArchiveFinishedParam{}}})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -48,3 +48,21 @@ func runCmd(logger *logrus.Logger, cmd *exec.Cmd) error {
|
||||
|
||||
return cmd.Run()
|
||||
}
|
||||
|
||||
func (e *Executor) reportPath(barcode string) (string, string) {
|
||||
return path.Join(e.workDirectory, "write-reports"), fmt.Sprintf("%s.log", barcode)
|
||||
}
|
||||
|
||||
func (e *Executor) newReportWriter(barcode string) (*os.File, error) {
|
||||
dir, filename := e.reportPath(barcode)
|
||||
if err := os.MkdirAll(dir, 0755); err != nil {
|
||||
return nil, fmt.Errorf("make job log dir fail, path= '%s', err= %w", dir, err)
|
||||
}
|
||||
|
||||
file, err := os.OpenFile(path.Join(dir, filename), os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0644)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create file fail, path= '%s', err= %w", path.Join(dir, filename), err)
|
||||
}
|
||||
|
||||
return file, nil
|
||||
}
|
||||
|
||||
+48
-1
@@ -1,8 +1,55 @@
|
||||
package executor
|
||||
|
||||
import (
|
||||
"sync/atomic"
|
||||
"time"
|
||||
)
|
||||
|
||||
const SpeedLen = 30
|
||||
|
||||
type speedEvent struct {
|
||||
bytes int64
|
||||
time time.Time
|
||||
}
|
||||
|
||||
type progress struct {
|
||||
speed int64
|
||||
speedEvents []speedEvent
|
||||
speedLen int
|
||||
speedIdx int
|
||||
|
||||
startTime time.Time
|
||||
speed int64
|
||||
|
||||
totalBytes, totalFiles int64
|
||||
bytes, files int64
|
||||
}
|
||||
|
||||
func newProgress() *progress {
|
||||
return &progress{speedEvents: make([]speedEvent, SpeedLen), speedLen: SpeedLen, startTime: time.Now()}
|
||||
}
|
||||
|
||||
func (p *progress) setBytes(bytes int64) {
|
||||
atomic.StoreInt64(&p.bytes, bytes)
|
||||
now := time.Now()
|
||||
|
||||
p.speedEvents[p.speedIdx] = speedEvent{bytes: bytes, time: now}
|
||||
for earliest := p.speedIdx; ; {
|
||||
earliest++
|
||||
if earliest >= p.speedLen {
|
||||
earliest = 0
|
||||
}
|
||||
if earliest == p.speedIdx {
|
||||
break
|
||||
}
|
||||
|
||||
if !p.speedEvents[earliest].time.IsZero() {
|
||||
p.speed = (bytes - p.speedEvents[earliest].bytes) * 1e9 / now.Sub(p.speedEvents[earliest].time).Nanoseconds()
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
p.speedIdx++
|
||||
if p.speedIdx >= p.speedLen {
|
||||
p.speedIdx = 0
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user