Adding initial working project, ported from pico-thermo
This commit is contained in:
+3
-1
@@ -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
|
||||
|
||||
|
||||
@@ -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"
|
||||
)
|
||||
@@ -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
|
||||
)
|
||||
@@ -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=
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
+249
@@ -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)
|
||||
}
|
||||
Reference in New Issue
Block a user