handler.go 1.1 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364
  1. package handler
  2. import (
  3. "time"
  4. "zero/stash/es"
  5. "zero/stash/filter"
  6. jsoniter "github.com/json-iterator/go"
  7. )
  8. const (
  9. timestampFormat = "2006-01-02T15:04:05.000Z"
  10. timestampKey = "@timestamp"
  11. )
  12. type MessageHandler struct {
  13. writer *es.Writer
  14. filters []filter.FilterFunc
  15. }
  16. func NewHandler(writer *es.Writer) *MessageHandler {
  17. return &MessageHandler{
  18. writer: writer,
  19. }
  20. }
  21. func (mh *MessageHandler) AddFilters(filters ...filter.FilterFunc) {
  22. for _, f := range filters {
  23. mh.filters = append(mh.filters, f)
  24. }
  25. }
  26. func (mh *MessageHandler) Consume(_, val string) error {
  27. m := make(map[string]interface{})
  28. if err := jsoniter.Unmarshal([]byte(val), &m); err != nil {
  29. return err
  30. }
  31. for _, proc := range mh.filters {
  32. if m = proc(m); m == nil {
  33. return nil
  34. }
  35. }
  36. bs, err := jsoniter.Marshal(m)
  37. if err != nil {
  38. return err
  39. }
  40. return mh.writer.Write(mh.getTime(m), string(bs))
  41. }
  42. func (mh *MessageHandler) getTime(m map[string]interface{}) time.Time {
  43. if ti, ok := m[timestampKey]; ok {
  44. if ts, ok := ti.(string); ok {
  45. if t, err := time.Parse(timestampFormat, ts); err == nil {
  46. return t
  47. }
  48. }
  49. }
  50. return time.Now()
  51. }