jobrun: error handling through notification channel
This commit is contained in:
+9
-1
@@ -10,11 +10,13 @@ type Job struct {
|
|||||||
Name string
|
Name string
|
||||||
RunFunc func() (err error)
|
RunFunc func() (err error)
|
||||||
LastStart time.Time
|
LastStart time.Time
|
||||||
|
LastError error
|
||||||
Interval time.Duration
|
Interval time.Duration
|
||||||
Repeats bool
|
Repeats bool
|
||||||
}
|
}
|
||||||
|
|
||||||
type JobRunner struct {
|
type JobRunner struct {
|
||||||
|
notificationChan chan Job
|
||||||
newJobChan chan Job
|
newJobChan chan Job
|
||||||
finishedJobChan chan Job
|
finishedJobChan chan Job
|
||||||
scheduleTimer <-chan time.Time
|
scheduleTimer <-chan time.Time
|
||||||
@@ -24,6 +26,7 @@ type JobRunner struct {
|
|||||||
|
|
||||||
func NewJobRunner() *JobRunner {
|
func NewJobRunner() *JobRunner {
|
||||||
return &JobRunner{
|
return &JobRunner{
|
||||||
|
notificationChan: make(chan Job),
|
||||||
newJobChan: make(chan Job),
|
newJobChan: make(chan Job),
|
||||||
finishedJobChan: make(chan Job),
|
finishedJobChan: make(chan Job),
|
||||||
pending: make(map[string]Job),
|
pending: make(map[string]Job),
|
||||||
@@ -39,6 +42,10 @@ func (r *JobRunner) AddJob(j Job) {
|
|||||||
r.newJobChan <- j
|
r.newJobChan <- j
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (r *JobRunner) NotificationChan() <-chan Job {
|
||||||
|
return r.notificationChan
|
||||||
|
}
|
||||||
|
|
||||||
func (r *JobRunner) Start() {
|
func (r *JobRunner) Start() {
|
||||||
|
|
||||||
loop:
|
loop:
|
||||||
@@ -104,7 +111,8 @@ loop:
|
|||||||
|
|
||||||
go func(job Job) {
|
go func(job Job) {
|
||||||
if err := job.RunFunc(); err != nil {
|
if err := job.RunFunc(); err != nil {
|
||||||
panic(fmt.Sprintf("%#v", err)) // TODO better policy, store in job + notification channel?
|
job.LastError = err
|
||||||
|
r.notificationChan <- job
|
||||||
}
|
}
|
||||||
r.finishedJobChan <- job
|
r.finishedJobChan <- job
|
||||||
}(job)
|
}(job)
|
||||||
|
|||||||
Reference in New Issue
Block a user