summaryrefslogtreecommitdiff
path: root/pkg
diff options
context:
space:
mode:
Diffstat (limited to 'pkg')
-rw-r--r--pkg/assert/assert.go2
-rw-r--r--pkg/snowflake/snowflake.go7
-rw-r--r--pkg/util/io.go57
-rw-r--r--pkg/util/io_test.go59
4 files changed, 4 insertions, 121 deletions
diff --git a/pkg/assert/assert.go b/pkg/assert/assert.go
index 6c4a4f5..7281a95 100644
--- a/pkg/assert/assert.go
+++ b/pkg/assert/assert.go
@@ -14,7 +14,7 @@ func NoError(err error, message string, a ...any) {
}
}
-func Unreachable(message string, a ...any) {
+func Never(message string, a ...any) {
log.Fatalf(message+"\n", a...)
}
diff --git a/pkg/snowflake/snowflake.go b/pkg/snowflake/snowflake.go
index 5b58cbb..bb4a07b 100644
--- a/pkg/snowflake/snowflake.go
+++ b/pkg/snowflake/snowflake.go
@@ -13,11 +13,10 @@ const (
// Epoch is set to the twitter snowflake epoch of Nov 04 2010 01:42:54 UTC in milliseconds
// TODO: change this to eko epoch when eko is production ready
Epoch int64 = 1288834974657
-
nodeBits = 10
stepBits = 12
- nodeMax = 1<<nodeBits - 1
- nodeMask = nodeMax << stepBits
+ NodeMax = 1<<nodeBits - 1
+ nodeMask = NodeMax << stepBits
stepMask = 1<<stepBits - 1
timeShift = nodeBits + stepBits
nodeShift = stepBits
@@ -35,7 +34,7 @@ type ID int64
func NewNode(node int64) *Node {
assert.Assert(nodeBits+stepBits <= 22, "node and step bits must add up to 22 or less")
- assert.Assert(0 <= node && node <= nodeMax, "node and step bits must add up to 22 or less")
+ assert.Assert(0 <= node && node <= NodeMax, "node and step bits must add up to 22 or less")
// Credit to https://github.com/bwmarrin/snowflake
currentTime := time.Now()
diff --git a/pkg/util/io.go b/pkg/util/io.go
deleted file mode 100644
index 5c489ed..0000000
--- a/pkg/util/io.go
+++ /dev/null
@@ -1,57 +0,0 @@
-package util
-
-import (
- "context"
- "io"
-)
-
-const bufferSize = 512
-
-type ChannelReader struct {
- Out <-chan []byte
- Err <-chan error
-}
-
-func NewChannelReader(ctx context.Context, reader io.Reader) ChannelReader {
- outCh := make(chan []byte)
- errCh := make(chan error)
-
- go func(in chan<- []byte, inErr chan<- error) {
- defer close(in)
- defer close(inErr)
- buffer := make([]byte, 2 * bufferSize)
- isD1 := true
- d1 := buffer[:bufferSize]
- d2 := buffer[bufferSize:]
- outer:
- for {
- select {
- case <-ctx.Done():
- inErr <- ctx.Err()
- break outer
- default:
- var data []byte
- if isD1 {
- data = d1
- } else {
- data = d2
- }
- isD1 = !isD1
-
- n, err := reader.Read(data)
- if err != nil {
- if err != io.EOF {
- inErr <- err
- }
- break outer
- }
- in <- data[:n]
- }
- }
- }(outCh, errCh)
-
- return ChannelReader{
- Out: outCh,
- Err: errCh,
- }
-}
diff --git a/pkg/util/io_test.go b/pkg/util/io_test.go
deleted file mode 100644
index f05b833..0000000
--- a/pkg/util/io_test.go
+++ /dev/null
@@ -1,59 +0,0 @@
-package util
-
-import (
- "bytes"
- "context"
- "testing"
- "time"
-)
-
-func TestChannelReader(t *testing.T) {
- ctx, cancel := context.WithTimeout(context.Background(), time.Second)
- defer cancel()
- b := []byte{1, 2, 3, 4, 5, 6, 7, 8}
- reader := NewChannelReader(ctx, bytes.NewReader(b))
-
- select {
- case data := <-reader.Out:
- if !bytes.Equal(b, data) {
- t.Errorf("%v != %v", b, data)
- }
- case err := <-reader.Err:
- t.Errorf("reading err: %v", err)
- }
-}
-
-func TestChannelReaderMultiPartRead(t *testing.T) {
- ctx, cancel := context.WithTimeout(context.Background(), time.Second)
- defer cancel()
-
- b := []byte("Testing FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting FramerTesting Framer")
- b1 := b[:512]
- b2 := b[512:]
- reader := NewChannelReader(ctx, bytes.NewReader(b))
-
- counter := 0
-outer:
- for {
- select {
- case data := <-reader.Out:
- if counter == 0 {
- if !bytes.Equal(b1, data) {
- t.Errorf("%v != %v", b, data)
- }
- } else {
- if !bytes.Equal(b2[:len(data)], data) {
- t.Errorf("%v != %v", b, data)
- }
- }
- counter++
- if counter == 2 {
- break outer
- }
-
- case err := <-reader.Err:
- t.Errorf("reading err: %v", err)
- break outer
- }
- }
-}