From 7dcbfbecc6a02ebb8c6f886eb61ebe3c4a3af22f Mon Sep 17 00:00:00 2001 From: Ron Evans Date: Wed, 4 Sep 2019 12:13:07 +0200 Subject: [PATCH] espat: implement MQTT subscribe functionality via blocking select/channels. also refactor response processing for greater speed and efficiency. Signed-off-by: Ron Evans --- espat/espat.go | 150 +++++++++--------------- espat/mqtt/mqtt.go | 174 +++++++++++++++++++++++----- espat/mqtt/paho.go | 13 +++ espat/mqtt/router.go | 182 ++++++++++++++++++++++++++++++ espat/mqtt/token.go | 19 ++++ espat/net/net.go | 34 ++++-- espat/tcp.go | 74 +++++++----- espat/wifi.go | 70 +++++++----- examples/espat/espconsole/main.go | 75 +++++++----- examples/espat/esphub/main.go | 49 +++++--- examples/espat/espstation/main.go | 37 +++++- examples/espat/mqttclient/main.go | 16 ++- examples/espat/mqttsub/main.go | 163 ++++++++++++++++++++++++++ examples/espat/tcpclient/main.go | 42 ++++++- 14 files changed, 851 insertions(+), 247 deletions(-) create mode 100644 espat/mqtt/router.go create mode 100644 espat/mqtt/token.go create mode 100644 examples/espat/mqttsub/main.go diff --git a/espat/espat.go b/espat/espat.go index d8b7a50..3a5ff39 100644 --- a/espat/espat.go +++ b/espat/espat.go @@ -19,6 +19,7 @@ package espat // import "tinygo.org/x/drivers/espat" import ( + "errors" "machine" "strconv" "strings" @@ -54,11 +55,11 @@ func (d *Device) Connected() bool { d.Execute(Test) // handle response here, should include "OK" - r := d.Response(100) - if strings.Contains(string(r), "OK") { - return true + _, err := d.Response(100) + if err != nil { + return false } - return false + return true } // Write raw bytes to the UART. @@ -72,7 +73,7 @@ func (d *Device) Read(b []byte) (n int, err error) { } // how long in milliseconds to pause after sending AT commands -const pause = 100 +const pause = 300 // Execute sends an AT command to the ESP8266/ESP32. func (d Device) Execute(cmd string) error { @@ -97,7 +98,11 @@ func (d Device) Set(cmd, params string) error { // Version returns the ESP8266/ESP32 firmware version info. func (d Device) Version() []byte { d.Execute(Version) - return d.Response(100) + r, err := d.Response(100) + if err != nil { + return []byte("unknown") + } + return r } // Echo sets the ESP8266/ESP32 echo setting. @@ -122,7 +127,7 @@ func (d Device) Reset() { // ReadSocket returns the data that has already been read in from the responses. func (d *Device) ReadSocket(b []byte) (n int, err error) { // make sure no data in buffer - d.Response(100) + d.Response(300) count := len(b) if len(b) >= len(d.socketdata) { @@ -142,118 +147,73 @@ func (d *Device) ReadSocket(b []byte) (n int, err error) { // Response gets the next response bytes from the ESP8266/ESP32. // The call will retry for up to timeout milliseconds before returning nothing. -func (d *Device) Response(timeout int) []byte { - var i int - pause := 10 // pause to wait for 10 ms +func (d *Device) Response(timeout int) ([]byte, error) { + // read data + var size int + var start, end int + pause := 100 // pause to wait for 100 ms retries := timeout / pause - header := make([]byte, 2) for { - for d.bus.Buffered() > 0 { - // get the first 2 bytes - header[0], _ = d.bus.ReadByte() - header[1], _ = d.bus.ReadByte() + size = d.bus.Buffered() - if d.isLeadingCRLF(header) { - // skip it - header[0], _ = d.bus.ReadByte() - header[1], _ = d.bus.ReadByte() + if size > 0 { + end += size + d.bus.Read(d.response[start:end]) + + // if "+IPD" then read socket data + if strings.Contains(string(d.response[:end]), "+IPD") { + // handle socket data + return nil, d.parseIPD(end) } - if d.isIPD(header) { - // is socket data packet - d.parseIPD() - } else { - // no, so put into response - d.response[i] = header[0] - i++ - d.response[i] = header[1] - i++ + // if "OK" then the command worked + if strings.Contains(string(d.response[:end]), "OK") { + return d.response[start:end], nil } - // read the rest of normal command response - for d.bus.Buffered() > 0 { - data, err := d.bus.ReadByte() - if err != nil { - return nil - } - d.response[i] = data - i++ + // if "Error" then the command failed + if strings.Contains(string(d.response[:end]), "ERROR") { + return d.response[start:end], errors.New("response error:" + string(d.response[start:end])) } + + // if anything else, then keep reading data in? + start = end } + + // wait longer? retries-- if retries == 0 { - break + return nil, errors.New("response timeout error:" + string(d.response[start:end])) } - // pause to make sure is no more data to be read time.Sleep(time.Duration(pause) * time.Millisecond) } - return d.response[:i] } -func (d *Device) isLeadingCRLF(b []byte) bool { - if len(b) < 2 { - return false - } - if b[0] == 13 && b[1] == 10 { - return true - } - return false -} +func (d *Device) parseIPD(end int) error { + // find the "+IPD," to get length + s := strings.Index(string(d.response[:end]), "+IPD,") -func (d *Device) isIPD(b []byte) bool { - if len(b) < 2 { - return false - } - if b[0] == '+' && b[1] == 'I' { - return true - } - return false -} + // find the ":" + e := strings.Index(string(d.response[:end]), ":") -func (d *Device) parseIPD() bool { - data, _ := d.bus.ReadByte() - if data != 'P' { - // error - return false - } - data, _ = d.bus.ReadByte() - if data != 'D' { - // error - return false - } - data, _ = d.bus.ReadByte() - if data != ',' { - // error - return false - } + // find the data length + val := string(d.response[s+5 : e]) - // get the expected data length - // skip remaining header up to the ":" - buf := []byte{} - data, _ = d.bus.ReadByte() - for data != ':' { - // put into the buffer with int value here - buf = append(buf, data) - - // read next value - data, _ = d.bus.ReadByte() - } - - val := string(buf) - count, err := strconv.Atoi(val) + // TODO: verify count + _, err := strconv.Atoi(val) if err != nil { // not expected data here. what to do? - return false + return err } // load up the socket data - // only read the expected amount of data - for m := 0; m < count; m++ { - data, _ = d.bus.ReadByte() - d.socketdata = append(d.socketdata, data) - } - - return true + d.socketdata = append(d.socketdata, d.response[e+1:end]...) + return nil +} + +// IsSocketDataAvailable returns of there is socket data available +func (d *Device) IsSocketDataAvailable() bool { + return len(d.socketdata) > 0 || d.bus.Buffered() > 0 } diff --git a/espat/mqtt/mqtt.go b/espat/mqtt/mqtt.go index 57bda0b..65d9edd 100644 --- a/espat/mqtt/mqtt.go +++ b/espat/mqtt/mqtt.go @@ -19,15 +19,21 @@ import ( // 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 *espat.Device - conn net.Conn - connected bool - opts *ClientOptions - mid uint16 + adaptor *espat.Device + 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 @@ -71,6 +77,12 @@ func (c *mqttclient) Connect() Token { 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 @@ -84,7 +96,7 @@ func (c *mqttclient) Connect() Token { connectPkt.PasswordFlag = true } - connectPkt.ClientIdentifier = c.opts.ClientID //"tinygo-client-" + randomString(10) + connectPkt.ClientIdentifier = c.opts.ClientID connectPkt.ProtocolVersion = byte(c.opts.ProtocolVersion) connectPkt.ProtocolName = "MQTT" connectPkt.Keepalive = 30 @@ -94,26 +106,25 @@ func (c *mqttclient) Connect() Token { return &mqtttoken{err: err} } - // TODO: handle timeout - for { - packet, _ := packets.ReadPacket(c.conn) - - if packet != nil { - ack, ok := packet.(*packets.ConnackPacket) - if ok { - if ack.ReturnCode == 0 { - // success - return &mqtttoken{} - } - // otherwise something went wrong + // 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 } - - time.Sleep(100 * time.Millisecond) } - c.connected = true + go readMessages(c) + go processInbound(c) + return &mqtttoken{} } @@ -129,6 +140,10 @@ func (c *mqttclient) Disconnect(quiesce uint) { // 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 @@ -144,12 +159,37 @@ func (c *mqttclient) Publish(topic string, qos byte, retained bool, payload inte c.mid++ err := pub.Write(c.conn) - return &mqtttoken{err: err} + 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{} } @@ -173,18 +213,92 @@ func (c *mqttclient) OptionsReader() ClientOptionsReader { return r } -type mqtttoken struct { - err error +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 + } + } } -func (t *mqtttoken) Wait() bool { - return true +// 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 (t *mqtttoken) WaitTimeout(time.Duration) bool { - return true +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 + } + } } -func (t *mqtttoken) Error() error { - return t.err +// 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 espat.ActiveDevice.IsSocketDataAvailable() { + return packets.ReadPacket(c.conn) + } + return nil, nil } diff --git a/espat/mqtt/paho.go b/espat/mqtt/paho.go index 7f0a99b..4a9eaf3 100644 --- a/espat/mqtt/paho.go +++ b/espat/mqtt/paho.go @@ -24,6 +24,7 @@ import ( "strings" "time" + "github.com/eclipse/paho.mqtt.golang/packets" "tinygo.org/x/drivers/espat" ) @@ -155,6 +156,18 @@ 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 diff --git a/espat/mqtt/router.go b/espat/mqtt/router.go new file mode 100644 index 0000000..6d6c59f --- /dev/null +++ b/espat/mqtt/router.go @@ -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 + } + } + }() +} diff --git a/espat/mqtt/token.go b/espat/mqtt/token.go new file mode 100644 index 0000000..f4225ff --- /dev/null +++ b/espat/mqtt/token.go @@ -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 +} diff --git a/espat/net/net.go b/espat/net/net.go index 58b12e7..4ed784c 100644 --- a/espat/net/net.go +++ b/espat/net/net.go @@ -23,7 +23,10 @@ func DialUDP(network string, laddr, raddr *UDPAddr) (*UDPSerialConn, error) { espat.ActiveDevice.DisconnectSocket() // connect new socket - espat.ActiveDevice.ConnectUDPSocket(addr, sendport, listenport) + err := espat.ActiveDevice.ConnectUDPSocket(addr, sendport, listenport) + if err != nil { + return nil, err + } return &UDPSerialConn{SerialConn: SerialConn{Adaptor: espat.ActiveDevice}, laddr: laddr, raddr: raddr}, nil } @@ -38,7 +41,10 @@ func ListenUDP(network string, laddr *UDPAddr) (*UDPSerialConn, error) { espat.ActiveDevice.DisconnectSocket() // connect new socket - espat.ActiveDevice.ConnectUDPSocket(addr, sendport, listenport) + err := espat.ActiveDevice.ConnectUDPSocket(addr, sendport, listenport) + if err != nil { + return nil, err + } return &UDPSerialConn{SerialConn: SerialConn{Adaptor: espat.ActiveDevice}, laddr: laddr}, nil } @@ -50,11 +56,14 @@ func DialTCP(network string, laddr, raddr *TCPAddr) (*TCPSerialConn, error) { addr := raddr.IP.String() sendport := strconv.Itoa(raddr.Port) - // disconnect any old socket - espat.ActiveDevice.DisconnectSocket() + // disconnect any old socket? + //espat.ActiveDevice.DisconnectSocket() // connect new socket - espat.ActiveDevice.ConnectTCPSocket(addr, sendport) + err := espat.ActiveDevice.ConnectTCPSocket(addr, sendport) + if err != nil { + return nil, err + } return &TCPSerialConn{SerialConn: SerialConn{Adaptor: espat.ActiveDevice}, laddr: laddr, raddr: raddr}, nil } @@ -132,8 +141,19 @@ func (c *SerialConn) Read(b []byte) (n int, err error) { 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. - c.Adaptor.StartSocketSend(len(b)) - return c.Adaptor.Write(b) + err = c.Adaptor.StartSocketSend(len(b)) + if err != nil { + return + } + n, err = c.Adaptor.Write(b) + if err != nil { + return n, err + } + _, err = c.Adaptor.Response(1000) + if err != nil { + return n, err + } + return n, err } // Close closes the connection. diff --git a/espat/tcp.go b/espat/tcp.go index 463c895..94506ea 100644 --- a/espat/tcp.go +++ b/espat/tcp.go @@ -17,7 +17,14 @@ const ( // GetDNS returns the IP address for a domain name. func (d *Device) GetDNS(domain string) (string, error) { d.Set(TCPDNSLookup, "\""+domain+"\"") - r := strings.Split(string(d.Response(1000)), ":") + resp, err := d.Response(1000) + if err != nil { + return "", err + } + if !strings.Contains(string(resp), ":") { + return "", errors.New("GetDNS error:" + string(resp)) + } + r := strings.Split(string(resp), ":") if len(r) != 2 { return "", errors.New("Invalid domain lookup result") } @@ -30,24 +37,30 @@ func (d *Device) GetDNS(domain string) (string, error) { func (d *Device) ConnectTCPSocket(addr, port string) error { protocol := "TCP" val := "\"" + protocol + "\",\"" + addr + "\"," + port + ",120" - d.Set(TCPConnect, val) - r := d.Response(1000) - if strings.Contains(string(r), "OK") { - return nil + err := d.Set(TCPConnect, val) + if err != nil { + return err } - return errors.New("ConnectTCPSocket error:" + string(r)) + _, e := d.Response(3000) + if e != nil { + return e + } + return nil } // ConnectUDPSocket creates a new UDP connection for the ESP8266/ESP32. func (d *Device) ConnectUDPSocket(addr, sendport, listenport string) error { protocol := "UDP" val := "\"" + protocol + "\",\"" + addr + "\"," + sendport + "," + listenport + ",2" - d.Set(TCPConnect, val) - r := d.Response(pause) - if strings.Contains(string(r), "OK") { - return nil + err := d.Set(TCPConnect, val) + if err != nil { + return err } - return errors.New("ConnectUDPSocket error:" + string(r)) + _, e := d.Response(3000) + if e != nil { + return e + } + return nil } // ConnectSSLSocket creates a new SSL socket connection for the ESP8266/ESP32. @@ -57,17 +70,23 @@ func (d *Device) ConnectSSLSocket(addr, port string) error { val := "\"" + protocol + "\",\"" + addr + "\"," + port + ",120" d.Set(TCPConnect, val) // this operation takes longer, so wait up to 6 seconds to complete. - r := d.Response(6000) - if strings.Contains(string(r), "CONNECT") { - return nil + _, err := d.Response(6000) + if err != nil { + return err } - return errors.New("ConnectSSLSocket error:" + string(r)) + return nil } // DisconnectSocket disconnects the ESP8266/ESP32 from the current TCP/UDP connection. func (d *Device) DisconnectSocket() error { - d.Execute(TCPClose) - d.Response(pause) + err := d.Execute(TCPClose) + if err != nil { + return err + } + _, e := d.Response(pause) + if e != nil { + return e + } return nil } @@ -76,14 +95,14 @@ func (d *Device) DisconnectSocket() error { func (d *Device) SetMux(mode int) error { val := strconv.Itoa(mode) d.Set(TCPMultiple, val) - d.Response(pause) - return nil + _, err := d.Response(pause) + return err } // GetMux returns the ESP8266/ESP32 current client TCP/UDP configuration for concurrent connections. func (d *Device) GetMux() ([]byte, error) { d.Query(TCPMultiple) - return d.Response(pause), nil + return d.Response(pause) } // SetTCPTransferMode sets the ESP8266/ESP32 current client TCP/UDP transfer mode. @@ -91,12 +110,12 @@ func (d *Device) GetMux() ([]byte, error) { func (d *Device) SetTCPTransferMode(mode int) error { val := strconv.Itoa(mode) d.Set(TransmissionMode, val) - d.Response(pause) - return nil + _, err := d.Response(pause) + return err } // GetTCPTransferMode returns the ESP8266/ESP32 current client TCP/UDP transfer mode. -func (d *Device) GetTCPTransferMode() []byte { +func (d *Device) GetTCPTransferMode() ([]byte, error) { d.Query(TransmissionMode) return d.Response(pause) } @@ -108,7 +127,10 @@ func (d *Device) StartSocketSend(size int) error { // when ">" is received, it indicates // ready to receive data - r := d.Response(pause) + r, err := d.Response(2000) + if err != nil { + return err + } if strings.Contains(string(r), ">") { return nil } @@ -120,6 +142,6 @@ func (d *Device) StartSocketSend(size int) error { func (d *Device) EndSocketSend() error { d.Write([]byte("+++")) - d.Response(pause) - return nil + _, err := d.Response(pause) + return err } diff --git a/espat/wifi.go b/espat/wifi.go index 48c734b..ea37a8e 100644 --- a/espat/wifi.go +++ b/espat/wifi.go @@ -16,7 +16,7 @@ const ( ) // GetWifiMode returns the ESP8266/ESP32 wifi mode. -func (d *Device) GetWifiMode() []byte { +func (d *Device) GetWifiMode() ([]byte, error) { d.Query(WifiMode) return d.Response(100) } @@ -25,14 +25,14 @@ func (d *Device) GetWifiMode() []byte { func (d *Device) SetWifiMode(mode int) error { val := strconv.Itoa(mode) d.Set(WifiMode, val) - d.Response(pause) - return nil + _, err := d.Response(pause) + return err } // Wifi Client // GetConnectedAP returns the ESP8266/ESP32 is currently connected to as a client. -func (d *Device) GetConnectedAP() []byte { +func (d *Device) GetConnectedAP() ([]byte, error) { d.Query(ConnectAP) return d.Response(100) } @@ -42,37 +42,43 @@ func (d *Device) GetConnectedAP() []byte { func (d *Device) ConnectToAP(ssid, pwd string, ws int) error { val := "\"" + ssid + "\",\"" + pwd + "\"" d.Set(ConnectAP, val) - d.Response(ws * 1000) + + _, err := d.Response(ws * 1000) + if err != nil { + return err + } return nil } // DisconnectFromAP disconnects the ESP8266/ESP32 from the current access point. func (d *Device) DisconnectFromAP() error { d.Execute(Disconnect) - d.Response(1000) - return nil + _, err := d.Response(1000) + return err } // GetClientIP returns the ESP8266/ESP32 current client IP addess when connected to an Access Point. -func (d *Device) GetClientIP() string { +func (d *Device) GetClientIP() (string, error) { d.Query(SetStationIP) - return string(d.Response(100)) + r, err := d.Response(1000) + return string(r), err } // SetClientIP sets the ESP8266/ESP32 current client IP addess when connected to an Access Point. -func (d *Device) SetClientIP(ipaddr string) []byte { +func (d *Device) SetClientIP(ipaddr string) error { val := "\"" + ipaddr + "\"" d.Set(ConnectAP, val) - d.Response(500) - return nil + _, err := d.Response(500) + return err } // Access Point // GetAPConfig returns the ESP8266/ESP32 current configuration when acting as an Access Point. -func (d *Device) GetAPConfig() string { +func (d *Device) GetAPConfig() (string, error) { d.Query(SoftAPConfigCurrent) - return string(d.Response(100)) + r, err := d.Response(100) + return string(r), err } // SetAPConfig sets the ESP8266/ESP32 current configuration when acting as an Access Point. @@ -83,35 +89,38 @@ func (d *Device) SetAPConfig(ssid, pwd string, ch, security int) error { ecnval := strconv.Itoa(security) val := "\"" + ssid + "\",\"" + pwd + "\"," + chval + "," + ecnval d.Set(SoftAPConfigCurrent, val) - d.Response(1000) - return nil + _, err := d.Response(1000) + return err } // GetAPClients returns the ESP8266/ESP32 current clients when acting as an Access Point. -func (d *Device) GetAPClients() string { +func (d *Device) GetAPClients() (string, error) { d.Query(ListConnectedIP) - return string(d.Response(100)) + r, err := d.Response(100) + return string(r), err } // GetAPIP returns the ESP8266/ESP32 current IP addess when configured as an Access Point. -func (d *Device) GetAPIP() string { +func (d *Device) GetAPIP() (string, error) { d.Query(SetSoftAPIPCurrent) - return string(d.Response(100)) + r, err := d.Response(100) + return string(r), err } // SetAPIP sets the ESP8266/ESP32 current IP addess when configured as an Access Point. func (d *Device) SetAPIP(ipaddr string) error { val := "\"" + ipaddr + "\"" d.Set(SetSoftAPIPCurrent, val) - d.Response(500) - return nil + _, err := d.Response(500) + return err } // GetAPConfigFlash returns the ESP8266/ESP32 current configuration acting as an Access Point // from flash storage. These settings are those used after a reset. -func (d *Device) GetAPConfigFlash() string { +func (d *Device) GetAPConfigFlash() (string, error) { d.Query(SoftAPConfigFlash) - return string(d.Response(100)) + r, err := d.Response(100) + return string(r), err } // SetAPConfigFlash sets the ESP8266/ESP32 current configuration acting as an Access Point, @@ -123,15 +132,16 @@ func (d *Device) SetAPConfigFlash(ssid, pwd string, ch, security int) error { ecnval := strconv.Itoa(security) val := "\"" + ssid + "\",\"" + pwd + "\"," + chval + "," + ecnval d.Set(SoftAPConfigFlash, val) - d.Response(1000) - return nil + _, err := d.Response(1000) + return err } // GetAPIPFlash returns the ESP8266/ESP32 IP address as saved to flash storage. // This is the IP address that will be used after a reset. -func (d *Device) GetAPIPFlash() string { +func (d *Device) GetAPIPFlash() (string, error) { d.Query(SetSoftAPIPFlash) - return string(d.Response(100)) + r, err := d.Response(100) + return string(r), err } // SetAPIPFlash sets the ESP8266/ESP32 current IP addess when configured as an Access Point. @@ -139,6 +149,6 @@ func (d *Device) GetAPIPFlash() string { func (d *Device) SetAPIPFlash(ipaddr string) error { val := "\"" + ipaddr + "\"" d.Set(SetSoftAPIPFlash, val) - d.Response(500) - return nil + _, err := d.Response(500) + return err } diff --git a/examples/espat/espconsole/main.go b/examples/espat/espconsole/main.go index ec20aea..c97b631 100644 --- a/examples/espat/espconsole/main.go +++ b/examples/espat/espconsole/main.go @@ -26,7 +26,7 @@ const pass = "YOURPASS" // these are the default pins for the Arduino Nano33 IoT. // change these to connect to a different UART or pins for the ESP8266/ESP32 var ( - uart = machine.UART2 + uart = machine.UART1 tx = machine.PA22 rx = machine.PA23 @@ -43,28 +43,20 @@ func main() { adaptor.Configure() // first check if connected - if adaptor.Connected() { + if connectToESP() { + println("Connected to wifi adaptor.") adaptor.Echo(false) - console.Write([]byte("\r\n")) - console.Write([]byte("ESP-AT console enabled.\r\n")) - console.Write([]byte("Firmware version:\r\n")) - console.Write(adaptor.Version()) - console.Write([]byte("\r\n")) - if actAsAP { - provideAP() - } else { - connectToAP() - } - - console.Write([]byte("Type an AT command then press enter:\r\n")) - prompt() + connectToAP() } else { - console.Write([]byte("\r\n")) - console.Write([]byte("Unable to connect to wifi adaptor.\r\n")) + println("") + failMessage("Unable to connect to wifi adaptor.") return } + println("Type an AT command then press enter:") + prompt() + input := make([]byte, 64) i := 0 for { @@ -82,7 +74,8 @@ func main() { adaptor.Write(input[:i+2]) // display response - console.Write(adaptor.Response(100)) + r, _ := adaptor.Response(500) + console.Write(r) // prompt prompt() @@ -101,28 +94,50 @@ func main() { } func prompt() { - console.Write([]byte("ESPAT>")) + print("ESPAT>") +} + +// connect to ESP8266/ESP32 +func connectToESP() bool { + for i := 0; i < 5; i++ { + println("Connecting to wifi adaptor...") + if adaptor.Connected() { + return true + } + time.Sleep(1 * time.Second) + } + return false } // connect to access point func connectToAP() { - console.Write([]byte("Connecting to wifi network...\r\n")) + println("Connecting to wifi network '" + ssid + "'") + adaptor.SetWifiMode(espat.WifiModeClient) adaptor.ConnectToAP(ssid, pass, 10) - console.Write([]byte("Connected.\r\n")) - console.Write([]byte(adaptor.GetClientIP())) - console.Write([]byte("\r\n")) + + println("Connected.") + ip, err := adaptor.GetClientIP() + if err != nil { + failMessage(err.Error()) + } + + println(ip) } // provide access point func provideAP() { - time.Sleep(500 * time.Millisecond) - console.Write([]byte("Starting wifi network as access point '")) - console.Write([]byte(ssid)) - console.Write([]byte("'...\r\n")) + println("Starting wifi network as access point '" + ssid + "'...") adaptor.SetWifiMode(espat.WifiModeAP) adaptor.SetAPConfig(ssid, pass, 7, espat.WifiAPSecurityWPA2_PSK) - console.Write([]byte("Ready.\r\n")) - console.Write([]byte(adaptor.GetAPIP())) - console.Write([]byte("\r\n")) + println("Ready.") + ip, _ := adaptor.GetAPIP() + println(ip) +} + +func failMessage(msg string) { + for { + println(msg) + time.Sleep(1 * time.Second) + } } diff --git a/examples/espat/esphub/main.go b/examples/espat/esphub/main.go index 7d2f450..bba6427 100644 --- a/examples/espat/esphub/main.go +++ b/examples/espat/esphub/main.go @@ -24,7 +24,7 @@ const pass = "YOURPASS" // these are the default pins for the Arduino Nano33 IoT. // change these to connect to a different UART or pins for the ESP8266/ESP32 var ( - uart = machine.UART2 + uart = machine.UART1 tx = machine.PA22 rx = machine.PA23 @@ -43,17 +43,14 @@ func main() { readyled.High() // first check if connected - if adaptor.Connected() { + if connectToESP() { println("Connected to wifi adaptor.") adaptor.Echo(false) - if actAsAP { - provideAP() - } else { - connectToAP() - } + connectToAP() } else { - println("Unable to connect to wifi adaptor.") + println("") + failMessage("Unable to connect to wifi adaptor.") return } @@ -86,21 +83,47 @@ func main() { println("Done.") } +// connect to ESP8266/ESP32 +func connectToESP() bool { + for i := 0; i < 5; i++ { + println("Connecting to wifi adaptor...") + if adaptor.Connected() { + return true + } + time.Sleep(1 * time.Second) + } + return false +} + // connect to access point func connectToAP() { - println("Connecting to wifi network...") + println("Connecting to wifi network '" + ssid + "'") + adaptor.SetWifiMode(espat.WifiModeClient) adaptor.ConnectToAP(ssid, pass, 10) + println("Connected.") - println(adaptor.GetClientIP()) + ip, err := adaptor.GetClientIP() + if err != nil { + failMessage(err.Error()) + } + + println(ip) } // provide access point func provideAP() { - println("Starting wifi network as access point:") - println(ssid) + println("Starting wifi network as access point '" + ssid + "'...") adaptor.SetWifiMode(espat.WifiModeAP) adaptor.SetAPConfig(ssid, pass, 7, espat.WifiAPSecurityWPA2_PSK) println("Ready.") - println(adaptor.GetAPIP()) + ip, _ := adaptor.GetAPIP() + println(ip) +} + +func failMessage(msg string) { + for { + println(msg) + time.Sleep(1 * time.Second) + } } diff --git a/examples/espat/espstation/main.go b/examples/espat/espstation/main.go index 38c4007..2cdb380 100644 --- a/examples/espat/espstation/main.go +++ b/examples/espat/espstation/main.go @@ -24,7 +24,7 @@ const hubIP = "0.0.0.0" // these are the default pins for the Arduino Nano33 IoT. // change these to connect to a different UART or pins for the ESP8266/ESP32 var ( - uart = machine.UART2 + uart = machine.UART1 tx = machine.PA22 rx = machine.PA23 @@ -39,13 +39,14 @@ func main() { adaptor.Configure() // first check if connected - if adaptor.Connected() { + if connectToESP() { println("Connected to wifi adaptor.") adaptor.Echo(false) connectToAP() } else { - println("Unable to connect to wifi adaptor.") + println("") + failMessage("Unable to connect to wifi adaptor.") return } @@ -71,11 +72,37 @@ func main() { println("Done.") } +// connect to ESP8266/ESP32 +func connectToESP() bool { + for i := 0; i < 5; i++ { + println("Connecting to wifi adaptor...") + if adaptor.Connected() { + return true + } + time.Sleep(1 * time.Second) + } + return false +} + // connect to access point func connectToAP() { - println("Connecting to wifi network...") + println("Connecting to wifi network '" + ssid + "'") + adaptor.SetWifiMode(espat.WifiModeClient) adaptor.ConnectToAP(ssid, pass, 10) + println("Connected.") - println(adaptor.GetClientIP()) + ip, err := adaptor.GetClientIP() + if err != nil { + failMessage(err.Error()) + } + + println(ip) +} + +func failMessage(msg string) { + for { + println(msg) + time.Sleep(1 * time.Second) + } } diff --git a/examples/espat/mqttclient/main.go b/examples/espat/mqttclient/main.go index b43ca0b..acbe829 100644 --- a/examples/espat/mqttclient/main.go +++ b/examples/espat/mqttclient/main.go @@ -25,8 +25,9 @@ const ssid = "YOURSSID" const pass = "YOURPASS" // IP address of the MQTT broker to use. Replace with your own info. -//const server = "tcp://test.mosquitto.org:1883" -const server = "ssl://test.mosquitto.org:8883" +const server = "tcp://test.mosquitto.org:1883" + +//const server = "ssl://test.mosquitto.org:8883" // these are the default pins for the Arduino Nano33 IoT. // change these to connect to a different UART or pins for the ESP8266/ESP32 @@ -66,7 +67,7 @@ func main() { opts := mqtt.NewClientOptions(adaptor) opts.AddBroker(server).SetClientID("tinygo-client-" + randomString(10)) - println("Connectng to MQTT...") + println("Connecting to MQTT broker at", server) cl := mqtt.NewClient(opts) if token := cl.Connect(); token.Wait() && token.Error() != nil { failMessage(token.Error().Error()) @@ -105,13 +106,18 @@ func connectToESP() bool { // connect to access point func connectToAP() { - println("Connecting to wifi network...") + println("Connecting to wifi network '" + ssid + "'") adaptor.SetWifiMode(espat.WifiModeClient) adaptor.ConnectToAP(ssid, pass, 10) println("Connected.") - println(adaptor.GetClientIP()) + ip, err := adaptor.GetClientIP() + if err != nil { + failMessage(err.Error()) + } + + println(ip) } // Returns an int >= min, < max diff --git a/examples/espat/mqttsub/main.go b/examples/espat/mqttsub/main.go new file mode 100644 index 0000000..1a429f5 --- /dev/null +++ b/examples/espat/mqttsub/main.go @@ -0,0 +1,163 @@ +// This is a sensor station that uses a ESP8266 or ESP32 running on the device UART1. +// It creates an MQTT connection that publishes a message every second +// to an MQTT broker. +// +// In other words: +// Your computer <--> UART0 <--> MCU <--> UART1 <--> ESP8266 <--> Internet <--> MQTT broker. +// +// You must also install the Paho MQTT package to build this program: +// +// go get -u github.com/eclipse/paho.mqtt.golang +// +package main + +import ( + "fmt" + "machine" + "math/rand" + "time" + + "tinygo.org/x/drivers/espat" + "tinygo.org/x/drivers/espat/mqtt" +) + +// access point info +const ssid = "YOURSSID" +const pass = "YOURPASS" + +// IP address of the MQTT broker to use. Replace with your own info. +//const server = "tcp://test.mosquitto.org:1883" + +const server = "ssl://test.mosquitto.org:8883" + +// change these to connect to a different UART or pins for the ESP8266/ESP32 +var ( + // these are defaults for the Arduino Nano33 IoT. + uart = machine.UART1 + tx = machine.PA22 + rx = machine.PA23 + + console = machine.UART0 + + adaptor *espat.Device + cl mqtt.Client + topicTx = "tinygo/tx" + topicRx = "tinygo/rx" +) + +func subHandler(client mqtt.Client, msg mqtt.Message) { + fmt.Printf("[%s] ", msg.Topic()) + fmt.Printf("%s\r\n", msg.Payload()) +} + +func main() { + time.Sleep(3000 * time.Millisecond) + + uart.Configure(machine.UARTConfig{TX: tx, RX: rx}) + rand.Seed(time.Now().UnixNano()) + + // Init esp8266/esp32 + adaptor = espat.New(uart) + adaptor.Configure() + + // first check if connected + if connectToESP() { + println("Connected to wifi adaptor.") + adaptor.Echo(false) + + connectToAP() + } else { + println("") + failMessage("Unable to connect to wifi adaptor.") + return + } + + opts := mqtt.NewClientOptions(adaptor) + opts.AddBroker(server).SetClientID("tinygo-client-" + randomString(10)) + + println("Connecting to MQTT broker at", server) + cl = mqtt.NewClient(opts) + if token := cl.Connect(); token.Wait() && token.Error() != nil { + failMessage(token.Error().Error()) + } + + // subscribe + token := cl.Subscribe(topicRx, 0, subHandler) + token.Wait() + if token.Error() != nil { + failMessage(token.Error().Error()) + } + + go publishing() + + select {} + + // Right now this code is never reached. Need a way to trigger it... + println("Disconnecting MQTT...") + cl.Disconnect(100) + + println("Done.") +} + +func publishing() { + for { + println("Publishing MQTT message...") + data := []byte("{\"e\":[{ \"n\":\"hello\", \"v\":101 }]}") + token := cl.Publish(topicTx, 0, false, data) + token.Wait() + if token.Error() != nil { + println(token.Error().Error()) + } + + time.Sleep(1000 * time.Millisecond) + } +} + +// connect to ESP8266/ESP32 +func connectToESP() bool { + for i := 0; i < 5; i++ { + println("Connecting to wifi adaptor...") + if adaptor.Connected() { + return true + } + time.Sleep(1 * time.Second) + } + return false +} + +// connect to access point +func connectToAP() { + println("Connecting to wifi network '" + ssid + "'") + + adaptor.SetWifiMode(espat.WifiModeClient) + adaptor.ConnectToAP(ssid, pass, 10) + + println("Connected.") + ip, err := adaptor.GetClientIP() + if err != nil { + failMessage(err.Error()) + } + + println(ip) +} + +// Returns an int >= min, < max +func randomInt(min, max int) int { + return min + rand.Intn(max-min) +} + +// Generate a random string of A-Z chars with len = l +func randomString(len int) string { + bytes := make([]byte, len) + for i := 0; i < len; i++ { + bytes[i] = byte(randomInt(65, 90)) + } + return string(bytes) +} + +func failMessage(msg string) { + for { + println(msg) + time.Sleep(1 * time.Second) + } +} diff --git a/examples/espat/tcpclient/main.go b/examples/espat/tcpclient/main.go index 831e440..f85a62d 100644 --- a/examples/espat/tcpclient/main.go +++ b/examples/espat/tcpclient/main.go @@ -24,7 +24,7 @@ const serverIP = "0.0.0.0" // these are the default pins for the Arduino Nano33 IoT. // change these to connect to a different UART or pins for the ESP8266/ESP32 var ( - uart = machine.UART2 + uart = machine.UART1 tx = machine.PA22 rx = machine.PA23 @@ -39,13 +39,14 @@ func main() { adaptor.Configure() // first check if connected - if adaptor.Connected() { + if connectToESP() { println("Connected to wifi adaptor.") adaptor.Echo(false) connectToAP() } else { - println("Unable to connect to wifi adaptor.") + println("") + failMessage("Unable to connect to wifi adaptor.") return } @@ -55,7 +56,10 @@ func main() { laddr := &net.TCPAddr{Port: 8080} println("Dialing TCP connection...") - conn, _ := net.DialTCP("tcp", laddr, raddr) + conn, err := net.DialTCP("tcp", laddr, raddr) + if err != nil { + failMessage(err.Error()) + } for { // send data @@ -71,11 +75,37 @@ func main() { println("Done.") } +// connect to ESP8266/ESP32 +func connectToESP() bool { + for i := 0; i < 5; i++ { + println("Connecting to wifi adaptor...") + if adaptor.Connected() { + return true + } + time.Sleep(1 * time.Second) + } + return false +} + // connect to access point func connectToAP() { - println("Connecting to wifi network...") + println("Connecting to wifi network '" + ssid + "'") + adaptor.SetWifiMode(espat.WifiModeClient) adaptor.ConnectToAP(ssid, pass, 10) + println("Connected.") - println(adaptor.GetClientIP()) + ip, err := adaptor.GetClientIP() + if err != nil { + failMessage(err.Error()) + } + + println(ip) +} + +func failMessage(msg string) { + for { + println(msg) + time.Sleep(1 * time.Second) + } }