stash/internal/api/resolver_subscription_job.go

44 lines
1.0 KiB
Go
Raw Normal View History

2021-05-24 04:24:18 +00:00
package api
import (
"context"
"github.com/stashapp/stash/internal/manager"
2021-05-24 04:24:18 +00:00
"github.com/stashapp/stash/pkg/job"
)
func makeJobStatusUpdate(t JobStatusUpdateType, j job.Job) *JobStatusUpdate {
return &JobStatusUpdate{
2021-05-24 04:24:18 +00:00
Type: t,
Job: jobToJobModel(j),
}
}
func (r *subscriptionResolver) JobsSubscribe(ctx context.Context) (<-chan *JobStatusUpdate, error) {
msg := make(chan *JobStatusUpdate, 100)
2021-05-24 04:24:18 +00:00
subscription := manager.GetInstance().JobManager.Subscribe(ctx)
go func() {
for {
select {
case j := <-subscription.NewJob:
msg <- makeJobStatusUpdate(JobStatusUpdateTypeAdd, j)
2021-05-24 04:24:18 +00:00
case j := <-subscription.RemovedJob:
msg <- makeJobStatusUpdate(JobStatusUpdateTypeRemove, j)
2021-05-24 04:24:18 +00:00
case j := <-subscription.UpdatedJob:
msg <- makeJobStatusUpdate(JobStatusUpdateTypeUpdate, j)
2021-05-24 04:24:18 +00:00
case <-ctx.Done():
close(msg)
return
}
}
}()
return msg, nil
}
func (r *subscriptionResolver) ScanCompleteSubscribe(ctx context.Context) (<-chan bool, error) {
return manager.GetInstance().ScanSubscribe(ctx), nil
}