aboutsummaryrefslogtreecommitdiffstats
path: root/peer.go
blob: a9f88b1e12b0dd6be54197d457dbe3eb0bc6b1b5 (plain) (blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
package main

import (
    "github.com/ethereum/ethwire-go"
    "log"
    "net"
)

type Peer struct {
    // Server interface
    server *Server
    // Net connection
    conn net.Conn
    // Output queue which is used to communicate and handle messages
    outputQueue chan ethwire.InOutMsg
    // Quit channel
    quit chan bool
}

func NewPeer(conn net.Conn, server *Server) *Peer {
    return &Peer{
        outputQueue: make(chan ethwire.InOutMsg, 1), // Buffered chan of 1 is enough
        quit:        make(chan bool),

        server: server,
        conn:   conn,
    }
}

// Outputs any RLP encoded data to the peer
func (p *Peer) QueueMessage(msgType string, data []byte) {
    p.outputQueue <- ethwire.InOutMsg{MsgType: msgType, Data: data}
}

// Outbound message handler. Outbound messages are handled here
func (p *Peer) HandleOutbound() {
out:
    for {
        select {
        // Main message queue. All outbound messages are processed through here
        case msg := <-p.outputQueue:
            // TODO Message checking and handle accordingly
            err := ethwire.WriteMessage(p.conn, msg)
            if err != nil {
                log.Println(err)

                // Stop the client if there was an error writing to it
                p.Stop()
            }

        // Break out of the for loop if a quit message is posted
        case <-p.quit:
            break out
        }
    }
}

// Inbound handler. Inbound messages are received here and passed to the appropriate methods
func (p *Peer) HandleInbound() {
    defer p.Stop()

out:
    for {
        // Wait for a message from the peer
        msg, err := ethwire.ReadMessage(p.conn)
        if err != nil {
            log.Println(err)

            break out
        }

        // TODO
        data, _ := Decode(msg.Data, 0)
        log.Printf("%s, %s\n", msg.MsgType, data)
    }

    // Notify the out handler we're quiting
    p.quit <- true
}

func (p *Peer) Start() {
    // Run the outbound handler in a new goroutine
    go p.HandleOutbound()
    // Run the inbound handler in a new goroutine
    go p.HandleInbound()
}

func (p *Peer) Stop() {
    p.conn.Close()

    p.quit <- true
}