@@ -16,6 +16,21 @@ import (
1616 "github.com/sirupsen/logrus"
1717)
1818
19+ type MvEngine struct {
20+ scheduler * gocron.Scheduler
21+ firstRunDone chan struct {}
22+ once sync.Once
23+ cfg util.Config
24+ }
25+
26+ func NewMvEngine (cfg util.Config ) * MvEngine {
27+ return & MvEngine {
28+ scheduler : gocron .NewScheduler (time .UTC ),
29+ firstRunDone : make (chan struct {}),
30+ cfg : cfg ,
31+ }
32+ }
33+
1934func TriggerMVE (cfg util.Config ) error {
2035 db , err := NewDb (cfg )
2136 if err != nil {
@@ -25,18 +40,16 @@ func TriggerMVE(cfg util.Config) error {
2540 return runInBackground (db , MVProcedures ).Wait ()
2641}
2742
28- func StartMVEScheduler (cfg util.Config ) {
29- mve := getMVE ()
30-
31- periodMinutes := cfg .DBMvCalcPeriodMinutes
43+ func (mve * MvEngine ) Start () {
44+ periodMinutes := mve .cfg .DBMvCalcPeriodMinutes
3245 if periodMinutes <= 0 {
3346 periodMinutes = 200
3447 }
3548
3649 logrus .Debugf ("MVE scheduling period set to %d minutes" , periodMinutes )
3750
3851 _ , err := mve .scheduler .Every (periodMinutes ).Minutes ().SingletonMode ().Do (func () {
39- err := TriggerMVE (cfg )
52+ err := TriggerMVE (mve . cfg )
4053 if err != nil {
4154 logrus .WithError (err ).Error ("MVE Trigger error" )
4255 }
@@ -55,47 +68,18 @@ func StartMVEScheduler(cfg util.Config) {
5568 mve .scheduler .StartAsync ()
5669}
5770
58- func StopMVEScheduler () {
59- mve := getMVE ()
60- mve .Shutdown ()
71+ func (mve * MvEngine ) Stop () {
72+ mve .scheduler .Clear ()
73+ // The following method is not advisory as it may hang for a long time:
74+ // mve.scheduler.Stop()
6175}
6276
63- func WaitMVEForFirstRun () {
64- mve := getMVE ()
77+ func (mve * MvEngine ) WaitForFirstRun () {
6578 <- mve .firstRunDone
6679}
6780
6881////////// Internals
6982
70- var mvEngine * MvEngine
71-
72- type MvEngine struct {
73- scheduler * gocron.Scheduler
74- firstRunDone chan struct {}
75- once sync.Once
76- }
77-
78- func NewMvEngine () * MvEngine {
79- return & MvEngine {
80- scheduler : gocron .NewScheduler (time .UTC ),
81- firstRunDone : make (chan struct {}),
82- }
83- }
84-
85- func getMVE () * MvEngine {
86- if mvEngine == nil {
87- mvEngine = NewMvEngine ()
88- }
89-
90- return mvEngine
91- }
92-
93- func (mve * MvEngine ) Shutdown () {
94- mve .scheduler .Clear ()
95- // This is not advisory as it may hang for a long time:
96- // mve.scheduler.Stop()
97- }
98-
9983type mveCtx struct {
10084 wg sync.WaitGroup
10185 mu sync.Mutex
@@ -126,13 +110,16 @@ func (mc *mveCtx) Wait() error {
126110 return nil
127111}
128112
129- func runInBackground (db Db , procs []MVProcedure ) * mveCtx {
113+ func runInBackground (db Db , procs [][] MVProcedure ) * mveCtx {
130114 mc := & mveCtx {}
131115
132- for i , p := range procs {
116+ for i , pl := range procs {
133117 mc .wg .Go (func () {
134- if err := TxCall (p , context .Background (), db ); err != nil {
135- mc .appendErrorMessage (fmt .Sprintf ("(procIdx: %d): %v" , i , err ))
118+ for j , p := range pl {
119+ if err := TxCall (p , context .Background (), db ); err != nil {
120+ mc .appendErrorMessage (fmt .Sprintf ("(procIdx: %d:%d): %v" , i , j , err ))
121+ break
122+ }
136123 }
137124 })
138125 }
0 commit comments