From ffa955a21b91a11a517dc0c0ef3b4567fea3541e Mon Sep 17 00:00:00 2001 From: Tom Wilkie Date: Mon, 4 Jan 2016 19:09:34 +0000 Subject: [PATCH] Split and move xfer package. --- app/controls.go | 2 +- app/controls_test.go | 7 +++-- app/pipes.go | 2 +- app/pipes_internal_test.go | 9 +++--- app/router.go | 2 +- common/xfer/constants.go | 19 +++++++++++++ {xfer => common/xfer}/controls.go | 0 {xfer => common/xfer}/pipes.go | 0 experimental/demoprobe/main.go | 7 +++-- experimental/fixprobe/main.go | 7 +++-- {xfer => probe/appclient}/app_client.go | 28 ++++++++----------- .../appclient}/app_client_internal_test.go | 7 +++-- {xfer => probe/appclient}/multi_client.go | 16 ++++++++--- .../appclient}/multi_client_internal_test.go | 2 +- .../appclient}/multi_client_test.go | 11 ++++---- {xfer => probe/appclient}/probe_config.go | 9 ++---- {xfer => probe/appclient}/report_publisher.go | 2 +- {xfer => probe/appclient}/resolver.go | 6 ++-- .../appclient}/resolver_internal_test.go | 12 ++++---- probe/controls/controls.go | 2 +- probe/controls/controls_test.go | 2 +- probe/controls/pipes.go | 2 +- probe/docker/controls.go | 2 +- probe/docker/controls_test.go | 2 +- probe/probe.go | 8 +++--- prog/app.go | 2 +- prog/probe.go | 11 ++++---- xfer/ports.go | 8 ------ xfer/publisher.go | 10 ------- 29 files changed, 104 insertions(+), 93 deletions(-) create mode 100644 common/xfer/constants.go rename {xfer => common/xfer}/controls.go (100%) rename {xfer => common/xfer}/pipes.go (100%) rename {xfer => probe/appclient}/app_client.go (92%) rename {xfer => probe/appclient}/app_client_internal_test.go (95%) rename {xfer => probe/appclient}/multi_client.go (92%) rename {xfer => probe/appclient}/multi_client_internal_test.go (97%) rename {xfer => probe/appclient}/multi_client_test.go (88%) rename {xfer => probe/appclient}/probe_config.go (82%) rename {xfer => probe/appclient}/report_publisher.go (97%) rename {xfer => probe/appclient}/resolver.go (95%) rename {xfer => probe/appclient}/resolver_internal_test.go (89%) delete mode 100644 xfer/ports.go delete mode 100644 xfer/publisher.go diff --git a/app/controls.go b/app/controls.go index 0a77b2e1d..354afd703 100644 --- a/app/controls.go +++ b/app/controls.go @@ -9,7 +9,7 @@ import ( "github.com/gorilla/mux" - "github.com/weaveworks/scope/xfer" + "github.com/weaveworks/scope/common/xfer" ) // RegisterControlRoutes registers the various control routes with a http mux. diff --git a/app/controls_test.go b/app/controls_test.go index 82e51bf64..679f73a22 100644 --- a/app/controls_test.go +++ b/app/controls_test.go @@ -12,7 +12,8 @@ import ( "github.com/gorilla/mux" "github.com/weaveworks/scope/app" - "github.com/weaveworks/scope/xfer" + "github.com/weaveworks/scope/common/xfer" + "github.com/weaveworks/scope/probe/appclient" ) func TestControl(t *testing.T) { @@ -26,7 +27,7 @@ func TestControl(t *testing.T) { t.Fatal(err) } - probeConfig := xfer.ProbeConfig{ + probeConfig := appclient.ProbeConfig{ ProbeID: "foo", } controlHandler := xfer.ControlHandlerFunc(func(req xfer.Request) xfer.Response { @@ -42,7 +43,7 @@ func TestControl(t *testing.T) { Value: "foo", } }) - client, err := xfer.NewAppClient(probeConfig, ip+":"+port, ip+":"+port, controlHandler) + client, err := appclient.NewAppClient(probeConfig, ip+":"+port, ip+":"+port, controlHandler) if err != nil { t.Fatal(err) } diff --git a/app/pipes.go b/app/pipes.go index dc6a1b69e..f3dcadbd8 100644 --- a/app/pipes.go +++ b/app/pipes.go @@ -10,7 +10,7 @@ import ( "github.com/gorilla/mux" "github.com/weaveworks/scope/common/mtime" - "github.com/weaveworks/scope/xfer" + "github.com/weaveworks/scope/common/xfer" ) const ( diff --git a/app/pipes_internal_test.go b/app/pipes_internal_test.go index 312e02bbe..69615c5ac 100644 --- a/app/pipes_internal_test.go +++ b/app/pipes_internal_test.go @@ -14,9 +14,10 @@ import ( "github.com/gorilla/websocket" "github.com/weaveworks/scope/common/mtime" + "github.com/weaveworks/scope/common/xfer" + "github.com/weaveworks/scope/probe/appclient" "github.com/weaveworks/scope/probe/controls" "github.com/weaveworks/scope/test" - "github.com/weaveworks/scope/xfer" ) func TestPipeTimeout(t *testing.T) { @@ -50,7 +51,7 @@ func TestPipeTimeout(t *testing.T) { } type adapter struct { - c xfer.AppClient + c appclient.AppClient } func (a adapter) PipeConnection(_, pipeID string, pipe xfer.Pipe) error { @@ -75,10 +76,10 @@ func TestPipeClose(t *testing.T) { t.Fatal(err) } - probeConfig := xfer.ProbeConfig{ + probeConfig := appclient.ProbeConfig{ ProbeID: "foo", } - client, err := xfer.NewAppClient(probeConfig, ip+":"+port, ip+":"+port, nil) + client, err := appclient.NewAppClient(probeConfig, ip+":"+port, ip+":"+port, nil) if err != nil { t.Fatal(err) } diff --git a/app/router.go b/app/router.go index 36dd835fd..31456ea40 100644 --- a/app/router.go +++ b/app/router.go @@ -12,8 +12,8 @@ import ( "github.com/gorilla/mux" "github.com/weaveworks/scope/common/hostname" + "github.com/weaveworks/scope/common/xfer" "github.com/weaveworks/scope/report" - "github.com/weaveworks/scope/xfer" ) var ( diff --git a/common/xfer/constants.go b/common/xfer/constants.go new file mode 100644 index 000000000..6bbc0d7c6 --- /dev/null +++ b/common/xfer/constants.go @@ -0,0 +1,19 @@ +package xfer + +const ( + // AppPort is the default port that the app will use for its HTTP server. + // The app publishes the API and user interface, and receives reports from + // probes, on this port. + AppPort = 4040 + + // ScopeProbeIDHeader is the header we use to carry the probe's unique ID. The + // ID is currently set to the a random string on probe startup. + ScopeProbeIDHeader = "X-Scope-Probe-ID" +) + +// Details are some generic details that can be fetched from /api +type Details struct { + ID string `json:"id"` + Version string `json:"version"` + Hostname string `json:"hostname"` +} diff --git a/xfer/controls.go b/common/xfer/controls.go similarity index 100% rename from xfer/controls.go rename to common/xfer/controls.go diff --git a/xfer/pipes.go b/common/xfer/pipes.go similarity index 100% rename from xfer/pipes.go rename to common/xfer/pipes.go diff --git a/experimental/demoprobe/main.go b/experimental/demoprobe/main.go index 57dbc8e4c..792b5df89 100644 --- a/experimental/demoprobe/main.go +++ b/experimental/demoprobe/main.go @@ -9,10 +9,11 @@ import ( "strconv" "time" + "github.com/weaveworks/scope/common/xfer" + "github.com/weaveworks/scope/probe/appclient" "github.com/weaveworks/scope/probe/docker" "github.com/weaveworks/scope/probe/process" "github.com/weaveworks/scope/report" - "github.com/weaveworks/scope/xfer" ) func main() { @@ -23,7 +24,7 @@ func main() { ) flag.Parse() - client, err := xfer.NewAppClient(xfer.ProbeConfig{ + client, err := appclient.NewAppClient(appclient.ProbeConfig{ Token: "demoprobe", ProbeID: "demoprobe", Insecure: false, @@ -31,7 +32,7 @@ func main() { if err != nil { log.Fatal(err) } - rp := xfer.NewReportPublisher(client) + rp := appclient.NewReportPublisher(client) rand.Seed(time.Now().UnixNano()) for range time.Tick(*publishInterval) { diff --git a/experimental/fixprobe/main.go b/experimental/fixprobe/main.go index d19eb99b4..d5c29b2ad 100644 --- a/experimental/fixprobe/main.go +++ b/experimental/fixprobe/main.go @@ -9,8 +9,9 @@ import ( "os" "time" + "github.com/weaveworks/scope/common/xfer" + "github.com/weaveworks/scope/probe/appclient" "github.com/weaveworks/scope/report" - "github.com/weaveworks/scope/xfer" ) func main() { @@ -34,7 +35,7 @@ func main() { } f.Close() - client, err := xfer.NewAppClient(xfer.ProbeConfig{ + client, err := appclient.NewAppClient(appclient.ProbeConfig{ Token: "fixprobe", ProbeID: "fixprobe", Insecure: false, @@ -43,7 +44,7 @@ func main() { log.Fatal(err) } - rp := xfer.NewReportPublisher(client) + rp := appclient.NewReportPublisher(client) for range time.Tick(*publishInterval) { rp.Publish(fixedReport) } diff --git a/xfer/app_client.go b/probe/appclient/app_client.go similarity index 92% rename from xfer/app_client.go rename to probe/appclient/app_client.go index 379b93604..43cac073d 100644 --- a/xfer/app_client.go +++ b/probe/appclient/app_client.go @@ -1,4 +1,4 @@ -package xfer +package appclient import ( "encoding/json" @@ -13,6 +13,7 @@ import ( "github.com/gorilla/websocket" "github.com/weaveworks/scope/common/sanitize" + "github.com/weaveworks/scope/common/xfer" ) const ( @@ -20,18 +21,11 @@ const ( maxBackoff = 60 * time.Second ) -// Details are some generic details that can be fetched from /api -type Details struct { - ID string `json:"id"` - Version string `json:"version"` - Hostname string `json:"hostname"` -} - // AppClient is a client to an app for dealing with controls. type AppClient interface { - Details() (Details, error) + Details() (xfer.Details, error) ControlConnection() - PipeConnection(string, Pipe) + PipeConnection(string, xfer.Pipe) PipeClose(string) error Publish(r io.Reader) error Stop() @@ -58,11 +52,11 @@ type appClient struct { readers chan io.Reader // For controls - control ControlHandler + control xfer.ControlHandler } // NewAppClient makes a new appClient. -func NewAppClient(pc ProbeConfig, hostname, target string, control ControlHandler) (AppClient, error) { +func NewAppClient(pc ProbeConfig, hostname, target string, control xfer.ControlHandler) (AppClient, error) { httpTransport, err := pc.getHTTPTransport(hostname) if err != nil { return nil, err @@ -144,8 +138,8 @@ func (c *appClient) Stop() { } // Details fetches the details (version, id) of the app. -func (c *appClient) Details() (Details, error) { - result := Details{} +func (c *appClient) Details() (xfer.Details, error) { + result := xfer.Details{} req, err := c.ProbeConfig.authorizedRequest("GET", sanitize.URL("", 0, "/api")(c.target), nil) if err != nil { return result, err @@ -202,7 +196,7 @@ func (c *appClient) controlConnection() (bool, error) { conn.Close() }() - codec := NewJSONWebsocketCodec(conn) + codec := xfer.NewJSONWebsocketCodec(conn) server := rpc.NewServer() if err := server.RegisterName("control", c.control); err != nil { return false, err @@ -271,7 +265,7 @@ func (c *appClient) Publish(r io.Reader) error { return nil } -func (c *appClient) pipeConnection(id string, pipe Pipe) (bool, error) { +func (c *appClient) pipeConnection(id string, pipe xfer.Pipe) (bool, error) { headers := http.Header{} c.ProbeConfig.authorizeHeaders(headers) url := sanitize.URL("ws://", 0, fmt.Sprintf("/api/pipe/%s/probe", id))(c.target) @@ -295,7 +289,7 @@ func (c *appClient) pipeConnection(id string, pipe Pipe) (bool, error) { return false, pipe.CopyToWebsocket(remote, conn) } -func (c *appClient) PipeConnection(id string, pipe Pipe) { +func (c *appClient) PipeConnection(id string, pipe xfer.Pipe) { go func() { log.Printf("Pipe %s connection to %s starting", id, c.target) defer log.Printf("Pipe %s connection to %s exiting", id, c.target) diff --git a/xfer/app_client_internal_test.go b/probe/appclient/app_client_internal_test.go similarity index 95% rename from xfer/app_client_internal_test.go rename to probe/appclient/app_client_internal_test.go index d13b1d7cc..bd316c45d 100644 --- a/xfer/app_client_internal_test.go +++ b/probe/appclient/app_client_internal_test.go @@ -1,4 +1,4 @@ -package xfer +package appclient import ( "compress/gzip" @@ -15,6 +15,7 @@ import ( "time" "github.com/gorilla/handlers" + "github.com/weaveworks/scope/common/xfer" "github.com/weaveworks/scope/report" "github.com/weaveworks/scope/test" ) @@ -33,7 +34,7 @@ func dummyServer(t *testing.T, expectedToken, expectedID string, expectedReport t.Errorf("want %q, have %q", expectedToken, have) } - if have := r.Header.Get(ScopeProbeIDHeader); expectedID != have { + if have := r.Header.Get(xfer.ScopeProbeIDHeader); expectedID != have { t.Errorf("want %q, have %q", expectedID, have) } @@ -151,7 +152,7 @@ func TestAppClientPublish(t *testing.T) { var ( id = "foobarbaz" version = "imalittleteapot" - want = Details{ID: id, Version: version} + want = xfer.Details{ID: id, Version: version} ) handler := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { diff --git a/xfer/multi_client.go b/probe/appclient/multi_client.go similarity index 92% rename from xfer/multi_client.go rename to probe/appclient/multi_client.go index d4eceb9fc..7ee020bba 100644 --- a/xfer/multi_client.go +++ b/probe/appclient/multi_client.go @@ -1,4 +1,4 @@ -package xfer +package appclient import ( "bytes" @@ -10,6 +10,7 @@ import ( "strings" "sync" + "github.com/weaveworks/scope/common/xfer" "github.com/weaveworks/scope/report" ) @@ -29,15 +30,22 @@ type multiClient struct { } type clientTuple struct { - Details + xfer.Details AppClient } +// Publisher is something which can send a stream of data somewhere, probably +// to a remote collector. +type Publisher interface { + Publish(io.Reader) error + Stop() +} + // MultiAppClient maintains a set of upstream apps, and ensures we have an // AppClient for each one. type MultiAppClient interface { Set(hostname string, endpoints []string) - PipeConnection(appID, pipeID string, pipe Pipe) error + PipeConnection(appID, pipeID string, pipe xfer.Pipe) error PipeClose(appID, pipeID string) error Stop() Publish(io.Reader) error @@ -122,7 +130,7 @@ func (c *multiClient) withClient(appID string, f func(AppClient) error) error { return f(client) } -func (c *multiClient) PipeConnection(appID, pipeID string, pipe Pipe) error { +func (c *multiClient) PipeConnection(appID, pipeID string, pipe xfer.Pipe) error { return c.withClient(appID, func(client AppClient) error { client.PipeConnection(pipeID, pipe) return nil diff --git a/xfer/multi_client_internal_test.go b/probe/appclient/multi_client_internal_test.go similarity index 97% rename from xfer/multi_client_internal_test.go rename to probe/appclient/multi_client_internal_test.go index 66523adc4..0b9683f91 100644 --- a/xfer/multi_client_internal_test.go +++ b/probe/appclient/multi_client_internal_test.go @@ -1,4 +1,4 @@ -package xfer +package appclient import ( "testing" diff --git a/xfer/multi_client_test.go b/probe/appclient/multi_client_test.go similarity index 88% rename from xfer/multi_client_test.go rename to probe/appclient/multi_client_test.go index b4e495503..98126c1cd 100644 --- a/xfer/multi_client_test.go +++ b/probe/appclient/multi_client_test.go @@ -1,4 +1,4 @@ -package xfer_test +package appclient_test import ( "bytes" @@ -6,7 +6,8 @@ import ( "runtime" "testing" - "github.com/weaveworks/scope/xfer" + "github.com/weaveworks/scope/common/xfer" + "github.com/weaveworks/scope/probe/appclient" ) type mockClient struct { @@ -41,7 +42,7 @@ var ( a2 = &mockClient{id: "2"} // hostname a, app id 2 b2 = &mockClient{id: "2"} // hostname b, app id 2 (duplicate) b3 = &mockClient{id: "3"} // hostname b, app id 3 - factory = func(hostname, target string) (xfer.AppClient, error) { + factory = func(hostname, target string) (appclient.AppClient, error) { switch target { case "a1": return a1, nil @@ -66,7 +67,7 @@ func TestMultiClient(t *testing.T) { } ) - mp := xfer.NewMultiAppClient(factory) + mp := appclient.NewMultiAppClient(factory) defer mp.Stop() // Add two hostnames with overlapping apps, check we don't add the same app twice @@ -88,7 +89,7 @@ func TestMultiClient(t *testing.T) { } func TestMultiClientPublish(t *testing.T) { - mp := xfer.NewMultiAppClient(factory) + mp := appclient.NewMultiAppClient(factory) defer mp.Stop() sum := func() int { return a1.publish + a2.publish + b2.publish + b3.publish } diff --git a/xfer/probe_config.go b/probe/appclient/probe_config.go similarity index 82% rename from xfer/probe_config.go rename to probe/appclient/probe_config.go index 58c1093d4..c0fcaf2eb 100644 --- a/xfer/probe_config.go +++ b/probe/appclient/probe_config.go @@ -1,4 +1,4 @@ -package xfer +package appclient import ( "crypto/tls" @@ -9,12 +9,9 @@ import ( "net/http" "github.com/certifi/gocertifi" + "github.com/weaveworks/scope/common/xfer" ) -// ScopeProbeIDHeader is the header we use to carry the probe's unique ID. The -// ID is currently set to the a random string on probe startup. -const ScopeProbeIDHeader = "X-Scope-Probe-ID" - var certPool *x509.CertPool func init() { @@ -34,7 +31,7 @@ type ProbeConfig struct { func (pc ProbeConfig) authorizeHeaders(headers http.Header) { headers.Set("Authorization", fmt.Sprintf("Scope-Probe token=%s", pc.Token)) - headers.Set(ScopeProbeIDHeader, pc.ProbeID) + headers.Set(xfer.ScopeProbeIDHeader, pc.ProbeID) } func (pc ProbeConfig) authorizedRequest(method string, urlStr string, body io.Reader) (*http.Request, error) { diff --git a/xfer/report_publisher.go b/probe/appclient/report_publisher.go similarity index 97% rename from xfer/report_publisher.go rename to probe/appclient/report_publisher.go index efde88da4..83098491a 100644 --- a/xfer/report_publisher.go +++ b/probe/appclient/report_publisher.go @@ -1,4 +1,4 @@ -package xfer +package appclient import ( "bytes" diff --git a/xfer/resolver.go b/probe/appclient/resolver.go similarity index 95% rename from xfer/resolver.go rename to probe/appclient/resolver.go index c2524a8f1..f37909ada 100644 --- a/xfer/resolver.go +++ b/probe/appclient/resolver.go @@ -1,4 +1,4 @@ -package xfer +package appclient import ( "log" @@ -6,6 +6,8 @@ import ( "strconv" "strings" "time" + + "github.com/weaveworks/scope/common/xfer" ) const ( @@ -76,7 +78,7 @@ func prepare(strs []string) []target { continue } } else { - host, port = s, strconv.Itoa(AppPort) + host, port = s, strconv.Itoa(xfer.AppPort) } targets = append(targets, target{host, port}) } diff --git a/xfer/resolver_internal_test.go b/probe/appclient/resolver_internal_test.go similarity index 89% rename from xfer/resolver_internal_test.go rename to probe/appclient/resolver_internal_test.go index c1d57b948..f86abfd00 100644 --- a/xfer/resolver_internal_test.go +++ b/probe/appclient/resolver_internal_test.go @@ -1,4 +1,4 @@ -package xfer +package appclient import ( "fmt" @@ -7,6 +7,8 @@ import ( "sync" "testing" "time" + + "github.com/weaveworks/scope/common/xfer" ) func TestResolver(t *testing.T) { @@ -68,22 +70,22 @@ func TestResolver(t *testing.T) { } // Initial resolve should just give us IPs - assertAdd(ip1+port, fmt.Sprintf("%s:%d", ip2, AppPort)) + assertAdd(ip1+port, fmt.Sprintf("%s:%d", ip2, xfer.AppPort)) // Trigger another resolve with a tick; again, // just want ips. c <- time.Now() - assertAdd(ip1+port, fmt.Sprintf("%s:%d", ip2, AppPort)) + assertAdd(ip1+port, fmt.Sprintf("%s:%d", ip2, xfer.AppPort)) ip3 := "1.2.3.4" updateIPs("symbolic.name", makeIPs(ip3)) c <- time.Now() // trigger a resolve - assertAdd(ip3+port, ip1+port, fmt.Sprintf("%s:%d", ip2, AppPort)) + assertAdd(ip3+port, ip1+port, fmt.Sprintf("%s:%d", ip2, xfer.AppPort)) ip4 := "10.10.10.10" updateIPs("symbolic.name", makeIPs(ip3, ip4)) c <- time.Now() // trigger another resolve, this time with 2 adds - assertAdd(ip3+port, ip4+port, ip1+port, fmt.Sprintf("%s:%d", ip2, AppPort)) + assertAdd(ip3+port, ip4+port, ip1+port, fmt.Sprintf("%s:%d", ip2, xfer.AppPort)) done := make(chan struct{}) go func() { r.Stop(); close(done) }() diff --git a/probe/controls/controls.go b/probe/controls/controls.go index ed2f61bd6..2d12daa69 100644 --- a/probe/controls/controls.go +++ b/probe/controls/controls.go @@ -3,7 +3,7 @@ package controls import ( "sync" - "github.com/weaveworks/scope/xfer" + "github.com/weaveworks/scope/common/xfer" ) var ( diff --git a/probe/controls/controls_test.go b/probe/controls/controls_test.go index 04bb66268..b7c153a6c 100644 --- a/probe/controls/controls_test.go +++ b/probe/controls/controls_test.go @@ -4,9 +4,9 @@ import ( "reflect" "testing" + "github.com/weaveworks/scope/common/xfer" "github.com/weaveworks/scope/probe/controls" "github.com/weaveworks/scope/test" - "github.com/weaveworks/scope/xfer" ) func TestControls(t *testing.T) { diff --git a/probe/controls/pipes.go b/probe/controls/pipes.go index eec4da28f..cf540419f 100644 --- a/probe/controls/pipes.go +++ b/probe/controls/pipes.go @@ -4,7 +4,7 @@ import ( "fmt" "math/rand" - "github.com/weaveworks/scope/xfer" + "github.com/weaveworks/scope/common/xfer" ) // PipeClient is the type of the thing the probe uses to make pipe connections. diff --git a/probe/docker/controls.go b/probe/docker/controls.go index fdb63a19a..2057069e4 100644 --- a/probe/docker/controls.go +++ b/probe/docker/controls.go @@ -5,9 +5,9 @@ import ( docker_client "github.com/fsouza/go-dockerclient" + "github.com/weaveworks/scope/common/xfer" "github.com/weaveworks/scope/probe/controls" "github.com/weaveworks/scope/report" - "github.com/weaveworks/scope/xfer" ) // Control IDs used by the docker intergation. diff --git a/probe/docker/controls_test.go b/probe/docker/controls_test.go index 67a658958..97ab285f6 100644 --- a/probe/docker/controls_test.go +++ b/probe/docker/controls_test.go @@ -8,11 +8,11 @@ import ( "github.com/gorilla/websocket" + "github.com/weaveworks/scope/common/xfer" "github.com/weaveworks/scope/probe/controls" "github.com/weaveworks/scope/probe/docker" "github.com/weaveworks/scope/report" "github.com/weaveworks/scope/test" - "github.com/weaveworks/scope/xfer" ) func TestControls(t *testing.T) { diff --git a/probe/probe.go b/probe/probe.go index c28ac688b..bf244a7b0 100644 --- a/probe/probe.go +++ b/probe/probe.go @@ -7,8 +7,8 @@ import ( "github.com/armon/go-metrics" + "github.com/weaveworks/scope/probe/appclient" "github.com/weaveworks/scope/report" - "github.com/weaveworks/scope/xfer" ) const ( @@ -18,7 +18,7 @@ const ( // Probe sits there, generating and publishing reports. type Probe struct { spyInterval, publishInterval time.Duration - publisher *xfer.ReportPublisher + publisher *appclient.ReportPublisher tickers []Ticker reporters []Reporter @@ -52,11 +52,11 @@ type Ticker interface { } // New makes a new Probe. -func New(spyInterval, publishInterval time.Duration, publisher xfer.Publisher) *Probe { +func New(spyInterval, publishInterval time.Duration, publisher appclient.Publisher) *Probe { result := &Probe{ spyInterval: spyInterval, publishInterval: publishInterval, - publisher: xfer.NewReportPublisher(publisher), + publisher: appclient.NewReportPublisher(publisher), quit: make(chan struct{}), spiedReports: make(chan report.Report, reportBufferSize), shortcutReports: make(chan report.Report, reportBufferSize), diff --git a/prog/app.go b/prog/app.go index f7986e069..8d8a15090 100644 --- a/prog/app.go +++ b/prog/app.go @@ -14,7 +14,7 @@ import ( "github.com/weaveworks/weave/common" "github.com/weaveworks/scope/app" - "github.com/weaveworks/scope/xfer" + "github.com/weaveworks/scope/common/xfer" ) // Router creates the mux for all the various app components. diff --git a/prog/probe.go b/prog/probe.go index e6f46536d..d7dba92e5 100644 --- a/prog/probe.go +++ b/prog/probe.go @@ -17,7 +17,9 @@ import ( "github.com/weaveworks/weave/common" "github.com/weaveworks/scope/common/hostname" + "github.com/weaveworks/scope/common/xfer" "github.com/weaveworks/scope/probe" + "github.com/weaveworks/scope/probe/appclient" "github.com/weaveworks/scope/probe/controls" "github.com/weaveworks/scope/probe/docker" "github.com/weaveworks/scope/probe/endpoint" @@ -26,7 +28,6 @@ import ( "github.com/weaveworks/scope/probe/overlay" "github.com/weaveworks/scope/probe/process" "github.com/weaveworks/scope/report" - "github.com/weaveworks/scope/xfer" ) // Main runs the probe @@ -94,20 +95,20 @@ func probeMain() { } log.Printf("publishing to: %s", strings.Join(targets, ", ")) - probeConfig := xfer.ProbeConfig{ + probeConfig := appclient.ProbeConfig{ Token: *token, ProbeID: probeID, Insecure: *insecure, } - clients := xfer.NewMultiAppClient(func(hostname, endpoint string) (xfer.AppClient, error) { - return xfer.NewAppClient( + clients := appclient.NewMultiAppClient(func(hostname, endpoint string) (appclient.AppClient, error) { + return appclient.NewAppClient( probeConfig, hostname, endpoint, xfer.ControlHandlerFunc(controls.HandleControlRequest), ) }) defer clients.Stop() - resolver := xfer.NewStaticResolver(targets, clients.Set) + resolver := appclient.NewStaticResolver(targets, clients.Set) defer resolver.Stop() processCache := process.NewCachingWalker(process.NewWalker(*procRoot)) diff --git a/xfer/ports.go b/xfer/ports.go deleted file mode 100644 index 11e1e074a..000000000 --- a/xfer/ports.go +++ /dev/null @@ -1,8 +0,0 @@ -package xfer - -var ( - // AppPort is the default port that the app will use for its HTTP server. - // The app publishes the API and user interface, and receives reports from - // probes, on this port. - AppPort = 4040 -) diff --git a/xfer/publisher.go b/xfer/publisher.go deleted file mode 100644 index 7ac56a93e..000000000 --- a/xfer/publisher.go +++ /dev/null @@ -1,10 +0,0 @@ -package xfer - -import "io" - -// Publisher is something which can send a stream of data somewhere, probably -// to a remote collector. -type Publisher interface { - Publish(io.Reader) error - Stop() -}