storage

package
v0.0.0-...-8ca6afe Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 3, 2026 License: MIT Imports: 24 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var ErrCorruptBlob = errors.New("blob does not match its content hash")

ErrCorruptBlob is returned when a fetched blob does not hash to its key.

View Source
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")
)
View Source
var (
	// ErrNotFound indicates a key is not in the store.
	ErrNotFound = errors.New("not found")
)
View Source
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")
)
View Source
var (
	// ErrTimeout is the error returned for when requests time out
	ErrTimeout = errors.New("request timed out")
)
View Source
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

func Clear(store Store) (removed int, err error)

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

func ContentOf(key []byte) func(value []byte) bool

ContentOf returns the check for a content-addressed key: a value is correct when it hashes to the key.

func IsNotLeader

func IsNotLeader(err error) (string, bool)

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.

func PutBatch

func PutBatch(store Store, pairs []KeyValue) error

PutBatch writes several pairs, as one unit where the store can and one at a time where it cannot.

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.

func (*BlobStoreWrapper) Put

func (s *BlobStoreWrapper) Put(value []byte) (key []byte, err error)

Put implements the BlobStore interface

type BlobSync

type BlobSync struct {
	Interval time.Duration
	Never    bool
}

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 ChangeListener func(message.Message)

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

type Deletable interface {
	Delete(key []byte) error
}

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

func NewDiskStore(dir string) *DiskStore

NewDiskStore constructs a new Disk backed store

func (*DiskStore) Close

func (s *DiskStore) Close() error

Close syncs anything a deferred policy still owes and stops its timer.

func (*DiskStore) Delete

func (s *DiskStore) Delete(key []byte) error

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

func (s *DiskStore) ForEachKey(fn func(key []byte, size int64, modified time.Time) error) error

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) Get

func (s *DiskStore) Get(key []byte) (value []byte, err error)

Get implements the BlobStore interface

func (*DiskStore) Has

func (s *DiskStore) Has(key []byte) (bool, error)

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

func (s *DiskStore) Put(key, value []byte) (err error)

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

func (s *DiskStore) SetSync(policy BlobSync)

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.

func (*DiskStore) Stat

func (s *DiskStore) Stat(key []byte) (int64, bool, error)

Stat is Has, also saying how large the held copy is -- which is what lets a node that is about to drop its own copy tell a peer's whole copy from a truncated one.

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

type InMemoryStore struct {
	sync.Mutex
	// contains filtered or unexported fields
}

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) Get

func (s *InMemoryStore) Get(key []byte) (value []byte, err error)

func (*InMemoryStore) Put

func (s *InMemoryStore) Put(key, value []byte) (err error)

type Iterable

type Iterable interface {
	ForEach(fn func(key, value []byte) error) error
}

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 KeyValue

type KeyValue struct {
	Key   []byte
	Value []byte
}

KeyValue is a pair for a batched write.

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

func NewLRUStore(disk *DiskStore, max int64) *LRUStore

NewLRUStore wraps disk, counting what it already holds and trimming it to max straight away if it holds more.

func (*LRUStore) Delete

func (s *LRUStore) Delete(key []byte) error

Delete implements Deletable, and stops counting the blob.

func (*LRUStore) Get

func (s *LRUStore) Get(key []byte) ([]byte, error)

Get implements Store, marking the blob as just used.

func (*LRUStore) Put

func (s *LRUStore) Put(key, value []byte) error

Put implements Store, then lets go of the least recently used blobs until the cache is back under its limit. The blob just written is never the one let go.

func (*LRUStore) Used

func (s *LRUStore) Used() int64

Used is how many bytes the cache is counted as holding.

type MembersFunc

type MembersFunc func() []placement.Member

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

type Notifier interface {
	OnMutation(fn func(key, value []byte, version uint64))
}

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 WithAuthKey(value string) Option

func WithChangeListener

func WithChangeListener(value ChangeListener) Option

func WithMaxRedirects

func WithMaxRedirects(value int) Option

WithMaxRedirects bounds how many times a request will follow a redirect before giving up.

func WithRedirectBackoff

func WithRedirectBackoff(value time.Duration) Option

WithRedirectBackoff sets how long to wait before retrying when the cluster has no leader yet.

func WithRequestTimeout

func WithRequestTimeout(value time.Duration) Option

func WithResponseBackoff

func WithResponseBackoff(value time.Duration) Option

func WithResyncListener

func WithResyncListener(value ResyncListener) Option

WithResyncListener sets what to call when a change notification was dropped.

func WithRetryTimeout

func WithRetryTimeout(value time.Duration) Option

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 NewPaired

func NewPaired(fast, slow Store) Paired

func (Paired) Get

func (s Paired) Get(key []byte) (value []byte, err error)

func (Paired) GetChecked

func (s Paired) GetChecked(key []byte, valid func([]byte) bool) (value []byte, err error)

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.

func (Paired) Put

func (s Paired) Put(key, value []byte) error

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) Get

func (r *RemoteStore) Get(key []byte) (value []byte, err error)

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

func (s *ReplicatedStore) GetChecked(key []byte, valid func([]byte) bool) ([]byte, error)

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

func (s *ReplicatedStore) Stat(address string, key []byte) (int64, bool, error)

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

func NewBitcaskStore(dbPath string) (Store, error)

NewBitcaskStore creates a new store using Bitcask

func NewStore

func NewStore(store string) (Store, error)

NewStore constructs a new store from the `store` uri and returns a `Store` interfaces matching the store type in `://...`

type StoreURI

type StoreURI struct {
	Type string
	Path string
}

StoreURI holds configuration parameters for a store parsed from a string such as ?=

func ParseStoreURI

func ParseStoreURI(uri string) (*StoreURI, error)

func (StoreURI) IsZero

func (u StoreURI) IsZero() bool

func (StoreURI) String

func (u StoreURI) String() string

type VersionedEntry

type VersionedEntry struct {
	Key     []byte
	Value   []byte
	Version uint64
}

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

type VersionedWrapper struct {
	sync.Mutex
	// contains filtered or unexported fields
}

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.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL