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}