diff options
| author | Jeff Halter <868228+jhalter@users.noreply.github.com> | 2026-07-10 09:48:18 -0700 |
|---|---|---|
| committer | Jeff Halter <868228+jhalter@users.noreply.github.com> | 2026-07-10 09:48:18 -0700 |
| commit | 21f24d24fd6f501b32f15a2bef41c89cc461f623 (patch) | |
| tree | ba4952fb82f3ef9c265228633f27536866e23931 /internal | |
| parent | ae44fb222ec73cae8441f5a5d9a21f8585fe587f (diff) | |
Overhaul regression testing: e2e suite, fuzzing, CI, and bug fixes
Add a protocol-level end-to-end suite (in-process fully wired server on
an ephemeral port pair, driven by hotline.Client over TCP) covering
handshake, login, public and private chat, message board, threaded
news, file list/download/upload, account admin, disconnect
notification, and shutdown broadcast. The harness retries on a fresh
port pair when another process steals a probed port before
ListenAndServe binds it, and Server gains WithConnectionRateLimit so
tests can disable the per-IP connection throttle.
Add native fuzz tests for Transaction, Field, and flattened file
object decoding, and fix the bugs the new tests surfaced:
- Transaction.Write panicked on out-of-range attacker-controlled size
fields, and transactionScanner's uint32 length addition could wrap
and yield a truncated token. The information fork size declared in
an untrusted fork header is now bounded too.
- Client keepalive read c.done unsynchronized while Disconnect
replaces it under the mutex.
- The shared Agreement's Seek+ReadAll login path raced concurrent
logins; the server now prefers an AgreementBytes() snapshot.
Fill unit-test gaps (main's config-copy helpers, file resume data,
ReloaderFunc, R2 error injection and env validation) and add a CI test
workflow (build/vet + race-enabled shuffled suite), fixed lint
workflow triggers with golangci-lint v2.6, and Makefile test/cover/
lint/fuzz targets.
Diffstat (limited to 'internal')
| -rw-r--r-- | internal/mobius/account_manager.go | 12 | ||||
| -rw-r--r-- | internal/mobius/agreement.go | 14 | ||||
| -rw-r--r-- | internal/mobius/config_test.go | 14 | ||||
| -rw-r--r-- | internal/mobius/e2e_test.go | 439 | ||||
| -rw-r--r-- | internal/mobius/integration_test.go | 429 | ||||
| -rw-r--r-- | internal/mobius/reload_test.go | 32 |
6 files changed, 917 insertions, 23 deletions
diff --git a/internal/mobius/account_manager.go b/internal/mobius/account_manager.go index 1265d9c..44c054a 100644 --- a/internal/mobius/account_manager.go +++ b/internal/mobius/account_manager.go @@ -13,18 +13,6 @@ import ( "gopkg.in/yaml.v3" ) -// loadFromYAMLFile loads data from a YAML file into the provided data structure. -func loadFromYAMLFile(path string, data interface{}) error { - fh, err := os.Open(path) - if err != nil { - return err - } - defer func() { _ = fh.Close() }() - - decoder := yaml.NewDecoder(fh) - return decoder.Decode(data) -} - // YAMLAccountManager implements AccountManager interface using YAML files for persistence. // It maintains an in-memory cache of accounts and synchronizes with YAML files on disk. type YAMLAccountManager struct { diff --git a/internal/mobius/agreement.go b/internal/mobius/agreement.go index d9cf2d4..0f652e5 100644 --- a/internal/mobius/agreement.go +++ b/internal/mobius/agreement.go @@ -63,7 +63,21 @@ func (a *Agreement) Read(p []byte) (int, error) { } func (a *Agreement) Seek(offset int64, _ int) (int64, error) { + a.mu.Lock() + defer a.mu.Unlock() + a.readOffset = int(offset) return 0, nil } + +// AgreementBytes returns a private copy of the agreement text. Unlike Seek+Read it touches no +// shared read offset, so concurrent logins can each obtain the agreement without racing. +func (a *Agreement) AgreementBytes() []byte { + a.mu.RLock() + defer a.mu.RUnlock() + + out := make([]byte, len(a.data)) + copy(out, a.data) + return out +} diff --git a/internal/mobius/config_test.go b/internal/mobius/config_test.go index bed739d..f090b5c 100644 --- a/internal/mobius/config_test.go +++ b/internal/mobius/config_test.go @@ -9,11 +9,7 @@ import ( func TestLoadConfig_InvalidBannerFileExtension(t *testing.T) { // Create a temporary directory for test files - tmpDir, err := os.MkdirTemp("", "mobius-config-test") - if err != nil { - t.Fatalf("Failed to create temp dir: %v", err) - } - defer os.RemoveAll(tmpDir) + tmpDir := t.TempDir() // Create a test config file with an invalid banner file extension configContent := ` @@ -28,7 +24,7 @@ FileRoot: "files" } // Attempt to load the config - _, err = LoadConfig(configPath) + _, err := LoadConfig(configPath) // Verify that we get the improved error message if err == nil { @@ -94,11 +90,7 @@ func TestLoadConfig_ValidBannerFileExtensions(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { // Create a temporary directory for test files - tmpDir, err := os.MkdirTemp("", "mobius-config-test") - if err != nil { - t.Fatalf("Failed to create temp dir: %v", err) - } - defer os.RemoveAll(tmpDir) + tmpDir := t.TempDir() // Create files subdirectory filesDir := filepath.Join(tmpDir, "files") diff --git a/internal/mobius/e2e_test.go b/internal/mobius/e2e_test.go new file mode 100644 index 0000000..ccd56f3 --- /dev/null +++ b/internal/mobius/e2e_test.go @@ -0,0 +1,439 @@ +package mobius + +// End-to-end protocol regression tests. Each top-level test starts a fresh in-process server (see +// integration_test.go for the harness) and drives it with the real hotline.Client over TCP. + +import ( + "bytes" + "encoding/binary" + "io" + "net" + "os" + "path/filepath" + "testing" + "time" + + "github.com/jhalter/mobius/hotline" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func mustReadAll(r io.Reader) []byte { + b, err := io.ReadAll(r) + if err != nil { + panic(err) + } + return b +} + +// readUploadedFile polls for an uploaded file to appear under the server's Files tree (the upload +// completes asynchronously on the transfer connection) and returns its data-fork content. +func readUploadedFile(t *testing.T, s *e2eServer, name string) []byte { + t.Helper() + path := filepath.Join(s.cfgDir, "Files", name) + deadline := time.Now().Add(5 * time.Second) + for time.Now().Before(deadline) { + if data, err := os.ReadFile(path); err == nil { + return data + } + time.Sleep(20 * time.Millisecond) + } + t.Fatalf("uploaded file %s never appeared", path) + return nil +} + +func newReq(tranType [2]byte, fields ...hotline.Field) hotline.Transaction { + return hotline.NewTransaction(tranType, [2]byte{0, 0}, fields...) +} + +func TestE2E_Handshake(t *testing.T) { + s := startE2EServer(t) + + t.Run("valid handshake succeeds", func(t *testing.T) { + c := hotline.NewClient("probe", NewTestLogger()) + conn, err := net.DialTimeout("tcp", s.addr, 5*time.Second) + require.NoError(t, err) + defer func() { _ = conn.Close() }() + c.Connection = conn + require.NoError(t, c.Handshake()) + }) + + t.Run("garbage handshake is rejected", func(t *testing.T) { + conn, err := net.DialTimeout("tcp", s.addr, 5*time.Second) + require.NoError(t, err) + defer func() { _ = conn.Close() }() + + _, err = conn.Write([]byte("not a valid hotline handshake!!!")) + require.NoError(t, err) + + // The server closes the connection on a bad handshake; the read returns EOF. + require.NoError(t, conn.SetReadDeadline(time.Now().Add(3*time.Second))) + buf := make([]byte, 8) + _, err = conn.Read(buf) + require.Error(t, err) + }) +} + +func TestE2E_Login(t *testing.T) { + s := startE2EServer(t) + + t.Run("guest login and agreement flow", func(t *testing.T) { + c := connectE2E(t, s.addr, "guest", "") + // A successful, fully-established session answers a request. The user name list is not + // access-gated, so a non-error reply proves login (and the agreement auto-reply) completed. + reply := c.roundTrip(newReq(hotline.TranGetUserNameList)) + assert.False(t, isError(reply), "an established session should answer TranGetUserNameList") + }) + + t.Run("wrong password is rejected", func(t *testing.T) { + c := hotline.NewClient("baduser", NewTestLogger()) + require.NoError(t, c.Connect(s.addr, e2eAdminLogin, "wrong-password")) + // The server replies to the login with an error, then closes the connection. + require.NoError(t, c.Connection.SetReadDeadline(time.Now().Add(3*time.Second))) + buf := make([]byte, 512) + n, err := c.Connection.Read(buf) + if err == nil { + // If a reply was delivered it must be an error reply. + var reply hotline.Transaction + _, werr := reply.Write(buf[:n]) + if werr == nil { + assert.True(t, isError(reply), "login with wrong password should return an error reply") + } + } + _ = c.Disconnect() + }) +} + +func TestE2E_Chat(t *testing.T) { + s := startE2EServer(t) + + alice := connectE2E(t, s.addr, "guest", "") + bob := connectE2E(t, s.addr, "guest", "") + + // A completed round trip proves a session is published to the client manager (publication + // happens before the server answers requests), so after these both peers are guaranteed to be + // in the chat broadcast set. Waiting for a login notification instead would be order-dependent: + // alice only hears about bob if bob logs in after alice is established. + require.False(t, isError(alice.roundTrip(newReq(hotline.TranGetUserNameList)))) + require.False(t, isError(bob.roundTrip(newReq(hotline.TranGetUserNameList)))) + + require.NoError(t, alice.client.Send(newReq(hotline.TranChatSend, + hotline.NewField(hotline.FieldData, []byte("hello from alice")), + ))) + + msg := bob.waitFor(hotline.TranChatMsg) + got := msg.GetField(hotline.FieldData) + require.NotNil(t, got) + assert.Contains(t, string(got.Data), "hello from alice") +} + +// findUserID picks the user with the given display name out of a TranGetUserNameList reply. +func findUserID(t *testing.T, reply hotline.Transaction, name string) [2]byte { + t.Helper() + for _, f := range reply.Fields { + if f.Type != hotline.FieldUsernameWithInfo { + continue + } + var u hotline.User + _, err := u.Write(f.Data) + require.NoError(t, err) + if u.Name == name { + return u.ID + } + } + t.Fatalf("user %q not in user list", name) + return [2]byte{} +} + +func TestE2E_PrivateChat(t *testing.T) { + s := startE2EServer(t) + + alice := connectE2ENamed(t, s.addr, "alice", "guest", "") + bob := connectE2ENamed(t, s.addr, "bob", "guest", "") + + // A completed round trip proves a session is published to the client manager, so bob is + // guaranteed to be in alice's user list below. + require.False(t, isError(bob.roundTrip(newReq(hotline.TranGetUserNameList)))) + list := alice.roundTrip(newReq(hotline.TranGetUserNameList)) + require.False(t, isError(list)) + bobID := findUserID(t, list, "bob") + + // Alice opens a private chat with bob; the reply carries the new chat ID. + invite := alice.roundTrip(newReq(hotline.TranInviteNewChat, + hotline.NewField(hotline.FieldUserID, bobID[:]), + )) + require.False(t, isError(invite)) + chatIDField := invite.GetField(hotline.FieldChatID) + require.NotNil(t, chatIDField) + chatID := chatIDField.Data + + // Bob receives the invitation, carrying the same chat ID. + inv := bob.waitFor(hotline.TranInviteToChat) + require.Equal(t, chatID, inv.GetField(hotline.FieldChatID).Data) + + // Bob joins; the join reply lists alice as an existing member. + join := bob.roundTrip(newReq(hotline.TranJoinChat, + hotline.NewField(hotline.FieldChatID, chatID), + )) + require.False(t, isError(join)) + var members []string + for _, f := range join.Fields { + if f.Type == hotline.FieldUsernameWithInfo { + var u hotline.User + _, err := u.Write(f.Data) + require.NoError(t, err) + members = append(members, u.Name) + } + } + assert.Contains(t, members, "alice", "join reply should list the chat's existing members") + + // Alice is notified that bob joined. + joined := alice.waitFor(hotline.TranNotifyChatChangeUser) + assert.Equal(t, chatID, joined.GetField(hotline.FieldChatID).Data) + + // Alice sets the subject; bob is notified. + require.NoError(t, alice.client.Send(newReq(hotline.TranSetChatSubject, + hotline.NewField(hotline.FieldChatID, chatID), + hotline.NewField(hotline.FieldChatSubject, []byte("secret plans")), + ))) + subj := bob.waitFor(hotline.TranNotifyChatSubject) + assert.Equal(t, "secret plans", string(subj.GetField(hotline.FieldChatSubject).Data)) + + // A message sent with the chat ID goes to the room's members (including the sender's echo), + // tagged with the chat ID. + require.NoError(t, alice.client.Send(newReq(hotline.TranChatSend, + hotline.NewField(hotline.FieldChatID, chatID), + hotline.NewField(hotline.FieldData, []byte("psst")), + ))) + msg := bob.waitFor(hotline.TranChatMsg) + assert.Equal(t, chatID, msg.GetField(hotline.FieldChatID).Data) + assert.Contains(t, string(msg.GetField(hotline.FieldData).Data), "psst") + echo := alice.waitFor(hotline.TranChatMsg) + assert.Contains(t, string(echo.GetField(hotline.FieldData).Data), "psst") + + // Bob leaves; alice is notified. + require.NoError(t, bob.client.Send(newReq(hotline.TranLeaveChat, + hotline.NewField(hotline.FieldChatID, chatID), + ))) + left := alice.waitFor(hotline.TranNotifyChatDeleteUser) + assert.Equal(t, chatID, left.GetField(hotline.FieldChatID).Data) + + // A declined invitation is announced to the chat's members. + invite2 := alice.roundTrip(newReq(hotline.TranInviteNewChat, + hotline.NewField(hotline.FieldUserID, bobID[:]), + )) + require.False(t, isError(invite2)) + inv2 := bob.waitFor(hotline.TranInviteToChat) + require.NoError(t, bob.client.Send(newReq(hotline.TranRejectChatInvite, + hotline.NewField(hotline.FieldChatID, inv2.GetField(hotline.FieldChatID).Data), + ))) + decline := alice.waitFor(hotline.TranChatMsg) + assert.Contains(t, string(decline.GetField(hotline.FieldData).Data), "bob declined invitation to chat") +} + +func TestE2E_MessageBoard(t *testing.T) { + s := startE2EServer(t) + c := connectE2E(t, s.addr, "guest", "") + + reply := c.roundTrip(newReq(hotline.TranGetMsgs)) + require.False(t, isError(reply)) + data := reply.GetField(hotline.FieldData) + require.NotNil(t, data) + assert.Contains(t, string(data.Data), "Test News Post", "message board should return the fixture content") +} + +func TestE2E_ThreadedNews(t *testing.T) { + s := startE2EServer(t) + c := connectE2E(t, s.addr, "guest", "") + + reply := c.roundTrip(newReq(hotline.TranGetNewsCatNameList, + hotline.NewField(hotline.FieldNewsPath, []byte{}), + )) + require.False(t, isError(reply)) + // The fixture ThreadedNews.yaml defines a "TestBundle" category at the root. + var found bool + for _, f := range reply.Fields { + if bytes.Contains(f.Data, []byte("TestBundle")) { + found = true + } + } + assert.True(t, found, "root news category list should include the fixture bundle") +} + +func TestE2E_FileList(t *testing.T) { + s := startE2EServer(t) + c := connectE2E(t, s.addr, "guest", "") + + reply := c.roundTrip(newReq(hotline.TranGetFileNameList, + hotline.NewField(hotline.FieldFilePath, []byte{}), + )) + require.False(t, isError(reply)) + + var found bool + for _, f := range reply.Fields { + if f.Type == hotline.FieldFileNameWithInfo && bytes.Contains(f.Data, []byte("testfile.txt")) { + found = true + } + } + assert.True(t, found, "file list should include the fixture testfile.txt") +} + +func TestE2E_FileDownload(t *testing.T) { + s := startE2EServer(t) + c := connectE2E(t, s.addr, "guest", "") + + reply := c.roundTrip(newReq(hotline.TranDownloadFile, + hotline.NewField(hotline.FieldFileName, []byte("testfile.txt")), + hotline.NewField(hotline.FieldFilePath, []byte{}), + )) + require.False(t, isError(reply)) + + refField := reply.GetField(hotline.FieldRefNum) + require.NotNil(t, refField, "download reply must carry a reference number") + var refNum [4]byte + copy(refNum[:], refField.Data) + + data := downloadOverTransferPort(t, s.xferAddr, refNum) + assert.Contains(t, string(data), "Hello, I'm a test file!", "downloaded payload should contain the fixture data fork") +} + +func TestE2E_FileUpload(t *testing.T) { + s := startE2EServer(t) + // Upload to the file root requires AccessUploadAnywhere, which the admin account has and guest + // does not. + c := connectE2E(t, s.addr, e2eAdminLogin, e2eAdminPass) + + content := []byte("uploaded by the e2e suite") + payload := buildUploadFFO("uploaded.txt", content) + + reply := c.roundTrip(newReq(hotline.TranUploadFile, + hotline.NewField(hotline.FieldFileName, []byte("uploaded.txt")), + hotline.NewField(hotline.FieldFilePath, []byte{}), + hotline.NewField(hotline.FieldTransferSize, u32(uint32(len(payload)))), + )) + require.False(t, isError(reply)) + + refField := reply.GetField(hotline.FieldRefNum) + require.NotNil(t, refField) + var refNum [4]byte + copy(refNum[:], refField.Data) + + uploadOverTransferPort(t, s.xferAddr, refNum, payload) + + // The uploaded file should land under the server's Files tree with the expected content. + uploaded := readUploadedFile(t, s, "uploaded.txt") + assert.Equal(t, content, uploaded) +} + +func TestE2E_AccountAdmin(t *testing.T) { + s := startE2EServer(t) + + t.Run("admin can create, read, and delete a user", func(t *testing.T) { + admin := connectE2E(t, s.addr, e2eAdminLogin, e2eAdminPass) + + create := admin.roundTrip(newReq(hotline.TranNewUser, + hotline.NewField(hotline.FieldUserLogin, hotline.EncodeString([]byte("newbie"))), + hotline.NewField(hotline.FieldUserName, []byte("New Bie")), + hotline.NewField(hotline.FieldUserPassword, hotline.EncodeString([]byte("pw"))), + hotline.NewField(hotline.FieldUserAccess, make([]byte, 8)), + )) + require.False(t, isError(create), "admin should be allowed to create a user") + + get := admin.roundTrip(newReq(hotline.TranGetUser, + hotline.NewField(hotline.FieldUserLogin, []byte("newbie")), + )) + require.False(t, isError(get)) + name := get.GetField(hotline.FieldUserName) + require.NotNil(t, name) + assert.Equal(t, "New Bie", string(name.Data)) + + del := admin.roundTrip(newReq(hotline.TranDeleteUser, + hotline.NewField(hotline.FieldUserLogin, hotline.EncodeString([]byte("newbie"))), + )) + assert.False(t, isError(del)) + }) + + t.Run("guest cannot create a user", func(t *testing.T) { + guest := connectE2E(t, s.addr, "guest", "") + reply := guest.roundTrip(newReq(hotline.TranNewUser, + hotline.NewField(hotline.FieldUserLogin, hotline.EncodeString([]byte("sneaky"))), + hotline.NewField(hotline.FieldUserName, []byte("Sneaky")), + hotline.NewField(hotline.FieldUserPassword, hotline.EncodeString([]byte("pw"))), + hotline.NewField(hotline.FieldUserAccess, make([]byte, 8)), + )) + assert.True(t, isError(reply), "guest lacks CreateUser access and must be rejected") + }) +} + +func TestE2E_DisconnectNotifiesPeers(t *testing.T) { + s := startE2EServer(t) + + alice := connectE2E(t, s.addr, "guest", "") + bob := connectE2E(t, s.addr, "guest", "") + + // Ensure both sessions are fully established (and thus both in the client manager) before alice + // leaves, so bob is guaranteed to be a peer that receives the delete notification. + require.False(t, isError(alice.roundTrip(newReq(hotline.TranGetUserNameList)))) + require.False(t, isError(bob.roundTrip(newReq(hotline.TranGetUserNameList)))) + + alice.cancel() + _ = alice.client.Disconnect() + + del := bob.waitFor(hotline.TranNotifyDeleteUser) + assert.Equal(t, hotline.TranNotifyDeleteUser, hotline.TranType(del.Type)) +} + +func TestE2E_ShutdownBroadcast(t *testing.T) { + if testing.Short() { + t.Skip("Shutdown pays a 3s broadcast-flush sleep") + } + s := startE2EServer(t) + c := connectE2E(t, s.addr, "guest", "") + // Ensure login has completed before triggering shutdown. + require.False(t, isError(c.roundTrip(newReq(hotline.TranGetUserNameList)))) + + go s.srv.Shutdown([]byte("server going down")) + + msg := c.waitFor(hotline.TranDisconnectMsg) + data := msg.GetField(hotline.FieldData) + require.NotNil(t, data) + assert.Contains(t, string(data.Data), "server going down") +} + +// --- upload payload construction ------------------------------------------------------------- + +func u32(v uint32) []byte { + b := make([]byte, 4) + binary.BigEndian.PutUint32(b, v) + return b +} + +// buildUploadFFO assembles the flattened file object bytes a client streams during an upload: +// FlatFileHeader + INFO fork header + info fork + DATA fork header + data fork. It mirrors the +// layout hotline.flattenedFileObject.ReadFrom expects (that type is unexported, so we build the +// bytes by hand from the exported information fork). +func buildUploadFFO(name string, content []byte) []byte { + infoFork := hotline.NewFlatFileInformationFork(name, [8]byte{}, "TEXT", "TTXT") + infoBody := mustReadAll(&infoFork) + + var buf bytes.Buffer + // FlatFileHeader: "FILP" + version 1 + 16 reserved + fork count 2. + buf.WriteString("FILP") + buf.Write([]byte{0, 1}) + buf.Write(make([]byte, 16)) + buf.Write([]byte{0, 2}) + + // INFO fork header + body. + buf.WriteString("INFO") + buf.Write(make([]byte, 8)) // compression + reserved + buf.Write(u32(uint32(len(infoBody)))) + buf.Write(infoBody) + + // DATA fork header + body. + buf.WriteString("DATA") + buf.Write(make([]byte, 8)) + buf.Write(u32(uint32(len(content)))) + buf.Write(content) + + return buf.Bytes() +} diff --git a/internal/mobius/integration_test.go b/internal/mobius/integration_test.go new file mode 100644 index 0000000..f1485eb --- /dev/null +++ b/internal/mobius/integration_test.go @@ -0,0 +1,429 @@ +package mobius + +// This file provides an in-process, protocol-level end-to-end harness: it builds a fully wired +// Hotline server (real YAML managers on a t.TempDir copy of test/config) on ephemeral ports and +// drives it with the in-repo hotline.Client. The individual regression tests live in e2e_test.go. + +import ( + "context" + "fmt" + "io" + "net" + "os" + "path" + "path/filepath" + "sync" + "testing" + "time" + + "github.com/jhalter/mobius/hotline" + "github.com/stretchr/testify/require" + "golang.org/x/time/rate" +) + +// e2eServer is a running in-process server plus the metadata tests need to connect to it. +type e2eServer struct { + srv *hotline.Server + addr string // session port, host:port + xferAddr string // file transfer port (session port + 1), host:port + cfgDir string // tempdir holding the copied config + Files tree +} + +// startE2EServer builds and starts a fully-wired server on a free ephemeral port pair, returning +// once the session port accepts connections. The server is shut down via t.Cleanup. +// +// Binding is inherently racy: between probing for a free port pair and ListenAndServe claiming it, +// another process (a concurrently-tested package, or anything else on the machine) can steal a +// port. A stolen port surfaces as an early ListenAndServe error, so the serve attempt is retried +// on a fresh pair rather than failing the test. +func startE2EServer(t *testing.T) *e2eServer { + t.Helper() + + cfgDir := t.TempDir() + copyDirRecursiveTest(t, "test/config", cfgDir) + + cfg, err := LoadConfig(filepath.Join(cfgDir, "config.yaml")) + require.NoError(t, err) + // The fixture's FileRoot is a placeholder; point it at the copied Files tree. + cfg.FileRoot = filepath.Join(cfgDir, "Files") + + // Wire the concrete managers exactly as cmd/mobius-hotline-server/main.go does. + messageBoard, err := NewFlatNews(path.Join(cfgDir, "MessageBoard.txt")) + require.NoError(t, err) + + banFile, err := NewBanFile(path.Join(cfgDir, "Banlist.yaml")) + require.NoError(t, err) + + threadedNews, err := NewThreadedNewsYAML(path.Join(cfgDir, "ThreadedNews.yaml")) + require.NoError(t, err) + + am, err := NewYAMLAccountManager(path.Join(cfgDir, "Users/")) + require.NoError(t, err) + + agreement, err := NewAgreement(cfgDir, "\r") + require.NoError(t, err) + + // Create an admin account with a known password so tests that need elevated rights can log in + // deterministically. The client obfuscates passwords with EncodeString and the server bcrypt- + // compares that obfuscated form directly (see ClientConn.Authenticate), so the stored account + // must hash the obfuscated password — exactly what HandleNewUser does for real accounts. + var adminAccess hotline.AccessBitmap + for i := 0; i <= hotline.AccessSendPrivMsg; i++ { + adminAccess.Set(i) + } + obfuscatedPass := string(hotline.EncodeString([]byte(e2eAdminPass))) + require.NoError(t, am.Create(*hotline.NewAccount(e2eAdminLogin, "e2e admin", obfuscatedPass, adminAccess))) + + for attempt := 0; attempt < 5; attempt++ { + port := findFreePortPairTest(t) + + srv, err := hotline.NewServer( + hotline.WithInterface("127.0.0.1"), + hotline.WithPort(port), + hotline.WithConfig(*cfg), + hotline.WithLogger(NewTestLogger()), + // Disable per-IP connection throttling so multi-client tests don't pay 2s per connection. + hotline.WithConnectionRateLimit(rate.Inf, 1), + ) + require.NoError(t, err) + + srv.MessageBoard = messageBoard + srv.BanList = banFile + srv.ThreadedNewsMgr = threadedNews + srv.AccountManager = am + srv.Agreement = agreement + RegisterHandlers(srv) + + ctx, cancel := context.WithCancel(context.Background()) + serveErr := make(chan error, 1) + go func() { serveErr <- srv.ListenAndServe(ctx) }() + + addr := fmt.Sprintf("127.0.0.1:%d", port) + if !waitForListener(t, addr, serveErr) { + cancel() + continue + } + + t.Cleanup(func() { + cancel() + select { + case <-serveErr: + case <-time.After(10 * time.Second): + t.Error("server did not shut down within 10s") + } + }) + + return &e2eServer{ + srv: srv, + addr: addr, + xferAddr: fmt.Sprintf("127.0.0.1:%d", port+1), + cfgDir: cfgDir, + } + } + + t.Fatal("could not start e2e server: probed port pairs kept being claimed by other processes") + return nil +} + +const ( + e2eAdminLogin = "e2e-admin" + e2eAdminPass = "e2e-admin-pw" +) + +// waitForListener retry-dials until the server accepts a TCP connection, returning true on +// success. It returns false if ListenAndServe returned early instead — a bind failure, meaning +// another process claimed a probed port first — which the caller treats as a retry signal. A +// dial can even "succeed" against that other process's listener, so a successful dial only counts +// after the bind error has had a moment to surface. A server that neither accepts nor errors +// within the deadline fails the test. +func waitForListener(t *testing.T, addr string, serveErr <-chan error) bool { + t.Helper() + deadline := time.Now().Add(5 * time.Second) + for time.Now().Before(deadline) { + select { + case <-serveErr: + return false + default: + } + conn, err := net.DialTimeout("tcp", addr, 200*time.Millisecond) + if err == nil { + _ = conn.Close() + select { + case <-serveErr: + return false + case <-time.After(20 * time.Millisecond): + return true + } + } + time.Sleep(10 * time.Millisecond) + } + t.Fatalf("server at %s never became ready", addr) + return false +} + +// findFreePortPairTest finds a port p such that both p and p+1 are free on the loopback interface. +// Duplicated from hotline/server_test.go's findFreePortPair (that copy is in-package and unexported). +func findFreePortPairTest(t *testing.T) int { + t.Helper() + for attempt := 0; attempt < 50; attempt++ { + l0, err := net.Listen("tcp", "127.0.0.1:0") + require.NoError(t, err) + p := l0.Addr().(*net.TCPAddr).Port + l1, err := net.Listen("tcp", fmt.Sprintf("127.0.0.1:%d", p+1)) + _ = l0.Close() + if err != nil { + continue + } + _ = l1.Close() + return p + } + t.Fatal("could not find a free consecutive port pair") + return 0 +} + +// copyDirRecursiveTest copies a directory tree from src into dst, creating dst subdirectories as +// needed. It mirrors copyDirRecursive in cmd/mobius-hotline-server (not importable from here). +func copyDirRecursiveTest(t *testing.T, src, dst string) { + t.Helper() + entries, err := os.ReadDir(src) + require.NoError(t, err) + require.NoError(t, os.MkdirAll(dst, 0755)) + + for _, entry := range entries { + srcPath := filepath.Join(src, entry.Name()) + dstPath := filepath.Join(dst, entry.Name()) + if entry.IsDir() { + copyDirRecursiveTest(t, srcPath, dstPath) + continue + } + data, err := os.ReadFile(srcPath) + require.NoError(t, err) + require.NoError(t, os.WriteFile(dstPath, data, 0644)) + } +} + +// --- E2E client driver ----------------------------------------------------------------------- + +// e2eClient wraps hotline.Client with request/reply correlation and collection of unsolicited +// server-pushed transactions, so tests can drive the protocol synchronously. +type e2eClient struct { + t *testing.T + client *hotline.Client + cancel context.CancelFunc + + mu sync.Mutex + waiters map[[4]byte]chan hotline.Transaction // reply ID -> waiter + incoming []hotline.Transaction // unsolicited (IsReply==0) transactions + incCh chan hotline.Transaction // signals a new unsolicited transaction +} + +// connectE2E performs handshake + login and starts the read loop. It auto-replies to the agreement +// prompt so the session becomes fully established. The login doubles as the display name. +func connectE2E(t *testing.T, addr, login, pass string) *e2eClient { + t.Helper() + return connectE2ENamed(t, addr, login, login, pass) +} + +// connectE2ENamed is connectE2E with a display name distinct from the login, so tests that run +// several sessions on the same account can tell them apart in user lists and notifications. +func connectE2ENamed(t *testing.T, addr, name, login, pass string) *e2eClient { + t.Helper() + + c := hotline.NewClient(name, NewTestLogger()) + ec := &e2eClient{ + t: t, + client: c, + waiters: map[[4]byte]chan hotline.Transaction{}, + incCh: make(chan hotline.Transaction, 64), + } + + // Register the same dispatcher for every transaction type the tests care about. Replies are + // re-typed to their request type by the client before dispatch, so one handler covers both + // request-reply and server-push flows. + for _, tt2 := range dispatchTypes { + c.HandleFunc(tt2, ec.dispatch) + } + + require.NoError(t, c.Connect(addr, login, pass)) + + ctx, cancel := context.WithCancel(context.Background()) + ec.cancel = cancel + go func() { _ = c.HandleTransactions(ctx) }() + + t.Cleanup(func() { + cancel() + _ = c.Disconnect() + }) + + return ec +} + +// dispatchTypes is the set of transaction types the e2e dispatcher is registered for: every request +// type the tests send (replies arrive re-typed to these) plus server-pushed types. +var dispatchTypes = [][2]byte{ + hotline.TranLogin, hotline.TranAgreed, hotline.TranShowAgreement, + hotline.TranChatSend, hotline.TranChatMsg, + hotline.TranGetMsgs, hotline.TranOldPostNews, + hotline.TranGetNewsCatNameList, hotline.TranGetNewsArtNameList, + hotline.TranGetNewsArtData, hotline.TranPostNewsArt, + hotline.TranGetFileNameList, hotline.TranDownloadFile, hotline.TranUploadFile, + hotline.TranGetUser, hotline.TranNewUser, hotline.TranSetUser, hotline.TranDeleteUser, + hotline.TranListUsers, hotline.TranGetUserNameList, + hotline.TranNotifyChangeUser, hotline.TranNotifyDeleteUser, + hotline.TranServerMsg, hotline.TranDisconnectMsg, hotline.TranUserAccess, + hotline.TranInviteNewChat, hotline.TranInviteToChat, hotline.TranJoinChat, + hotline.TranNotifyChatChangeUser, hotline.TranNotifyChatDeleteUser, hotline.TranNotifyChatSubject, +} + +func (ec *e2eClient) dispatch(_ context.Context, _ *hotline.Client, t *hotline.Transaction) ([]hotline.Transaction, error) { + if t.IsReply == 1 { + ec.mu.Lock() + ch := ec.waiters[t.ID] + delete(ec.waiters, t.ID) + ec.mu.Unlock() + if ch != nil { + ch <- *t + } + return nil, nil + } + + // Unsolicited server push. + if t.Type == hotline.TranShowAgreement { + // Accept the agreement so login completes. + return []hotline.Transaction{ + hotline.NewTransaction(hotline.TranAgreed, [2]byte{}, + hotline.NewField(hotline.FieldUserName, []byte(ec.client.Pref.Username)), + hotline.NewField(hotline.FieldUserIconID, ec.client.Pref.IconBytes()), + hotline.NewField(hotline.FieldOptions, []byte{0, 0}), + ), + }, nil + } + + ec.mu.Lock() + ec.incoming = append(ec.incoming, *t) + ec.mu.Unlock() + select { + case ec.incCh <- *t: + default: + } + return nil, nil +} + +// roundTrip sends a request and waits for the matching reply (correlated by transaction ID). +func (ec *e2eClient) roundTrip(req hotline.Transaction) hotline.Transaction { + ec.t.Helper() + + ch := make(chan hotline.Transaction, 1) + ec.mu.Lock() + ec.waiters[req.ID] = ch + ec.mu.Unlock() + + require.NoError(ec.t, ec.client.Send(req)) + + select { + case reply := <-ch: + return reply + case <-time.After(15 * time.Second): + ec.t.Fatalf("timed out waiting for reply to %v", req.Type) + return hotline.Transaction{} + } +} + +// waitFor blocks until an unsolicited transaction of the given type arrives, or the test fails. +// It polls the buffer on a short interval as well as waking on incCh, so a signal that races with +// the buffer scan can never cause a lost wakeup. +func (ec *e2eClient) waitFor(tranType [2]byte) hotline.Transaction { + ec.t.Helper() + + deadline := time.After(15 * time.Second) + ticker := time.NewTicker(20 * time.Millisecond) + defer ticker.Stop() + + for { + // Scan anything already buffered. + ec.mu.Lock() + for i, tr := range ec.incoming { + if tr.Type == tranType { + ec.incoming = append(ec.incoming[:i], ec.incoming[i+1:]...) + ec.mu.Unlock() + return tr + } + } + ec.mu.Unlock() + + select { + case <-ec.incCh: + // Woke on a new arrival; loop and re-scan. + case <-ticker.C: + // Periodic re-scan guards against a dropped incCh signal. + case <-deadline: + ec.t.Fatalf("timed out waiting for unsolicited %v", tranType) + return hotline.Transaction{} + } + } +} + +// isError reports whether a reply carries a Hotline error (non-zero ErrorCode). +func isError(t hotline.Transaction) bool { + return t.ErrorCode != [4]byte{0, 0, 0, 0} +} + +// --- transfer-port helpers ------------------------------------------------------------------- + +// htxfHeader builds the 16-byte HTXF transfer header (protocol + reference number + data size). +func htxfHeader(refNum [4]byte, dataSize uint32) []byte { + b := make([]byte, 16) + copy(b[0:4], hotline.HTXF[:]) + copy(b[4:8], refNum[:]) + b[8] = byte(dataSize >> 24) + b[9] = byte(dataSize >> 16) + b[10] = byte(dataSize >> 8) + b[11] = byte(dataSize) + return b +} + +// downloadOverTransferPort dials the transfer port, sends the HTXF header for refNum, and returns +// all bytes the server streams back (the flattened file object). +func downloadOverTransferPort(t *testing.T, xferAddr string, refNum [4]byte) []byte { + t.Helper() + conn, err := net.DialTimeout("tcp", xferAddr, 5*time.Second) + require.NoError(t, err) + defer func() { _ = conn.Close() }() + + _, err = conn.Write(htxfHeader(refNum, 0)) + require.NoError(t, err) + + require.NoError(t, conn.SetReadDeadline(time.Now().Add(5*time.Second))) + data, err := io.ReadAll(conn) + // A read deadline or server-initiated close both surface here; the payload collected so far is + // what we assert on. + if err != nil && !isTimeout(err) { + require.ErrorIs(t, err, io.EOF) + } + return data +} + +// uploadOverTransferPort dials the transfer port, sends the HTXF header, then streams payload. +func uploadOverTransferPort(t *testing.T, xferAddr string, refNum [4]byte, payload []byte) { + t.Helper() + conn, err := net.DialTimeout("tcp", xferAddr, 5*time.Second) + require.NoError(t, err) + defer func() { _ = conn.Close() }() + + _, err = conn.Write(htxfHeader(refNum, uint32(len(payload)))) + require.NoError(t, err) + _, err = conn.Write(payload) + require.NoError(t, err) + + // Give the server a moment to persist before the caller asserts. + require.NoError(t, conn.SetReadDeadline(time.Now().Add(2*time.Second))) + _, _ = io.ReadAll(conn) +} + +func isTimeout(err error) bool { + var ne net.Error + if e, ok := err.(net.Error); ok { + ne = e + } + return ne != nil && ne.Timeout() +} diff --git a/internal/mobius/reload_test.go b/internal/mobius/reload_test.go new file mode 100644 index 0000000..d864fa7 --- /dev/null +++ b/internal/mobius/reload_test.go @@ -0,0 +1,32 @@ +package mobius + +import ( + "errors" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// The concrete Reload implementations (FlatNews, Agreement, BanFile, ThreadedNewsYAML) are +// exercised in their own *_test.go files. This covers the ReloaderFunc adapter, which is the +// only unit defined in reload.go without direct coverage. +func TestReloaderFunc_Reload(t *testing.T) { + t.Run("invokes the wrapped func", func(t *testing.T) { + called := false + var r Reloader = ReloaderFunc(func() error { + called = true + return nil + }) + + require.NoError(t, r.Reload()) + assert.True(t, called) + }) + + t.Run("propagates the wrapped error", func(t *testing.T) { + sentinel := errors.New("reload failed") + var r Reloader = ReloaderFunc(func() error { return sentinel }) + + assert.ErrorIs(t, r.Reload(), sentinel) + }) +} |