|
1 | 1 | package handler |
2 | 2 |
|
3 | 3 | import ( |
4 | | - "encoding/json" |
5 | 4 | "errors" |
6 | | - "log" |
7 | 5 | "strings" |
8 | 6 |
|
9 | 7 | "github.com/gofiber/fiber/v2" |
10 | 8 | "github.com/memphisdev/memphis.go" |
11 | | - "google.golang.org/protobuf/proto" |
12 | 9 | ) |
13 | 10 |
|
14 | | -type Handler struct{ P *memphis.Producer } |
| 11 | +var producers = make(map[string]*memphis.Producer) |
15 | 12 |
|
16 | | -func (p Handler) produce(message []byte, hdrs memphis.Headers) error { |
17 | | - if err := p.P.Produce(message, memphis.MsgHeaders(hdrs)); err != nil { |
18 | | - return err |
19 | | - } |
20 | | - return nil |
21 | | - |
22 | | -} |
| 13 | +func handleHeaders(headers map[string]string) (memphis.Headers, error) { |
| 14 | + hdrs := memphis.Headers{} |
| 15 | + hdrs.New() |
23 | 16 |
|
24 | | -func handleMessageHdrs(bodyReq []byte, hdrs memphis.Headers) ([]byte, memphis.Headers, error) { |
25 | | - type body struct { |
26 | | - Message string `json:"message"` |
27 | | - Headers string `json:"headers"` |
28 | | - } |
29 | | - var b body |
30 | | - err := json.Unmarshal(bodyReq, &b) |
31 | | - if err != nil { |
32 | | - return nil, memphis.Headers{}, err |
33 | | - } |
34 | | - |
35 | | - var headers map[string]string |
36 | | - err = json.Unmarshal([]byte(b.Headers), &headers) |
37 | | - if err != nil { |
38 | | - return nil, memphis.Headers{}, err |
39 | | - } |
40 | | - |
41 | | - var k, v string |
42 | 17 | for key, value := range headers { |
43 | | - k = key |
44 | | - v = value |
45 | | - |
46 | | - err = hdrs.Add(k, v) |
| 18 | + err := hdrs.Add(key, value) |
47 | 19 | if err != nil { |
48 | | - return nil, memphis.Headers{}, err |
| 20 | + return memphis.Headers{}, err |
49 | 21 | } |
50 | 22 | } |
| 23 | + return hdrs, nil |
| 24 | +} |
51 | 25 |
|
52 | | - message, err := json.Marshal(b.Message) |
53 | | - if err != nil { |
54 | | - return nil, memphis.Headers{}, err |
| 26 | +func createProducer(conn *memphis.Conn, producers map[string]*memphis.Producer, stationName string) (*memphis.Producer, error) { |
| 27 | + producerName := "http_proxy" |
| 28 | + var producer *memphis.Producer |
| 29 | + var err error |
| 30 | + if _, ok := producers[stationName]; !ok { |
| 31 | + producer, err = conn.CreateProducer(stationName, producerName, memphis.ProducerGenUniqueSuffix()) |
| 32 | + if err != nil { |
| 33 | + return nil, err |
| 34 | + } |
| 35 | + producers[stationName] = producer |
| 36 | + } else { |
| 37 | + producer = producers[stationName] |
55 | 38 | } |
56 | | - return message, hdrs, nil |
| 39 | + |
| 40 | + return producer, nil |
57 | 41 | } |
58 | 42 |
|
59 | | -func CreateHandleMessage(p *memphis.Producer) func(*fiber.Ctx) error { |
| 43 | +func CreateHandleMessage(conn *memphis.Conn) func(*fiber.Ctx) error { |
60 | 44 | return func(c *fiber.Ctx) error { |
| 45 | + stationName := c.Params("stationName") |
| 46 | + var producer *memphis.Producer |
| 47 | + |
| 48 | + producer, err := createProducer(conn, producers, stationName) |
| 49 | + if err != nil { |
| 50 | + return err |
| 51 | + } |
| 52 | + |
61 | 53 | bodyReq := c.Body() |
| 54 | + headers := c.GetReqHeaders() |
62 | 55 | contentType := string(c.Request().Header.ContentType()) |
63 | | - var message []byte |
64 | | - var err error |
65 | | - hdrs := memphis.Headers{} |
66 | | - hdrs.New() |
67 | 56 | caseText := strings.Contains(contentType, "text") |
68 | | - |
69 | 57 | if caseText { |
70 | 58 | contentType = "text/" |
71 | 59 | } |
72 | 60 |
|
73 | 61 | switch contentType { |
74 | | - case "application/json": |
75 | | - message, hdrs, err = handleMessageHdrs(bodyReq, hdrs) |
| 62 | + case "application/json", "text/", "application/x-protobuf": |
| 63 | + message := bodyReq |
| 64 | + hdrs, err := handleHeaders(headers) |
76 | 65 | if err != nil { |
77 | 66 | return err |
78 | 67 | } |
79 | | - case "text/": |
80 | | - message, hdrs, err = handleMessageHdrs(bodyReq, hdrs) |
81 | | - if err != nil { |
82 | | - return err |
83 | | - } |
84 | | - case "application/x-protobuf": |
85 | | - msg := &Msg{} |
86 | | - err := proto.Unmarshal(bodyReq, msg) |
87 | | - if err != nil { |
88 | | - log.Fatal("unmarshaling error: ", err) |
89 | | - } |
90 | | - |
91 | | - message, err = json.Marshal(msg.Message) |
92 | | - if err != nil { |
93 | | - return err |
94 | | - } |
95 | | - |
96 | | - var headers map[string]string |
97 | | - err = json.Unmarshal([]byte(msg.Headers), &headers) |
98 | | - if err != nil { |
99 | | - return err |
100 | | - } |
101 | | - |
102 | | - var k, v string |
103 | | - for key, value := range headers { |
104 | | - k = key |
105 | | - v = value |
106 | | - err = hdrs.Add(k, v) |
107 | | - if err != nil { |
108 | | - return err |
109 | | - } |
| 68 | + if err := producer.Produce(message, memphis.MsgHeaders(hdrs)); err != nil { |
| 69 | + c.Status(400) |
| 70 | + return c.JSON(&fiber.Map{ |
| 71 | + "success": false, |
| 72 | + "error": err.Error(), |
| 73 | + }) |
110 | 74 | } |
111 | 75 | default: |
112 | 76 | return errors.New("unsupported content type") |
113 | 77 | } |
114 | 78 |
|
115 | | - if err := p.Produce(message, memphis.MsgHeaders(hdrs)); err != nil { |
116 | | - c.Status(400) |
117 | | - return c.JSON(&fiber.Map{ |
118 | | - "success": false, |
119 | | - "error": err.Error(), |
120 | | - }) |
121 | | - } |
122 | | - |
123 | 79 | c.Status(200) |
124 | 80 | return c.JSON(&fiber.Map{ |
125 | 81 | "success": true, |
|
0 commit comments