- Sort Score
- Result 10 results
- Languages All
Results 1 - 10 of 30 for NewTicker (0.2 sec)
-
internal/store/store.go
func replayItems[I any](store Store[I], doneCh <-chan struct{}, log logger, id string) <-chan Key { keyCh := make(chan Key) go func() { defer xioutil.SafeClose(keyCh) retryTicker := time.NewTicker(retryInterval) defer retryTicker.Stop() for { names, err := store.List() if err != nil { log(context.Background(), fmt.Errorf("store.List() failed with: %w", err), id) } else {
Go - Registered: Sun Apr 21 19:28:08 GMT 2024 - Last Modified: Mon Mar 25 16:44:20 GMT 2024 - 3.5K bytes - Viewed (0) -
cmd/admin-handlers.go
Go - Registered: Sun Apr 28 19:28:10 GMT 2024 - Last Modified: Sun Apr 21 11:43:18 GMT 2024 - 97.3K bytes - Viewed (2) -
cmd/metrics-v2_test.go
label: labels[1], }, { val: 0.31, label: labels[1], }, { val: 0.61, label: labels[3], }, { val: 0.79, label: labels[2], }, } ticker := time.NewTicker(1 * time.Millisecond) defer ticker.Stop() for _, obs := range observations { // Send observations once every 1ms, to simulate delay between // observations. This is to test the channel based
Go - Registered: Sun Apr 28 19:28:10 GMT 2024 - Last Modified: Mon Mar 04 18:05:56 GMT 2024 - 2.3K bytes - Viewed (0) -
internal/logger/target/http/http.go
// before this method is launched. logChLock.Lock() globalBuffer := logChBuffers[name] logChLock.Unlock() newTicker := time.NewTicker(time.Second) isTick := false for { isTick = false select { case _ = <-newTicker.C: isTick = true case entry, _ = <-globalBuffer: case entry, ok = <-h.logCh: if !ok { return } case <-ctx.Done():
Go - Registered: Sun Apr 21 19:28:08 GMT 2024 - Last Modified: Mon Mar 25 16:44:20 GMT 2024 - 14.9K bytes - Viewed (0) -
cmd/listen-notification-handlers.go
} if pingInterval < 1 { writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrInvalidQueryParams), r.URL) return } t := time.NewTicker(time.Duration(pingInterval) * time.Second) defer t.Stop() emptyEventTicker = t.C } else { // Deprecated Apr 2023 t := time.NewTicker(500 * time.Millisecond) defer t.Stop() keepAliveTicker = t.C } enc := json.NewEncoder(w) for { select {
Go - Registered: Sun Apr 28 19:28:10 GMT 2024 - Last Modified: Thu Apr 04 12:04:40 GMT 2024 - 6K bytes - Viewed (0) -
cmd/bucket-replication-stats.go
qCache: newQueueCache(r), pCache: newProxyStatsCache(), srStats: newSRStats(), movingAvgTicker: time.NewTicker(2 * time.Second), wTimer: time.NewTicker(2 * time.Second), qTimer: time.NewTicker(2 * time.Second), workers: newActiveWorkerStat(r), registry: r, } go rs.collectWorkerMetrics(ctx) go rs.collectQueueMetrics(ctx) return &rs
Go - Registered: Sun Apr 28 19:28:10 GMT 2024 - Last Modified: Thu Feb 22 06:26:06 GMT 2024 - 13.4K bytes - Viewed (0) -
internal/bucket/bandwidth/monitor.go
m := &Monitor{ bucketsMeasurement: make(map[BucketOptions]*bucketMeasurement), bucketsThrottle: make(map[BucketOptions]*bucketThrottle), bucketMovingAvgTicker: time.NewTicker(2 * time.Second), ctx: ctx, NodeCount: numNodes, } go m.trackEWMA() return m } func (m *Monitor) updateMeasurement(opts BucketOptions, bytes uint64) { m.mlock.Lock()
Go - Registered: Sun Apr 28 19:28:10 GMT 2024 - Last Modified: Mon Feb 19 22:54:46 GMT 2024 - 6K bytes - Viewed (0) -
cmd/erasure.go
saverWg.Add(1) go func() { // Add jitter to the update time so multiple sets don't sync up. updateTime := 30*time.Second + time.Duration(float64(10*time.Second)*rand.Float64()) t := time.NewTicker(updateTime) defer t.Stop() defer saverWg.Done() var lastSave time.Time for { select { case <-t.C: if cache.Info.LastUpdate.Equal(lastSave) { continue }
Go - Registered: Sun Apr 28 19:28:10 GMT 2024 - Last Modified: Fri Apr 26 06:32:14 GMT 2024 - 16K bytes - Viewed (1) -
cmd/site-replication-metrics.go
} else { srs.XferRateSml.addSize(sz, duration) } } func newSRStats() *SRStats { s := SRStats{ M: make(map[string]*SRStatus), movingAvgTicker: time.NewTicker(time.Second * 2), } go s.trackEWMA() return &s } func (sr *SRStats) trackEWMA() { for { select { case <-sr.movingAvgTicker.C: sr.updateMovingAvg() case <-GlobalContext.Done():
Go - Registered: Sun Apr 28 19:28:10 GMT 2024 - Last Modified: Tue Feb 06 06:00:45 GMT 2024 - 8.2K bytes - Viewed (0) -
istioctl/pkg/wait/wait.go
w = withContext(ctx) w.Go(func(result chan string) error { result <- generation return nil }) } // wait for all deployed versions to be contained in generations t := time.NewTicker(pollInterval) printVerbosef(cmd, "getting first version from chan") firstVersion, err := w.BlockingRead() if err != nil { return fmt.Errorf("unable to retrieve Kubernetes resource %s: %v", "", err) }
Go - Registered: Wed May 01 22:53:12 GMT 2024 - Last Modified: Sat Feb 17 12:24:17 GMT 2024 - 10.1K bytes - Viewed (0)