Files
backend_fiesta/messaging/posmqtt_test.go
Suriya cd2459dbb6 Let an admin create till staff from the console, through the same code
An admin sets a shop up from a browser; a supervisor adds a cashier at the
counter. Both had to be possible, and only the second one was.

So the console gets createposuser / updateposuser / getposusers /
deleteposuser, under both /v1/web/tenants and /v1/mob/tenants — calling the
same service methods `/pos/users` calls. Not a parallel implementation: a
supervisor created from a browser is the same row, with the same PIN rules, the
same duplicate check and the same identity-column allocation, as one created at
a till. Two paths writing one table is precisely how the two stop matching, and
this codebase already had that happen once.

configid is inferred rather than asked for. It is a number nobody looks up, it
varies per tenant — 1087's accounts are spread across 1, 6 and 15 — and getting
it wrong creates somebody who cannot sign into the portal their colleagues use
and is invisible to half the platform's queries.

/posroles is served rather than left to the console to hardcode. A console that
knew supervisor was 7 would be wrong the day that changed and would have no way
to find out.

The outlet is the real difference between the two doors. A terminal proves it
with a signed token; the console asserts it, and is checked against the tenant
before anything is written. That is weaker, and it is worth being plain about:
these mint till credentials on an unauthenticated request, exactly like every
other route in the /v1/web and /v1/mob groups, because there is no auth
middleware on the web API at all. Documented as the weakest point in the design
and flagged to move behind a session guard once the console can hold one. The
terminal routes are untouched by it.

Proven in a rolled-back transaction against live data: the console creates a
supervisor at 1135, that supervisor signs in by PIN with can_manage_staff true,
the till's /pos/staff sees them alongside the two created at the counter, and
0451, 1234 and a duplicate PIN are each refused with the same message the
terminal gives.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-06 20:42:27 +05:30

463 lines
15 KiB
Go

package messaging
import (
"context"
"encoding/json"
"errors"
"strings"
"sync"
"testing"
"time"
"nearle/models"
mqtt "github.com/eclipse/paho.mqtt.golang"
)
// These drive the real handlers through paho's own interfaces, so what is under
// test is the code that runs in production rather than a parallel
// reimplementation of it.
//
// No embedded broker: the infrastructure audit established that the broker is
// Mosquitto 2.1.2 and that it works. What was never established is whether
// *this* code acks the right terminal, and refuses to ack when the ingest
// failed — which is where a bug would cost a shop its takings.
// fakePosService lets a test decide what the ingest did.
type fakePosService struct {
ack *models.PosAck
err error
batches []models.PosOrderBatch
custBatch []models.PosCustomerBatch
heartbeats []models.PosHealth
healthErr error
}
func (f *fakePosService) IngestOrders(batch models.PosOrderBatch) (*models.PosAck, error) {
f.batches = append(f.batches, batch)
return f.ack, f.err
}
func (f *fakePosService) IngestCustomers(batch models.PosCustomerBatch) (*models.PosAck, error) {
f.custBatch = append(f.custBatch, batch)
return f.ack, f.err
}
func (f *fakePosService) Catalogue(string, string, int, int) (*models.PosCatalogueResponse, error) {
return nil, nil
}
func (f *fakePosService) RecordHealth(_ context.Context, health models.PosHealth) error {
f.heartbeats = append(f.heartbeats, health)
return f.healthErr
}
func (f *fakePosService) TerminalHealth(context.Context, string) (map[string]string, error) {
return nil, nil
}
func (f *fakePosService) Sales(models.PosSalesFilter) (*models.PosSalesPage, error) {
return nil, nil
}
func (f *fakePosService) SaleDetail(int, string) (*models.PosOrders, error) {
return nil, nil
}
func (f *fakePosService) SalesSummary(models.PosSalesFilter) (*models.PosSalesSummary, error) {
return nil, nil
}
// Sign-in plays no part over the broker: a terminal on MQTT authenticates to
// the broker itself, and the topic it publishes on already names its store.
// These exist to satisfy the interface, and returning "denied" is the safer
// stub — a fake that waved authorisation through could hide a real regression.
func (f *fakePosService) Login(models.PosLoginRequest) (*models.PosSession, error) {
return nil, nil
}
func (f *fakePosService) LocationAllowed(int, int) (bool, error) {
return false, nil
}
func (f *fakePosService) Staff(int, int) ([]models.PosStaffMember, error) {
return nil, nil
}
// Staff management plays no part over the broker — a terminal on MQTT publishes
// bills and nothing else. Denied rather than permitted, so a fake cannot hide a
// regression by waving authorisation through.
func (f *fakePosService) CreateUser(int, int, int, models.PosUserRequest) (*models.PosUser, error) {
return nil, nil
}
func (f *fakePosService) UpdateUser(int, int, models.PosUserRequest) (*models.PosUser, error) {
return nil, nil
}
func (f *fakePosService) ListUsers(int, int, bool) ([]models.PosUser, error) {
return nil, nil
}
func (f *fakePosService) DeactivateUser(int, int, int) error { return nil }
func (f *fakePosService) LoginWithPin(int, int, string) (*models.PosSession, error) {
return nil, nil
}
func (f *fakePosService) ConfigidFor(int) int { return 0 }
func (f *fakePosService) LocationHealth(context.Context, string) ([]map[string]string, error) {
return nil, nil
}
// ---------------------------------------------------------------- paho fakes
type published struct {
topic string
qos byte
retained bool
payload []byte
}
// fakeClient records what was published and nothing else.
type fakeClient struct {
mu sync.Mutex
sent []published
}
func (c *fakeClient) Publish(topic string, qos byte, retained bool, payload any) mqtt.Token {
c.mu.Lock()
defer c.mu.Unlock()
body, _ := payload.([]byte)
c.sent = append(c.sent, published{topic: topic, qos: qos, retained: retained, payload: body})
return doneToken{}
}
func (c *fakeClient) publishes() []published {
c.mu.Lock()
defer c.mu.Unlock()
return append([]published(nil), c.sent...)
}
func (c *fakeClient) IsConnected() bool { return true }
func (c *fakeClient) IsConnectionOpen() bool { return true }
func (c *fakeClient) Connect() mqtt.Token { return doneToken{} }
func (c *fakeClient) Disconnect(uint) {}
func (c *fakeClient) Subscribe(string, byte, mqtt.MessageHandler) mqtt.Token {
return doneToken{}
}
func (c *fakeClient) SubscribeMultiple(map[string]byte, mqtt.MessageHandler) mqtt.Token {
return doneToken{}
}
func (c *fakeClient) Unsubscribe(...string) mqtt.Token { return doneToken{} }
func (c *fakeClient) AddRoute(string, mqtt.MessageHandler) {}
func (c *fakeClient) OptionsReader() mqtt.ClientOptionsReader { return mqtt.ClientOptionsReader{} }
type doneToken struct{}
func (doneToken) Wait() bool { return true }
func (doneToken) WaitTimeout(time.Duration) bool { return true }
func (doneToken) Done() <-chan struct{} {
ch := make(chan struct{})
close(ch)
return ch
}
func (doneToken) Error() error { return nil }
type fakeMessage struct {
topic string
payload []byte
}
func (m fakeMessage) Duplicate() bool { return false }
func (m fakeMessage) Qos() byte { return 1 }
func (m fakeMessage) Retained() bool { return false }
func (m fakeMessage) Topic() string { return m.topic }
func (m fakeMessage) MessageID() uint16 { return 1 }
func (m fakeMessage) Payload() []byte { return m.payload }
func (m fakeMessage) Ack() {}
func consumerFor(svc *fakePosService) (*PosMqttConsumer, *fakeClient) {
client := &fakeClient{}
return &PosMqttConsumer{client: client, svc: svc}, client
}
func orderBatch(batchID string, ids ...string) []byte {
orders := make([]models.PosOrder, 0, len(ids))
for _, id := range ids {
orders = append(orders, models.PosOrder{Id: id, Invoicenumber: "INV-" + id})
}
body, _ := json.Marshal(models.PosOrderBatch{Schema: 1, Batchid: batchID, Orders: orders})
return body
}
// ---------------------------------------------------------------------- tests
func TestAckGoesBackToTheTerminalThatSent(t *testing.T) {
ack := models.NewPosAck("batch-1")
ack.Accept("order-a")
c, client := consumerFor(&fakePosService{ack: ack})
c.handleOrders(nil, fakeMessage{
topic: "nearle/pos/12/T4A9/order",
payload: orderBatch("batch-1", "order-a"),
})
sent := client.publishes()
if len(sent) != 1 {
t.Fatalf("published %d messages, want exactly 1", len(sent))
}
// Addressed to the till that sent it. A store-wide ack would tell every
// other counter that bills they never sent had landed.
if sent[0].topic != "nearle/pos/12/T4A9/ack" {
t.Errorf("ack topic = %q, want nearle/pos/12/T4A9/ack", sent[0].topic)
}
if sent[0].qos != 1 {
t.Errorf("ack qos = %d, want 1", sent[0].qos)
}
if sent[0].retained {
t.Error("the ack was retained; a stale ack replayed to a new session would retire bills that were never sent")
}
var got models.PosAck
if err := json.Unmarshal(sent[0].payload, &got); err != nil {
t.Fatalf("decode ack: %v", err)
}
if got.Batchid != "batch-1" {
t.Errorf("batch_id = %q, want batch-1", got.Batchid)
}
if len(got.Accepted) != 1 || got.Accepted[0] != "order-a" {
t.Errorf("accepted = %v, want [order-a]", got.Accepted)
}
}
func TestAFailedIngestIsNotAcked(t *testing.T) {
// The single most important behaviour here. An ack the ingest did not earn
// tells a terminal to delete a bill that was never banked.
c, client := consumerFor(&fakePosService{err: errors.New("database is having a bad minute")})
c.handleOrders(nil, fakeMessage{
topic: "nearle/pos/12/T4A9/order",
payload: orderBatch("batch-2", "order-a"),
})
if sent := client.publishes(); len(sent) != 0 {
t.Fatalf("a batch that failed to commit was acknowledged: %s", sent[0].payload)
}
}
func TestStoreAndTerminalComeFromTheTopic(t *testing.T) {
// The body is authoritative for nothing about identity. A till that could
// name a store in its payload could post sales into another shop's books.
svc := &fakePosService{ack: models.NewPosAck("batch-3")}
c, _ := consumerFor(svc)
c.handleOrders(nil, fakeMessage{
topic: "nearle/pos/44/T0001/order",
payload: orderBatch("batch-3", "order-x"),
})
if len(svc.batches) != 1 {
t.Fatalf("the batch never reached the ingest")
}
if got := svc.batches[0].Storeid; got != "44" {
t.Errorf("store_id = %q, want 44 (from the topic)", got)
}
if got := svc.batches[0].Terminalid; got != "T0001" {
t.Errorf("terminal_id = %q, want T0001", got)
}
}
func TestABodyCannotOverrideTheTopicIdentity(t *testing.T) {
// A till claiming to be somewhere else must not be believed.
svc := &fakePosService{ack: models.NewPosAck("batch-4")}
c, client := consumerFor(svc)
body, _ := json.Marshal(models.PosOrderBatch{
Schema: 1,
Batchid: "batch-4",
Storeid: "99", // a shop this till has no claim on
Orders: []models.PosOrder{{Id: "order-a"}},
})
c.handleOrders(nil, fakeMessage{topic: "nearle/pos/12/T4A9/order", payload: body})
// The ingest still resolves the store it was *told*, which is a known gap —
// but the ack must go back to the real terminal, so a forged store id
// cannot redirect another till's acknowledgements.
sent := client.publishes()
if len(sent) != 1 || sent[0].topic != "nearle/pos/12/T4A9/ack" {
t.Fatalf("ack went to %v, want nearle/pos/12/T4A9/ack", sent)
}
}
func TestAMalformedBatchIsDroppedWithoutAcking(t *testing.T) {
// Nothing to key an ack on, and nothing committed. Silence is correct: the
// till times out and re-sends.
svc := &fakePosService{ack: models.NewPosAck("x")}
c, client := consumerFor(svc)
c.handleOrders(nil, fakeMessage{
topic: "nearle/pos/12/T4A9/order",
payload: []byte("not json at all"),
})
if len(svc.batches) != 0 {
t.Error("an unreadable batch reached the ingest")
}
if sent := client.publishes(); len(sent) != 0 {
t.Errorf("an unreadable batch was acknowledged: %s", sent[0].payload)
}
}
func TestRegistrationsAckOnTheSameTopic(t *testing.T) {
ack := models.NewPosAck("cust-1")
ack.Accept("customer-a")
c, client := consumerFor(&fakePosService{ack: ack})
body, _ := json.Marshal(models.PosCustomerBatch{
Schema: 1,
Batchid: "cust-1",
Customers: []models.PosCustomer{{Id: "customer-a", Mobile: "9840012345", Name: "Meena"}},
})
c.handleCustomers(nil, fakeMessage{topic: "nearle/pos/12/T4A9/customer", payload: body})
sent := client.publishes()
if len(sent) != 1 || sent[0].topic != "nearle/pos/12/T4A9/ack" {
t.Fatalf("registration ack went to %v", sent)
}
}
func TestAHeartbeatIsRecordedAndNeverAcked(t *testing.T) {
// Presence is fire-and-forget. A till waiting on an ack for its heartbeat
// would be a till that a busy dashboard can block.
svc := &fakePosService{}
c, client := consumerFor(svc)
body, _ := json.Marshal(models.PosHealth{Status: "online", Pendingbills: 4, Todaybills: 37})
c.handleHealth(nil, fakeMessage{topic: "nearle/pos/12/T4A9/health", payload: body})
if len(svc.heartbeats) != 1 {
t.Fatalf("the heartbeat never reached the presence store")
}
got := svc.heartbeats[0]
if got.Terminalid != "T4A9" || got.Locationid != "12" {
t.Errorf("identity = %s/%s, want 12/T4A9 (from the topic)", got.Locationid, got.Terminalid)
}
if got.Pendingbills != 4 {
t.Errorf("pending_bills = %d, want 4", got.Pendingbills)
}
if sent := client.publishes(); len(sent) != 0 {
t.Error("a heartbeat was acknowledged; presence must be fire-and-forget")
}
}
func TestALastWillIsRecordedAsOffline(t *testing.T) {
// The broker publishes this on the till's behalf when it loses power. It
// carries nothing but a status, and that is the point — it is the only way
// to tell "closed for the night" from "unplugged".
svc := &fakePosService{}
c, _ := consumerFor(svc)
c.handleHealth(nil, fakeMessage{
topic: "nearle/pos/12/T4A9/health",
payload: []byte(`{"status":"offline"}`),
})
if len(svc.heartbeats) != 1 {
t.Fatal("the will never reached the presence store")
}
if got := svc.heartbeats[0].Status; got != "offline" {
t.Errorf("status = %q, want offline — a will must not be defaulted to online", got)
}
}
func TestAFailedPresenceWriteDoesNotStopTheTill(t *testing.T) {
// Redis being unreachable must cost the board, never a sale.
svc := &fakePosService{healthErr: errors.New("redis is down")}
c, client := consumerFor(svc)
c.handleHealth(nil, fakeMessage{
topic: "nearle/pos/12/T4A9/health",
payload: []byte(`{"status":"online"}`),
})
if sent := client.publishes(); len(sent) != 0 {
t.Error("a failed heartbeat produced a message back to the till")
}
}
func TestAnEmptyAckSerialisesAsAListNotNull(t *testing.T) {
// A terminal reading `null` for accepted treats the whole batch as
// unconfirmed and sends it again for ever.
body, err := json.Marshal(models.NewPosAck("batch-5"))
if err != nil {
t.Fatalf("marshal: %v", err)
}
if want := `"accepted":[]`; !strings.Contains(string(body), want) {
t.Errorf("ack serialised as %s, want it to contain %s", body, want)
}
}
func TestTopicIdentityRejectsShortTopics(t *testing.T) {
// A topic that names no terminal must yield nothing rather than a
// plausible-looking wrong answer that sends an ack to the wrong place.
for _, topic := range []string{"nearle/pos/order", "pos/12/T4A9/order", "", "nearle"} {
store, terminal := topicIdentity(topic)
if store != "" || terminal != "" {
t.Errorf("topicIdentity(%q) = %q/%q, want empty", topic, store, terminal)
}
}
store, terminal := topicIdentity("nearle/pos/12/T4A9/order")
if store != "12" || terminal != "T4A9" {
t.Errorf("topicIdentity = %q/%q, want 12/T4A9", store, terminal)
}
}
// MQTT has no queue groups: every subscriber gets every message. Three replicas
// all consuming would commit the same bill three times and publish three acks —
// harmless, because the ingest is idempotent, but three times the work.
func TestOnlyTheFirstStatefulSetReplicaConsumes(t *testing.T) {
cases := []struct {
name string
hostname string
override string
want bool
}{
{"statefulset ordinal 0", "fiesta-0", "", true},
{"statefulset ordinal 1", "fiesta-1", "", false},
{"statefulset ordinal 2", "fiesta-2", "", false},
{"double-digit ordinal", "fiesta-10", "", false},
// A Deployment pod has a random suffix, not an ordinal. Refusing to
// consume there would be a far more confusing failure than consuming.
{"deployment pod", "fiesta-7d4f9c8b6d-x2k9p", "", true},
{"bare container", "a1b2c3d4e5f6", "", true},
{"no hostname", "", "", true},
// The override settles it outright either way.
{"forced on", "fiesta-2", "always", true},
{"forced off", "fiesta-0", "never", false},
{"forced on via true", "fiesta-5", "true", true},
{"forced off via false", "fiesta-0", "false", false},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
t.Setenv("HOSTNAME", c.hostname)
t.Setenv("POS_MQTT_CONSUMER", c.override)
if got := posConsumerElected(); got != c.want {
t.Errorf("posConsumerElected() = %v, want %v (hostname %q, override %q)",
got, c.want, c.hostname, c.override)
}
})
}
}