Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ BUILD_VERSIONS_DEBUG = $(shell jq -r '.versions|map("build-\(.)-debug")[]' ${TAR
STORE_MOD_VERSIONS = $(shell jq -r '.versions|map("store-mod-\(.)")[]' ${TARGETS})
TEST_VERSIONS = $(shell jq -r '.versions|map("test-\(.)")[]' ${TARGETS})
COVERAGE_VERSIONS = $(shell jq -r '.versions|map("coverage-\(.)")[]' ${TARGETS})
BENCHMARK_VERSIONS = $(shell jq -r '.versions|map("benchmark-\(.)")[]' ${TARGETS})
BRANCH := $(shell git rev-parse --abbrev-ref HEAD)
COMMIT := $(shell git log -1 --format='%H')

Expand Down Expand Up @@ -56,6 +57,11 @@ $(TEST_VERSIONS):
go test -v -failfast -race -count=1 \
-tags $(shell echo $@ | sed -e 's/test-/sdk_/g' -e 's/-/_/g'),muslc \
./...

$(BENCHMARK_VERSIONS):
go test -v -failfast -bench=. -run=^# -benchmem -count=1 \
-tags $(shell echo $@ | sed -e 's/benchmark-/sdk_/g' -e 's/-/_/g'),muslc \
./...

$(COVERAGE_VERSIONS):
go test -v -failfast -coverprofile=coverage.out -covermode=atomic -count=1\
Expand Down
2 changes: 1 addition & 1 deletion cmd/tracelistener/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,7 @@ func main() {
if ca.existingDatabasePath != "" {
importer := bulk.Importer{
Path: ca.existingDatabasePath,
TraceWatcher: watcher,
TraceWatcher: &watcher,
Processor: dpi,
Logger: logger,
Database: di,
Expand Down
112 changes: 112 additions & 0 deletions cmd/tracestats/main.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
package main

import (
"bufio"
"encoding/csv"
"encoding/json"
"fmt"
"os"
"strconv"

"github.com/allinbits/tracelistener/tracelistener"
)

type traceInfo struct {
BlockHeight uint64
KeyLength uint64
ValueLength uint64
Length uint64
}

type traceInfos []traceInfo

func (ti traceInfos) CSV() [][]string {
ret := make([][]string, 0, 1+len(ti)) // add 1 row for title

ret = append(ret, []string{"block_height", "key_lengt", "value_length", "length"})

for _, t := range ti {
ret = append(ret, []string{
strconv.FormatUint(t.BlockHeight, 10),
strconv.FormatUint(t.KeyLength, 10),
strconv.FormatUint(t.ValueLength, 10),
strconv.FormatUint(t.Length, 10),
})
}

return ret
}

func main() {
fname := os.Args[1]

rows, err := loadTestFile(fname)
if err != nil {
panic(err)
}

ti, err := getTraceInfo(rows)
if err != nil {
panic(err)
}

o, err := os.OpenFile("tracestats.csv", os.O_CREATE|os.O_APPEND|os.O_RDWR, 0755)
if err != nil {
panic(err)
}

defer func() {
if err := o.Close(); err != nil {
panic(err)
}
}()

w := csv.NewWriter(o)
for _, record := range ti.CSV() {
if err := w.Write(record); err != nil {
panic(err)
}
}
}

func getTraceInfo(traces []string) (traceInfos, error) {
ret := make([]traceInfo, 0, len(traces))

for _, t := range traces {
tr := tracelistener.TraceOperation{}
if err := json.Unmarshal([]byte(t), &tr); err != nil {
return nil, err
}

ret = append(ret, traceInfo{
BlockHeight: tr.Metadata.BlockHeight,
KeyLength: uint64(len(tr.Key)),
ValueLength: uint64(len(tr.Value)),
Length: uint64(len(t)),
})
}

return ret, nil
}

func loadTestFile(fname string) ([]string, error) {
file, err := os.Open(fname)
if err != nil {
return nil, fmt.Errorf("cannot open file %s, %w", fname, err)
}

scanner := bufio.NewScanner(file)
buf := make([]byte, 1000000) // a very high capacity
scanner.Buffer(buf, 1000000)

ret := []string{}
for scanner.Scan() {
ret = append(ret, scanner.Text())
}

if err := scanner.Err(); err != nil {
return nil, fmt.Errorf("scanning error, %w", err)
}

return ret, nil
}
176 changes: 176 additions & 0 deletions tracelistener/benchmark_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,176 @@
package tracelistener_test

import (
"bufio"
"context"
"fmt"
"io"
"os"
"syscall"
"testing"

"github.com/allinbits/tracelistener/tracelistener"
"github.com/containerd/fifo"
"go.uber.org/zap"
)

func setup(b *testing.B) (io.ReadWriteCloser, string) {
b.Helper()
f, err := os.CreateTemp("", "test_data")
if err != nil {
panic(err)
}

err = f.Close()
if err != nil {
panic(err)
}

dataChan := make(chan tracelistener.TraceOperation)
errChan := make(chan error)
l := zap.NewNop()
tw := tracelistener.TraceWatcher{
DataSourcePath: f.Name(),
WatchedOps: []tracelistener.Operation{
tracelistener.WriteOp,
tracelistener.DeleteOp,
},
DataChan: dataChan,
ErrorChan: errChan,
Logger: l.Sugar(),
}

go func() {
// drain data channel
for range dataChan {
}
}()

go func() {
tw.Watch()
}()

ff, err := fifo.OpenFifo(context.Background(), f.Name(), syscall.O_WRONLY, 0655)
if err != nil {
panic(err)
}

return ff, f.Name()
}

func runBenchmark(b *testing.B, amount int, kind string) {
fileWriter, fifoName := setup(b)
defer func() {
if err := fileWriter.Close(); err != nil {
panic(err)
}
}()

b.ResetTimer()

for i := 0; i < amount; i++ {
err := loadTest(b, i, fileWriter, kind)
if err != nil {
panic(err)
}
}

os.Remove(fifoName)
}

func BenchmarkTracelistenerRealTraces(b *testing.B) {
b.Log("reading test traces file...")
lines, err := loadTestFile(b)
if err != nil {
b.Fatal(err)
}
b.Log("finished reading test traces file!")

fileWriter, fifoName := setup(b)
defer func() {
if err := fileWriter.Close(); err != nil {
panic(err)
}
}()

b.ResetTimer()

for _, line := range lines {
fmt.Fprintf(fileWriter, line+"\n")
}

os.Remove(fifoName)
}

func BenchmarkTraceListenerKindWrite(b *testing.B) {
runBenchmark(b, b.N, "write")
}

func BenchmarkTraceListener100KKindWrite(b *testing.B) {
runBenchmark(b, 100000, "write")
}

func BenchmarkTraceListener1MKindWrite(b *testing.B) {
runBenchmark(b, 1000000, "write")
}

func BenchmarkTraceListener1MKindIterRange(b *testing.B) {
runBenchmark(b, 1000000, "IterRange")
}

func BenchmarkTraceListener10MKindWrite(b *testing.B) {
runBenchmark(b, 10000000, "write")
}

func loadTest(b *testing.B, height int, ff io.Writer, kind string) error {
b.Helper()

// trace := tracelistener.TraceOperation{
// Operation: string(tracelistener.WriteOp),
// Key: []byte{0x68, 0x65, 0x6c, 0x6c, 0x6f, 0xa},
// Value: []byte{0x68, 0x65, 0x6c, 0x6c, 0x6f, 0xa},
// BlockHeight: uint64(height),
// TxHash: "A5CF62609D62ADDE56816681B6191F5F0252D2800FC2C312EB91D962AB7A97CB",
// }
// data, err := json.Marshal(trace)
// if err != nil {
// return err
// }

// println(string(data))

s := `{"operation":"%s","key":"aGVsbG8K","value":"aGVsbG8K","block_height":158284,"tx_hash":"A5CF62609D62ADDE56816681B6191F5F0252D2800FC2C312EB91D962AB7A97CB","SuggestedProcessor":""}`

fmt.Fprintf(ff, s+"\n", kind)

return nil
}

func loadTestFile(b *testing.B) ([]string, error) {
b.Helper()

fname := os.Getenv("TRACELISTENER_BENCH_TRACEFILE")
if fname == "" {
return nil, fmt.Errorf("TRACELISTENER_BENCH_TRACEFILE environment variable not defined")
}

file, err := os.Open(fname)
if err != nil {
return nil, fmt.Errorf("cannot open file %s, %w", fname, err)
}

scanner := bufio.NewScanner(file)
buf := make([]byte, 1000000) // a very high capacity
scanner.Buffer(buf, 1000000)

ret := []string{}
for scanner.Scan() {
ret = append(ret, scanner.Text())
}

if err := scanner.Err(); err != nil {
return nil, fmt.Errorf("scanning error, %w", err)
}

return ret, nil
}
16 changes: 9 additions & 7 deletions tracelistener/bulk/bulkimport.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ import (

type Importer struct {
Path string
TraceWatcher tracelistener.TraceWatcher
TraceWatcher *tracelistener.TraceWatcher
Processor tracelistener.DataProcessor
Logger *zap.SugaredLogger
Database *database.Instance
Expand All @@ -41,7 +41,7 @@ func ImportableModulesList() []string {
return ml
}

func (i Importer) validateModulesList() error {
func (i *Importer) validateModulesList() error {
for _, m := range i.Modules {
if _, ok := tracelistener.SupportedSDKModuleList[tracelistener.SDKModuleName(m)]; !ok {
return fmt.Errorf("unknown bulk import module %s", m)
Expand Down Expand Up @@ -155,14 +155,16 @@ func (i *Importer) Do() error {

for ; ii.Valid(); ii.Next() {
to := tracelistener.TraceOperation{
Operation: tracelistener.WriteOp.String(),
Key: ii.Key(),
Value: ii.Value(),
BlockHeight: uint64(latestBlockHeight),
Operation: tracelistener.WriteOp.String(),
Key: ii.Key(),
Value: ii.Value(),
Metadata: tracelistener.TraceMetadata{
BlockHeight: uint64(latestBlockHeight),
},
SuggestedProcessor: tracelistener.SDKModuleName(key.Name()),
}

if err := i.TraceWatcher.ParseOperation(to); err != nil {
if err := i.TraceWatcher.ParseOperation(&to); err != nil {
return fmt.Errorf("cannot parse operation %v, %w", to, err)
}

Expand Down
16 changes: 10 additions & 6 deletions tracelistener/processor/auth_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -77,9 +77,11 @@ func TestAuthProcess(t *testing.T) {
Sequence: 11,
},
tracelistener.TraceOperation{
Operation: string(tracelistener.WriteOp),
Key: []byte("cosmos1xrnner9s783446"),
BlockHeight: 1,
Operation: string(tracelistener.WriteOp),
Key: []byte("cosmos1xrnner9s783446"),
Metadata: tracelistener.TraceMetadata{
BlockHeight: 1,
},
},
false,
1,
Expand All @@ -92,9 +94,11 @@ func TestAuthProcess(t *testing.T) {
Sequence: 11,
},
tracelistener.TraceOperation{
Operation: string(tracelistener.WriteOp),
Key: []byte("cosmos1xrnner9s783446"),
BlockHeight: 1,
Operation: string(tracelistener.WriteOp),
Key: []byte("cosmos1xrnner9s783446"),
Metadata: tracelistener.TraceMetadata{
BlockHeight: 1,
},
},
true,
0,
Expand Down
Loading