Refactor push messaging (#5861)
This commit is contained in:
parent
ad9a5d0fed
commit
e7d89fa595
8 changed files with 150 additions and 243 deletions
|
|
@ -215,7 +215,7 @@ func configureMessengers(conf messagingConfig, cache *util.Cache) (chan push.Eve
|
|||
}
|
||||
|
||||
for _, service := range conf.Services {
|
||||
impl, err := push.NewMessengerFromConfig(service.Type, service.Other)
|
||||
impl, err := push.NewFromConfig(service.Type, service.Other)
|
||||
if err != nil {
|
||||
return messageChan, fmt.Errorf("failed configuring push service %s: %w", service.Type, err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -3,48 +3,42 @@ package push
|
|||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"github.com/evcc-io/evcc/util"
|
||||
)
|
||||
|
||||
// Sender implements message sending
|
||||
type Sender interface {
|
||||
// Messenger implements message sending
|
||||
type Messenger interface {
|
||||
Send(title, msg string)
|
||||
}
|
||||
|
||||
var log = util.NewLogger("push")
|
||||
type senderRegistry map[string]func(map[string]interface{}) (Messenger, error)
|
||||
|
||||
// NewMessengerFromConfig creates a new messenger
|
||||
func NewMessengerFromConfig(typ string, other map[string]interface{}) (res Sender, err error) {
|
||||
switch strings.ToLower(typ) {
|
||||
case "pushover":
|
||||
var cc pushOverConfig
|
||||
if err = util.DecodeOther(other, &cc); err == nil {
|
||||
res, err = NewPushOverMessenger(cc.App, cc.Recipients)
|
||||
func (r senderRegistry) Add(name string, factory func(map[string]interface{}) (Messenger, error)) {
|
||||
if _, exists := r[name]; exists {
|
||||
panic(fmt.Sprintf("cannot register duplicate messenger type: %s", name))
|
||||
}
|
||||
r[name] = factory
|
||||
}
|
||||
|
||||
func (r senderRegistry) Get(name string) (func(map[string]interface{}) (Messenger, error), error) {
|
||||
factory, exists := r[name]
|
||||
if !exists {
|
||||
return nil, fmt.Errorf("messenger type not registered: %s", name)
|
||||
}
|
||||
return factory, nil
|
||||
}
|
||||
|
||||
var registry senderRegistry = make(map[string]func(map[string]interface{}) (Messenger, error))
|
||||
|
||||
// NewFromConfig creates messenger from configuration
|
||||
func NewFromConfig(typ string, other map[string]interface{}) (v Messenger, err error) {
|
||||
factory, err := registry.Get(strings.ToLower(typ))
|
||||
if err == nil {
|
||||
if v, err = factory(other); err != nil {
|
||||
err = fmt.Errorf("cannot create messenger '%s': %w", typ, err)
|
||||
}
|
||||
case "telegram":
|
||||
var cc telegramConfig
|
||||
if err = util.DecodeOther(other, &cc); err == nil {
|
||||
res, err = NewTelegramMessenger(cc.Token, cc.Chats)
|
||||
}
|
||||
case "email", "shout":
|
||||
var cc shoutrrrConfig
|
||||
if err = util.DecodeOther(other, &cc); err == nil {
|
||||
res, err = NewShoutrrrMessenger(cc.URI)
|
||||
}
|
||||
case "script":
|
||||
var cc scriptConfig
|
||||
if err = util.DecodeOther(other, &cc); err == nil {
|
||||
res, err = NewScriptMessenger(cc.CmdLine, cc.Timeout, cc.Scale, cc.Cache)
|
||||
}
|
||||
case "ntfy":
|
||||
var cc ntfyConfig
|
||||
if err = util.DecodeOther(other, &cc); err == nil {
|
||||
res, err = NewNtfyMessenger(cc.URI, cc.Priority, cc.Tags)
|
||||
}
|
||||
default:
|
||||
err = fmt.Errorf("unknown messenger type: %s", typ)
|
||||
} else {
|
||||
err = fmt.Errorf("invalid messenger type: %s", typ)
|
||||
}
|
||||
|
||||
return res, err
|
||||
return
|
||||
}
|
||||
|
|
|
|||
|
|
@ -28,7 +28,7 @@ type EventTemplate struct {
|
|||
// Hub subscribes to event notifications and sends them to client devices
|
||||
type Hub struct {
|
||||
definitions map[string]EventTemplate
|
||||
sender []Sender
|
||||
sender []Messenger
|
||||
cache *util.Cache
|
||||
}
|
||||
|
||||
|
|
@ -62,7 +62,7 @@ func NewHub(cc map[string]EventTemplateConfig, cache *util.Cache) (*Hub, error)
|
|||
}
|
||||
|
||||
// Add adds a sender to the list of senders
|
||||
func (h *Hub) Add(sender Sender) {
|
||||
func (h *Hub) Add(sender Messenger) {
|
||||
h.sender = append(h.sender, sender)
|
||||
}
|
||||
|
||||
|
|
@ -91,6 +91,8 @@ func (h *Hub) apply(ev Event, tmpl *template.Template) (string, error) {
|
|||
|
||||
// Run is the Hub's main publishing loop
|
||||
func (h *Hub) Run(events <-chan Event) {
|
||||
log := util.NewLogger("push")
|
||||
|
||||
for ev := range events {
|
||||
if len(h.sender) == 0 {
|
||||
continue
|
||||
|
|
|
|||
39
push/ntfy.go
39
push/ntfy.go
|
|
@ -5,32 +5,43 @@ import (
|
|||
"net/http"
|
||||
"strings"
|
||||
|
||||
"github.com/evcc-io/evcc/util"
|
||||
"github.com/evcc-io/evcc/util/request"
|
||||
)
|
||||
|
||||
func init() {
|
||||
registry.Add("ntfy", NewNtfyFromConfig)
|
||||
}
|
||||
|
||||
// Ntfy implements the ntfy messaging aggregator
|
||||
type Ntfy struct {
|
||||
log *util.Logger
|
||||
uri string
|
||||
priority string
|
||||
tags string
|
||||
}
|
||||
|
||||
type ntfyConfig struct {
|
||||
URI string
|
||||
Priority string
|
||||
Tags string
|
||||
}
|
||||
// NewNtfyFromConfig creates new Ntfy messenger
|
||||
func NewNtfyFromConfig(other map[string]interface{}) (Messenger, error) {
|
||||
var cc struct {
|
||||
URI string
|
||||
Priority string
|
||||
Tags string
|
||||
}
|
||||
|
||||
// NewNtfyMessenger creates new Ntfy messenger
|
||||
func NewNtfyMessenger(uri string, priority string, tags string) (*Ntfy, error) {
|
||||
if uri == "" {
|
||||
return nil, errors.New("ntfy: missing uri")
|
||||
if err := util.DecodeOther(other, &cc); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if cc.URI == "" {
|
||||
return nil, errors.New("missing uri")
|
||||
}
|
||||
|
||||
m := &Ntfy{
|
||||
uri: uri,
|
||||
priority: priority,
|
||||
tags: tags,
|
||||
log: util.NewLogger("ntfy"),
|
||||
uri: cc.URI,
|
||||
priority: cc.Priority,
|
||||
tags: cc.Tags,
|
||||
}
|
||||
|
||||
return m, nil
|
||||
|
|
@ -44,10 +55,10 @@ func (m *Ntfy) Send(title, msg string) {
|
|||
"Tags": m.tags,
|
||||
})
|
||||
if err != nil {
|
||||
log.ERROR.Printf("ntfy: %v", err)
|
||||
m.log.ERROR.Printf("ntfy: %v", err)
|
||||
}
|
||||
|
||||
if _, err := http.DefaultClient.Do(req); err != nil {
|
||||
log.ERROR.Printf("ntfy: %v", err)
|
||||
m.log.ERROR.Printf("ntfy: %v", err)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -3,30 +3,41 @@ package push
|
|||
import (
|
||||
"errors"
|
||||
|
||||
"github.com/evcc-io/evcc/util"
|
||||
"github.com/gregdel/pushover"
|
||||
)
|
||||
|
||||
func init() {
|
||||
registry.Add("pushover", NewPushOverFromConfig)
|
||||
}
|
||||
|
||||
// PushOver implements the pushover messenger
|
||||
type PushOver struct {
|
||||
log *util.Logger
|
||||
app *pushover.Pushover
|
||||
recipients []string
|
||||
}
|
||||
|
||||
type pushOverConfig struct {
|
||||
App string
|
||||
Recipients []string
|
||||
Events map[string]EventTemplate
|
||||
}
|
||||
// NewPushOverFromConfig creates new pushover messenger
|
||||
func NewPushOverFromConfig(other map[string]interface{}) (Messenger, error) {
|
||||
var cc struct {
|
||||
App string
|
||||
Recipients []string
|
||||
Events map[string]EventTemplate
|
||||
}
|
||||
|
||||
// NewPushOverMessenger creates new pushover messenger
|
||||
func NewPushOverMessenger(app string, recipients []string) (*PushOver, error) {
|
||||
if app == "" {
|
||||
return nil, errors.New("pushover: missing app name")
|
||||
if err := util.DecodeOther(other, &cc); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if cc.App == "" {
|
||||
return nil, errors.New("missing app name")
|
||||
}
|
||||
|
||||
m := &PushOver{
|
||||
app: pushover.New(app),
|
||||
recipients: recipients,
|
||||
log: util.NewLogger("pushover"),
|
||||
app: pushover.New(cc.App),
|
||||
recipients: cc.Recipients,
|
||||
}
|
||||
|
||||
return m, nil
|
||||
|
|
@ -38,11 +49,11 @@ func (m *PushOver) Send(title, msg string) {
|
|||
|
||||
for _, id := range m.recipients {
|
||||
go func(id string) {
|
||||
log.DEBUG.Printf("pushover: sending to %s", id)
|
||||
m.log.DEBUG.Printf("sending to %s", id)
|
||||
|
||||
recipient := pushover.NewRecipient(id)
|
||||
if _, err := m.app.SendMessage(message, recipient); err != nil {
|
||||
log.ERROR.Print(err)
|
||||
m.log.ERROR.Print(err)
|
||||
}
|
||||
}(id)
|
||||
}
|
||||
|
|
|
|||
174
push/script.go
174
push/script.go
|
|
@ -3,85 +3,57 @@ package push
|
|||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"math"
|
||||
"os/exec"
|
||||
"regexp"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/evcc-io/evcc/util"
|
||||
"github.com/evcc-io/evcc/util/jq"
|
||||
"github.com/itchyny/gojq"
|
||||
"github.com/evcc-io/evcc/util/request"
|
||||
"github.com/kballard/go-shellquote"
|
||||
)
|
||||
|
||||
func init() {
|
||||
registry.Add("script", NewScriptFromConfig)
|
||||
}
|
||||
|
||||
// Script implements shell script-based message service and setters
|
||||
type Script struct {
|
||||
log *util.Logger
|
||||
script string
|
||||
timeout time.Duration
|
||||
cache time.Duration
|
||||
updated time.Time
|
||||
val string
|
||||
err error
|
||||
re *regexp.Regexp
|
||||
jq *gojq.Query
|
||||
scale float64
|
||||
}
|
||||
|
||||
type scriptConfig struct {
|
||||
CmdLine string
|
||||
Timeout time.Duration
|
||||
Scale float64
|
||||
Cache time.Duration
|
||||
}
|
||||
// NewScriptFromConfig creates a Script messenger. Script execution is aborted after given timeout.
|
||||
func NewScriptFromConfig(other map[string]interface{}) (Messenger, error) {
|
||||
cc := struct {
|
||||
CmdLine string
|
||||
Timeout time.Duration
|
||||
}{
|
||||
Timeout: request.Timeout,
|
||||
}
|
||||
|
||||
if err := util.DecodeOther(other, &cc); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// NewScriptMessenger creates a Script messenger. Script execution is aborted after given timeout.
|
||||
func NewScriptMessenger(script string, timeout time.Duration, scale float64, cache time.Duration) (*Script, error) {
|
||||
s := &Script{
|
||||
log: util.NewLogger("script"),
|
||||
script: script,
|
||||
timeout: timeout,
|
||||
scale: scale,
|
||||
cache: cache,
|
||||
script: cc.CmdLine,
|
||||
timeout: cc.Timeout,
|
||||
}
|
||||
|
||||
return s, nil
|
||||
}
|
||||
|
||||
func (p *Script) WithRegex(regex string) (*Script, error) {
|
||||
re, err := regexp.Compile(regex)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("invalid regex '%s': %w", re, err)
|
||||
}
|
||||
|
||||
p.re = re
|
||||
|
||||
return p, nil
|
||||
}
|
||||
|
||||
func (p *Script) WithJq(jq string) (*Script, error) {
|
||||
op, err := gojq.Parse(jq)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("invalid jq query '%s': %w", jq, err)
|
||||
}
|
||||
|
||||
p.jq = op
|
||||
|
||||
return p, nil
|
||||
}
|
||||
|
||||
// Send calls the script
|
||||
func (m *Script) Send(title, msg string) {
|
||||
_, err := m.exec(m.script + " '" + title + "' '" + msg + "'")
|
||||
_, err := m.exec(m.script, title, msg)
|
||||
if err != nil {
|
||||
m.log.ERROR.Printf("Script message error: %v", err)
|
||||
m.log.ERROR.Printf("exec: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func (m *Script) exec(script string) (string, error) {
|
||||
func (m *Script) exec(script, title, msg string) (string, error) {
|
||||
args, err := shellquote.Split(script)
|
||||
if err != nil {
|
||||
return "", err
|
||||
|
|
@ -90,6 +62,7 @@ func (m *Script) exec(script string) (string, error) {
|
|||
ctx, cancel := context.WithTimeout(context.Background(), m.timeout)
|
||||
defer cancel()
|
||||
|
||||
args = append(args, title, msg)
|
||||
cmd := exec.CommandContext(ctx, args[0], args[1:]...)
|
||||
b, err := cmd.Output()
|
||||
|
||||
|
|
@ -110,104 +83,3 @@ func (m *Script) exec(script string) (string, error) {
|
|||
|
||||
return s, nil
|
||||
}
|
||||
|
||||
// StringGetter returns string from exec result. Only STDOUT is considered.
|
||||
func (m *Script) StringGetter() func() (string, error) {
|
||||
return func() (string, error) {
|
||||
if time.Since(m.updated) > m.cache {
|
||||
m.val, m.err = m.exec(m.script)
|
||||
m.updated = time.Now()
|
||||
|
||||
if m.err == nil && m.re != nil {
|
||||
ma := m.re.FindStringSubmatch(m.val)
|
||||
if len(ma) > 1 {
|
||||
m.val = ma[1] // first submatch
|
||||
}
|
||||
}
|
||||
|
||||
if m.err == nil && m.jq != nil {
|
||||
var v interface{}
|
||||
if v, m.err = jq.Query(m.jq, []byte(m.val)); m.err == nil {
|
||||
m.val = fmt.Sprintf("%v", v)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return m.val, m.err
|
||||
}
|
||||
}
|
||||
|
||||
// FloatGetter parses float from exec result
|
||||
func (m *Script) FloatGetter() func() (float64, error) {
|
||||
g := m.StringGetter()
|
||||
|
||||
return func() (float64, error) {
|
||||
s, err := g()
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
f, err := strconv.ParseFloat(s, 64)
|
||||
if err == nil {
|
||||
f *= m.scale
|
||||
}
|
||||
|
||||
return f, err
|
||||
}
|
||||
}
|
||||
|
||||
// IntGetter parses int64 from exec result
|
||||
func (m *Script) IntGetter() func() (int64, error) {
|
||||
g := m.FloatGetter()
|
||||
|
||||
return func() (int64, error) {
|
||||
f, err := g()
|
||||
return int64(math.Round(f)), err
|
||||
}
|
||||
}
|
||||
|
||||
// BoolGetter parses bool from exec result. "on", "true" and 1 are considered truish.
|
||||
func (m *Script) BoolGetter() func() (bool, error) {
|
||||
g := m.StringGetter()
|
||||
|
||||
return func() (bool, error) {
|
||||
s, err := g()
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
|
||||
return util.Truish(s), nil
|
||||
}
|
||||
}
|
||||
|
||||
// IntSetter invokes script with parameter replaced by int value
|
||||
func (m *Script) IntSetter(param string) func(int64) error {
|
||||
// return func to access cached value
|
||||
return func(i int64) error {
|
||||
cmd, err := util.ReplaceFormatted(m.script, map[string]interface{}{
|
||||
param: i,
|
||||
})
|
||||
|
||||
if err == nil {
|
||||
_, err = m.exec(cmd)
|
||||
}
|
||||
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
// BoolSetter invokes script with parameter replaced by bool value
|
||||
func (m *Script) BoolSetter(param string) func(bool) error {
|
||||
// return func to access cached value
|
||||
return func(b bool) error {
|
||||
cmd, err := util.ReplaceFormatted(m.script, map[string]interface{}{
|
||||
param: b,
|
||||
})
|
||||
|
||||
if err == nil {
|
||||
_, err = m.exec(cmd)
|
||||
}
|
||||
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,30 +1,40 @@
|
|||
package push
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/containrrr/shoutrrr"
|
||||
"github.com/containrrr/shoutrrr/pkg/router"
|
||||
"github.com/containrrr/shoutrrr/pkg/types"
|
||||
"github.com/evcc-io/evcc/util"
|
||||
)
|
||||
|
||||
func init() {
|
||||
registry.Add("email", NewShoutrrrFromConfig)
|
||||
registry.Add("shout", NewShoutrrrFromConfig)
|
||||
}
|
||||
|
||||
// Shoutrrr implements the shoutrrr messaging aggregator
|
||||
type Shoutrrr struct {
|
||||
log *util.Logger
|
||||
app *router.ServiceRouter
|
||||
}
|
||||
|
||||
type shoutrrrConfig struct {
|
||||
URI string
|
||||
}
|
||||
// NewShoutrrrFromConfig creates new Shoutrrr messenger
|
||||
func NewShoutrrrFromConfig(other map[string]interface{}) (Messenger, error) {
|
||||
var cc struct {
|
||||
URI string
|
||||
}
|
||||
|
||||
// NewShoutrrrMessenger creates new Shoutrrr messenger
|
||||
func NewShoutrrrMessenger(uri string) (*Shoutrrr, error) {
|
||||
app, err := shoutrrr.CreateSender(uri)
|
||||
if err := util.DecodeOther(other, &cc); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
app, err := shoutrrr.CreateSender(cc.URI)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("shoutrrr: %v", err)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
m := &Shoutrrr{
|
||||
log: util.NewLogger("shoutrrr"),
|
||||
app: app,
|
||||
}
|
||||
|
||||
|
|
@ -39,7 +49,7 @@ func (m *Shoutrrr) Send(title, msg string) {
|
|||
|
||||
for _, err := range m.app.Send(msg, params) {
|
||||
if err != nil {
|
||||
log.ERROR.Printf("shoutrrr: %v", err)
|
||||
m.log.ERROR.Println("send:", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -4,40 +4,48 @@ import (
|
|||
"errors"
|
||||
"sync"
|
||||
|
||||
"github.com/evcc-io/evcc/util"
|
||||
tgbotapi "github.com/go-telegram-bot-api/telegram-bot-api/v5"
|
||||
)
|
||||
|
||||
func init() {
|
||||
registry.Add("telegram", NewTelegramFromConfig)
|
||||
}
|
||||
|
||||
// Telegram implements the Telegram messenger
|
||||
type Telegram struct {
|
||||
log *util.Logger
|
||||
sync.Mutex
|
||||
bot *tgbotapi.BotAPI
|
||||
chats map[int64]struct{}
|
||||
}
|
||||
|
||||
type telegramConfig struct {
|
||||
Token string
|
||||
Chats []int64
|
||||
}
|
||||
|
||||
func init() {
|
||||
if err := tgbotapi.SetLogger(log.ERROR); err != nil {
|
||||
log.ERROR.Printf("telegram: %v", err)
|
||||
// NewTelegramFromConfig creates new pushover messenger
|
||||
func NewTelegramFromConfig(other map[string]interface{}) (Messenger, error) {
|
||||
var cc struct {
|
||||
Token string
|
||||
Chats []int64
|
||||
}
|
||||
}
|
||||
|
||||
// NewTelegramMessenger creates new pushover messenger
|
||||
func NewTelegramMessenger(token string, chats []int64) (*Telegram, error) {
|
||||
bot, err := tgbotapi.NewBotAPI(token)
|
||||
if err := util.DecodeOther(other, &cc); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
bot, err := tgbotapi.NewBotAPI(cc.Token)
|
||||
if err != nil {
|
||||
return nil, errors.New("telegram: invalid bot token")
|
||||
}
|
||||
|
||||
log := util.NewLogger("telegram")
|
||||
_ = tgbotapi.SetLogger(log.ERROR)
|
||||
|
||||
m := &Telegram{
|
||||
log: log,
|
||||
bot: bot,
|
||||
chats: make(map[int64]struct{}),
|
||||
}
|
||||
|
||||
for _, chat := range chats {
|
||||
for _, chat := range cc.Chats {
|
||||
m.chats[chat] = struct{}{}
|
||||
}
|
||||
|
||||
|
|
@ -54,8 +62,7 @@ func (m *Telegram) trackChats() {
|
|||
for update := range m.bot.GetUpdatesChan(conf) {
|
||||
m.Lock()
|
||||
if _, ok := m.chats[update.Message.Chat.ID]; !ok {
|
||||
log.INFO.Printf("telegram: new chat id: %d", update.Message.Chat.ID)
|
||||
// m.chats[update.Message.Chat.ID] = struct{}{}
|
||||
m.log.INFO.Printf("new chat id: %d", update.Message.Chat.ID)
|
||||
}
|
||||
m.Unlock()
|
||||
}
|
||||
|
|
@ -65,11 +72,11 @@ func (m *Telegram) trackChats() {
|
|||
func (m *Telegram) Send(title, msg string) {
|
||||
m.Lock()
|
||||
for chat := range m.chats {
|
||||
log.DEBUG.Printf("telegram: sending to %d", chat)
|
||||
m.log.DEBUG.Printf("sending to %d", chat)
|
||||
|
||||
msg := tgbotapi.NewMessage(chat, msg)
|
||||
if _, err := m.bot.Send(msg); err != nil {
|
||||
log.ERROR.Print(err)
|
||||
m.log.ERROR.Println("send:", err)
|
||||
}
|
||||
}
|
||||
m.Unlock()
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue