Skip to content
2 changes: 0 additions & 2 deletions .golangci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -34,5 +34,3 @@ issues:

run:
timeout: 5m
skip-dirs:
- vendor
9 changes: 7 additions & 2 deletions pkg/lib/algorithm/full_sync/sync.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package full_sync

import (
"fmt"

"github.com/String-Reconciliation-Ditributed-System/RCDS_GO/pkg/lib/genSync"
"github.com/String-Reconciliation-Ditributed-System/RCDS_GO/pkg/set"
"github.com/String-Reconciliation-Ditributed-System/RCDS_GO/pkg/util"
Expand Down Expand Up @@ -120,7 +121,9 @@ func (f *fullSync) SyncClient(ip string, port int) error {
return err
}
f.additionals.InsertKey(d)
f.AddElement(d)
if err = f.AddElement(d); err != nil {
return err
}
}
return nil
}
Expand Down Expand Up @@ -181,7 +184,9 @@ func (f *fullSync) SyncServer(ip string, port int) error {
if !f.FreezeLocal {
for elem := range *tempSet.Difference(f.Set) {
f.additionals.InsertKey(elem)
f.AddElement(elem)
if err = f.AddElement(elem); err != nil {
return err
}
}
} else {
logrus.Info("Server is freezing local set and skipping set update.")
Expand Down
6 changes: 3 additions & 3 deletions pkg/lib/algorithm/full_sync/sync_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -71,11 +71,11 @@ func TestNewFullSetSync(t *testing.T) {
var wg sync.WaitGroup
wg.Add(1)
go func() {
err := client.SyncServer("", 8080)
assert.NoError(t, err)
syncErr := client.SyncServer("", 9080)
assert.NoError(t, syncErr)
wg.Done()
}()
err = server.SyncClient("", 8080)
err = server.SyncClient("", 9080)
assert.NoError(t, err)
wg.Wait()

Expand Down
3 changes: 2 additions & 1 deletion pkg/lib/algorithm/hash_function_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,9 @@ package algorithm

import (
"crypto"
"github.com/stretchr/testify/assert"
"testing"

"github.com/stretchr/testify/assert"
)

func TestHashBytesWithCryptoFunc(t *testing.T) {
Expand Down
16 changes: 12 additions & 4 deletions pkg/lib/algorithm/iblt/sync.go
Original file line number Diff line number Diff line change
Expand Up @@ -65,15 +65,19 @@ func (i *ibltSync) AddElement(elem interface{}) error {
i.Set.Insert(key, elem)

for j := range i.resyncIBLTs {
i.resyncIBLTs[j].Insert(key)
if err := i.resyncIBLTs[j].Insert(key); err != nil {
return err
}
}
return i.Table.Insert(key)
} else {
i.Set.InsertKey(elem)
}
key := elem.([]byte)
for j := range i.resyncIBLTs {
i.resyncIBLTs[j].Insert(key)
if err := i.resyncIBLTs[j].Insert(key); err != nil {
return err
}
}
return i.Table.Insert(key)
}
Expand All @@ -86,14 +90,18 @@ func (i *ibltSync) DeleteElement(elem interface{}) error {
}
i.Set.Remove(key)
for j := range i.resyncIBLTs {
i.resyncIBLTs[j].Delete(key)
if err := i.resyncIBLTs[j].Delete(key); err != nil {
return err
}
}
return i.Table.Delete(key)
}
i.Set.Remove(elem)
key := elem.([]byte)
for j := range i.resyncIBLTs {
i.resyncIBLTs[j].Delete(key)
if err := i.resyncIBLTs[j].Delete(key); err != nil {
return err
}
}
return i.Table.Delete(key)
}
Expand Down
24 changes: 22 additions & 2 deletions pkg/lib/algorithm/rcds/backtracking.go
Original file line number Diff line number Diff line change
Expand Up @@ -76,7 +76,12 @@ func (s *hashShingleSet) interactiveBacktracking(previousEdge, currentEdge uint6
i = i - 1
}
previousEdge = currentEdge
for tail, count := range tails {

// Sort tail keys to ensure deterministic iteration order
tailKeys := getSortedTailKeys(tails)

for _, tail := range tailKeys {
count := tails[tail]
currentEdge = tail
// get the changed shingles from last layer.
tailChangeHistory[i] = tailChangeHistory[i-1]
Expand Down Expand Up @@ -152,7 +157,12 @@ func (s *hashShingleSet) reverseInteractiveBacktracking(hashArray []uint64) (*Cy
i = i - 1
}
previousEdge = currentEdge
for tail, count := range tails {

// Sort tail keys to ensure deterministic iteration order
tailKeys := getSortedTailKeys(tails)

for _, tail := range tailKeys {
count := tails[tail]
currentEdge = tail
// get the changed shingles from last layer.
tailChangeHistory[i] = tailChangeHistory[i-1]
Expand Down Expand Up @@ -192,6 +202,16 @@ func (s *hashShingleSet) getTailEdges(firstEdge uint64) shingleTailCount {
return nil
}

// getSortedTailKeys returns tail keys sorted in ascending order for deterministic iteration
func getSortedTailKeys(tails shingleTailCount) []uint64 {
tailKeys := make([]uint64, 0, len(tails))
for tail := range tails {
tailKeys = append(tailKeys, tail)
}
sort.Slice(tailKeys, func(i, j int) bool { return tailKeys[i] < tailKeys[j] })
return tailKeys
}

// sort sorts the hash shingle set and its tail counts.
func (s *hashShingleSet) sort() error {
size := s.Size()
Expand Down
8 changes: 6 additions & 2 deletions pkg/lib/algorithm/rcds/hashShingling.go
Original file line number Diff line number Diff line change
Expand Up @@ -95,7 +95,9 @@ func (s *hashShingleSet) addChunksToShingleSet(chunks *[]string) (*algorithm.Dic
if err != nil {
return nil, err
}
s.AddShingle(0, hash, 1)
if err = s.AddShingle(0, hash, 1); err != nil {
return nil, err
}

for i := 1; i < len(*chunks); i++ {
first, err := dict.AddToDict((*chunks)[i-1])
Expand All @@ -107,7 +109,9 @@ func (s *hashShingleSet) addChunksToShingleSet(chunks *[]string) (*algorithm.Dic
return nil, err
}
if _, err = s.getShingleCount(first, second); err != nil {
s.AddShingle(first, second, 1)
if err = s.AddShingle(first, second, 1); err != nil {
return nil, err
}
} else {
if _, err = s.addShingleCount(first, second, 1); err != nil {
return nil, err
Expand Down
42 changes: 34 additions & 8 deletions pkg/lib/genSync/connection.go
Original file line number Diff line number Diff line change
@@ -1,13 +1,17 @@
package genSync

import (
"context"
"fmt"
"github.com/String-Reconciliation-Ditributed-System/RCDS_GO/pkg/util"
"github.com/sirupsen/logrus"
"k8s.io/client-go/util/retry"
"net"
"strconv"
"strings"
"syscall"

"github.com/sirupsen/logrus"
"k8s.io/client-go/util/retry"

"github.com/String-Reconciliation-Ditributed-System/RCDS_GO/pkg/util"
)

type Connection interface {
Expand Down Expand Up @@ -40,7 +44,7 @@
// Original TCP buffer size for slower networks.
const bufferSize int = 65535

func NewTcpConnection(ipAddr string, port int) (Connection, error) {

Check failure on line 47 in pkg/lib/genSync/connection.go

View workflow job for this annotation

GitHub Actions / Lint

var-naming: func NewTcpConnection should be NewTCPConnection (revive)
if ipAddr == "" {
ipAddr = "localhost"
}
Expand Down Expand Up @@ -80,12 +84,29 @@

func (s *socketConnection) Listen() error {
var err error
s.listener, err = net.ListenTCP("tcp", s.tcpAddress)
logrus.Infof("listening on: %v", s.tcpAddress)

// Create ListenConfig with SO_REUSEADDR to allow quick port reuse
lc := &net.ListenConfig{
Control: func(_, _ string, c syscall.RawConn) error {
var sockOptErr error
if controlErr := c.Control(func(fd uintptr) {
// Set SO_REUSEADDR to allow quick port reuse
sockOptErr = syscall.SetsockoptInt(int(fd), syscall.SOL_SOCKET, syscall.SO_REUSEADDR, 1)
}); controlErr != nil {
return controlErr
}
return sockOptErr
},
}

listener, err := lc.Listen(context.Background(), "tcp", s.tcpAddress.String())
if err != nil {
return fmt.Errorf("failed to listen: %v", err)
}

s.listener = listener.(*net.TCPListener)
logrus.Infof("listening on: %v", s.tcpAddress)

s.connection, err = s.listener.AcceptTCP()
return err
}
Expand Down Expand Up @@ -201,13 +222,18 @@
}

func (s *socketConnection) Close() error {
if err := s.listener.Close(); err != nil {
logrus.Debugf("failed to close listener, %v", err)
if s.listener != nil {
if err := s.listener.Close(); err != nil {
logrus.Debugf("failed to close listener, %v", err)
}
}
if s.connection != nil {
return s.connection.Close()
}
return s.connection.Close()
return nil
}

func (s *socketConnection) GetIp() string {

Check failure on line 236 in pkg/lib/genSync/connection.go

View workflow job for this annotation

GitHub Actions / Lint

var-naming: method GetIp should be GetIP (revive)
return s.tcpAddress.IP.String()
}

Expand Down
5 changes: 3 additions & 2 deletions pkg/lib/genSync/conversion_test.go
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
package genSync

Check failure on line 1 in pkg/lib/genSync/conversion_test.go

View workflow job for this annotation

GitHub Actions / Lint

var-naming: don't use MixedCaps in package name; genSync should be gensync (revive)

import (
"crypto/rand"
"math"
"math/rand"
"testing"

"github.com/stretchr/testify/assert"
Expand Down Expand Up @@ -43,7 +43,8 @@
make([]byte, 512),
make([]byte, 1024),
} {
rand.Read(test)
_, err := rand.Read(test)
assert.NoError(t, err)
bytes, err := ToBigInt(test)
assert.NoError(t, err)
assert.Equal(t, test, bytes.ToBytes())
Expand Down
3 changes: 2 additions & 1 deletion pkg/set/set.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,9 @@ package set

import (
"fmt"
"github.com/String-Reconciliation-Ditributed-System/RCDS_GO/pkg/lib/algorithm"
"reflect"

"github.com/String-Reconciliation-Ditributed-System/RCDS_GO/pkg/lib/algorithm"
)

type Set map[interface{}]interface{}
Expand Down
5 changes: 3 additions & 2 deletions pkg/set/set_test.go
Original file line number Diff line number Diff line change
@@ -1,11 +1,12 @@
package set

import (
"k8s.io/apimachinery/pkg/util/rand"
"testing"

"k8s.io/apimachinery/pkg/util/rand"
)

func TestSet_Insert(t *testing.T) {
func TestSet_Insert(_ *testing.T) {
s := New()
s.InsertKey([]byte(rand.String(10)))
}
3 changes: 2 additions & 1 deletion pkg/util/auxiliary_test.go
Original file line number Diff line number Diff line change
@@ -1,9 +1,10 @@
package util

import (
"github.com/stretchr/testify/assert"
"math"
"testing"

"github.com/stretchr/testify/assert"
)

func TestBytesAndIntConversion(t *testing.T) {
Expand Down
Loading