From f13953790b3bfb77cd2a636182b178df1ad2c12a Mon Sep 17 00:00:00 2001 From: Scott Knight <4534275+knightsc@users.noreply.github.com> Date: Thu, 21 Jun 2018 11:34:19 -0400 Subject: [PATCH] Create pushinfo Worker separate from DB implementation (#445) --- cmd/micromdm/serve.go | 3 ++ platform/apns/builtin/db.go | 42 --------------------- platform/apns/worker.go | 74 +++++++++++++++++++++++++++++++++++++ 3 files changed, 77 insertions(+), 42 deletions(-) create mode 100644 platform/apns/worker.go diff --git a/cmd/micromdm/serve.go b/cmd/micromdm/serve.go index 3cd83219..68ac7824 100644 --- a/cmd/micromdm/serve.go +++ b/cmd/micromdm/serve.go @@ -599,6 +599,9 @@ after: c.pushService = apns.LoggingMiddleware( log.With(level.Info(logger), "component", "apns"), )(service) + + pushinfoWorker := apns.NewWorker(db, c.pubclient, logger) + go pushinfoWorker.Run(context.Background()) } func (c *server) setupEnrollmentService() { diff --git a/platform/apns/builtin/db.go b/platform/apns/builtin/db.go index e5584b20..0f355f27 100644 --- a/platform/apns/builtin/db.go +++ b/platform/apns/builtin/db.go @@ -1,13 +1,11 @@ package builtin import ( - "context" "fmt" "github.com/boltdb/bolt" "github.com/pkg/errors" - "github.com/micromdm/micromdm/mdm" "github.com/micromdm/micromdm/platform/apns" "github.com/micromdm/micromdm/platform/pubsub" ) @@ -29,9 +27,6 @@ func NewDB(db *bolt.DB, sub pubsub.Subscriber) (*DB, error) { datastore := &DB{ DB: db, } - if err := datastore.pollCheckin(sub); err != nil { - return nil, err - } return datastore, nil } @@ -79,40 +74,3 @@ func (db *DB) Save(info *apns.PushInfo) error { } return tx.Commit() } - -func (db *DB) pollCheckin(sub pubsub.Subscriber) error { - tokenUpdateEvents, err := sub.Subscribe(context.TODO(), "push-info", mdm.TokenUpdateTopic) - if err != nil { - return errors.Wrapf(err, - "subscribing push to %s topic", mdm.TokenUpdateTopic) - } - go func() { - for { - select { - case event := <-tokenUpdateEvents: - var ev mdm.CheckinEvent - if err := mdm.UnmarshalCheckinEvent(event.Message, &ev); err != nil { - fmt.Println(err) - continue - } - info := apns.PushInfo{ - UDID: ev.Command.UDID, - Token: ev.Command.Token.String(), - PushMagic: ev.Command.PushMagic, - MDMTopic: ev.Command.Topic, - } - if ev.Command.UserID != "" { - // use the GUID if this is a user TokenUpdate. - info.UDID = ev.Command.UserID - } - if err := db.Save(&info); err != nil { - fmt.Println(err) - continue - } - fmt.Printf("updated pushinfo for udid %s\n", info.UDID) - } - } - }() - - return nil -} diff --git a/platform/apns/worker.go b/platform/apns/worker.go new file mode 100644 index 00000000..24bc4c3e --- /dev/null +++ b/platform/apns/worker.go @@ -0,0 +1,74 @@ +package apns + +import ( + "context" + + "github.com/go-kit/kit/log" + "github.com/go-kit/kit/log/level" + "github.com/micromdm/micromdm/mdm" + "github.com/micromdm/micromdm/platform/pubsub" + "github.com/pkg/errors" +) + +type WorkerStore interface { + Save(*PushInfo) error +} + +type Worker struct { + db WorkerStore + sub pubsub.Subscriber + logger log.Logger +} + +func NewWorker(db WorkerStore, subscriber pubsub.Subscriber, logger log.Logger) *Worker { + return &Worker{ + db: db, + sub: subscriber, + logger: logger, + } +} + +func (w *Worker) Run(ctx context.Context) error { + const subscription = "pushinfo_worker" + tokenUpdateEvents, err := w.sub.Subscribe(ctx, subscription, mdm.TokenUpdateTopic) + if err != nil { + return errors.Wrapf(err, + "subscribing %s to %s topic", subscription, mdm.TokenUpdateTopic) + } + + for { + var err error + select { + case <-ctx.Done(): + return ctx.Err() + case event := <-tokenUpdateEvents: + err = w.updatePushInfoFromTokenUpdate(ctx, event.Message) + } + if err != nil { + level.Info(w.logger).Log( + "msg", "update pushinfo from event", + "err", err, + ) + continue + } + } +} + +func (w *Worker) updatePushInfoFromTokenUpdate(ctx context.Context, message []byte) error { + var ev mdm.CheckinEvent + if err := mdm.UnmarshalCheckinEvent(message, &ev); err != nil { + return errors.Wrap(err, "unmarshal pushinfo event") + } + info := PushInfo{ + UDID: ev.Command.UDID, + Token: ev.Command.Token.String(), + PushMagic: ev.Command.PushMagic, + MDMTopic: ev.Command.Topic, + } + if ev.Command.UserID != "" { + // use the GUID if this is a user TokenUpdate. + info.UDID = ev.Command.UserID + } + err := w.db.Save(&info) + return errors.Wrapf(err, "saving pushinfo for udid=%s", info.UDID) +}