Skip to content

Commit f408e4a

Browse files
Merge pull request #150 from niclic/federation-additions
API for runtime parameters, e.g. for configuring federation upstream components.
2 parents fb19cc2 + 0259cac commit f408e4a

7 files changed

Lines changed: 793 additions & 61 deletions

File tree

README.md

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -326,6 +326,64 @@ resp, err := rmqc.DeleteShovel("/", "a.shovel")
326326

327327
```
328328

329+
### Operations on Runtime (vhost-scoped) Parameters
330+
331+
```golang
332+
// list all runtime parameters
333+
params, err := rmqc.ListRuntimeParameters()
334+
// => []RuntimeParameter, error
335+
336+
// list all runtime parameters for a component
337+
params, err := rmqc.ListRuntimeParametersFor("federation-upstream")
338+
// => []RuntimeParameter, error
339+
340+
// list runtime parameters in a vhost
341+
params, err := rmqc.ListRuntimeParametersIn("federation-upstream", "/")
342+
// => []RuntimeParameter, error
343+
344+
// information about a runtime parameter
345+
p, err := rmqc.GetRuntimeParameter("federation-upstream", "/", "name")
346+
// => *RuntimeParameter, error
347+
348+
// declare or update a runtime parameter
349+
resp, err := rmqc.PutRuntimeParameter("federation-upstream", "/", "name", FederationDefinition{
350+
Uri: "amqp://server-name",
351+
})
352+
// => *http.Response, error
353+
354+
// remove a runtime parameter
355+
resp, err := rmqc.DeleteRuntimeParameter("federation-upstream", "/", "name")
356+
// => *http.Response, error
357+
358+
```
359+
360+
### Operations on Federation Upstreams
361+
362+
```golang
363+
// list all federation upstreams
364+
ups, err := rmqc.ListFederationUpstreams()
365+
// => []FederationUpstream, error
366+
367+
// list federation upstreams in a vhost
368+
ups, err := rmqc.ListFederationUpstreamsIn("/")
369+
// => []FederationUpstream, error
370+
371+
// information about a federated upstream
372+
up, err := rmqc.GetFederationUpstream("/", "name")
373+
// => *FederationUpstream, error
374+
375+
// declare or update a federation upstream
376+
resp, err := rmqc.PutFederationUpstream("/", "name", FederationDefinition{
377+
Uri: "amqp://server-name",
378+
})
379+
// => *http.Response, error
380+
381+
// delete an upstream
382+
resp, err := rmqc.DeleteFederationUpstream("/", "name")
383+
// => *http.Response, error
384+
385+
```
386+
329387
### Operations on cluster name
330388
``` go
331389
// Get cluster name

bin/ci/before_build.bat

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,3 +25,7 @@ call %RABBITHOLE_RABBITMQCTL% set_permissions -p "rabbit/hole" guest ".*" ".*" "
2525
REM Enable shovel plugin
2626
call %RABBITHOLE_RABBITMQ_PLUGINS% enable rabbitmq_shovel
2727
call %RABBITHOLE_RABBITMQ_PLUGINS% enable rabbitmq_shovel_management
28+
29+
REM Enable shovel plugin
30+
call %RABBITHOLE_RABBITMQ_PLUGINS% enable rabbitmq_federation
31+
call %RABBITHOLE_RABBITMQ_PLUGINS% enable rabbitmq_federation_management

bin/ci/before_build.sh

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,3 +29,7 @@ $CTL set_cluster_name rabbitmq@localhost
2929
# Enable shovel plugin
3030
$PLUGINS enable rabbitmq_shovel
3131
$PLUGINS enable rabbitmq_shovel_management
32+
33+
# Enable federation plugin
34+
$PLUGINS enable rabbitmq_federation
35+
$PLUGINS enable rabbitmq_federation_management

doc.go

Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -199,6 +199,58 @@ Managing Topic Permissions
199199
resp, err := rmqc.DeleteTopicPermissionsIn("/", "my.user", "exchange")
200200
// => *http.Response, err
201201
202+
Managing Runtime Parameters
203+
204+
// list all runtime parameters
205+
params, err := rmqc.ListRuntimeParameters()
206+
// => []RuntimeParameter, error
207+
208+
// list all runtime parameters for a component
209+
params, err := rmqc.ListRuntimeParametersFor("federation-upstream")
210+
// => []RuntimeParameter, error
211+
212+
// list runtime parameters in a vhost
213+
params, err := rmqc.ListRuntimeParametersIn("federation-upstream", "/")
214+
// => []RuntimeParameter, error
215+
216+
// information about a runtime parameter
217+
p, err := rmqc.GetRuntimeParameter("federation-upstream", "/", "name")
218+
// => *RuntimeParameter, error
219+
220+
// declare or update a runtime parameter
221+
resp, err := rmqc.PutRuntimeParameter("federation-upstream", "/", "name", FederationDefinition{
222+
Uri: "amqp://server-name",
223+
})
224+
// => *http.Response, error
225+
226+
// remove a runtime parameter
227+
resp, err := rmqc.DeleteRuntimeParameter("federation-upstream", "/", "name")
228+
// => *http.Response, error
229+
230+
Managing Federation Upstreams
231+
232+
// list all federation upstreams
233+
ups, err := rmqc.ListFederationUpstreams()
234+
// => []FederationUpstream, error
235+
236+
// list federation upstreams in a vhost
237+
ups, err := rmqc.ListFederationUpstreamsIn("/")
238+
// => []FederationUpstream, error
239+
240+
// information about a federated upstream
241+
up, err := rmqc.GetFederationUpstream("/", "upstream-name")
242+
// => *FederationUpstream, error
243+
244+
// declare or update a federation upstream
245+
resp, err := rmqc.PutFederationUpstream("/", "upstream-name", FederationDefinition{
246+
Uri: "amqp://server-name",
247+
})
248+
// => *http.Response, error
249+
250+
// delete an upstream
251+
resp, err := rmqc.DeleteFederationUpstream("/", "upstream-name")
252+
// => *http.Response, error
253+
202254
Operations on cluster name
203255
// Get cluster name
204256
cn, err := rmqc.GetClusterName()

federation.go

Lines changed: 101 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,7 @@
11
package rabbithole
22

33
import (
4-
"encoding/json"
54
"net/http"
6-
"net/url"
75
)
86

97
// Federation definition: additional arguments
@@ -24,49 +22,132 @@ type FederationDefinition struct {
2422

2523
// Represents a configured Federation upstream.
2624
type FederationUpstream struct {
25+
Name string `json:"name"`
26+
Vhost string `json:"vhost"`
27+
Component string `json:"component"`
2728
Definition FederationDefinition `json:"value"`
2829
}
2930

31+
const FederationUpstreamComponent string = "federation-upstream"
32+
3033
//
31-
// PUT /api/parameters/federation-upstream/{vhost}/{upstream}
34+
// GET /api/parameters/federation-upstream
3235
//
3336

34-
// Updates a federation upstream
35-
func (c *Client) PutFederationUpstream(vhost string, upstreamName string, fDef FederationDefinition) (res *http.Response, err error) {
36-
fedUp := FederationUpstream{
37-
Definition: fDef,
38-
}
39-
body, err := json.Marshal(fedUp)
37+
// ListFederationUpstreams returns a list of all federation upstreams.
38+
func (c *Client) ListFederationUpstreams() (ups []FederationUpstream, err error) {
39+
params, err := c.ListRuntimeParametersFor(FederationUpstreamComponent)
4040
if err != nil {
4141
return nil, err
4242
}
4343

44-
req, err := newRequestWithBody(c, "PUT", "parameters/federation-upstream/"+url.PathEscape(vhost)+"/"+url.PathEscape(upstreamName), body)
44+
for _, p := range params {
45+
up := paramToUpstream(&p)
46+
ups = append(ups, *up)
47+
}
48+
return ups, nil
49+
}
50+
51+
//
52+
// GET /api/parameters/federation-upstream/{vhost}
53+
//
54+
55+
// ListFederationUpstreamsIn returns a list of all federation upstreams in a vhost.
56+
func (c *Client) ListFederationUpstreamsIn(vhost string) (ups []FederationUpstream, err error) {
57+
params, err := c.ListRuntimeParametersIn(FederationUpstreamComponent, vhost)
4558
if err != nil {
4659
return nil, err
4760
}
4861

49-
if res, err = executeRequest(c, req); err != nil {
62+
for _, p := range params {
63+
up := paramToUpstream(&p)
64+
ups = append(ups, *up)
65+
}
66+
return ups, nil
67+
}
68+
69+
//
70+
// GET /api/parameters/federation-upstream/{vhost}/{upstream}
71+
//
72+
73+
// GetFederationUpstream returns information about a federation upstream.
74+
func (c *Client) GetFederationUpstream(vhost, name string) (up *FederationUpstream, err error) {
75+
p, err := c.GetRuntimeParameter(FederationUpstreamComponent, vhost, name)
76+
if err != nil {
5077
return nil, err
5178
}
79+
return paramToUpstream(p), nil
80+
}
5281

53-
return res, nil
82+
//
83+
// PUT /api/parameters/federation-upstream/{vhost}/{upstream}
84+
//
85+
86+
// PutFederationUpstream creates or updates a federation upstream configuration.
87+
func (c *Client) PutFederationUpstream(vhost, name string, def FederationDefinition) (res *http.Response, err error) {
88+
return c.PutRuntimeParameter(FederationUpstreamComponent, vhost, name, def)
5489
}
5590

5691
//
5792
// DELETE /api/parameters/federation-upstream/{vhost}/{name}
5893
//
5994

60-
// Deletes a federation upstream.
61-
func (c *Client) DeleteFederationUpstream(vhost, upstreamName string) (res *http.Response, err error) {
62-
req, err := newRequestWithBody(c, "DELETE", "parameters/federation-upstream/"+url.PathEscape(vhost)+"/"+url.PathEscape(upstreamName), nil)
63-
if err != nil {
64-
return nil, err
95+
// DeleteFederationUpstream removes a federation upstream.
96+
func (c *Client) DeleteFederationUpstream(vhost, name string) (res *http.Response, err error) {
97+
return c.DeleteRuntimeParameter(FederationUpstreamComponent, vhost, name)
98+
}
99+
100+
// paramToUpstream maps from a RuntimeParameter structure to a FederationUpstream structure.
101+
func paramToUpstream(p *RuntimeParameter) (up *FederationUpstream) {
102+
up = &FederationUpstream{
103+
Name: p.Name,
104+
Vhost: p.Vhost,
105+
Component: p.Component,
65106
}
66107

67-
if res, err = executeRequest(c, req); err != nil {
68-
return nil, err
108+
def := FederationDefinition{}
109+
m := p.Value.(map[string]interface{})
110+
111+
if v, ok := m["uri"].(string); ok {
112+
def.Uri = v
113+
}
114+
115+
if v, ok := m["expires"].(float64); ok {
116+
def.Expires = int(v)
117+
}
118+
119+
if v, ok := m["message-ttl"].(float64); ok {
120+
def.MessageTTL = int32(v)
121+
}
122+
123+
if v, ok := m["max-hops"].(float64); ok {
124+
def.MaxHops = int(v)
125+
}
126+
127+
if v, ok := m["prefetch-count"].(float64); ok {
128+
def.PrefetchCount = int(v)
129+
}
130+
131+
if v, ok := m["reconnect-delay"].(float64); ok {
132+
def.ReconnectDelay = int(v)
133+
}
134+
135+
if v, ok := m["ack-mode"].(string); ok {
136+
def.AckMode = v
137+
}
138+
139+
if v, ok := m["trust-user-id"].(bool); ok {
140+
def.TrustUserId = v
141+
}
142+
143+
if v, ok := m["exchange"].(string); ok {
144+
def.Exchange = v
145+
}
146+
147+
if v, ok := m["queue"].(string); ok {
148+
def.Queue = v
69149
}
70150

71-
return res, nil
151+
up.Definition = def
152+
return up
72153
}

0 commit comments

Comments
 (0)