Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Support
Keyboard shortcuts
?
Submit feedback
Contribute to GitLab
Sign in / Register
Toggle navigation
G
gitlab-workhorse
Project overview
Project overview
Details
Activity
Releases
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Issues
0
Issues
0
List
Boards
Labels
Milestones
Merge Requests
1
Merge Requests
1
Analytics
Analytics
Repository
Value Stream
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Create a new issue
Commits
Issue Boards
Open sidebar
nexedi
gitlab-workhorse
Commits
8a4f8840
Commit
8a4f8840
authored
Feb 01, 2017
by
Tomasz Maczukin
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
Add metrics for queueing mechanism
parent
7378c05b
Changes
1
Show whitespace changes
Inline
Side-by-side
Showing
1 changed file
with
51 additions
and
0 deletions
+51
-0
internal/queueing/queue.go
internal/queueing/queue.go
+51
-0
No files found.
internal/queueing/queue.go
View file @
8a4f8840
...
...
@@ -3,6 +3,8 @@ package queueing
import
(
"errors"
"time"
"github.com/prometheus/client_golang/prometheus"
)
type
errTooManyRequests
struct
{
error
}
...
...
@@ -11,16 +13,57 @@ type errQueueingTimedout struct{ error }
var
ErrTooManyRequests
=
&
errTooManyRequests
{
errors
.
New
(
"too many requests queued"
)}
var
ErrQueueingTimedout
=
&
errQueueingTimedout
{
errors
.
New
(
"queueing timedout"
)}
var
(
queueingLimit
=
prometheus
.
NewGauge
(
prometheus
.
GaugeOpts
{
Name
:
"gitlab_workhorse_queueing_limit"
,
Help
:
"Current limit set for the queueing mechanism"
,
})
queueingQueueLimit
=
prometheus
.
NewGauge
(
prometheus
.
GaugeOpts
{
Name
:
"gitlab_workhorse_queueing_queue_limit"
,
Help
:
"Current queueLimit set for the queueing mechanism"
,
})
queueingBusy
=
prometheus
.
NewGauge
(
prometheus
.
GaugeOpts
{
Name
:
"gitlab_workhorse_queueing_busy"
,
Help
:
"How many queued requests are now processed"
,
})
queueingWaiting
=
prometheus
.
NewGauge
(
prometheus
.
GaugeOpts
{
Name
:
"gitlab_workhorse_queueing_waiting"
,
Help
:
"How many requests are now queued"
,
})
queueingErrors
=
prometheus
.
NewCounterVec
(
prometheus
.
CounterOpts
{
Name
:
"gitlab_workhorse_queueing_errors"
,
Help
:
"How many times the TooManyRequests or QueueintTimedout errors were returned while queueing, partitioned by error type"
,
},
[]
string
{
"type"
},
)
)
type
Queue
struct
{
busyCh
chan
struct
{}
waitingCh
chan
struct
{}
}
func
init
()
{
prometheus
.
MustRegister
(
queueingErrors
)
prometheus
.
MustRegister
(
queueingLimit
)
prometheus
.
MustRegister
(
queueingBusy
)
prometheus
.
MustRegister
(
queueingWaiting
)
prometheus
.
MustRegister
(
queueingQueueLimit
)
}
// NewQueue creates a new queue
// limit specifies number of requests run concurrently
// queueLimit specifies maximum number of requests that can be queued
// if the number of requests is above the limit
func
NewQueue
(
limit
,
queueLimit
uint
)
*
Queue
{
queueingLimit
.
Set
(
float64
(
limit
))
queueingQueueLimit
.
Set
(
float64
(
queueLimit
))
return
&
Queue
{
busyCh
:
make
(
chan
struct
{},
limit
),
waitingCh
:
make
(
chan
struct
{},
limit
+
queueLimit
),
...
...
@@ -35,20 +78,24 @@ func (s *Queue) Acquire(timeout time.Duration) (err error) {
// push item to a queue to claim your own slot (non-blocking)
select
{
case
s
.
waitingCh
<-
struct
{}{}
:
queueingWaiting
.
Inc
()
break
default
:
queueingErrors
.
WithLabelValues
(
"too_many_requests"
)
.
Inc
()
return
ErrTooManyRequests
}
defer
func
()
{
if
err
!=
nil
{
<-
s
.
waitingCh
queueingWaiting
.
Dec
()
}
}()
// fast path: push item to current processed items (non-blocking)
select
{
case
s
.
busyCh
<-
struct
{}{}
:
queueingBusy
.
Inc
()
return
nil
default
:
break
...
...
@@ -60,9 +107,11 @@ func (s *Queue) Acquire(timeout time.Duration) (err error) {
// push item to current processed items (blocking)
select
{
case
s
.
busyCh
<-
struct
{}{}
:
queueingBusy
.
Inc
()
return
nil
case
<-
timer
.
C
:
queueingErrors
.
WithLabelValues
(
"queueing_timedout"
)
.
Inc
()
return
ErrQueueingTimedout
}
}
...
...
@@ -72,5 +121,7 @@ func (s *Queue) Acquire(timeout time.Duration) (err error) {
func
(
s
*
Queue
)
Release
()
{
// dequeue from queue to allow next request to be processed
<-
s
.
waitingCh
queueingWaiting
.
Dec
()
<-
s
.
busyCh
queueingBusy
.
Dec
()
}
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment