diff --git a/app/boot_levels.go b/app/boot_levels.go index b3f92b432..435662355 100644 --- a/app/boot_levels.go +++ b/app/boot_levels.go @@ -281,7 +281,7 @@ func (app *CortezaApp) InitServices(ctx context.Context) (err error) { // // Note: this is a legacy approach, all services from all 3 apps // will most likely be merged in the future - err = sysService.Initialize(ctx, app.Log, app.Store, sysService.Config{ + err = sysService.Initialize(ctx, app.Log, app.Store, app.WsServer, sysService.Config{ ActionLog: app.Opt.ActionLog, Storage: app.Opt.ObjStore, Template: app.Opt.Template, diff --git a/system/service/reminder.go b/system/service/reminder.go index 8e12be3ed..432807335 100644 --- a/system/service/reminder.go +++ b/system/service/reminder.go @@ -6,15 +6,24 @@ import ( intAuth "github.com/cortezaproject/corteza-server/pkg/auth" "github.com/cortezaproject/corteza-server/store" "github.com/cortezaproject/corteza-server/system/types" + "github.com/cortezaproject/corteza-server/websocket" + "github.com/getsentry/sentry-go" + "go.uber.org/zap" "time" ) type ( + reminderSender interface { + Send(kind string, payload interface{}, userIDs ...uint64) error + } + reminder struct { ac reminderAccessController - actionlog actionlog.Recorder - store store.Reminders + log *zap.Logger + actionlog actionlog.Recorder + store store.Reminders + reminderSender reminderSender } reminderAccessController interface { @@ -33,13 +42,17 @@ type ( Snooze(context.Context, uint64, *time.Time) error Delete(context.Context, uint64) error + + Watch(ctx context.Context) } ) -func Reminder(ctx context.Context) ReminderService { +func Reminder(ctx context.Context, log *zap.Logger, rs reminderSender) ReminderService { return &reminder{ - ac: DefaultAccessControl, - store: DefaultStore, + ac: DefaultAccessControl, + log: log, + store: DefaultStore, + reminderSender: rs, } } @@ -263,3 +276,47 @@ func (svc reminder) Delete(ctx context.Context, ID uint64) (err error) { return svc.recordAction(ctx, raProps, ReminderActionDelete, err) } + +func (svc reminder) Watch(ctx context.Context) { + if svc.reminderSender != nil { + rTicker := time.NewTicker(time.Second) + + go func() { + defer sentry.Recover() + defer rTicker.Stop() + defer svc.log.Info("stopped") + + for { + select { + case <-ctx.Done(): + return + case t := <-rTicker.C: + // Get scheduled reminders of users + rr, _, err := svc.Find(ctx, types.ReminderFilter{ + ExcludeDismissed: true, + ScheduledOnly: true, + }) + if err != nil { + svc.log.Error("failed to get reminders of users", zap.Error(err)) + } + + // sendReminderNow checks time is equal to current time or not + sendReminderNow := func(tt time.Time) bool { + timeLayout := time.RFC3339 + return tt.Format(timeLayout) == t.Format(timeLayout) + } + + // Send scheduled reminders to users + _ = rr.Walk(func(r *types.Reminder) error { + if r.RemindAt != nil && sendReminderNow(*r.RemindAt) { + if err := svc.reminderSender.Send(websocket.StatusOK, r, r.AssignedTo); err != nil { + svc.log.Error("failed to send reminder to user", zap.Error(err)) + } + } + return nil + }) + } + } + }() + } +} diff --git a/system/service/service.go b/system/service/service.go index cf8ed9be5..c5a7b06d9 100644 --- a/system/service/service.go +++ b/system/service/service.go @@ -24,6 +24,10 @@ import ( ) type ( + websocketSender interface { + Send(kind string, payload interface{}, userIDs ...uint64) error + } + RBACServicer interface { accessControlRBACServicer Watch(ctx context.Context) @@ -91,7 +95,7 @@ var ( } ) -func Initialize(ctx context.Context, log *zap.Logger, s store.Storer, c Config) (err error) { +func Initialize(ctx context.Context, log *zap.Logger, s store.Storer, ws websocketSender, c Config) (err error) { var ( hcd = healthcheck.Defaults() ) @@ -164,7 +168,7 @@ func Initialize(ctx context.Context, log *zap.Logger, s store.Storer, c Config) DefaultUser = User(ctx) DefaultRole = Role(ctx) DefaultApplication = Application(DefaultStore, DefaultAccessControl, DefaultActionlog, eventbus.Service()) - DefaultReminder = Reminder(ctx) + DefaultReminder = Reminder(ctx, DefaultLogger.Named("reminder"), ws) DefaultSink = Sink() DefaultStatistics = Statistics() DefaultAttachment = Attachment(DefaultObjectStore) @@ -209,7 +213,8 @@ func Activate(ctx context.Context) (err error) { } func Watchers(ctx context.Context) { - // + DefaultReminder.Watch(ctx) + return } // isGeneric returns true if given error is generic