feature: api for sending emails with attachments (#15201)

Co-authored-by: Qiu Jian <qiujian@yunionyun.com>
This commit is contained in:
Jian Qiu
2022-10-17 14:30:03 +08:00
committed by GitHub
co-authored by Qiu Jian
parent bc8b181d9f
commit df373cdaf6
43 changed files with 3124 additions and 151 deletions
+107
View File
@@ -0,0 +1,107 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package notify
import (
"time"
"yunion.io/x/jsonutils"
"yunion.io/x/onecloud/pkg/apis"
)
type SEmailMessage struct {
To []string `json:"to"`
Subject string `json:"subject"`
Body string `json:"body"`
Attachments []SEmailAttachment `json:"attachments"`
}
type SEmailAttachment struct {
Filename string `json:"filename"`
Mime string `json:"mime"`
Base64Content string `json:"content"`
}
type SEmailConfig struct {
Hostname string `json:"hostname"`
Hostport int `json:"hostport"`
Username string `json:"username"`
Password string `json:"password"`
SenderAddress string `json:"sender_address"`
SslGlobal bool `json:"ssl_global"`
}
type EmailQueueCreateInput struct {
SEmailMessage
// swagger: ignore
Dest string `json:"dest"`
// swagger: ignore
Content jsonutils.JSONObject `json:"content"`
// swagger: ignore
ProjectId string `json:"project_id"`
// swagger: ignore
Project string `json:"project"`
// swagger: ignore
ProjectDomainId string `json:"project_domain_id"`
// swagger: ignore
ProjectDomain string `json:"project_domain"`
// swagger: ignore
UserId string `json:"user_id"`
// swagger: ignore
User string `json:"user"`
// swagger: ignore
DomainId string `json:"domain_id"`
// swagger: ignore
Domain string `json:"domain"`
// swagger: ignore
Roles string `json:"roles"`
SessionId string `json:"session_id"`
}
type EmailQueueListInput struct {
apis.ModelBaseListInput
Id []int `json:"id"`
To []string `json:"to"`
Subject string `json:"subject"`
SessionId []string `json:"session_id"`
}
const (
EmailQueued = "queued"
EmailSending = "sending"
EmailSuccess = "success"
EmailFail = "fail"
)
type EmailQueueSendInput struct {
Sync bool `json:"sync"`
}
type EmailQueueDetails struct {
apis.ModelBaseDetails
SentAt time.Time `json:"sent_at"`
Status string `json:"status"`
Results string `json:"results"`
}
+147
View File
@@ -0,0 +1,147 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package db
import (
"context"
"fmt"
"strconv"
"time"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/sqlchemy"
"yunion.io/x/sqlchemy/backends/clickhouse"
"yunion.io/x/onecloud/pkg/cloudcommon/consts"
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
)
type SLogBaseManager struct {
SModelBaseManager
}
type SLogBase struct {
SModelBase
Id int64 `primary:"true" auto_increment:"true" list:"user" clickhouse_partition_by:"toInt64(divide(id,100000000000))"`
}
func NewLogBaseManager(model interface{}, table string, keyword, keywordPlural string, timeCol string, useClickHouse bool) SLogBaseManager {
if useClickHouse {
man := SLogBaseManager{NewModelBaseManagerWithDBName(
model,
table,
keyword,
keywordPlural,
ClickhouseDB,
)}
col := man.TableSpec().ColumnSpec("timeCol")
if clickCol, ok := col.(clickhouse.IClickhouseColumnSpec); ok {
clickCol.SetTTL(consts.SplitableMaxKeepMonths(), "MONTH")
}
return man
} else {
return SLogBaseManager{NewModelBaseManagerWithSplitable(
model,
table,
keyword,
keywordPlural,
"id",
timeCol,
consts.SplitableMaxDuration(),
consts.SplitableMaxKeepMonths(),
)}
}
}
func (manager *SLogBaseManager) CreateByInsertOrUpdate() bool {
return false
}
func CurrentTimestamp(t time.Time) int64 {
ret := int64(0)
const (
yOffset = 10000000000000
mOffset = 100000000000
dOffset = 1000000000
hOffset = 10000000
iOffset = 100000
sOffset = 1000
)
ret += int64(t.Year()) * yOffset
ret += int64(t.Month()) * mOffset
ret += int64(t.Day()) * dOffset
ret += int64(t.Hour()) * hOffset
ret += int64(t.Minute()) * iOffset
ret += int64(t.Second()) * sOffset
ret += int64(t.Nanosecond()) / 1000000
return ret
}
func (opslog *SLogBase) BeforeInsert() {
t := time.Now().UTC()
opslog.Id = CurrentTimestamp(t)
}
func (opslog *SLogBase) GetId() string {
return fmt.Sprintf("%d", opslog.Id)
}
func (self *SLogBase) ValidateDeleteCondition(ctx context.Context, info jsonutils.JSONObject) error {
return httperrors.NewForbiddenError("not allow to delete log")
}
func (self *SLogBaseManager) FilterById(q *sqlchemy.SQuery, idStr string) *sqlchemy.SQuery {
id, _ := strconv.Atoi(idStr)
return q.Equals("id", id)
}
func (self *SLogBaseManager) FilterByNotId(q *sqlchemy.SQuery, idStr string) *sqlchemy.SQuery {
id, _ := strconv.Atoi(idStr)
return q.NotEquals("id", id)
}
func (self *SLogBaseManager) FilterByName(q *sqlchemy.SQuery, name string) *sqlchemy.SQuery {
return q
}
func (manager *SLogBaseManager) GetPagingConfig() *SPagingConfig {
return &SPagingConfig{
Order: sqlchemy.SQL_ORDER_DESC,
MarkerFields: []string{"id"},
DefaultLimit: 20,
}
}
func (lb *SLogBase) GetRecordTime() time.Time {
log.Fatalf("not implemented yet!")
return time.Time{}
}
func (manager *SLogBaseManager) FetchById(idStr string) (IModel, error) {
return FetchById(manager.GetIModelManager(), idStr)
}
func (l *SLogBase) ValidateUpdateData(
ctx context.Context,
userCred mcclient.TokenCredential,
query jsonutils.JSONObject,
data jsonutils.JSONObject,
) (jsonutils.JSONObject, error) {
return nil, errors.Wrap(httperrors.ErrForbidden, "not allow")
}
+11 -88
View File
@@ -19,7 +19,6 @@ import (
"database/sql"
"fmt"
"runtime/debug"
"strconv"
"strings"
"time"
@@ -30,7 +29,6 @@ import (
"yunion.io/x/pkg/util/stringutils"
"yunion.io/x/pkg/util/timeutils"
"yunion.io/x/sqlchemy"
"yunion.io/x/sqlchemy/backends/clickhouse"
"yunion.io/x/onecloud/pkg/apis"
"yunion.io/x/onecloud/pkg/appsrv"
@@ -42,13 +40,12 @@ import (
)
type SOpsLogManager struct {
SModelBaseManager
SLogBaseManager
}
type SOpsLog struct {
SModelBase
SLogBase
Id int64 `primary:"true" auto_increment:"true" list:"user" clickhouse_partition_by:"toInt64(divide(id,100000000000))"`
ObjType string `width:"40" charset:"ascii" nullable:"false" list:"user" create:"required"`
ObjId string `width:"128" charset:"ascii" nullable:"false" list:"user" create:"required" index:"true"`
ObjName string `width:"128" charset:"utf8" nullable:"false" list:"user" create:"required"`
@@ -81,41 +78,21 @@ var _ IModel = (*SOpsLog)(nil)
var opslogQueryWorkerMan *appsrv.SWorkerManager
var opslogWriteWorkerMan *appsrv.SWorkerManager
func InitOpsLog() {
if consts.OpsLogWithClickhouse {
OpsLog = &SOpsLogManager{NewModelBaseManagerWithDBName(
SOpsLog{},
"opslog_tbl",
"event",
"events",
ClickhouseDB,
)}
col := OpsLog.TableSpec().ColumnSpec("ops_time")
if clickCol, ok := col.(clickhouse.IClickhouseColumnSpec); ok {
clickCol.SetTTL(consts.SplitableMaxKeepMonths(), "MONTH")
}
} else {
OpsLog = &SOpsLogManager{NewModelBaseManagerWithSplitable(
SOpsLog{},
"opslog_tbl",
"event",
"events",
"id",
"ops_time",
consts.SplitableMaxDuration(),
consts.SplitableMaxKeepMonths(),
)}
func NewOpsLogManager(opslog interface{}, tblName string, keyword, keywordPlural string, timeField string, clickhouse bool) SOpsLogManager {
return SOpsLogManager{
SLogBaseManager: NewLogBaseManager(opslog, tblName, keyword, keywordPlural, timeField, clickhouse),
}
}
func InitOpsLog() {
tmp := NewOpsLogManager(SOpsLog{}, "opslog_tbl", "event", "events", "ops_time", consts.OpsLogWithClickhouse)
OpsLog = &tmp
OpsLog.SetVirtualObject(OpsLog)
opslogQueryWorkerMan = appsrv.NewWorkerManager("opslog_query_worker", 2, 512, true)
opslogWriteWorkerMan = appsrv.NewWorkerManager("opslog_write_worker", 1, 2048, true)
}
func (manager *SOpsLogManager) CreateByInsertOrUpdate() bool {
return false
}
func (manager *SOpsLogManager) CustomizeHandlerInfo(info *appsrv.SHandlerInfo) {
manager.SModelBaseManager.CustomizeHandlerInfo(info)
@@ -125,35 +102,6 @@ func (manager *SOpsLogManager) CustomizeHandlerInfo(info *appsrv.SHandlerInfo) {
}
}
func CurrentTimestamp(t time.Time) int64 {
ret := int64(0)
const (
yOffset = 10000000000000
mOffset = 100000000000
dOffset = 1000000000
hOffset = 10000000
iOffset = 100000
sOffset = 1000
)
ret += int64(t.Year()) * yOffset
ret += int64(t.Month()) * mOffset
ret += int64(t.Day()) * dOffset
ret += int64(t.Hour()) * hOffset
ret += int64(t.Minute()) * iOffset
ret += int64(t.Second()) * sOffset
ret += int64(t.Nanosecond()) / 1000000
return ret
}
func (opslog *SOpsLog) BeforeInsert() {
t := time.Now().UTC()
opslog.Id = CurrentTimestamp(t)
}
func (opslog *SOpsLog) GetId() string {
return fmt.Sprintf("%d", opslog.Id)
}
func (opslog *SOpsLog) GetName() string {
return fmt.Sprintf("%s-%s", opslog.ObjType, opslog.Action)
}
@@ -410,24 +358,6 @@ func (manager *SOpsLogManager) LogSyncUpdate(m IModel, uds sqlchemy.UpdateDiffs,
}
}
func (self *SOpsLog) ValidateDeleteCondition(ctx context.Context, info jsonutils.JSONObject) error {
return httperrors.NewForbiddenError("not allow to delete log")
}
func (self *SOpsLogManager) FilterById(q *sqlchemy.SQuery, idStr string) *sqlchemy.SQuery {
id, _ := strconv.Atoi(idStr)
return q.Equals("id", id)
}
func (self *SOpsLogManager) FilterByNotId(q *sqlchemy.SQuery, idStr string) *sqlchemy.SQuery {
id, _ := strconv.Atoi(idStr)
return q.NotEquals("id", id)
}
func (self *SOpsLogManager) FilterByName(q *sqlchemy.SQuery, name string) *sqlchemy.SQuery {
return q
}
func (self *SOpsLogManager) FilterByOwner(q *sqlchemy.SQuery, ownerId mcclient.IIdentityProvider, scope rbacutils.TRbacScope) *sqlchemy.SQuery {
if ownerId != nil {
switch scope {
@@ -487,14 +417,6 @@ func (manager *SOpsLogManager) ResourceScope() rbacutils.TRbacScope {
return rbacutils.ScopeUser
}
func (manager *SOpsLogManager) GetPagingConfig() *SPagingConfig {
return &SPagingConfig{
Order: sqlchemy.SQL_ORDER_DESC,
MarkerFields: []string{"id"},
DefaultLimit: 20,
}
}
func (manager *SOpsLogManager) FetchOwnerId(ctx context.Context, data jsonutils.JSONObject) (mcclient.IIdentityProvider, error) {
ownerId := SOwnerId{}
err := data.Unmarshal(&ownerId)
@@ -533,6 +455,7 @@ func (log *SOpsLog) CustomizeCreate(ctx context.Context,
return log.SModelBase.CustomizeCreate(ctx, userCred, ownerId, query, data)
}
// override
func (log *SOpsLog) GetRecordTime() time.Time {
return log.OpsTime
}
+6 -38
View File
@@ -26,7 +26,6 @@ import (
"yunion.io/x/pkg/util/timeutils"
"yunion.io/x/pkg/utils"
"yunion.io/x/sqlchemy"
"yunion.io/x/sqlchemy/backends/clickhouse"
"yunion.io/x/onecloud/pkg/apis"
api "yunion.io/x/onecloud/pkg/apis/logger"
@@ -85,45 +84,14 @@ var logQueue = make(chan *SActionlog, 50)
func InitActionLog() {
InitActionWhiteList()
var initTable func(tbname string) *SActionlogManager
if consts.OpsLogWithClickhouse {
initTable = func(tbname string) *SActionlogManager {
tbl := &SActionlogManager{
SOpsLogManager: db.SOpsLogManager{
SModelBaseManager: db.NewModelBaseManagerWithDBName(
SActionlog{},
tbname,
"action",
"actions",
db.ClickhouseDB,
),
},
}
col := tbl.TableSpec().ColumnSpec("ops_time")
if clickCol, ok := col.(clickhouse.IClickhouseColumnSpec); ok {
clickCol.SetTTL(consts.SplitableMaxKeepMonths(), "MONTH")
}
return tbl
initTable := func(tbname string) *SActionlogManager {
tbl := &SActionlogManager{
SOpsLogManager: db.NewOpsLogManager(SActionlog{}, tbname, "action", "actions", "ops_time", consts.OpsLogWithClickhouse),
}
} else {
initTable = func(tbname string) *SActionlogManager {
tbl := &SActionlogManager{
SOpsLogManager: db.SOpsLogManager{
SModelBaseManager: db.NewModelBaseManagerWithSplitable(
SActionlog{},
tbname,
"action",
"actions",
"id",
"start_time",
consts.SplitableMaxDuration(),
consts.SplitableMaxKeepMonths(),
),
},
SRecordChecksumResourceBaseManager: *db.NewRecordChecksumResourceBaseManager(),
}
return tbl
if consts.OpsLogWithClickhouse {
tbl.SRecordChecksumResourceBaseManager = *db.NewRecordChecksumResourceBaseManager()
}
return tbl
}
ActionLog = initTable("action_tbl")
ActionLog.SetVirtualObject(ActionLog)
@@ -0,0 +1,32 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package notify
import (
"yunion.io/x/onecloud/pkg/mcclient/modulebase"
"yunion.io/x/onecloud/pkg/mcclient/modules"
)
var (
EmailQueues modulebase.ResourceManager
)
func init() {
EmailQueues = modules.NewNotifyv2Manager("emailqueue", "emailqueues",
[]string{"id", "subject", "dest", "status", "recv_at"},
[]string{},
)
modules.Register(&EmailQueues)
}
@@ -0,0 +1,96 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package notify
import (
"encoding/base64"
"io/ioutil"
"path/filepath"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
api "yunion.io/x/onecloud/pkg/apis/notify"
"yunion.io/x/onecloud/pkg/mcclient/options"
)
type EmailQueueListOptions struct {
options.BaseListOptions
Id []int `json:"id"`
To []string `json:"to"`
Subject string `json:"subject"`
SessionId []string `json:"session_id"`
}
func (rl *EmailQueueListOptions) Params() (jsonutils.JSONObject, error) {
return options.ListStructToParams(rl)
}
type EmailQueueCreateOptions struct {
SUBJECT string `help:"email subject"`
BODY string `help:"email body"`
TO []string `json:"to" help:"receiver email"`
SessionId string `help:"session id of sending email"`
Attach []string `help:"path to attachment"`
}
func (rc *EmailQueueCreateOptions) Params() (jsonutils.JSONObject, error) {
input := api.EmailQueueCreateInput{}
input.To = rc.TO
input.Subject = rc.SUBJECT
body, err := ioutil.ReadFile(rc.BODY)
if err != nil {
return nil, errors.Wrap(err, "Read content")
}
input.Body = string(body)
for _, attach := range rc.Attach {
contBytes, err := ioutil.ReadFile(attach)
if err != nil {
return nil, errors.Wrapf(err, "read %s", attach)
}
input.Attachments = append(input.Attachments, api.SEmailAttachment{
Filename: filepath.Base(attach),
Base64Content: base64.StdEncoding.EncodeToString(contBytes),
})
}
log.Debugf("%s", jsonutils.Marshal(input))
return jsonutils.Marshal(input), nil
}
type EmailQueueOptions struct {
ID string `help:"Id of email queue" json:"-"`
}
func (r *EmailQueueOptions) GetId() string {
return r.ID
}
func (r *EmailQueueOptions) Params() (jsonutils.JSONObject, error) {
return jsonutils.Marshal(r), nil
}
type EmailQueueSendOptions struct {
EmailQueueOptions
Sync bool `json:"sync" help:"send email synchronously"`
}
func (r *EmailQueueSendOptions) Params() (jsonutils.JSONObject, error) {
return jsonutils.Marshal(r), nil
}
+16
View File
@@ -595,6 +595,22 @@ func (self *SConfigManager) GetConfigs(contactType string) ([]notifyv2.SConfig,
return ret, nil
}
func (manager *SConfigManager) getEmailConfig() (*api.SEmailConfig, error) {
confs, err := manager.GetConfigs(api.EMAIL)
if err != nil {
return nil, errors.Wrap(err, "GetConfigs")
}
if len(confs) == 0 {
return nil, errors.Wrap(errors.ErrNotSupported, "email not supported")
}
conf := api.SEmailConfig{}
err = jsonutils.Marshal(confs[0].Config).Unmarshal(&conf)
if err != nil {
return nil, errors.Wrap(err, "Unmarshal")
}
return &conf, nil
}
func (self *SConfigManager) SetConfig(contactType string, config notifyv2.SConfig) error {
content := jsonutils.Marshal(config.Config)
sConfig := &SConfig{
+284
View File
@@ -0,0 +1,284 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package models
import (
"context"
"fmt"
"strings"
"time"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/util/regutils"
"yunion.io/x/sqlchemy"
api "yunion.io/x/onecloud/pkg/apis/notify"
"yunion.io/x/onecloud/pkg/cloudcommon/consts"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/notify/sender"
"yunion.io/x/onecloud/pkg/util/stringutils2"
)
type SEmailQueueManager struct {
db.SLogBaseManager
}
type SEmailQueue struct {
db.SLogBase
RecvAt time.Time `nullable:"false" created_at:"true" index:"true" get:"user" list:"user" json:"recv_at"`
Dest string `width:"256" charset:"ascii" nullable:"false" list:"user" create:"admin_required"`
Subject string `width:"256" charset:"utf8" nullable:"false" list:"user" create:"admin_required"`
SessionId string `width:"256" charset:"utf8" nullable:"false" list:"user" create:"admin_optional"`
Content jsonutils.JSONObject `length:"long" charset:"utf8" nullable:"false" list:"user" create:"admin_required"`
ProjectId string `width:"128" charset:"ascii" list:"user" create:"admin_optional" index:"true"`
Project string `width:"128" charset:"utf8" list:"user" create:"admin_optional"`
ProjectDomainId string `name:"project_domain_id" default:"default" width:"128" charset:"ascii" list:"user" create:"admin_optional"`
ProjectDomain string `name:"project_domain" default:"Default" width:"128" charset:"utf8" list:"user" create:"admin_optional"`
UserId string `width:"128" charset:"ascii" list:"user" create:"admin_required"`
User string `width:"128" charset:"utf8" list:"user" create:"admin_required"`
DomainId string `width:"128" charset:"ascii" list:"user" create:"admin_optional"`
Domain string `width:"128" charset:"utf8" list:"user" create:"admin_optional"`
Roles string `width:"64" charset:"utf8" list:"user" create:"admin_optional"`
}
var EmailQueueManager *SEmailQueueManager
func InitEmailQueue() {
EmailQueueManager = &SEmailQueueManager{
SLogBaseManager: db.NewLogBaseManager(SEmailQueue{}, "emailqueue_tbl", "emailqueue", "emailqueues", "recv_at", consts.OpsLogWithClickhouse),
}
EmailQueueManager.SetVirtualObject(EmailQueueManager)
}
func (e *SEmailQueue) GetRecordTime() time.Time {
return e.RecvAt
}
func (manager *SEmailQueueManager) ValidateCreateData(
ctx context.Context,
userCred mcclient.TokenCredential,
ownerId mcclient.IIdentityProvider,
query jsonutils.JSONObject,
input api.EmailQueueCreateInput,
) (api.EmailQueueCreateInput, error) {
// check permission
if db.IsAdminAllowCreate(userCred, manager).Result.IsDeny() {
return input, errors.Wrap(httperrors.ErrForbidden, "only admin can send email")
}
// validate data
if len(input.To) == 0 {
return input, errors.Wrap(httperrors.ErrInputParameter, "empty receiver")
}
invalidTos := make([]string, 0)
for _, to := range input.To {
if !regutils.MatchEmail(to) {
invalidTos = append(invalidTos, to)
}
}
if len(invalidTos) > 0 {
return input, errors.Wrapf(httperrors.ErrInputParameter, "invalid email %s", strings.Join(invalidTos, ","))
}
input.Dest = strings.Join(input.To, ",")
msg := api.SEmailMessage{
Body: input.Body,
Attachments: input.Attachments,
}
input.Content = jsonutils.Marshal(msg)
input.Project = userCred.GetProjectName()
input.ProjectId = userCred.GetProjectId()
input.ProjectDomain = userCred.GetProjectDomain()
input.ProjectDomainId = userCred.GetProjectDomainId()
input.User = userCred.GetUserName()
input.UserId = userCred.GetUserId()
input.Domain = userCred.GetDomainName()
input.DomainId = userCred.GetDomainId()
input.Roles = strings.Join(userCred.GetRoles(), ",")
return input, nil
}
func (eq *SEmailQueue) PostCreate(
ctx context.Context,
userCred mcclient.TokenCredential,
ownerId mcclient.IIdentityProvider,
query jsonutils.JSONObject,
data jsonutils.JSONObject,
) {
eq.SLogBase.PostCreate(ctx, userCred, ownerId, query, data)
eq.setStatus(ctx, api.EmailQueued, nil)
eq.doSendAsync()
}
func (eq *SEmailQueue) doSendAsync() {
sender.Worker.Run(eq, nil, nil)
}
func (eq *SEmailQueue) Dump() string {
return fmt.Sprintf("send email %s", eq.Subject)
}
func (eq *SEmailQueue) Run() {
log.Debugf("send email")
eq.doSend(context.TODO())
}
func (eq *SEmailQueue) doSend(ctx context.Context) {
conf, err := ConfigManager.getEmailConfig()
if err != nil {
eq.setStatus(ctx, api.EmailFail, err)
return
}
log.Debugf("conf: %s", jsonutils.Marshal(conf))
msg, err := eq.getMessage()
if err != nil {
eq.setStatus(ctx, api.EmailFail, err)
return
}
log.Debugf("msg: %s", jsonutils.Marshal(msg))
eq.setStatus(ctx, api.EmailSending, nil)
err = sender.SendEmail(conf, msg)
if err != nil {
eq.setStatus(ctx, api.EmailFail, err)
return
}
eq.setStatus(ctx, api.EmailSuccess, nil)
return
}
func (eq *SEmailQueue) getMessage() (*api.SEmailMessage, error) {
msg := api.SEmailMessage{}
err := eq.Content.Unmarshal(&msg)
if err != nil {
return nil, errors.Wrap(err, "Unmarshal")
}
msg.To = strings.Split(eq.Dest, ",")
msg.Subject = eq.Subject
return &msg, nil
}
func (eq *SEmailQueue) setStatus(ctx context.Context, status string, results error) {
eqs := SEmailQueueStatus{
Id: eq.Id,
Status: status,
}
if results != nil {
eqs.Results = results.Error()
}
if status == api.EmailSuccess || status == api.EmailFail {
eqs.SentAt = time.Now()
}
EmailQueueStatusManager.TableSpec().InsertOrUpdate(ctx, &eqs)
}
// 宿主机/物理机列表
func (manager *SEmailQueueManager) ListItemFilter(
ctx context.Context,
q *sqlchemy.SQuery,
userCred mcclient.TokenCredential,
query api.EmailQueueListInput,
) (*sqlchemy.SQuery, error) {
var err error
q, err = manager.SLogBaseManager.ListItemFilter(ctx, q, userCred, query.ModelBaseListInput)
if err != nil {
return q, errors.Wrap(err, "SLogBaseManager.ListItemFilter")
}
if len(query.Id) > 0 {
q = q.In("id", query.Id)
}
if len(query.To) > 0 {
cond := make([]sqlchemy.ICondition, 0)
for _, to := range query.To {
cond = append(cond, sqlchemy.Contains(q.Field("dest"), to))
}
q = q.Filter(sqlchemy.OR(cond...))
}
if len(query.Subject) > 0 {
q = q.Contains("subject", query.Subject)
}
if len(query.SessionId) > 0 {
q = q.In("session_id", query.SessionId)
}
return q, nil
}
func (eq *SEmailQueue) PerformSend(
ctx context.Context,
userCred mcclient.TokenCredential,
query jsonutils.JSONObject,
input api.EmailQueueSendInput,
) (jsonutils.JSONObject, error) {
eq.setStatus(ctx, api.EmailQueued, nil)
if input.Sync {
log.Debugf("send email synchronously")
eq.doSend(ctx)
} else {
log.Debugf("send email Asynchronously")
eq.doSendAsync()
}
return nil, nil
}
func (manager *SEmailQueueManager) FetchCustomizeColumns(
ctx context.Context,
userCred mcclient.TokenCredential,
query jsonutils.JSONObject,
objs []interface{},
fields stringutils2.SSortedStrings,
isList bool,
) []api.EmailQueueDetails {
rows := make([]api.EmailQueueDetails, len(objs))
baseRows := manager.SModelBaseManager.FetchCustomizeColumns(ctx, userCred, query, objs, fields, isList)
emailIds := make([]int64, len(objs))
for i := range rows {
rows[i] = api.EmailQueueDetails{
ModelBaseDetails: baseRows[i],
}
eq := objs[i].(*SEmailQueue)
emailIds[i] = eq.Id
}
rets, err := EmailQueueStatusManager.fetchEmailQueueStatus(emailIds)
if err != nil {
log.Errorf("EmailQueueStatusManager.fetchEmailQueueStatus fail %s", err)
return rows
}
for i := range rows {
eq := objs[i].(*SEmailQueue)
if eqs, ok := rets[eq.Id]; ok {
rows[i].Status = eqs.Status
rows[i].SentAt = eqs.SentAt
rows[i].Results = eqs.Results
}
}
return rows
}
+63
View File
@@ -0,0 +1,63 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package models
import (
"time"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
)
type SEmailQueueStatusManager struct {
db.SModelBaseManager
}
type SEmailQueueStatus struct {
db.SModelBase
Id int64 `primary:"true" list:"user"`
SentAt time.Time `list:"user"`
Status string `width:"16" charset:"ascii" default:"queued" list:"user"`
Results string `list:"user" charset:"utf8"`
}
var EmailQueueStatusManager *SEmailQueueStatusManager
func init() {
EmailQueueStatusManager = &SEmailQueueStatusManager{
SModelBaseManager: db.NewModelBaseManager(SEmailQueueStatus{}, "emailqueue_status_tbl", "emailqueue_status", "emailqueue_status"),
}
EmailQueueStatusManager.SetVirtualObject(EmailQueueStatusManager)
}
func (manager *SEmailQueueStatusManager) fetchEmailQueueStatus(ids []int64) (map[int64]SEmailQueueStatus, error) {
q := manager.Query().In("id", ids)
results := make([]SEmailQueueStatus, 0)
err := q.All(&results)
if err != nil {
return nil, errors.Wrap(err, "query.All")
}
ret := make(map[int64]SEmailQueueStatus)
for i := range results {
eqs := results[i]
ret[eqs.Id] = eqs
}
return ret, nil
}
+13 -9
View File
@@ -16,30 +16,30 @@ package models
import (
"context"
"time"
"yunion.io/x/onecloud/pkg/cloudcommon/consts"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
)
type SEventManager struct {
db.SStandaloneAnonResourceBaseManager
db.SLogBaseManager
}
var EventManager *SEventManager
func init() {
func InitEventLog() {
EventManager = &SEventManager{
SStandaloneAnonResourceBaseManager: db.NewStandaloneAnonResourceBaseManager(
SEvent{},
"events_tbl",
"notifyevent",
"notifyevents",
),
SLogBaseManager: db.NewLogBaseManager(SEvent{}, "events2_tbl", "notifyevent", "notifyevents", "created_at", consts.OpsLogWithClickhouse),
}
EventManager.SetVirtualObject(EventManager)
}
type SEvent struct {
db.SStandaloneAnonResourceBase
db.SLogBase
// 资源创建时间
CreatedAt time.Time `nullable:"false" created_at:"true" index:"true" get:"user" list:"user" json:"created_at"`
Message string
Event string `width:"64" nullable:"true"`
@@ -68,3 +68,7 @@ func (e *SEventManager) GetEvent(id string) (*SEvent, error) {
}
return model.(*SEvent), nil
}
func (e *SEvent) GetRecordTime() time.Time {
return e.CreatedAt
}
+4 -4
View File
@@ -303,7 +303,7 @@ func (nm *SNotificationManager) PerformEventNotify(ctx context.Context, userCred
if nm.needWebconsole([]STopic{*topic}) {
// webconsole
err = nm.create(ctx, userCred, api.WEBCONSOLE, receiverIds, webconsoleContacts.UnsortedList(), input.Priority, event.Id)
err = nm.create(ctx, userCred, api.WEBCONSOLE, receiverIds, webconsoleContacts.UnsortedList(), input.Priority, event.GetId())
if err != nil {
output.FailedList = append(output.FailedList, api.FailedElem{
ContactType: api.WEBCONSOLE,
@@ -316,7 +316,7 @@ func (nm *SNotificationManager) PerformEventNotify(ctx context.Context, userCred
if ct == api.MOBILE {
continue
}
err := nm.create(ctx, userCred, ct, receiverIds, nil, input.Priority, event.Id)
err := nm.create(ctx, userCred, ct, receiverIds, nil, input.Priority, event.GetId())
if err != nil {
output.FailedList = append(output.FailedList, api.FailedElem{
ContactType: ct,
@@ -324,7 +324,7 @@ func (nm *SNotificationManager) PerformEventNotify(ctx context.Context, userCred
})
}
}
err = nm.createWithWebhookRobots(ctx, userCred, webhookRobots, input.Priority, event.Id)
err = nm.createWithWebhookRobots(ctx, userCred, webhookRobots, input.Priority, event.GetId())
if err != nil {
output.FailedList = append(output.FailedList, api.FailedElem{
ContactType: api.WEBHOOK,
@@ -332,7 +332,7 @@ func (nm *SNotificationManager) PerformEventNotify(ctx context.Context, userCred
})
}
// robot
err = nm.createWithRobots(ctx, userCred, robots, input.Priority, event.Id)
err = nm.createWithRobots(ctx, userCred, robots, input.Priority, event.GetId())
if err != nil {
output.FailedList = append(output.FailedList, api.FailedElem{
ContactType: api.ROBOT,
+15
View File
@@ -0,0 +1,15 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package sender // import "yunion.io/x/onecloud/pkg/notify/sender"
+110
View File
@@ -0,0 +1,110 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package sender
import (
"crypto/tls"
"encoding/base64"
"io"
"time"
"gopkg.in/mail.v2"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
api "yunion.io/x/onecloud/pkg/apis/notify"
)
type errorMap map[string]error
func (em errorMap) Error() string {
msg := make(map[string]string)
for k, e := range em {
msg[k] = e.Error()
}
return jsonutils.Marshal(msg).String()
}
func SendEmail(conf *api.SEmailConfig, msg *api.SEmailMessage) error {
dialer := mail.NewDialer(conf.Hostname, conf.Hostport, conf.Username, conf.Password)
if conf.SslGlobal {
dialer.SSL = true
} else {
dialer.SSL = false
dialer.TLSConfig = &tls.Config{
InsecureSkipVerify: true,
}
}
sender, err := dialer.Dial()
if err != nil {
return errors.Wrap(err, "dialer.Dial")
}
retErr := errorMap{}
for _, to := range msg.To {
gmsg := mail.NewMessage()
gmsg.SetHeader("From", conf.SenderAddress)
gmsg.SetHeader("To", to)
gmsg.SetHeader("Subject", msg.Subject)
gmsg.SetBody("text/html", msg.Body)
for _, attach := range msg.Attachments {
gmsg.Attach(attach.Filename,
mail.SetCopyFunc(func(w io.Writer) error {
mime := attach.Mime
if len(mime) == 0 {
mime = "application/octet-stream"
}
_, err := w.Write([]byte("Content-Type: " + attach.Mime))
return errors.Wrap(err, "WriteMime")
}),
mail.SetCopyFunc(func(w io.Writer) error {
contBytes, err := base64.StdEncoding.DecodeString(attach.Base64Content)
if err != nil {
return errors.Wrap(err, "base64.StdEncoding.DecodeString")
}
_, err = w.Write(contBytes)
return errors.Wrap(err, "WriteContent")
}),
)
}
errs := make([]error, 0)
for tryTime := 3; tryTime > 0; tryTime-- {
err = mail.Send(sender, gmsg)
log.Debugf("send email ...")
if err != nil {
errs = append(errs, err)
time.Sleep(time.Second * 10)
continue
}
errs = errs[0:0]
break
}
if len(errs) > 0 {
retErr[to] = errors.NewAggregate(errs)
}
}
if len(retErr) > 0 {
return retErr
}
return nil
}
+23
View File
@@ -0,0 +1,23 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package sender
import "yunion.io/x/onecloud/pkg/appsrv"
var Worker *appsrv.SWorkerManager
func init() {
Worker = appsrv.NewWorkerManager("notify_sender_worker", 1, 2048, true)
}
+5
View File
@@ -30,6 +30,9 @@ const (
func InitHandlers(app *appsrv.Application) {
db.InitAllManagers()
models.InitEventLog()
models.InitEmailQueue()
db.RegistUserCredCacheUpdater()
db.AddScopeResourceCountHandler(API_VERSION, app)
@@ -54,6 +57,7 @@ func InitHandlers(app *appsrv.Application) {
db.SharedResourceManager,
models.VerificationManager,
models.EventManager,
models.EmailQueueStatusManager,
} {
db.RegisterModelManager(manager)
}
@@ -68,6 +72,7 @@ func InitHandlers(app *appsrv.Application) {
models.TopicManager,
models.RobotManager,
models.SubscriberManager,
models.EmailQueueManager,
} {
db.RegisterModelManager(manager)
handler := db.NewModelHandler(manager)
+1 -12
View File
@@ -45,18 +45,7 @@ func GetCommandLogManager() *SCommandLogManager {
return commandLogManager
}
commandLogManager = &SCommandLogManager{
SOpsLogManager: db.SOpsLogManager{
SModelBaseManager: db.NewModelBaseManagerWithSplitable(
SCommandLog{},
"command_log_tbl",
"commandlog",
"commandlogs",
"id",
"start_time",
consts.SplitableMaxDuration(),
consts.SplitableMaxKeepMonths(),
),
},
SOpsLogManager: db.NewOpsLogManager(SCommandLog{}, "command_log_tbl", "commandlog", "commandlogs", "start_time", consts.OpsLogWithClickhouse),
}
commandLogManager.SetVirtualObject(commandLogManager)
return commandLogManager