aboutsummaryrefslogtreecommitdiff
path: root/internal
diff options
context:
space:
mode:
authorJeff Halter <868228+jhalter@users.noreply.github.com>2026-07-10 09:48:18 -0700
committerJeff Halter <868228+jhalter@users.noreply.github.com>2026-07-10 09:48:18 -0700
commit21f24d24fd6f501b32f15a2bef41c89cc461f623 (patch)
treeba4952fb82f3ef9c265228633f27536866e23931 /internal
parentae44fb222ec73cae8441f5a5d9a21f8585fe587f (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.go12
-rw-r--r--internal/mobius/agreement.go14
-rw-r--r--internal/mobius/config_test.go14
-rw-r--r--internal/mobius/e2e_test.go439
-rw-r--r--internal/mobius/integration_test.go429
-rw-r--r--internal/mobius/reload_test.go32
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)
+ })
+}