Documentation
¶
Index ¶
- Variables
- func ApplyMessage(store VersionedStore, in message.Message) (out message.Message)
- func Clear(store Store) (removed int, err error)
- func ContentOf(key []byte) func(value []byte) bool
- func IsNotLeader(err error) (string, bool)
- func PutAll(store VersionedStore, entries []VersionedEntry) error
- func PutBatch(store Store, pairs []KeyValue) error
- type AtomicStore
- type BitcaskStore
- func (s *BitcaskStore) Close() error
- func (s *BitcaskStore) Delete(key []byte) error
- func (s *BitcaskStore) ForEach(fn func(key, value []byte) error) error
- func (s *BitcaskStore) Get(key []byte) (value []byte, err error)
- func (s *BitcaskStore) Merge() error
- func (s *BitcaskStore) Put(key, value []byte) (err error)
- func (s *BitcaskStore) PutBatch(pairs []KeyValue) error
- func (s *BitcaskStore) Reclaimable() (reclaimable, size int64, err error)
- type BlobStore
- type BlobStoreWrapper
- type BlobSync
- type ChangeListener
- type CheckedGetter
- type Deletable
- type DiskStore
- func (s *DiskStore) Close() error
- func (s *DiskStore) Delete(key []byte) error
- func (s *DiskStore) ForEachKey(fn func(key []byte, size int64, modified time.Time) error) error
- func (s *DiskStore) Get(key []byte) (value []byte, err error)
- func (s *DiskStore) Has(key []byte) (bool, error)
- func (s *DiskStore) Put(key, value []byte) (err error)
- func (s *DiskStore) SetSync(policy BlobSync)
- func (s *DiskStore) Stat(key []byte) (int64, bool, error)
- type Enumerable
- type ErrNotLeader
- type InMemoryStore
- type Iterable
- type KeyValue
- type LRUStore
- type MembersFunc
- type Mergeable
- type MultiVersionedStore
- type Notifier
- type Option
- func WithAuthKey(value string) Option
- func WithChangeListener(value ChangeListener) Option
- func WithMaxRedirects(value int) Option
- func WithRedirectBackoff(value time.Duration) Option
- func WithRequestTimeout(value time.Duration) Option
- func WithResponseBackoff(value time.Duration) Option
- func WithResyncListener(value ResyncListener) Option
- func WithRetryTimeout(value time.Duration) Option
- type Paired
- type RemoteStore
- type RemoteVersionedStore
- func (rs *RemoteVersionedStore) Get(key []byte) (version uint64, value []byte, err error)
- func (rs *RemoteVersionedStore) GetMulti(keys [][]byte) ([]VersionedEntry, error)
- func (rs *RemoteVersionedStore) Put(version uint64, key []byte, value []byte) (err error)
- func (rs *RemoteVersionedStore) PutMulti(entries []VersionedEntry) (err error)
- func (rs *RemoteVersionedStore) Start()
- func (rs *RemoteVersionedStore) Stop()
- type ReplicatedStore
- func (s *ReplicatedStore) Get(key []byte) ([]byte, error)
- func (s *ReplicatedStore) GetChecked(key []byte, valid func([]byte) bool) ([]byte, error)
- func (s *ReplicatedStore) Has(address string, key []byte) (bool, error)
- func (s *ReplicatedStore) Preferred(key []byte) []string
- func (s *ReplicatedStore) Put(key, value []byte) error
- func (s *ReplicatedStore) PutTo(address string, key, value []byte) error
- func (s *ReplicatedStore) Replicas() int
- func (s *ReplicatedStore) Stat(address string, key []byte) (int64, bool, error)
- func (s *ReplicatedStore) WithTLS(cfg *tls.Config) *ReplicatedStore
- type ResyncListener
- type Store
- type StoreURI
- type VersionedEntry
- type VersionedStore
- type VersionedWrapper
- func (s *VersionedWrapper) Get(key []byte) (version uint64, value []byte, err error)
- func (s *VersionedWrapper) GetMulti(keys [][]byte) ([]VersionedEntry, error)
- func (s *VersionedWrapper) Put(version uint64, key []byte, value []byte) error
- func (s *VersionedWrapper) PutMulti(entries []VersionedEntry) error
Constants ¶
This section is empty.
Variables ¶
var ErrCorruptBlob = errors.New("blob does not match its content hash")
ErrCorruptBlob is returned when a fetched blob does not hash to its key.
var ( // ErrInvalidStore is returned when calling NewStore() with an invalid or unspproted // store type. Support stores are: memory, disk, bitcask ErrInvalidStore = errors.New("error: invalid or unsupproted store") )
var ( // ErrNotFound indicates a key is not in the store. ErrNotFound = errors.New("not found") )
var ( // ErrStalePut indicates that some client has not see the latest version of the // key-value pair being put. The client should get the current version, decide // if it still wants to do the put, and in that case do the put with the // correct version. ErrStalePut = errors.New("stale put") )
var ( // ErrTimeout is the error returned for when requests time out ErrTimeout = errors.New("request timed out") )
var ErrUnanswered = errors.New("the node could not say whether it holds the blob")
ErrUnanswered is Has's answer when the node could not say whether it holds a blob -- most often because it predates the question and refuses the method. It is not "absent": a caller that would delete something on the strength of the answer must treat it as a refusal to act (R148).
Functions ¶
func ApplyMessage ¶
func ApplyMessage(store VersionedStore, in message.Message) (out message.Message)
ApplyMessage applies the message to the store
func Clear ¶
Clear removes every key from the store and reports how many it removed. It is used when restoring a raft snapshot, which must replace the state machine's contents rather than merge into them, and when discarding a state machine that an unclean stop may have left ahead of the log.
The count is returned because those two callers both say what they did, and "nothing was there" and "eleven thousand records were thrown away" are not the same event however similar the code path.
func ContentOf ¶
ContentOf returns the check for a content-addressed key: a value is correct when it hashes to the key.
func IsNotLeader ¶
IsNotLeader reports whether err is an ErrNotLeader and returns the leader address it carries.
func PutAll ¶
func PutAll(store VersionedStore, entries []VersionedEntry) error
PutAll writes several pairs, as one unit where the store can do that and one at a time where it cannot. The fallback is not atomic, and is there so that a caller does not have to know which kind of store it holds: a single-node store built from VersionedWrapper is atomic, a remote one is atomic, and a bare VersionedStore from somewhere else still works.
Types ¶
type AtomicStore ¶
type AtomicStore interface {
Store
// PutBatch writes every pair or none of them.
PutBatch(pairs []KeyValue) error
}
AtomicStore is a Store whose writes can be grouped so that either all of them land or none does. Bitcask, which is dinofs's default metadata backend, can do this through a transaction; a directory of files cannot, which is why this is a separate interface rather than part of Store.
Callers should reach for PutBatch, which uses this where it exists and writes one at a time where it does not.
type BitcaskStore ¶
type BitcaskStore struct {
// contains filtered or unexported fields
}
BitcaskStore is a bitcask based storage engine
func (*BitcaskStore) Close ¶
func (s *BitcaskStore) Close() error
Close releases the underlying database.
func (*BitcaskStore) Delete ¶
func (s *BitcaskStore) Delete(key []byte) error
Delete implements the Deletable interface.
func (*BitcaskStore) ForEach ¶
func (s *BitcaskStore) ForEach(fn func(key, value []byte) error) error
ForEach implements the Iterable interface.
func (*BitcaskStore) Get ¶
func (s *BitcaskStore) Get(key []byte) (value []byte, err error)
Get implements the Store interface
func (*BitcaskStore) Merge ¶
func (s *BitcaskStore) Merge() error
Merge implements the Mergeable interface, reclaiming the space taken by values that have been overwritten or deleted.
Bitcask is log-structured: every write is an append, so a record rewritten many times occupies the store many times over. A filesystem does exactly that -- a directory record is rewritten on every create in it -- so the store grows with the number of writes rather than with the amount of data, and nothing gave that space back before this existed.
It is safe to call on a live store: bitcask takes its lock briefly to roll the file it is writing to, drops it for the whole rewrite, and re-takes it only to swap the compacted files in and reopen. Reads and writes continue through the long part.
That was not true before bitcask v2.2.0. Merge and Stats both read the key trie without taking the database's own lock, so this store held an exclusive mutex of its own across each of them -- which made a merge block every read and write for its whole duration, 3.5s at fifty thousand keys. The library locks them now (go.mills.io/bitcask#278, and Stats, Sync and Iterator after it), so the workaround is gone.
What is left is the swap, and it is not free: closing, deleting, renaming and reopening scales with the key count, and it measures about a fifth of the merge -- 1.1s at fifty thousand keys. So where a merge runs is still a decision, not a detail; see internal/raftstore/merge.go. TestMergeReclaimsAndStallsOnlyForTheSwap is what holds both halves of that to the measurement.
func (*BitcaskStore) Put ¶
func (s *BitcaskStore) Put(key, value []byte) (err error)
Put implements the Store interface
func (*BitcaskStore) PutBatch ¶
func (s *BitcaskStore) PutBatch(pairs []KeyValue) error
PutBatch implements the AtomicStore interface, writing every pair or none.
Bitcask gives us a transaction -- a snapshot of the key space, isolated from other transactions, whose writes are batched and land together on commit -- and this is the one place dinofs has a use for it. The replicated state machine applies a multi-key command (a create writes the new inode record and its parent directory), and applying half of one would leave a record nothing points at. Committing the whole batch or discarding it removes that window.
Note what this is not: it is not a durability barrier. The store is opened without WithSync, so bitcask buffers rather than fsyncing, deliberately -- durability comes from the raft log, and the state machine can be rebuilt from it. This makes the write atomic, not synchronous.
func (*BitcaskStore) Reclaimable ¶
func (s *BitcaskStore) Reclaimable() (reclaimable, size int64, err error)
Reclaimable implements the Mergeable interface.
type BlobStore ¶
type BlobStore interface {
Get(key []byte) (value []byte, err error)
Put(value []byte) (key []byte, err error)
}
BlobStore is the interface for storing and retrieving blogs of data
type BlobStoreWrapper ¶
type BlobStoreWrapper struct {
// contains filtered or unexported fields
}
BlobStoreWrapper wraps a Store to make sure content is never overwritten, by using as key for a value the Blake2b hash of the value. Even if there are concurrent writes for the same key, those would write the same contents (with very high probability).
func NewBlobStore ¶
func NewBlobStore(delegate Store) *BlobStoreWrapper
NewBlobStore creates a new blob store with the provided delegate store
func (*BlobStoreWrapper) Get ¶
func (s *BlobStoreWrapper) Get(key []byte) (value []byte, err error)
Get implements the BlobStore interface. It verifies that what came back hashes to the key it asked for: a blob is content-addressed, so a mismatch means corruption on disk or in flight, or a forged blob from a peer, and returning it silently would let that spread. The check is the same hash the key already is, so it costs one hash per read and nothing in trust.
A delegate that holds more than one copy -- the replicas, or a cache in front of them -- is told what a correct value looks like, so that it can pass over a bad copy for a good one rather than hand back the first it found. Refusing is not enough on its own: one rotten replica would otherwise make a blob unreadable while two good ones sat beside it.
type BlobSync ¶
BlobSync is how a DiskStore pushes what it writes to stable storage. It is the blob plane's copy of the raft log's policy, so that one --fsync means the same thing for both: the zero value syncs each blob and its directory before the write is acknowledged, Interval syncs what was written on a timer, and Never leaves it to the operating system.
type ChangeListener ¶
type CheckedGetter ¶
type CheckedGetter interface {
GetChecked(key []byte, valid func(value []byte) bool) ([]byte, error)
}
CheckedGetter is a store that holds more than one copy of a value and can be told how to recognise a correct one, so that it passes over a bad copy instead of returning it.
type Deletable ¶
Deletable is an optional capability of a Store: removing a key. Restoring a raft snapshot needs it, in order to discard state that the snapshot does not contain, and so does collection.
type DiskStore ¶
type DiskStore struct {
// contains filtered or unexported fields
}
DiskStore implements Store.
func NewDiskStore ¶
NewDiskStore constructs a new Disk backed store
func (*DiskStore) Delete ¶
Delete implements the Deletable interface. A key that is not there is not an error: collection races with nothing in particular, but two sweeps of the same store, or a sweep and a cache eviction, should not turn a blob that is already gone into a failure.
func (*DiskStore) ForEachKey ¶
ForEachKey implements the Enumerable interface, walking the fan-out directories and turning each filename back into the key it was made from.
A name that is not a key is skipped rather than refused. The store's directory is not a private format -- it is files on a disk, and an editor's backup or a half-written temporary file should not stop a sweep that is about to delete things.
func (*DiskStore) Has ¶
Has reports whether the store holds key, without reading it. It answers the question a rebalance asks of every blob it considers, where reading the bytes to find out would cost the transfer it is trying to avoid.
func (*DiskStore) Put ¶
Put implements the BlobStore interface.
The value goes to a temporary file beside the key and is renamed over it, so a reader sees the old file or the new one and never a partial one. That matters because a key is written more than once: two files sharing a chunk, a retried flush and a repair all put a blob that is already there, and writing in place truncates it first -- a reader in that moment got a short copy with no error. A write that fails part way, a full disk included, leaves the temporary file rather than a truncated blob at the key.
Renaming also refreshes the modification time, which is what the collection grace period reads: a blob that a new write refers to again is young again.
func (*DiskStore) SetSync ¶
SetSync gives the store a durability policy. Call it once, before the store takes writes.
Before this, nothing synced a blob at all: a write was acknowledged once a quorum of replicas held it in their filesystem cache, while the metadata pointing at it was synced by raft. A power cut across the replicas could keep the pointer and lose what it pointed at.
type Enumerable ¶
type Enumerable interface {
// ForEachKey visits every key, with the size of the value and the time it
// was last written. Visiting order is unspecified.
ForEachKey(fn func(key []byte, size int64, modified time.Time) error) error
}
Enumerable is an optional capability of a Store: listing what it holds without reading any of it.
It is separate from Iterable because the difference is the whole point. Collection asks a blob store what keys it has and how old they are, and a blob is the large half of this filesystem: reading every value to answer that would mean reading the entire store off disk to decide what to delete from it.
type ErrNotLeader ¶
type ErrNotLeader struct {
Leader string
}
ErrNotLeader is returned by a VersionedStore that replicates through a leader, when the operation reached a node that is not currently the leader. Leader is the address a client should talk to instead, or empty when no leader is known yet -- during an election, say -- in which case the client should retry rather than redirect.
It lives here, rather than in the package that implements replication, so that the message layer can translate it without importing that package.
func (ErrNotLeader) Error ¶
func (e ErrNotLeader) Error() string
Error implements the error interface.
type InMemoryStore ¶
InMemoryStore is a Store implementation powered by a map, to be used for testing or caches.
func NewInMemoryStore ¶
func NewInMemoryStore() *InMemoryStore
func (*InMemoryStore) Delete ¶
func (s *InMemoryStore) Delete(key []byte) error
Delete implements the Deletable interface.
func (*InMemoryStore) ForEach ¶
func (s *InMemoryStore) ForEach(fn func(key, value []byte) error) error
ForEach implements the Iterable interface.
func (*InMemoryStore) Put ¶
func (s *InMemoryStore) Put(key, value []byte) (err error)
type Iterable ¶
Iterable is an optional capability of a Store: visiting every key/value pair it holds. Raft snapshotting needs it, so any store used as a replicated state machine must implement it.
type LRUStore ¶
type LRUStore struct {
// contains filtered or unexported fields
}
LRUStore bounds a DiskStore used as a cache: it holds at most max bytes and lets go of the blobs used least recently.
A mount's cache was never trimmed. Every chunk read or written through a node stayed on its disk for the life of the mount, the chunks of deleted files included, and nothing but the disk filling up would ever stop it. Letting go of a cached blob is always safe -- the replicas hold it -- so the only question is which, and the one used longest ago is the least likely to be wanted again.
Recency is kept in memory, seeded at start from the files' modification times, rather than by touching a file on every read.
func NewLRUStore ¶
NewLRUStore wraps disk, counting what it already holds and trimming it to max straight away if it holds more.
type MembersFunc ¶
MembersFunc reports the blob servers currently believed to be alive, with the weight and zone each states. It is called on every operation, so that the set can change underneath as nodes join and leave.
func Unweighted ¶
func Unweighted(addrs func() []string) MembersFunc
Unweighted is a MembersFunc over plain addresses: equal weights, no zones, which is every cluster from before either existed.
type Mergeable ¶
type Mergeable interface {
// Merge rewrites the store keeping only the current value of each key.
//
// It is safe to call on a live store but it is not cheap: the store stops
// serving for the whole merge, measured at about seventy microseconds a
// key. Where that is affordable is the caller's problem, and for dinofs it
// is the reason a leader hands leadership over first.
Merge() error
// Reclaimable reports how many bytes the store would give back, and how
// many it currently occupies. A policy needs both: the ratio is what says
// whether a merge is worth its cost.
Reclaimable() (reclaimable, size int64, err error)
}
Mergeable is a store that can compact itself, reclaiming the space taken by values that have been overwritten or deleted.
Only the log-structured backend has anything to reclaim. Bitcask appends every write, so the store grows with the number of writes rather than with the amount of data in it: one 4KiB record rewritten twenty thousand times is 78.8MiB on disk for 4KiB of live data. Nothing reclaimed that until this existed.
type MultiVersionedStore ¶
type MultiVersionedStore interface {
VersionedStore
// PutMulti applies every entry or none. It returns ErrStalePut if any entry
// names a version that is not the one it would be replacing.
PutMulti(entries []VersionedEntry) error
// GetMulti reads several keys in one round trip. Keys that are not there
// are absent from the result rather than an error: a caller asking for many
// usually expects some to be gone, because another writer may have removed
// one between listing it and reading it.
GetMulti(keys [][]byte) ([]VersionedEntry, error)
}
MultiVersionedStore is a VersionedStore that can write several pairs as one unit: either all of them are written or none is, and a stale version on any of them rejects the whole write.
It exists because a filesystem operation is rarely one record. Creating a file writes the new inode and then its parent directory, and through raft each of those is a commit, which is an fsync -- so two of them is twice the cost. It is also not atomic: a crash between the two leaves an inode record nothing points at.
It is a separate interface rather than a method on VersionedStore because not every store can offer the guarantee. Callers should type-assert and fall back to a sequence of Puts, which is what PutAll does.
type Notifier ¶
Notifier is an optional capability of a VersionedStore: reporting every mutation it applies, in a total order. A store that replicates through a log can provide this, and the metadata server uses it to fan changes out to connected clients in the order the log committed them.
The listener is called on the store's apply path and must not block.
type Option ¶
type Option func(*options)
func WithAuthKey ¶
func WithChangeListener ¶
func WithChangeListener(value ChangeListener) Option
func WithMaxRedirects ¶
WithMaxRedirects bounds how many times a request will follow a redirect before giving up.
func WithRedirectBackoff ¶
WithRedirectBackoff sets how long to wait before retrying when the cluster has no leader yet.
func WithRequestTimeout ¶
func WithResponseBackoff ¶
func WithResyncListener ¶
func WithResyncListener(value ResyncListener) Option
WithResyncListener sets what to call when a change notification was dropped.
func WithRetryTimeout ¶
WithRetryTimeout bounds how long a request retries across redirects and node failures before giving up. It should comfortably exceed an election.
type Paired ¶
type Paired struct {
// contains filtered or unexported fields
}
Paired implements Store wrapping a pair of stores, one fast, one slow. A put writes through to the slow store before returning, so a write is durable once it succeeds; the fast store is a read cache in front of it, populated on write and on a read miss.
It used to acknowledge a put as soon as the fast (local) store had it and replicate to the slow store on a background queue that retried forever. That made fsync a lie: it returned success with the content on one machine, while the metadata about to be committed through raft was genuinely durable, so a node dying between the two lost the file on a cluster that reported itself healthy. The queue could also fill and wedge the mount. The write is now synchronous.
func (Paired) GetChecked ¶
GetChecked is Get, passing over any copy valid rejects. A cached copy that fails is dropped and fetched again; a fetched copy that fails is refused and never cached, because the cache answers every later read and would turn one bad read into a file that fails for as long as the cache lives.
type RemoteStore ¶
type RemoteStore struct {
// contains filtered or unexported fields
}
RemoteStore implements Store. It requires to connect to a blobserver.
func NewRemoteStore ¶
func NewRemoteStore(address, token string) *RemoteStore
NewRemoteStore connects to a blob server. token, when non-empty, is sent as a bearer credential; pass "" for an unauthenticated server.
func (*RemoteStore) Has ¶
func (r *RemoteStore) Has(key []byte) (bool, error)
Has asks the blob server whether it holds key, without transferring it.
func (*RemoteStore) Put ¶
func (r *RemoteStore) Put(key, value []byte) (err error)
func (*RemoteStore) Stat ¶
func (r *RemoteStore) Stat(key []byte) (int64, bool, error)
Stat is Has, also returning the size of the copy the server holds, or -1 when the server did not say.
func (*RemoteStore) WithTLS ¶
func (r *RemoteStore) WithTLS(cfg *tls.Config) *RemoteStore
WithTLS reaches the server inside TLS with this configuration. A nil one leaves the store in plaintext, which is a cluster with no secret.
type RemoteVersionedStore ¶
type RemoteVersionedStore struct {
// contains filtered or unexported fields
}
RemoteVersionedStore is an implementation of VersionedStore, via a client to a remote metadataserver process.
func NewRemoteVersionedStore ¶
func NewRemoteVersionedStore(remote *client.Client, options ...Option) *RemoteVersionedStore
func (*RemoteVersionedStore) Get ¶
func (rs *RemoteVersionedStore) Get(key []byte) (version uint64, value []byte, err error)
func (*RemoteVersionedStore) GetMulti ¶
func (rs *RemoteVersionedStore) GetMulti(keys [][]byte) ([]VersionedEntry, error)
GetMulti reads several keys in one round trip. Reading a directory means reading every record it names, and one request per name is one wait per name: a thousand files was a thousand of them.
Keys that are not there come back absent rather than as an error, because a caller asking for many usually expects some to be gone.
func (*RemoteVersionedStore) Put ¶
func (rs *RemoteVersionedStore) Put(version uint64, key []byte, value []byte) (err error)
func (*RemoteVersionedStore) PutMulti ¶
func (rs *RemoteVersionedStore) PutMulti(entries []VersionedEntry) (err error)
PutMulti writes several pairs as one unit. The server applies them through one raft commit, so an operation that touches two records -- creating a file writes the new inode and then its parent -- costs one fsync rather than two, and cannot be interrupted between them.
func (*RemoteVersionedStore) Start ¶
func (rs *RemoteVersionedStore) Start()
func (*RemoteVersionedStore) Stop ¶
func (rs *RemoteVersionedStore) Stop()
type ReplicatedStore ¶
type ReplicatedStore struct {
// contains filtered or unexported fields
}
ReplicatedStore spreads blobs across several blob servers.
Replicating content-addressed data needs no coordination: a key is the hash of its own value, so two nodes can never disagree about what a key means and there is no write conflict to resolve. That makes this far simpler than the metadata plane -- there is no consensus here, only copies.
func NewReplicatedStore ¶
func NewReplicatedStore(members MembersFunc, replicas int, token string) *ReplicatedStore
NewReplicatedStore creates a store that keeps replicas copies of each blob. A replicas of 1 gives no redundancy; 3 tolerates losing one node while still accepting writes.
func (*ReplicatedStore) Get ¶
func (s *ReplicatedStore) Get(key []byte) ([]byte, error)
Get fetches the blob from the first member that has it, searching in preference order. A blob found outside its preferred replicas is copied back to them, so that a node which missed a write catches up the first time anybody reads the data.
func (*ReplicatedStore) GetChecked ¶
GetChecked is Get, treating a copy valid rejects as a bad replica: it is passed over for the next one, and overwritten with the good copy when one is found -- the same repair a missing copy gets.
func (*ReplicatedStore) Has ¶
func (s *ReplicatedStore) Has(address string, key []byte) (bool, error)
Has asks one blob server whether it holds key. See RemoteStore.Has for what an answer of ErrUnanswered means.
func (*ReplicatedStore) Preferred ¶
func (s *ReplicatedStore) Preferred(key []byte) []string
Preferred is every member in the order a blob with this key is placed: the first Replicas are where it should live. Every node and every client computes the same order from the same member list, which is what lets a rebalance on one node decide where a blob belongs without asking anyone.
func (*ReplicatedStore) Put ¶
func (s *ReplicatedStore) Put(key, value []byte) error
Put stores the blob on the preferred members, returning once a quorum of them has accepted it. Copies that failed are not retried here: the caller wraps this in a Paired store whose write-back retries, and a later read repairs any replica that is still missing.
A preferred member that refuses the write -- full, unreachable, erroring -- is replaced by the next member down the preference list, up to as many spares as there are replicas (R150). Without that the first full node fails its share of every write with nothing to catch them. A read finds a spilled copy because it walks the whole list, and the push pass returns the blob to the preferred member once that member takes writes again. Every refusal is replaced, not only enough to make a quorum: with three replicas, a quorum of two would absorb one full node silently and leave every blob it should have held one copy short.
func (*ReplicatedStore) PutTo ¶
func (s *ReplicatedStore) PutTo(address string, key, value []byte) error
PutTo stores a blob on one blob server, whatever the placement says.
func (*ReplicatedStore) Replicas ¶
func (s *ReplicatedStore) Replicas() int
Replicas is how many copies of each blob this store keeps.
func (*ReplicatedStore) Stat ¶
Stat asks one blob server whether it holds key and how large its copy is.
func (*ReplicatedStore) WithTLS ¶
func (s *ReplicatedStore) WithTLS(cfg *tls.Config) *ReplicatedStore
WithTLS makes every blob server this store talks to be reached inside TLS with this configuration. It must be called before the store is used.
type ResyncListener ¶
type ResyncListener func()
ResyncListener is told that a change notification could not be delivered, so that whatever the listener caches can be reloaded wholesale. It is called once per run of drops rather than once per drop: the caller has no way to know what it missed, so the answer is always the same and doing it many times costs the same as doing it once.
type Store ¶
type Store interface {
Put(key, value []byte) (err error)
// Get should return ErrNotFound if the key is not in the store.
Get(key []byte) (value []byte, err error)
}
Store represents a key-value store.
func NewBitcaskStore ¶
NewBitcaskStore creates a new store using Bitcask
type StoreURI ¶
StoreURI holds configuration parameters for a store parsed from a string
such as
func ParseStoreURI ¶
type VersionedEntry ¶
VersionedEntry is one key-value pair of a multi-key write, with the version the writer believes it is replacing.
func GetAll ¶
func GetAll(store VersionedStore, keys [][]byte) ([]VersionedEntry, error)
GetAll reads several keys, in one round trip where the store can do that and one at a time where it cannot. Missing keys are omitted.
type VersionedStore ¶
type VersionedStore interface {
// Put should return ErrStalePut if the current version is not the version
// passed as argument minus one. The client should have to prove that they've
// seen the most current version before trying to update it.
Put(version uint64, key []byte, value []byte) (err error)
// Get should return ErrNotFound if the key is not in the store.
Get(key []byte) (version uint64, value []byte, err error)
}
type VersionedWrapper ¶
VersionedWrapper is a VersionedStore implementation wraping a given Store implementation. This is the quickest way of building a VersionedStore, but it's alos the slowest, as it serializes all calls to the underlying Store.
func NewVersionedWrapper ¶
func NewVersionedWrapper(delegate Store) *VersionedWrapper
func (*VersionedWrapper) Get ¶
func (s *VersionedWrapper) Get(key []byte) (version uint64, value []byte, err error)
Get retrieves the value associated with a key and its version number.
func (*VersionedWrapper) GetMulti ¶
func (s *VersionedWrapper) GetMulti(keys [][]byte) ([]VersionedEntry, error)
GetMulti reads several keys under one lock.
func (*VersionedWrapper) Put ¶
func (s *VersionedWrapper) Put(version uint64, key []byte, value []byte) error
Put stores the given value at the given key, provided the passed version number is the current version number. If the put is successful, the version number is incremented by one.
func (*VersionedWrapper) PutMulti ¶
func (s *VersionedWrapper) PutMulti(entries []VersionedEntry) error
PutMulti applies every entry or none. The check for all of them happens before any of them is written, under the one lock, which is what makes it a unit.