From 158b96f1736794531c49676c196e687d321d1611 Mon Sep 17 00:00:00 2001 From: Radu Ioan Fericean Date: Wed, 2 Dec 2015 10:58:41 +0200 Subject: [PATCH 1/8] migrate action plans --- cmd/cgr-loader/cgr-loader.go | 7 +++++- cmd/cgr-loader/migrator_rc8.go | 42 ++++++++++++++++++++++++++++++++++ 2 files changed, 48 insertions(+), 1 deletion(-) diff --git a/cmd/cgr-loader/cgr-loader.go b/cmd/cgr-loader/cgr-loader.go index ac359032f..db41efc3d 100644 --- a/cmd/cgr-loader/cgr-loader.go +++ b/cmd/cgr-loader/cgr-loader.go @@ -36,7 +36,7 @@ import ( var ( //separator = flag.String("separator", ",", "Default field separator") cgrConfig, _ = config.NewDefaultCGRConfig() - migrateRC8 = flag.String("migrate_rc8", "", "Migrate Accounts, Actions, ActionTriggers and DerivedChargers to RC8 structures, possible values: *all,acc,atr,act,dcs") + migrateRC8 = flag.String("migrate_rc8", "", "Migrate Accounts, Actions, ActionTriggers and DerivedChargers to RC8 structures, possible values: *all,acc,atr,act,dcs,apl") tpdb_type = flag.String("tpdb_type", cgrConfig.TpDbType, "The type of the TariffPlan database ") tpdb_host = flag.String("tpdb_host", cgrConfig.TpDbHost, "The TariffPlan host to connect to.") tpdb_port = flag.String("tpdb_port", cgrConfig.TpDbPort, "The TariffPlan port to bind to.") @@ -142,6 +142,11 @@ func main() { log.Print(err.Error()) } } + if strings.Contains(*migrateRC8, "apl") || strings.Contains(*migrateRC8, "*all") { + if err := migratorRC8rat.migrateActionPlans(); err != nil { + log.Print(err.Error()) + } + } log.Print("Done!") return } diff --git a/cmd/cgr-loader/migrator_rc8.go b/cmd/cgr-loader/migrator_rc8.go index e237f1dbb..a1f11378b 100644 --- a/cmd/cgr-loader/migrator_rc8.go +++ b/cmd/cgr-loader/migrator_rc8.go @@ -425,3 +425,45 @@ func (mig MigratorRC8) migrateDerivedChargers() error { } return nil } + +func (mig MigratorRC8) migrateActionPlans() error { + keys, err := mig.db.Cmd("KEYS", utils.ACTION_PLAN_PREFIX+"*").List() + if err != nil { + return err + } + aplsMap := make(map[string]engine.ActionPlans, len(keys)) + for _, key := range keys { + log.Printf("Migrating action plans: %s...", key) + var apls engine.ActionPlans + var values []byte + if values, err = mig.db.Cmd("GET", key).Bytes(); err == nil { + if err := mig.ms.Unmarshal(values, &apls); err != nil { + return err + } + } + // change all AccountIds + for _, apl := range apls { + for idx, actionId := range apl.AccountIds { + // fix id + idElements := strings.Split(actionId, utils.CONCATENATED_KEY_SEP) + if len(idElements) != 3 { + log.Printf("Malformed account ID %s", actionId) + continue + } + apl.AccountIds[idx] = fmt.Sprintf("%s:%s", idElements[1], idElements[2]) + } + } + aplsMap[key] = apls + } + // write data back + for key, apl := range aplsMap { + result, err := mig.ms.Marshal(apl) + if err != nil { + return err + } + if err = mig.db.Cmd("SET", key, result).Err; err != nil { + return err + } + } + return nil +} From b78a70ea21d859ecb09557f624fb455a6f550a62 Mon Sep 17 00:00:00 2001 From: Radu Ioan Fericean Date: Wed, 2 Dec 2015 17:15:38 +0200 Subject: [PATCH 2/8] added scheduler queue console command + other scheduler improvements --- apier/v1/scheduler.go | 59 +++++++++++++++++------------------ cmd/cgr-engine/rater.go | 2 +- console/scheduler_queue.go | 63 ++++++++++++++++++++++++++++++++++++++ scheduler/scheduler.go | 14 ++++++--- 4 files changed, 103 insertions(+), 35 deletions(-) create mode 100644 console/scheduler_queue.go diff --git a/apier/v1/scheduler.go b/apier/v1/scheduler.go index ec02807a5..5fbcb7591 100644 --- a/apier/v1/scheduler.go +++ b/apier/v1/scheduler.go @@ -110,22 +110,11 @@ type ScheduledActions struct { } func (self *ApierV1) GetScheduledActions(attrs AttrsGetScheduledActions, reply *[]*ScheduledActions) error { - schedActions := make([]*ScheduledActions, 0) if self.Sched == nil { return errors.New("SCHEDULER_NOT_ENABLED") } + var schedActions []*ScheduledActions scheduledActions := self.Sched.GetQueue() - var min, max int - if attrs.Paginator.Offset != nil { - min = *attrs.Paginator.Offset - } - if attrs.Paginator.Limit != nil { - max = *attrs.Paginator.Limit - } - if max > len(scheduledActions) { - max = len(scheduledActions) - } - scheduledActions = scheduledActions[min : min+max] for _, qActions := range scheduledActions { sas := &ScheduledActions{ActionsId: qActions.ActionsId, ActionPlanId: qActions.Id, ActionPlanUuid: qActions.Uuid} if attrs.SearchTerm != "" && @@ -140,28 +129,40 @@ func (self *ApierV1) GetScheduledActions(attrs AttrsGetScheduledActions, reply * if !attrs.TimeEnd.IsZero() && (sas.NextRunTime.After(attrs.TimeEnd) || sas.NextRunTime.Equal(attrs.TimeEnd)) { continue } - acntFiltersMatch := false - for _, acntKey := range qActions.AccountIds { - tenantMatched := len(attrs.Tenant) == 0 - accountMatched := len(attrs.Account) == 0 - dta, _ := utils.NewTAFromAccountKey(acntKey) - sas.Accounts = append(sas.Accounts, dta) - // One member matching - if !tenantMatched && attrs.Tenant == dta.Tenant { - tenantMatched = true + // filter on account + if attrs.Tenant != "" || attrs.Account != "" { + found := false + for _, accID := range qActions.AccountIds { + split := strings.Split(accID, utils.CONCATENATED_KEY_SEP) + if len(split) != 2 { + continue // malformed account id + } + if attrs.Tenant != "" && attrs.Tenant != split[0] { + continue + } + if attrs.Account != "" && attrs.Account != split[1] { + continue + } + found = true + break } - if !accountMatched && attrs.Account == dta.Account { - accountMatched = true - } - if tenantMatched && accountMatched { - acntFiltersMatch = true + if !found { + continue } } - if !acntFiltersMatch { - continue - } + // we have a winner schedActions = append(schedActions, sas) } + if attrs.Paginator.Offset != nil { + if *attrs.Paginator.Offset <= len(schedActions) { + schedActions = schedActions[*attrs.Paginator.Offset:] + } + } + if attrs.Paginator.Limit != nil { + if *attrs.Paginator.Limit <= len(schedActions) { + schedActions = schedActions[:*attrs.Paginator.Limit] + } + } *reply = schedActions return nil } diff --git a/cmd/cgr-engine/rater.go b/cmd/cgr-engine/rater.go index f5a5a655f..b03594864 100644 --- a/cmd/cgr-engine/rater.go +++ b/cmd/cgr-engine/rater.go @@ -45,7 +45,7 @@ func startRater(internalRaterChan chan *engine.Responder, internalBalancerChan c server *utils.Server, ratingDb engine.RatingStorage, accountDb engine.AccountingStorage, loadDb engine.LoadStorage, cdrDb engine.CdrStorage, logDb engine.LogStorage, stopHandled *bool, exitChan chan bool) { - waitTasks := make([]chan struct{}, 0) + var waitTasks []chan struct{} //Cache load cacheTaskChan := make(chan struct{}) diff --git a/console/scheduler_queue.go b/console/scheduler_queue.go new file mode 100644 index 000000000..a409c3a68 --- /dev/null +++ b/console/scheduler_queue.go @@ -0,0 +1,63 @@ +/* +Rating system designed to be used in VoIP Carriers World +Copyright (C) 2012-2015 ITsysCOM + +This program is free software: you can redistribute it and/or modify +it under the terms of the GNU General Public License as published by +the Free Software Foundation, either version 3 of the License, or +(at your option) any later version. + +This program is distributed in the hope that it will be useful, +but WITHOUT ANY WARRANTY; without even the implied warranty of +MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +GNU General Public License for more details. + +You should have received a copy of the GNU General Public License +along with this program. If not, see +*/ + +package console + +import "github.com/cgrates/cgrates/apier/v1" + +func init() { + c := &CmdGetScheduledActions{ + name: "scheduler_queue", + rpcMethod: "ApierV1.GetScheduledActions", + rpcParams: &v1.AttrsGetScheduledActions{}, + } + commands[c.Name()] = c + c.CommandExecuter = &CommandExecuter{c} +} + +// Commander implementation +type CmdGetScheduledActions struct { + name string + rpcMethod string + rpcParams *v1.AttrsGetScheduledActions + *CommandExecuter +} + +func (self *CmdGetScheduledActions) Name() string { + return self.name +} + +func (self *CmdGetScheduledActions) RpcMethod() string { + return self.rpcMethod +} + +func (self *CmdGetScheduledActions) RpcParams(reset bool) interface{} { + if reset || self.rpcParams == nil { + self.rpcParams = &v1.AttrsGetScheduledActions{} + } + return self.rpcParams +} + +func (self *CmdGetScheduledActions) PostprocessRpcParams() error { + return nil +} + +func (self *CmdGetScheduledActions) RpcResult() interface{} { + s := v1.ScheduledActions{} + return &s +} diff --git a/scheduler/scheduler.go b/scheduler/scheduler.go index db9728bb0..dfae55c04 100644 --- a/scheduler/scheduler.go +++ b/scheduler/scheduler.go @@ -66,7 +66,7 @@ func (s *Scheduler) Loop() { } else { s.Unlock() d := a0.GetNextStartTime(now).Sub(now) - //utils.Logger.Info(fmt.Sprintf("Timer set to wait for %v", d)) + utils.Logger.Info(fmt.Sprintf("Time to next action (%s): %v", a0.Id, d)) s.timer = time.NewTimer(d) select { case <-s.timer.C: @@ -84,27 +84,30 @@ func (s *Scheduler) LoadActionPlans(storage engine.RatingStorage) { if err != nil && err != utils.ErrNotFound { utils.Logger.Warning(fmt.Sprintf("Cannot get action plans: %v", err)) } + utils.Logger.Info(fmt.Sprintf(" processing %d action plans", len(actionPlans))) // recreate the queue s.Lock() s.queue = engine.ActionPlanPriotityList{} for key, aps := range actionPlans { toBeSaved := false isAsap := false - newApls := make([]*engine.ActionPlan, 0) // will remove the one time runs from the database + var newApls []*engine.ActionPlan // will remove the one time runs from the database for _, ap := range aps { if ap.Timing == nil { utils.Logger.Warning(fmt.Sprintf(" Nil timing on action plan: %+v, discarding!", ap)) continue } + if len(ap.AccountIds) == 0 { // no accounts just ignore + continue + } isAsap = ap.IsASAP() toBeSaved = toBeSaved || isAsap if isAsap { - if len(ap.AccountIds) > 0 { - utils.Logger.Info(fmt.Sprintf("Time for one time action on %v", key)) - } + utils.Logger.Info(fmt.Sprintf("Time for one time action on %v", key)) ap.Execute() ap.AccountIds = make([]string, 0) } else { + now := time.Now() if ap.GetNextStartTime(now).Before(now) { // the task is obsolete, do not add it to the queue @@ -124,6 +127,7 @@ func (s *Scheduler) LoadActionPlans(storage engine.RatingStorage) { } } sort.Sort(s.queue) + utils.Logger.Info(fmt.Sprintf(" queued %d action plans", len(s.queue))) s.Unlock() } From ea935183ed9f546271c685f9c21f8bbe3c21493b Mon Sep 17 00:00:00 2001 From: DanB Date: Wed, 2 Dec 2015 16:26:58 +0100 Subject: [PATCH 3/8] Scheduler to wait for cache to load before starting --- cmd/cgr-engine/cgr-engine.go | 10 +++++++--- cmd/cgr-engine/rater.go | 4 ++-- 2 files changed, 9 insertions(+), 5 deletions(-) diff --git a/cmd/cgr-engine/cgr-engine.go b/cmd/cgr-engine/cgr-engine.go index d519c7c9c..42eed606a 100644 --- a/cmd/cgr-engine/cgr-engine.go +++ b/cmd/cgr-engine/cgr-engine.go @@ -484,7 +484,10 @@ func startCDRS(internalCdrSChan chan *engine.CdrServer, logDb engine.LogStorage, internalCdrSChan <- cdrServer // Signal that cdrS is operational } -func startScheduler(internalSchedulerChan chan *scheduler.Scheduler, ratingDb engine.RatingStorage, exitChan chan bool) { +func startScheduler(internalSchedulerChan chan *scheduler.Scheduler, cacheDoneChan chan struct{}, ratingDb engine.RatingStorage, exitChan chan bool) { + // Wait for cache to load data before starting + cacheDone := <- cacheDoneChan + cacheDoneChan <- cacheDone utils.Logger.Info("Starting CGRateS Scheduler.") sched := scheduler.NewScheduler() go reloadSchedulerSingnalHandler(sched, ratingDb) @@ -666,6 +669,7 @@ func main() { // Define internal connections via channels internalBalancerChan := make(chan *balancer2go.Balancer, 1) internalRaterChan := make(chan *engine.Responder, 1) + cacheDoneChan := make(chan struct{}, 1) internalSchedulerChan := make(chan *scheduler.Scheduler, 1) internalCdrSChan := make(chan *engine.CdrServer, 1) internalCdrStatSChan := make(chan engine.StatsInterface, 1) @@ -681,13 +685,13 @@ func main() { // Start rater service if cfg.RaterEnabled { - go startRater(internalRaterChan, internalBalancerChan, internalSchedulerChan, internalCdrStatSChan, internalHistorySChan, internalPubSubSChan, internalUserSChan, internalAliaseSChan, + go startRater(internalRaterChan, cacheDoneChan, internalBalancerChan, internalSchedulerChan, internalCdrStatSChan, internalHistorySChan, internalPubSubSChan, internalUserSChan, internalAliaseSChan, server, ratingDb, accountDb, loadDb, cdrDb, logDb, &stopHandled, exitChan) } // Start Scheduler if cfg.SchedulerEnabled { - go startScheduler(internalSchedulerChan, ratingDb, exitChan) + go startScheduler(internalSchedulerChan, cacheDoneChan, ratingDb, exitChan) } // Start CDR Server diff --git a/cmd/cgr-engine/rater.go b/cmd/cgr-engine/rater.go index f5a5a655f..08a8e6d0d 100644 --- a/cmd/cgr-engine/rater.go +++ b/cmd/cgr-engine/rater.go @@ -39,7 +39,7 @@ func startBalancer(internalBalancerChan chan *balancer2go.Balancer, stopHandled } // Starts rater and reports on chan -func startRater(internalRaterChan chan *engine.Responder, internalBalancerChan chan *balancer2go.Balancer, internalSchedulerChan chan *scheduler.Scheduler, +func startRater(internalRaterChan chan *engine.Responder, cacheDoneChan chan struct{}, internalBalancerChan chan *balancer2go.Balancer, internalSchedulerChan chan *scheduler.Scheduler, internalCdrStatSChan chan engine.StatsInterface, internalHistorySChan chan history.Scribe, internalPubSubSChan chan engine.PublisherSubscriber, internalUserSChan chan engine.UserService, internalAliaseSChan chan engine.AliasService, server *utils.Server, @@ -62,7 +62,7 @@ func startRater(internalRaterChan chan *engine.Responder, internalBalancerChan c exitChan <- true return } - + cacheDoneChan <- struct{}{} }() // Retrieve scheduler for it's API methods From a2081516a28b3f3e468b85e8eb00b701aaaa1736 Mon Sep 17 00:00:00 2001 From: Radu Ioan Fericean Date: Wed, 2 Dec 2015 18:26:12 +0200 Subject: [PATCH 4/8] added mongo as supported tarrifplan db --- cmd/cgr-engine/cgr-engine.go | 2 +- engine/storage_utils.go | 3 +++ 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/cmd/cgr-engine/cgr-engine.go b/cmd/cgr-engine/cgr-engine.go index 42eed606a..048c09119 100644 --- a/cmd/cgr-engine/cgr-engine.go +++ b/cmd/cgr-engine/cgr-engine.go @@ -594,7 +594,7 @@ func main() { writePid() } if *singlecpu { - runtime.GOMAXPROCS(1) // Having multiple cpus slows down computing due to CPU management, to be reviewed in future Go releases + runtime.GOMAXPROCS(1) // Having multiple cpus may slow down computing due to CPU management, to be reviewed in future Go releases } if *cpuprofile != "" { f, err := os.Create(*cpuprofile) diff --git a/engine/storage_utils.go b/engine/storage_utils.go index da1292c3e..934e99378 100644 --- a/engine/storage_utils.go +++ b/engine/storage_utils.go @@ -41,6 +41,9 @@ func ConfigureRatingStorage(db_type, host, port, name, user, pass, marshaler str host += ":" + port } d, err = NewRedisStorage(host, db_nb, pass, marshaler, utils.REDIS_MAX_CONNS) + case utils.MONGO: + d, err = NewMongoStorage(host, port, name, user, pass) + db = d.(RatingStorage) default: err = errors.New("unknown db") } From 211cbd404332965d830a40f0f665ad2bb5504fa6 Mon Sep 17 00:00:00 2001 From: Radu Ioan Fericean Date: Wed, 2 Dec 2015 18:32:48 +0200 Subject: [PATCH 5/8] fix local test --- apier/v1/apier_local_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apier/v1/apier_local_test.go b/apier/v1/apier_local_test.go index b4fee06fe..25b753aa7 100644 --- a/apier/v1/apier_local_test.go +++ b/apier/v1/apier_local_test.go @@ -1624,7 +1624,7 @@ func TestApierLocalGetScheduledActions(t *testing.T) { } var rply []*ScheduledActions if err := rater.Call("ApierV1.GetScheduledActions", AttrsGetScheduledActions{}, &rply); err != nil { - t.Error("Unexpected error: ", err.Error) + t.Error("Unexpected error: ", err.Error()) } } From 0c8f2fd1bb998e102a25487e2b83de93f51e5f30 Mon Sep 17 00:00:00 2001 From: Radu Ioan Fericean Date: Wed, 2 Dec 2015 18:38:29 +0200 Subject: [PATCH 6/8] ignore test files --- .gitignore | 1 + 1 file changed, 1 insertion(+) diff --git a/.gitignore b/.gitignore index dfd8ef924..bd8808a87 100644 --- a/.gitignore +++ b/.gitignore @@ -14,3 +14,4 @@ data/vagrant/.vagrant data/vagrant/vagrant_ansible_inventory_default data/tutorials/fs_evsock/freeswitch/etc/freeswitch/ vendor +*.test From a01cad90394c589aae8d55641501a3f5954d7cbd Mon Sep 17 00:00:00 2001 From: Radu Ioan Fericean Date: Wed, 2 Dec 2015 18:44:39 +0200 Subject: [PATCH 7/8] scheduler logging --- scheduler/scheduler.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/scheduler/scheduler.go b/scheduler/scheduler.go index dfae55c04..0493fb6dd 100644 --- a/scheduler/scheduler.go +++ b/scheduler/scheduler.go @@ -66,12 +66,12 @@ func (s *Scheduler) Loop() { } else { s.Unlock() d := a0.GetNextStartTime(now).Sub(now) - utils.Logger.Info(fmt.Sprintf("Time to next action (%s): %v", a0.Id, d)) + utils.Logger.Info(fmt.Sprintf(" Time to next action (%s): %v", a0.Id, d)) s.timer = time.NewTimer(d) select { case <-s.timer.C: // timer has expired - utils.Logger.Info(fmt.Sprintf("Time for action on %v", a0)) + utils.Logger.Info(fmt.Sprintf(" Time for action on %v", a0)) case <-s.restartLoop: // nothing to do, just continue the loop } @@ -82,7 +82,7 @@ func (s *Scheduler) Loop() { func (s *Scheduler) LoadActionPlans(storage engine.RatingStorage) { actionPlans, err := storage.GetAllActionPlans() if err != nil && err != utils.ErrNotFound { - utils.Logger.Warning(fmt.Sprintf("Cannot get action plans: %v", err)) + utils.Logger.Warning(fmt.Sprintf(" Cannot get action plans: %v", err)) } utils.Logger.Info(fmt.Sprintf(" processing %d action plans", len(actionPlans))) // recreate the queue @@ -103,7 +103,7 @@ func (s *Scheduler) LoadActionPlans(storage engine.RatingStorage) { isAsap = ap.IsASAP() toBeSaved = toBeSaved || isAsap if isAsap { - utils.Logger.Info(fmt.Sprintf("Time for one time action on %v", key)) + utils.Logger.Info(fmt.Sprintf(" Time for one time action on %v", key)) ap.Execute() ap.AccountIds = make([]string, 0) } else { From b45581968a1c0c19d1756c2fc32d4ac95fa1b3b8 Mon Sep 17 00:00:00 2001 From: Radu Ioan Fericean Date: Wed, 2 Dec 2015 19:25:51 +0200 Subject: [PATCH 8/8] initialize action plan api array --- apier/v1/apier_local_test.go | 2 +- apier/v1/scheduler.go | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/apier/v1/apier_local_test.go b/apier/v1/apier_local_test.go index 25b753aa7..10ed9ec9d 100644 --- a/apier/v1/apier_local_test.go +++ b/apier/v1/apier_local_test.go @@ -1624,7 +1624,7 @@ func TestApierLocalGetScheduledActions(t *testing.T) { } var rply []*ScheduledActions if err := rater.Call("ApierV1.GetScheduledActions", AttrsGetScheduledActions{}, &rply); err != nil { - t.Error("Unexpected error: ", err.Error()) + t.Error("Unexpected error: ", err) } } diff --git a/apier/v1/scheduler.go b/apier/v1/scheduler.go index 5fbcb7591..69871f61a 100644 --- a/apier/v1/scheduler.go +++ b/apier/v1/scheduler.go @@ -113,7 +113,7 @@ func (self *ApierV1) GetScheduledActions(attrs AttrsGetScheduledActions, reply * if self.Sched == nil { return errors.New("SCHEDULER_NOT_ENABLED") } - var schedActions []*ScheduledActions + schedActions := make([]*ScheduledActions, 0) // needs to be initialized if remains empty scheduledActions := self.Sched.GetQueue() for _, qActions := range scheduledActions { sas := &ScheduledActions{ActionsId: qActions.ActionsId, ActionPlanId: qActions.Id, ActionPlanUuid: qActions.Uuid}