diff --git a/CHANGELOG.md b/CHANGELOG.md index a42e7939..10df955a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,10 @@ *??? ??, ????* +### CHANGES + +- `pkg/client/message,pkg/network/message`: move to pkg/message/layer2,pkg/message/layer1 + ## v1.7.9 diff --git a/images/go-peer_layer1_net_message.jpg b/images/go-peer_layer1_message.jpg similarity index 100% rename from images/go-peer_layer1_net_message.jpg rename to images/go-peer_layer1_message.jpg diff --git a/pkg/anonymity/adapters/adapter.go b/pkg/anonymity/adapters/adapter.go index dff0dfc0..d4e4fdcd 100644 --- a/pkg/anonymity/adapters/adapter.go +++ b/pkg/anonymity/adapters/adapter.go @@ -3,7 +3,7 @@ package adapters import ( "context" - net_message "github.com/number571/go-peer/pkg/network/message" + "github.com/number571/go-peer/pkg/message/layer1" ) var ( @@ -11,8 +11,8 @@ var ( ) type ( - iProducerF func(context.Context, net_message.IMessage) error - iConsumerF func(context.Context) (net_message.IMessage, error) + iProducerF func(context.Context, layer1.IMessage) error + iConsumerF func(context.Context) (layer1.IMessage, error) ) type sAdapter struct { @@ -27,10 +27,10 @@ func NewAdapterByFuncs(pProduce iProducerF, pConsume iConsumerF) IAdapter { } } -func (p *sAdapter) Produce(pCtx context.Context, pMsg net_message.IMessage) error { +func (p *sAdapter) Produce(pCtx context.Context, pMsg layer1.IMessage) error { return p.fProduce(pCtx, pMsg) } -func (p *sAdapter) Consume(pCtx context.Context) (net_message.IMessage, error) { +func (p *sAdapter) Consume(pCtx context.Context) (layer1.IMessage, error) { return p.fConsume(pCtx) } diff --git a/pkg/anonymity/adapters/adapters_test.go b/pkg/anonymity/adapters/adapters_test.go index ed035b65..0edaa102 100644 --- a/pkg/anonymity/adapters/adapters_test.go +++ b/pkg/anonymity/adapters/adapters_test.go @@ -5,7 +5,7 @@ import ( "context" "testing" - "github.com/number571/go-peer/pkg/network/message" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/payload" ) @@ -14,22 +14,22 @@ const ( ) func TestAdapter(t *testing.T) { - msgChan := make(chan message.IMessage, 1) + msgChan := make(chan layer1.IMessage, 1) adapter := NewAdapterByFuncs( - func(_ context.Context, msg message.IMessage) error { + func(_ context.Context, msg layer1.IMessage) error { msgChan <- msg return nil }, - func(_ context.Context) (message.IMessage, error) { + func(_ context.Context) (layer1.IMessage, error) { return <-msgChan, nil }, ) ctx := context.Background() - err := adapter.Produce(ctx, message.NewMessage( - message.NewConstructSettings(&message.SConstructSettings{ - FSettings: message.NewSettings(&message.SSettings{}), + err := adapter.Produce(ctx, layer1.NewMessage( + layer1.NewConstructSettings(&layer1.SConstructSettings{ + FSettings: layer1.NewSettings(&layer1.SSettings{}), }), payload.NewPayload32(0x01, []byte(tcMessage)), )) diff --git a/pkg/anonymity/adapters/types.go b/pkg/anonymity/adapters/types.go index 5ebe9ff0..617f98e5 100644 --- a/pkg/anonymity/adapters/types.go +++ b/pkg/anonymity/adapters/types.go @@ -3,7 +3,7 @@ package adapters import ( "context" - net_message "github.com/number571/go-peer/pkg/network/message" + "github.com/number571/go-peer/pkg/message/layer1" ) type IAdapter interface { @@ -12,9 +12,9 @@ type IAdapter interface { } type IProducer interface { - Produce(context.Context, net_message.IMessage) error + Produce(context.Context, layer1.IMessage) error } type IConsumer interface { - Consume(context.Context) (net_message.IMessage, error) + Consume(context.Context) (layer1.IMessage, error) } diff --git a/pkg/anonymity/anonymity.go b/pkg/anonymity/anonymity.go index e0738270..38970846 100644 --- a/pkg/anonymity/anonymity.go +++ b/pkg/anonymity/anonymity.go @@ -9,17 +9,17 @@ import ( "github.com/number571/go-peer/pkg/anonymity/adapters" "github.com/number571/go-peer/pkg/anonymity/queue" - "github.com/number571/go-peer/pkg/client/message" "github.com/number571/go-peer/pkg/crypto/asymmetric" "github.com/number571/go-peer/pkg/crypto/hashing" "github.com/number571/go-peer/pkg/crypto/random" "github.com/number571/go-peer/pkg/logger" + "github.com/number571/go-peer/pkg/message/layer1" + "github.com/number571/go-peer/pkg/message/layer2" "github.com/number571/go-peer/pkg/payload" "github.com/number571/go-peer/pkg/state" "github.com/number571/go-peer/pkg/storage/database" anon_logger "github.com/number571/go-peer/pkg/anonymity/logger" - net_message "github.com/number571/go-peer/pkg/network/message" ) var ( @@ -225,7 +225,7 @@ func (p *sNode) recvResponse(pCtx context.Context, pActionKey string) ([]byte, e } } -func (p *sNode) consumeMessage(pCtx context.Context, pNetMsg net_message.IMessage) error { +func (p *sNode) consumeMessage(pCtx context.Context, pNetMsg layer1.IMessage) error { logBuilder := anon_logger.NewLogBuilder(p.fSettings.GetServiceName()) // update logger state @@ -235,7 +235,7 @@ func (p *sNode) consumeMessage(pCtx context.Context, pNetMsg net_message.IMessag encMsg := pNetMsg.GetPayload().GetBody() // load encrypted message without decryption try - if _, err := message.LoadMessage(client.GetMessageSize(), encMsg); err != nil { + if _, err := layer2.LoadMessage(client.GetMessageSize(), encMsg); err != nil { // problem from sender's side p.fLogger.PushWarn(logBuilder.WithType(anon_logger.CLogWarnMessageNull)) return errors.Join(ErrLoadMessage, err) @@ -377,7 +377,7 @@ func (p *sNode) enqueuePayload( return nil } -func (p *sNode) enrichLogger(pLogBuilder anon_logger.ILogBuilder, pNetMsg net_message.IMessage) anon_logger.ILogBuilder { +func (p *sNode) enrichLogger(pLogBuilder anon_logger.ILogBuilder, pNetMsg layer1.IMessage) anon_logger.ILogBuilder { var ( size = len(pNetMsg.ToBytes()) hash = pNetMsg.GetHash() @@ -391,7 +391,7 @@ func (p *sNode) enrichLogger(pLogBuilder anon_logger.ILogBuilder, pNetMsg net_me func (p *sNode) produceMessage( pCtx context.Context, - pNetMsg net_message.IMessage, + pNetMsg layer1.IMessage, ) (bool, error) { serviceName := p.fSettings.GetServiceName() diff --git a/pkg/anonymity/anonymity_test.go b/pkg/anonymity/anonymity_test.go index 9d258e19..86db95c7 100644 --- a/pkg/anonymity/anonymity_test.go +++ b/pkg/anonymity/anonymity_test.go @@ -19,6 +19,7 @@ import ( "github.com/number571/go-peer/pkg/crypto/random" "github.com/number571/go-peer/pkg/encoding" "github.com/number571/go-peer/pkg/logger" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/network" "github.com/number571/go-peer/pkg/payload" "github.com/number571/go-peer/pkg/storage/cache" @@ -27,7 +28,6 @@ import ( anon_logger "github.com/number571/go-peer/pkg/anonymity/logger" "github.com/number571/go-peer/pkg/network/conn" - net_message "github.com/number571/go-peer/pkg/network/message" ) const ( @@ -345,8 +345,8 @@ func TestHandleWrapper(t *testing.T) { pubKey := privKey.GetPubKey() node.GetMapPubKeys().SetPubKey(privKey.GetPubKey()) - sett := net_message.NewConstructSettings(&net_message.SConstructSettings{ - FSettings: net_message.NewSettings(&net_message.SSettings{}), + sett := layer1.NewConstructSettings(&layer1.SConstructSettings{ + FSettings: layer1.NewSettings(&layer1.SSettings{}), }) msg, err := client.EncryptMessage( @@ -483,8 +483,8 @@ func TestStoreHashWithBroadcastMessage(t *testing.T) { return } - sett := net_message.NewConstructSettings(&net_message.SConstructSettings{ - FSettings: net_message.NewSettings(&net_message.SSettings{}), + sett := layer1.NewConstructSettings(&layer1.SConstructSettings{ + FSettings: layer1.NewSettings(&layer1.SSettings{}), }) netMsg := node.testNewNetworkMessage(sett, msg) @@ -675,7 +675,7 @@ func (p *stLogging) HasErro() bool { */ func testRunNodeWithDB(ctx context.Context, timeWait time.Duration, addr string, db database.IKVDatabase) (INode, network.INode) { - msgChan := make(chan net_message.IMessage) + msgChan := make(chan layer1.IMessage) parallel := uint64(1) networkMask := uint32(1) limitVoidSize := uint64(10_000) @@ -686,7 +686,7 @@ func testRunNodeWithDB(ctx context.Context, timeWait time.Duration, addr string, FReadTimeout: timeWait, FWriteTimeout: timeWait, FConnSettings: conn.NewSettings(&conn.SSettings{ - FMessageSettings: net_message.NewSettings(&net_message.SSettings{ + FMessageSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: tcWorkSize, }), FLimitMessageSizeBytes: tcMsgSize + limitVoidSize, @@ -697,7 +697,7 @@ func testRunNodeWithDB(ctx context.Context, timeWait time.Duration, addr string, }), }), cache.NewLRUCache(1024), - ).HandleFunc(networkMask, func(_ context.Context, _ network.INode, _ conn.IConn, msg net_message.IMessage) error { + ).HandleFunc(networkMask, func(_ context.Context, _ network.INode, _ conn.IConn, msg layer1.IMessage) error { msgChan <- msg return nil }) @@ -712,10 +712,10 @@ func testRunNodeWithDB(ctx context.Context, timeWait time.Duration, addr string, func(_ logger.ILogArg) string { return "" }, ), adapters.NewAdapterByFuncs( - func(ctx context.Context, msg net_message.IMessage) error { + func(ctx context.Context, msg layer1.IMessage) error { return networkNode.BroadcastMessage(ctx, msg) }, - func(ctx context.Context) (net_message.IMessage, error) { + func(ctx context.Context) (layer1.IMessage, error) { select { case <-ctx.Done(): return nil, ctx.Err() @@ -728,8 +728,8 @@ func testRunNodeWithDB(ctx context.Context, timeWait time.Duration, addr string, db, queue.NewQBProblemProcessor( queue.NewSettings(&queue.SSettings{ - FMessageConstructSettings: net_message.NewConstructSettings(&net_message.SConstructSettings{ - FSettings: net_message.NewSettings(&net_message.SSettings{ + FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ + FSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: tcWorkSize, }), FParallel: parallel, @@ -770,8 +770,8 @@ func testDeleteDB(typeDB int) { } } -func (p *sNode) testNewNetworkMessage(pSett net_message.IConstructSettings, pMsgBytes []byte) net_message.IMessage { - return net_message.NewMessage( +func (p *sNode) testNewNetworkMessage(pSett layer1.IConstructSettings, pMsgBytes []byte) layer1.IMessage { + return layer1.NewMessage( pSett, payload.NewPayload32( p.fQBProcessor.GetSettings().GetNetworkMask(), diff --git a/pkg/anonymity/examples/echo/construct.go b/pkg/anonymity/examples/echo/construct.go index 0191c0e7..0a996011 100644 --- a/pkg/anonymity/examples/echo/construct.go +++ b/pkg/anonymity/examples/echo/construct.go @@ -13,9 +13,9 @@ import ( "github.com/number571/go-peer/pkg/client" "github.com/number571/go-peer/pkg/crypto/asymmetric" "github.com/number571/go-peer/pkg/logger" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/network" "github.com/number571/go-peer/pkg/network/conn" - net_message "github.com/number571/go-peer/pkg/network/message" "github.com/number571/go-peer/pkg/storage/cache" "github.com/number571/go-peer/pkg/storage/database" ) @@ -29,7 +29,7 @@ const ( ) func newNode(serviceName, address string) (network.INode, anonymity.INode) { - msgChan := make(chan net_message.IMessage) + msgChan := make(chan layer1.IMessage) networkNode := network.NewNode( network.NewSettings(&network.SSettings{ FAddress: address, @@ -37,7 +37,7 @@ func newNode(serviceName, address string) (network.INode, anonymity.INode) { FReadTimeout: time.Minute, FWriteTimeout: time.Minute, FConnSettings: conn.NewSettings(&conn.SSettings{ - FMessageSettings: net_message.NewSettings(&net_message.SSettings{ + FMessageSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: workSize, }), FLimitMessageSizeBytes: msgSize, @@ -50,7 +50,7 @@ func newNode(serviceName, address string) (network.INode, anonymity.INode) { cache.NewLRUCache(1024), ).HandleFunc( networkMask, - func(ctx context.Context, _ network.INode, _ conn.IConn, msg net_message.IMessage) error { + func(ctx context.Context, _ network.INode, _ conn.IConn, msg layer1.IMessage) error { msgChan <- msg return nil }, @@ -82,10 +82,10 @@ func newNode(serviceName, address string) (network.INode, anonymity.INode) { }, ), adapters.NewAdapterByFuncs( - func(ctx context.Context, msg net_message.IMessage) error { + func(ctx context.Context, msg layer1.IMessage) error { return networkNode.BroadcastMessage(ctx, msg) }, - func(ctx context.Context) (net_message.IMessage, error) { + func(ctx context.Context) (layer1.IMessage, error) { select { case <-ctx.Done(): return nil, ctx.Err() @@ -104,8 +104,8 @@ func newNode(serviceName, address string) (network.INode, anonymity.INode) { }(), queue.NewQBProblemProcessor( queue.NewSettings(&queue.SSettings{ - FMessageConstructSettings: net_message.NewConstructSettings(&net_message.SConstructSettings{ - FSettings: net_message.NewSettings(&net_message.SSettings{ + FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ + FSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: workSize, }), }), diff --git a/pkg/anonymity/examples/ping-pong/construct.go b/pkg/anonymity/examples/ping-pong/construct.go index 0191c0e7..0a996011 100644 --- a/pkg/anonymity/examples/ping-pong/construct.go +++ b/pkg/anonymity/examples/ping-pong/construct.go @@ -13,9 +13,9 @@ import ( "github.com/number571/go-peer/pkg/client" "github.com/number571/go-peer/pkg/crypto/asymmetric" "github.com/number571/go-peer/pkg/logger" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/network" "github.com/number571/go-peer/pkg/network/conn" - net_message "github.com/number571/go-peer/pkg/network/message" "github.com/number571/go-peer/pkg/storage/cache" "github.com/number571/go-peer/pkg/storage/database" ) @@ -29,7 +29,7 @@ const ( ) func newNode(serviceName, address string) (network.INode, anonymity.INode) { - msgChan := make(chan net_message.IMessage) + msgChan := make(chan layer1.IMessage) networkNode := network.NewNode( network.NewSettings(&network.SSettings{ FAddress: address, @@ -37,7 +37,7 @@ func newNode(serviceName, address string) (network.INode, anonymity.INode) { FReadTimeout: time.Minute, FWriteTimeout: time.Minute, FConnSettings: conn.NewSettings(&conn.SSettings{ - FMessageSettings: net_message.NewSettings(&net_message.SSettings{ + FMessageSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: workSize, }), FLimitMessageSizeBytes: msgSize, @@ -50,7 +50,7 @@ func newNode(serviceName, address string) (network.INode, anonymity.INode) { cache.NewLRUCache(1024), ).HandleFunc( networkMask, - func(ctx context.Context, _ network.INode, _ conn.IConn, msg net_message.IMessage) error { + func(ctx context.Context, _ network.INode, _ conn.IConn, msg layer1.IMessage) error { msgChan <- msg return nil }, @@ -82,10 +82,10 @@ func newNode(serviceName, address string) (network.INode, anonymity.INode) { }, ), adapters.NewAdapterByFuncs( - func(ctx context.Context, msg net_message.IMessage) error { + func(ctx context.Context, msg layer1.IMessage) error { return networkNode.BroadcastMessage(ctx, msg) }, - func(ctx context.Context) (net_message.IMessage, error) { + func(ctx context.Context) (layer1.IMessage, error) { select { case <-ctx.Done(): return nil, ctx.Err() @@ -104,8 +104,8 @@ func newNode(serviceName, address string) (network.INode, anonymity.INode) { }(), queue.NewQBProblemProcessor( queue.NewSettings(&queue.SSettings{ - FMessageConstructSettings: net_message.NewConstructSettings(&net_message.SConstructSettings{ - FSettings: net_message.NewSettings(&net_message.SSettings{ + FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ + FSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: workSize, }), }), diff --git a/pkg/anonymity/queue/examples/queue/main.go b/pkg/anonymity/queue/examples/queue/main.go index 13aca472..22ebb0c5 100644 --- a/pkg/anonymity/queue/examples/queue/main.go +++ b/pkg/anonymity/queue/examples/queue/main.go @@ -9,9 +9,9 @@ import ( "github.com/number571/go-peer/pkg/anonymity/queue" "github.com/number571/go-peer/pkg/client" - "github.com/number571/go-peer/pkg/client/message" "github.com/number571/go-peer/pkg/crypto/asymmetric" - net_message "github.com/number571/go-peer/pkg/network/message" + "github.com/number571/go-peer/pkg/message/layer1" + "github.com/number571/go-peer/pkg/message/layer2" "github.com/number571/go-peer/pkg/payload" ) @@ -23,8 +23,8 @@ const ( func main() { q := queue.NewQBProblemProcessor( queue.NewSettings(&queue.SSettings{ - FMessageConstructSettings: net_message.NewConstructSettings(&net_message.SConstructSettings{ - FSettings: net_message.NewSettings(&net_message.SSettings{}), + FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ + FSettings: layer1.NewSettings(&layer1.SSettings{}), }), FQueuePeriod: time.Second, FQueuePoolCap: [2]uint64{1 << 5, 1 << 5}, @@ -60,7 +60,7 @@ func main() { if netMsg == nil { panic("net message is nil") } - msg, err := message.LoadMessage(q.GetClient().GetMessageSize(), netMsg.GetPayload().GetBody()) + msg, err := layer2.LoadMessage(q.GetClient().GetMessageSize(), netMsg.GetPayload().GetBody()) if err != nil { panic(err) } diff --git a/pkg/anonymity/queue/queue.go b/pkg/anonymity/queue/queue.go index 7bca910e..a030c4f2 100644 --- a/pkg/anonymity/queue/queue.go +++ b/pkg/anonymity/queue/queue.go @@ -11,10 +11,9 @@ import ( "github.com/number571/go-peer/pkg/crypto/asymmetric" "github.com/number571/go-peer/pkg/crypto/random" "github.com/number571/go-peer/pkg/encoding" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/payload" "github.com/number571/go-peer/pkg/state" - - net_message "github.com/number571/go-peer/pkg/network/message" ) var ( @@ -34,14 +33,14 @@ type sQBProblemProcessor struct { type sMainPool struct { fMutex sync.Mutex fCount int64 // atomic variable - fQueue chan net_message.IMessage + fQueue chan layer1.IMessage fRawQueue map[uint64]chan []byte fConsumers map[string]uint64 } type sRandPool struct { fCount int64 // atomic variable - fQueue chan net_message.IMessage + fQueue chan layer1.IMessage fReceiver asymmetric.IPubKey } @@ -53,7 +52,7 @@ func NewQBProblemProcessor(pSettings ISettings, pClient client.IClient) IQBProbl fSettings: pSettings, fClient: pClient, fMainPool: &sMainPool{ - fQueue: make(chan net_message.IMessage, queuePoolCap[0]*consumersCap), + fQueue: make(chan layer1.IMessage, queuePoolCap[0]*consumersCap), fConsumers: make(map[string]uint64, 128), fRawQueue: func() map[uint64]chan []byte { m := make(map[uint64]chan []byte, consumersCap) @@ -64,7 +63,7 @@ func NewQBProblemProcessor(pSettings ISettings, pClient client.IClient) IQBProbl }(), }, fRandPool: &sRandPool{ - fQueue: make(chan net_message.IMessage, queuePoolCap[1]), + fQueue: make(chan layer1.IMessage, queuePoolCap[1]), fReceiver: asymmetric.NewPrivKey().GetPubKey(), }, } @@ -159,7 +158,7 @@ func (p *sQBProblemProcessor) EnqueueMessage(pPubKey asymmetric.IPubKey, pBytes return nil } -func (p *sQBProblemProcessor) DequeueMessage(pCtx context.Context) net_message.IMessage { +func (p *sQBProblemProcessor) DequeueMessage(pCtx context.Context) layer1.IMessage { for { select { case <-pCtx.Done(): @@ -208,10 +207,10 @@ func (p *sQBProblemProcessor) fillRandPool(pCtx context.Context) error { return p.pushMessage(pCtx, p.fRandPool.fQueue, msg) } -func (p *sQBProblemProcessor) pushMessage(pCtx context.Context, pQueue chan<- net_message.IMessage, pMsg []byte) error { - chNetMsg := make(chan net_message.IMessage) +func (p *sQBProblemProcessor) pushMessage(pCtx context.Context, pQueue chan<- layer1.IMessage, pMsg []byte) error { + chNetMsg := make(chan layer1.IMessage) go func() { - chNetMsg <- net_message.NewMessage( + chNetMsg <- layer1.NewMessage( p.fSettings.GetMessageConstructSettings(), payload.NewPayload32(p.fSettings.GetNetworkMask(), pMsg), ) diff --git a/pkg/anonymity/queue/queue_test.go b/pkg/anonymity/queue/queue_test.go index 12f93e97..8c40123c 100644 --- a/pkg/anonymity/queue/queue_test.go +++ b/pkg/anonymity/queue/queue_test.go @@ -11,7 +11,7 @@ import ( "github.com/number571/go-peer/pkg/client" "github.com/number571/go-peer/pkg/crypto/asymmetric" - net_message "github.com/number571/go-peer/pkg/network/message" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/payload" testutils "github.com/number571/go-peer/test/utils" ) @@ -51,16 +51,16 @@ func testSettings(t *testing.T, n int) { switch n { case 0: _ = NewSettings(&SSettings{ - FMessageConstructSettings: net_message.NewConstructSettings(&net_message.SConstructSettings{ - FSettings: net_message.NewSettings(&net_message.SSettings{}), + FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ + FSettings: layer1.NewSettings(&layer1.SSettings{}), }), FQueuePeriod: 500 * time.Millisecond, FConsumersCap: 1, }) case 1: _ = NewSettings(&SSettings{ - FMessageConstructSettings: net_message.NewConstructSettings(&net_message.SConstructSettings{ - FSettings: net_message.NewSettings(&net_message.SSettings{}), + FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ + FSettings: layer1.NewSettings(&layer1.SSettings{}), }), FQueuePoolCap: [2]uint64{tcQueueCap, tcQueueCap}, FConsumersCap: 1, @@ -73,8 +73,8 @@ func testSettings(t *testing.T, n int) { }) case 3: _ = NewSettings(&SSettings{ - FMessageConstructSettings: net_message.NewConstructSettings(&net_message.SConstructSettings{ - FSettings: net_message.NewSettings(&net_message.SSettings{}), + FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ + FSettings: layer1.NewSettings(&layer1.SSettings{}), }), FQueuePoolCap: [2]uint64{tcQueueCap, tcQueueCap}, FQueuePeriod: 500 * time.Millisecond, @@ -91,8 +91,8 @@ func TestRunStopQueue(t *testing.T) { ) queue := NewQBProblemProcessor( NewSettings(&SSettings{ - FMessageConstructSettings: net_message.NewConstructSettings(&net_message.SConstructSettings{ - FSettings: net_message.NewSettings(&net_message.SSettings{}), + FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ + FSettings: layer1.NewSettings(&layer1.SSettings{}), }), FQueuePoolCap: [2]uint64{tcQueueCap, tcQueueCap}, FQueuePeriod: 100 * time.Millisecond, @@ -158,8 +158,8 @@ func TestQueue(t *testing.T) { queue := NewQBProblemProcessor( NewSettings(&SSettings{ - FMessageConstructSettings: net_message.NewConstructSettings(&net_message.SConstructSettings{ - FSettings: net_message.NewSettings(&net_message.SSettings{ + FMessageConstructSettings: layer1.NewConstructSettings(&layer1.SConstructSettings{ + FSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: 10, }), }), @@ -210,7 +210,7 @@ func testQueue(queue IQBProblemProcessor) error { time.Sleep(300 * time.Millisecond) // auto fill queue enabled only if QB=true - msgs := make([]net_message.IMessage, 0, 3) + msgs := make([]layer1.IMessage, 0, 3) for i := 0; i < 3; i++ { msgs = append(msgs, queue.DequeueMessage(ctx)) } diff --git a/pkg/anonymity/queue/settings.go b/pkg/anonymity/queue/settings.go index a9c78ae8..93765b32 100644 --- a/pkg/anonymity/queue/settings.go +++ b/pkg/anonymity/queue/settings.go @@ -3,7 +3,7 @@ package queue import ( "time" - net_message "github.com/number571/go-peer/pkg/network/message" + "github.com/number571/go-peer/pkg/message/layer1" ) var ( @@ -12,7 +12,7 @@ var ( type SSettings sSettings type sSettings struct { - FMessageConstructSettings net_message.IConstructSettings + FMessageConstructSettings layer1.IConstructSettings FNetworkMask uint32 FConsumersCap uint64 FQueuePoolCap [2]uint64 @@ -46,7 +46,7 @@ func (p *sSettings) mustNotNull() ISettings { return p } -func (p *sSettings) GetMessageConstructSettings() net_message.IConstructSettings { +func (p *sSettings) GetMessageConstructSettings() layer1.IConstructSettings { return p.FMessageConstructSettings } diff --git a/pkg/anonymity/queue/types.go b/pkg/anonymity/queue/types.go index 58a3c064..89cf3cf2 100644 --- a/pkg/anonymity/queue/types.go +++ b/pkg/anonymity/queue/types.go @@ -6,9 +6,8 @@ import ( "github.com/number571/go-peer/pkg/client" "github.com/number571/go-peer/pkg/crypto/asymmetric" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/types" - - net_message "github.com/number571/go-peer/pkg/network/message" ) type IQBProblemProcessor interface { @@ -18,11 +17,11 @@ type IQBProblemProcessor interface { GetClient() client.IClient EnqueueMessage(asymmetric.IPubKey, []byte) error - DequeueMessage(context.Context) net_message.IMessage + DequeueMessage(context.Context) layer1.IMessage } type ISettings interface { - GetMessageConstructSettings() net_message.IConstructSettings + GetMessageConstructSettings() layer1.IConstructSettings GetNetworkMask() uint32 GetConsumersCap() uint64 GetQueuePeriod() time.Duration diff --git a/pkg/client/client.go b/pkg/client/client.go index a8f6a44f..a289d19f 100644 --- a/pkg/client/client.go +++ b/pkg/client/client.go @@ -3,11 +3,11 @@ package client import ( "bytes" - "github.com/number571/go-peer/pkg/client/message" "github.com/number571/go-peer/pkg/crypto/asymmetric" "github.com/number571/go-peer/pkg/crypto/hashing" "github.com/number571/go-peer/pkg/crypto/random" "github.com/number571/go-peer/pkg/crypto/symmetric" + "github.com/number571/go-peer/pkg/message/layer2" "github.com/number571/go-peer/pkg/payload/joiner" ) @@ -104,7 +104,7 @@ func (p *sClient) encryptWithParams( } cipher := symmetric.NewCipher(sk) - return message.NewMessage( + return layer2.NewMessage( ct, cipher.EncryptBytes(joiner.NewBytesJoiner32([][]byte{ pkey.GetHasher().ToBytes(), @@ -119,7 +119,7 @@ func (p *sClient) encryptWithParams( // Decrypt message with private key of receiver. // No one else except the sender will be able to decrypt the message. func (p *sClient) DecryptMessage(pMapPubKeys asymmetric.IMapPubKeys, pMsg []byte) (asymmetric.IPubKey, []byte, error) { - msg, err := message.LoadMessage(p.fMessageSize, pMsg) + msg, err := layer2.LoadMessage(p.fMessageSize, pMsg) if err != nil { return nil, nil, ErrInitCheckMessage } diff --git a/pkg/client/client_test.go b/pkg/client/client_test.go index 2511daff..41b41ffe 100644 --- a/pkg/client/client_test.go +++ b/pkg/client/client_test.go @@ -6,12 +6,12 @@ import ( "errors" "testing" - "github.com/number571/go-peer/pkg/client/message" "github.com/number571/go-peer/pkg/crypto/asymmetric" "github.com/number571/go-peer/pkg/crypto/hashing" "github.com/number571/go-peer/pkg/crypto/random" "github.com/number571/go-peer/pkg/crypto/symmetric" "github.com/number571/go-peer/pkg/encoding" + "github.com/number571/go-peer/pkg/message/layer2" "github.com/number571/go-peer/pkg/payload/joiner" ) @@ -268,7 +268,7 @@ func tcEncryptWithParamsInvalidPKID( } cipher := symmetric.NewCipher(sk) - return message.NewMessage( + return layer2.NewMessage( ct, cipher.EncryptBytes(joiner.NewBytesJoiner32([][]byte{ []byte("123"), diff --git a/pkg/client/doc.go b/pkg/client/doc.go index 6158f8d9..4f75f105 100644 --- a/pkg/client/doc.go +++ b/pkg/client/doc.go @@ -42,6 +42,6 @@ IF ≠, than protocol is interrupted. More information in article: https://github.com/number571/go-peer/blob/master/docs/monolithic_cryptographic_protocol.pdf - Scheme: https://github.com/number571/go-peer/blob/master/images/go-peer_layer2_message.jpg + Scheme: https://github.com/number571/go-peer/blob/master/images/go-peer_layer1_message.jpg */ package client diff --git a/pkg/client/message/doc.go b/pkg/client/message/doc.go deleted file mode 100644 index aa844c28..00000000 --- a/pkg/client/message/doc.go +++ /dev/null @@ -1,5 +0,0 @@ -// Package message used as a storage and loading of encrypted messages. -// -// The package allows initializing verification of the correctness of -// the message by the hash length and proof of work. -package message diff --git a/pkg/network/message/doc.go b/pkg/message/layer1/doc.go similarity index 91% rename from pkg/network/message/doc.go rename to pkg/message/layer1/doc.go index be2fba45..de5691e5 100644 --- a/pkg/network/message/doc.go +++ b/pkg/message/layer1/doc.go @@ -14,6 +14,6 @@ P - proof of work E - encrypt - Scheme: https://github.com/number571/go-peer/blob/master/images/go-peer_layer1_net_message.jpg + Scheme: https://github.com/number571/go-peer/blob/master/images/go-peer_layer1_message.jpg */ -package message +package layer1 diff --git a/pkg/network/message/errors.go b/pkg/message/layer1/errors.go similarity index 92% rename from pkg/network/message/errors.go rename to pkg/message/layer1/errors.go index bc29dcad..582a57c3 100644 --- a/pkg/network/message/errors.go +++ b/pkg/message/layer1/errors.go @@ -1,7 +1,7 @@ -package message +package layer1 const ( - errPrefix = "pkg/network/message = " + errPrefix = "pkg/message/layer1 = " ) type SMessageError struct { diff --git a/pkg/network/message/message.go b/pkg/message/layer1/message.go similarity index 99% rename from pkg/network/message/message.go rename to pkg/message/layer1/message.go index bf8a4889..fcc5cca2 100644 --- a/pkg/network/message/message.go +++ b/pkg/message/layer1/message.go @@ -1,4 +1,4 @@ -package message +package layer1 import ( "bytes" diff --git a/pkg/network/message/message_bench_test.go b/pkg/message/layer1/message_bench_test.go similarity index 99% rename from pkg/network/message/message_bench_test.go rename to pkg/message/layer1/message_bench_test.go index 16fe52e9..410c9f1f 100644 --- a/pkg/network/message/message_bench_test.go +++ b/pkg/message/layer1/message_bench_test.go @@ -1,4 +1,4 @@ -package message +package layer1 import ( "runtime" diff --git a/pkg/network/message/message_test.go b/pkg/message/layer1/message_test.go similarity index 99% rename from pkg/network/message/message_test.go rename to pkg/message/layer1/message_test.go index 9429474e..d822a6f4 100644 --- a/pkg/network/message/message_test.go +++ b/pkg/message/layer1/message_test.go @@ -1,4 +1,4 @@ -package message +package layer1 import ( "bytes" diff --git a/pkg/network/message/settings.go b/pkg/message/layer1/settings.go similarity index 98% rename from pkg/network/message/settings.go rename to pkg/message/layer1/settings.go index cc9e6345..0a01faee 100644 --- a/pkg/network/message/settings.go +++ b/pkg/message/layer1/settings.go @@ -1,4 +1,4 @@ -package message +package layer1 var ( _ IConstructSettings = &sConstructSettings{} diff --git a/pkg/network/message/types.go b/pkg/message/layer1/types.go similarity index 96% rename from pkg/network/message/types.go rename to pkg/message/layer1/types.go index 80a89b0a..1e11e175 100644 --- a/pkg/network/message/types.go +++ b/pkg/message/layer1/types.go @@ -1,4 +1,4 @@ -package message +package layer1 import ( "github.com/number571/go-peer/pkg/payload" diff --git a/pkg/message/layer2/doc.go b/pkg/message/layer2/doc.go new file mode 100644 index 00000000..4c3a7122 --- /dev/null +++ b/pkg/message/layer2/doc.go @@ -0,0 +1,20 @@ +// Package message used as a storage and loading of encrypted messages. +// +// The package allows initializing verification of the correctness of +// the message by the hash length and proof of work. +/* + NETWORK MESSAGE FORMAT + + E( K, P(HM) || HM || M ) + where + HM = H( K, M ) + where + H - HMAC + K - network key + M - message bytes + P - proof of work + E - encrypt + + Scheme: https://github.com/number571/go-peer/blob/master/images/go-peer_layer2_message.jpg +*/ +package layer2 diff --git a/pkg/client/message/errors.go b/pkg/message/layer2/errors.go similarity index 86% rename from pkg/client/message/errors.go rename to pkg/message/layer2/errors.go index ed531ca2..03241e97 100644 --- a/pkg/client/message/errors.go +++ b/pkg/message/layer2/errors.go @@ -1,7 +1,7 @@ -package message +package layer2 const ( - errPrefix = "pkg/client/message = " + errPrefix = "pkg/message/layer2 = " ) type SMessageError struct { diff --git a/pkg/client/message/message.go b/pkg/message/layer2/message.go similarity index 98% rename from pkg/client/message/message.go rename to pkg/message/layer2/message.go index 24440820..a5820102 100644 --- a/pkg/client/message/message.go +++ b/pkg/message/layer2/message.go @@ -1,4 +1,4 @@ -package message +package layer2 import ( "bytes" diff --git a/pkg/client/message/message_test.go b/pkg/message/layer2/message_test.go similarity index 99% rename from pkg/client/message/message_test.go rename to pkg/message/layer2/message_test.go index 9b63eafc..685c3ee5 100644 --- a/pkg/client/message/message_test.go +++ b/pkg/message/layer2/message_test.go @@ -1,4 +1,4 @@ -package message +package layer2 import ( "bytes" diff --git a/pkg/client/message/test_binary.msg b/pkg/message/layer2/test_binary.msg similarity index 100% rename from pkg/client/message/test_binary.msg rename to pkg/message/layer2/test_binary.msg diff --git a/pkg/client/message/test_string.msg b/pkg/message/layer2/test_string.msg similarity index 100% rename from pkg/client/message/test_string.msg rename to pkg/message/layer2/test_string.msg diff --git a/pkg/client/message/types.go b/pkg/message/layer2/types.go similarity index 89% rename from pkg/client/message/types.go rename to pkg/message/layer2/types.go index 28b177b2..04de370a 100644 --- a/pkg/client/message/types.go +++ b/pkg/message/layer2/types.go @@ -1,4 +1,4 @@ -package message +package layer2 import ( "github.com/number571/go-peer/pkg/types" diff --git a/pkg/network/conn/conn.go b/pkg/network/conn/conn.go index a521228e..e00f8c22 100644 --- a/pkg/network/conn/conn.go +++ b/pkg/network/conn/conn.go @@ -9,9 +9,8 @@ import ( "time" "github.com/number571/go-peer/pkg/encoding" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/payload/joiner" - - net_message "github.com/number571/go-peer/pkg/network/message" ) var ( @@ -52,7 +51,7 @@ func (p *sConn) Close() error { return p.fSocket.Close() } -func (p *sConn) WriteMessage(pCtx context.Context, pMsg net_message.IMessage) error { +func (p *sConn) WriteMessage(pCtx context.Context, pMsg layer1.IMessage) error { p.fMutex.Lock() defer p.fMutex.Unlock() @@ -64,7 +63,7 @@ func (p *sConn) WriteMessage(pCtx context.Context, pMsg net_message.IMessage) er return nil } -func (p *sConn) ReadMessage(pCtx context.Context, pChRead chan<- struct{}) (net_message.IMessage, error) { +func (p *sConn) ReadMessage(pCtx context.Context, pChRead chan<- struct{}) (layer1.IMessage, error) { // large wait read deadline => the connection has not sent anything yet msgSize, err := p.recvHeadBytes(pCtx, pChRead, p.fSettings.GetWaitReadTimeout()) if err != nil { @@ -77,7 +76,7 @@ func (p *sConn) ReadMessage(pCtx context.Context, pChRead chan<- struct{}) (net_ } // try unpack message from bytes - msg, err := net_message.LoadMessage(p.fSettings.GetMessageSettings(), dataBytes) + msg, err := layer1.LoadMessage(p.fSettings.GetMessageSettings(), dataBytes) if err != nil { return nil, errors.Join(ErrInvalidMessageBytes, err) } @@ -143,10 +142,10 @@ func (p *sConn) recvHeadBytes( copy(msgSizeBytes[:], headBytes) gotMsgSize := encoding.BytesToUint32(msgSizeBytes) - fullMsgSize := p.fSettings.GetLimitMessageSizeBytes() + net_message.CMessageHeadSize + fullMsgSize := p.fSettings.GetLimitMessageSizeBytes() + layer1.CMessageHeadSize switch { - case gotMsgSize < net_message.CMessageHeadSize: + case gotMsgSize < layer1.CMessageHeadSize: fallthrough case uint64(gotMsgSize) > fullMsgSize: return 0, ErrInvalidMsgSize diff --git a/pkg/network/conn/conn_test.go b/pkg/network/conn/conn_test.go index 20c1eb67..1c72ea4c 100644 --- a/pkg/network/conn/conn_test.go +++ b/pkg/network/conn/conn_test.go @@ -13,7 +13,7 @@ import ( "github.com/number571/go-peer/pkg/crypto/random" "github.com/number571/go-peer/pkg/encoding" - "github.com/number571/go-peer/pkg/network/message" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/payload" testutils "github.com/number571/go-peer/test/utils" ) @@ -91,7 +91,7 @@ func testSettings(t *testing.T, n int) { switch n { case 0: _ = NewSettings(&SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{}), + FMessageSettings: layer1.NewSettings(&layer1.SSettings{}), FWaitReadTimeout: time.Hour, FDialTimeout: time.Minute, FReadTimeout: time.Minute, @@ -99,7 +99,7 @@ func testSettings(t *testing.T, n int) { }) case 1: _ = NewSettings(&SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{}), + FMessageSettings: layer1.NewSettings(&layer1.SSettings{}), FLimitMessageSizeBytes: tcMsgSize, FDialTimeout: time.Minute, FReadTimeout: time.Minute, @@ -107,7 +107,7 @@ func testSettings(t *testing.T, n int) { }) case 2: _ = NewSettings(&SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{}), + FMessageSettings: layer1.NewSettings(&layer1.SSettings{}), FLimitMessageSizeBytes: tcMsgSize, FWaitReadTimeout: time.Hour, FDialTimeout: time.Minute, @@ -115,7 +115,7 @@ func testSettings(t *testing.T, n int) { }) case 3: _ = NewSettings(&SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{}), + FMessageSettings: layer1.NewSettings(&layer1.SSettings{}), FLimitMessageSizeBytes: tcMsgSize, FWaitReadTimeout: time.Hour, FDialTimeout: time.Minute, @@ -123,7 +123,7 @@ func testSettings(t *testing.T, n int) { }) case 4: _ = NewSettings(&SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{}), + FMessageSettings: layer1.NewSettings(&layer1.SSettings{}), FLimitMessageSizeBytes: tcMsgSize, FWaitReadTimeout: time.Hour, FReadTimeout: time.Minute, @@ -149,7 +149,7 @@ func TestClosedConn(t *testing.T) { conn, err := Connect( context.Background(), NewSettings(&SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{ + FMessageSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: tcWorkSize, }), FLimitMessageSizeBytes: tcMsgSize, @@ -170,12 +170,12 @@ func TestClosedConn(t *testing.T) { return } - sett := message.NewConstructSettings(&message.SConstructSettings{ + sett := layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: conn.GetSettings().GetMessageSettings(), }) pld := payload.NewPayload32(1, []byte("aaa")) - msg := message.NewMessage(sett, pld) + msg := layer1.NewMessage(sett, pld) ctx, cancel := context.WithCancel(context.Background()) defer cancel() @@ -211,7 +211,7 @@ func TestInvalidConn(t *testing.T) { _, err := Connect( context.Background(), NewSettings(&SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{ + FMessageSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: tcWorkSize, }), FLimitMessageSizeBytes: tcMsgSize, @@ -234,7 +234,7 @@ func TestReadMessage(t *testing.T) { rawConn := &tsConn{} conn := LoadConn( NewSettings(&SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{ + FMessageSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: tcWorkSize, }), FLimitMessageSizeBytes: tcMsgSize, @@ -251,8 +251,8 @@ func TestReadMessage(t *testing.T) { ch := make(chan struct{}) rawConn.bodyPart = false - rawConn.headSize = message.CMessageHeadSize + 10 - rawConn.bodySize = message.CMessageHeadSize + 10 + rawConn.headSize = layer1.CMessageHeadSize + 10 + rawConn.bodySize = layer1.CMessageHeadSize + 10 go func() { defer wg.Done() ctx := context.Background() @@ -267,8 +267,8 @@ func TestReadMessage(t *testing.T) { wg.Add(1) rawConn.cancelBody = true rawConn.bodyPart = false - rawConn.headSize = message.CMessageHeadSize + 10 - rawConn.bodySize = message.CMessageHeadSize + 10 + rawConn.headSize = layer1.CMessageHeadSize + 10 + rawConn.bodySize = layer1.CMessageHeadSize + 10 go func() { defer wg.Done() ctx := context.Background() @@ -287,7 +287,7 @@ func TestRecvDataBytes(t *testing.T) { rawConn := &tsConn{} conn := LoadConn( NewSettings(&SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{ + FMessageSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: tcWorkSize, }), FLimitMessageSizeBytes: tcMsgSize, @@ -310,8 +310,8 @@ func TestRecvDataBytes(t *testing.T) { } rawConn.bodyPart = false - rawConn.headSize = message.CMessageHeadSize + 10 - rawConn.bodySize = message.CMessageHeadSize + 10 + rawConn.headSize = layer1.CMessageHeadSize + 10 + rawConn.bodySize = layer1.CMessageHeadSize + 10 rawConn.readDlError = true go func() { ctx := context.Background() @@ -329,7 +329,7 @@ func TestSendBytes(t *testing.T) { rawConn := &tsConn{} conn := LoadConn( NewSettings(&SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{ + FMessageSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: tcWorkSize, }), FLimitMessageSizeBytes: tcMsgSize, @@ -361,7 +361,7 @@ func TestRecvHeadBytes(t *testing.T) { rawConn := &tsConn{} conn := LoadConn( NewSettings(&SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{ + FMessageSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: tcWorkSize, }), FLimitMessageSizeBytes: tcMsgSize, @@ -412,7 +412,7 @@ func testConn(t *testing.T, pAddr, pNetworkKey string) { conn, err := Connect( context.Background(), NewSettings(&SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{ + FMessageSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: tcWorkSize, }), FLimitMessageSizeBytes: tcMsgSize, @@ -435,12 +435,12 @@ func testConn(t *testing.T, pAddr, pNetworkKey string) { return } - msgSett := message.NewConstructSettings(&message.SConstructSettings{ + msgSett := layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: conn.GetSettings().GetMessageSettings(), }) pld := payload.NewPayload32(tcHead, []byte(tcBody)) - msg := message.NewMessage(msgSett, pld) + msg := layer1.NewMessage(msgSett, pld) ctx := context.Background() if err := conn.WriteMessage(ctx, msg); err != nil { t.Error(err) @@ -478,7 +478,7 @@ func testNewService(t *testing.T, pAddr, pNetworkKey string) net.Listener { conn := LoadConn( NewSettings(&SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{ + FMessageSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: tcWorkSize, FNetworkKey: pNetworkKey, }), diff --git a/pkg/network/conn/settings.go b/pkg/network/conn/settings.go index 5f5fb5c6..2625585c 100644 --- a/pkg/network/conn/settings.go +++ b/pkg/network/conn/settings.go @@ -3,7 +3,7 @@ package conn import ( "time" - net_message "github.com/number571/go-peer/pkg/network/message" + "github.com/number571/go-peer/pkg/message/layer1" ) var ( @@ -17,7 +17,7 @@ type sSettings struct { FDialTimeout time.Duration FReadTimeout time.Duration FWriteTimeout time.Duration - FMessageSettings net_message.ISettings + FMessageSettings layer1.ISettings } func NewSettings(pSett *SSettings) ISettings { @@ -53,7 +53,7 @@ func (p *sSettings) mustNotNull() ISettings { return p } -func (p *sSettings) GetMessageSettings() net_message.ISettings { +func (p *sSettings) GetMessageSettings() layer1.ISettings { return p.FMessageSettings } diff --git a/pkg/network/conn/types.go b/pkg/network/conn/types.go index 43736499..982d6874 100644 --- a/pkg/network/conn/types.go +++ b/pkg/network/conn/types.go @@ -6,7 +6,7 @@ import ( "net" "time" - net_message "github.com/number571/go-peer/pkg/network/message" + "github.com/number571/go-peer/pkg/message/layer1" ) type IConn interface { @@ -15,12 +15,12 @@ type IConn interface { GetSettings() ISettings GetSocket() net.Conn - WriteMessage(context.Context, net_message.IMessage) error - ReadMessage(context.Context, chan<- struct{}) (net_message.IMessage, error) + WriteMessage(context.Context, layer1.IMessage) error + ReadMessage(context.Context, chan<- struct{}) (layer1.IMessage, error) } type ISettings interface { - GetMessageSettings() net_message.ISettings + GetMessageSettings() layer1.ISettings GetLimitMessageSizeBytes() uint64 GetDialTimeout() time.Duration GetReadTimeout() time.Duration diff --git a/pkg/network/connkeeper/connkeeper_test.go b/pkg/network/connkeeper/connkeeper_test.go index 4aba535d..96aa5709 100644 --- a/pkg/network/connkeeper/connkeeper_test.go +++ b/pkg/network/connkeeper/connkeeper_test.go @@ -8,9 +8,9 @@ import ( "testing" "time" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/network" "github.com/number571/go-peer/pkg/network/conn" - "github.com/number571/go-peer/pkg/network/message" "github.com/number571/go-peer/pkg/storage/cache" testutils "github.com/number571/go-peer/test/utils" ) @@ -132,7 +132,7 @@ func newTestConnKeeper(pDuration time.Duration) IConnKeeper { FReadTimeout: time.Minute, FWriteTimeout: time.Minute, FConnSettings: conn.NewSettings(&conn.SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{ + FMessageSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: 10, }), FLimitMessageSizeBytes: (8 << 10), diff --git a/pkg/network/examples/broadcast/main.go b/pkg/network/examples/broadcast/main.go index 645cb6fd..47babf66 100644 --- a/pkg/network/examples/broadcast/main.go +++ b/pkg/network/examples/broadcast/main.go @@ -5,9 +5,9 @@ import ( "fmt" "time" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/network" "github.com/number571/go-peer/pkg/network/conn" - "github.com/number571/go-peer/pkg/network/message" "github.com/number571/go-peer/pkg/payload" ) @@ -20,7 +20,7 @@ const ( ) var handler = func(serviceName string) network.IHandlerF { - return func(ctx context.Context, node network.INode, _ conn.IConn, msg message.IMessage) error { + return func(ctx context.Context, node network.INode, _ conn.IConn, msg layer1.IMessage) error { defer node.BroadcastMessage(ctx, msg) // send this message to other connections fmt.Printf("'%s' got '%s'\n", serviceName, string(msg.GetPayload().GetBody())) return nil @@ -40,8 +40,8 @@ func main() { node4.BroadcastMessage( context.Background(), - message.NewMessage( - message.NewConstructSettings(&message.SConstructSettings{ + layer1.NewMessage( + layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: node4.GetSettings().GetConnSettings().GetMessageSettings(), }), payload.NewPayload32( diff --git a/pkg/network/examples/echo/main.go b/pkg/network/examples/echo/main.go index 60b86a02..c3d66ed1 100644 --- a/pkg/network/examples/echo/main.go +++ b/pkg/network/examples/echo/main.go @@ -5,9 +5,9 @@ import ( "fmt" "time" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/network" "github.com/number571/go-peer/pkg/network/conn" - "github.com/number571/go-peer/pkg/network/message" "github.com/number571/go-peer/pkg/payload" ) @@ -18,12 +18,12 @@ const ( serviceAddress = "127.0.0.1:8080" ) -var handler = func(ctx context.Context, node network.INode, c conn.IConn, msg message.IMessage) error { +var handler = func(ctx context.Context, node network.INode, c conn.IConn, msg layer1.IMessage) error { resp := fmt.Sprintf("echo: [%s]", string(msg.GetPayload().GetBody())) _ = c.WriteMessage( ctx, - message.NewMessage( - message.NewConstructSettings(&message.SConstructSettings{ + layer1.NewMessage( + layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: node.GetSettings().GetConnSettings().GetMessageSettings(), }), payload.NewPayload32(serviceHeader, []byte(resp)), @@ -43,8 +43,8 @@ func main() { _ = conn.WriteMessage( ctx, - message.NewMessage( - message.NewConstructSettings(&message.SConstructSettings{ + layer1.NewMessage( + layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: conn.GetSettings().GetMessageSettings(), }), payload.NewPayload32(serviceHeader, []byte("hello, world!")), diff --git a/pkg/network/examples/ping-pong/main.go b/pkg/network/examples/ping-pong/main.go index 9a65c4ce..3655fd7f 100644 --- a/pkg/network/examples/ping-pong/main.go +++ b/pkg/network/examples/ping-pong/main.go @@ -6,9 +6,9 @@ import ( "strconv" "time" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/network" "github.com/number571/go-peer/pkg/network/conn" - "github.com/number571/go-peer/pkg/network/message" "github.com/number571/go-peer/pkg/payload" ) @@ -18,7 +18,7 @@ const ( ) var handler = func(id string) network.IHandlerF { - return func(ctx context.Context, n network.INode, _ conn.IConn, msg message.IMessage) error { + return func(ctx context.Context, n network.INode, _ conn.IConn, msg layer1.IMessage) error { time.Sleep(time.Second) // delay for view "ping-pong" game num, err := strconv.Atoi(string(msg.GetPayload().GetBody())) @@ -34,8 +34,8 @@ var handler = func(id string) network.IHandlerF { fmt.Printf("'%s' got '%s#%d'\n", id, val, num) n.BroadcastMessage( ctx, - message.NewMessage( - message.NewConstructSettings(&message.SConstructSettings{ + layer1.NewMessage( + layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: n.GetSettings().GetConnSettings().GetMessageSettings(), }), payload.NewPayload32(serviceHeader, []byte(fmt.Sprintf("%d", num+1))), @@ -55,8 +55,8 @@ func main() { node1 = runClientNode("node2") ) - msg := message.NewMessage( - message.NewConstructSettings(&message.SConstructSettings{ + msg := layer1.NewMessage( + layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: node1.GetSettings().GetConnSettings().GetMessageSettings(), }), payload.NewPayload32(serviceHeader, []byte("0")), diff --git a/pkg/network/network.go b/pkg/network/network.go index 9006e757..dae978ce 100644 --- a/pkg/network/network.go +++ b/pkg/network/network.go @@ -7,10 +7,9 @@ import ( "sync" "time" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/network/conn" "github.com/number571/go-peer/pkg/storage/cache" - - net_message "github.com/number571/go-peer/pkg/network/message" ) var ( @@ -51,7 +50,7 @@ func (p *sNode) GetCacheSetter() cache.ICacheSetter { } // Puts the hash of the message in the buffer and sends the message to all connections of the node. -func (p *sNode) BroadcastMessage(pCtx context.Context, pMsg net_message.IMessage) error { +func (p *sNode) BroadcastMessage(pCtx context.Context, pMsg layer1.IMessage) error { connections := p.GetConnections() lenConnections := len(connections) @@ -220,7 +219,7 @@ func (p *sNode) handleConn(pCtx context.Context, pAddress string, pConn conn.ICo var ( readHeadCh = make(chan struct{}) - readFullCh = make(chan net_message.IMessage) + readFullCh = make(chan layer1.IMessage) ) go p.messageReader( @@ -257,7 +256,7 @@ func (p *sNode) messageReader( pCtx context.Context, pConn conn.IConn, readHeadCh chan<- struct{}, - readFullCh chan<- net_message.IMessage, + readFullCh chan<- layer1.IMessage, ) { for { select { @@ -277,7 +276,7 @@ func (p *sNode) messageReader( // Processes the message for correctness and redirects it to the handler function. // Returns true if the message was successfully redirected to the handler function // > or if the message already existed in the hash value store. -func (p *sNode) handleMessage(pCtx context.Context, pConn conn.IConn, pMsg net_message.IMessage) bool { +func (p *sNode) handleMessage(pCtx context.Context, pConn conn.IConn, pMsg layer1.IMessage) bool { if !p.fCacheSetter.Set(pMsg.GetHash(), []byte{}) { return true // hash of message already in queue } diff --git a/pkg/network/network_test.go b/pkg/network/network_test.go index 6b2cb9e4..e91c4e91 100644 --- a/pkg/network/network_test.go +++ b/pkg/network/network_test.go @@ -10,8 +10,8 @@ import ( "testing" "time" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/network/conn" - "github.com/number571/go-peer/pkg/network/message" "github.com/number571/go-peer/pkg/payload" "github.com/number571/go-peer/pkg/storage/cache" testutils "github.com/number571/go-peer/test/utils" @@ -56,7 +56,7 @@ func testSettings(t *testing.T, n int) { FReadTimeout: tcTimeWait, FWriteTimeout: tcTimeWait, FConnSettings: conn.NewSettings(&conn.SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{}), + FMessageSettings: layer1.NewSettings(&layer1.SSettings{}), FLimitMessageSizeBytes: (8 << 10), FWaitReadTimeout: time.Hour, FDialTimeout: time.Minute, @@ -70,7 +70,7 @@ func testSettings(t *testing.T, n int) { FMaxConnects: 16, FWriteTimeout: tcTimeWait, FConnSettings: conn.NewSettings(&conn.SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{}), + FMessageSettings: layer1.NewSettings(&layer1.SSettings{}), FLimitMessageSizeBytes: (8 << 10), FWaitReadTimeout: time.Hour, FDialTimeout: time.Minute, @@ -84,7 +84,7 @@ func testSettings(t *testing.T, n int) { FMaxConnects: 16, FReadTimeout: tcTimeWait, FConnSettings: conn.NewSettings(&conn.SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{}), + FMessageSettings: layer1.NewSettings(&layer1.SSettings{}), FLimitMessageSizeBytes: (8 << 10), FWaitReadTimeout: time.Hour, FDialTimeout: time.Minute, @@ -120,7 +120,7 @@ func TestBroadcast(t *testing.T) { wg.Add(4 * tcIter) headHandle := uint32(123) - handleF := func(pCtx context.Context, node INode, _ conn.IConn, pMsg message.IMessage) error { + handleF := func(pCtx context.Context, node INode, _ conn.IConn, pMsg layer1.IMessage) error { defer func() { _ = node.BroadcastMessage(pCtx, pMsg) wg.Done() @@ -158,10 +158,10 @@ func TestBroadcast(t *testing.T) { headHandle, []byte(fmt.Sprintf(tcBodyTemplate, i)), ) - sett := message.NewConstructSettings(&message.SConstructSettings{ + sett := layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: nodes[0].GetSettings().GetConnSettings().GetMessageSettings(), }) - _ = nodes[0].BroadcastMessage(ctx, message.NewMessage(sett, pld)) + _ = nodes[0].BroadcastMessage(ctx, layer1.NewMessage(sett, pld)) }(i) } @@ -286,30 +286,30 @@ func TestHandleMessage(t *testing.T) { node := newTestNode("", 16).(*sNode) ctx := context.Background() - sett := message.NewConstructSettings(&message.SConstructSettings{ + sett := layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: node.GetSettings().GetConnSettings().GetMessageSettings(), }) node.HandleFunc(1, nil) - msg1 := message.NewMessage(sett, payload.NewPayload32(1, []byte{1})) + msg1 := layer1.NewMessage(sett, payload.NewPayload32(1, []byte{1})) if ok := node.handleMessage(ctx, nil, msg1); ok { t.Error("success handle message with nil function") return } - node.HandleFunc(1, func(_ context.Context, _ INode, _ conn.IConn, _ message.IMessage) error { + node.HandleFunc(1, func(_ context.Context, _ INode, _ conn.IConn, _ layer1.IMessage) error { return errors.New("some error") }) - msg2 := message.NewMessage(sett, payload.NewPayload32(1, []byte{2})) + msg2 := layer1.NewMessage(sett, payload.NewPayload32(1, []byte{2})) if ok := node.handleMessage(ctx, nil, msg2); ok { t.Error("success handle message with got error from function") return } - node.HandleFunc(1, func(_ context.Context, _ INode, _ conn.IConn, _ message.IMessage) error { + node.HandleFunc(1, func(_ context.Context, _ INode, _ conn.IConn, _ layer1.IMessage) error { return nil }) - msg3 := message.NewMessage(sett, payload.NewPayload32(1, []byte{3})) + msg3 := layer1.NewMessage(sett, payload.NewPayload32(1, []byte{3})) if ok := node.handleMessage(ctx, nil, msg3); !ok { t.Error("failed handle message with correct function") return @@ -345,7 +345,7 @@ func TestContextCancel(t *testing.T) { } headHandle := uint32(123) - sett := message.NewConstructSettings(&message.SConstructSettings{ + sett := layer1.NewConstructSettings(&layer1.SConstructSettings{ FSettings: node2.GetSettings().GetConnSettings().GetMessageSettings(), }) @@ -355,7 +355,7 @@ func TestContextCancel(t *testing.T) { headHandle, []byte(fmt.Sprintf(tcBodyTemplate, i)), ) - if err := node2.BroadcastMessage(ctx, message.NewMessage(sett, pld)); err != nil { + if err := node2.BroadcastMessage(ctx, layer1.NewMessage(sett, pld)); err != nil { return } } @@ -416,7 +416,7 @@ func newTestNode(pAddr string, pMaxConns uint64) INode { FReadTimeout: timeout, FWriteTimeout: timeout, FConnSettings: conn.NewSettings(&conn.SSettings{ - FMessageSettings: message.NewSettings(&message.SSettings{ + FMessageSettings: layer1.NewSettings(&layer1.SSettings{ FWorkSizeBits: 10, }), FLimitMessageSizeBytes: (8 << 10), diff --git a/pkg/network/types.go b/pkg/network/types.go index 06dbf87a..e3467261 100644 --- a/pkg/network/types.go +++ b/pkg/network/types.go @@ -4,18 +4,17 @@ import ( "context" "time" + "github.com/number571/go-peer/pkg/message/layer1" "github.com/number571/go-peer/pkg/network/conn" "github.com/number571/go-peer/pkg/storage/cache" "github.com/number571/go-peer/pkg/types" - - net_message "github.com/number571/go-peer/pkg/network/message" ) type IHandlerF func( context.Context, INode, conn.IConn, - net_message.IMessage, + layer1.IMessage, ) error type INode interface { @@ -29,7 +28,7 @@ type INode interface { AddConnection(context.Context, string) error DelConnection(string) error - BroadcastMessage(context.Context, net_message.IMessage) error + BroadcastMessage(context.Context, layer1.IMessage) error } type ISettings interface { diff --git a/test/result/badge_codelines.svg b/test/result/badge_codelines.svg index 1dabc5a3..42963782 100644 --- a/test/result/badge_codelines.svg +++ b/test/result/badge_codelines.svg @@ -1 +1 @@ -code lines: 12551code lines12551 \ No newline at end of file +code lines: 12561code lines12561 \ No newline at end of file diff --git a/test/result/coverage.svg b/test/result/coverage.svg index 8d09958c..6e5ae3c9 100644 --- a/test/result/coverage.svg +++ b/test/result/coverage.svg @@ -7,7 +7,7 @@ > - + - + anonymity @@ -52,12 +52,12 @@ - + client @@ -65,12 +65,12 @@ - + crypto @@ -78,12 +78,12 @@ - + encoding @@ -91,12 +91,12 @@ - + logger @@ -104,12 +104,25 @@ - + message + + + + + + + +network @@ -117,12 +130,12 @@ - + payload @@ -130,12 +143,12 @@ - + state @@ -143,12 +156,12 @@ - + storage @@ -156,12 +169,12 @@ - + action.go @@ -169,18 +182,18 @@ - + - + anonymity.go @@ -188,18 +201,18 @@ - + - + head.go @@ -207,12 +220,12 @@ - + logger/log_builder.go @@ -220,12 +233,12 @@ - + queue @@ -233,12 +246,12 @@ - + settings.go @@ -246,12 +259,12 @@ - + client.go @@ -259,25 +272,12 @@ - + message - - - - - - - -asymmetric @@ -285,12 +285,12 @@ - + hashing @@ -298,18 +298,18 @@ - + - + puzzle/puzzle.go @@ -317,12 +317,12 @@ - + random/random.go @@ -330,12 +330,12 @@ - + symmetric/symmetric.go @@ -343,12 +343,12 @@ - + bytes.go @@ -356,50 +356,43 @@ - + - + + + + + + + + + + + + + hex.go + data-math="N">serialize_yaml.go - + serialize_json.go - - - - - - - - - - - - - -logger.go @@ -407,18 +400,44 @@ - + - + layer1 + + + + + + + +layer2 + + + + + + + +conn @@ -426,12 +445,12 @@ - + connkeeper @@ -439,31 +458,12 @@ - - - - - - - + message - - - - - - - -network.go @@ -471,12 +471,12 @@ - + settings.go @@ -484,12 +484,12 @@ - + joiner @@ -497,12 +497,12 @@ - + payload32.go @@ -510,12 +510,12 @@ - + payload64.go @@ -523,24 +523,18 @@ - + - - - - - - - + cache/lru.go @@ -548,12 +542,12 @@ - + database @@ -561,18 +555,18 @@ - + - + queue.go @@ -580,31 +574,25 @@ - - - - - - - + message.go + data-math="N">settings.go - + dsa.go @@ -612,12 +600,12 @@ - + kem.go @@ -625,12 +613,12 @@ - + key.go @@ -638,108 +626,36 @@ - + + + + + + + + + + + + + + + + + + + + + + + + + map_pubkeys.go - - - - - - - - - - - - - -hmac.go - - - - - - - -conn.go - - - - - - - -settings.go - - - - - - - -connkeeper.go - - - - - - - - - - - - - -settings.go - - - - - - - - - - - - - -message.go @@ -747,24 +663,95 @@ - - - - - - - - - - - - - + settings.go + + + + + + + +message.go + + + + + + + +conn.go + + + + + + + +settings.go + + + + + + + +connkeeper.go + + + + + + + + + + + + + + + + + + + + + + + + + +database.go