providers.go 7.25 KB
Newer Older
1 2 3
package providers

import (
Jeromy's avatar
Jeromy committed
4
	"context"
5 6 7 8 9
	"encoding/binary"
	"fmt"
	"strings"
	"time"

10
	lru "github.com/hashicorp/golang-lru"
11
	cid "github.com/ipfs/go-cid"
12
	ds "github.com/ipfs/go-datastore"
Hector Sanjuan's avatar
Hector Sanjuan committed
13
	autobatch "github.com/ipfs/go-datastore/autobatch"
14 15 16 17
	dsq "github.com/ipfs/go-datastore/query"
	logging "github.com/ipfs/go-log"
	goprocess "github.com/jbenet/goprocess"
	goprocessctx "github.com/jbenet/goprocess/context"
18
	peer "github.com/libp2p/go-libp2p-peer"
19
	base32 "github.com/whyrusleeping/base32"
20 21
)

22 23
var batchBufferSize = 256

24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48
var log = logging.Logger("providers")

var lruCacheSize = 256
var ProvideValidity = time.Hour * 24
var defaultCleanupInterval = time.Hour

type ProviderManager struct {
	// all non channel fields are meant to be accessed only within
	// the run method
	providers *lru.Cache
	dstore    ds.Datastore

	newprovs chan *addProv
	getprovs chan *getProv
	proc     goprocess.Process

	cleanupInterval time.Duration
}

type providerSet struct {
	providers []peer.ID
	set       map[peer.ID]time.Time
}

type addProv struct {
49
	k   cid.Cid
50 51 52 53
	val peer.ID
}

type getProv struct {
54
	k    cid.Cid
55 56 57
	resp chan []peer.ID
}

58
func NewProviderManager(ctx context.Context, local peer.ID, dstore ds.Batching) *ProviderManager {
59 60 61
	pm := new(ProviderManager)
	pm.getprovs = make(chan *getProv)
	pm.newprovs = make(chan *addProv)
62
	pm.dstore = autobatch.NewAutoBatching(dstore, batchBufferSize)
63 64 65 66 67 68 69 70 71 72 73 74 75 76 77
	cache, err := lru.New(lruCacheSize)
	if err != nil {
		panic(err) //only happens if negative value is passed to lru constructor
	}
	pm.providers = cache

	pm.proc = goprocessctx.WithContext(ctx)
	pm.cleanupInterval = defaultCleanupInterval
	pm.proc.Go(func(p goprocess.Process) { pm.run() })

	return pm
}

const providersKeyPrefix = "/providers/"

78
func mkProvKey(k cid.Cid) string {
79
	return providersKeyPrefix + base32.RawStdEncoding.EncodeToString(k.Bytes())
80 81 82 83 84 85
}

func (pm *ProviderManager) Process() goprocess.Process {
	return pm.proc
}

86
func (pm *ProviderManager) providersForKey(k cid.Cid) ([]peer.ID, error) {
87 88 89 90 91 92 93
	pset, err := pm.getProvSet(k)
	if err != nil {
		return nil, err
	}
	return pset.providers, nil
}

94
func (pm *ProviderManager) getProvSet(k cid.Cid) (*providerSet, error) {
95
	cached, ok := pm.providers.Get(k)
96 97 98 99 100 101 102 103 104 105
	if ok {
		return cached.(*providerSet), nil
	}

	pset, err := loadProvSet(pm.dstore, k)
	if err != nil {
		return nil, err
	}

	if len(pset.providers) > 0 {
106
		pm.providers.Add(k, pset)
107 108 109 110 111
	}

	return pset, nil
}

112
func loadProvSet(dstore ds.Datastore, k cid.Cid) (*providerSet, error) {
113
	res, err := dstore.Query(dsq.Query{Prefix: mkProvKey(k)})
114 115 116 117 118
	if err != nil {
		return nil, err
	}

	out := newProviderSet()
119 120 121 122 123
	for {
		e, ok := res.NextSync()
		if !ok {
			break
		}
124
		if e.Error != nil {
125
			log.Error("got an error: ", e.Error)
126 127 128
			continue
		}

129 130 131
		lix := strings.LastIndex(e.Key, "/")

		decstr, err := base32.RawStdEncoding.DecodeString(e.Key[lix+1:])
132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150
		if err != nil {
			log.Error("base32 decoding error: ", err)
			continue
		}

		pid := peer.ID(decstr)

		t, err := readTimeValue(e.Value)
		if err != nil {
			log.Warning("parsing providers record from disk: ", err)
			continue
		}

		out.setVal(pid, t)
	}

	return out, nil
}

151 152 153 154
func readTimeValue(data []byte) (time.Time, error) {
	nsec, n := binary.Varint(data)
	if n <= 0 {
		return time.Time{}, fmt.Errorf("failed to parse time")
155 156 157 158 159
	}

	return time.Unix(0, nsec), nil
}

160
func (pm *ProviderManager) addProv(k cid.Cid, p peer.ID) error {
161
	now := time.Now()
162 163 164
	if provs, ok := pm.providers.Get(k); ok {
		provs.(*providerSet).setVal(p, now)
	} // else not cached, just write through
165 166 167 168

	return writeProviderEntry(pm.dstore, k, p, now)
}

169
func writeProviderEntry(dstore ds.Datastore, k cid.Cid, p peer.ID, t time.Time) error {
170
	dsk := mkProvKey(k) + "/" + base32.RawStdEncoding.EncodeToString([]byte(p))
171 172 173 174

	buf := make([]byte, 16)
	n := binary.PutVarint(buf, t.UnixNano())

175
	return dstore.Put(ds.NewKey(dsk), buf[:n])
176 177
}

178
func (pm *ProviderManager) deleteProvSet(k cid.Cid) error {
179
	pm.providers.Remove(k)
180 181 182

	res, err := pm.dstore.Query(dsq.Query{
		KeysOnly: true,
183
		Prefix:   mkProvKey(k),
184
	})
185 186 187
	if err != nil {
		return err
	}
188 189 190 191 192 193 194

	entries, err := res.Rest()
	if err != nil {
		return err
	}

	for _, e := range entries {
195
		err := pm.dstore.Delete(ds.NewKey(e.Key))
196 197 198 199 200 201 202
		if err != nil {
			log.Error("deleting provider set: ", err)
		}
	}
	return nil
}

203
func (pm *ProviderManager) getProvKeys() (func() (cid.Cid, bool), error) {
204
	res, err := pm.dstore.Query(dsq.Query{
205
		KeysOnly: true,
206 207 208 209 210 211
		Prefix:   providersKeyPrefix,
	})
	if err != nil {
		return nil, err
	}

212
	iter := func() (cid.Cid, bool) {
213 214 215
		for e := range res.Next() {
			parts := strings.Split(e.Key, "/")
			if len(parts) != 4 {
216
				log.Warningf("incorrectly formatted provider entry in datastore: %s", e.Key)
217 218 219 220
				continue
			}
			decoded, err := base32.RawStdEncoding.DecodeString(parts[2])
			if err != nil {
221
				log.Warning("error decoding base32 provider key: %s: %s", parts[2], err)
222 223
				continue
			}
224

225 226
			c, err := cid.Cast(decoded)
			if err != nil {
227
				log.Warning("error casting key to cid from datastore key: %s", err)
228 229
				continue
			}
230

231
			return c, true
232
		}
233
		return cid.Cid{}, false
234 235
	}

236
	return iter, nil
237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255
}

func (pm *ProviderManager) run() {
	tick := time.NewTicker(pm.cleanupInterval)
	for {
		select {
		case np := <-pm.newprovs:
			err := pm.addProv(np.k, np.val)
			if err != nil {
				log.Error("error adding new providers: ", err)
			}
		case gp := <-pm.getprovs:
			provs, err := pm.providersForKey(gp.k)
			if err != nil && err != ds.ErrNotFound {
				log.Error("error reading providers: ", err)
			}

			gp.resp <- provs
		case <-tick.C:
256
			keys, err := pm.getProvKeys()
257 258 259 260
			if err != nil {
				log.Error("Error loading provider keys: ", err)
				continue
			}
Steven Allen's avatar
Steven Allen committed
261
			now := time.Now()
262 263 264 265 266 267
			for {
				k, ok := keys()
				if !ok {
					break
				}

268 269 270 271 272 273
				provs, err := pm.getProvSet(k)
				if err != nil {
					log.Error("error loading known provset: ", err)
					continue
				}
				for p, t := range provs.set {
Steven Allen's avatar
Steven Allen committed
274
					if now.Sub(t) > ProvideValidity {
275 276 277
						delete(provs.set, p)
					}
				}
Steven Allen's avatar
Steven Allen committed
278 279 280
				// have we run out of providers?
				if len(provs.set) == 0 {
					provs.providers = nil
281 282 283 284
					err := pm.deleteProvSet(k)
					if err != nil {
						log.Error("error deleting provider set: ", err)
					}
Steven Allen's avatar
Steven Allen committed
285 286
				} else if len(provs.set) < len(provs.providers) {
					// We must have modified the providers set, recompute.
Steven Allen's avatar
Steven Allen committed
287 288 289 290
					provs.providers = make([]peer.ID, 0, len(provs.set))
					for p := range provs.set {
						provs.providers = append(provs.providers, p)
					}
291 292 293
				}
			}
		case <-pm.proc.Closing():
294
			tick.Stop()
295 296 297 298 299
			return
		}
	}
}

300
func (pm *ProviderManager) AddProvider(ctx context.Context, k cid.Cid, val peer.ID) {
301 302 303 304 305 306 307 308 309 310
	prov := &addProv{
		k:   k,
		val: val,
	}
	select {
	case pm.newprovs <- prov:
	case <-ctx.Done():
	}
}

311
func (pm *ProviderManager) GetProviders(ctx context.Context, k cid.Cid) []peer.ID {
312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346
	gp := &getProv{
		k:    k,
		resp: make(chan []peer.ID, 1), // buffered to prevent sender from blocking
	}
	select {
	case <-ctx.Done():
		return nil
	case pm.getprovs <- gp:
	}
	select {
	case <-ctx.Done():
		return nil
	case peers := <-gp.resp:
		return peers
	}
}

func newProviderSet() *providerSet {
	return &providerSet{
		set: make(map[peer.ID]time.Time),
	}
}

func (ps *providerSet) Add(p peer.ID) {
	ps.setVal(p, time.Now())
}

func (ps *providerSet) setVal(p peer.ID, t time.Time) {
	_, found := ps.set[p]
	if !found {
		ps.providers = append(ps.providers, p)
	}

	ps.set[p] = t
}