@@ -45,9 +45,8 @@ type Kick struct {
4545 SaveCfgFn func (Config ) error
4646 PrepareLocker xsync.Mutex
4747
48- lazyInitOnce sync.Once
49- getAccessTokenLocker xsync.Mutex
50- lastGetMutexSuccessAt time.Time
48+ lazyInitOnce sync.Once
49+ getAccessTokenLocker xsync.Mutex
5150}
5251
5352var _ streamcontrol.StreamController [StreamProfile ] = (* Kick )(nil )
@@ -123,6 +122,11 @@ func (k *Kick) onUserAccessTokenRefreshed(
123122 if err != nil {
124123 logger .Errorf (ctx , "unable to save the config: %v" , err )
125124 }
125+ client := k .GetClient ()
126+ if client != nil {
127+ client .SetUserAccessToken (userAccessToken )
128+ client .SetUserRefreshToken (refreshToken )
129+ }
126130 })
127131}
128132
@@ -140,13 +144,10 @@ func (k *Kick) keepAliveLoop(
140144 time .Sleep (time .Second )
141145 continue
142146 }
143- client := k .GetClient ()
144- if client == nil {
145- logger .Errorf (ctx , "client is not initialized" )
146- time .Sleep (time .Second )
147- continue
148- }
149- _ , err := client .GetLivestreams (k .CloseCtx , gokick .NewLivestreamListFilter ().SetBroadcasterUserIDs (int (k .Channel .UserID )))
147+ _ , err := k .getLivestreams (
148+ k .CloseCtx ,
149+ gokick .NewLivestreamListFilter ().SetBroadcasterUserIDs (int (k .Channel .UserID )),
150+ )
150151 if err != nil {
151152 logger .Errorf (ctx , "unable to get my stream status: %v" , err )
152153 time .Sleep (time .Second )
@@ -245,13 +246,8 @@ func (k *Kick) getAccessToken(
245246) (_err error ) {
246247 logger .Tracef (ctx , "getAccessToken" )
247248 defer func () { logger .Tracef (ctx , "/getAccessToken: %v" , _err ) }()
248- now := time .Now ()
249249 return xsync .DoR1 (ctx , & k .getAccessTokenLocker , func () error {
250- if k .lastGetMutexSuccessAt .After (now ) {
251- return nil
252- }
253- err := k .getAccessTokenNoLock (ctx )
254- return err
250+ return k .getAccessTokenNoLock (ctx )
255251 })
256252}
257253
@@ -261,12 +257,6 @@ func (k *Kick) getAccessTokenNoLock(
261257 logger .Tracef (ctx , "getAccessTokenNoLock" )
262258 defer func () { logger .Tracef (ctx , "/getAccessTokenNoLock: %v" , _err ) }()
263259
264- defer func () {
265- if _err == nil {
266- k .lastGetMutexSuccessAt = time .Now ()
267- }
268- }()
269-
270260 getPortsFn := k .CurrentConfig .Config .GetOAuthListenPorts
271261 if getPortsFn == nil {
272262 // TODO: find a way to adjust the OAuth ports dynamically without re-creating the Kick client.
@@ -391,10 +381,34 @@ func (k *Kick) Flush(ctx context.Context) error {
391381}
392382
393383func (k * Kick ) EndStream (ctx context.Context ) error {
394- logger . Warnf ( ctx , "not implemented yet" )
384+ // Kick ends a stream automatically, nothing to do:
395385 return nil
396386}
397387
388+ func (k * Kick ) getLivestreams (
389+ ctx context.Context ,
390+ filter gokick.LivestreamListFilter ,
391+ ) (_ret * gokick.LivestreamsResponseWrapper , _err error ) {
392+ logger .Debugf (ctx , "getLivestreams" )
393+ defer func () { logger .Debugf (ctx , "/getLivestreams: %v, %v" , _ret , _err ) }()
394+
395+ if err := k .prepare (ctx ); err != nil {
396+ return nil , fmt .Errorf ("unable to get a prepared client: %w" , err )
397+ }
398+
399+ client := k .GetClient ()
400+ if client == nil {
401+ return nil , fmt .Errorf ("client is not initialized" )
402+ }
403+
404+ resp , err := client .GetLivestreams (k .CloseCtx , filter )
405+ if err != nil {
406+ return nil , fmt .Errorf ("unable to get my stream status: %w" , err )
407+ }
408+
409+ return & resp , nil
410+ }
411+
398412func (k * Kick ) GetStreamStatus (
399413 ctx context.Context ,
400414) (_ret * streamcontrol.StreamStatus , _err error ) {
@@ -765,6 +779,14 @@ func (k *Kick) getAccessTokenIfNeeded(
765779 logger .Tracef (ctx , "getAccessTokenIfNeeded" )
766780 defer func () { logger .Tracef (ctx , "/getAccessTokenIfNeeded: %v" , _err ) }()
767781
782+ if time .Now ().After (k .CurrentConfig .Config .UserAccessTokenExpiresAt .Add (- 30 * time .Second )) {
783+ if k .CurrentConfig .Config .RefreshToken .Get () != "" {
784+ if err := k .refreshAccessToken (ctx ); err != nil {
785+ return fmt .Errorf ("unable to refresh the access token: %w" , err )
786+ }
787+ }
788+ }
789+
768790 if k .CurrentConfig .Config .UserAccessToken .Get () != "" {
769791 return nil
770792 }
@@ -777,6 +799,29 @@ func (k *Kick) getAccessTokenIfNeeded(
777799 return nil
778800}
779801
802+ func (k * Kick ) refreshAccessToken (
803+ ctx context.Context ,
804+ ) (_err error ) {
805+ logger .Tracef (ctx , "refreshAccessToken" )
806+ defer func () { logger .Tracef (ctx , "/refreshAccessToken: %v" , _err ) }()
807+
808+ resp , err := k .GetClient ().RefreshToken (ctx , k .CurrentConfig .Config .RefreshToken .Get ())
809+ if err != nil {
810+ logger .Errorf (ctx , "unable to refresh the token: %v" , err )
811+ if getErr := k .getAccessToken (ctx ); getErr != nil {
812+ return fmt .Errorf ("unable to refresh access token (%w); and unable to get a new access token (%w)" , err , getErr )
813+ }
814+ return nil
815+ }
816+
817+ err = k .setToken (ctx , resp , time .Now ())
818+ if err != nil {
819+ return fmt .Errorf ("unable to set access token: %w" )
820+ }
821+
822+ return nil
823+ }
824+
780825func (k * Kick ) IsCapable (
781826 ctx context.Context ,
782827 cap streamcontrol.Capability ,
0 commit comments