Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Feature/wip #4

Merged
merged 8 commits into from
Sep 5, 2023
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
45 changes: 38 additions & 7 deletions cmd/main.go
Original file line number Diff line number Diff line change
@@ -1,35 +1,66 @@
package main

import (
"fmt"
"os"
"os/signal"
"syscall"

"github.com/0xPolygonHermez/zkevm-data-streamer/datastreamer"
"github.com/0xPolygonHermez/zkevm-data-streamer/log"
)

func main() {
fmt.Println(">> App begin")
log.Info(">> App begin")

// Create server stream
ss, err := datastreamer.New(1337, "streamfile.bin")
s, err := datastreamer.New(1337, "streamfile.bin")
if err != nil {
os.Exit(1)
}
err = ss.Start()
err = s.Start()
if err != nil {
os.Exit(1)
log.Error(">> App error! Start")
return
}

// Create clients
go datastreamer.NewClient(1)

// Fake stream data
data := make([]byte, 32)
for i := 0; i < 32; i++ {
data[i] = byte(i)
}

// Start tx
err = s.StartStreamTx()
if err != nil {
log.Error(">> App error! StartStreamTx")
return
}

// Add stream entries
for i := 1; i <= 3; i++ {
entry, err := s.AddStreamEntry(1, data)
if err != nil {
log.Error(">> App error! AddStreamEntry:", err)
return
}
log.Info(">> App info. Added entry:", entry)
}

// Commit tx
err = s.CommitStreamTx()
if err != nil {
log.Error(">> App error! CommitStreamTx")
return
}

// Wait for ctl+c
fmt.Println(">> Press Control+C to finish...")
log.Info(">> Press Control+C to finish...")
interruptSignal := make(chan os.Signal, 1)
signal.Notify(interruptSignal, os.Interrupt, syscall.SIGTERM)
<-interruptSignal

fmt.Println(">> App end")
log.Info(">> App end")
}
23 changes: 11 additions & 12 deletions datastreamer/streamclient.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@ import (
"bufio"
"encoding/binary"
"errors"
"fmt"
"io"
"net"
"time"
Expand All @@ -21,12 +20,12 @@ func NewClient(i int) {
// Connect to server
conn, err := net.Dial("tcp", server)
if err != nil {
fmt.Println("**Error connecting to server:", server, err)
log.Error("**Error connecting to server:", server, err)
return
}

defer conn.Close()
fmt.Println("**Connected to server:", server)
log.Info("**Connected to server:", server)

// Send the command and stream type
err = writeFullUint64(uint64(i), conn)
Expand Down Expand Up @@ -58,15 +57,15 @@ func readFromServer(conn net.Conn) {
n, err := conn.Read(buffer)
if err != nil {
if err == io.EOF {
fmt.Printf("client %s: no more data\n", client)
log.Errorf("**client %s: no more data", client)
return
}
fmt.Printf("client %s:error reading from server:%s\n", client, err)
log.Errorf("**client %s:error reading from server:%s", client, err)
time.Sleep(2 * time.Second)
continue
}

fmt.Printf("client %s:message from server:[%s]\n", client, buffer[:n])
log.Infof("**client %s:message from server:[%s]", client, buffer[:n])
}
}

Expand All @@ -76,7 +75,7 @@ func writeFullUint64(value uint64, conn net.Conn) error {

_, err := conn.Write(buffer)
if err != nil {
fmt.Println("**Error sending to server:", err)
log.Error("**Error sending to server:", err)
return err
}
return nil
Expand All @@ -91,27 +90,27 @@ func readResultEntry(conn net.Conn) (ResultEntry, error) {
_, err := io.ReadFull(reader, buffer)
if err != nil {
if err == io.EOF {
fmt.Println("**Server close connection")
log.Warn("**Server close connection")
} else {
fmt.Println("**Error reading from server:", err)
log.Error("**Error reading from server:", err)
}
return e, err
}

// Read variable field (errStr)
length := binary.BigEndian.Uint32(buffer[1:5])
if length < 10 {
fmt.Println("**Error reading result entry")
log.Error("**Error reading result entry")
return e, errors.New("error reading result entry")
}

bufferAux := make([]byte, length-9)
_, err = io.ReadFull(reader, bufferAux)
if err != nil {
if err == io.EOF {
fmt.Println("**Server close connection")
log.Warn("**Server close connection")
} else {
fmt.Println("**Error reading from server:", err)
log.Error("**Error reading from server:", err)
}
return e, err
}
Expand Down
Loading