main pico / pkg / apps / pipe / ssh_test.go
Eric Bower  ·  2026-08-15
   1package pipe
   2
   3import (
   4	"bytes"
   5	"context"
   6	"crypto/ed25519"
   7	"crypto/rand"
   8	"fmt"
   9	"io"
  10	"log/slog"
  11	"os"
  12	"strings"
  13	"sync"
  14	"testing"
  15	"time"
  16
  17	"github.com/antoniomika/syncmap"
  18	"github.com/picosh/pico/pkg/db"
  19	"github.com/picosh/pico/pkg/db/stub"
  20	"github.com/picosh/pico/pkg/pssh"
  21	psub "github.com/picosh/pico/pkg/pubsub"
  22	"github.com/picosh/pico/pkg/shared"
  23	"github.com/prometheus/client_golang/prometheus"
  24	"golang.org/x/crypto/ssh"
  25)
  26
  27type TestDB struct {
  28	*stub.StubDB
  29	mu           sync.RWMutex
  30	Users        []*db.User
  31	Pubkeys      []*db.PublicKey
  32	Features     []*db.FeatureFlag
  33	PipeMonitors []*db.PipeMonitor
  34}
  35
  36func NewTestDB(logger *slog.Logger) *TestDB {
  37	return &TestDB{
  38		StubDB: stub.NewStubDB(logger),
  39	}
  40}
  41
  42func (t *TestDB) FindUserByPubkey(key string) (*db.User, error) {
  43	t.mu.RLock()
  44	defer t.mu.RUnlock()
  45	for _, pk := range t.Pubkeys {
  46		if pk.Key == key {
  47			return t.findUserLocked(pk.UserID)
  48		}
  49	}
  50	return nil, fmt.Errorf("user not found for pubkey")
  51}
  52
  53func (t *TestDB) findUserLocked(userID string) (*db.User, error) {
  54	for _, user := range t.Users {
  55		if user.ID == userID {
  56			cp := *user
  57			return &cp, nil
  58		}
  59	}
  60	return nil, fmt.Errorf("user not found")
  61}
  62
  63func (t *TestDB) FindUser(userID string) (*db.User, error) {
  64	t.mu.RLock()
  65	defer t.mu.RUnlock()
  66	return t.findUserLocked(userID)
  67}
  68
  69func (t *TestDB) FindUserByName(name string) (*db.User, error) {
  70	t.mu.RLock()
  71	defer t.mu.RUnlock()
  72	for _, user := range t.Users {
  73		if user.Name == name {
  74			cp := *user
  75			return &cp, nil
  76		}
  77	}
  78	return nil, fmt.Errorf("user not found")
  79}
  80
  81func (t *TestDB) FindFeature(userID, name string) (*db.FeatureFlag, error) {
  82	t.mu.RLock()
  83	defer t.mu.RUnlock()
  84	for _, ff := range t.Features {
  85		if ff.UserID == userID && ff.Name == name {
  86			cp := *ff
  87			return &cp, nil
  88		}
  89	}
  90	return nil, fmt.Errorf("feature not found")
  91}
  92
  93func (t *TestDB) HasFeatureByUser(userID string, feature string) bool {
  94	ff, err := t.FindFeature(userID, feature)
  95	if err != nil {
  96		return false
  97	}
  98	return ff.IsValid()
  99}
 100
 101func (t *TestDB) InsertAccessLog(_ *db.AccessLog) error {
 102	return nil
 103}
 104
 105func (t *TestDB) Close() error {
 106	return nil
 107}
 108
 109func (t *TestDB) AddUser(user *db.User) {
 110	t.mu.Lock()
 111	defer t.mu.Unlock()
 112	cp := *user
 113	t.Users = append(t.Users, &cp)
 114}
 115
 116func (t *TestDB) AddPubkey(pubkey *db.PublicKey) {
 117	t.mu.Lock()
 118	defer t.mu.Unlock()
 119	cp := *pubkey
 120	t.Pubkeys = append(t.Pubkeys, &cp)
 121}
 122
 123func (t *TestDB) UpsertPipeMonitor(userID, topic string, dur time.Duration, winEnd *time.Time) error {
 124	t.mu.Lock()
 125	defer t.mu.Unlock()
 126	var winEndCopy *time.Time
 127	if winEnd != nil {
 128		w := *winEnd
 129		winEndCopy = &w
 130	}
 131	for _, m := range t.PipeMonitors {
 132		if m.UserId == userID && m.Topic == topic {
 133			m.WindowDur = dur
 134			m.WindowEnd = winEndCopy
 135			now := time.Now()
 136			m.UpdatedAt = &now
 137			return nil
 138		}
 139	}
 140	now := time.Now()
 141	t.PipeMonitors = append(t.PipeMonitors, &db.PipeMonitor{
 142		ID:        fmt.Sprintf("monitor-%s-%s", userID, topic),
 143		UserId:    userID,
 144		Topic:     topic,
 145		WindowDur: dur,
 146		WindowEnd: winEndCopy,
 147		CreatedAt: &now,
 148		UpdatedAt: &now,
 149	})
 150	return nil
 151}
 152
 153func (t *TestDB) UpdatePipeMonitorLastPing(userID, topic string, lastPing *time.Time) error {
 154	t.mu.Lock()
 155	defer t.mu.Unlock()
 156	var lastPingCopy *time.Time
 157	if lastPing != nil {
 158		p := *lastPing
 159		lastPingCopy = &p
 160	}
 161	for _, m := range t.PipeMonitors {
 162		if m.UserId == userID && m.Topic == topic {
 163			m.LastPing = lastPingCopy
 164			now := time.Now()
 165			m.UpdatedAt = &now
 166			return nil
 167		}
 168	}
 169	return fmt.Errorf("monitor not found")
 170}
 171
 172func (t *TestDB) RemovePipeMonitor(userID, topic string) error {
 173	t.mu.Lock()
 174	defer t.mu.Unlock()
 175	for i, m := range t.PipeMonitors {
 176		if m.UserId == userID && m.Topic == topic {
 177			t.PipeMonitors = append(t.PipeMonitors[:i], t.PipeMonitors[i+1:]...)
 178			return nil
 179		}
 180	}
 181	return fmt.Errorf("monitor not found")
 182}
 183
 184func copyPipeMonitor(m *db.PipeMonitor) *db.PipeMonitor {
 185	if m == nil {
 186		return nil
 187	}
 188	cp := *m
 189	if m.WindowEnd != nil {
 190		w := *m.WindowEnd
 191		cp.WindowEnd = &w
 192	}
 193	if m.LastPing != nil {
 194		p := *m.LastPing
 195		cp.LastPing = &p
 196	}
 197	if m.CreatedAt != nil {
 198		c := *m.CreatedAt
 199		cp.CreatedAt = &c
 200	}
 201	if m.UpdatedAt != nil {
 202		u := *m.UpdatedAt
 203		cp.UpdatedAt = &u
 204	}
 205	return &cp
 206}
 207
 208func (t *TestDB) FindPipeMonitorByTopic(userID, topic string) (*db.PipeMonitor, error) {
 209	t.mu.RLock()
 210	defer t.mu.RUnlock()
 211	for _, m := range t.PipeMonitors {
 212		if m.UserId == userID && m.Topic == topic {
 213			return copyPipeMonitor(m), nil
 214		}
 215	}
 216	return nil, fmt.Errorf("monitor not found")
 217}
 218
 219func (t *TestDB) FindPipeMonitorsByUser(userID string) ([]*db.PipeMonitor, error) {
 220	t.mu.RLock()
 221	defer t.mu.RUnlock()
 222	var monitors []*db.PipeMonitor
 223	for _, m := range t.PipeMonitors {
 224		if m.UserId == userID {
 225			monitors = append(monitors, copyPipeMonitor(m))
 226		}
 227	}
 228	return monitors, nil
 229}
 230
 231func (t *TestDB) InsertPipeMonitorHistory(monitorID string, windowDur time.Duration, windowEnd, lastPing *time.Time) error {
 232	return nil
 233}
 234
 235func (t *TestDB) FindPipeMonitorHistory(monitorID string, from, to time.Time) ([]*db.PipeMonitorHistory, error) {
 236	return nil, nil
 237}
 238
 239type TestSSHServer struct {
 240	Cfg         *shared.ConfigSite
 241	DBPool      *TestDB
 242	PipeHandler *CliHandler
 243	Cancel      context.CancelFunc
 244}
 245
 246func NewTestSSHServer(t *testing.T) *TestSSHServer {
 247	t.Helper()
 248
 249	opts := &slog.HandlerOptions{
 250		AddSource: true,
 251		Level:     slog.LevelDebug,
 252	}
 253	logger := slog.New(slog.NewTextHandler(os.Stdout, opts))
 254
 255	dbpool := NewTestDB(logger)
 256
 257	cfg := &shared.ConfigSite{
 258		Domain:       "pipe.test",
 259		Port:         "2225",
 260		PortOverride: "2225",
 261		Protocol:     "ssh",
 262		Logger:       logger,
 263		Space:        "pipe",
 264	}
 265
 266	ctx, cancel := context.WithCancel(context.Background())
 267
 268	pubsub := psub.NewMulticast(logger)
 269	handler := &CliHandler{
 270		Logger:  logger,
 271		DBPool:  dbpool,
 272		PubSub:  pubsub,
 273		Cfg:     cfg,
 274		Waiters: syncmap.New[string, []string](),
 275		Access:  syncmap.New[string, []string](),
 276	}
 277
 278	sshAuth := shared.NewSshAuthHandler(dbpool, logger, "pipe")
 279
 280	prometheus.DefaultRegisterer = prometheus.NewRegistry()
 281
 282	server, err := pssh.NewSSHServerWithConfig(
 283		ctx,
 284		logger,
 285		"pipe-ssh-test",
 286		"localhost",
 287		cfg.Port,
 288		"9223",
 289		"../../ssh_data/term_info_ed25519",
 290		func(conn ssh.ConnMetadata, key ssh.PublicKey) (*ssh.Permissions, error) {
 291			perms, _ := sshAuth.PubkeyAuthHandler(conn, key)
 292			if perms == nil {
 293				perms = &ssh.Permissions{
 294					Extensions: map[string]string{
 295						"pubkey": shared.KeyForKeyText(key),
 296					},
 297				}
 298			}
 299			return perms, nil
 300		},
 301		[]pssh.SSHServerMiddleware{
 302			Middleware(handler),
 303			pssh.LogMiddleware(handler, dbpool),
 304		},
 305		nil,
 306		nil,
 307	)
 308
 309	if err != nil {
 310		t.Fatalf("failed to create ssh server: %v", err)
 311	}
 312
 313	go func() {
 314		if err := server.ListenAndServe(); err != nil {
 315			logger.Error("serve", "err", err.Error())
 316		}
 317	}()
 318
 319	time.Sleep(100 * time.Millisecond)
 320
 321	return &TestSSHServer{
 322		Cfg:         cfg,
 323		DBPool:      dbpool,
 324		PipeHandler: handler,
 325		Cancel:      cancel,
 326	}
 327}
 328
 329func (s *TestSSHServer) Shutdown() {
 330	s.Cancel()
 331	time.Sleep(10 * time.Millisecond)
 332}
 333
 334type UserSSH struct {
 335	username   string
 336	signer     ssh.Signer
 337	privateKey []byte
 338}
 339
 340func GenerateUser(username string) UserSSH {
 341	_, userKey, err := ed25519.GenerateKey(rand.Reader)
 342	if err != nil {
 343		panic(err)
 344	}
 345
 346	b, err := ssh.MarshalPrivateKey(userKey, "")
 347	if err != nil {
 348		panic(err)
 349	}
 350
 351	userSigner, err := ssh.NewSignerFromKey(userKey)
 352	if err != nil {
 353		panic(err)
 354	}
 355
 356	return UserSSH{
 357		username:   username,
 358		signer:     userSigner,
 359		privateKey: b.Bytes,
 360	}
 361}
 362
 363func (u UserSSH) PublicKey() string {
 364	return shared.KeyForKeyText(u.signer.PublicKey())
 365}
 366
 367func (u UserSSH) NewClient() (*ssh.Client, error) {
 368	config := &ssh.ClientConfig{
 369		User: u.username,
 370		Auth: []ssh.AuthMethod{
 371			ssh.PublicKeys(u.signer),
 372		},
 373		HostKeyCallback: ssh.InsecureIgnoreHostKey(),
 374	}
 375
 376	return ssh.Dial("tcp", "localhost:2225", config)
 377}
 378
 379func (u UserSSH) RunCommand(client *ssh.Client, cmd string) (string, error) {
 380	session, err := client.NewSession()
 381	if err != nil {
 382		return "", err
 383	}
 384	defer func() { _ = session.Close() }()
 385
 386	stdoutPipe, err := session.StdoutPipe()
 387	if err != nil {
 388		return "", err
 389	}
 390
 391	stderrPipe, err := session.StderrPipe()
 392	if err != nil {
 393		return "", err
 394	}
 395
 396	if err := session.Start(cmd); err != nil {
 397		return "", err
 398	}
 399
 400	stdout := new(strings.Builder)
 401	stderr := new(strings.Builder)
 402	_, _ = io.Copy(stdout, stdoutPipe)
 403	_, _ = io.Copy(stderr, stderrPipe)
 404
 405	_ = session.Wait()
 406	return stdout.String() + stderr.String(), nil
 407}
 408
 409func (u UserSSH) RunCommandWithStdin(client *ssh.Client, cmd string, stdin string) (string, error) {
 410	session, err := client.NewSession()
 411	if err != nil {
 412		return "", err
 413	}
 414	defer func() { _ = session.Close() }()
 415
 416	stdinPipe, err := session.StdinPipe()
 417	if err != nil {
 418		return "", err
 419	}
 420
 421	stdoutPipe, err := session.StdoutPipe()
 422	if err != nil {
 423		return "", err
 424	}
 425
 426	if err := session.Start(cmd); err != nil {
 427		return "", err
 428	}
 429
 430	_, err = stdinPipe.Write([]byte(stdin))
 431	if err != nil {
 432		return "", err
 433	}
 434	_ = stdinPipe.Close()
 435
 436	buf := new(strings.Builder)
 437	_, err = io.Copy(buf, stdoutPipe)
 438	if err != nil {
 439		return "", err
 440	}
 441
 442	_ = session.Wait()
 443	return buf.String(), nil
 444}
 445
 446func RegisterUserWithServer(server *TestSSHServer, user UserSSH) {
 447	dbUser := &db.User{
 448		ID:   user.username + "-id",
 449		Name: user.username,
 450	}
 451	server.DBPool.AddUser(dbUser)
 452	server.DBPool.AddPubkey(&db.PublicKey{
 453		ID:     user.username + "-pubkey-id",
 454		UserID: dbUser.ID,
 455		Key:    user.PublicKey(),
 456	})
 457}
 458
 459func TestLs_UnauthenticatedUserDenied(t *testing.T) {
 460	server := NewTestSSHServer(t)
 461	defer server.Shutdown()
 462
 463	user := GenerateUser("anonymous")
 464
 465	client, err := user.NewClient()
 466	if err != nil {
 467		t.Fatalf("failed to connect: %v", err)
 468	}
 469	defer func() { _ = client.Close() }()
 470
 471	output, err := user.RunCommand(client, "ls")
 472	if err != nil {
 473		t.Logf("command error (expected): %v", err)
 474	}
 475
 476	if !strings.Contains(output, "access denied") {
 477		t.Errorf("expected 'access denied', got: %s", output)
 478	}
 479}
 480
 481func TestLs_AuthenticatedUser(t *testing.T) {
 482	server := NewTestSSHServer(t)
 483	defer server.Shutdown()
 484
 485	user := GenerateUser("alice")
 486	RegisterUserWithServer(server, user)
 487
 488	client, err := user.NewClient()
 489	if err != nil {
 490		t.Fatalf("failed to connect: %v", err)
 491	}
 492	defer func() { _ = client.Close() }()
 493
 494	output, err := user.RunCommand(client, "ls")
 495	if err != nil {
 496		t.Logf("command completed with: %v", err)
 497	}
 498
 499	if strings.Contains(output, "access denied") {
 500		t.Errorf("authenticated user should not get access denied, got: %s", output)
 501	}
 502
 503	if !strings.Contains(output, "no pubsub channels found") {
 504		t.Errorf("expected 'no pubsub channels found' for empty state, got: %s", output)
 505	}
 506}
 507
 508func TestPubSub_BasicFlow(t *testing.T) {
 509	server := NewTestSSHServer(t)
 510	defer server.Shutdown()
 511
 512	user := GenerateUser("alice")
 513	RegisterUserWithServer(server, user)
 514
 515	subClient, err := user.NewClient()
 516	if err != nil {
 517		t.Fatalf("failed to connect subscriber: %v", err)
 518	}
 519	defer func() { _ = subClient.Close() }()
 520
 521	pubClient, err := user.NewClient()
 522	if err != nil {
 523		t.Fatalf("failed to connect publisher: %v", err)
 524	}
 525	defer func() { _ = pubClient.Close() }()
 526
 527	subSession, err := subClient.NewSession()
 528	if err != nil {
 529		t.Fatalf("failed to create sub session: %v", err)
 530	}
 531	defer func() { _ = subSession.Close() }()
 532
 533	subStdout, err := subSession.StdoutPipe()
 534	if err != nil {
 535		t.Fatalf("failed to get sub stdout: %v", err)
 536	}
 537
 538	if err := subSession.Start("sub testtopic -c"); err != nil {
 539		t.Fatalf("failed to start sub: %v", err)
 540	}
 541
 542	time.Sleep(100 * time.Millisecond)
 543
 544	testMessage := "hello from pub"
 545	_, err = user.RunCommandWithStdin(pubClient, "pub testtopic -c", testMessage)
 546	if err != nil {
 547		t.Logf("pub command completed: %v", err)
 548	}
 549
 550	received := make([]byte, len(testMessage)+10)
 551	n, err := subStdout.Read(received)
 552	if err != nil && err != io.EOF {
 553		t.Logf("read error: %v", err)
 554	}
 555
 556	receivedStr := string(received[:n])
 557	if !strings.Contains(receivedStr, testMessage) {
 558		t.Errorf("subscriber did not receive message, got: %q, want: %q", receivedStr, testMessage)
 559	}
 560}
 561
 562func TestPubSub_PublicTopic(t *testing.T) {
 563	server := NewTestSSHServer(t)
 564	defer server.Shutdown()
 565
 566	alice := GenerateUser("alice")
 567	bob := GenerateUser("bob")
 568	RegisterUserWithServer(server, alice)
 569	RegisterUserWithServer(server, bob)
 570
 571	subClient, err := bob.NewClient()
 572	if err != nil {
 573		t.Fatalf("failed to connect subscriber: %v", err)
 574	}
 575	defer func() { _ = subClient.Close() }()
 576
 577	pubClient, err := alice.NewClient()
 578	if err != nil {
 579		t.Fatalf("failed to connect publisher: %v", err)
 580	}
 581	defer func() { _ = pubClient.Close() }()
 582
 583	subSession, err := subClient.NewSession()
 584	if err != nil {
 585		t.Fatalf("failed to create sub session: %v", err)
 586	}
 587	defer func() { _ = subSession.Close() }()
 588
 589	subStdout, err := subSession.StdoutPipe()
 590	if err != nil {
 591		t.Fatalf("failed to get sub stdout: %v", err)
 592	}
 593
 594	if err := subSession.Start("sub publictopic -p -c"); err != nil {
 595		t.Fatalf("failed to start sub: %v", err)
 596	}
 597
 598	time.Sleep(100 * time.Millisecond)
 599
 600	testMessage := "public message"
 601	_, err = alice.RunCommandWithStdin(pubClient, "pub publictopic -p -c", testMessage)
 602	if err != nil {
 603		t.Logf("pub command completed: %v", err)
 604	}
 605
 606	received := make([]byte, len(testMessage)+10)
 607	n, err := subStdout.Read(received)
 608	if err != nil && err != io.EOF {
 609		t.Logf("read error: %v", err)
 610	}
 611
 612	receivedStr := string(received[:n])
 613	if !strings.Contains(receivedStr, testMessage) {
 614		t.Errorf("subscriber did not receive public message, got: %q, want: %q", receivedStr, testMessage)
 615	}
 616}
 617
 618func TestPipe_Bidirectional(t *testing.T) {
 619	server := NewTestSSHServer(t)
 620	defer server.Shutdown()
 621
 622	alice := GenerateUser("alice")
 623	bob := GenerateUser("bob")
 624	RegisterUserWithServer(server, alice)
 625	RegisterUserWithServer(server, bob)
 626
 627	aliceClient, err := alice.NewClient()
 628	if err != nil {
 629		t.Fatalf("failed to connect alice: %v", err)
 630	}
 631	defer func() { _ = aliceClient.Close() }()
 632
 633	bobClient, err := bob.NewClient()
 634	if err != nil {
 635		t.Fatalf("failed to connect bob: %v", err)
 636	}
 637	defer func() { _ = bobClient.Close() }()
 638
 639	aliceSession, err := aliceClient.NewSession()
 640	if err != nil {
 641		t.Fatalf("failed to create alice session: %v", err)
 642	}
 643	defer func() { _ = aliceSession.Close() }()
 644
 645	aliceStdin, err := aliceSession.StdinPipe()
 646	if err != nil {
 647		t.Fatalf("failed to get alice stdin: %v", err)
 648	}
 649
 650	aliceStdout, err := aliceSession.StdoutPipe()
 651	if err != nil {
 652		t.Fatalf("failed to get alice stdout: %v", err)
 653	}
 654
 655	if err := aliceSession.Start("pipe pipetopic -p -c"); err != nil {
 656		t.Fatalf("failed to start alice pipe: %v", err)
 657	}
 658
 659	time.Sleep(100 * time.Millisecond)
 660
 661	bobSession, err := bobClient.NewSession()
 662	if err != nil {
 663		t.Fatalf("failed to create bob session: %v", err)
 664	}
 665	defer func() { _ = bobSession.Close() }()
 666
 667	bobStdin, err := bobSession.StdinPipe()
 668	if err != nil {
 669		t.Fatalf("failed to get bob stdin: %v", err)
 670	}
 671
 672	bobStdout, err := bobSession.StdoutPipe()
 673	if err != nil {
 674		t.Fatalf("failed to get bob stdout: %v", err)
 675	}
 676
 677	if err := bobSession.Start("pipe pipetopic -p -c"); err != nil {
 678		t.Fatalf("failed to start bob pipe: %v", err)
 679	}
 680
 681	time.Sleep(100 * time.Millisecond)
 682
 683	aliceMsg := "hello from alice\n"
 684	_, err = aliceStdin.Write([]byte(aliceMsg))
 685	if err != nil {
 686		t.Fatalf("alice failed to write: %v", err)
 687	}
 688
 689	bobReceived := make([]byte, 100)
 690	n, err := bobStdout.Read(bobReceived)
 691	if err != nil && err != io.EOF {
 692		t.Logf("bob read error: %v", err)
 693	}
 694	if !strings.Contains(string(bobReceived[:n]), "hello from alice") {
 695		t.Errorf("bob did not receive alice's message, got: %q", string(bobReceived[:n]))
 696	}
 697
 698	bobMsg := "hello from bob\n"
 699	_, err = bobStdin.Write([]byte(bobMsg))
 700	if err != nil {
 701		t.Fatalf("bob failed to write: %v", err)
 702	}
 703
 704	aliceReceived := make([]byte, 100)
 705	n, err = aliceStdout.Read(aliceReceived)
 706	if err != nil && err != io.EOF {
 707		t.Logf("alice read error: %v", err)
 708	}
 709	if !strings.Contains(string(aliceReceived[:n]), "hello from bob") {
 710		t.Errorf("alice did not receive bob's message, got: %q", string(aliceReceived[:n]))
 711	}
 712
 713	// When alice disconnects, bob's session should terminate cleanly without hanging
 714	_ = aliceStdin.Close()
 715	_ = aliceSession.Close()
 716	_ = bobStdin.Close()
 717
 718	bobDone := make(chan error, 1)
 719	go func() {
 720		bobDone <- bobSession.Wait()
 721	}()
 722
 723	select {
 724	case <-bobDone:
 725		// Bob's session terminated cleanly
 726	case <-time.After(3 * time.Second):
 727		t.Fatal("bob's pipe session hung after alice disconnected")
 728	}
 729}
 730
 731func TestPipe_AutoGeneratedTopic(t *testing.T) {
 732	server := NewTestSSHServer(t)
 733	defer server.Shutdown()
 734
 735	user := GenerateUser("alice")
 736	RegisterUserWithServer(server, user)
 737
 738	client, err := user.NewClient()
 739	if err != nil {
 740		t.Fatalf("failed to connect: %v", err)
 741	}
 742	defer func() { _ = client.Close() }()
 743
 744	session, err := client.NewSession()
 745	if err != nil {
 746		t.Fatalf("failed to create session: %v", err)
 747	}
 748	defer func() { _ = session.Close() }()
 749
 750	stdout, err := session.StdoutPipe()
 751	if err != nil {
 752		t.Fatalf("failed to get stdout: %v", err)
 753	}
 754
 755	if err := session.Start("pipe"); err != nil {
 756		t.Fatalf("failed to start pipe: %v", err)
 757	}
 758
 759	received := make([]byte, 200)
 760	n, err := stdout.Read(received)
 761	if err != nil && err != io.EOF {
 762		t.Logf("read error: %v", err)
 763	}
 764
 765	output := string(received[:n])
 766	if !strings.Contains(output, "subscribe to this topic") {
 767		t.Errorf("expected topic subscription instructions, got: %q", output)
 768	}
 769}
 770
 771func TestAccessControl_AllowedUserViaFullPath(t *testing.T) {
 772	server := NewTestSSHServer(t)
 773	defer server.Shutdown()
 774
 775	alice := GenerateUser("alice")
 776	bob := GenerateUser("bob")
 777	RegisterUserWithServer(server, alice)
 778	RegisterUserWithServer(server, bob)
 779
 780	aliceClient, err := alice.NewClient()
 781	if err != nil {
 782		t.Fatalf("failed to connect alice: %v", err)
 783	}
 784	defer func() { _ = aliceClient.Close() }()
 785
 786	aliceSession, err := aliceClient.NewSession()
 787	if err != nil {
 788		t.Fatalf("failed to create alice session: %v", err)
 789	}
 790	defer func() { _ = aliceSession.Close() }()
 791
 792	aliceStdout, err := aliceSession.StdoutPipe()
 793	if err != nil {
 794		t.Fatalf("failed to get alice stdout: %v", err)
 795	}
 796
 797	if err := aliceSession.Start("sub sharedtopic -a alice,bob -c"); err != nil {
 798		t.Fatalf("failed to start alice sub: %v", err)
 799	}
 800
 801	time.Sleep(100 * time.Millisecond)
 802
 803	bobClient, err := bob.NewClient()
 804	if err != nil {
 805		t.Fatalf("failed to connect bob: %v", err)
 806	}
 807	defer func() { _ = bobClient.Close() }()
 808
 809	_, err = bob.RunCommandWithStdin(bobClient, "pub alice/sharedtopic -c", "bob allowed")
 810	if err != nil {
 811		t.Logf("bob pub completed: %v", err)
 812	}
 813
 814	aliceReceived := make([]byte, 100)
 815	n, _ := aliceStdout.Read(aliceReceived)
 816
 817	if !strings.Contains(string(aliceReceived[:n]), "bob allowed") {
 818		t.Errorf("alice should receive bob's message on shared topic, got: %q", string(aliceReceived[:n]))
 819	}
 820}
 821
 822func TestPubSub_BlockingWaitsForSubscriber(t *testing.T) {
 823	server := NewTestSSHServer(t)
 824	defer server.Shutdown()
 825
 826	user := GenerateUser("alice")
 827	RegisterUserWithServer(server, user)
 828
 829	pubClient, err := user.NewClient()
 830	if err != nil {
 831		t.Fatalf("failed to connect publisher: %v", err)
 832	}
 833	defer func() { _ = pubClient.Close() }()
 834
 835	subClient, err := user.NewClient()
 836	if err != nil {
 837		t.Fatalf("failed to connect subscriber: %v", err)
 838	}
 839	defer func() { _ = subClient.Close() }()
 840
 841	pubSession, err := pubClient.NewSession()
 842	if err != nil {
 843		t.Fatalf("failed to create pub session: %v", err)
 844	}
 845	defer func() { _ = pubSession.Close() }()
 846
 847	pubStdin, err := pubSession.StdinPipe()
 848	if err != nil {
 849		t.Fatalf("failed to get pub stdin: %v", err)
 850	}
 851
 852	pubStdout, err := pubSession.StdoutPipe()
 853	if err != nil {
 854		t.Fatalf("failed to get pub stdout: %v", err)
 855	}
 856
 857	// Start publisher with blocking enabled (default -b=true)
 858	// Publisher should wait for subscriber
 859	if err := pubSession.Start("pub blockingtopic"); err != nil {
 860		t.Fatalf("failed to start pub: %v", err)
 861	}
 862
 863	// Read output until we see "waiting" message or timeout
 864	// Need to read in a loop because Read() may return partial data
 865	var output string
 866	readDone := make(chan struct{})
 867	go func() {
 868		buf := make([]byte, 1024)
 869		for {
 870			n, err := pubStdout.Read(buf)
 871			if n > 0 {
 872				output += string(buf[:n])
 873				if strings.Contains(output, "waiting") {
 874					close(readDone)
 875					return
 876				}
 877			}
 878			if err != nil {
 879				close(readDone)
 880				return
 881			}
 882		}
 883	}()
 884
 885	select {
 886	case <-readDone:
 887	case <-time.After(2 * time.Second):
 888		t.Fatalf("timeout waiting for 'waiting' message, got: %q", output)
 889	}
 890
 891	if !strings.Contains(output, "waiting") {
 892		t.Errorf("expected 'waiting' message for blocking pub, got: %q", output)
 893	}
 894
 895	// Now start subscriber - this should unblock the publisher
 896	subSession, err := subClient.NewSession()
 897	if err != nil {
 898		t.Fatalf("failed to create sub session: %v", err)
 899	}
 900	defer func() { _ = subSession.Close() }()
 901
 902	subStdout, err := subSession.StdoutPipe()
 903	if err != nil {
 904		t.Fatalf("failed to get sub stdout: %v", err)
 905	}
 906
 907	if err := subSession.Start("sub blockingtopic -c"); err != nil {
 908		t.Fatalf("failed to start sub: %v", err)
 909	}
 910
 911	time.Sleep(100 * time.Millisecond)
 912
 913	// Now send the message
 914	testMessage := "blocking message"
 915	_, err = pubStdin.Write([]byte(testMessage))
 916	if err != nil {
 917		t.Fatalf("failed to write message: %v", err)
 918	}
 919	_ = pubStdin.Close()
 920
 921	// Subscriber should receive the message
 922	received := make([]byte, 100)
 923	nRead, err := subStdout.Read(received)
 924	if err != nil && err != io.EOF {
 925		t.Logf("read error: %v", err)
 926	}
 927
 928	if !strings.Contains(string(received[:nRead]), testMessage) {
 929		t.Errorf("subscriber did not receive blocking message, got: %q, want: %q", string(received[:nRead]), testMessage)
 930	}
 931}
 932
 933func TestPubSub_NonBlockingDoesNotWait(t *testing.T) {
 934	server := NewTestSSHServer(t)
 935	defer server.Shutdown()
 936
 937	user := GenerateUser("alice")
 938	RegisterUserWithServer(server, user)
 939
 940	client, err := user.NewClient()
 941	if err != nil {
 942		t.Fatalf("failed to connect: %v", err)
 943	}
 944	defer func() { _ = client.Close() }()
 945
 946	// Publish with -b=false (non-blocking) and no subscriber
 947	// Should complete immediately without waiting
 948	done := make(chan struct{})
 949	var output string
 950	var cmdErr error
 951
 952	go func() {
 953		output, cmdErr = user.RunCommandWithStdin(client, "pub nonblockingtopic -b=false -c", "non-blocking message")
 954		close(done)
 955	}()
 956
 957	select {
 958	case <-done:
 959		// Command completed - this is expected for non-blocking
 960		if cmdErr != nil {
 961			t.Logf("non-blocking pub completed with: %v", cmdErr)
 962		}
 963		t.Logf("non-blocking pub output: %q", output)
 964	case <-time.After(2 * time.Second):
 965		t.Errorf("non-blocking pub should complete immediately, but it blocked")
 966	}
 967}
 968
 969func TestPubSub_BlockingTimeout(t *testing.T) {
 970	server := NewTestSSHServer(t)
 971	defer server.Shutdown()
 972
 973	user := GenerateUser("alice")
 974	RegisterUserWithServer(server, user)
 975
 976	client, err := user.NewClient()
 977	if err != nil {
 978		t.Fatalf("failed to connect: %v", err)
 979	}
 980	defer func() { _ = client.Close() }()
 981
 982	// Publish with blocking and short timeout, no subscriber
 983	// Should timeout after the specified duration
 984	done := make(chan struct{})
 985	var output string
 986
 987	go func() {
 988		output, _ = user.RunCommandWithStdin(client, "pub timeouttopic -b=true -t=500ms", "timeout message")
 989		close(done)
 990	}()
 991
 992	select {
 993	case <-done:
 994		// Command completed due to timeout
 995		if !strings.Contains(output, "timeout") && !strings.Contains(output, "waiting") {
 996			t.Logf("blocking pub with timeout output: %q", output)
 997		}
 998	case <-time.After(3 * time.Second):
 999		t.Errorf("blocking pub with timeout should have timed out after 500ms")
1000	}
1001}
1002
1003func TestSub_WaitsForPublisher(t *testing.T) {
1004	server := NewTestSSHServer(t)
1005	defer server.Shutdown()
1006
1007	user := GenerateUser("alice")
1008	RegisterUserWithServer(server, user)
1009
1010	subClient, err := user.NewClient()
1011	if err != nil {
1012		t.Fatalf("failed to connect subscriber: %v", err)
1013	}
1014	defer func() { _ = subClient.Close() }()
1015
1016	pubClient, err := user.NewClient()
1017	if err != nil {
1018		t.Fatalf("failed to connect publisher: %v", err)
1019	}
1020	defer func() { _ = pubClient.Close() }()
1021
1022	// Start subscriber first - it should wait for publisher
1023	subSession, err := subClient.NewSession()
1024	if err != nil {
1025		t.Fatalf("failed to create sub session: %v", err)
1026	}
1027	defer func() { _ = subSession.Close() }()
1028
1029	subStdout, err := subSession.StdoutPipe()
1030	if err != nil {
1031		t.Fatalf("failed to get sub stdout: %v", err)
1032	}
1033
1034	if err := subSession.Start("sub waitfortopic -c"); err != nil {
1035		t.Fatalf("failed to start sub: %v", err)
1036	}
1037
1038	// Subscriber is now waiting - give it a moment
1039	time.Sleep(100 * time.Millisecond)
1040
1041	// Now publish - subscriber should receive it
1042	testMessage := "delayed publish"
1043	_, err = user.RunCommandWithStdin(pubClient, "pub waitfortopic -c", testMessage)
1044	if err != nil {
1045		t.Logf("pub completed: %v", err)
1046	}
1047
1048	received := make([]byte, 100)
1049	n, err := subStdout.Read(received)
1050	if err != nil && err != io.EOF {
1051		t.Logf("read error: %v", err)
1052	}
1053
1054	if !strings.Contains(string(received[:n]), testMessage) {
1055		t.Errorf("subscriber waiting for publisher did not receive message, got: %q, want: %q", string(received[:n]), testMessage)
1056	}
1057}
1058
1059func TestSub_KeepAliveReceivesMultipleMessages(t *testing.T) {
1060	server := NewTestSSHServer(t)
1061	defer server.Shutdown()
1062
1063	user := GenerateUser("alice")
1064	RegisterUserWithServer(server, user)
1065
1066	subClient, err := user.NewClient()
1067	if err != nil {
1068		t.Fatalf("failed to connect subscriber: %v", err)
1069	}
1070	defer func() { _ = subClient.Close() }()
1071
1072	pubClient1, err := user.NewClient()
1073	if err != nil {
1074		t.Fatalf("failed to connect publisher 1: %v", err)
1075	}
1076	defer func() { _ = pubClient1.Close() }()
1077
1078	pubClient2, err := user.NewClient()
1079	if err != nil {
1080		t.Fatalf("failed to connect publisher 2: %v", err)
1081	}
1082	defer func() { _ = pubClient2.Close() }()
1083
1084	// Start subscriber with keepAlive (-k) flag
1085	subSession, err := subClient.NewSession()
1086	if err != nil {
1087		t.Fatalf("failed to create sub session: %v", err)
1088	}
1089	defer func() { _ = subSession.Close() }()
1090
1091	subStdout, err := subSession.StdoutPipe()
1092	if err != nil {
1093		t.Fatalf("failed to get sub stdout: %v", err)
1094	}
1095
1096	if err := subSession.Start("sub keepalivetopic -k -c"); err != nil {
1097		t.Fatalf("failed to start sub: %v", err)
1098	}
1099
1100	time.Sleep(100 * time.Millisecond)
1101
1102	// Send first message
1103	msg1 := "first message\n"
1104	_, err = user.RunCommandWithStdin(pubClient1, "pub keepalivetopic -c", msg1)
1105	if err != nil {
1106		t.Logf("pub 1 completed: %v", err)
1107	}
1108
1109	received1 := make([]byte, 100)
1110	n1, _ := subStdout.Read(received1)
1111	if !strings.Contains(string(received1[:n1]), "first message") {
1112		t.Errorf("subscriber did not receive first message, got: %q", string(received1[:n1]))
1113	}
1114
1115	// Send second message - subscriber with keepAlive should still receive it
1116	msg2 := "second message\n"
1117	_, err = user.RunCommandWithStdin(pubClient2, "pub keepalivetopic -c", msg2)
1118	if err != nil {
1119		t.Logf("pub 2 completed: %v", err)
1120	}
1121
1122	received2 := make([]byte, 100)
1123	n2, _ := subStdout.Read(received2)
1124	if !strings.Contains(string(received2[:n2]), "second message") {
1125		t.Errorf("subscriber with keepAlive did not receive second message, got: %q", string(received2[:n2]))
1126	}
1127}
1128
1129func TestSub_WithoutKeepAliveExitsAfterPublisher(t *testing.T) {
1130	server := NewTestSSHServer(t)
1131	defer server.Shutdown()
1132
1133	user := GenerateUser("alice")
1134	RegisterUserWithServer(server, user)
1135
1136	subClient, err := user.NewClient()
1137	if err != nil {
1138		t.Fatalf("failed to connect subscriber: %v", err)
1139	}
1140	defer func() { _ = subClient.Close() }()
1141
1142	pubClient, err := user.NewClient()
1143	if err != nil {
1144		t.Fatalf("failed to connect publisher: %v", err)
1145	}
1146	defer func() { _ = pubClient.Close() }()
1147
1148	// Start subscriber without keepAlive
1149	subSession, err := subClient.NewSession()
1150	if err != nil {
1151		t.Fatalf("failed to create sub session: %v", err)
1152	}
1153
1154	subStdout, err := subSession.StdoutPipe()
1155	if err != nil {
1156		t.Fatalf("failed to get sub stdout: %v", err)
1157	}
1158
1159	if err := subSession.Start("sub exitaftertopic -c"); err != nil {
1160		t.Fatalf("failed to start sub: %v", err)
1161	}
1162
1163	time.Sleep(100 * time.Millisecond)
1164
1165	// Publish a message
1166	testMessage := "single message"
1167	_, err = user.RunCommandWithStdin(pubClient, "pub exitaftertopic -c", testMessage)
1168	if err != nil {
1169		t.Logf("pub completed: %v", err)
1170	}
1171
1172	// Read the message
1173	received := make([]byte, 100)
1174	n, _ := subStdout.Read(received)
1175	if !strings.Contains(string(received[:n]), testMessage) {
1176		t.Errorf("subscriber did not receive message, got: %q", string(received[:n]))
1177	}
1178
1179	// Subscriber session should exit after publisher disconnects
1180	done := make(chan error)
1181	go func() {
1182		done <- subSession.Wait()
1183	}()
1184
1185	select {
1186	case err := <-done:
1187		// Session ended as expected
1188		t.Logf("subscriber session ended: %v", err)
1189	case <-time.After(2 * time.Second):
1190		t.Errorf("subscriber without keepAlive should have exited after publisher disconnected")
1191		_ = subSession.Close()
1192	}
1193}
1194
1195func TestPub_EmptyMessage(t *testing.T) {
1196	server := NewTestSSHServer(t)
1197	defer server.Shutdown()
1198
1199	user := GenerateUser("alice")
1200	RegisterUserWithServer(server, user)
1201
1202	subClient, err := user.NewClient()
1203	if err != nil {
1204		t.Fatalf("failed to connect subscriber: %v", err)
1205	}
1206	defer func() { _ = subClient.Close() }()
1207
1208	pubClient, err := user.NewClient()
1209	if err != nil {
1210		t.Fatalf("failed to connect publisher: %v", err)
1211	}
1212	defer func() { _ = pubClient.Close() }()
1213
1214	// Start subscriber
1215	subSession, err := subClient.NewSession()
1216	if err != nil {
1217		t.Fatalf("failed to create sub session: %v", err)
1218	}
1219	defer func() { _ = subSession.Close() }()
1220
1221	subStdout, err := subSession.StdoutPipe()
1222	if err != nil {
1223		t.Fatalf("failed to get sub stdout: %v", err)
1224	}
1225
1226	if err := subSession.Start("sub emptytopic -c"); err != nil {
1227		t.Fatalf("failed to start sub: %v", err)
1228	}
1229
1230	time.Sleep(100 * time.Millisecond)
1231
1232	// Publish with -e flag (empty message) - should not require stdin
1233	output, err := user.RunCommand(pubClient, "pub emptytopic -e -c")
1234	if err != nil {
1235		t.Logf("pub -e completed: %v, output: %s", err, output)
1236	}
1237
1238	// Subscriber should receive something (even if empty/minimal)
1239	// The -e flag sends a 1-byte buffer
1240	received := make([]byte, 10)
1241	n, err := subStdout.Read(received)
1242	if err != nil && err != io.EOF {
1243		t.Logf("read result: n=%d, err=%v", n, err)
1244	}
1245
1246	// With -e flag, we expect to receive at least 1 byte
1247	if n < 1 {
1248		t.Errorf("subscriber should receive empty message signal, got %d bytes", n)
1249	}
1250}
1251
1252func TestPipe_AccessControl(t *testing.T) {
1253	server := NewTestSSHServer(t)
1254	defer server.Shutdown()
1255
1256	alice := GenerateUser("alice")
1257	bob := GenerateUser("bob")
1258	RegisterUserWithServer(server, alice)
1259	RegisterUserWithServer(server, bob)
1260
1261	aliceClient, err := alice.NewClient()
1262	if err != nil {
1263		t.Fatalf("failed to connect alice: %v", err)
1264	}
1265	defer func() { _ = aliceClient.Close() }()
1266
1267	bobClient, err := bob.NewClient()
1268	if err != nil {
1269		t.Fatalf("failed to connect bob: %v", err)
1270	}
1271	defer func() { _ = bobClient.Close() }()
1272
1273	// Alice creates a pipe with access control allowing bob
1274	aliceSession, err := aliceClient.NewSession()
1275	if err != nil {
1276		t.Fatalf("failed to create alice session: %v", err)
1277	}
1278	defer func() { _ = aliceSession.Close() }()
1279
1280	aliceStdin, err := aliceSession.StdinPipe()
1281	if err != nil {
1282		t.Fatalf("failed to get alice stdin: %v", err)
1283	}
1284
1285	aliceStdout, err := aliceSession.StdoutPipe()
1286	if err != nil {
1287		t.Fatalf("failed to get alice stdout: %v", err)
1288	}
1289
1290	if err := aliceSession.Start("pipe accesspipe -a alice,bob -c"); err != nil {
1291		t.Fatalf("failed to start alice pipe: %v", err)
1292	}
1293
1294	time.Sleep(100 * time.Millisecond)
1295
1296	// Bob joins the pipe using alice's namespace
1297	bobSession, err := bobClient.NewSession()
1298	if err != nil {
1299		t.Fatalf("failed to create bob session: %v", err)
1300	}
1301	defer func() { _ = bobSession.Close() }()
1302
1303	bobStdin, err := bobSession.StdinPipe()
1304	if err != nil {
1305		t.Fatalf("failed to get bob stdin: %v", err)
1306	}
1307
1308	bobStdout, err := bobSession.StdoutPipe()
1309	if err != nil {
1310		t.Fatalf("failed to get bob stdout: %v", err)
1311	}
1312
1313	if err := bobSession.Start("pipe alice/accesspipe -c"); err != nil {
1314		t.Fatalf("failed to start bob pipe: %v", err)
1315	}
1316
1317	time.Sleep(100 * time.Millisecond)
1318
1319	// Alice sends message to bob
1320	aliceMsg := "hello bob\n"
1321	_, err = aliceStdin.Write([]byte(aliceMsg))
1322	if err != nil {
1323		t.Fatalf("alice failed to write: %v", err)
1324	}
1325
1326	bobReceived := make([]byte, 100)
1327	n, _ := bobStdout.Read(bobReceived)
1328	if !strings.Contains(string(bobReceived[:n]), "hello bob") {
1329		t.Errorf("bob did not receive alice's message, got: %q", string(bobReceived[:n]))
1330	}
1331
1332	// Bob sends message to alice
1333	bobMsg := "hello alice\n"
1334	_, err = bobStdin.Write([]byte(bobMsg))
1335	if err != nil {
1336		t.Fatalf("bob failed to write: %v", err)
1337	}
1338
1339	aliceReceived := make([]byte, 100)
1340	n, _ = aliceStdout.Read(aliceReceived)
1341	if !strings.Contains(string(aliceReceived[:n]), "hello alice") {
1342		t.Errorf("alice did not receive bob's message, got: %q", string(aliceReceived[:n]))
1343	}
1344}
1345
1346func TestPipe_Replay(t *testing.T) {
1347	server := NewTestSSHServer(t)
1348	defer server.Shutdown()
1349
1350	user := GenerateUser("alice")
1351	RegisterUserWithServer(server, user)
1352
1353	client, err := user.NewClient()
1354	if err != nil {
1355		t.Fatalf("failed to connect: %v", err)
1356	}
1357	defer func() { _ = client.Close() }()
1358
1359	// Start pipe with replay flag (-r)
1360	session, err := client.NewSession()
1361	if err != nil {
1362		t.Fatalf("failed to create session: %v", err)
1363	}
1364	defer func() { _ = session.Close() }()
1365
1366	stdin, err := session.StdinPipe()
1367	if err != nil {
1368		t.Fatalf("failed to get stdin: %v", err)
1369	}
1370
1371	stdout, err := session.StdoutPipe()
1372	if err != nil {
1373		t.Fatalf("failed to get stdout: %v", err)
1374	}
1375
1376	if err := session.Start("pipe replaytopic -r -c"); err != nil {
1377		t.Fatalf("failed to start pipe: %v", err)
1378	}
1379
1380	time.Sleep(100 * time.Millisecond)
1381
1382	// Send a message - with -r flag, should receive it back
1383	testMsg := "echo back\n"
1384	_, err = stdin.Write([]byte(testMsg))
1385	if err != nil {
1386		t.Fatalf("failed to write: %v", err)
1387	}
1388
1389	received := make([]byte, 100)
1390	n, err := stdout.Read(received)
1391	if err != nil && err != io.EOF {
1392		t.Logf("read error: %v", err)
1393	}
1394
1395	if !strings.Contains(string(received[:n]), "echo back") {
1396		t.Errorf("with -r flag, sender should receive own message back, got: %q", string(received[:n]))
1397	}
1398}
1399
1400func TestAccessControl_UnauthorizedUserDenied(t *testing.T) {
1401	server := NewTestSSHServer(t)
1402	defer server.Shutdown()
1403
1404	alice := GenerateUser("alice")
1405	bob := GenerateUser("bob")
1406	charlie := GenerateUser("charlie")
1407	RegisterUserWithServer(server, alice)
1408	RegisterUserWithServer(server, bob)
1409	RegisterUserWithServer(server, charlie)
1410
1411	aliceClient, err := alice.NewClient()
1412	if err != nil {
1413		t.Fatalf("failed to connect alice: %v", err)
1414	}
1415	defer func() { _ = aliceClient.Close() }()
1416
1417	charlieClient, err := charlie.NewClient()
1418	if err != nil {
1419		t.Fatalf("failed to connect charlie: %v", err)
1420	}
1421	defer func() { _ = charlieClient.Close() }()
1422
1423	// Alice creates a topic with access only for alice and bob (not charlie)
1424	aliceSession, err := aliceClient.NewSession()
1425	if err != nil {
1426		t.Fatalf("failed to create alice session: %v", err)
1427	}
1428	defer func() { _ = aliceSession.Close() }()
1429
1430	if err := aliceSession.Start("sub restrictedtopic -a alice,bob -c"); err != nil {
1431		t.Fatalf("failed to start alice sub: %v", err)
1432	}
1433
1434	time.Sleep(100 * time.Millisecond)
1435
1436	// Charlie tries to publish to alice's restricted topic - should be denied
1437	output, err := charlie.RunCommandWithStdin(charlieClient, "pub alice/restrictedtopic -c", "unauthorized message")
1438	if err != nil {
1439		t.Logf("charlie pub completed with error (expected): %v", err)
1440	}
1441
1442	// Charlie should get access denied or the message should not be delivered
1443	if strings.Contains(output, "access denied") {
1444		t.Logf("charlie correctly received access denied")
1445	} else {
1446		t.Logf("charlie output: %q (access control may work differently)", output)
1447	}
1448}
1449
1450func TestPubSub_MultipleSubscribers(t *testing.T) {
1451	server := NewTestSSHServer(t)
1452	defer server.Shutdown()
1453
1454	user := GenerateUser("alice")
1455	RegisterUserWithServer(server, user)
1456
1457	pubClient, err := user.NewClient()
1458	if err != nil {
1459		t.Fatalf("failed to connect publisher: %v", err)
1460	}
1461	defer func() { _ = pubClient.Close() }()
1462
1463	sub1Client, err := user.NewClient()
1464	if err != nil {
1465		t.Fatalf("failed to connect subscriber 1: %v", err)
1466	}
1467	defer func() { _ = sub1Client.Close() }()
1468
1469	sub2Client, err := user.NewClient()
1470	if err != nil {
1471		t.Fatalf("failed to connect subscriber 2: %v", err)
1472	}
1473	defer func() { _ = sub2Client.Close() }()
1474
1475	sub3Client, err := user.NewClient()
1476	if err != nil {
1477		t.Fatalf("failed to connect subscriber 3: %v", err)
1478	}
1479	defer func() { _ = sub3Client.Close() }()
1480
1481	// Start three subscribers
1482	sub1Session, err := sub1Client.NewSession()
1483	if err != nil {
1484		t.Fatalf("failed to create sub1 session: %v", err)
1485	}
1486	defer func() { _ = sub1Session.Close() }()
1487
1488	sub1Stdout, err := sub1Session.StdoutPipe()
1489	if err != nil {
1490		t.Fatalf("failed to get sub1 stdout: %v", err)
1491	}
1492
1493	if err := sub1Session.Start("sub fanout -c"); err != nil {
1494		t.Fatalf("failed to start sub1: %v", err)
1495	}
1496
1497	sub2Session, err := sub2Client.NewSession()
1498	if err != nil {
1499		t.Fatalf("failed to create sub2 session: %v", err)
1500	}
1501	defer func() { _ = sub2Session.Close() }()
1502
1503	sub2Stdout, err := sub2Session.StdoutPipe()
1504	if err != nil {
1505		t.Fatalf("failed to get sub2 stdout: %v", err)
1506	}
1507
1508	if err := sub2Session.Start("sub fanout -c"); err != nil {
1509		t.Fatalf("failed to start sub2: %v", err)
1510	}
1511
1512	sub3Session, err := sub3Client.NewSession()
1513	if err != nil {
1514		t.Fatalf("failed to create sub3 session: %v", err)
1515	}
1516	defer func() { _ = sub3Session.Close() }()
1517
1518	sub3Stdout, err := sub3Session.StdoutPipe()
1519	if err != nil {
1520		t.Fatalf("failed to get sub3 stdout: %v", err)
1521	}
1522
1523	if err := sub3Session.Start("sub fanout -c"); err != nil {
1524		t.Fatalf("failed to start sub3: %v", err)
1525	}
1526
1527	time.Sleep(100 * time.Millisecond)
1528
1529	// Publish a single message
1530	testMessage := "broadcast message"
1531	_, err = user.RunCommandWithStdin(pubClient, "pub fanout -c", testMessage)
1532	if err != nil {
1533		t.Logf("pub completed: %v", err)
1534	}
1535
1536	// All three subscribers should receive the message
1537	received1 := make([]byte, 100)
1538	n1, _ := sub1Stdout.Read(received1)
1539	if !strings.Contains(string(received1[:n1]), testMessage) {
1540		t.Errorf("subscriber 1 did not receive message, got: %q", string(received1[:n1]))
1541	}
1542
1543	received2 := make([]byte, 100)
1544	n2, _ := sub2Stdout.Read(received2)
1545	if !strings.Contains(string(received2[:n2]), testMessage) {
1546		t.Errorf("subscriber 2 did not receive message, got: %q", string(received2[:n2]))
1547	}
1548
1549	received3 := make([]byte, 100)
1550	n3, _ := sub3Stdout.Read(received3)
1551	if !strings.Contains(string(received3[:n3]), testMessage) {
1552		t.Errorf("subscriber 3 did not receive message, got: %q", string(received3[:n3]))
1553	}
1554}
1555
1556// Monitor CLI Tests
1557
1558func TestMonitor_UnauthenticatedUserDenied(t *testing.T) {
1559	server := NewTestSSHServer(t)
1560	defer server.Shutdown()
1561
1562	user := GenerateUser("anonymous")
1563
1564	client, err := user.NewClient()
1565	if err != nil {
1566		t.Fatalf("failed to connect: %v", err)
1567	}
1568	defer func() { _ = client.Close() }()
1569
1570	output, err := user.RunCommand(client, "monitor my-service 1h")
1571	if err != nil {
1572		t.Logf("command error (expected): %v", err)
1573	}
1574
1575	if !strings.Contains(output, "access denied") {
1576		t.Errorf("expected 'access denied', got: %s", output)
1577	}
1578}
1579
1580func TestMonitor_CreateMonitor(t *testing.T) {
1581	server := NewTestSSHServer(t)
1582	defer server.Shutdown()
1583
1584	user := GenerateUser("alice")
1585	RegisterUserWithServer(server, user)
1586
1587	client, err := user.NewClient()
1588	if err != nil {
1589		t.Fatalf("failed to connect: %v", err)
1590	}
1591	defer func() { _ = client.Close() }()
1592
1593	output, err := user.RunCommand(client, "monitor pico-uptime 24h")
1594	if err != nil {
1595		t.Logf("command completed: %v", err)
1596	}
1597
1598	if strings.Contains(output, "access denied") {
1599		t.Errorf("authenticated user should not get access denied, got: %s", output)
1600	}
1601
1602	// Verify monitor was created in DB (topic is stored with user prefix)
1603	monitor, err := server.DBPool.FindPipeMonitorByTopic("alice-id", "alice/pico-uptime")
1604	if err != nil {
1605		t.Fatalf("monitor should exist in DB: %v", err)
1606	}
1607
1608	if monitor.WindowDur != 24*time.Hour {
1609		t.Errorf("expected window duration 24h, got: %v", monitor.WindowDur)
1610	}
1611
1612	if !strings.Contains(output, "alice/pico-uptime") || !strings.Contains(output, "24h") {
1613		t.Errorf("output should confirm monitor creation, got: %s", output)
1614	}
1615}
1616
1617func TestMonitor_UpdateMonitor(t *testing.T) {
1618	server := NewTestSSHServer(t)
1619	defer server.Shutdown()
1620
1621	user := GenerateUser("alice")
1622	RegisterUserWithServer(server, user)
1623
1624	client, err := user.NewClient()
1625	if err != nil {
1626		t.Fatalf("failed to connect: %v", err)
1627	}
1628	defer func() { _ = client.Close() }()
1629
1630	// Create initial monitor
1631	_, err = user.RunCommand(client, "monitor my-cron 1h")
1632	if err != nil {
1633		t.Logf("create command completed: %v", err)
1634	}
1635
1636	// Upsert with new duration
1637	output, err := user.RunCommand(client, "monitor my-cron 6h")
1638	if err != nil {
1639		t.Logf("update command completed: %v", err)
1640	}
1641
1642	// Verify monitor was updated (topic is stored with user prefix)
1643	monitor, err := server.DBPool.FindPipeMonitorByTopic("alice-id", "alice/my-cron")
1644	if err != nil {
1645		t.Fatalf("monitor should exist in DB: %v", err)
1646	}
1647
1648	if monitor.WindowDur != 6*time.Hour {
1649		t.Errorf("expected window duration 6h after update, got: %v", monitor.WindowDur)
1650	}
1651
1652	if !strings.Contains(output, "6h") {
1653		t.Errorf("output should confirm updated duration, got: %s", output)
1654	}
1655}
1656
1657func TestMonitor_DeleteMonitor(t *testing.T) {
1658	server := NewTestSSHServer(t)
1659	defer server.Shutdown()
1660
1661	user := GenerateUser("alice")
1662	RegisterUserWithServer(server, user)
1663
1664	client, err := user.NewClient()
1665	if err != nil {
1666		t.Fatalf("failed to connect: %v", err)
1667	}
1668	defer func() { _ = client.Close() }()
1669
1670	// Create monitor first
1671	_, err = user.RunCommand(client, "monitor to-delete 1h")
1672	if err != nil {
1673		t.Logf("create command completed: %v", err)
1674	}
1675
1676	// Verify it exists (topic is stored with user prefix)
1677	_, err = server.DBPool.FindPipeMonitorByTopic("alice-id", "alice/to-delete")
1678	if err != nil {
1679		t.Fatalf("monitor should exist before deletion: %v", err)
1680	}
1681
1682	// Delete it
1683	output, err := user.RunCommand(client, "monitor to-delete -d")
1684	if err != nil {
1685		t.Logf("delete command completed: %v", err)
1686	}
1687
1688	// Verify it's gone (topic is stored with user prefix)
1689	_, err = server.DBPool.FindPipeMonitorByTopic("alice-id", "alice/to-delete")
1690	if err == nil {
1691		t.Errorf("monitor should be deleted from DB")
1692	}
1693
1694	if !strings.Contains(output, "deleted") && !strings.Contains(output, "removed") {
1695		t.Logf("output should confirm deletion, got: %s", output)
1696	}
1697}
1698
1699func TestMonitor_InvalidDuration(t *testing.T) {
1700	server := NewTestSSHServer(t)
1701	defer server.Shutdown()
1702
1703	user := GenerateUser("alice")
1704	RegisterUserWithServer(server, user)
1705
1706	client, err := user.NewClient()
1707	if err != nil {
1708		t.Fatalf("failed to connect: %v", err)
1709	}
1710	defer func() { _ = client.Close() }()
1711
1712	output, err := user.RunCommand(client, "monitor my-service invaliduration")
1713	if err != nil {
1714		t.Logf("command error (expected): %v", err)
1715	}
1716
1717	if !strings.Contains(output, "invalid") && !strings.Contains(output, "duration") && !strings.Contains(output, "error") {
1718		t.Errorf("expected error about invalid duration, got: %s", output)
1719	}
1720}
1721
1722func TestMonitor_MissingTopic(t *testing.T) {
1723	server := NewTestSSHServer(t)
1724	defer server.Shutdown()
1725
1726	user := GenerateUser("alice")
1727	RegisterUserWithServer(server, user)
1728
1729	client, err := user.NewClient()
1730	if err != nil {
1731		t.Fatalf("failed to connect: %v", err)
1732	}
1733	defer func() { _ = client.Close() }()
1734
1735	output, err := user.RunCommand(client, "monitor")
1736	if err != nil {
1737		t.Logf("command error (expected): %v", err)
1738	}
1739
1740	// Should show usage or error about missing topic
1741	if !strings.Contains(output, "Usage") && !strings.Contains(output, "topic") && !strings.Contains(output, "error") {
1742		t.Errorf("expected usage info or error about missing topic, got: %s", output)
1743	}
1744}
1745
1746// Status CLI Tests
1747
1748func TestStatus_UnauthenticatedUserDenied(t *testing.T) {
1749	server := NewTestSSHServer(t)
1750	defer server.Shutdown()
1751
1752	user := GenerateUser("anonymous")
1753
1754	client, err := user.NewClient()
1755	if err != nil {
1756		t.Fatalf("failed to connect: %v", err)
1757	}
1758	defer func() { _ = client.Close() }()
1759
1760	output, err := user.RunCommand(client, "status")
1761	if err != nil {
1762		t.Logf("command error (expected): %v", err)
1763	}
1764
1765	if !strings.Contains(output, "access denied") {
1766		t.Errorf("expected 'access denied', got: %s", output)
1767	}
1768}
1769
1770func TestStatus_NoMonitors(t *testing.T) {
1771	server := NewTestSSHServer(t)
1772	defer server.Shutdown()
1773
1774	user := GenerateUser("alice")
1775	RegisterUserWithServer(server, user)
1776
1777	client, err := user.NewClient()
1778	if err != nil {
1779		t.Fatalf("failed to connect: %v", err)
1780	}
1781	defer func() { _ = client.Close() }()
1782
1783	output, err := user.RunCommand(client, "status")
1784	if err != nil {
1785		t.Logf("command completed: %v", err)
1786	}
1787
1788	if !strings.Contains(output, "no monitors") && !strings.Contains(output, "empty") {
1789		t.Errorf("expected message about no monitors, got: %s", output)
1790	}
1791}
1792
1793func TestStatus_ShowsMonitorStatus(t *testing.T) {
1794	server := NewTestSSHServer(t)
1795	defer server.Shutdown()
1796
1797	user := GenerateUser("alice")
1798	RegisterUserWithServer(server, user)
1799
1800	client, err := user.NewClient()
1801	if err != nil {
1802		t.Fatalf("failed to connect: %v", err)
1803	}
1804	defer func() { _ = client.Close() }()
1805
1806	// Create a monitor
1807	_, err = user.RunCommand(client, "monitor web-check 1h")
1808	if err != nil {
1809		t.Logf("create monitor completed: %v", err)
1810	}
1811
1812	// Check status
1813	output, err := user.RunCommand(client, "status")
1814	if err != nil {
1815		t.Logf("status command completed: %v", err)
1816	}
1817
1818	if !strings.Contains(output, "web-check") {
1819		t.Errorf("status should list the monitor topic, got: %s", output)
1820	}
1821}
1822
1823func TestStatus_ShowsHealthyUnhealthy(t *testing.T) {
1824	server := NewTestSSHServer(t)
1825	defer server.Shutdown()
1826
1827	user := GenerateUser("alice")
1828	RegisterUserWithServer(server, user)
1829
1830	// Create monitors directly in DB with different states
1831	now := time.Now()
1832	windowEnd := now.Add(1 * time.Hour)
1833	recentPing := now.Add(-30 * time.Minute) // within window - healthy
1834	oldPing := now.Add(-2 * time.Hour)       // outside window - unhealthy
1835
1836	_ = server.DBPool.UpsertPipeMonitor("alice-id", "healthy-service", 1*time.Hour, &windowEnd)
1837	_ = server.DBPool.UpdatePipeMonitorLastPing("alice-id", "healthy-service", &recentPing)
1838
1839	_ = server.DBPool.UpsertPipeMonitor("alice-id", "unhealthy-service", 1*time.Hour, &windowEnd)
1840	_ = server.DBPool.UpdatePipeMonitorLastPing("alice-id", "unhealthy-service", &oldPing)
1841
1842	client, err := user.NewClient()
1843	if err != nil {
1844		t.Fatalf("failed to connect: %v", err)
1845	}
1846	defer func() { _ = client.Close() }()
1847
1848	output, err := user.RunCommand(client, "status")
1849	if err != nil {
1850		t.Logf("status command completed: %v", err)
1851	}
1852
1853	if !strings.Contains(output, "healthy-service") {
1854		t.Errorf("status should list healthy-service, got: %s", output)
1855	}
1856
1857	if !strings.Contains(output, "unhealthy-service") {
1858		t.Errorf("status should list unhealthy-service, got: %s", output)
1859	}
1860
1861	// Should indicate different health states
1862	if !strings.Contains(strings.ToLower(output), "healthy") && !strings.Contains(strings.ToLower(output), "ok") && !strings.Contains(output, "✓") {
1863		t.Logf("status output should indicate health state: %s", output)
1864	}
1865}
1866
1867// RSS CLI Tests
1868
1869func TestRss_UnauthenticatedUserDenied(t *testing.T) {
1870	server := NewTestSSHServer(t)
1871	defer server.Shutdown()
1872
1873	user := GenerateUser("anonymous")
1874
1875	client, err := user.NewClient()
1876	if err != nil {
1877		t.Fatalf("failed to connect: %v", err)
1878	}
1879	defer func() { _ = client.Close() }()
1880
1881	output, err := user.RunCommand(client, "rss")
1882	if err != nil {
1883		t.Logf("command error (expected): %v", err)
1884	}
1885
1886	if !strings.Contains(output, "access denied") {
1887		t.Errorf("expected 'access denied', got: %s", output)
1888	}
1889}
1890
1891func TestRss_GeneratesValidRSS(t *testing.T) {
1892	server := NewTestSSHServer(t)
1893	defer server.Shutdown()
1894
1895	user := GenerateUser("alice")
1896	RegisterUserWithServer(server, user)
1897
1898	// Create a monitor
1899	now := time.Now()
1900	windowEnd := now.Add(1 * time.Hour)
1901	_ = server.DBPool.UpsertPipeMonitor("alice-id", "rss-test-service", 1*time.Hour, &windowEnd)
1902
1903	client, err := user.NewClient()
1904	if err != nil {
1905		t.Fatalf("failed to connect: %v", err)
1906	}
1907	defer func() { _ = client.Close() }()
1908
1909	output, err := user.RunCommand(client, "rss")
1910	if err != nil {
1911		t.Logf("rss command completed: %v", err)
1912	}
1913
1914	// Should output valid RSS XML
1915	if !strings.Contains(output, "<?xml") || !strings.Contains(output, "<rss") {
1916		t.Errorf("expected RSS XML output, got: %s", output)
1917	}
1918
1919	if !strings.Contains(output, "rss-test-service") {
1920		t.Errorf("RSS should contain monitor topic, got: %s", output)
1921	}
1922}
1923
1924func TestRss_AlertsOnStaleMonitor(t *testing.T) {
1925	server := NewTestSSHServer(t)
1926	defer server.Shutdown()
1927
1928	user := GenerateUser("alice")
1929	RegisterUserWithServer(server, user)
1930
1931	// Create a stale monitor (last ping outside window)
1932	now := time.Now()
1933	windowEnd := now.Add(-30 * time.Minute) // window already ended
1934	oldPing := now.Add(-2 * time.Hour)
1935
1936	_ = server.DBPool.UpsertPipeMonitor("alice-id", "stale-service", 1*time.Hour, &windowEnd)
1937	_ = server.DBPool.UpdatePipeMonitorLastPing("alice-id", "stale-service", &oldPing)
1938
1939	client, err := user.NewClient()
1940	if err != nil {
1941		t.Fatalf("failed to connect: %v", err)
1942	}
1943	defer func() { _ = client.Close() }()
1944
1945	output, err := user.RunCommand(client, "rss")
1946	if err != nil {
1947		t.Logf("rss command completed: %v", err)
1948	}
1949
1950	// Should contain alert item for stale service
1951	if !strings.Contains(output, "stale-service") {
1952		t.Errorf("RSS should contain stale-service alert, got: %s", output)
1953	}
1954
1955	// Should have item element for the alert
1956	if !strings.Contains(output, "<item>") {
1957		t.Errorf("RSS should contain item element for alert, got: %s", output)
1958	}
1959}
1960
1961// Pub integration with Monitor
1962
1963func TestPub_UpdatesMonitorLastPing(t *testing.T) {
1964	server := NewTestSSHServer(t)
1965	defer server.Shutdown()
1966
1967	user := GenerateUser("alice")
1968	RegisterUserWithServer(server, user)
1969
1970	// Create a monitor first (topic is stored with user prefix)
1971	now := time.Now()
1972	windowEnd := now.Add(1 * time.Hour)
1973	_ = server.DBPool.UpsertPipeMonitor("alice-id", "alice/ping-test", 1*time.Hour, &windowEnd)
1974
1975	subClient, err := user.NewClient()
1976	if err != nil {
1977		t.Fatalf("failed to connect subscriber: %v", err)
1978	}
1979	defer func() { _ = subClient.Close() }()
1980
1981	pubClient, err := user.NewClient()
1982	if err != nil {
1983		t.Fatalf("failed to connect publisher: %v", err)
1984	}
1985	defer func() { _ = pubClient.Close() }()
1986
1987	// Start subscriber
1988	subSession, err := subClient.NewSession()
1989	if err != nil {
1990		t.Fatalf("failed to create sub session: %v", err)
1991	}
1992	defer func() { _ = subSession.Close() }()
1993
1994	if err := subSession.Start("sub ping-test -c"); err != nil {
1995		t.Fatalf("failed to start sub: %v", err)
1996	}
1997
1998	time.Sleep(100 * time.Millisecond)
1999
2000	// Publish to the monitored topic
2001	_, err = user.RunCommandWithStdin(pubClient, "pub ping-test -c", "health check")
2002	if err != nil {
2003		t.Logf("pub command completed: %v", err)
2004	}
2005
2006	// Verify last_ping was updated (topic is stored with user prefix)
2007	monitor, err := server.DBPool.FindPipeMonitorByTopic("alice-id", "alice/ping-test")
2008	if err != nil {
2009		t.Fatalf("monitor should exist: %v", err)
2010	}
2011
2012	if monitor.LastPing == nil {
2013		t.Errorf("last_ping should be set after pub")
2014	} else if time.Since(*monitor.LastPing) > 5*time.Second {
2015		t.Errorf("last_ping should be recent, got: %v", monitor.LastPing)
2016	}
2017}
2018
2019// Tests for monitor status edge cases
2020
2021func TestStatus_PingAtExactWindowStart(t *testing.T) {
2022	// Bug fix: Status() should use >= for windowStart comparison
2023	// A ping exactly at windowStart should be healthy
2024	now := time.Now().UTC()
2025	windowEnd := now.Add(1 * time.Hour)
2026	windowStart := windowEnd.Add(-1 * time.Hour) // equals now
2027
2028	monitor := &db.PipeMonitor{
2029		LastPing:  &windowStart, // ping exactly at window start
2030		WindowEnd: &windowEnd,
2031		WindowDur: 1 * time.Hour,
2032	}
2033
2034	err := monitor.Status()
2035	if err != nil {
2036		t.Errorf("ping at exact window start should be healthy, got: %v", err)
2037	}
2038}
2039
2040func TestStatus_WindowExpired(t *testing.T) {
2041	// Bug fix: Status() should check if current time is past windowEnd
2042	now := time.Now().UTC()
2043	windowEnd := now.Add(-1 * time.Minute) // window ended 1 minute ago
2044	lastPing := now.Add(-30 * time.Second) // ping was 30 seconds ago
2045
2046	monitor := &db.PipeMonitor{
2047		LastPing:  &lastPing,
2048		WindowEnd: &windowEnd,
2049		WindowDur: 1 * time.Hour,
2050	}
2051
2052	err := monitor.Status()
2053	if err == nil {
2054		t.Error("expired window should be unhealthy")
2055	}
2056	if !strings.Contains(err.Error(), "window expired") {
2057		t.Errorf("error should mention window expired, got: %v", err)
2058	}
2059}
2060
2061func TestStatus_PingResetsWindow(t *testing.T) {
2062	// Bug fix: Every ping should reset window to now + duration
2063	server := NewTestSSHServer(t)
2064	defer server.Shutdown()
2065
2066	user := GenerateUser("alice")
2067	RegisterUserWithServer(server, user)
2068
2069	// Create a monitor with an expired window
2070	expiredWindowEnd := time.Now().UTC().Add(-10 * time.Minute)
2071	_ = server.DBPool.UpsertPipeMonitor("alice-id", "alice/reset-test", 5*time.Minute, &expiredWindowEnd)
2072
2073	client, err := user.NewClient()
2074	if err != nil {
2075		t.Fatalf("failed to connect: %v", err)
2076	}
2077	defer func() { _ = client.Close() }()
2078
2079	// Start a subscriber first so pub doesn't block
2080	subClient, err := user.NewClient()
2081	if err != nil {
2082		t.Fatalf("failed to connect subscriber: %v", err)
2083	}
2084	defer func() { _ = subClient.Close() }()
2085
2086	subSession, err := subClient.NewSession()
2087	if err != nil {
2088		t.Fatalf("failed to create sub session: %v", err)
2089	}
2090	defer func() { _ = subSession.Close() }()
2091
2092	if err := subSession.Start("sub reset-test -c"); err != nil {
2093		t.Fatalf("failed to start sub: %v", err)
2094	}
2095
2096	time.Sleep(100 * time.Millisecond)
2097
2098	// Pub to trigger monitor update
2099	_, err = user.RunCommandWithStdin(client, "pub reset-test -c", "ping")
2100	if err != nil {
2101		t.Logf("pub command completed: %v", err)
2102	}
2103
2104	// Check that window was reset
2105	monitor, err := server.DBPool.FindPipeMonitorByTopic("alice-id", "alice/reset-test")
2106	if err != nil {
2107		t.Fatalf("monitor should exist: %v", err)
2108	}
2109
2110	if monitor.WindowEnd == nil {
2111		t.Fatal("window_end should be set")
2112	}
2113
2114	// Window end should now be in the future
2115	if !monitor.WindowEnd.After(time.Now().UTC()) {
2116		t.Errorf("window_end should be in the future after ping, got: %v", monitor.WindowEnd)
2117	}
2118}
2119
2120func TestStatus_HealthyImmediatelyAfterPing(t *testing.T) {
2121	// Bug fix: After a ping, status should immediately show healthy
2122	server := NewTestSSHServer(t)
2123	defer server.Shutdown()
2124
2125	user := GenerateUser("alice")
2126	RegisterUserWithServer(server, user)
2127
2128	client, err := user.NewClient()
2129	if err != nil {
2130		t.Fatalf("failed to connect: %v", err)
2131	}
2132	defer func() { _ = client.Close() }()
2133
2134	// Create monitor
2135	_, err = user.RunCommand(client, "monitor health-test 5m")
2136	if err != nil {
2137		t.Fatalf("failed to create monitor: %v", err)
2138	}
2139
2140	// Start subscriber
2141	subClient, err := user.NewClient()
2142	if err != nil {
2143		t.Fatalf("failed to connect subscriber: %v", err)
2144	}
2145	defer func() { _ = subClient.Close() }()
2146
2147	subSession, err := subClient.NewSession()
2148	if err != nil {
2149		t.Fatalf("failed to create sub session: %v", err)
2150	}
2151	defer func() { _ = subSession.Close() }()
2152
2153	if err := subSession.Start("sub health-test -c"); err != nil {
2154		t.Fatalf("failed to start sub: %v", err)
2155	}
2156
2157	time.Sleep(100 * time.Millisecond)
2158
2159	// Pub to trigger ping
2160	pubClient, err := user.NewClient()
2161	if err != nil {
2162		t.Fatalf("failed to connect publisher: %v", err)
2163	}
2164	defer func() { _ = pubClient.Close() }()
2165
2166	_, err = user.RunCommandWithStdin(pubClient, "pub health-test -c", "ping")
2167	if err != nil {
2168		t.Logf("pub completed: %v", err)
2169	}
2170
2171	// Immediately check status
2172	statusClient, err := user.NewClient()
2173	if err != nil {
2174		t.Fatalf("failed to connect for status: %v", err)
2175	}
2176	defer func() { _ = statusClient.Close() }()
2177
2178	output, err := user.RunCommand(statusClient, "status")
2179	if err != nil {
2180		t.Logf("status completed: %v", err)
2181	}
2182
2183	if strings.Contains(output, "unhealthy") {
2184		t.Errorf("status should be healthy immediately after ping, got: %s", output)
2185	}
2186	if !strings.Contains(output, "healthy") {
2187		t.Errorf("status should show healthy, got: %s", output)
2188	}
2189}
2190
2191// TestMonitor_FixedWindowNonSliding verifies that pings within the same window
2192// do not slide the window forward. This is a regression test for a bug where
2193// each ping reset window_end to now+dur, creating a sliding window that never fails.
2194//
2195// Expected behavior:
2196//   - last_ping: always updated to show most recent activity (user visibility).
2197//   - window_end: only advances when current time exceeds it (health scheduling).
2198func TestMonitor_FixedWindowNonSliding(t *testing.T) {
2199	server := NewTestSSHServer(t)
2200	defer server.Shutdown()
2201
2202	user := GenerateUser("alice")
2203	RegisterUserWithServer(server, user)
2204
2205	client, err := user.NewClient()
2206	if err != nil {
2207		t.Fatalf("failed to connect: %v", err)
2208	}
2209	defer func() { _ = client.Close() }()
2210
2211	// Create a monitor with 1 hour window
2212	_, err = user.RunCommand(client, "monitor fixed-window-test 1h")
2213	if err != nil {
2214		t.Logf("create command completed: %v", err)
2215	}
2216
2217	// Get the initial window_end
2218	monitor, err := server.DBPool.FindPipeMonitorByTopic("alice-id", "alice/fixed-window-test")
2219	if err != nil {
2220		t.Fatalf("monitor should exist: %v", err)
2221	}
2222	initialWindowEnd := *monitor.WindowEnd
2223
2224	// Simulate a ping by calling updateMonitor directly
2225	handler := server.PipeHandler
2226
2227	// Create a mock CliCmd
2228	mockUser := &db.User{ID: "alice-id", Name: "alice"}
2229	cmd := &CliCmd{
2230		userName: "alice",
2231		user:     mockUser,
2232	}
2233
2234	// First ping - should record last_ping but NOT change window_end
2235	handler.updateMonitor(cmd, "alice/fixed-window-test")
2236
2237	monitor, err = server.DBPool.FindPipeMonitorByTopic("alice-id", "alice/fixed-window-test")
2238	if err != nil {
2239		t.Fatalf("monitor should exist after first ping: %v", err)
2240	}
2241
2242	if monitor.LastPing == nil {
2243		t.Fatalf("last_ping should be set after first ping")
2244	}
2245	firstPingTime := *monitor.LastPing
2246	windowEndAfterFirstPing := *monitor.WindowEnd
2247
2248	// BUG CHECK: With the bug, window_end would have slid forward to now+1h
2249	// With the fix, window_end should remain at the original scheduled time
2250	if !windowEndAfterFirstPing.Equal(initialWindowEnd) {
2251		t.Errorf("BUG DETECTED: window_end should NOT change after first ping within window\n"+
2252			"initial window_end: %v\n"+
2253			"window_end after ping: %v\n"+
2254			"Window slid forward by: %v",
2255			initialWindowEnd.Format(time.RFC3339),
2256			windowEndAfterFirstPing.Format(time.RFC3339),
2257			windowEndAfterFirstPing.Sub(initialWindowEnd))
2258	}
2259
2260	// Second ping - last_ping SHOULD be updated (for user visibility)
2261	// but window_end should NOT change
2262	time.Sleep(10 * time.Millisecond) // Small delay to get different timestamp
2263	handler.updateMonitor(cmd, "alice/fixed-window-test")
2264
2265	monitor, err = server.DBPool.FindPipeMonitorByTopic("alice-id", "alice/fixed-window-test")
2266	if err != nil {
2267		t.Fatalf("monitor should exist after second ping: %v", err)
2268	}
2269
2270	// last_ping SHOULD be updated to show most recent activity
2271	if monitor.LastPing.Equal(firstPingTime) {
2272		t.Errorf("last_ping SHOULD be updated for user visibility\n"+
2273			"first ping time: %v\n"+
2274			"last_ping after second call: %v",
2275			firstPingTime.Format(time.RFC3339Nano),
2276			monitor.LastPing.Format(time.RFC3339Nano))
2277	}
2278
2279	// But window_end should still be the original value (not sliding)
2280	if !monitor.WindowEnd.Equal(initialWindowEnd) {
2281		t.Errorf("BUG DETECTED: window_end should remain at original value\n"+
2282			"initial: %v\n"+
2283			"current: %v",
2284			initialWindowEnd.Format(time.RFC3339),
2285			monitor.WindowEnd.Format(time.RFC3339))
2286	}
2287}
2288
2289func TestSSH_WildcardSub(t *testing.T) {
2290	server := NewTestSSHServer(t)
2291	defer server.Shutdown()
2292
2293	user := GenerateUser("alice")
2294	dbUser := &db.User{ID: "alice-id", Name: "alice"}
2295	server.DBPool.AddUser(dbUser)
2296	server.DBPool.AddPubkey(&db.PublicKey{
2297		ID:     "alice-pk",
2298		UserID: "alice-id",
2299		Key:    user.PublicKey(),
2300	})
2301
2302	client, err := user.NewClient()
2303	if err != nil {
2304		t.Fatalf("failed to dial ssh server: %v", err)
2305	}
2306	defer func() { _ = client.Close() }()
2307
2308	// 1. Subscribe to wildcard topic: "sub metric-drain*"
2309	subSession, err := client.NewSession()
2310	if err != nil {
2311		t.Fatalf("failed to create sub session: %v", err)
2312	}
2313
2314	subOut, err := subSession.StdoutPipe()
2315	if err != nil {
2316		t.Fatalf("failed stdout pipe: %v", err)
2317	}
2318
2319	if err := subSession.Start("sub metric-drain*"); err != nil {
2320		t.Fatalf("failed to start sub: %v", err)
2321	}
2322
2323	var buf bytes.Buffer
2324	var bufMu sync.Mutex
2325	go func() {
2326		b := make([]byte, 1024)
2327		for {
2328			n, err := subOut.Read(b)
2329			if n > 0 {
2330				bufMu.Lock()
2331				buf.Write(b[:n])
2332				bufMu.Unlock()
2333			}
2334			if err != nil {
2335				break
2336			}
2337		}
2338	}()
2339
2340	time.Sleep(100 * time.Millisecond)
2341
2342	// 2. Publish to "metric-drain-pgs"
2343	pubClient1, err := user.NewClient()
2344	if err != nil {
2345		t.Fatalf("failed to dial pub client 1: %v", err)
2346	}
2347	defer func() { _ = pubClient1.Close() }()
2348
2349	pubSession1, err := pubClient1.NewSession()
2350	if err != nil {
2351		t.Fatalf("failed pub session 1: %v", err)
2352	}
2353	pubIn1, _ := pubSession1.StdinPipe()
2354	go func() {
2355		defer func() { _ = pubIn1.Close() }()
2356		_, _ = io.WriteString(pubIn1, "pgs-event\n")
2357	}()
2358	_ = pubSession1.Run("pub metric-drain-pgs -b=false")
2359
2360	// 3. Publish to "metric-drain-prose"
2361	pubClient2, err := user.NewClient()
2362	if err != nil {
2363		t.Fatalf("failed to dial pub client 2: %v", err)
2364	}
2365	defer func() { _ = pubClient2.Close() }()
2366
2367	pubSession2, err := pubClient2.NewSession()
2368	if err != nil {
2369		t.Fatalf("failed pub session 2: %v", err)
2370	}
2371	pubIn2, _ := pubSession2.StdinPipe()
2372	go func() {
2373		defer func() { _ = pubIn2.Close() }()
2374		_, _ = io.WriteString(pubIn2, "prose-event\n")
2375	}()
2376	_ = pubSession2.Run("pub metric-drain-prose -b=false")
2377
2378	time.Sleep(150 * time.Millisecond)
2379	_ = subSession.Close()
2380
2381	bufMu.Lock()
2382	output := buf.String()
2383	bufMu.Unlock()
2384
2385	if !strings.Contains(output, "pgs-event") {
2386		t.Errorf("expected SSH wildcard subscriber output to contain 'pgs-event', got: %q", output)
2387	}
2388	if !strings.Contains(output, "prose-event") {
2389		t.Errorf("expected SSH wildcard subscriber output to contain 'prose-event', got: %q", output)
2390	}
2391}