/
out.go
110 lines (104 loc) · 2.78 KB
/
out.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
package sidecar
import (
"bufio"
"context"
"fmt"
"io/ioutil"
"net/http"
"os"
"time"
dfv1 "github.com/argoproj-labs/argo-dataflow/api/v1alpha1"
"github.com/google/uuid"
"github.com/opentracing/opentracing-go"
runtimeutil "k8s.io/apimachinery/pkg/util/runtime"
)
func connectOut(ctx context.Context, sink func(context.Context, []byte) error) {
connectOutFIFO(ctx, sink)
connectOutHTTP(sink)
}
func connectOutHTTP(sink func(context.Context, []byte) error) {
logger.Info("HTTP out interface configured")
v, err := ioutil.ReadFile(dfv1.PathAuthorization)
if err != nil {
panic(fmt.Errorf("failed to read authorization file: %w", err))
}
authorization := string(v)
http.HandleFunc("/messages", func(w http.ResponseWriter, r *http.Request) {
span, ctx := opentracing.StartSpanFromContext(r.Context(), "/messages")
defer span.Finish()
if r.Header.Get("Authorization") != authorization {
w.WriteHeader(403)
return
}
data, err := ioutil.ReadAll(r.Body)
_ = r.Body.Close()
if err != nil {
logger.Error(err, "failed to read message body from main via HTTP")
w.WriteHeader(400)
_, _ = w.Write([]byte(err.Error()))
return
}
id := r.Header.Get(dfv1.MetaID)
if id == "" {
id = uuid.New().String()
}
if err := sink(
dfv1.ContextWithMeta(
ctx,
dfv1.Meta{
Source: fmt.Sprintf("urn:dataflow:pod:%s.pod.%s.%s:messages", pod, namespace, cluster),
ID: id,
Time: time.Now().Unix(),
},
),
data,
); err != nil {
logger.Error(err, "failed to send message from main to sink")
w.WriteHeader(500)
_, _ = w.Write([]byte(err.Error()))
} else {
w.WriteHeader(204)
}
})
}
func connectOutFIFO(ctx context.Context, sink func(context.Context, []byte) error) {
logger.Info("FIFO out interface configured")
go func() {
defer runtimeutil.HandleCrash()
err := func() error {
fifo, err := os.OpenFile(dfv1.PathFIFOOut, os.O_RDONLY, os.ModeNamedPipe)
if err != nil {
return fmt.Errorf("failed to open output FIFO: %w", err)
}
addStopHook(func(ctx context.Context) error {
logger.Info("closing out FIFO")
return fifo.Close()
})
logger.Info("opened output FIFO")
scanner := bufio.NewScanner(fifo)
for scanner.Scan() {
if err := sink(
dfv1.ContextWithMeta(
ctx,
dfv1.Meta{
Source: fmt.Sprintf("urn:dataflow:pod:%s.pod.%s.%s:fifo", pod, namespace, cluster),
ID: uuid.New().String(),
Time: time.Now().Unix(),
},
),
scanner.Bytes(),
); err != nil {
return fmt.Errorf("failed to send message from main to sink: %w", err)
}
}
if err = scanner.Err(); err != nil {
return fmt.Errorf("scanner error: %w", err)
}
return nil
}()
if err != nil {
logger.Error(err, "failed to received message from FIFO")
os.Exit(1)
}
}()
}