2018-01-29 14:37:26 -05:00
|
|
|
package reflector
|
2017-08-10 18:25:42 -04:00
|
|
|
|
|
|
|
import (
|
|
|
|
"encoding/json"
|
2020-01-08 22:12:40 -05:00
|
|
|
"log"
|
2017-08-10 18:25:42 -04:00
|
|
|
"net"
|
2018-08-20 17:50:39 -04:00
|
|
|
|
2019-11-13 19:11:35 -05:00
|
|
|
"github.com/lbryio/lbry.go/v2/extras/errors"
|
|
|
|
"github.com/lbryio/lbry.go/v2/stream"
|
2017-08-10 18:25:42 -04:00
|
|
|
)
|
|
|
|
|
2018-08-09 14:56:49 -04:00
|
|
|
// ErrBlobExists is a default error for when a blob already exists on the reflector server.
|
|
|
|
var ErrBlobExists = errors.Base("blob exists on server")
|
|
|
|
|
2018-05-29 21:38:55 -04:00
|
|
|
// Client is an instance of a client connected to a server.
|
2017-08-10 18:25:42 -04:00
|
|
|
type Client struct {
|
2017-08-15 16:02:18 -04:00
|
|
|
conn net.Conn
|
|
|
|
connected bool
|
2017-08-10 18:25:42 -04:00
|
|
|
}
|
|
|
|
|
2018-05-29 21:38:55 -04:00
|
|
|
// Connect connects to a specific clients and errors if it cannot be contacted.
|
2017-08-10 18:25:42 -04:00
|
|
|
func (c *Client) Connect(address string) error {
|
|
|
|
var err error
|
2018-08-09 14:56:49 -04:00
|
|
|
c.conn, err = net.Dial(network, address)
|
2017-08-10 18:25:42 -04:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2017-08-15 16:02:18 -04:00
|
|
|
c.connected = true
|
2017-08-10 18:25:42 -04:00
|
|
|
return c.doHandshake(protocolVersion1)
|
|
|
|
}
|
2018-05-29 21:38:55 -04:00
|
|
|
|
|
|
|
// Close closes the connection with the client.
|
2017-08-10 18:25:42 -04:00
|
|
|
func (c *Client) Close() error {
|
2017-08-15 16:02:18 -04:00
|
|
|
c.connected = false
|
2017-08-10 18:25:42 -04:00
|
|
|
return c.conn.Close()
|
|
|
|
}
|
|
|
|
|
2020-01-08 22:12:40 -05:00
|
|
|
// SendBlob sends a blob to the server.
|
2018-08-20 17:50:39 -04:00
|
|
|
func (c *Client) SendBlob(blob stream.Blob) error {
|
2020-01-08 22:12:40 -05:00
|
|
|
return c.sendBlob(blob, false)
|
|
|
|
}
|
|
|
|
|
|
|
|
// SendSDBlob sends an SD blob request to the server.
|
|
|
|
func (c *Client) SendSDBlob(blob stream.Blob) error {
|
|
|
|
return c.sendBlob(blob, true)
|
|
|
|
}
|
|
|
|
|
|
|
|
// sendBlob does the actual blob sending
|
|
|
|
func (c *Client) sendBlob(blob stream.Blob, isSDBlob bool) error {
|
2017-08-15 16:02:18 -04:00
|
|
|
if !c.connected {
|
2018-01-24 11:45:18 -05:00
|
|
|
return errors.Err("not connected")
|
2017-08-15 16:02:18 -04:00
|
|
|
}
|
|
|
|
|
2018-08-20 17:50:39 -04:00
|
|
|
if err := blob.ValidForSend(); err != nil {
|
|
|
|
return errors.Err(err)
|
2017-08-10 18:25:42 -04:00
|
|
|
}
|
|
|
|
|
2018-08-20 17:50:39 -04:00
|
|
|
blobHash := blob.HashHex()
|
2020-01-08 22:12:40 -05:00
|
|
|
var req sendBlobRequest
|
|
|
|
if isSDBlob {
|
|
|
|
req.SdBlobSize = blob.Size()
|
|
|
|
req.SdBlobHash = blobHash
|
|
|
|
} else {
|
|
|
|
req.BlobSize = blob.Size()
|
|
|
|
req.BlobHash = blobHash
|
|
|
|
}
|
|
|
|
sendRequest, err := json.Marshal(req)
|
2017-08-10 18:25:42 -04:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2018-08-09 14:56:49 -04:00
|
|
|
|
2017-08-10 18:25:42 -04:00
|
|
|
_, err = c.conn.Write(sendRequest)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
dec := json.NewDecoder(c.conn)
|
|
|
|
|
2020-01-08 22:12:40 -05:00
|
|
|
if isSDBlob {
|
|
|
|
var sendResp sendSdBlobResponse
|
|
|
|
err = dec.Decode(&sendResp)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
if !sendResp.SendSdBlob {
|
|
|
|
return errors.Prefix(blobHash[:8], ErrBlobExists)
|
|
|
|
}
|
|
|
|
log.Println("Sending SD blob " + blobHash[:8])
|
|
|
|
} else {
|
|
|
|
var sendResp sendBlobResponse
|
|
|
|
err = dec.Decode(&sendResp)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
if !sendResp.SendBlob {
|
|
|
|
return errors.Prefix(blobHash[:8], ErrBlobExists)
|
|
|
|
}
|
|
|
|
log.Println("Sending blob " + blobHash[:8])
|
2017-08-10 18:25:42 -04:00
|
|
|
}
|
|
|
|
|
|
|
|
_, err = c.conn.Write(blob)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
2020-01-08 22:12:40 -05:00
|
|
|
if isSDBlob {
|
|
|
|
var transferResp sdBlobTransferResponse
|
|
|
|
err = dec.Decode(&transferResp)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
if !transferResp.ReceivedSdBlob {
|
|
|
|
return errors.Err("server did not received SD blob")
|
|
|
|
}
|
|
|
|
} else {
|
|
|
|
var transferResp blobTransferResponse
|
|
|
|
err = dec.Decode(&transferResp)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
if !transferResp.ReceivedBlob {
|
|
|
|
return errors.Err("server did not received blob")
|
|
|
|
}
|
2017-08-10 18:25:42 -04:00
|
|
|
}
|
|
|
|
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *Client) doHandshake(version int) error {
|
2017-08-15 16:02:18 -04:00
|
|
|
if !c.connected {
|
2018-01-24 11:45:18 -05:00
|
|
|
return errors.Err("not connected")
|
2017-08-15 16:02:18 -04:00
|
|
|
}
|
|
|
|
|
2019-06-05 11:03:55 -04:00
|
|
|
handshake, err := json.Marshal(handshakeRequestResponse{Version: &version})
|
2017-08-10 18:25:42 -04:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
_, err = c.conn.Write(handshake)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
var resp handshakeRequestResponse
|
2018-08-09 14:56:49 -04:00
|
|
|
err = json.NewDecoder(c.conn).Decode(&resp)
|
2017-08-10 18:25:42 -04:00
|
|
|
if err != nil {
|
|
|
|
return err
|
2019-06-05 11:03:55 -04:00
|
|
|
} else if resp.Version == nil {
|
|
|
|
return errors.Err("invalid handshake")
|
|
|
|
} else if *resp.Version != version {
|
2018-01-24 11:45:18 -05:00
|
|
|
return errors.Err("handshake version mismatch")
|
2017-08-10 18:25:42 -04:00
|
|
|
}
|
|
|
|
|
|
|
|
return nil
|
|
|
|
}
|