package twitch import ( "context" "crypto/rand" "embed" "encoding/json" "errors" "fmt" "io" "log" "net/http" "net/url" "os" "os/signal" "path" "strings" "time" "codeberg.org/arimelody/ari-stream-tools/broadcast" "codeberg.org/arimelody/ari-stream-tools/config" "codeberg.org/arimelody/ari-stream-tools/twitch/api" "github.com/gin-gonic/gin" "github.com/gorilla/websocket" "golang.org/x/oauth2" "golang.org/x/oauth2/twitch" ) type ( twitchLabel struct { Text string C chan string Broadcast broadcast.BroadcastChannel[string] } twitchLabels struct { LatestFollower *twitchLabel LatestSubscriber *twitchLabel LatestCheer *twitchLabel } ChatMessageModifier string chatMessage struct { ID string `json:"id"` Username string `json:"username"` Colour string `json:"colour,omitempty"` Text string `json:"text"` Badges []api.ChatBadge `json:"badges"` Fragments []api.ChatMessageFragment `json:"fragments"` Modifier ChatMessageModifier `json:"modifier,omitempty"` AnnouncementColour string `json:"announcement_colour,omitempty"` } SystemChatMessageType string systemChatMessage struct { Text string `json:"text"` Type SystemChatMessageType `json:"type,omitempty"` } deleteChatMessage struct { Username string `json:"username"` MessageID string `json:"message_id"` } ServiceOptions struct { Host string Port int16 ChannelName string ClientID string ClientSecret string } Service struct { host string port int16 channelName string channelID string clientID string clientSecret string oauthConfig *oauth2.Config oauthState string oauthToken *oauth2.Token eventSubSession *api.EventSubSession labels *twitchLabels ChatC chan *chatMessage ChatBroadcast broadcast.BroadcastChannel[*chatMessage] SystemChatC chan *systemChatMessage SystemChatBroadcast broadcast.BroadcastChannel[*systemChatMessage] DeleteChatC chan *deleteChatMessage DeleteChatBroadcast broadcast.BroadcastChannel[*deleteChatMessage] } ) const ( LABEL_LATEST_FOLLOWER string = "latest-follower" LABEL_LATEST_SUBSCRIBER string = "latest-subscriber" LABEL_LATEST_CHEER string = "latest-cheer" CHAT_MODIFIER_NONE ChatMessageModifier = "" CHAT_MODIFIER_SUBSCRIBE ChatMessageModifier = "subscribe" CHAT_MODIFIER_CHEER ChatMessageModifier = "cheer" CHAT_MODIFIER_HIGHLIGHT ChatMessageModifier = "highlight" CHAT_MODIFIER_ANNOUNCEMENT ChatMessageModifier = "announcement" SYSTEM_CHAT_FOLLOW SystemChatMessageType = "follow" SYSTEM_CHAT_SUBSCRIBE SystemChatMessageType = "subscribe" SYSTEM_CHAT_CHEER SystemChatMessageType = "cheer" SYSTEM_CHAT_RAID SystemChatMessageType = "raid" SYSTEM_CHAT_POLL SystemChatMessageType = "poll" SYSTEM_CHAT_SHOUTOUT SystemChatMessageType = "shoutout" SYSTEM_CHAT_POINT_REDEEM SystemChatMessageType = "point-redeem" ) var DATA_PATH string = path.Join(config.CONFIG_DIR, "twitch") var CONFIG_FILEPATH string = path.Join(DATA_PATH, "twitch-config.json") var AUTH_FILEPATH string = path.Join(DATA_PATH, "twitch-auth") //go:embed public var publicFS embed.FS //go:embed pages var pagesFS embed.FS func New(ctx context.Context, opts ServiceOptions) (*Service, error) { if err := os.MkdirAll(DATA_PATH, 0750); err != nil { panic(err) } if err := os.MkdirAll(path.Join(DATA_PATH, "state"), 0750); err != nil { panic(err) } if len(opts.ChannelName) == 0 { return nil, errors.New("config: channel_name cannot be empty") } if len(opts.ClientID) == 0 { return nil, errors.New("config: client_id cannot be empty") } if len(opts.ClientSecret) == 0 { return nil, errors.New("config: client_secret cannot be empty") } latestFollowerC := make(chan string) latestFollowerBroadcast := broadcast.NewBroadcastChannel( ctx, latestFollowerC) latestSubscriberC := make(chan string) latestSubscriberBroadcast := broadcast.NewBroadcastChannel( ctx, latestSubscriberC) latestCheerC := make(chan string) latestCheerBroadcast := broadcast.NewBroadcastChannel( ctx, latestCheerC) chatC := make(chan *chatMessage) chatBroadcast := broadcast.NewBroadcastChannel( ctx, chatC) systemChatC := make(chan *systemChatMessage) systemChatBroadcast := broadcast.NewBroadcastChannel( ctx, systemChatC) deleteChatC := make(chan *deleteChatMessage) deleteChatBroadcast := broadcast.NewBroadcastChannel( ctx, deleteChatC) srv := &Service{ host: opts.Host, port: opts.Port, labels: &twitchLabels{ LatestFollower: &twitchLabel{ C: latestFollowerC, Broadcast: latestFollowerBroadcast, }, LatestSubscriber: &twitchLabel{ C: latestSubscriberC, Broadcast: latestSubscriberBroadcast, }, LatestCheer: &twitchLabel{ C: latestCheerC, Broadcast: latestCheerBroadcast, }, }, channelName: opts.ChannelName, clientID: opts.ClientID, clientSecret: opts.ClientSecret, oauthConfig: &oauth2.Config{ ClientID: opts.ClientID, ClientSecret: opts.ClientSecret, Endpoint: twitch.Endpoint, Scopes: []string{ "moderator:read:followers", "user:read:chat", "user:bot", "channel:bot", "channel:read:subscriptions", "bits:read", "channel:read:redemptions", "channel:read:polls", "channel:read:predictions", "channel:read:hype_train", "moderator:read:shoutouts", }, RedirectURL: fmt.Sprintf("http://%s:%d/twitch/auth", opts.Host, opts.Port), }, ChatC: chatC, ChatBroadcast: chatBroadcast, SystemChatC: systemChatC, SystemChatBroadcast: systemChatBroadcast, DeleteChatC: deleteChatC, DeleteChatBroadcast: deleteChatBroadcast, } if authFile, err := os.OpenFile(AUTH_FILEPATH, os.O_RDONLY, 0600); err != nil { if !os.IsNotExist(err) { log.Fatalf("open %s: %v", AUTH_FILEPATH, err) } } else { defer authFile.Close() srv.oauthToken = &oauth2.Token{} err = json.NewDecoder(authFile).Decode(srv.oauthToken) if err != nil { log.Printf("read %s: %v", AUTH_FILEPATH, err) } } return srv, nil } func (srv *Service) BindRoutes(group *gin.RouterGroup) { group.GET("/public/*path", func(ctx *gin.Context) { filepath := strings.TrimPrefix(ctx.Request.URL.Path, "/twitch/") http.ServeFileFS(ctx.Writer, ctx.Request, publicFS, filepath) }) group.GET("/login", func(ctx *gin.Context) { srv.oauthState = rand.Text() authCodeURL := srv.oauthConfig.AuthCodeURL(srv.oauthState) ctx.Redirect(http.StatusTemporaryRedirect, authCodeURL) }) group.GET("/auth", func(ctx *gin.Context) { code := ctx.Query("code") scope := ctx.Query("scope") resState := ctx.Query("state") if len(code) == 0 { ctx.String(http.StatusBadRequest, "code cannot be empty"); return } if len(scope) == 0 { ctx.String(http.StatusBadRequest, "scope cannot be empty"); return } if resState != srv.oauthState { ctx.String(http.StatusBadRequest, "state mismatch"); return } token, err := srv.oauthConfig.Exchange(ctx, code) if err != nil { log.Printf("Could not exchange OAuth2 code: %v", err) ctx.String(http.StatusBadRequest, "Could not exchange OAuth2 code.") return } srv.oauthToken = token authFile, err := os.OpenFile(AUTH_FILEPATH, os.O_CREATE | os.O_RDWR, 0600) if err != nil { log.Printf("open %s: %v", AUTH_FILEPATH, err) ctx.String(http.StatusInternalServerError, http.StatusText(http.StatusInternalServerError)) return } defer authFile.Close() authFile.Truncate(0) json.NewEncoder(authFile).Encode(srv.oauthToken) ctx.String( http.StatusOK, "Authentication successful! You may now close this tab.", ) go srv.start(ctx) }) group.GET("/sse/labels/:label", func(ctx *gin.Context) { ctx.Header("connection", "keep-alive") labelParam := ctx.Param("label") var label *twitchLabel switch labelParam { case LABEL_LATEST_FOLLOWER: label = srv.labels.LatestFollower case LABEL_LATEST_SUBSCRIBER: label = srv.labels.LatestSubscriber case LABEL_LATEST_CHEER: label = srv.labels.LatestCheer default: ctx.String(http.StatusBadRequest, "Unknown label %s", labelParam) return } labelUpdate := label.Broadcast.Subscribe() defer label.Broadcast.Cancel(labelUpdate) ctx.SSEvent("update", label.Text) ticker := time.NewTicker(10 * time.Millisecond) ctx.Stream(func(w io.Writer) bool { select { case text := <-labelUpdate: ctx.SSEvent("update", text) case <-ticker.C: } return true }) }) group.GET("/labels/:label", func(ctx *gin.Context) { labelParam := ctx.Param("label") switch labelParam { case LABEL_LATEST_FOLLOWER: case LABEL_LATEST_SUBSCRIBER: case LABEL_LATEST_CHEER: default: http.NotFound(ctx.Writer, ctx.Request) return } http.ServeFileFS(ctx.Writer, ctx.Request, pagesFS, "pages/labels.html") }) group.GET("/sse/chat", func(ctx *gin.Context) { ctx.Header("connection", "keep-alive") chatUpdate := srv.ChatBroadcast.Subscribe() defer srv.ChatBroadcast.Cancel(chatUpdate) systemChatUpdate := srv.SystemChatBroadcast.Subscribe() defer srv.SystemChatBroadcast.Cancel(systemChatUpdate) deleteChatUpdate := srv.DeleteChatBroadcast.Subscribe() defer srv.DeleteChatBroadcast.Cancel(deleteChatUpdate) ctx.SSEvent("hello", "hello") ticker := time.NewTicker(10 * time.Millisecond) ctx.Stream(func(w io.Writer) bool { select { case update := <-chatUpdate: ctx.SSEvent("chat", update) case update := <-systemChatUpdate: ctx.SSEvent("system", update) case update := <-deleteChatUpdate: ctx.SSEvent("delete", update) case <-ticker.C: } return true }) }) group.GET("/chat", func(ctx *gin.Context) { http.ServeFileFS(ctx.Writer, ctx.Request, pagesFS, "pages/chat.html") }) } func (srv *Service) Run(ctx context.Context) { if srv.oauthToken == nil || !srv.oauthToken.Valid() { log.Printf("Log in with Twitch: http://%s:%d/twitch/login", srv.host, srv.port) } else { if err := srv.start(ctx); err != nil { log.Printf("Failed to start Twitch service: %v", err) } } } func (srv *Service) start(ctx context.Context) error { if res, err := srv.getUsers(nil, []string{ srv.channelName }); err == nil { if len(res.Data) == 0 { return fmt.Errorf("resolve username \"%s\", it may not exist?", srv.channelName) } srv.channelID = res.Data[0].ID } else { return fmt.Errorf("fetch Twitch users: %v", err) } if res, err := srv.getFollowers(ctx, srv.channelID, 1); err == nil { if len(res.Data) > 0 { srv.labels.LatestFollower.Text = res.Data[0].UserName log.Printf("Loaded latest follower: %s", res.Data[0].UserName) } } else { log.Printf("Failed to fetch latest follower: %v", err) } // TODO: build satellite service for offline-fetching subscribers and cheers if data, err := os.ReadFile(path.Join(DATA_PATH, "state", LABEL_LATEST_SUBSCRIBER)); err == nil { srv.labels.LatestSubscriber.Text = string(data) log.Printf("Loaded latest subscriber: %s", string(data)) } if data, err := os.ReadFile(path.Join(DATA_PATH, "state", LABEL_LATEST_CHEER)); err == nil { srv.labels.LatestCheer.Text = string(data) log.Printf("Loaded latest cheer: %s", string(data)) } srv.startWebsocketListener(ctx) return nil } func (srv *Service) startWebsocketListener(ctx context.Context) { interrupt := make(chan os.Signal, 1) signal.Notify(interrupt, os.Interrupt) u, err := url.Parse(api.EVENTSUB_URL + "?keepalive_timeout_seconds=600") if err != nil { panic(err) } c, _, err := websocket.DefaultDialer.Dial(u.String(), http.Header{ "Authorization": []string{ "Bearer " + srv.clientSecret }, }) if err != nil { log.Printf("Failed to connect to Twitch: %v", err) } defer c.Close() failed := make(chan error) go func() { for { _, rawData, err := c.ReadMessage() if err != nil { failed <- err; return } var message api.EventSubMessage err = json.Unmarshal(rawData, &message) if err != nil { failed <- fmt.Errorf("parse JSON: %v", err) return } if message.Payload.Session != nil { if err := srv.registerEventSubSession(message.Payload.Session); err != nil { failed <- fmt.Errorf("register eventsub session: %v", err) return } srv.subscribeToDefaultEvents(ctx) log.Printf("Connected to Twitch for channel %s.", srv.channelName) } if message.Metadata.MessageType == string(api.NOTIFICATION) { if message.Payload.Subscription == nil { continue } if err := srv.handleNotification(&message.Payload); err != nil { log.Printf("Failed to handle %s event: %v", message.Payload.Subscription.Type, err) } continue } } }() select { case err := <-failed: log.Printf("Twitch error: %v", err) case <-ctx.Done(): } }