package lighthouse_test import ( "context" "fmt" "strconv" "testing" "time" "git.golem15.com/golem15/summercms/modules/backpack" "git.golem15.com/golem15/summercms/modules/compass" "git.golem15.com/golem15/summercms/modules/lagoon" "git.golem15.com/golem15/summercms/modules/lighthouse" "gorm.io/gorm" ) // newApp returns an app whose config selects a realtime driver. func newApp(settings map[string]any) (*backpack.App, error) { cfg, err := compass.Open(compass.Options{Dir: "config", Env: "development", Environ: []string{}}) if err != nil { return nil, err } for k, v := range settings { if err := cfg.Set(k, v); err != nil { return nil, err } } return backpack.New(cfg), nil } // isMember stands in for the application's own membership query. func isMember(ctx context.Context, userID uint, blogID int64) bool { return userID == 42 && blogID == 7 } func ExampleRegistry_Register() { app, err := newApp(map[string]any{"realtime.driver": "memory"}) if err != nil { fmt.Println(err) return } svc, err := lighthouse.From(app) if err != nil { fmt.Println(err) return } // blog:{entity}:{id} channels are open to members of the blog only. The // authorizer runs on every subscribe; nothing is cached. err = svc.Registry().Register("blog", lighthouse.AuthorizerFunc( func(ctx context.Context, userID uint, channel string) lighthouse.Result { if isMember(ctx, userID, lighthouse.ChannelID(channel)) { return lighthouse.Allowed(nil) } return lighthouse.Denied("not a member of the blog") })) if err != nil { fmt.Println(err) return } for _, sub := range []struct { user uint channel string }{{42, "blog:7"}, {42, "blog:8"}, {42, "presence:blog:7"}, {42, "shop:7"}} { ns, presence := lighthouse.ParseChannel(sub.channel) auth, ok := svc.Registry().Get(ns) if !ok { fmt.Println(sub.channel, "no authorizer") continue } res := auth.Authorize(context.Background(), sub.user, sub.channel) fmt.Printf("%d %s namespace=%s presence=%v allowed=%v reason=%q\n", sub.user, sub.channel, ns, presence, res.Allowed, res.Reason()) } fmt.Println(lighthouse.ChannelID("blog:12abc"), lighthouse.FormatChannels("acme", []string{"Blog:7"})) // Output: // 42 blog:7 namespace=blog presence=false allowed=true reason="" // 42 blog:8 namespace=blog presence=false allowed=false reason="not a member of the blog" // 42 presence:blog:7 namespace=blog presence=true allowed=false reason="not a member of the blog" // shop:7 no authorizer // 12 [acme:blog:7] } // Post is the acme.blog post model. It knows nothing about realtime. type Post struct { ID uint `gorm:"column:id;primaryKey"` BlogID uint `gorm:"column:blog_id"` Title string `gorm:"column:title"` } func (Post) TableName() string { return "acme_blog_posts" } // bindPosts registers the broadcast contract of Post from the plugin's Boot. func bindPosts(svc *lighthouse.Service) error { // docs:start bind return lighthouse.Bind[Post](svc, lighthouse.Binding[Post]{ Alias: "blog.post", Channels: func(ctx context.Context, tx *gorm.DB, p *Post) ([]string, error) { return []string{"blog:" + strconv.FormatUint(uint64(p.BlogID), 10)}, nil }, }) // docs:end bind } // publishPost creates a post; its broadcast is published after commit. func publishPost(ctx context.Context, db *gorm.DB, fail bool) error { // docs:start write return lagoon.Transaction(ctx, db, func(ctx context.Context, tx *gorm.DB) error { if err := tx.Create(&Post{BlogID: 7, Title: "Hello"}).Error; err != nil { return err } // The broadcast job is now queued in this transaction. It is // published only if the transaction commits. if fail { return fmt.Errorf("rolled back") } return nil }) // docs:end write } // importPosts writes many posts and publishes one summary event. func importPosts(ctx context.Context, svc *lighthouse.Service, db *gorm.DB, titles []string) error { // docs:start bulk return lighthouse.WithoutBroadcasting[Post](ctx, func(ctx context.Context) error { return lagoon.Transaction(ctx, db, func(ctx context.Context, tx *gorm.DB) error { for _, title := range titles { if err := tx.Create(&Post{BlogID: 7, Title: title}).Error; err != nil { return err } } return svc.Emit(ctx, tx, lighthouse.Broadcast{ Channels: []string{"blog:7"}, Event: "blog.posts_imported", Payload: struct { Count int `json:"count"` }{len(titles)}, }) }) }) // docs:end bulk } // waitFor waits until the memory driver has recorded n publications. func waitFor(t *testing.T, mem *lighthouse.MemoryDriver, n int) []lighthouse.Publication { t.Helper() deadline := time.Now().Add(10 * time.Second) for len(mem.Publications()) < n { if time.Now().After(deadline) { t.Fatalf("got %d publications, want %d", len(mem.Publications()), n) } time.Sleep(20 * time.Millisecond) } time.Sleep(300 * time.Millisecond) // catch extra publications pubs := mem.Publications() if len(pubs) != n { t.Fatalf("got %d publications, want exactly %d: %+v", len(pubs), n, pubs) } return pubs } // TestDocsBroadcast runs the bind, write and bulk regions of the Realtime // page on the package's Postgres harness with a job worker. func TestDocsBroadcast(t *testing.T) { _, svc, db := lighthouse.DocsEnv(t) if err := db.Exec(`CREATE TABLE acme_blog_posts (id SERIAL PRIMARY KEY, blog_id INTEGER NOT NULL, title TEXT NOT NULL)`).Error; err != nil { t.Fatal(err) } if err := bindPosts(svc); err != nil { t.Fatal(err) } mem, ok := svc.Driver().(*lighthouse.MemoryDriver) if !ok { t.Fatalf("driver is %T", svc.Driver()) } ctx := t.Context() if err := publishPost(ctx, db, true); err == nil { t.Fatal("rolled-back publish succeeded") } if err := publishPost(ctx, db, false); err != nil { t.Fatal(err) } pubs := waitFor(t, mem, 1) if pubs[0].Event != "created.blog.post" || len(pubs[0].Channels) != 1 || pubs[0].Channels[0] != "blog:7" { t.Fatalf("publication = %+v", pubs[0]) } if err := importPosts(ctx, svc, db, []string{"A", "B", "C"}); err != nil { t.Fatal(err) } pubs = waitFor(t, mem, 2) if pubs[1].Event != "blog.posts_imported" || string(pubs[1].Payload) != `{"count":3}` { t.Fatalf("summary = %+v %s", pubs[1], pubs[1].Payload) } }