package runner_test
import (
"context"
"errors"
"time"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
"github.com/opencloud-eu/opencloud/pkg/runner"
)
func TimedTask(ch chan error, dur time.Duration) runner.Runable {
return func() error {
timer := time.NewTimer(dur)
defer timer.Stop()
var result error
select {
case <-timer.C:
case result = <-ch:
}
return result
}
}
var _ = Describe("Runner", func() {
Describe("Run", func() {
It("Context is done", func(ctx SpecContext) {
ch := make(chan error)
r := runner.New("run001", TimedTask(ch, 15*time.Second), func() {
ch <- nil
close(ch)
})
myCtx, cancel := context.WithTimeout(ctx, 1*time.Second)
defer cancel()
ch2 := make(chan *runner.Result)
go func(ch2 chan *runner.Result) {
ch2 <- r.Run(myCtx)
close(ch2)
}(ch2)
expectedResult := &runner.Result{
RunnerID: "run001",
RunnerError: nil,
}
Eventually(ctx, ch2).Should(Receive(Equal(expectedResult)))
}, SpecTimeout(5*time.Second))
It("Context is done and interrupt after", func(ctx SpecContext) {
ch := make(chan error)
r := runner.New("run001", TimedTask(ch, 15*time.Second), func() {
ch <- nil
close(ch)
})
myCtx, cancel := context.WithTimeout(ctx, 1*time.Second)
defer cancel()
ch2 := make(chan *runner.Result)
go func(ch2 chan *runner.Result) {
ch2 <- r.Run(myCtx)
close(ch2)
}(ch2)
expectedResult := &runner.Result{
RunnerID: "run001",
RunnerError: nil,
}
Eventually(ctx, ch2).Should(Receive(Equal(expectedResult)))
r.Interrupt()
}, SpecTimeout(5*time.Second))
It("Task finishes naturally", func(ctx SpecContext) {
e := errors.New("overslept!")
r := runner.New("run002", func() error {
time.Sleep(50 * time.Millisecond)
return e
}, func() {
})
myCtx, cancel := context.WithTimeout(ctx, 1*time.Second)
defer cancel()
ch2 := make(chan *runner.Result)
go func(ch2 chan *runner.Result) {
ch2 <- r.Run(myCtx)
close(ch2)
}(ch2)
expectedResult := &runner.Result{
RunnerID: "run002",
RunnerError: e,
}
Eventually(ctx, ch2).Should(Receive(Equal(expectedResult)))
}, SpecTimeout(5*time.Second))
It("Task doesn't finish", func(ctx SpecContext) {
r := runner.New("run003", func() error {
time.Sleep(20 * time.Second)
return nil
}, func() {
})
myCtx, cancel := context.WithTimeout(ctx, 1*time.Second)
defer cancel()
ch2 := make(chan *runner.Result)
go func(ch2 chan *runner.Result) {
ch2 <- r.Run(myCtx)
close(ch2)
}(ch2)
Consistently(ctx, ch2).WithTimeout(4500 * time.Millisecond).ShouldNot(Receive())
}, SpecTimeout(5*time.Second))
It("Task doesn't finish and times out", func(ctx SpecContext) {
r := runner.New("run003", func() error {
time.Sleep(20 * time.Second)
return nil
}, func() {
}, runner.WithInterruptDuration(3*time.Second))
myCtx, cancel := context.WithTimeout(ctx, 1*time.Second)
defer cancel()
ch2 := make(chan *runner.Result)
go func(ch2 chan *runner.Result) {
ch2 <- r.Run(myCtx)
close(ch2)
}(ch2)
var expectedResult *runner.Result
Eventually(ctx, ch2).Should(Receive(&expectedResult))
Expect(expectedResult.RunnerID).To(Equal("run003"))
var timeoutError *runner.TimeoutError
Expect(errors.As(expectedResult.RunnerError, &timeoutError)).To(BeTrue())
Expect(timeoutError.RunnerID).To(Equal("run003"))
Expect(timeoutError.Duration).To(Equal(3 * time.Second))
}, SpecTimeout(5*time.Second))
It("Run mutiple times panics", func(ctx SpecContext) {
e := errors.New("overslept!")
r := runner.New("run002", func() error {
time.Sleep(50 * time.Millisecond)
return e
}, func() {
})
myCtx, cancel := context.WithTimeout(ctx, 1*time.Second)
defer cancel()
Expect(func() {
r.Run(myCtx)
r.Run(myCtx)
}).To(Panic())
}, SpecTimeout(5*time.Second))
})
Describe("RunAsync", func() {
It("Wait in channel", func(ctx SpecContext) {
ch := make(chan *runner.Result)
e := errors.New("Task has finished")
r := runner.New("run004", func() error {
time.Sleep(50 * time.Millisecond)
return e
}, func() {
})
r.RunAsync(ch)
expectedResult := &runner.Result{
RunnerID: "run004",
RunnerError: e,
}
Eventually(ctx, ch).Should(Receive(Equal(expectedResult)))
}, SpecTimeout(5*time.Second))
It("Run multiple times panics", func(ctx SpecContext) {
ch := make(chan *runner.Result)
e := errors.New("Task has finished")
r := runner.New("run004", func() error {
time.Sleep(50 * time.Millisecond)
return e
}, func() {
})
r.RunAsync(ch)
Expect(func() {
r.RunAsync(ch)
}).To(Panic())
}, SpecTimeout(5*time.Second))
It("Interrupt async", func(ctx SpecContext) {
ch := make(chan *runner.Result)
e := errors.New("Task interrupted")
taskCh := make(chan error)
r := runner.New("run005", TimedTask(taskCh, 20*time.Second), func() {
taskCh <- e
close(taskCh)
})
r.RunAsync(ch)
r.Interrupt()
expectedResult := &runner.Result{
RunnerID: "run005",
RunnerError: e,
}
Eventually(ctx, ch).Should(Receive(Equal(expectedResult)))
}, SpecTimeout(5*time.Second))
It("Interrupt async times out", func(ctx SpecContext) {
ch := make(chan *runner.Result)
e := errors.New("Task interrupted")
r := runner.New("run005", func() error {
time.Sleep(30 * time.Second)
return e
}, func() {
}, runner.WithInterruptDuration(3*time.Second))
r.RunAsync(ch)
r.Interrupt()
var expectedResult *runner.Result
Eventually(ctx, ch).Should(Receive(&expectedResult))
Expect(expectedResult.RunnerID).To(Equal("run005"))
Expect(expectedResult.RunnerError.Error()).To(ContainSubstring("timed out"))
}, SpecTimeout(5*time.Second))
It("Interrupt async multiple times", func(ctx SpecContext) {
ch := make(chan *runner.Result)
e := errors.New("Task interrupted")
taskCh := make(chan error)
r := runner.New("run005", TimedTask(taskCh, 20*time.Second), func() {
taskCh <- e
close(taskCh)
})
r.RunAsync(ch)
r.Interrupt()
r.Interrupt()
r.Interrupt()
expectedResult := &runner.Result{
RunnerID: "run005",
RunnerError: e,
}
Eventually(ctx, ch).Should(Receive(Equal(expectedResult)))
}, SpecTimeout(5*time.Second))
})
Describe("Finished", func() {
It("Finish channel closes", func(ctx SpecContext) {
r := runner.New("run006", func() error {
time.Sleep(50 * time.Millisecond)
return nil
}, func() {
})
myCtx, cancel := context.WithTimeout(ctx, 1*time.Second)
defer cancel()
ch2 := make(chan *runner.Result)
go func(ch2 chan *runner.Result) {
ch2 <- r.Run(myCtx)
close(ch2)
}(ch2)
finishedCh := r.Finished()
Eventually(ctx, finishedCh).Should(BeClosed())
}, SpecTimeout(5*time.Second))
})
})