Files

966 lines
30 KiB
Go

package robot
import (
"context"
"encoding/json"
"fmt"
"strings"
"sync"
"testing"
"time"
"git.i3omb.com/gronod/ha-n95-local-control/internal/ctl"
"git.i3omb.com/gronod/ha-n95-local-control/internal/session"
)
const (
testBotJID = "e20123456789@155.ecorobot.net/atom"
testSerial = "e20123456789"
testCtlJID = "n95bridge@ecouser.net/homeassistant"
)
// --- fakes and recorders ---
type fakeClock struct {
mu sync.Mutex
now time.Time
timers []*fakeTimer
}
type fakeTimer struct {
at time.Time
ch chan time.Time
fired bool
}
func newFakeClock() *fakeClock {
return &fakeClock{now: time.Unix(1790194386, 0)}
}
func (c *fakeClock) Now() time.Time {
c.mu.Lock()
defer c.mu.Unlock()
return c.now
}
func (c *fakeClock) After(d time.Duration) <-chan time.Time {
c.mu.Lock()
defer c.mu.Unlock()
t := &fakeTimer{at: c.now.Add(d), ch: make(chan time.Time, 1)}
c.timers = append(c.timers, t)
return t.ch
}
func (c *fakeClock) advance(d time.Duration) {
c.mu.Lock()
defer c.mu.Unlock()
c.now = c.now.Add(d)
for _, t := range c.timers {
if !t.fired && !t.at.After(c.now) {
t.fired = true
t.ch <- c.now
}
}
}
type sendRec struct {
mu sync.Mutex
xml []string
err error
}
func (s *sendRec) send(b []byte) error {
s.mu.Lock()
defer s.mu.Unlock()
s.xml = append(s.xml, string(b))
return s.err
}
func (s *sendRec) all() []string {
s.mu.Lock()
defer s.mu.Unlock()
return append([]string(nil), s.xml...)
}
type pubRec struct {
mu sync.Mutex
snaps []Snapshot
ch chan struct{}
}
func newPubRec() *pubRec { return &pubRec{ch: make(chan struct{}, 128)} }
func (r *pubRec) pub(_ context.Context, s Snapshot) {
r.mu.Lock()
r.snaps = append(r.snaps, s)
r.mu.Unlock()
r.ch <- struct{}{}
}
func (r *pubRec) last() Snapshot {
r.mu.Lock()
defer r.mu.Unlock()
if len(r.snaps) == 0 {
return Snapshot{}
}
return r.snaps[len(r.snaps)-1]
}
func (r *pubRec) count() int {
r.mu.Lock()
defer r.mu.Unlock()
return len(r.snaps)
}
func (r *pubRec) wait(t *testing.T) {
t.Helper()
select {
case <-r.ch:
case <-time.After(3 * time.Second):
t.Fatal("timed out waiting for republish")
}
}
type traceRec struct {
mu sync.Mutex
traces []ctl.Trace
ch chan struct{}
}
func newTraceRec() *traceRec { return &traceRec{ch: make(chan struct{}, 128)} }
func (r *traceRec) sink(_ string, tr ctl.Trace) {
r.mu.Lock()
r.traces = append(r.traces, tr)
r.mu.Unlock()
r.ch <- struct{}{}
}
func (r *traceRec) wait(t *testing.T) {
t.Helper()
select {
case <-r.ch:
case <-time.After(3 * time.Second):
t.Fatal("timed out waiting for trace")
}
}
func (r *traceRec) all() []ctl.Trace {
r.mu.Lock()
defer r.mu.Unlock()
return append([]ctl.Trace(nil), r.traces...)
}
type diagRec struct {
mu sync.Mutex
list []session.Diagnostic
}
func (r *diagRec) sink(d session.Diagnostic) {
r.mu.Lock()
r.list = append(r.list, d)
r.mu.Unlock()
}
func (r *diagRec) all() []session.Diagnostic {
r.mu.Lock()
defer r.mu.Unlock()
return append([]session.Diagnostic(nil), r.list...)
}
type rig struct {
actor *Actor
send *sendRec
pub *pubRec
trace *traceRec
diag *diagRec
clock *fakeClock
}
func newRig(t *testing.T) *rig {
t.Helper()
send := &sendRec{}
pub := newPubRec()
tr := newTraceRec()
dg := &diagRec{}
clk := newFakeClock()
ctx, cancel := context.WithCancel(context.Background())
t.Cleanup(cancel)
a := NewActor(ctx, testBotJID, testSerial, testCtlJID, pub.pub, tr.sink, dg.sink, clk)
a.SessionReady(session.ReadyEvent{
Generation: 1,
JID: testBotJID,
Serial: testSerial,
Send: send.send,
})
flushActor(t, a)
return &rig{actor: a, send: send, pub: pub, trace: tr, diag: dg, clock: clk}
}
// flushActor blocks until every event queued so far has been processed.
func flushActor(t *testing.T, a *Actor) {
t.Helper()
ch := make(chan struct{})
select {
case a.mailbox <- evBarrier(ch):
case <-time.After(3 * time.Second):
t.Fatal("actor mailbox blocked")
}
select {
case <-ch:
case <-time.After(3 * time.Second):
t.Fatal("actor did not process barrier")
}
}
func (r *rig) stanza(t *testing.T, xml string) {
t.Helper()
r.actor.Stanza(session.StanzaEvent{
Generation: 1,
JID: testBotJID,
Serial: testSerial,
Stanza: []byte(xml),
})
}
// --- stanza builders ---
func iqResult(sid string) string {
return fmt.Sprintf(`<iq to="%s" type="result" id="%s"/>`, testCtlJID, sid)
}
func ctlResult(cid, ret, inner string) string {
return fmt.Sprintf(`<iq to="%s" type="set" id="9"><query xmlns="com:ctl"><ctl id="%s" ret="%s">%s</ctl></query></iq>`,
testCtlJID, cid, ret, inner)
}
func iqSet(frag string) string {
return `<iq to="` + testCtlJID + `" type="set" id="9"><query xmlns="com:ctl">` + frag + `</query></iq>`
}
func stanzaSID(t *testing.T, stanza string) string {
t.Helper()
i := strings.Index(stanza, `id="`)
if i < 0 {
t.Fatalf("no iq id in %s", stanza)
}
rest := stanza[i+4:]
j := strings.IndexByte(rest, '"')
return rest[:j]
}
func stanzaCID(t *testing.T, stanza string) string {
t.Helper()
in, err := ctl.Parse([]byte(stanza))
if err != nil {
t.Fatalf("parse outbound: %v", err)
}
cid := in.Attrs["id"]
if cid == "" {
t.Fatalf("no ctl id in %s", stanza)
}
return cid
}
func applyPush(t *testing.T, snap *Snapshot, frag string) {
t.Helper()
in, err := ctl.Parse([]byte(iqSet(frag)))
if err != nil {
t.Fatalf("parse push: %v", err)
}
if err := Apply(snap, in, ""); err != nil {
t.Fatalf("Apply(%s): %v", frag, err)
}
}
// --- named phase 03 tests ---
func TestCtlSidVersusCid(t *testing.T) {
r := newRig(t)
r.actor.AnnounceOK(session.ReadyEvent{Generation: 1, JID: testBotJID, Serial: testSerial, Send: r.send.send})
flushActor(t, r.actor)
if err := r.actor.Submit(context.Background(), Command{Name: "start"}); err != nil {
t.Fatalf("Submit start: %v", err)
}
sends := r.send.all()
clean := sends[len(sends)-1]
sid, cid := stanzaSID(t, clean), stanzaCID(t, clean)
if sid == cid {
t.Fatalf("iq id %q must differ from ctl id", sid)
}
// The stanza ack completes the ack phase only.
r.stanza(t, iqResult(sid))
flushActor(t, r.actor)
traces := r.trace.all()
if len(traces) != 1 || traces[0].Phase != "ack" || traces[0].SID != sid {
t.Fatalf("after ack, traces = %+v", traces)
}
if got := r.pub.last().Attributes.LastCommandError; got != nil {
t.Fatalf("ack must clear last_command_error, got %q", *got)
}
// A later iq-set whose ctl id matches completes the result phase once.
r.stanza(t, ctlResult(cid, "ok", ""))
r.stanza(t, ctlResult(cid, "ok", ""))
flushActor(t, r.actor)
traces = r.trace.all()
if len(traces) != 2 || traces[1].Phase != "result" || traces[1].Ret != "ok" {
t.Fatalf("after result, traces = %+v", traces)
}
// SetTime result without errno is success.
setTimeCID := stanzaCID(t, sends[0])
r.stanza(t, ctlResult(setTimeCID, "ok", ""))
flushActor(t, r.actor)
if got := r.pub.last().Attributes.LastCommandError; got != nil {
t.Fatalf("ret ok without errno must clear last_command_error, got %q", *got)
}
// ret=fail sets last_command_error and leaves last_error null.
if err := r.actor.Submit(context.Background(), Command{Name: "start"}); err != nil {
t.Fatalf("Submit start: %v", err)
}
sends = r.send.all()
failCID := stanzaCID(t, sends[len(sends)-1])
r.stanza(t, fmt.Sprintf(`<iq to="%s" type="set" id="9"><query xmlns="com:ctl"><ctl id="%s" ret="fail" errno="5"/></query></iq>`, testCtlJID, failCID))
flushActor(t, r.actor)
attrs := r.pub.last().Attributes
if got := attrs.LastCommandError; got == nil || *got != "fail:5" {
t.Fatalf("last_command_error = %v, want fail:5", got)
}
if attrs.LastError != nil {
t.Fatalf("last_error must stay null, got %q", *attrs.LastError)
}
}
func TestReadyFanOutOrder(t *testing.T) {
r := newRig(t)
if n := len(r.send.all()); n != 0 {
t.Fatalf("%d sends before AnnounceOK, want 0", n)
}
r.actor.AnnounceOK(session.ReadyEvent{Generation: 1, JID: testBotJID, Serial: testSerial, Send: r.send.send})
flushActor(t, r.actor)
wantTD := []string{"SetTime", "GetBatteryInfo", "GetCleanState", "GetChargeState", "GetCleanSpeed", "GetSched", "GetLifeSpan", "GetLifeSpan", "GetLifeSpan"}
wantType := []string{"", "", "", "", "", "", "SideBrush", "Brush", "DustCaseHeap"}
sends := r.send.all()
if len(sends) != len(wantTD) {
t.Fatalf("%d fan-out stanzas, want %d", len(sends), len(wantTD))
}
seen := map[string]bool{}
for i, s := range sends {
in, err := ctl.Parse([]byte(s))
if err != nil {
t.Fatalf("stanza %d parse: %v", i, err)
}
if in.TD != wantTD[i] {
t.Errorf("stanza %d td = %q, want %q", i, in.TD, wantTD[i])
}
if wantType[i] != "" && in.Attrs["type"] != wantType[i] {
t.Errorf("stanza %d type = %q, want %q", i, in.Attrs["type"], wantType[i])
}
cid := in.Attrs["id"]
if cid == "" {
t.Errorf("stanza %d has no ctl id", i)
}
if seen[cid] {
t.Errorf("cid %q reused", cid)
}
seen[cid] = true
}
if !strings.Contains(sends[0], "<time t=") {
t.Errorf("SetTime missing time element: %s", sends[0])
}
}
func TestStatePrecedenceTable(t *testing.T) {
e103 := "103"
rows := []struct {
name string
facts Facts
want string
}{
{"init", Facts{Fan: "standard"}, "idle"},
{"slot charging", Facts{Charge: "SlotCharging", Fan: "standard"}, "docked"},
{"going", Facts{Charge: "going", Fan: "standard"}, "returning"},
{"going beats clean auto", Facts{Charge: "going", CleanType: "auto", Fan: "standard"}, "returning"},
{"error latch", Facts{ErrorLatched: true, LastError: &e103, CleanType: "auto"}, "error"},
{"error beats docked stop", Facts{ErrorLatched: true, LastError: &e103, Charge: "SlotCharging", CleanType: "stop"}, "error"},
{"cleared error shows docked", Facts{Charge: "SlotCharging", CleanType: "stop"}, "docked"},
{"cleaning no charge", Facts{CleanType: "auto"}, "cleaning"},
{"idle charge with stop", Facts{Charge: "Idle", CleanType: "stop"}, "idle"},
{"idle charge with auto", Facts{Charge: "Idle", CleanType: "auto"}, "cleaning"},
{"going with stop", Facts{Charge: "going", CleanType: "stop"}, "returning"},
{"paused", Facts{Paused: true, CleanType: "auto"}, "paused"},
{"returning beats paused", Facts{Paused: true, Charge: "going"}, "returning"},
{"docked beats everything below", Facts{Charge: "SlotCharging", Paused: true, CleanType: "auto"}, "docked"},
}
for _, row := range rows {
if got := Derive(row.facts); got != row.want {
t.Errorf("%s: Derive = %q, want %q", row.name, got, row.want)
}
}
}
func TestErrno100ClearsLastErrorOnly(t *testing.T) {
snap := NewSnapshot()
applyPush(t, &snap, `<ctl td="ChargeState"><charge type="Idle"/></ctl>`)
applyPush(t, &snap, `<ctl td="CleanReport"><clean type="auto" speed="standard" st=" " rsn=" "/></ctl>`)
if got := snap.State.State; got != "cleaning" {
t.Fatalf("state = %q, want cleaning", got)
}
applyPush(t, &snap, `<ctl td="error" errno="103"/>`)
if got := snap.State.State; got != "error" {
t.Fatalf("state = %q, want error", got)
}
if le := snap.Attributes.LastError; le == nil || *le != "103" {
t.Fatalf("last_error = %v, want 103", le)
}
// Reports keep updating facts while the latch holds.
applyPush(t, &snap, `<ctl td="ChargeState"><charge type="SlotCharging"/></ctl>`)
applyPush(t, &snap, `<ctl td="CleanReport"><clean type="stop" speed="standard" st=" " rsn=" "/></ctl>`)
if got := snap.Facts.CleanType; got != "stop" {
t.Fatalf("clean_type = %q, want stop", got)
}
if got := snap.State.State; got != "error" {
t.Fatalf("state = %q, want error while latched", got)
}
applyPush(t, &snap, `<ctl td="error" errno="100"/>`)
if snap.Attributes.LastError != nil {
t.Fatalf("last_error = %v, want null after errno 100", *snap.Attributes.LastError)
}
if got := snap.State.State; got != "docked" {
t.Fatalf("state = %q, want docked from SlotCharging + stop", got)
}
if got := snap.Facts.Charge; got != "SlotCharging" {
t.Fatalf("charge = %q, errno 100 must not clear facts", got)
}
if got := snap.Facts.Fan; got != "standard" {
t.Fatalf("fan = %q, errno 100 must not clear fan", got)
}
}
func TestIdleReDerivesWithoutStartingClean(t *testing.T) {
snap := NewSnapshot()
applyPush(t, &snap, `<ctl td="ChargeState"><charge type="SlotCharging"/></ctl>`)
applyPush(t, &snap, `<ctl td="CleanReport"><clean type="stop" speed="standard" st=" " rsn=" "/></ctl>`)
applyPush(t, &snap, `<ctl td="ChargeState"><charge type="Idle"/></ctl>`)
if got := snap.State.State; got != "idle" {
t.Fatalf("Idle + stop: state = %q, want idle", got)
}
if got := snap.Facts.CleanType; got != "stop" {
t.Fatalf("Idle must not write clean_type, got %q", got)
}
snap2 := NewSnapshot()
applyPush(t, &snap2, `<ctl td="ChargeState"><charge type="SlotCharging"/></ctl>`)
applyPush(t, &snap2, `<ctl td="CleanReport"><clean type="auto" speed="standard" st=" " rsn=" "/></ctl>`)
applyPush(t, &snap2, `<ctl td="ChargeState"><charge type="Idle"/></ctl>`)
if got := snap2.State.State; got != "cleaning" {
t.Fatalf("Idle + remembered auto: state = %q, want cleaning", got)
}
// A direct td="Idle" push is the same report: it rewrites charge_state
// and re-derives without writing a clean type.
snap3 := NewSnapshot()
applyPush(t, &snap3, `<ctl td="ChargeState"><charge type="SlotCharging"/></ctl>`)
applyPush(t, &snap3, `<ctl td="CleanReport"><clean type="stop" speed="standard" st=" " rsn=" "/></ctl>`)
applyPush(t, &snap3, `<ctl td="Idle"/>`)
if got := snap3.State.State; got != "idle" {
t.Fatalf("td=Idle + stop: state = %q, want idle", got)
}
if got := snap3.Facts.CleanType; got != "stop" {
t.Fatalf("td=Idle must not write clean_type, got %q", got)
}
if got := snap3.Facts.Charge; got != "Idle" {
t.Fatalf("td=Idle must set charge_state, got %q", got)
}
snap4 := NewSnapshot()
applyPush(t, &snap4, `<ctl td="ChargeState"><charge type="SlotCharging"/></ctl>`)
applyPush(t, &snap4, `<ctl td="CleanReport"><clean type="auto" speed="standard" st=" " rsn=" "/></ctl>`)
applyPush(t, &snap4, `<ctl td="Idle"/>`)
if got := snap4.State.State; got != "cleaning" {
t.Fatalf("td=Idle + remembered auto: state = %q, want cleaning", got)
}
}
func TestStopDockedOnlyWhileSlotCharging(t *testing.T) {
for _, tc := range []struct {
charge string
want string
}{
{"SlotCharging", "docked"},
{"Idle", "idle"},
{"going", "returning"},
} {
snap := NewSnapshot()
applyPush(t, &snap, `<ctl td="ChargeState"><charge type="`+tc.charge+`"/></ctl>`)
applyPush(t, &snap, `<ctl td="CleanReport"><clean type="stop" speed="standard" st=" " rsn=" "/></ctl>`)
if got := snap.State.State; got != tc.want {
t.Errorf("stop + %s: state = %q, want %q", tc.charge, got, tc.want)
}
}
}
func TestSetCleanSpeedStoresRequested(t *testing.T) {
snap := NewSnapshot()
// A SetCleanSpeed result ret="ok" errno="" has no speed echo; the
// requested fan is stored anyway.
in := ctl.Inbound{Kind: ctl.KindResult, TD: "SetCleanSpeed", Ret: "ok", Errno: strPtr(""), Attrs: map[string]string{}}
if err := Apply(&snap, in, "strong"); err != nil {
t.Fatalf("Apply SetCleanSpeed result: %v", err)
}
if snap.Facts.Fan != "strong" || snap.State.FanSpeed != "strong" {
t.Fatalf("fan = %q/%q, want strong", snap.Facts.Fan, snap.State.FanSpeed)
}
// A later CleanReport at standard replaces it.
applyPush(t, &snap, `<ctl td="CleanReport"><clean type="auto" speed="standard" st=" " rsn=" "/></ctl>`)
if snap.Facts.Fan != "standard" {
t.Fatalf("fan = %q, want standard from CleanReport", snap.Facts.Fan)
}
// An unknown speed is rejected and leaves the fan untouched.
bad := ctl.Inbound{Kind: ctl.KindPush, TD: "CleanReport", CleanAttrs: map[string]string{"type": "auto", "speed": "ludicrous"}}
if err := Apply(&snap, bad, ""); err == nil {
t.Fatal("Apply with unknown speed returned nil error")
}
if snap.Facts.Fan != "standard" {
t.Fatalf("fan = %q after invalid speed, want standard", snap.Facts.Fan)
}
}
func TestBareBatteryNotAckedAndCidCompletesOnce(t *testing.T) {
r := newRig(t)
r.stanza(t, `<iq to="`+testCtlJID+`" type="set" id="46"><query xmlns="com:ctl"><battery power="076"/></query></iq>`)
flushActor(t, r.actor)
if n := len(r.send.all()); n != 0 {
t.Fatalf("bare battery produced %d outbound stanzas, want 0", n)
}
if b := r.pub.last().Attributes.BatteryLevel; b == nil || *b != 76 {
t.Fatalf("battery_level = %v, want 76", b)
}
c := ctl.NewCorrelator(nil)
c.Register("1", "00000042", "GetCleanState", true)
c.Register("2", "00000042", "GetCleanState", true)
if _, ok := c.CompleteCID("00000042"); !ok {
t.Fatal("first cid completion failed")
}
if _, ok := c.CompleteCID("00000042"); ok {
t.Fatal("cid completed twice")
}
c.Register("3", "00000042", "GetCleanState", true)
if _, ok := c.CompleteCID("00000042"); !ok {
t.Fatal("cid was not reusable after completion")
}
}
func TestNoAckForUnsolicitedIQSet(t *testing.T) {
r := newRig(t)
r.actor.AnnounceOK(session.ReadyEvent{Generation: 1, JID: testBotJID, Serial: testSerial, Send: r.send.send})
flushActor(t, r.actor)
r.stanza(t, iqSet(`<ctl td="CleanReport"><clean type="auto" speed="standard" st=" " rsn=" "/></ctl>`))
r.stanza(t, ctlResult(stanzaCID(t, r.send.all()[0]), "ok", ""))
flushActor(t, r.actor)
for _, s := range r.send.all() {
if strings.Contains(s, `type="result"`) {
t.Fatalf("actor sent an iq result for an unsolicited iq-set: %s", s)
}
}
}
func TestOutstandingCidFailedOnReplace(t *testing.T) {
r := newRig(t)
if err := r.actor.Submit(context.Background(), Command{Name: "start"}); err != nil {
t.Fatalf("Submit: %v", err)
}
sends := r.send.all()
cid := stanzaCID(t, sends[len(sends)-1])
r.actor.SessionDown(session.DownEvent{Generation: 1, JID: testBotJID, Serial: testSerial, Reason: session.ReasonReplaced})
flushActor(t, r.actor)
traces := r.trace.all()
if len(traces) != 1 || traces[0].Ret != "connection-lost" || traces[0].CID != cid {
t.Fatalf("traces after SessionDown = %+v", traces)
}
if got := r.pub.last().Attributes.LastCommandError; got == nil || *got != "connection-lost" {
t.Fatalf("last_command_error = %v, want connection-lost", got)
}
// A new generation does not resurrect the old cid: a late result for it
// is ignored.
r.actor.SessionReady(session.ReadyEvent{Generation: 2, JID: testBotJID, Serial: testSerial, Send: r.send.send})
flushActor(t, r.actor)
r.stanza(t, ctlResult(cid, "ok", ""))
r.actor.Stanza(session.StanzaEvent{Generation: 2, JID: testBotJID, Serial: testSerial, Stanza: []byte(ctlResult(cid, "ok", ""))})
flushActor(t, r.actor)
if n := len(r.trace.all()); n != 1 {
t.Fatalf("late result produced %d traces total, want 1", n)
}
if got := r.pub.last().Attributes.LastCommandError; got == nil || *got != "connection-lost" {
t.Fatalf("last_command_error = %v, want connection-lost still", got)
}
}
func TestCommandTimeoutExpiresCorrelation(t *testing.T) {
r := newRig(t)
if err := r.actor.Submit(context.Background(), Command{Name: "start"}); err != nil {
t.Fatalf("Submit: %v", err)
}
r.clock.advance(ctl.CommandTimeout + time.Second)
r.trace.wait(t)
tr := r.trace.all()[0]
if tr.Ret != "timeout" || tr.Phase != "result" {
t.Fatalf("expire trace = %+v, want timeout/result", tr)
}
if got := r.pub.last().Attributes.LastCommandError; got == nil || *got != "timeout" {
t.Fatalf("last_command_error = %v, want timeout", got)
}
}
// A command registered 29s after actor start must expire at its own 30-second
// deadline, not at the next coarse timer tick.
func TestCorrelationExpiresAtOwnDeadline(t *testing.T) {
r := newRig(t)
r.clock.advance(29 * time.Second)
if err := r.actor.Submit(context.Background(), Command{Name: "start"}); err != nil {
t.Fatalf("Submit: %v", err)
}
// 29s after registration: not yet expired.
r.clock.advance(29 * time.Second)
flushActor(t, r.actor)
if traces := r.trace.all(); len(traces) != 0 {
t.Fatalf("command expired early: %+v", traces)
}
// Registration + 30s: timeout fires on the command's own deadline.
r.clock.advance(time.Second)
r.trace.wait(t)
tr := r.trace.all()[0]
if tr.Ret != "timeout" || tr.Phase != "result" {
t.Fatalf("expire trace = %+v, want timeout/result", tr)
}
if got := r.pub.last().Attributes.LastCommandError; got == nil || *got != "timeout" {
t.Fatalf("last_command_error = %v, want timeout", got)
}
}
// A SetCleanSpeed result ret="ok" with no speed echo stores the requested fan,
// and the cid entry is consumed. This drives the actor path end to end; the
// Apply-level case lives in TestSetCleanSpeedStoresRequested.
func TestSetCleanSpeedResultStoresFan(t *testing.T) {
r := newRig(t)
if err := r.actor.Submit(context.Background(), Command{Name: "set_fan_speed", Args: map[string]string{"speed": "strong"}}); err != nil {
t.Fatalf("Submit: %v", err)
}
sends := r.send.all()
cid := stanzaCID(t, sends[len(sends)-1])
flushActor(t, r.actor)
if got := r.actor.requestedFan[cid]; got != "strong" {
t.Fatalf("requestedFan[%s] = %q, want strong", cid, got)
}
r.stanza(t, fmt.Sprintf(`<iq to="%s" type="set" id="9"><query xmlns="com:ctl"><ctl id="%s" ret="ok" errno=""/></query></iq>`, testCtlJID, cid))
flushActor(t, r.actor)
if got := r.pub.last().Facts.Fan; got != "strong" {
t.Fatalf("fan = %q, want strong", got)
}
if got := r.pub.last().State.FanSpeed; got != "strong" {
t.Fatalf("fan_speed = %q, want strong", got)
}
if _, ok := r.actor.requestedFan[cid]; ok {
t.Fatal("requestedFan not cleared after result")
}
}
// A SetCleanSpeed result that never arrives must not pin the requested fan:
// the timeout clears the cid entry and leaves the fan unchanged.
func TestSetCleanSpeedExpireClearsRequestedFan(t *testing.T) {
r := newRig(t)
if err := r.actor.Submit(context.Background(), Command{Name: "set_fan_speed", Args: map[string]string{"speed": "strong"}}); err != nil {
t.Fatalf("Submit: %v", err)
}
sends := r.send.all()
cid := stanzaCID(t, sends[len(sends)-1])
flushActor(t, r.actor)
if got := r.actor.requestedFan[cid]; got != "strong" {
t.Fatalf("requestedFan[%s] = %q, want strong", cid, got)
}
r.clock.advance(ctl.CommandTimeout + time.Second)
r.trace.wait(t)
flushActor(t, r.actor)
if _, ok := r.actor.requestedFan[cid]; ok {
t.Fatal("requestedFan retained after timeout")
}
if got := r.pub.last().Attributes.LastCommandError; got == nil || *got != "timeout" {
t.Fatalf("last_command_error = %v, want timeout", got)
}
if got := r.pub.last().Facts.Fan; got != "standard" {
t.Fatalf("fan = %q, want unchanged standard", got)
}
}
func TestSubmitOfflineAndRejected(t *testing.T) {
r := newRig(t)
// Drop the generation so the actor is offline.
r.actor.SessionDown(session.DownEvent{Generation: 1, JID: testBotJID, Serial: testSerial, Reason: session.ReasonTCPClose})
flushActor(t, r.actor)
if err := r.actor.Submit(context.Background(), Command{Name: "start"}); err == nil || err.Error() != "offline" {
t.Fatalf("offline Submit err = %v, want offline", err)
}
if got := r.pub.last().Attributes.LastCommandError; got == nil || *got != "offline" {
t.Fatalf("last_command_error = %v, want offline", got)
}
// Unknown names and bad fan speeds are rejected without a write.
r.actor.SessionReady(session.ReadyEvent{Generation: 2, JID: testBotJID, Serial: testSerial, Send: r.send.send})
flushActor(t, r.actor)
before := len(r.send.all())
if err := r.actor.Submit(context.Background(), Command{Name: "dance"}); err == nil || err.Error() != "rejected:dance" {
t.Fatalf("unknown Submit err = %v, want rejected:dance", err)
}
if err := r.actor.Submit(context.Background(), Command{Name: "set_fan_speed", Args: map[string]string{"speed": "turbo"}}); err == nil || err.Error() != "rejected:set_fan_speed" {
t.Fatalf("bad speed Submit err = %v, want rejected:set_fan_speed", err)
}
if got := r.pub.last().Attributes.LastCommandError; got == nil || *got != "rejected:set_fan_speed" {
t.Fatalf("last_command_error = %v, want rejected:set_fan_speed", got)
}
if n := len(r.send.all()); n != before {
t.Fatalf("rejected commands wrote %d stanzas", n-before)
}
}
func TestStaleSendFailsCommandOnce(t *testing.T) {
r := newRig(t)
r.send.err = session.ErrStale
err := r.actor.Submit(context.Background(), Command{Name: "start"})
if err == nil || err.Error() != "connection-lost" {
t.Fatalf("stale Submit err = %v, want connection-lost", err)
}
// One write attempt, no retry on the stale generation.
if n := len(r.send.all()); n != 1 {
t.Fatalf("%d write attempts, want 1", n)
}
traces := r.trace.all()
if len(traces) != 1 || traces[0].Ret != "connection-lost" {
t.Fatalf("traces = %+v, want one connection-lost", traces)
}
if got := r.pub.last().Attributes.LastCommandError; got == nil || *got != "connection-lost" {
t.Fatalf("last_command_error = %v, want connection-lost", got)
}
}
func TestRepublishFullObjects(t *testing.T) {
r := newRig(t)
r.stanza(t, iqSet(`<ctl td="BatteryInfo"><battery power="076"/></ctl>`))
r.pub.wait(t)
snap := r.pub.last()
stateJSON, err := json.Marshal(snap.State)
if err != nil {
t.Fatalf("marshal state: %v", err)
}
var state map[string]any
if err := json.Unmarshal(stateJSON, &state); err != nil {
t.Fatalf("unmarshal state: %v", err)
}
if state["state"] != "idle" || state["fan_speed"] != "standard" {
t.Fatalf("state doc = %s", stateJSON)
}
if len(state) != 2 {
t.Fatalf("state doc has %d keys, want state and fan_speed: %s", len(state), stateJSON)
}
attrJSON, err := json.Marshal(snap.Attributes)
if err != nil {
t.Fatalf("marshal attributes: %v", err)
}
var attrs map[string]any
if err := json.Unmarshal(attrJSON, &attrs); err != nil {
t.Fatalf("unmarshal attributes: %v", err)
}
for _, key := range []string{
"battery_level", "side_brush", "main_brush", "filter",
"lifespan_total", "clean_type", "charge_state",
"last_error", "last_command_error", "schedules",
} {
if _, ok := attrs[key]; !ok {
t.Errorf("attributes missing %q: %s", key, attrJSON)
}
}
if attrs["battery_level"] != float64(76) {
t.Errorf("battery_level = %v, want 76", attrs["battery_level"])
}
if !strings.Contains(string(attrJSON), `"schedules":[]`) {
t.Errorf("schedules must marshal as []: %s", attrJSON)
}
lt, ok := attrs["lifespan_total"].(map[string]any)
if !ok {
t.Fatalf("lifespan_total = %v, want object", attrs["lifespan_total"])
}
for _, key := range []string{"side_brush", "main_brush", "filter"} {
if v, ok := lt[key]; !ok || v != nil {
t.Errorf("lifespan_total.%s = %v, want null", key, v)
}
}
}
// requireDiag finds the diagnostic emitted for an exact inbound stanza whose
// reason contains want.
func requireDiag(t *testing.T, dg *diagRec, xml, want string) session.Diagnostic {
t.Helper()
for _, d := range dg.all() {
if string(d.XML) == xml && strings.Contains(d.Reason, want) {
if d.Kind != "unparsed" || d.Direction != session.DirectionIn {
t.Fatalf("diagnostic = %+v, want kind unparsed direction in", d)
}
return d
}
}
t.Fatalf("no diagnostic for %q with reason containing %q in %+v", xml, want, dg.all())
return session.Diagnostic{}
}
func TestUnparsedDiagnostics(t *testing.T) {
r := newRig(t)
// Unknown ctl td: diagnosed as unparsed, no mutation, no republish.
mystery := iqSet(`<ctl td="Mystery"><zap/></ctl>`)
r.stanza(t, mystery)
flushActor(t, r.actor)
requireDiag(t, r.diag, mystery, "Mystery")
if n := r.pub.count(); n != 0 {
t.Fatalf("unknown ctl produced %d republishes, want 0", n)
}
// Malformed XML.
malformed := `<iq type="set"><query><ctl td="x"`
r.stanza(t, malformed)
flushActor(t, r.actor)
requireDiag(t, r.diag, malformed, "malformed")
// Well-formed stanza that is not a ctl shape.
unknownShape := `<iq type="get" id="5"><query/></iq>`
r.stanza(t, unknownShape)
flushActor(t, r.actor)
requireDiag(t, r.diag, unknownShape, "unparsed")
// Out-of-range battery: diagnosed, attribute untouched, no republish.
badBattery := iqSet(`<ctl td="BatteryInfo"><battery power="999"/></ctl>`)
pubs := r.pub.count()
r.stanza(t, badBattery)
flushActor(t, r.actor)
requireDiag(t, r.diag, badBattery, "999")
if n := r.pub.count(); n != pubs {
t.Fatalf("invalid battery produced a republish")
}
if b := r.pub.last().Attributes.BatteryLevel; b != nil {
t.Fatalf("battery_level = %v, want null after invalid report", *b)
}
// Nonblank CleanReport st/rsn: diagnosed with key/value, while the
// canonical clean type and fan still apply and republish.
report := iqSet(`<ctl td="CleanReport"><clean type="auto" speed="strong" st="x" rsn="wheels"/></ctl>`)
pubs = r.pub.count()
r.stanza(t, report)
flushActor(t, r.actor)
d := requireDiag(t, r.diag, report, "st")
if !strings.Contains(d.Reason, "rsn") || !strings.Contains(d.Reason, "wheels") {
t.Fatalf("diagnostic reason %q missing rsn detail", d.Reason)
}
if n := r.pub.count(); n <= pubs {
t.Fatal("valid canonical fields were not republished")
}
snap := r.pub.last()
if got := snap.Facts.CleanType; got != "auto" {
t.Fatalf("clean_type = %q, want auto", got)
}
if got := snap.Facts.Fan; got != "strong" {
t.Fatalf("fan = %q, want strong", got)
}
// The session is untouched throughout: a valid push still applies.
r.stanza(t, iqSet(`<ctl td="BatteryInfo"><battery power="050"/></ctl>`))
flushActor(t, r.actor)
if b := r.pub.last().Attributes.BatteryLevel; b == nil || *b != 50 {
t.Fatalf("battery_level = %v after diagnostics, want 50", b)
}
}
func TestNonblankGetCleanStateFieldsDiagnose(t *testing.T) {
snap := NewSnapshot()
in := ctl.Inbound{
Kind: ctl.KindResult,
TD: "GetCleanState",
CleanAttrs: map[string]string{
"type": "stop", "speed": "standard", "st": "h", "t": "123", "a": " ",
},
}
err := Apply(&snap, in, "")
if err == nil {
t.Fatal("Apply with nonblank t returned nil error")
}
if !strings.Contains(err.Error(), "t=") || !strings.Contains(err.Error(), "123") {
t.Fatalf("Apply error = %v, want t= key/value", err)
}
if strings.Contains(err.Error(), "a=") {
t.Fatalf("whitespace a must be ignored, got %v", err)
}
if snap.Facts.CleanType != "stop" {
t.Fatalf("clean_type = %q, canonical update must still apply", snap.Facts.CleanType)
}
}
func TestFleetCreatesAndForwards(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
pub := newPubRec()
send := &sendRec{}
f := NewFleet(ctx, testCtlJID, pub.pub, nil, nil, nil)
f.SessionReady(session.ReadyEvent{Generation: 1, JID: testBotJID, Serial: testSerial, Send: send.send})
a, ok := f.Actor(testBotJID)
if !ok {
t.Fatal("Fleet did not create an actor on SessionReady")
}
if _, ok := f.Actor("other@155.ecorobot.net/atom"); ok {
t.Fatal("Fleet returned an actor for an unknown JID")
}
// Events for the JID reach the actor; events for other JIDs are dropped.
f.Stanza(session.StanzaEvent{Generation: 1, JID: testBotJID, Serial: testSerial, Stanza: []byte(iqSet(`<ctl td="BatteryInfo"><battery power="060"/></ctl>`))})
flushActor(t, a)
if b := pub.last().Attributes.BatteryLevel; b == nil || *b != 60 {
t.Fatalf("battery_level = %v, want 60", b)
}
}
func strPtr(v string) *string { return &v }