diff --git a/swarm/network/stream/common_test.go b/swarm/network/stream/common_test.go index cfbfa37b2d..f6e89d9010 100644 --- a/swarm/network/stream/common_test.go +++ b/swarm/network/stream/common_test.go @@ -78,7 +78,7 @@ func NewStreamerService(ctx *adapters.ServiceContext) (node.Service, error) { db := storage.NewDBAPI(store) delivery := NewDelivery(kad, db) deliveries[id] = delivery - r := NewRegistry(addr, delivery, db, state.NewMemStore(), defaultSkipCheck, false) + r := NewRegistry(addr, delivery, db, state.NewMemStore(), defaultSkipCheck, false, false) RegisterSwarmSyncerServer(r, db) RegisterSwarmSyncerClient(r, db) go func() { @@ -109,7 +109,7 @@ func newStreamerTester(t *testing.T) (*p2ptest.ProtocolTester, *Registry, *stora db := storage.NewDBAPI(localStore) delivery := NewDelivery(to, db) - streamer := NewRegistry(addr, delivery, db, state.NewMemStore(), defaultSkipCheck, false) + streamer := NewRegistry(addr, delivery, db, state.NewMemStore(), defaultSkipCheck, false, false) teardown := func() { streamer.Close() removeDataDir() diff --git a/swarm/network/stream/intervals_test.go b/swarm/network/stream/intervals_test.go index 2823b7e855..25cdfe917f 100644 --- a/swarm/network/stream/intervals_test.go +++ b/swarm/network/stream/intervals_test.go @@ -50,7 +50,7 @@ func newIntervalsStreamerService(ctx *adapters.ServiceContext) (node.Service, er db := storage.NewDBAPI(store) delivery := NewDelivery(kad, db) deliveries[id] = delivery - r := NewRegistry(addr, delivery, db, state.NewMemStore(), defaultSkipCheck, false) + r := NewRegistry(addr, delivery, db, state.NewMemStore(), defaultSkipCheck, false, false) r.RegisterClientFunc(externalStreamName, func(p *Peer, t []byte, live bool) (Client, error) { return newTestExternalClient(t, db), nil diff --git a/swarm/network/stream/stream.go b/swarm/network/stream/stream.go index a145b6f180..4bdedd9102 100644 --- a/swarm/network/stream/stream.go +++ b/swarm/network/stream/stream.go @@ -58,10 +58,11 @@ type Registry struct { delivery *Delivery intervalsStore state.Store doSync bool + doRetrieve bool } // NewRegistry is Streamer constructor -func NewRegistry(addr *network.BzzAddr, delivery *Delivery, db *storage.DBAPI, intervalsStore state.Store, skipCheck, doSync bool) *Registry { +func NewRegistry(addr *network.BzzAddr, delivery *Delivery, db *storage.DBAPI, intervalsStore state.Store, skipCheck, doSync, doRetrieve bool) *Registry { streamer := &Registry{ addr: addr, skipCheck: skipCheck, @@ -71,6 +72,7 @@ func NewRegistry(addr *network.BzzAddr, delivery *Delivery, db *storage.DBAPI, i delivery: delivery, intervalsStore: intervalsStore, doSync: doSync, + doRetrieve: doRetrieve, } streamer.api = NewAPI(streamer) delivery.getPeer = streamer.getPeer @@ -312,6 +314,12 @@ func (r *Registry) Run(p *network.BzzPeer) error { return true }) } + if r.doRetrieve { + err := r.Subscribe(p.ID(), NewStream(swarmChunkServerStreamName, nil, false), nil, Top) + if err != nil { + return err + } + } return sp.Run(sp.HandleMsg) } diff --git a/swarm/swarm.go b/swarm/swarm.go index dac54f5270..09d5a2bb79 100644 --- a/swarm/swarm.go +++ b/swarm/swarm.go @@ -153,7 +153,7 @@ func NewSwarm(ctx *node.ServiceContext, backend chequebook.Backend, config *api. if err != nil { return } - self.streamer = stream.NewRegistry(addr, delivery, db, stateStore, false, true) + self.streamer = stream.NewRegistry(addr, delivery, db, stateStore, false, true, true) self.bzz = network.NewBzz(bzzconfig, to, stateStore, stream.Spec, self.streamer.Run)