Add initial sink endpoint & service impl.
This commit is contained in:
@@ -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))
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user