From cc13c0e27dd8967c022c13240392bc3771279364 Mon Sep 17 00:00:00 2001 From: Denis Arh Date: Tue, 3 Sep 2019 14:59:12 +0200 Subject: [PATCH] Add initial sink endpoint & service impl. --- system/internal/service/automation_runner.go | 2 +- system/internal/service/mailproc.go | 7 +-- system/internal/service/service.go | 4 ++ system/internal/service/sink.go | 53 ++++++++++++++++++ system/proto/mail_message.go | 43 +++++++++++++-- system/rest/router.go | 7 +++ system/rest/sink.go | 57 ++++++++++++++++++++ 7 files changed, 166 insertions(+), 7 deletions(-) create mode 100644 system/internal/service/sink.go create mode 100644 system/rest/sink.go diff --git a/system/internal/service/automation_runner.go b/system/internal/service/automation_runner.go index e9d746c07..eba7c09bb 100644 --- a/system/internal/service/automation_runner.go +++ b/system/internal/service/automation_runner.go @@ -115,7 +115,7 @@ func (svc automationRunner) RecordScriptTester(ctx context.Context, source strin func (svc automationRunner) makeMailScriptRunner(ctx context.Context, mail *types.MailMessage) func(script *automation.Script) error { // Static request params (record gets updated var req = &corredor.RunMailMessageRequest{ - MailMessage: proto.FromMailMessage(mail), + MailMessage: proto.NewMailMessage(mail), } svc.logger.Debug("executing script", zap.Any("mail", mail)) diff --git a/system/internal/service/mailproc.go b/system/internal/service/mailproc.go index 672e77804..418133635 100644 --- a/system/internal/service/mailproc.go +++ b/system/internal/service/mailproc.go @@ -19,12 +19,13 @@ type ( } mailprocScriptsRunner interface { - OnMailMessage(ctx context.Context, message *types.MailMessage) (err error) + OnReceiveMailMessage(ctx context.Context, message *types.MailMessage) (err error) } ) func Mailproc() *mailproc { return &mailproc{ + sr: DefaultAutomationRunner, logger: DefaultLogger.Named("mailproc"), } } @@ -34,10 +35,10 @@ func Mailproc() *mailproc { // return logger.AddRequestID(svc.ctx, svc.logger).With(fields...) // } -func (svc mailproc) RawMessage(ctx context.Context, m io.Reader) error { +func (svc mailproc) ContentProcessor(ctx context.Context, m io.Reader) error { if m, err := mailProcMessage(m); err != nil { return err - } else if err = svc.sr.OnMailMessage(ctx, m); err != nil { + } else if err = svc.sr.OnReceiveMailMessage(ctx, m); err != nil { return err } diff --git a/system/internal/service/service.go b/system/internal/service/service.go index 2ab83272a..27d36bcdb 100644 --- a/system/internal/service/service.go +++ b/system/internal/service/service.go @@ -55,6 +55,8 @@ var ( DefaultAuthNotification AuthNotificationService DefaultAuthSettings AuthSettings + DefaultSink *sink + DefaultAuth AuthService DefaultUser UserService DefaultRole RoleService @@ -135,6 +137,8 @@ func Init(ctx context.Context, log *zap.Logger, c Config) (err error) { ) } + DefaultSink = Sink() + return } diff --git a/system/internal/service/sink.go b/system/internal/service/sink.go new file mode 100644 index 000000000..e98b59db1 --- /dev/null +++ b/system/internal/service/sink.go @@ -0,0 +1,53 @@ +package service + +import ( + "context" + "io" + "strings" + + "go.uber.org/zap" +) + +type ( + sink struct { + // processors + proc map[string]sinkContentProc + + logger *zap.Logger + } + + sinkContentProc interface { + ContentProcessor(context.Context, io.Reader) error + } +) + +const ( + ErrSinkContentTypeUnsupported serviceError = "SinkUnsupportedContentType" + ErrSinkContentProcessingFailed serviceError = "SinkProcessFailed" + + SinkContentTypeMail = "message/rfc822" +) + +func Sink() *sink { + return &sink{ + proc: map[string]sinkContentProc{ + SinkContentTypeMail: Mailproc(), + }, + logger: DefaultLogger, + } +} + +// Finds appropriate sink processor +func (svc *sink) Process(ctx context.Context, contentType string, r io.Reader) (err error) { + switch strings.ToLower(contentType) { + case SinkContentTypeMail, "rfc822", "email", "mail": + if err = svc.proc[SinkContentTypeMail].ContentProcessor(ctx, r); err != nil { + return ErrSinkContentProcessingFailed + } + + default: + return ErrSinkContentTypeUnsupported + } + + return +} diff --git a/system/proto/mail_message.go b/system/proto/mail_message.go index 909d3a2a5..72462344c 100644 --- a/system/proto/mail_message.go +++ b/system/proto/mail_message.go @@ -1,16 +1,53 @@ package proto import ( + mail2 "net/mail" + + "github.com/golang/protobuf/ptypes/timestamp" + "github.com/cortezaproject/corteza-server/system/types" ) -func FromMailMessage(mail *types.MailMessage) *MailMessage { +func NewMailMessage(mail *types.MailMessage) *MailMessage { if mail == nil { return nil } - var p = &MailMessage{} + addrConv := func(aa []*mail2.Address) []*MailMessage_Header_MailAddress { + out := make([]*MailMessage_Header_MailAddress, len(aa)) + + for i, a := range aa { + out[i] = &MailMessage_Header_MailAddress{ + Address: a.Address, + Name: a.Name, + } + } + + return out + } + + hConv := func(hh mail2.Header) map[string]*MailMessage_Header_HeaderValues { + out := make(map[string]*MailMessage_Header_HeaderValues, len(hh)) + + for k, vv := range hh { + out[k] = &MailMessage_Header_HeaderValues{Values: vv} + } + + return out + } + + var p = &MailMessage{ + Header: &MailMessage_Header{ + Date: ×tamp.Timestamp{Seconds: mail.Date.Unix()}, + To: addrConv(mail.Header.To), + Cc: addrConv(mail.Header.CC), + Bcc: addrConv(mail.Header.BCC), + From: addrConv(mail.Header.From), + ReplyTo: addrConv(mail.Header.ReplyTo), + Raw: hConv(mail.Header.Raw), + }, + RawBody: mail.RawBody, + } - panic("@todo implement mailmessage => proto conv!") return p } diff --git a/system/rest/router.go b/system/rest/router.go index d0bb6051b..ebd28f40b 100644 --- a/system/rest/router.go +++ b/system/rest/router.go @@ -4,6 +4,7 @@ import ( "github.com/go-chi/chi" "github.com/cortezaproject/corteza-server/internal/auth" + "github.com/cortezaproject/corteza-server/system/internal/service" "github.com/cortezaproject/corteza-server/system/rest/handlers" ) @@ -19,6 +20,12 @@ func MountRoutes(r chi.Router) { r.Group(func(r chi.Router) { r.Use(auth.MiddlewareValidOnly) + // A special case that, we do not add this through standard request, handlers & controllers + // combo but directly -- we need access to r.Body + r.Handle("/sink", &Sink{ + svc: service.DefaultSink, + }) + handlers.NewUser(User{}.New()).MountRoutes(r) handlers.NewRole(Role{}.New()).MountRoutes(r) handlers.NewOrganisation(Organisation{}.New()).MountRoutes(r) diff --git a/system/rest/sink.go b/system/rest/sink.go new file mode 100644 index 000000000..cbbeb04ac --- /dev/null +++ b/system/rest/sink.go @@ -0,0 +1,57 @@ +package rest + +import ( + "context" + "io" + "net/http" + "strings" + + "github.com/pkg/errors" + + "github.com/cortezaproject/corteza-server/system/internal/service" +) + +var _ = errors.Wrap + +type Sink struct { + // xxx service.XXXService + svc interface { + Process(context.Context, string, io.Reader) error + } +} + +func (ctrl *Sink) ServeHTTP(w http.ResponseWriter, r *http.Request) { + var ( + ctx = r.Context() + cType = r.URL.Query().Get("content-type") + + unsupported = func() { + http.Error(w, "unsupported content-type", http.StatusBadRequest) + } + ) + + if cType == "" { + // If content-type not explicitly set (via QS), + // try to get it from the headers + cType = r.Header.Get("content-type") + if i := strings.Index(cType, ";"); i > 0 { + // intentionally > 0 + cType = cType[0 : i-1] + } + } + + if cType == "" { + unsupported() + return + } + + defer r.Body.Close() + + switch ctrl.svc.Process(ctx, cType, r.Body) { + case service.ErrSinkContentProcessingFailed: + http.Error(w, "sink processing failed", http.StatusInternalServerError) + + case service.ErrSinkContentTypeUnsupported: + unsupported() + } +}