Adds watcher at service level and sends reminders via websocket
This commit is contained in:
+1
-1
@@ -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,
|
||||
|
||||
@@ -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
|
||||
})
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user