@@ -14,6 +14,10 @@ const (
1414 antigravityTokenRefreshSkew = 3 * time .Minute
1515 antigravityTokenCacheSkew = 5 * time .Minute
1616 antigravityBackfillCooldown = 5 * time .Minute
17+ // antigravityRequestRefreshTimeout 请求路径上 token 刷新的最大等待时间。
18+ // 超过此时间直接放弃刷新、标记账号临时不可调度并触发 failover,
19+ // 让后台 TokenRefreshService 在下个周期继续重试。
20+ antigravityRequestRefreshTimeout = 8 * time .Second
1721)
1822
1923// AntigravityTokenCache token cache interface.
@@ -28,6 +32,7 @@ type AntigravityTokenProvider struct {
2832 refreshAPI * OAuthRefreshAPI
2933 executor OAuthRefreshExecutor
3034 refreshPolicy ProviderRefreshPolicy
35+ tempUnschedCache TempUnschedCache // 用于同步更新 Redis 临时不可调度缓存
3136}
3237
3338func NewAntigravityTokenProvider (
@@ -54,6 +59,11 @@ func (p *AntigravityTokenProvider) SetRefreshPolicy(policy ProviderRefreshPolicy
5459 p .refreshPolicy = policy
5560}
5661
62+ // SetTempUnschedCache injects temp unschedulable cache for immediate scheduler sync.
63+ func (p * AntigravityTokenProvider ) SetTempUnschedCache (cache TempUnschedCache ) {
64+ p .tempUnschedCache = cache
65+ }
66+
5767// GetAccessToken returns a valid access_token.
5868func (p * AntigravityTokenProvider ) GetAccessToken (ctx context.Context , account * Account ) (string , error ) {
5969 if account == nil {
@@ -88,8 +98,13 @@ func (p *AntigravityTokenProvider) GetAccessToken(ctx context.Context, account *
8898 expiresAt := account .GetCredentialAsTime ("expires_at" )
8999 needsRefresh := expiresAt == nil || time .Until (* expiresAt ) <= antigravityTokenRefreshSkew
90100 if needsRefresh && p .refreshAPI != nil && p .executor != nil {
91- result , err := p .refreshAPI .RefreshIfNeeded (ctx , account , p .executor , antigravityTokenRefreshSkew )
101+ // 请求路径使用短超时,避免代理不通时阻塞过久(后台刷新服务会继续重试)
102+ refreshCtx , cancel := context .WithTimeout (ctx , antigravityRequestRefreshTimeout )
103+ defer cancel ()
104+ result , err := p .refreshAPI .RefreshIfNeeded (refreshCtx , account , p .executor , antigravityTokenRefreshSkew )
92105 if err != nil {
106+ // 标记账号临时不可调度,避免后续请求继续命中
107+ p .markTempUnschedulable (account , err )
93108 if p .refreshPolicy .OnRefreshError == ProviderRefreshErrorReturn {
94109 return "" , err
95110 }
@@ -172,6 +187,45 @@ func (p *AntigravityTokenProvider) shouldAttemptBackfill(accountID int64) bool {
172187 return true
173188}
174189
190+ // markTempUnschedulable 在请求路径上 token 刷新失败时标记账号临时不可调度。
191+ // 同时写 DB 和 Redis 缓存,确保调度器立即跳过该账号。
192+ // 使用 background context 因为请求 context 可能已超时。
193+ func (p * AntigravityTokenProvider ) markTempUnschedulable (account * Account , refreshErr error ) {
194+ if p .accountRepo == nil || account == nil {
195+ return
196+ }
197+ now := time .Now ()
198+ until := now .Add (tokenRefreshTempUnschedDuration )
199+ reason := "token refresh failed on request path: " + refreshErr .Error ()
200+ bgCtx := context .Background ()
201+ if err := p .accountRepo .SetTempUnschedulable (bgCtx , account .ID , until , reason ); err != nil {
202+ slog .Warn ("antigravity_token_provider.set_temp_unschedulable_failed" ,
203+ "account_id" , account .ID ,
204+ "error" , err ,
205+ )
206+ return
207+ }
208+ slog .Warn ("antigravity_token_provider.temp_unschedulable_set" ,
209+ "account_id" , account .ID ,
210+ "until" , until .Format (time .RFC3339 ),
211+ "reason" , reason ,
212+ )
213+ // 同步写 Redis 缓存,调度器立即生效
214+ if p .tempUnschedCache != nil {
215+ state := & TempUnschedState {
216+ UntilUnix : until .Unix (),
217+ TriggeredAtUnix : now .Unix (),
218+ ErrorMessage : reason ,
219+ }
220+ if err := p .tempUnschedCache .SetTempUnsched (bgCtx , account .ID , state ); err != nil {
221+ slog .Warn ("antigravity_token_provider.temp_unsched_cache_set_failed" ,
222+ "account_id" , account .ID ,
223+ "error" , err ,
224+ )
225+ }
226+ }
227+ }
228+
175229func (p * AntigravityTokenProvider ) markBackfillAttempted (accountID int64 ) {
176230 p .backfillCooldown .Store (accountID , time .Now ())
177231}
0 commit comments