mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-17 17:33:47 +00:00
swarm/network/stream: Subscribe to stream to retrieve chunks
This commit is contained in:
parent
b777b0dff8
commit
746d1d6a54
4 changed files with 13 additions and 5 deletions
|
|
@ -78,7 +78,7 @@ func NewStreamerService(ctx *adapters.ServiceContext) (node.Service, error) {
|
||||||
db := storage.NewDBAPI(store)
|
db := storage.NewDBAPI(store)
|
||||||
delivery := NewDelivery(kad, db)
|
delivery := NewDelivery(kad, db)
|
||||||
deliveries[id] = delivery
|
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)
|
RegisterSwarmSyncerServer(r, db)
|
||||||
RegisterSwarmSyncerClient(r, db)
|
RegisterSwarmSyncerClient(r, db)
|
||||||
go func() {
|
go func() {
|
||||||
|
|
@ -109,7 +109,7 @@ func newStreamerTester(t *testing.T) (*p2ptest.ProtocolTester, *Registry, *stora
|
||||||
|
|
||||||
db := storage.NewDBAPI(localStore)
|
db := storage.NewDBAPI(localStore)
|
||||||
delivery := NewDelivery(to, db)
|
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() {
|
teardown := func() {
|
||||||
streamer.Close()
|
streamer.Close()
|
||||||
removeDataDir()
|
removeDataDir()
|
||||||
|
|
|
||||||
|
|
@ -50,7 +50,7 @@ func newIntervalsStreamerService(ctx *adapters.ServiceContext) (node.Service, er
|
||||||
db := storage.NewDBAPI(store)
|
db := storage.NewDBAPI(store)
|
||||||
delivery := NewDelivery(kad, db)
|
delivery := NewDelivery(kad, db)
|
||||||
deliveries[id] = delivery
|
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) {
|
r.RegisterClientFunc(externalStreamName, func(p *Peer, t []byte, live bool) (Client, error) {
|
||||||
return newTestExternalClient(t, db), nil
|
return newTestExternalClient(t, db), nil
|
||||||
|
|
|
||||||
|
|
@ -58,10 +58,11 @@ type Registry struct {
|
||||||
delivery *Delivery
|
delivery *Delivery
|
||||||
intervalsStore state.Store
|
intervalsStore state.Store
|
||||||
doSync bool
|
doSync bool
|
||||||
|
doRetrieve bool
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewRegistry is Streamer constructor
|
// 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{
|
streamer := &Registry{
|
||||||
addr: addr,
|
addr: addr,
|
||||||
skipCheck: skipCheck,
|
skipCheck: skipCheck,
|
||||||
|
|
@ -71,6 +72,7 @@ func NewRegistry(addr *network.BzzAddr, delivery *Delivery, db *storage.DBAPI, i
|
||||||
delivery: delivery,
|
delivery: delivery,
|
||||||
intervalsStore: intervalsStore,
|
intervalsStore: intervalsStore,
|
||||||
doSync: doSync,
|
doSync: doSync,
|
||||||
|
doRetrieve: doRetrieve,
|
||||||
}
|
}
|
||||||
streamer.api = NewAPI(streamer)
|
streamer.api = NewAPI(streamer)
|
||||||
delivery.getPeer = streamer.getPeer
|
delivery.getPeer = streamer.getPeer
|
||||||
|
|
@ -312,6 +314,12 @@ func (r *Registry) Run(p *network.BzzPeer) error {
|
||||||
return true
|
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)
|
return sp.Run(sp.HandleMsg)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -153,7 +153,7 @@ func NewSwarm(ctx *node.ServiceContext, backend chequebook.Backend, config *api.
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return
|
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)
|
self.bzz = network.NewBzz(bzzconfig, to, stateStore, stream.Spec, self.streamer.Run)
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue