package ocpp import ( "context" "net/http" "net/http/httptest" "testing" "time" ) // fakeUpstream is a stand-in for Anker's cloud CSMS: it upgrades, records every // CALL it receives, and answers each with {status:"Accepted"} (plus boot fields). func fakeUpstream(t *testing.T) (*httptest.Server, <-chan Message) { t.Helper() recv := make(chan Message, 64) srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { conn, err := Upgrade(w, r) if err != nil { t.Errorf("upstream upgrade: %v", err) return } for { data, err := conn.ReadMessage() if err != nil { return } msg, err := DecodeMessage(data) if err != nil { continue } if msg.Type != MessageTypeCall { continue } select { case recv <- msg: default: } out, _ := EncodeCallResult(msg.ID, map[string]any{ "status": "Accepted", "currentTime": time.Now().UTC().Format(time.RFC3339), "interval": 300, }) _ = conn.WriteMessage(out) } })) return srv, recv } func TestProxyModeForwardsAndInjects(t *testing.T) { upSrv, upRecv := fakeUpstream(t) defer upSrv.Close() upURL := wsURL(upSrv.URL, "/ocpp/CP-A5191") csms := NewCSMS(nil) sessCh := make(chan *Session, 1) dvSrv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { cpConn, err := Upgrade(w, r) if err != nil { t.Errorf("dv upgrade: %v", err) return } upConn, err := DialUpstream(context.Background(), upURL, BasicAuthHeader("CP-A5191", "secret")) if err != nil { t.Errorf("dial upstream: %v", err) _ = cpConn.Close() return } sess, err := csms.Accept("CP-A5191", ModeProxy, cpConn, upConn) if err != nil { t.Errorf("accept proxy: %v", err) return } sessCh <- sess })) defer dvSrv.Close() cp := dialSimCP(t, wsURL(dvSrv.URL, "/ocpp/CP-A5191"), nil) defer cp.close() // Charger→upstream: BootNotification is forwarded and the upstream's answer // comes back to the charger. boot := cp.call(t, "BootNotification", map[string]any{"chargePointModel": "A5191"}) assertStatus(t, boot, "Accepted") if got := waitMsg(t, upRecv); got.Action != "BootNotification" { t.Fatalf("upstream received %q, want BootNotification", got.Action) } var sess *Session select { case sess = <-sessCh: case <-time.After(2 * time.Second): t.Fatal("proxy session never registered") } // A forwarded StatusNotification is tapped for the snapshot. cp.call(t, "StatusNotification", map[string]any{"connectorId": 1, "status": "Charging", "errorCode": "NoError"}) waitMsg(t, upRecv) // forwarded upstream if s := sess.Snapshot(); s.ConnectorStatus != "Charging" { t.Errorf("proxy snapshot status = %q, want Charging", s.ConnectorStatus) } // Injected control: DriverVault issues RemoteStart. The charger must receive // it and answer, and it must NOT leak to the upstream CSMS. status, err := sess.RemoteStartTransaction(context.Background(), "TAG", 1) if err != nil { t.Fatalf("inject RemoteStart: %v", err) } if status != "Accepted" { t.Errorf("injected RemoteStart status = %q, want Accepted", status) } if got := cp.waitRecv(t); got.Action != "RemoteStartTransaction" { t.Fatalf("charger received %q, want RemoteStartTransaction", got.Action) } select { case leaked := <-upRecv: t.Fatalf("injected call leaked to upstream: %q", leaked.Action) case <-time.After(300 * time.Millisecond): // good — nothing forwarded upstream } } func waitMsg(t *testing.T, ch <-chan Message) Message { t.Helper() select { case m := <-ch: return m case <-time.After(5 * time.Second): t.Fatal("timeout awaiting message") return Message{} } }