-
Notifications
You must be signed in to change notification settings - Fork 8
/
Copy pathhelper_test.go
110 lines (103 loc) · 2.74 KB
/
helper_test.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 gokini
import (
"os"
"time"
"fmt"
"github.com/aws/aws-sdk-go/aws"
awsclient "github.com/aws/aws-sdk-go/aws/client"
"github.com/aws/aws-sdk-go/aws/session"
"github.com/aws/aws-sdk-go/service/dynamodb"
"github.com/aws/aws-sdk-go/service/kinesis"
log "github.com/sirupsen/logrus"
)
func createStream(streamName string, shards int64) error {
session, err := session.NewSessionWithOptions(
session.Options{
SharedConfigState: session.SharedConfigEnable,
Config: aws.Config{
CredentialsChainVerboseErrors: aws.Bool(true),
Endpoint: aws.String(os.Getenv("KINESIS_ENDPOINT")),
Retryer: awsclient.DefaultRetryer{NumMaxRetries: 1},
},
},
)
if err != nil {
return fmt.Errorf("Error starting kinesis client %s", err)
}
svc := kinesis.New(session)
_, err = svc.CreateStream(&kinesis.CreateStreamInput{
ShardCount: aws.Int64(shards),
StreamName: aws.String(streamName),
})
if err != nil {
return err
}
time.Sleep(500 * time.Millisecond)
return nil
}
func pushRecordToKinesis(streamName string, record []byte, create bool) error {
session, err := session.NewSessionWithOptions(
session.Options{
SharedConfigState: session.SharedConfigEnable,
Config: aws.Config{
CredentialsChainVerboseErrors: aws.Bool(true),
Endpoint: aws.String(os.Getenv("KINESIS_ENDPOINT")),
Retryer: awsclient.DefaultRetryer{NumMaxRetries: 1},
},
},
)
if err != nil {
log.Errorf("Error starting kinesis client %s", err)
return err
}
svc := kinesis.New(session)
if create {
if err := createStream(streamName, 1); err != nil {
return err
}
}
_, err = svc.PutRecord(&kinesis.PutRecordInput{
Data: record,
PartitionKey: aws.String("abc123"),
StreamName: &streamName,
})
if err != nil {
log.Errorf("Error sending data to kinesis %s", err)
}
return err
}
func deleteStream(streamName string) {
session, _ := session.NewSessionWithOptions(
session.Options{
SharedConfigState: session.SharedConfigEnable,
Config: aws.Config{
CredentialsChainVerboseErrors: aws.Bool(true),
Endpoint: aws.String(os.Getenv("KINESIS_ENDPOINT")),
},
},
)
svc := kinesis.New(session)
_, err := svc.DeleteStream(&kinesis.DeleteStreamInput{
StreamName: &streamName,
})
if err != nil {
log.Errorln(err)
}
}
func deleteTable(tableName string) {
session, _ := session.NewSessionWithOptions(
session.Options{
Config: aws.Config{
Endpoint: aws.String(os.Getenv("DYNAMODB_ENDPOINT")),
},
SharedConfigState: session.SharedConfigEnable,
},
)
svc := dynamodb.New(session)
_, err := svc.DeleteTable(&dynamodb.DeleteTableInput{
TableName: &tableName,
})
if err != nil {
log.Errorln(err)
}
}