208 lines
5.3 KiB
Go
208 lines
5.3 KiB
Go
package asks
|
|
|
|
import (
|
|
"encoding/json"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
stdio "git.novox.be/novox/mesh-sdk/go"
|
|
"git.novox.be/novox/mesh-sdk/go/asker"
|
|
"git.novox.be/novox/mesh-sdk/go/asks"
|
|
"github.com/nats-io/nats.go"
|
|
)
|
|
|
|
// labAsker is the module `lab-asker` asking the operator, as any module does: the SDK's asker client
|
|
// (mesh-sdk go/asker) over the lab's bus, on that module's own composed credential. What it performs on a
|
|
// warrant is one switch flipped, bound when it asked, and counted: the proof that an approval does exactly
|
|
// that one thing, once.
|
|
type labAsker struct {
|
|
t *testing.T
|
|
client *asker.Client
|
|
nc *nats.Conn
|
|
|
|
mu sync.Mutex
|
|
warrants []asks.Warrant
|
|
handled []asker.Handled
|
|
acts []map[string]string // what was performed
|
|
}
|
|
|
|
// flip is the act an approve option binds.
|
|
func flip(which string) asks.Act {
|
|
return asks.Act{"verb": "lab.flip", "arg.switch": which}
|
|
}
|
|
|
|
func newLabAsker(t *testing.T, l *lab) *labAsker {
|
|
nc, _ := l.as(machine + ".lab-asker")
|
|
js, err := nc.JetStream()
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
a := &labAsker{t: t, nc: nc}
|
|
a.client = &asker.Client{Module: "lab-asker", State: &memState{vals: map[string]json.RawMessage{}, revs: map[string]uint64{}},
|
|
Now: time.Now,
|
|
Publish: func(subject string, body any, id string) error {
|
|
raw, err := json.Marshal(body)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
_, err = js.Publish(subject, raw, nats.MsgId("lab-asker."+id))
|
|
return err
|
|
}}
|
|
// Its warrants, on the one subject the bus lets it hear them on.
|
|
sub, err := nc.Subscribe(asks.DecidedSubject("lab-asker"), func(m *nats.Msg) {
|
|
var w asks.Warrant
|
|
if err := json.Unmarshal(m.Data, &w); err != nil {
|
|
t.Errorf("a word on the asker's decided subject is not a warrant: %v", err)
|
|
return
|
|
}
|
|
a.take(w)
|
|
})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
t.Cleanup(func() { _ = sub.Unsubscribe() })
|
|
return a
|
|
}
|
|
|
|
// take hands one warrant to the SDK's Take, which performs the option once, and only the act it bound.
|
|
func (a *labAsker) take(w asks.Warrant) asker.Handled {
|
|
h, err := a.client.Take(w, func(o asks.Option, w asks.Warrant) error {
|
|
if o.Level == asks.Acknowledge {
|
|
return nil // decline: nothing is performed
|
|
}
|
|
act := flip(o.ID)
|
|
if err := o.Performs(act); err != nil {
|
|
return err
|
|
}
|
|
a.mu.Lock()
|
|
a.acts = append(a.acts, act)
|
|
a.mu.Unlock()
|
|
return nil
|
|
})
|
|
a.mu.Lock()
|
|
a.warrants = append(a.warrants, w)
|
|
a.handled = append(a.handled, h)
|
|
a.mu.Unlock()
|
|
if err != nil {
|
|
a.t.Logf("the asker refused a warrant for %s: %v", w.Ask, err)
|
|
}
|
|
return h
|
|
}
|
|
|
|
// ask asks the operator: one approve option binding a switch, one decline that acknowledges.
|
|
func (a *labAsker) ask(id, headline string, expires time.Duration) asks.Ask {
|
|
a.t.Helper()
|
|
binds, err := asks.ActDigest(flip("approve"))
|
|
if err != nil {
|
|
a.t.Fatal(err)
|
|
}
|
|
q := asks.Ask{ID: id, Headline: headline,
|
|
Explanation: "Needs you: approve or decline. The lab asks whether to flip one switch.",
|
|
Who: asks.Operator, Expires: time.Now().Add(expires), OnExpiry: "nothing is flipped",
|
|
Options: []asks.Option{
|
|
{ID: "approve", Label: "Approve", Does: "the switch is flipped, once", Level: asks.Approve, Binds: binds},
|
|
{ID: "decline", Label: "Decline", Does: "nothing is done", Level: asks.Acknowledge},
|
|
}}
|
|
if err := a.client.Ask(q); err != nil {
|
|
a.t.Fatalf("the asker could not ask %s: %v", id, err)
|
|
}
|
|
return q
|
|
}
|
|
|
|
// warrantFor waits for the router's word on one ask.
|
|
func (a *labAsker) warrantFor(id string, within time.Duration) asks.Warrant {
|
|
a.t.Helper()
|
|
deadline := time.Now().Add(within)
|
|
for time.Now().Before(deadline) {
|
|
a.mu.Lock()
|
|
for _, w := range a.warrants {
|
|
if w.Ask == id {
|
|
a.mu.Unlock()
|
|
return w
|
|
}
|
|
}
|
|
a.mu.Unlock()
|
|
time.Sleep(100 * time.Millisecond)
|
|
}
|
|
a.t.Fatalf("no word on the ask %s within %s", id, within)
|
|
return asks.Warrant{}
|
|
}
|
|
|
|
// noWordOn says the router said nothing on an ask for a while.
|
|
func (a *labAsker) noWordOn(id string, wait time.Duration) bool {
|
|
time.Sleep(wait)
|
|
a.mu.Lock()
|
|
defer a.mu.Unlock()
|
|
for _, w := range a.warrants {
|
|
if w.Ask == id {
|
|
return false
|
|
}
|
|
}
|
|
return true
|
|
}
|
|
|
|
func (a *labAsker) performed() []map[string]string {
|
|
a.mu.Lock()
|
|
defer a.mu.Unlock()
|
|
return append([]map[string]string(nil), a.acts...)
|
|
}
|
|
|
|
func (a *labAsker) wordsOn(id string) int {
|
|
a.mu.Lock()
|
|
defer a.mu.Unlock()
|
|
n := 0
|
|
for _, w := range a.warrants {
|
|
if w.Ask == id {
|
|
n++
|
|
}
|
|
}
|
|
return n
|
|
}
|
|
|
|
// memState is the asker's own record, kept in memory: the lab's asker restarts never.
|
|
type memState struct {
|
|
mu sync.Mutex
|
|
vals map[string]json.RawMessage
|
|
revs map[string]uint64
|
|
n uint64
|
|
}
|
|
|
|
func (m *memState) Get(k string) (*stdio.StateEntry, error) {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
v, ok := m.vals[k]
|
|
if !ok {
|
|
return nil, nil
|
|
}
|
|
return &stdio.StateEntry{Key: k, Value: v, Revision: m.revs[k]}, nil
|
|
}
|
|
|
|
func (m *memState) Create(k string, v any) (uint64, error) {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
if _, ok := m.vals[k]; ok {
|
|
return 0, stdio.ErrChanged
|
|
}
|
|
return m.put(k, v)
|
|
}
|
|
|
|
func (m *memState) Update(k string, v any, rev uint64) (uint64, error) {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
if rev == 0 || m.revs[k] != rev {
|
|
return 0, stdio.ErrChanged
|
|
}
|
|
return m.put(k, v)
|
|
}
|
|
|
|
func (m *memState) put(k string, v any) (uint64, error) {
|
|
raw, err := json.Marshal(v)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
m.n++
|
|
m.vals[k], m.revs[k] = raw, m.n
|
|
return m.n, nil
|
|
}
|