From 12959be182a8b92200fbc1a264061fdee42ba961 Mon Sep 17 00:00:00 2001 From: Gregory Ballantine Date: Sun, 6 Sep 2026 23:26:11 -0400 Subject: [PATCH] Adding initial working project, ported from pico-thermo --- .gitignore | 4 +- config.go.example | 15 ++ go.mod | 14 ++ go.sum | 10 ++ main.go | 351 ++++++++++++++++++++++++++++++++++++++++++++++ network.go | 249 ++++++++++++++++++++++++++++++++ 6 files changed, 642 insertions(+), 1 deletion(-) create mode 100644 config.go.example create mode 100644 go.mod create mode 100644 go.sum create mode 100644 main.go create mode 100644 network.go diff --git a/.gitignore b/.gitignore index 5b90e79..b182cf6 100644 --- a/.gitignore +++ b/.gitignore @@ -1,3 +1,6 @@ +# Config file that may have secrets +config.go + # ---> Go # If you prefer the allow list template instead of the deny list, see community template: # https://github.com/github/gitignore/blob/main/community/Golang/Go.AllowList.gitignore @@ -24,4 +27,3 @@ go.work.sum # env file .env - diff --git a/config.go.example b/config.go.example new file mode 100644 index 0000000..54514f8 --- /dev/null +++ b/config.go.example @@ -0,0 +1,15 @@ +package main + +const ( + WifiSSID = "My Wifi Name" + WifiPass = "SecretWifiPassword" + + MQTTHost = "192.168.1.50" // Must be an IP address; DNS resolution is not supported currently + MQTTPort = "1883" + MQTTUser = "" // Leave empty if not required + MQTTPass = "" // Leave empty if not required + + NodeName = "Studio Weather" + NodeID = "studio_env" + ClientID = "pico2w_studio_env" +) diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..2cce03b --- /dev/null +++ b/go.mod @@ -0,0 +1,14 @@ +module git.metaunix.net/gballan/pico-weather + +go 1.25.2 + +require ( + github.com/soypat/cyw43439 v0.1.1 + github.com/soypat/lneto v0.3.2 +) + +require ( + github.com/soypat/seqs v0.0.0-20250124201400-0d65bc7c1710 // indirect + github.com/tinygo-org/pio v0.3.0 // indirect + golang.org/x/exp v0.0.0-20241204233417-43b7b7cde48d // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..3fc1ab5 --- /dev/null +++ b/go.sum @@ -0,0 +1,10 @@ +github.com/soypat/cyw43439 v0.1.1 h1:vcaTiVzfuz3keK7lJpVxStZ6tV8HCw7Ugzsh1k4mneE= +github.com/soypat/cyw43439 v0.1.1/go.mod h1:R2uSILRwSPmcmmKy5Z0FtK4ypgiPf5YqK+F+IKmXqxc= +github.com/soypat/lneto v0.3.2 h1:iUFeRSq2czT7Db6MMOsAnMCBlKCqvIr941zsNf9dcu0= +github.com/soypat/lneto v0.3.2/go.mod h1:Be5PjwoYukvHFiUXxpYi8+ppH2F/gw/vjGBvFdv+Ti8= +github.com/soypat/seqs v0.0.0-20250124201400-0d65bc7c1710 h1:Y9fBuiR/urFY/m76+SAZTxk2xAOS2n85f+H1CugajeA= +github.com/soypat/seqs v0.0.0-20250124201400-0d65bc7c1710/go.mod h1:oCVCNGCHMKoBj97Zp9znLbQ1nHxpkmOY9X+UAGzOxc8= +github.com/tinygo-org/pio v0.3.0 h1:opEnOtw58KGB4RJD3/n/Rd0/djYGX3DeJiXLI6y/yDI= +github.com/tinygo-org/pio v0.3.0/go.mod h1:wf6c6lKZp+pQOzKKcpzchmRuhiMc27ABRuo7KVnaMFU= +golang.org/x/exp v0.0.0-20241204233417-43b7b7cde48d h1:0olWaB5pg3+oychR51GUVCEsGkeCU/2JxjBgIo4f3M0= +golang.org/x/exp v0.0.0-20241204233417-43b7b7cde48d/go.mod h1:qj5a5QZpwLU2NLQudwIN5koi3beDhSAlJwa67PuM98c= diff --git a/main.go b/main.go new file mode 100644 index 0000000..321081c --- /dev/null +++ b/main.go @@ -0,0 +1,351 @@ +package main + +import ( + "bytes" + "context" + "encoding/binary" + "fmt" + "io" + "machine" + "strconv" + "time" +) + +var StateTopic = fmt.Sprintf("homeassistant/sensor/%s/state", NodeID) + +// ------------------------------------------------------------- +// Minimal MQTT 3.1.1 Client +// ------------------------------------------------------------- +type SimpleMQTT struct { + rw io.ReadWriter +} + +func NewSimpleMQTT(rw io.ReadWriter) *SimpleMQTT { + return &SimpleMQTT{rw: rw} +} + +func (m *SimpleMQTT) Connect(clientID, username, password string) error { + flags := byte(0x02) // Clean session + if username != "" { + flags |= 0x80 + } + if password != "" { + flags |= 0x40 + } + + // Variable header: MQTT (proto name), Level 4 (MQTT 3.1.1), flags, KeepAlive 60s + varHeader := []byte{0x00, 0x04, 'M', 'Q', 'T', 'T', 0x04, flags, 0x00, 0x3C} + + payload := encodeString(clientID) + if username != "" { + payload = append(payload, encodeString(username)...) + } + if password != "" { + payload = append(payload, encodeString(password)...) + } + + body := append(varHeader, payload...) + packet := append([]byte{0x10}, encodeLength(len(body))...) + packet = append(packet, body...) + + if _, err := m.rw.Write(packet); err != nil { + return err + } + + resp := make([]byte, 4) + if _, err := io.ReadFull(m.rw, resp); err != nil { + return err + } + if resp[0] != 0x20 || resp[3] != 0x00 { + return fmt.Errorf("MQTT connection refused: code %d", resp[3]) + } + return nil +} + +func (m *SimpleMQTT) Publish(topic string, payload []byte, retain bool) error { + cmd := byte(0x30) + if retain { + cmd |= 0x01 + } + + body := append(encodeString(topic), payload...) + packet := append([]byte{cmd}, encodeLength(len(body))...) + packet = append(packet, body...) + + _, err := m.rw.Write(packet) + return err +} + +func encodeString(s string) []byte { + b := []byte(s) + length := uint16(len(b)) + return append([]byte{byte(length >> 8), byte(length & 0xFF)}, b...) +} + +func encodeLength(length int) []byte { + var encoded []byte + for { + digit := byte(length % 128) + length /= 128 + if length > 0 { + digit |= 0x80 + } + encoded = append(encoded, digit) + if length == 0 { + break + } + } + return encoded +} + +// ------------------------------------------------------------- +// BME280 Driver +// ------------------------------------------------------------- +const ( + BME280Addr = 0x77 + regCalib00 = 0x88 + regCalib26 = 0xE1 + regReset = 0xE0 + regCtrlHum = 0xF2 + regCtrlMeas = 0xF4 + regConfig = 0xF5 + regData = 0xF7 +) + +type BME280Calib struct { + digT1 uint16 + digT2 int16 + digT3 int16 + digP1 uint16 + digP2 int16 + digP3 int16 + digP4 int16 + digP5 int16 + digP6 int16 + digP7 int16 + digP8 int16 + digP9 int16 + digH1 uint8 + digH2 int16 + digH3 uint8 + digH4 int16 + digH5 int16 + digH6 int8 +} + +type BME280 struct { + bus *machine.I2C + addr uint8 + calib BME280Calib +} + +func NewBME280(bus *machine.I2C, addr uint8) (*BME280, error) { + b := &BME280{bus: bus, addr: addr} + if err := bus.WriteRegister(addr, regReset, []byte{0xB6}); err != nil { + return nil, err + } + time.Sleep(100 * time.Millisecond) + + if err := b.readCalibration(); err != nil { + return nil, err + } + if err := bus.WriteRegister(addr, regCtrlHum, []byte{0x01}); err != nil { + return nil, err + } + if err := bus.WriteRegister(addr, regCtrlMeas, []byte{0x27}); err != nil { + return nil, err + } + if err := bus.WriteRegister(addr, regConfig, []byte{0xA0}); err != nil { + return nil, err + } + return b, nil +} + +func (b *BME280) readCalibration() error { + var buf24 [24]byte + if err := b.bus.ReadRegister(b.addr, regCalib00, buf24[:]); err != nil { + return err + } + r := bytes.NewReader(buf24[:]) + binary.Read(r, binary.LittleEndian, &b.calib.digT1) + binary.Read(r, binary.LittleEndian, &b.calib.digT2) + binary.Read(r, binary.LittleEndian, &b.calib.digT3) + binary.Read(r, binary.LittleEndian, &b.calib.digP1) + binary.Read(r, binary.LittleEndian, &b.calib.digP2) + binary.Read(r, binary.LittleEndian, &b.calib.digP3) + binary.Read(r, binary.LittleEndian, &b.calib.digP4) + binary.Read(r, binary.LittleEndian, &b.calib.digP5) + binary.Read(r, binary.LittleEndian, &b.calib.digP6) + binary.Read(r, binary.LittleEndian, &b.calib.digP7) + binary.Read(r, binary.LittleEndian, &b.calib.digP8) + binary.Read(r, binary.LittleEndian, &b.calib.digP9) + + var h1 [1]byte + if err := b.bus.ReadRegister(b.addr, 0xA1, h1[:]); err != nil { + return err + } + b.calib.digH1 = h1[0] + + var hBuf [7]byte + if err := b.bus.ReadRegister(b.addr, regCalib26, hBuf[:]); err != nil { + return err + } + b.calib.digH2 = int16(binary.LittleEndian.Uint16(hBuf[0:2])) + b.calib.digH3 = hBuf[2] + b.calib.digH4 = (int16(hBuf[3]) << 4) | (int16(hBuf[4]) & 0x0F) + b.calib.digH5 = (int16(hBuf[5]) << 4) | (int16(hBuf[4]) >> 4) + b.calib.digH6 = int8(hBuf[6]) + return nil +} + +func (b *BME280) ReadValues() (float64, float64, float64, error) { + var raw [8]byte + if err := b.bus.ReadRegister(b.addr, regData, raw[:]); err != nil { + return 0, 0, 0, err + } + + rawP := (int32(raw[0]) << 12) | (int32(raw[1]) << 4) | (int32(raw[2]) >> 4) + rawT := (int32(raw[3]) << 12) | (int32(raw[4]) << 4) | (int32(raw[5]) >> 4) + rawH := (int32(raw[6]) << 8) | int32(raw[7]) + + var1 := (((rawT >> 3) - (int32(b.calib.digT1) << 1)) * int32(b.calib.digT2)) >> 11 + var2 := (((((rawT >> 4) - int32(b.calib.digT1)) * ((rawT >> 4) - int32(b.calib.digT1))) >> 12) * int32(b.calib.digT3)) >> 14 + tFine := var1 + var2 + temp := float64((tFine*5+128)>>8) / 100.0 + + pVar1 := int64(tFine) - 128000 + pVar2 := pVar1*pVar1*int64(b.calib.digP6) + ((pVar1 * int64(b.calib.digP5)) << 17) + (int64(b.calib.digP4) << 35) + pVar1 = ((pVar1 * pVar1 * int64(b.calib.digP3)) >> 8) + ((pVar1 * int64(b.calib.digP2)) << 12) + pVar1 = (((int64(1) << 47) + pVar1) * int64(b.calib.digP1)) >> 33 + + var pres float64 + if pVar1 != 0 { + p := int64(1048576 - rawP) + p = (((p << 31) - pVar2) * 3125) / pVar1 + pVar1 = (int64(b.calib.digP9) * (p >> 13) * (p >> 13)) >> 25 + pVar2 = (int64(b.calib.digP8) * p) >> 19 + p = ((p + pVar1 + pVar2) >> 8) + (int64(b.calib.digP7) << 4) + pres = (float64(p) / 256.0) / 100.0 + } + + hVar := tFine - 76800 + hVar = (((((rawH << 14) - (int32(b.calib.digH4) << 20) - (int32(b.calib.digH5) * hVar)) + 16384) >> 15) * + (((((((hVar * int32(b.calib.digH6)) >> 10) * (((hVar * int32(b.calib.digH3)) >> 11) + 32768)) >> 10) + 2097152)* + int32(b.calib.digH2) + 8192) >> 14)) + hVar = hVar - (((((hVar >> 15) * (hVar >> 15)) >> 7) * int32(b.calib.digH1)) >> 4) + if hVar < 0 { + hVar = 0 + } else if hVar > 419430400 { + hVar = 419430400 + } + hum := float64(hVar>>12) / 1024.0 + + return temp, pres, hum, nil +} + +// ------------------------------------------------------------- +// Home Assistant Auto-Discovery +// ------------------------------------------------------------- +type haSensorDef struct { + id string + name string + unit string + class string + valTpl string +} + +func registerHADiscovery(mqtt *SimpleMQTT) error { + deviceJSON := fmt.Sprintf(`"device":{"identifiers":["%s"],"name":"%s","model":"Pico 2 W","manufacturer":"Raspberry Pi"}`, + NodeID, NodeName) + + sensors := []haSensorDef{ + {"temperature", "Temperature", "°C", "temperature", "{{ value_json.temperature }}"}, + {"humidity", "Humidity", "%", "humidity", "{{ value_json.humidity }}"}, + {"pressure", "Pressure", "hPa", "atmospheric_pressure", "{{ value_json.pressure }}"}, + } + + for _, s := range sensors { + topic := fmt.Sprintf("homeassistant/sensor/%s/%s/config", NodeID, s.id) + payload := fmt.Sprintf(`{"name":"%s","has_entity_name":true,"unique_id":"%s_%s","device_class":"%s","state_class":"measurement","unit_of_measurement":"%s","state_topic":"%s","value_template":"%s",%s}`, + s.name, NodeID, s.id, s.class, s.unit, StateTopic, s.valTpl, deviceJSON) + + if err := mqtt.Publish(topic, []byte(payload), true); err != nil { + return err + } + println("Registered HA discovery:", s.id) + } + return nil +} + +// ------------------------------------------------------------- +// Panic Handler & Entry Point +// ------------------------------------------------------------- +func panicErr(msg string, err error) { + if err != nil { + println("FATAL:", msg, "-", err.Error()) + } + for { + time.Sleep(500 * time.Millisecond) + } +} + +func main() { + time.Sleep(2 * time.Second) + + // Configure I2C (Pin 16 SDA, Pin 17 SCL -> I2C0) + i2c := machine.I2C0 + err := i2c.Configure(machine.I2CConfig{ + Frequency: 100 * machine.KHz, + SDA: machine.GPIO16, + SCL: machine.GPIO17, + }) + if err != nil { + panicErr("I2C configuration", err) + } + + bme, err := NewBME280(i2c, BME280Addr) + if err != nil { + panicErr("BME280 init", err) + } + + netStack, err := InitNetwork(WifiSSID, WifiPass) + if err != nil { + panicErr("network init", err) + } + + portNum, _ := strconv.Atoi(MQTTPort) + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + + conn, err := netStack.DialTCP(ctx, MQTTHost, uint16(portNum)) + if err != nil { + panicErr("TCP dial", err) + } + + mqtt := NewSimpleMQTT(conn) + println("Connecting MQTT session...") + if err := mqtt.Connect(ClientID, MQTTUser, MQTTPass); err != nil { + panicErr("MQTT connect", err) + } + println("MQTT connected.") + + if err := registerHADiscovery(mqtt); err != nil { + println("Warning: HA discovery failed:", err.Error()) + } + + for { + t, p, h, err := bme.ReadValues() + if err != nil { + println("BME280 read error:", err.Error()) + } else { + stateJSON := fmt.Sprintf(`{"temperature":%.2f,"humidity":%.2f,"pressure":%.2f}`, t, h, p) + if err := mqtt.Publish(StateTopic, []byte(stateJSON), false); err != nil { + println("MQTT publish error:", err.Error()) + } else { + println(fmt.Sprintf("[%s] Published: %s", NodeName, stateJSON)) + } + } + + time.Sleep(30 * time.Second) + } +} diff --git a/network.go b/network.go new file mode 100644 index 0000000..979e43a --- /dev/null +++ b/network.go @@ -0,0 +1,249 @@ +package main + +import ( + "context" + "errors" + "fmt" + "net" + "net/netip" + "syscall" + "time" + + "github.com/soypat/cyw43439" + "github.com/soypat/lneto" + "github.com/soypat/lneto/ethernet" + "github.com/soypat/lneto/x/xnet" +) + +const ( + pollTime = 5 * time.Millisecond + protoTimeout = 5 * time.Second + protoRetries = 3 + + tcpBufsize = 2048 + tcpPacketQueueSize = 4 + tcpConnPoolSize = 5 + tcpEstablishedTimeout = 4 * time.Second + tcpCloseTimeout = protoTimeout +) + +var nanotime = func() int64 { + return time.Now().UnixNano() +} + +// CywAdapter bridges the cyw43439 driver to the lneto interface requirements. +type CywAdapter struct { + dev *cyw43439.Device + rxBuf [1514]byte + rxLen int + hasPkt bool +} + +func NewCywAdapter(dev *cyw43439.Device) *CywAdapter { + adapter := &CywAdapter{dev: dev} + + dev.RecvEthHandle(func(pkt []byte) error { + if !adapter.hasPkt && len(pkt) <= len(adapter.rxBuf) { + copy(adapter.rxBuf[:], pkt) + adapter.rxLen = len(pkt) + adapter.hasPkt = true + } + return nil + }) + + return adapter +} + +// In CywAdapter: +func (a *CywAdapter) PollHardware() error { + _, err := a.dev.TryPoll() + return err +} + +func (a *CywAdapter) SendEth(frame []byte) error { + return a.dev.SendEth(frame) +} + +func (a *CywAdapter) RecvEth(dst []byte) (int, error) { + if !a.hasPkt { + return 0, nil + } + n := copy(dst, a.rxBuf[:a.rxLen]) + a.hasPkt = false + a.rxLen = 0 + return n, nil +} + +func (a *CywAdapter) HardwareAddress6() ([6]byte, error) { + return a.dev.HardwareAddr6() +} + +func (a *CywAdapter) MaxFrameLength() (int, error) { + return 1514, nil +} + +type NetworkStack struct { + adapter *CywAdapter + stack *xnet.StackAsync + gostack xnet.StackGo + localIP netip.Addr +} + +func InitNetwork(ssid, pass string) (*NetworkStack, error) { + println("Initializing CYW43439 Wi-Fi hardware...") + dev := cyw43439.NewPicoWDevice() + cfg := cyw43439.DefaultWifiConfig() + if err := dev.Init(cfg); err != nil { + return nil, fmt.Errorf("device init failed: %w", err) + } + + println("Associating with SSID:", ssid) + if err := dev.JoinWPA2(ssid, pass); err != nil { + return nil, fmt.Errorf("wifi association failed: %w", err) + } + println("Wi-Fi associated.") + + adapter := NewCywAdapter(dev) + hwaddr, err := adapter.HardwareAddress6() + if err != nil { + return nil, fmt.Errorf("read MAC: %w", err) + } + framelen, err := adapter.MaxFrameLength() + if err != nil { + return nil, fmt.Errorf("max frame len: %w", err) + } + + stack := &xnet.StackAsync{} + err = stack.Reset(xnet.StackConfig{ + Hostname: "pico2w-bme280", + RandSeed: time.Now().UnixNano(), + MaxActiveTCPPorts: 2, + MTU: uint16(framelen - ethernet.MaxOverheadSize), + HardwareAddress: hwaddr, + }) + if err != nil { + return nil, fmt.Errorf("stack config reset: %w", err) + } + + // Start background frame pump + ctx := context.Background() + go stackLoop(ctx, stack, adapter) + + println("Acquiring IP via DHCP...") + rstack := stack.StackRetrying(stackBackoff) + results, err := rstack.DoDHCPv4([4]byte{}, protoTimeout, protoRetries) + if err != nil { + return nil, fmt.Errorf("DHCP failed: %w", err) + } + + err = stack.AssimilateDHCPResults(results) + if err != nil { + return nil, fmt.Errorf("assimilate DHCP failed: %w", err) + } + + println("Resolving router MAC...") + gateway, err := rstack.DoResolveHardwareAddress6(results.Router, protoTimeout, protoRetries) + if err != nil { + return nil, fmt.Errorf("resolving router MAC failed: %w", err) + } + stack.SetGatewayHardwareAddr(gateway) + + localIP := netip.AddrFrom4(results.AssignedAddr4) + println("DHCP lease assigned! IP:", localIP.String()) + + gostack := stack.StackBlocking(stackBackoff).StackGo(xnet.StackGoConfig{ + ListenerPoolConfig: xnet.TCPPoolConfig{ + PoolSize: tcpConnPoolSize, + QueueSize: tcpPacketQueueSize, + TxBufSize: tcpBufsize, + RxBufSize: tcpBufsize, + NanoTime: nanotime, + EstablishedTimeout: tcpEstablishedTimeout, + ClosingTimeout: tcpCloseTimeout, + NewBackoff: func() lneto.BackoffStrategy { return tcpBackoff }, + }, + }) + + return &NetworkStack{ + adapter: adapter, + stack: stack, + gostack: gostack, + localIP: localIP, + }, nil +} + +func (ns *NetworkStack) DialTCP(ctx context.Context, hostIP string, port uint16) (net.Conn, error) { + rIP, err := netip.ParseAddr(hostIP) + if err != nil { + return nil, fmt.Errorf("invalid host IPv4: %w", err) + } + + laddr := net.TCPAddrFromAddrPort(netip.AddrPortFrom(ns.localIP, 0)) + raddr := net.TCPAddrFromAddrPort(netip.AddrPortFrom(rIP, port)) + + const sockstream = 0x1 + c, err := ns.gostack.Socket(ctx, "tcp", syscall.AF_INET, sockstream, laddr, raddr) + if err != nil { + return nil, fmt.Errorf("socket dial: %w", err) + } + + conn, ok := c.(net.Conn) + if !ok { + return nil, errors.New("socket did not return a stream connection") + } + + return conn, nil +} + +func stackLoop(ctx context.Context, stack *xnet.StackAsync, adapter *CywAdapter) { + frameLength, _ := adapter.MaxFrameLength() + buf := make([]byte, frameLength) + + for ctx.Err() == nil { + // 1. Pump the CYW43439 hardware over SPI to trigger RecvEthHandle + _ = adapter.PollHardware() + + // 2. Ingress: read from adapter into lneto + nread, err := adapter.RecvEth(buf[:]) + if err != nil { + println("recv err:", err.Error()) + } else if nread > 0 { + err = stack.IngressEthernet(buf[:nread]) + if err != nil && err != lneto.ErrPacketDrop { + println("ingress err:", err.Error()) + } + } + + // 3. Egress: send out any frames generated by lneto + nwrite, err := stack.EgressEthernet(buf[:]) + if err != nil { + println("egress err:", err.Error()) + } else if nwrite > 0 { + if err := adapter.SendEth(buf[:nwrite]); err != nil { + println("send eth err:", err.Error()) + } + } + + if nwrite == 0 && nread == 0 { + time.Sleep(pollTime) + } + } +} + +func stackBackoff(consecutiveBackoffs uint) time.Duration { + if consecutiveBackoffs < 10 { + return time.Millisecond + } + return 10 * time.Millisecond +} + +func tcpBackoff(consecutiveBackoffs uint) time.Duration { + const ( + minWait = uint32(time.Microsecond) + maxWait = 5 * uint32(time.Millisecond) + maxShift = 22 + ) + shifted := minWait << min(consecutiveBackoffs, maxShift) + wait := min(shifted, maxWait) + return time.Duration(wait) +}