package twitch import ( "bytes" "context" "encoding/json" "fmt" "io" "log" "net/http" "net/url" "os" "path" "strconv" "codeberg.org/arimelody/ari-stream-tools/twitch/api" ) type ( FullUser struct { ID string `json:"id"` Login string `json:"login"` DisplayName string `json:"display_name"` Type string `json:"type"` BroadcasterType string `json:"broadcaster_type"` Description string `json:"description"` ProfileImageURL string `json:"profile_image_url"` OfflineImageURL string `json:"offline_image_url"` ViewCount int `json:"view_count"` Email string `json:"email"` CreatedAt string `json:"created_at"` } GetUsersResponse struct { Data []FullUser `json:"data"` } ) func (srv *Service) getUsers(ids []string, logins []string) (*GetUsersResponse, error) { if len(ids) == 0 && len(logins) == 0 { return nil, nil } url, err := url.Parse(api.BASE_URL + "/users?" + url.Values{ "id": ids, "login": logins, }.Encode()) req, err := http.NewRequest("GET", url.String(), nil) if err != nil { return nil, err } req.Header.Set("Client-Id", srv.clientID) req.Header.Set("Authorization", "Bearer " + srv.oauthToken.AccessToken) res, err := http.DefaultClient.Do(req) if err != nil { return nil, err } if res.StatusCode != http.StatusOK { body, _ := io.ReadAll(res.Body) return nil, fmt.Errorf("%s: %s", res.Status, string(body)) } data := &GetUsersResponse{} if err := json.NewDecoder(res.Body).Decode(data); err != nil { return nil, err } return data, nil } func (srv *Service) registerEventSubSession(session *api.EventSubSession) error { srv.eventSubSession = session return nil } func (srv *Service) handleNotification(payload *api.EventSubPayload) error { var err error jsonData, err := json.Marshal(payload.Event) if err != nil { return err } switch payload.Subscription.Type { case "channel.follow": var event api.FollowEvent err := json.Unmarshal(jsonData, &event) if err != nil { return fmt.Errorf("Failed to cast to api.FollowEvent: %v", err) } log.Printf("New follow: %s", event.UserName) srv.labels.LatestFollower.Text = event.UserName srv.labels.LatestFollower.C <- event.UserName if err := os.WriteFile( path.Join(DATA_PATH, "state", LABEL_LATEST_FOLLOWER), []byte(event.UserName), 0640, ); err != nil { return fmt.Errorf("Failed to write %s state: %v", LABEL_LATEST_FOLLOWER, err) } case "channel.subscribe": var event api.SubscribeEvent err := json.Unmarshal(jsonData, &event) if err != nil { return fmt.Errorf("Failed to cast to api.SubscribeEvent: %v", err) } log.Printf("New subscription: %s", event.UserName) srv.labels.LatestSubscriber.Text = event.UserName srv.labels.LatestSubscriber.C <- event.UserName if err := os.WriteFile( path.Join(DATA_PATH, "state", LABEL_LATEST_SUBSCRIBER), []byte(event.UserName), 0640, ); err != nil { return fmt.Errorf("Failed to write %s state: %v", LABEL_LATEST_SUBSCRIBER, err) } case "channel.cheer": var event api.CheerEvent err := json.Unmarshal(jsonData, &event) if err != nil { return fmt.Errorf("Failed to cast to api.CheerEvent: %v", err) } log.Printf("%s cheered x%d bits: %s", event.UserName, event.Bits, event.Message) srv.labels.LatestCheer.Text = event.UserName srv.labels.LatestCheer.C <- event.UserName if err := os.WriteFile( path.Join(DATA_PATH, "state", LABEL_LATEST_CHEER), []byte(event.UserName), 0640, ); err != nil { return fmt.Errorf("Failed to write %s state: %v", LABEL_LATEST_CHEER, err) } case "channel.raid": var event api.RaidEvent err := json.Unmarshal(jsonData, &event) if err != nil { return fmt.Errorf("Failed to cast to api.RaidEvent: %v", err) } log.Printf("%s is now raiding with %d viewers!", event.FromUserName, event.Viewers) case "channel.channel_points_custom_reward_redemption.add": var event api.ChannelPointCustomRewardRedeemEvent err := json.Unmarshal(jsonData, &event) if err != nil { return fmt.Errorf("Failed to cast to api.ChannelPointCustomRewardRedeemEvent: %v", err) } log.Printf( "%s just redeemed %s for %d channel points.", event.UserName, event.Reward.Title, event.Reward.Cost, ) case "channel.shoutout.create": var event api.ShoutoutCreate err := json.Unmarshal(jsonData, &event) if err != nil { return fmt.Errorf("Failed to cast to api.ShoutoutCreate: %v", err) } log.Printf( "%s gave a shoutout to %s.", event.FromUserName, event.ToUserName, ) case "channel.chat.message": var event api.ChatEvent err := json.Unmarshal(jsonData, &event) if err != nil { return fmt.Errorf("Failed to cast to api.ChatEvent: %v", err) } if event.Cheer != nil { return nil } log.Printf( "[%s] %s: %s", event.MessageID, event.ChatterLogin, event.Message.Text, ) case "channel.chat.message_delete": var event api.ChatDeleteEvent err := json.Unmarshal(jsonData, &event) if err != nil { return fmt.Errorf("Failed to cast to api.ChatDeleteEvent: %v", err) } log.Printf( "Message %s by %s deleted.", event.MessageID, event.TargetLogin, ) } return nil } func (srv *Service) subscribeToEvent( ctx context.Context, subscriptionType string, version string, condition map[string]string, sessionID string, ) error { type ( Transport struct { Method string `json:"method"` SessionID string `json:"session_id"` } Request struct { Type string `json:"type"` Version string `json:"version"` Condition map[string]string `json:"condition"` Transport Transport `json:"transport"` } ResponseData struct { ID string `json:"id"` Status string `json:"status"` Type string `json:"type"` Version string `json:"version"` Condition map[string]string `json:"condition"` CreatedAt string `json:"created_at"` Transport Transport `json:"transport"` Cost int `json:"cost"` } Response struct { Data []ResponseData `json:"data"` Total int `json:"total"` TotalCost int `json:"total_cost"` MaxTotalCost int `json:"max_total_cost"` } ) bodyBytes, err := json.Marshal(Request{ Type: subscriptionType, Version: version, Condition: condition, Transport: Transport{ Method: "websocket", SessionID: sessionID, }, }) body := bytes.NewBuffer(bodyBytes) client := srv.oauthConfig.Client(ctx, srv.oauthToken) req, err := http.NewRequest( "POST", api.BASE_URL + "/eventsub/subscriptions", body, ) if err != nil { return err } req.Header.Set("Content-Type", "application/json") req.Header.Set("Client-Id", srv.clientID) req.Header.Set("Authorization", "Bearer " + srv.oauthToken.AccessToken) res, err := client.Do(req) if err != nil { return err } if res.StatusCode != http.StatusAccepted { body, _ := io.ReadAll(res.Body) return fmt.Errorf("%s: %s", res.Status, string(body)) } return nil } func (srv *Service) subscribeToDefaultEvents(ctx context.Context) { // channel.follow if err := srv.subscribeToEvent( ctx, "channel.follow", "2", map[string]string{ "broadcaster_user_id": srv.channelID, "moderator_user_id": srv.channelID, }, srv.eventSubSession.ID, ); err != nil { log.Printf("Failed to subscribe to channel.follow: %v", err) } // channel.subscribe if err := srv.subscribeToEvent( ctx, "channel.subscribe", "1", map[string]string{ "broadcaster_user_id": srv.channelID }, srv.eventSubSession.ID, ); err != nil { log.Printf("Failed to subscribe to channel.subscribe: %v", err) } // channel.cheer if err := srv.subscribeToEvent( ctx, "channel.cheer", "1", map[string]string{ "broadcaster_user_id": srv.channelID }, srv.eventSubSession.ID, ); err != nil { log.Printf("Failed to subscribe to channel.cheer: %v", err) } // channel.raid if err := srv.subscribeToEvent( ctx, "channel.raid", "1", map[string]string{ "to_broadcaster_user_id": srv.channelID }, srv.eventSubSession.ID, ); err != nil { log.Printf("Failed to subscribe to channel.raid: %v", err) } // channel.channel_points_custom_reward_redemption.add if err := srv.subscribeToEvent( ctx, "channel.channel_points_custom_reward_redemption.add", "1", map[string]string{ "broadcaster_user_id": srv.channelID }, srv.eventSubSession.ID, ); err != nil { log.Printf("Failed to subscribe to channel.channel_points_custom_reward_redemption.add: %v", err) } // channel.shoutout.create if err := srv.subscribeToEvent( ctx, "channel.shoutout.create", "1", map[string]string{ "broadcaster_user_id": srv.channelID, "moderator_user_id": srv.channelID, }, srv.eventSubSession.ID, ); err != nil { log.Printf("Failed to subscribe to channel.shoutout.create: %v", err) } // channel.chat.message if err := srv.subscribeToEvent( ctx, "channel.chat.message", "1", map[string]string{ "broadcaster_user_id": srv.channelID, "user_id": srv.channelID, }, srv.eventSubSession.ID, ); err != nil { log.Printf("Failed to subscribe to channel.chat.message: %v", err) } // channel.chat.message_delete if err := srv.subscribeToEvent( ctx, "channel.chat.message_delete", "1", map[string]string{ "broadcaster_user_id": srv.channelID, "user_id": srv.channelID, }, srv.eventSubSession.ID, ); err != nil { log.Printf("Failed to subscribe to channel.chat.message_delete: %v", err) } } type ( Follower struct { FollowedAt string `json:"followed_at"` UserID string `json:"user_id"` UserLogin string `json:"user_login"` UserName string `json:"user_name"` } Pagination struct { Cursor string `json:"cursor"` } GetFollowersResponse struct { Data []Follower `json:"data"` Pagination Pagination `json:"pagination"` Total int `json:"total"` Points int `json:"points"` } ) func (srv *Service) getFollowers( ctx context.Context, broadcasterID string, limit int, ) (*GetFollowersResponse, error) { client := srv.oauthConfig.Client(ctx, srv.oauthToken) u, err := url.Parse(api.BASE_URL + "/channels/followers?" + url.Values{ "broadcaster_id": []string{ broadcasterID }, "first": []string { strconv.Itoa(limit) }, }.Encode()) req, err := http.NewRequest("GET", u.String(), nil) if err != nil { return nil, err } req.Header.Set("Content-Type", "application/json") req.Header.Set("Client-Id", srv.clientID) req.Header.Set("Authorization", "Bearer " + srv.oauthToken.AccessToken) res, err := client.Do(req) if err != nil { return nil, err } if res.StatusCode != http.StatusOK { body, _ := io.ReadAll(res.Body) return nil, fmt.Errorf("%s: %s", res.Status, string(body)) } resData := &GetFollowersResponse{} if err := json.NewDecoder(res.Body).Decode(resData); err != nil { return nil, err } return resData, nil } type ( Subscriber struct { BroadcasterID string `json:"broadcaster_id"` BroadcasterLogin string `json:"broadcaster_login"` BroadcasterName string `json:"broadcaster_name"` GifterID string `json:"gifter_id"` GifterLogin string `json:"gifter_login"` GifterName string `json:"gifter_name"` IsGift bool `json:"is_gift"` Tier string `json:"tier"` PlanName string `json:"plan_name"` UserID string `json:"user_id"` UserLogin string `json:"user_login"` UserName string `json:"user_name"` } GetSubscribersResponse struct { Data []Subscriber `json:"data"` Pagination Pagination `json:"pagination"` Total int `json:"total"` Points int `json:"points"` } ) // Unfortunately, this function is very unreliable for pulling chronological // subscription records. Twitch API does not currently provide a mechanism for // fetching the latest subscriber; this will need to be tracked manually. func (srv *Service) getSubscribers( ctx context.Context, broadcasterID string, limit int, ) (*GetSubscribersResponse, error) { client := srv.oauthConfig.Client(ctx, srv.oauthToken) u, err := url.Parse(api.BASE_URL + "/subscriptions?" + url.Values{ "broadcaster_id": []string{ broadcasterID }, "first": []string { strconv.Itoa(limit) }, }.Encode()) req, err := http.NewRequest("GET", u.String(), nil) if err != nil { return nil, err } req.Header.Set("Content-Type", "application/json") req.Header.Set("Client-Id", srv.clientID) req.Header.Set("Authorization", "Bearer " + srv.oauthToken.AccessToken) res, err := client.Do(req) if err != nil { return nil, err } if res.StatusCode != http.StatusOK { body, _ := io.ReadAll(res.Body) return nil, fmt.Errorf("%s: %s", res.Status, string(body)) } resData := &GetSubscribersResponse{} if err := json.NewDecoder(res.Body).Decode(resData); err != nil { return nil, err } return resData, nil }