mirror of
https://github.com/tinygo-org/drivers.git
synced 2026-08-09 09:23:39 +00:00
Decoupled net package from espat
This commit is contained in:
@@ -0,0 +1,27 @@
|
||||
package net
|
||||
|
||||
type DeviceDriver interface {
|
||||
GetDNS(domain string) (string, error)
|
||||
ConnectTCPSocket(addr, port string) error
|
||||
ConnectSSLSocket(addr, port string) error
|
||||
ConnectUDPSocket(addr, sendport, listenport string) error
|
||||
DisconnectSocket() error
|
||||
StartSocketSend(size int) error
|
||||
Write(b []byte) (n int, err error)
|
||||
ReadSocket(b []byte) (n int, err error)
|
||||
IsSocketDataAvailable() bool
|
||||
|
||||
// FIXME: this is really specific to espat, and maybe shouldn't be part
|
||||
// of the driver interface
|
||||
Response(timeout int) ([]byte, error)
|
||||
}
|
||||
|
||||
var ActiveDevice DeviceDriver
|
||||
|
||||
func UseDriver(driver DeviceDriver) {
|
||||
// TODO: rethink and refactor this
|
||||
if ActiveDevice != nil {
|
||||
panic("net.ActiveDevice is already set")
|
||||
}
|
||||
ActiveDevice = driver
|
||||
}
|
||||
@@ -0,0 +1,303 @@
|
||||
// Package mqtt is intended to provide compatible interfaces with the
|
||||
// Paho mqtt library.
|
||||
package mqtt
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/eclipse/paho.mqtt.golang/packets"
|
||||
"tinygo.org/x/drivers/net"
|
||||
"tinygo.org/x/drivers/net/tls"
|
||||
)
|
||||
|
||||
// NewClient will create an MQTT v3.1.1 client with all of the options specified
|
||||
// in the provided ClientOptions. The client must have the Connect method called
|
||||
// on it before it may be used. This is to make sure resources (such as a net
|
||||
// connection) are created before the application is actually ready.
|
||||
func NewClient(o *ClientOptions) Client {
|
||||
c := &mqttclient{opts: o, adaptor: o.Adaptor}
|
||||
c.msgRouter, c.stopRouter = newRouter()
|
||||
return c
|
||||
}
|
||||
|
||||
type mqttclient struct {
|
||||
adaptor net.DeviceDriver
|
||||
conn net.Conn
|
||||
connected bool
|
||||
opts *ClientOptions
|
||||
mid uint16
|
||||
inbound chan packets.ControlPacket
|
||||
stop chan struct{}
|
||||
msgRouter *router
|
||||
stopRouter chan bool
|
||||
incomingPubChan chan *packets.PublishPacket
|
||||
}
|
||||
|
||||
// AddRoute allows you to add a handler for messages on a specific topic
|
||||
// without making a subscription. For example having a different handler
|
||||
// for parts of a wildcard subscription
|
||||
func (c *mqttclient) AddRoute(topic string, callback MessageHandler) {
|
||||
return
|
||||
}
|
||||
|
||||
// IsConnected returns a bool signifying whether
|
||||
// the client is connected or not.
|
||||
func (c *mqttclient) IsConnected() bool {
|
||||
return c.connected
|
||||
}
|
||||
|
||||
// IsConnectionOpen return a bool signifying whether the client has an active
|
||||
// connection to mqtt broker, i.e not in disconnected or reconnect mode
|
||||
func (c *mqttclient) IsConnectionOpen() bool {
|
||||
return c.connected
|
||||
}
|
||||
|
||||
// Connect will create a connection to the message broker.
|
||||
func (c *mqttclient) Connect() Token {
|
||||
var err error
|
||||
|
||||
// make connection
|
||||
if strings.Contains(c.opts.Servers, "ssl://") {
|
||||
url := strings.TrimPrefix(c.opts.Servers, "ssl://")
|
||||
c.conn, err = tls.Dial("tcp", url, nil)
|
||||
if err != nil {
|
||||
return &mqtttoken{err: err}
|
||||
}
|
||||
} else if strings.Contains(c.opts.Servers, "tcp://") {
|
||||
url := strings.TrimPrefix(c.opts.Servers, "tcp://")
|
||||
c.conn, err = net.Dial("tcp", url)
|
||||
if err != nil {
|
||||
return &mqtttoken{err: err}
|
||||
}
|
||||
} else {
|
||||
// invalid protocol
|
||||
return &mqtttoken{err: errors.New("invalid protocol")}
|
||||
}
|
||||
|
||||
c.mid = 1
|
||||
c.inbound = make(chan packets.ControlPacket)
|
||||
c.stop = make(chan struct{})
|
||||
c.incomingPubChan = make(chan *packets.PublishPacket)
|
||||
c.msgRouter.matchAndDispatch(c.incomingPubChan, c.opts.Order, c)
|
||||
|
||||
// send the MQTT connect message
|
||||
connectPkt := packets.NewControlPacket(packets.Connect).(*packets.ConnectPacket)
|
||||
connectPkt.Qos = 0
|
||||
if c.opts.Username != "" {
|
||||
connectPkt.Username = c.opts.Username
|
||||
connectPkt.UsernameFlag = true
|
||||
}
|
||||
|
||||
if c.opts.Password != "" {
|
||||
connectPkt.Password = []byte(c.opts.Password)
|
||||
connectPkt.PasswordFlag = true
|
||||
}
|
||||
|
||||
connectPkt.ClientIdentifier = c.opts.ClientID
|
||||
connectPkt.ProtocolVersion = byte(c.opts.ProtocolVersion)
|
||||
connectPkt.ProtocolName = "MQTT"
|
||||
connectPkt.Keepalive = 30
|
||||
|
||||
err = connectPkt.Write(c.conn)
|
||||
if err != nil {
|
||||
return &mqtttoken{err: err}
|
||||
}
|
||||
|
||||
// TODO: handle timeout as ReadPacket blocks until it gets a packet.
|
||||
// CONNECT response.
|
||||
packet, err := packets.ReadPacket(c.conn)
|
||||
if err != nil {
|
||||
return &mqtttoken{err: err}
|
||||
}
|
||||
if packet != nil {
|
||||
ack, ok := packet.(*packets.ConnackPacket)
|
||||
if ok {
|
||||
if ack.ReturnCode != 0 {
|
||||
return &mqtttoken{err: errors.New(packet.String())}
|
||||
}
|
||||
c.connected = true
|
||||
}
|
||||
}
|
||||
|
||||
go readMessages(c)
|
||||
go processInbound(c)
|
||||
|
||||
return &mqtttoken{}
|
||||
}
|
||||
|
||||
// Disconnect will end the connection with the server, but not before waiting
|
||||
// the specified number of milliseconds to wait for existing work to be
|
||||
// completed.
|
||||
func (c *mqttclient) Disconnect(quiesce uint) {
|
||||
c.conn.Close()
|
||||
return
|
||||
}
|
||||
|
||||
// Publish will publish a message with the specified QoS and content
|
||||
// to the specified topic.
|
||||
// Returns a token to track delivery of the message to the broker
|
||||
func (c *mqttclient) Publish(topic string, qos byte, retained bool, payload interface{}) Token {
|
||||
if !c.IsConnected() {
|
||||
return &mqtttoken{err: errors.New("MQTT client not connected")}
|
||||
}
|
||||
|
||||
pub := packets.NewControlPacket(packets.Publish).(*packets.PublishPacket)
|
||||
pub.Qos = qos
|
||||
pub.TopicName = topic
|
||||
switch payload.(type) {
|
||||
case string:
|
||||
pub.Payload = []byte(payload.(string))
|
||||
case []byte:
|
||||
pub.Payload = payload.([]byte)
|
||||
default:
|
||||
return &mqtttoken{err: errors.New("Unknown payload type")}
|
||||
}
|
||||
pub.MessageID = c.mid
|
||||
c.mid++
|
||||
|
||||
err := pub.Write(c.conn)
|
||||
if err != nil {
|
||||
return &mqtttoken{err: err}
|
||||
}
|
||||
|
||||
return &mqtttoken{}
|
||||
}
|
||||
|
||||
// Subscribe starts a new subscription. Provide a MessageHandler to be executed when
|
||||
// a message is published on the topic provided.
|
||||
func (c *mqttclient) Subscribe(topic string, qos byte, callback MessageHandler) Token {
|
||||
if !c.IsConnected() {
|
||||
return &mqtttoken{err: errors.New("MQTT client not connected")}
|
||||
}
|
||||
|
||||
sub := packets.NewControlPacket(packets.Subscribe).(*packets.SubscribePacket)
|
||||
sub.Topics = append(sub.Topics, topic)
|
||||
sub.Qoss = append(sub.Qoss, qos)
|
||||
|
||||
if callback != nil {
|
||||
c.msgRouter.addRoute(topic, callback)
|
||||
}
|
||||
|
||||
sub.MessageID = c.mid
|
||||
c.mid++
|
||||
|
||||
// drop in the channel to send
|
||||
err := sub.Write(c.conn)
|
||||
if err != nil {
|
||||
return &mqtttoken{err: err}
|
||||
}
|
||||
|
||||
return &mqtttoken{}
|
||||
}
|
||||
|
||||
// SubscribeMultiple starts a new subscription for multiple topics. Provide a MessageHandler to
|
||||
// be executed when a message is published on one of the topics provided.
|
||||
func (c *mqttclient) SubscribeMultiple(filters map[string]byte, callback MessageHandler) Token {
|
||||
return &mqtttoken{}
|
||||
}
|
||||
|
||||
// Unsubscribe will end the subscription from each of the topics provided.
|
||||
// Messages published to those topics from other clients will no longer be
|
||||
// received.
|
||||
func (c *mqttclient) Unsubscribe(topics ...string) Token {
|
||||
return &mqtttoken{}
|
||||
}
|
||||
|
||||
// OptionsReader returns a ClientOptionsReader which is a copy of the clientoptions
|
||||
// in use by the client.
|
||||
func (c *mqttclient) OptionsReader() ClientOptionsReader {
|
||||
r := ClientOptionsReader{}
|
||||
return r
|
||||
}
|
||||
|
||||
func processInbound(c *mqttclient) {
|
||||
for {
|
||||
select {
|
||||
case msg := <-c.inbound:
|
||||
switch m := msg.(type) {
|
||||
case *packets.PingrespPacket:
|
||||
// TODO: handle this
|
||||
case *packets.SubackPacket:
|
||||
// TODO: handle this
|
||||
case *packets.UnsubackPacket:
|
||||
// TODO: handle this
|
||||
case *packets.PublishPacket:
|
||||
// TODO: handle Qos
|
||||
c.incomingPubChan <- m
|
||||
case *packets.PubackPacket:
|
||||
// TODO: handle this
|
||||
case *packets.PubrecPacket:
|
||||
// TODO: handle this
|
||||
case *packets.PubrelPacket:
|
||||
// TODO: handle this
|
||||
case *packets.PubcompPacket:
|
||||
// TODO: handle this
|
||||
}
|
||||
case <-c.stop:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// readMessages reads incoming messages off the wire.
|
||||
// incoming messages are then send into inbound channel.
|
||||
func readMessages(c *mqttclient) {
|
||||
var err error
|
||||
var cp packets.ControlPacket
|
||||
|
||||
PROCESS:
|
||||
for {
|
||||
if cp, err = c.ReadPacket(); err != nil {
|
||||
break PROCESS
|
||||
}
|
||||
if cp != nil {
|
||||
c.inbound <- cp
|
||||
// TODO: Notify keepalive logic that we recently received a packet
|
||||
}
|
||||
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
}
|
||||
|
||||
// TODO: handle if we received an error on read.
|
||||
// If disconnect is in progress, swallow error and return
|
||||
}
|
||||
|
||||
func (c *mqttclient) ackFunc(packet *packets.PublishPacket) func() {
|
||||
return func() {
|
||||
switch packet.Qos {
|
||||
case 2:
|
||||
// pr := packets.NewControlPacket(packets.Pubrec).(*packets.PubrecPacket)
|
||||
// pr.MessageID = packet.MessageID
|
||||
// DEBUG.Println(NET, "putting pubrec msg on obound")
|
||||
// select {
|
||||
// case c.oboundP <- &PacketAndToken{p: pr, t: nil}:
|
||||
// case <-c.stop:
|
||||
// }
|
||||
// DEBUG.Println(NET, "done putting pubrec msg on obound")
|
||||
case 1:
|
||||
// pa := packets.NewControlPacket(packets.Puback).(*packets.PubackPacket)
|
||||
// pa.MessageID = packet.MessageID
|
||||
// DEBUG.Println(NET, "putting puback msg on obound")
|
||||
// persistOutbound(c.persist, pa)
|
||||
// select {
|
||||
// case c.oboundP <- &PacketAndToken{p: pa, t: nil}:
|
||||
// case <-c.stop:
|
||||
// }
|
||||
// DEBUG.Println(NET, "done putting puback msg on obound")
|
||||
case 0:
|
||||
// do nothing, since there is no need to send an ack packet back
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ReadPacket tries to read the next incoming packet from the MQTT broker.
|
||||
// If there is no data yet but also is no error, it returns nil for both values.
|
||||
func (c *mqttclient) ReadPacket() (packets.ControlPacket, error) {
|
||||
// check for data first...
|
||||
if net.ActiveDevice.IsSocketDataAvailable() {
|
||||
return packets.ReadPacket(c.conn)
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
@@ -0,0 +1,280 @@
|
||||
// The following code is a slightly modified version of code taken from the Paho MQTT library.
|
||||
// It is here until TinyGo can compile the "net" package from the standard library, at which time
|
||||
// it can be removed.
|
||||
|
||||
/*
|
||||
* Copyright (c) 2013 IBM Corp.
|
||||
*
|
||||
* All rights reserved. This program and the accompanying materials
|
||||
* are made available under the terms of the Eclipse Public License v1.0
|
||||
* which accompanies this distribution, and is available at
|
||||
* http://www.eclipse.org/legal/epl-v10.html
|
||||
*
|
||||
* Contributors:
|
||||
* Seth Hoenig
|
||||
* Allan Stockdill-Mander
|
||||
* Mike Robertson
|
||||
*/
|
||||
|
||||
// Portions copyright © 2018 TIBCO Software Inc.
|
||||
|
||||
package mqtt
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/eclipse/paho.mqtt.golang/packets"
|
||||
"tinygo.org/x/drivers/net"
|
||||
)
|
||||
|
||||
const (
|
||||
disconnected uint32 = iota
|
||||
connecting
|
||||
reconnecting
|
||||
connected
|
||||
)
|
||||
|
||||
// Client is the interface definition for a Client as used by this
|
||||
// library, the interface is primarily to allow mocking tests.
|
||||
//
|
||||
// It is an MQTT v3.1.1 client for communicating
|
||||
// with an MQTT server using non-blocking methods that allow work
|
||||
// to be done in the background.
|
||||
// An application may connect to an MQTT server using:
|
||||
// A plain TCP socket
|
||||
// A secure SSL/TLS socket
|
||||
// A websocket
|
||||
// To enable ensured message delivery at Quality of Service (QoS) levels
|
||||
// described in the MQTT spec, a message persistence mechanism must be
|
||||
// used. This is done by providing a type which implements the Store
|
||||
// interface. For convenience, FileStore and MemoryStore are provided
|
||||
// implementations that should be sufficient for most use cases. More
|
||||
// information can be found in their respective documentation.
|
||||
// Numerous connection options may be specified by configuring a
|
||||
// and then supplying a ClientOptions type.
|
||||
type Client interface {
|
||||
// IsConnected returns a bool signifying whether
|
||||
// the client is connected or not.
|
||||
IsConnected() bool
|
||||
// IsConnectionOpen return a bool signifying wether the client has an active
|
||||
// connection to mqtt broker, i.e not in disconnected or reconnect mode
|
||||
IsConnectionOpen() bool
|
||||
// Connect will create a connection to the message broker, by default
|
||||
// it will attempt to connect at v3.1.1 and auto retry at v3.1 if that
|
||||
// fails
|
||||
Connect() Token
|
||||
// Disconnect will end the connection with the server, but not before waiting
|
||||
// the specified number of milliseconds to wait for existing work to be
|
||||
// completed.
|
||||
Disconnect(quiesce uint)
|
||||
// Publish will publish a message with the specified QoS and content
|
||||
// to the specified topic.
|
||||
// Returns a token to track delivery of the message to the broker
|
||||
Publish(topic string, qos byte, retained bool, payload interface{}) Token
|
||||
// Subscribe starts a new subscription. Provide a MessageHandler to be executed when
|
||||
// a message is published on the topic provided, or nil for the default handler
|
||||
Subscribe(topic string, qos byte, callback MessageHandler) Token
|
||||
// SubscribeMultiple starts a new subscription for multiple topics. Provide a MessageHandler to
|
||||
// be executed when a message is published on one of the topics provided, or nil for the
|
||||
// default handler
|
||||
SubscribeMultiple(filters map[string]byte, callback MessageHandler) Token
|
||||
// Unsubscribe will end the subscription from each of the topics provided.
|
||||
// Messages published to those topics from other clients will no longer be
|
||||
// received.
|
||||
Unsubscribe(topics ...string) Token
|
||||
// AddRoute allows you to add a handler for messages on a specific topic
|
||||
// without making a subscription. For example having a different handler
|
||||
// for parts of a wildcard subscription
|
||||
AddRoute(topic string, callback MessageHandler)
|
||||
// OptionsReader returns a ClientOptionsReader which is a copy of the clientoptions
|
||||
// in use by the client.
|
||||
OptionsReader() ClientOptionsReader
|
||||
}
|
||||
|
||||
// Token defines the interface for the tokens used to indicate when
|
||||
// actions have completed.
|
||||
type Token interface {
|
||||
Wait() bool
|
||||
WaitTimeout(time.Duration) bool
|
||||
Error() error
|
||||
}
|
||||
|
||||
// MessageHandler is a callback type which can be set to be
|
||||
// executed upon the arrival of messages published to topics
|
||||
// to which the client is subscribed.
|
||||
type MessageHandler func(Client, Message)
|
||||
|
||||
// Message defines the externals that a message implementation must support
|
||||
// these are received messages that are passed to the callbacks, not internal
|
||||
// messages
|
||||
type Message interface {
|
||||
Duplicate() bool
|
||||
Qos() byte
|
||||
Retained() bool
|
||||
Topic() string
|
||||
MessageID() uint16
|
||||
Payload() []byte
|
||||
Ack()
|
||||
}
|
||||
|
||||
type message struct {
|
||||
duplicate bool
|
||||
qos byte
|
||||
retained bool
|
||||
topic string
|
||||
messageID uint16
|
||||
payload []byte
|
||||
ack func()
|
||||
}
|
||||
|
||||
func (m *message) Duplicate() bool {
|
||||
return m.duplicate
|
||||
}
|
||||
|
||||
func (m *message) Qos() byte {
|
||||
return m.qos
|
||||
}
|
||||
|
||||
func (m *message) Retained() bool {
|
||||
return m.retained
|
||||
}
|
||||
|
||||
func (m *message) Topic() string {
|
||||
return m.topic
|
||||
}
|
||||
|
||||
func (m *message) MessageID() uint16 {
|
||||
return m.messageID
|
||||
}
|
||||
|
||||
func (m *message) Payload() []byte {
|
||||
return m.payload
|
||||
}
|
||||
|
||||
func (m *message) Ack() {
|
||||
return
|
||||
}
|
||||
|
||||
func messageFromPublish(p *packets.PublishPacket, ack func()) Message {
|
||||
return &message{
|
||||
duplicate: p.Dup,
|
||||
qos: p.Qos,
|
||||
retained: p.Retain,
|
||||
topic: p.TopicName,
|
||||
messageID: p.MessageID,
|
||||
payload: p.Payload,
|
||||
ack: ack,
|
||||
}
|
||||
}
|
||||
|
||||
// ClientOptionsReader provides an interface for reading ClientOptions after the client has been initialized.
|
||||
type ClientOptionsReader struct {
|
||||
options *ClientOptions
|
||||
}
|
||||
|
||||
// ClientOptions contains configurable options for an MQTT Client.
|
||||
type ClientOptions struct {
|
||||
Adaptor net.DeviceDriver
|
||||
|
||||
//Servers []*url.URL
|
||||
Servers string
|
||||
ClientID string
|
||||
Username string
|
||||
Password string
|
||||
//CredentialsProvider CredentialsProvider
|
||||
CleanSession bool
|
||||
Order bool
|
||||
WillEnabled bool
|
||||
WillTopic string
|
||||
WillPayload []byte
|
||||
WillQos byte
|
||||
WillRetained bool
|
||||
ProtocolVersion uint
|
||||
protocolVersionExplicit bool
|
||||
//TLSConfig *tls.Config
|
||||
KeepAlive int64
|
||||
PingTimeout time.Duration
|
||||
ConnectTimeout time.Duration
|
||||
MaxReconnectInterval time.Duration
|
||||
AutoReconnect bool
|
||||
//Store Store
|
||||
//DefaultPublishHandler MessageHandler
|
||||
//OnConnect OnConnectHandler
|
||||
//OnConnectionLost ConnectionLostHandler
|
||||
WriteTimeout time.Duration
|
||||
MessageChannelDepth uint
|
||||
ResumeSubs bool
|
||||
//HTTPHeaders http.Header
|
||||
}
|
||||
|
||||
// NewClientOptions returns a new ClientOptions struct.
|
||||
func NewClientOptions() *ClientOptions {
|
||||
return &ClientOptions{Adaptor: net.ActiveDevice, ProtocolVersion: 4}
|
||||
}
|
||||
|
||||
// AddBroker adds a broker URI to the list of brokers to be used. The format should be
|
||||
// scheme://host:port
|
||||
// Where "scheme" is one of "tcp", "ssl", or "ws", "host" is the ip-address (or hostname)
|
||||
// and "port" is the port on which the broker is accepting connections.
|
||||
//
|
||||
// Default values for hostname is "127.0.0.1", for schema is "tcp://".
|
||||
//
|
||||
// An example broker URI would look like: tcp://foobar.com:1883
|
||||
func (o *ClientOptions) AddBroker(server string) *ClientOptions {
|
||||
if len(server) > 0 && server[0] == ':' {
|
||||
server = "127.0.0.1" + server
|
||||
}
|
||||
if !strings.Contains(server, "://") {
|
||||
server = "tcp://" + server
|
||||
}
|
||||
|
||||
o.Servers = server
|
||||
return o
|
||||
}
|
||||
|
||||
// SetClientID will set the client id to be used by this client when
|
||||
// connecting to the MQTT broker. According to the MQTT v3.1 specification,
|
||||
// a client id mus be no longer than 23 characters.
|
||||
func (o *ClientOptions) SetClientID(id string) *ClientOptions {
|
||||
o.ClientID = id
|
||||
return o
|
||||
}
|
||||
|
||||
// SetUsername will set the username to be used by this client when connecting
|
||||
// to the MQTT broker. Note: without the use of SSL/TLS, this information will
|
||||
// be sent in plaintext accross the wire.
|
||||
func (o *ClientOptions) SetUsername(u string) *ClientOptions {
|
||||
o.Username = u
|
||||
return o
|
||||
}
|
||||
|
||||
// SetPassword will set the password to be used by this client when connecting
|
||||
// to the MQTT broker. Note: without the use of SSL/TLS, this information will
|
||||
// be sent in plaintext accross the wire.
|
||||
func (o *ClientOptions) SetPassword(p string) *ClientOptions {
|
||||
o.Password = p
|
||||
return o
|
||||
}
|
||||
|
||||
// SetWill accepts a string will message to be set. When the client connects,
|
||||
// it will give this will message to the broker, which will then publish the
|
||||
// provided payload (the will) to any clients that are subscribed to the provided
|
||||
// topic.
|
||||
func (o *ClientOptions) SetWill(topic string, payload string, qos byte, retained bool) *ClientOptions {
|
||||
o.SetBinaryWill(topic, []byte(payload), qos, retained)
|
||||
return o
|
||||
}
|
||||
|
||||
// SetBinaryWill accepts a []byte will message to be set. When the client connects,
|
||||
// it will give this will message to the broker, which will then publish the
|
||||
// provided payload (the will) to any clients that are subscribed to the provided
|
||||
// topic.
|
||||
func (o *ClientOptions) SetBinaryWill(topic string, payload []byte, qos byte, retained bool) *ClientOptions {
|
||||
o.WillEnabled = true
|
||||
o.WillTopic = topic
|
||||
o.WillPayload = payload
|
||||
o.WillQos = qos
|
||||
o.WillRetained = retained
|
||||
return o
|
||||
}
|
||||
@@ -0,0 +1,182 @@
|
||||
// The following code is a slightly modified version of code taken from the Paho MQTT library.
|
||||
// It is here until TinyGo can compile the "net" package from the standard library, at which time
|
||||
// it can be removed.
|
||||
|
||||
/*
|
||||
* Copyright (c) 2013 IBM Corp.
|
||||
*
|
||||
* All rights reserved. This program and the accompanying materials
|
||||
* are made available under the terms of the Eclipse Public License v1.0
|
||||
* which accompanies this distribution, and is available at
|
||||
* http://www.eclipse.org/legal/epl-v10.html
|
||||
*
|
||||
* Contributors:
|
||||
* Seth Hoenig
|
||||
* Allan Stockdill-Mander
|
||||
* Mike Robertson
|
||||
*/
|
||||
|
||||
package mqtt
|
||||
|
||||
import (
|
||||
"container/list"
|
||||
"strings"
|
||||
|
||||
"github.com/eclipse/paho.mqtt.golang/packets"
|
||||
)
|
||||
|
||||
// route is a type which associates MQTT Topic strings with a
|
||||
// callback to be executed upon the arrival of a message associated
|
||||
// with a subscription to that topic.
|
||||
type route struct {
|
||||
topic string
|
||||
callback MessageHandler
|
||||
}
|
||||
|
||||
// match takes a slice of strings which represent the route being tested having been split on '/'
|
||||
// separators, and a slice of strings representing the topic string in the published message, similarly
|
||||
// split.
|
||||
// The function determines if the topic string matches the route according to the MQTT topic rules
|
||||
// and returns a boolean of the outcome
|
||||
func match(route []string, topic []string) bool {
|
||||
if len(route) == 0 {
|
||||
if len(topic) == 0 {
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
if len(topic) == 0 {
|
||||
if route[0] == "#" {
|
||||
return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
if route[0] == "#" {
|
||||
return true
|
||||
}
|
||||
|
||||
if (route[0] == "+") || (route[0] == topic[0]) {
|
||||
return match(route[1:], topic[1:])
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func routeIncludesTopic(route, topic string) bool {
|
||||
return match(routeSplit(route), strings.Split(topic, "/"))
|
||||
}
|
||||
|
||||
// removes $share and sharename when splitting the route to allow
|
||||
// shared subscription routes to correctly match the topic
|
||||
func routeSplit(route string) []string {
|
||||
var result []string
|
||||
if strings.HasPrefix(route, "$share") {
|
||||
result = strings.Split(route, "/")[2:]
|
||||
} else {
|
||||
result = strings.Split(route, "/")
|
||||
}
|
||||
return result
|
||||
}
|
||||
|
||||
// match takes the topic string of the published message and does a basic compare to the
|
||||
// string of the current Route, if they match it returns true
|
||||
func (r *route) match(topic string) bool {
|
||||
return r.topic == topic || routeIncludesTopic(r.topic, topic)
|
||||
}
|
||||
|
||||
type router struct {
|
||||
//sync.RWMutex
|
||||
routes *list.List
|
||||
defaultHandler MessageHandler
|
||||
messages chan *packets.PublishPacket
|
||||
stop chan bool
|
||||
}
|
||||
|
||||
// newRouter returns a new instance of a Router and channel which can be used to tell the Router
|
||||
// to stop
|
||||
func newRouter() (*router, chan bool) {
|
||||
router := &router{routes: list.New(), messages: make(chan *packets.PublishPacket), stop: make(chan bool)}
|
||||
stop := router.stop
|
||||
return router, stop
|
||||
}
|
||||
|
||||
// addRoute takes a topic string and MessageHandler callback. It looks in the current list of
|
||||
// routes to see if there is already a matching Route. If there is it replaces the current
|
||||
// callback with the new one. If not it add a new entry to the list of Routes.
|
||||
func (r *router) addRoute(topic string, callback MessageHandler) {
|
||||
for e := r.routes.Front(); e != nil; e = e.Next() {
|
||||
if e.Value.(*route).match(topic) {
|
||||
r := e.Value.(*route)
|
||||
r.callback = callback
|
||||
return
|
||||
}
|
||||
}
|
||||
r.routes.PushBack(&route{topic: topic, callback: callback})
|
||||
}
|
||||
|
||||
// deleteRoute takes a route string, looks for a matching Route in the list of Routes. If
|
||||
// found it removes the Route from the list.
|
||||
func (r *router) deleteRoute(topic string) {
|
||||
for e := r.routes.Front(); e != nil; e = e.Next() {
|
||||
if e.Value.(*route).match(topic) {
|
||||
r.routes.Remove(e)
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// setDefaultHandler assigns a default callback that will be called if no matching Route
|
||||
// is found for an incoming Publish.
|
||||
func (r *router) setDefaultHandler(handler MessageHandler) {
|
||||
r.defaultHandler = handler
|
||||
}
|
||||
|
||||
// matchAndDispatch takes a channel of Message pointers as input and starts a go routine that
|
||||
// takes messages off the channel, matches them against the internal route list and calls the
|
||||
// associated callback (or the defaultHandler, if one exists and no other route matched). If
|
||||
// anything is sent down the stop channel the function will end.
|
||||
func (r *router) matchAndDispatch(messages <-chan *packets.PublishPacket, order bool, client *mqttclient) {
|
||||
go func() {
|
||||
for {
|
||||
select {
|
||||
case message := <-messages:
|
||||
sent := false
|
||||
m := messageFromPublish(message, client.ackFunc(message))
|
||||
handlers := []MessageHandler{}
|
||||
for e := r.routes.Front(); e != nil; e = e.Next() {
|
||||
if e.Value.(*route).match(message.TopicName) {
|
||||
if order {
|
||||
handlers = append(handlers, e.Value.(*route).callback)
|
||||
} else {
|
||||
hd := e.Value.(*route).callback
|
||||
go func() {
|
||||
hd(client, m)
|
||||
//TODO: m.Ack()
|
||||
}()
|
||||
}
|
||||
sent = true
|
||||
}
|
||||
}
|
||||
if !sent && r.defaultHandler != nil {
|
||||
if order {
|
||||
handlers = append(handlers, r.defaultHandler)
|
||||
} else {
|
||||
go func() {
|
||||
r.defaultHandler(client, m)
|
||||
//TODO: m.Ack()
|
||||
}()
|
||||
}
|
||||
}
|
||||
for _, handler := range handlers {
|
||||
func() {
|
||||
handler(client, m)
|
||||
//TODO: m.Ack()
|
||||
}()
|
||||
}
|
||||
case <-r.stop:
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
package mqtt
|
||||
|
||||
import "time"
|
||||
|
||||
type mqtttoken struct {
|
||||
err error
|
||||
}
|
||||
|
||||
func (t *mqtttoken) Wait() bool {
|
||||
return true
|
||||
}
|
||||
|
||||
func (t *mqtttoken) WaitTimeout(time.Duration) bool {
|
||||
return true
|
||||
}
|
||||
|
||||
func (t *mqtttoken) Error() error {
|
||||
return t.err
|
||||
}
|
||||
+410
@@ -0,0 +1,410 @@
|
||||
// package net is intended to provide compatible interfaces with the
|
||||
// Go standard library's net package.
|
||||
package net
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// DialUDP makes a UDP network connection. raadr is the port that the messages will
|
||||
// be sent to, and laddr is the port that will be listened to in order to
|
||||
// receive incoming messages.
|
||||
func DialUDP(network string, laddr, raddr *UDPAddr) (*UDPSerialConn, error) {
|
||||
addr := raddr.IP.String()
|
||||
sendport := strconv.Itoa(raddr.Port)
|
||||
listenport := strconv.Itoa(laddr.Port)
|
||||
|
||||
// disconnect any old socket
|
||||
ActiveDevice.DisconnectSocket()
|
||||
|
||||
// connect new socket
|
||||
err := ActiveDevice.ConnectUDPSocket(addr, sendport, listenport)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &UDPSerialConn{SerialConn: SerialConn{Adaptor: ActiveDevice}, laddr: laddr, raddr: raddr}, nil
|
||||
}
|
||||
|
||||
// ListenUDP listens for UDP connections on the port listed in laddr.
|
||||
func ListenUDP(network string, laddr *UDPAddr) (*UDPSerialConn, error) {
|
||||
addr := "0"
|
||||
sendport := "0"
|
||||
listenport := strconv.Itoa(laddr.Port)
|
||||
|
||||
// disconnect any old socket
|
||||
ActiveDevice.DisconnectSocket()
|
||||
|
||||
// connect new socket
|
||||
err := ActiveDevice.ConnectUDPSocket(addr, sendport, listenport)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &UDPSerialConn{SerialConn: SerialConn{Adaptor: ActiveDevice}, laddr: laddr}, nil
|
||||
}
|
||||
|
||||
// DialTCP makes a TCP network connection. raadr is the port that the messages will
|
||||
// be sent to, and laddr is the port that will be listened to in order to
|
||||
// receive incoming messages.
|
||||
func DialTCP(network string, laddr, raddr *TCPAddr) (*TCPSerialConn, error) {
|
||||
addr := raddr.IP.String()
|
||||
sendport := strconv.Itoa(raddr.Port)
|
||||
|
||||
// disconnect any old socket?
|
||||
//ActiveDevice.DisconnectSocket()
|
||||
|
||||
// connect new socket
|
||||
err := ActiveDevice.ConnectTCPSocket(addr, sendport)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &TCPSerialConn{SerialConn: SerialConn{Adaptor: ActiveDevice}, laddr: laddr, raddr: raddr}, nil
|
||||
}
|
||||
|
||||
// Dial connects to the address on the named network.
|
||||
// It tries to provide a mostly compatible interface
|
||||
// to net.Dial().
|
||||
func Dial(network, address string) (Conn, error) {
|
||||
switch network {
|
||||
case "tcp":
|
||||
raddr, err := ResolveTCPAddr(network, address)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
c, e := DialTCP(network, &TCPAddr{}, raddr)
|
||||
return c.opConn(), e
|
||||
case "udp":
|
||||
raddr, err := ResolveUDPAddr(network, address)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
c, e := DialUDP(network, &UDPAddr{}, raddr)
|
||||
return c.opConn(), e
|
||||
default:
|
||||
return nil, errors.New("invalid network for dial")
|
||||
}
|
||||
}
|
||||
|
||||
// SerialConn is a loosely net.Conn compatible implementation
|
||||
type SerialConn struct {
|
||||
Adaptor DeviceDriver
|
||||
}
|
||||
|
||||
// UDPSerialConn is a loosely net.Conn compatible intended to support
|
||||
// UDP over serial.
|
||||
type UDPSerialConn struct {
|
||||
SerialConn
|
||||
laddr *UDPAddr
|
||||
raddr *UDPAddr
|
||||
}
|
||||
|
||||
// NewUDPSerialConn returns a new UDPSerialConn/
|
||||
func NewUDPSerialConn(c SerialConn, laddr, raddr *UDPAddr) *UDPSerialConn {
|
||||
return &UDPSerialConn{SerialConn: c, raddr: raddr}
|
||||
}
|
||||
|
||||
// TCPSerialConn is a loosely net.Conn compatible intended to support
|
||||
// TCP over serial.
|
||||
type TCPSerialConn struct {
|
||||
SerialConn
|
||||
laddr *TCPAddr
|
||||
raddr *TCPAddr
|
||||
}
|
||||
|
||||
// NewTCPSerialConn returns a new TCPSerialConn/
|
||||
func NewTCPSerialConn(c SerialConn, laddr, raddr *TCPAddr) *TCPSerialConn {
|
||||
return &TCPSerialConn{SerialConn: c, raddr: raddr}
|
||||
}
|
||||
|
||||
// Read reads data from the connection.
|
||||
// TODO: implement the full method functionality:
|
||||
// Read can be made to time out and return an Error with Timeout() == true
|
||||
// after a fixed time limit; see SetDeadline and SetReadDeadline.
|
||||
func (c *SerialConn) Read(b []byte) (n int, err error) {
|
||||
// read only the data that has been received via "+IPD" socket
|
||||
return c.Adaptor.ReadSocket(b)
|
||||
}
|
||||
|
||||
// Write writes data to the connection.
|
||||
// TODO: implement the full method functionality for timeouts.
|
||||
// Write can be made to time out and return an Error with Timeout() == true
|
||||
// after a fixed time limit; see SetDeadline and SetWriteDeadline.
|
||||
func (c *SerialConn) Write(b []byte) (n int, err error) {
|
||||
// specify that is a data transfer to the
|
||||
// currently open socket, not commands to the ESP8266/ESP32.
|
||||
err = c.Adaptor.StartSocketSend(len(b))
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
n, err = c.Adaptor.Write(b)
|
||||
if err != nil {
|
||||
return n, err
|
||||
}
|
||||
/* TODO(bcg): this is kind of specific to espat, should maybe refactor */
|
||||
_, err = c.Adaptor.Response(1000)
|
||||
if err != nil {
|
||||
return n, err
|
||||
}
|
||||
return n, err
|
||||
}
|
||||
|
||||
// Close closes the connection.
|
||||
// Currently only supports a single Read or Write operations without blocking.
|
||||
func (c *SerialConn) Close() error {
|
||||
c.Adaptor.DisconnectSocket()
|
||||
return nil
|
||||
}
|
||||
|
||||
// LocalAddr returns the local network address.
|
||||
func (c *UDPSerialConn) LocalAddr() Addr {
|
||||
return c.laddr.opAddr()
|
||||
}
|
||||
|
||||
// RemoteAddr returns the remote network address.
|
||||
func (c *UDPSerialConn) RemoteAddr() Addr {
|
||||
return c.laddr.opAddr()
|
||||
}
|
||||
|
||||
func (c *UDPSerialConn) opConn() Conn {
|
||||
if c == nil {
|
||||
return nil
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
// LocalAddr returns the local network address.
|
||||
func (c *TCPSerialConn) LocalAddr() Addr {
|
||||
return c.laddr.opAddr()
|
||||
}
|
||||
|
||||
// RemoteAddr returns the remote network address.
|
||||
func (c *TCPSerialConn) RemoteAddr() Addr {
|
||||
return c.laddr.opAddr()
|
||||
}
|
||||
|
||||
func (c *TCPSerialConn) opConn() Conn {
|
||||
if c == nil {
|
||||
return nil
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
// SetDeadline sets the read and write deadlines associated
|
||||
// with the connection. It is equivalent to calling both
|
||||
// SetReadDeadline and SetWriteDeadline.
|
||||
//
|
||||
// A deadline is an absolute time after which I/O operations
|
||||
// fail with a timeout (see type Error) instead of
|
||||
// blocking. The deadline applies to all future and pending
|
||||
// I/O, not just the immediately following call to Read or
|
||||
// Write. After a deadline has been exceeded, the connection
|
||||
// can be refreshed by setting a deadline in the future.
|
||||
//
|
||||
// An idle timeout can be implemented by repeatedly extending
|
||||
// the deadline after successful Read or Write calls.
|
||||
//
|
||||
// A zero value for t means I/O operations will not time out.
|
||||
func (c *SerialConn) SetDeadline(t time.Time) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// SetReadDeadline sets the deadline for future Read calls
|
||||
// and any currently-blocked Read call.
|
||||
// A zero value for t means Read will not time out.
|
||||
func (c *SerialConn) SetReadDeadline(t time.Time) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// SetWriteDeadline sets the deadline for future Write calls
|
||||
// and any currently-blocked Write call.
|
||||
// Even if write times out, it may return n > 0, indicating that
|
||||
// some of the data was successfully written.
|
||||
// A zero value for t means Write will not time out.
|
||||
func (c *SerialConn) SetWriteDeadline(t time.Time) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// ResolveTCPAddr returns an address of TCP end point.
|
||||
//
|
||||
// The network must be a TCP network name.
|
||||
//
|
||||
func ResolveTCPAddr(network, address string) (*TCPAddr, error) {
|
||||
// TODO: make sure network is 'tcp'
|
||||
// separate domain from port, if any
|
||||
r := strings.Split(address, ":")
|
||||
addr, err := ActiveDevice.GetDNS(r[0])
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
ip := IP(addr)
|
||||
if len(r) > 1 {
|
||||
port, e := strconv.Atoi(r[1])
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
return &TCPAddr{IP: ip, Port: port}, nil
|
||||
}
|
||||
return &TCPAddr{IP: ip}, nil
|
||||
}
|
||||
|
||||
// ResolveUDPAddr returns an address of UDP end point.
|
||||
//
|
||||
// The network must be a UDP network name.
|
||||
//
|
||||
func ResolveUDPAddr(network, address string) (*UDPAddr, error) {
|
||||
// TODO: make sure network is 'udp'
|
||||
// separate domain from port, if any
|
||||
r := strings.Split(address, ":")
|
||||
addr, err := ActiveDevice.GetDNS(r[0])
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
ip := IP(addr)
|
||||
if len(r) > 1 {
|
||||
port, e := strconv.Atoi(r[1])
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
return &UDPAddr{IP: ip, Port: port}, nil
|
||||
}
|
||||
|
||||
return &UDPAddr{IP: ip}, nil
|
||||
}
|
||||
|
||||
// The following definitions are here to support a Golang standard package
|
||||
// net-compatible interface for IP until TinyGo can compile the net package.
|
||||
|
||||
// IP is an IP address. Unlike the standard implementation, it is only
|
||||
// a buffer of bytes that contains the string form of the IP address, not the
|
||||
// full byte format used by the Go standard .
|
||||
type IP []byte
|
||||
|
||||
// UDPAddr here to serve as compatible type. until TinyGo can compile the net package.
|
||||
type UDPAddr struct {
|
||||
IP IP
|
||||
Port int
|
||||
Zone string // IPv6 scoped addressing zone; added in Go 1.1
|
||||
}
|
||||
|
||||
// Network returns the address's network name, "udp".
|
||||
func (a *UDPAddr) Network() string { return "udp" }
|
||||
|
||||
func (a *UDPAddr) String() string {
|
||||
if a == nil {
|
||||
return "<nil>"
|
||||
}
|
||||
if a.Port != 0 {
|
||||
return a.IP.String() + ":" + strconv.Itoa(a.Port)
|
||||
}
|
||||
return a.IP.String()
|
||||
}
|
||||
|
||||
func (a *UDPAddr) opAddr() Addr {
|
||||
if a == nil {
|
||||
return nil
|
||||
}
|
||||
return a
|
||||
}
|
||||
|
||||
// TCPAddr here to serve as compatible type. until TinyGo can compile the net package.
|
||||
type TCPAddr struct {
|
||||
IP IP
|
||||
Port int
|
||||
Zone string // IPv6 scoped addressing zone
|
||||
}
|
||||
|
||||
// Network returns the address's network name, "tcp".
|
||||
func (a *TCPAddr) Network() string { return "tcp" }
|
||||
|
||||
func (a *TCPAddr) String() string {
|
||||
if a == nil {
|
||||
return "<nil>"
|
||||
}
|
||||
if a.Port != 0 {
|
||||
return a.IP.String() + ":" + strconv.Itoa(a.Port)
|
||||
}
|
||||
return a.IP.String()
|
||||
}
|
||||
|
||||
func (a *TCPAddr) opAddr() Addr {
|
||||
if a == nil {
|
||||
return nil
|
||||
}
|
||||
return a
|
||||
}
|
||||
|
||||
// ParseIP parses s as an IP address, returning the result.
|
||||
func ParseIP(s string) IP {
|
||||
return IP([]byte(s))
|
||||
}
|
||||
|
||||
// String returns the string form of the IP address ip.
|
||||
func (ip IP) String() string {
|
||||
return string(ip)
|
||||
}
|
||||
|
||||
// Conn is a generic stream-oriented network connection.
|
||||
// This interface is from the Go standard library.
|
||||
type Conn interface {
|
||||
// Read reads data from the connection.
|
||||
// Read can be made to time out and return an Error with Timeout() == true
|
||||
// after a fixed time limit; see SetDeadline and SetReadDeadline.
|
||||
Read(b []byte) (n int, err error)
|
||||
|
||||
// Write writes data to the connection.
|
||||
// Write can be made to time out and return an Error with Timeout() == true
|
||||
// after a fixed time limit; see SetDeadline and SetWriteDeadline.
|
||||
Write(b []byte) (n int, err error)
|
||||
|
||||
// Close closes the connection.
|
||||
// Any blocked Read or Write operations will be unblocked and return errors.
|
||||
Close() error
|
||||
|
||||
// LocalAddr returns the local network address.
|
||||
LocalAddr() Addr
|
||||
|
||||
// RemoteAddr returns the remote network address.
|
||||
RemoteAddr() Addr
|
||||
|
||||
// SetDeadline sets the read and write deadlines associated
|
||||
// with the connection. It is equivalent to calling both
|
||||
// SetReadDeadline and SetWriteDeadline.
|
||||
//
|
||||
// A deadline is an absolute time after which I/O operations
|
||||
// fail with a timeout (see type Error) instead of
|
||||
// blocking. The deadline applies to all future and pending
|
||||
// I/O, not just the immediately following call to Read or
|
||||
// Write. After a deadline has been exceeded, the connection
|
||||
// can be refreshed by setting a deadline in the future.
|
||||
//
|
||||
// An idle timeout can be implemented by repeatedly extending
|
||||
// the deadline after successful Read or Write calls.
|
||||
//
|
||||
// A zero value for t means I/O operations will not time out.
|
||||
SetDeadline(t time.Time) error
|
||||
|
||||
// SetReadDeadline sets the deadline for future Read calls
|
||||
// and any currently-blocked Read call.
|
||||
// A zero value for t means Read will not time out.
|
||||
SetReadDeadline(t time.Time) error
|
||||
|
||||
// SetWriteDeadline sets the deadline for future Write calls
|
||||
// and any currently-blocked Write call.
|
||||
// Even if write times out, it may return n > 0, indicating that
|
||||
// some of the data was successfully written.
|
||||
// A zero value for t means Write will not time out.
|
||||
SetWriteDeadline(t time.Time) error
|
||||
}
|
||||
|
||||
// Addr represents a network end point address.
|
||||
type Addr interface {
|
||||
Network() string // name of the network (for example, "tcp", "udp")
|
||||
String() string // string form of address (for example, "192.0.2.1:25", "[2001:db8::1]:80")
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
// Package tls is intended to provide a minimal set of compatible interfaces with the
|
||||
// Go standard library's tls package.
|
||||
package tls
|
||||
|
||||
import (
|
||||
"strconv"
|
||||
|
||||
"tinygo.org/x/drivers/net"
|
||||
)
|
||||
|
||||
// Dial makes a TLS network connection. It tries to provide a mostly compatible interface
|
||||
// to tls.Dial().
|
||||
// Dial connects to the given network address.
|
||||
func Dial(network, address string, config *Config) (*net.TCPSerialConn, error) {
|
||||
raddr, err := net.ResolveTCPAddr(network, address)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
addr := raddr.IP.String()
|
||||
sendport := strconv.Itoa(raddr.Port)
|
||||
|
||||
// disconnect any old socket
|
||||
net.ActiveDevice.DisconnectSocket()
|
||||
|
||||
// connect new socket
|
||||
err = net.ActiveDevice.ConnectSSLSocket(addr, sendport)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return net.NewTCPSerialConn(net.SerialConn{Adaptor: net.ActiveDevice}, nil, raddr), nil
|
||||
}
|
||||
|
||||
// Config is a placeholder for future compatibility with
|
||||
// tls.Config.
|
||||
type Config struct {
|
||||
}
|
||||
Reference in New Issue
Block a user