stash/pkg/job/task.go
WithoutPants 5495d72849 File storage rewrite (#2676)
* Restructure data layer part 2 (#2599)
* Refactor and separate image model
* Refactor image query builder
* Handle relationships in image query builder
* Remove relationship management methods
* Refactor gallery model/query builder
* Add scenes to gallery model
* Convert scene model
* Refactor scene models
* Remove unused methods
* Add unit tests for gallery
* Add image tests
* Add scene tests
* Convert unnecessary scene value pointers to values
* Convert unnecessary pointer values to values
* Refactor scene partial
* Add scene partial tests
* Refactor ImagePartial
* Add image partial tests
* Refactor gallery partial update
* Add partial gallery update tests
* Use zero/null package for null values
* Add files and scan system
* Add sqlite implementation for files/folders
* Add unit tests for files/folders
* Image refactors
* Update image data layer
* Refactor gallery model and creation
* Refactor scene model
* Refactor scenes
* Don't set title from filename
* Allow galleries to freely add/remove images
* Add multiple scene file support to graphql and UI
* Add multiple file support for images in graphql/UI
* Add multiple file for galleries in graphql/UI
* Remove use of some deprecated fields
* Remove scene path usage
* Remove gallery path usage
* Remove path from image
* Move funscript to video file
* Refactor caption detection
* Migrate existing data
* Add post commit/rollback hook system
* Lint. Comment out import/export tests
* Add WithDatabase read only wrapper
* Prepend tasks to list
* Add 32 pre-migration
* Add warnings in release and migration notes
2022-09-06 07:03:42 +00:00

67 lines
1.1 KiB
Go

package job
import (
"context"
"github.com/remeh/sizedwaitgroup"
)
type taskExec struct {
task
fn func(ctx context.Context)
}
type TaskQueue struct {
p *Progress
wg sizedwaitgroup.SizedWaitGroup
tasks chan taskExec
done chan struct{}
}
func NewTaskQueue(ctx context.Context, p *Progress, queueSize int, processes int) *TaskQueue {
ret := &TaskQueue{
p: p,
wg: sizedwaitgroup.New(processes),
tasks: make(chan taskExec, queueSize),
done: make(chan struct{}),
}
go ret.executer(ctx)
return ret
}
func (tq *TaskQueue) Add(description string, fn func(ctx context.Context)) {
tq.tasks <- taskExec{
task: task{
description: description,
},
fn: fn,
}
}
func (tq *TaskQueue) Close() {
close(tq.tasks)
// wait for all tasks to finish
<-tq.done
}
func (tq *TaskQueue) executer(ctx context.Context) {
defer close(tq.done)
defer tq.wg.Wait()
for task := range tq.tasks {
if IsCancelled(ctx) {
return
}
tt := task
tq.wg.Add()
go func() {
defer tq.wg.Done()
tq.p.ExecuteTask(tt.description, func() {
tt.fn(ctx)
})
}()
}
}