owner

package
v0.1.45 Latest Latest
Warning

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

Go to latest
Published: Oct 2, 2026 License: Apache-2.0 Imports: 19 Imported by: 0

Documentation

Overview

A reader's writes: published to the spool, the owner nudged through its socket, and the outcome awaited in the ledger the owner commits it to.

Closing a store: an owner drains the spool and hands the store back; any process then ingests what was published after that, while nobody owns it.

The owner's control socket: one JSON request and response per connection, answering status and nudging ingestion. Nothing depends on it for correctness; a producer that cannot reach it waits for the owner's poll.

Refusing network filesystems: file locks and SQLite's WAL are unreliable there, so an owner elected on one could share its store with another.

The filesystem type of a directory on Linux, named from statfs's magic number.

The owner's ingestion: every published batch applied once, in each producer's order, then trashed; one it cannot read is moved to failed/.

One Store per path within a process: a second Open of a path shares the first's, so a process never waits on its own lock or spools to itself.

The lock that elects a store's owner, and Exclusive, which holds it for a backend with no read-only mode to fall back to.

Package owner elects one process to write a record store, the owner, by a file lock beside the store's configured path. Every other process reads the store and hands its writes to the owner through the store's spool, and the first of them to take the lock after the owner exits becomes the owner.

Beside a store configured at P the package keeps P.lock (the lock, the only truth about who owns the store), P.owner.json (the owner's advisory state), P.sock (the owner's control socket, which only speeds things up) and P.spool (the spool, see recordstore/spool).

Store: one process's hold on a record store — as its owner, writing the file and ingesting the spool, or as a reader handing its writes over the spool and waiting to take over.

A bulk Writer over a store: applied directly by the owner, or published to the spool by a reader, which may exit before the owner ingests it.

Index

Constants

This section is empty.

Variables

View Source
var (
	// ErrPending reports a write handed to the owner that the caller stopped
	// waiting for: it is durable in the spool and will be applied; Await
	// reports how. An unkeyed write must not be retried blindly.
	ErrPending = errors.New("record store write is pending")

	// ErrCatalogVersion reports a store whose owner reads another catalog
	// version than this build: writes still reach it through the spool, but
	// this process cannot read the file.
	ErrCatalogVersion = errors.New("record store owner runs another catalog version")
)
View Source
var ErrLocked = errors.New("record store is held by another process")

ErrLocked reports a store another process holds.

Functions

func ControlIngest

func ControlIngest(path string, ids []string, wait bool, timeout time.Duration) ([]string, error)

ControlIngest asks the owner serving the socket at path to ingest its spool now, and, with wait, to answer once the batches ids name are applied, with the ids it applied.

func Exclusive

func Exclusive(path, build string) (release func() error, err error)

Exclusive holds the store at path for this process, publishing its state so another process refused with a *LockedError can name it, and returns the function that releases it. It is for a store only one process can open at all, such as a directory of ndjson files.

func NetworkFilesystem

func NetworkFilesystem(name string) bool

NetworkFilesystem reports whether a filesystem of type name is one a store may not live on: a network filesystem, or any FUSE one, whose locking is whatever the userspace server makes of it.

func RefuseNetworkFilesystem

func RefuseNetworkFilesystem(dir string) error

RefuseNetworkFilesystem refuses dir when it is on a network filesystem. A platform that cannot name its filesystems accepts every dir.

func SocketPath

func SocketPath(path string) string

SocketPath is where the owner of the store configured at path serves its control socket: beside the store, or, when that address is too long for a socket, in the temp dir under a name derived from the store's path.

Types

type CatalogVersionError

type CatalogVersionError struct{ Owner, Build int }

CatalogVersionError is ErrCatalogVersion with both versions.

func (*CatalogVersionError) Error

func (e *CatalogVersionError) Error() string

func (*CatalogVersionError) Unwrap

func (e *CatalogVersionError) Unwrap() error

type ControlHandler

type ControlHandler struct {
	Status func() State
	Ingest func(ids []string, wait bool) ([]string, error)
}

ControlHandler answers control requests: Status with the owner's state, and Ingest by ingesting the spool now, returning — when asked to wait — the ids of the batches named that the store has applied.

type ControlServer

type ControlServer struct {
	// contains filtered or unexported fields
}

ControlServer serves a control socket until Close.

func ServeControl

func ServeControl(path string, handler ControlHandler) (*ControlServer, error)

ServeControl serves handler on the socket at path, replacing a socket file nothing answers on and refusing one something does.

func (*ControlServer) Close

func (s *ControlServer) Close() error

Close stops serving and unlinks the socket if it is still this server's: a later owner may already serve its own at the same path.

type FlushOptions

type FlushOptions struct {
	// Wait waits for the owner to apply every batch this writer published,
	// and reports the entries it refused. Without it Flush returns once the
	// batches are durable in the spool.
	Wait bool
}

FlushOptions say what Flush waits for.

type LockedError

type LockedError struct {
	Path  string
	Owner *State
}

LockedError is ErrLocked for the store at Path, naming the process holding it when that process published its state.

func (*LockedError) Error

func (e *LockedError) Error() string

func (*LockedError) Unwrap

func (e *LockedError) Unwrap() error

type Options

type Options[T Target] struct {
	// Path is the store as configured; every file beside it derives from it.
	Path string

	// Open opens the backend: writable for the owner, or read-only handing
	// its writes to submit.
	Open func(ctx context.Context, readOnly bool, submit recordstore.Submitter) (T, error)

	// Store is the file Open opens, and CatalogVersion the catalog this
	// build reads; the owner publishes both, and a reader of another catalog
	// version cannot read the file.
	Store          string
	CatalogVersion int

	// Build names this build in the state it publishes as owner.
	Build string

	// Format is the spool format this process publishes; empty is ndjson.
	Format string

	// StartWait bounds how long a reader waits for the owner to be ready;
	// zero is a minute.
	StartWait time.Duration

	// DrainOnClose bounds how long Close ingests what is left in the spool,
	// as owner, or when the owner is gone; zero is ten seconds.
	DrainOnClose time.Duration

	// Poll is the longest the owner waits between looks at the spool when
	// nothing nudges it; zero is 500ms.
	Poll time.Duration
}

Options configure Open.

type PendingError

type PendingError struct{ BatchID string }

PendingError is ErrPending for one batch.

func (*PendingError) Error

func (e *PendingError) Error() string

func (*PendingError) Unwrap

func (e *PendingError) Unwrap() error

type Phase

type Phase string

Phase is where the owner is in its life.

const (
	PhaseStarting Phase = "starting"
	PhaseReady    Phase = "ready"
	PhaseDraining Phase = "draining"
)

type Role

type Role string

Role is what a process does with a store.

const (
	RoleOwner  Role = "owner"
	RoleReader Role = "reader"
)

type SpoolState

type SpoolState struct {
	Dir            string   `json:"dir"`
	ManifestFormat int      `json:"manifestFormat"`
	Formats        []string `json:"formats"`
}

SpoolState is what a producer needs to publish batches the owner reads.

type State

type State struct {
	Instance       string      `json:"instance"`
	PID            int         `json:"pid"`
	Host           string      `json:"host,omitempty"`
	Build          string      `json:"build,omitempty"`
	Phase          Phase       `json:"phase"`
	StartedAt      time.Time   `json:"startedAt"`
	HeartbeatAt    time.Time   `json:"heartbeatAt"`
	Store          string      `json:"store,omitempty"`
	CatalogVersion int         `json:"catalogVersion,omitempty"`
	Spool          *SpoolState `json:"spool,omitempty"`
	Socket         string      `json:"socket,omitempty"`
	Backlog        int         `json:"backlog"`
	Failed         int         `json:"failed"`
	LastError      string      `json:"lastError,omitempty"`
}

State is the owner's advisory state, published beside the store. The lock, not this file, says who owns the store; a reader uses it to find the owner's socket and spool, and to name the owner in errors.

func ControlStatus

func ControlStatus(path string, timeout time.Duration) (State, error)

ControlStatus asks the owner serving the socket at path for its state.

func ReadState

func ReadState(path string) (State, bool, error)

ReadState reads the state the owner of the store configured at path published, reporting false when no owner published one.

type Store

type Store[T Target] struct {
	// contains filtered or unexported fields
}

Store is this process's hold on a record store.

func Open

func Open[T Target](ctx context.Context, options Options[T]) (*Store[T], error)

Open takes a hold on the store options.Path names: as its owner when no other process holds it, otherwise as a reader that takes over when the owner goes. Within one process every Open of a path shares one Store, released by the last Close.

func (*Store[T]) Await

func (s *Store[T]) Await(ctx context.Context, id string) (recordstore.BatchResult, error)

Await waits for the batch id to be applied and returns what it did. A batch of this process's that the owner could not read at all is an error; when ctx ends first the batch is still pending, and the error is a *PendingError.

func (*Store[T]) Backend

func (s *Store[T]) Backend() (T, error)

Backend is the store's backend: writable for the owner, read-only for a reader. A reader of an owner on another catalog version has none, and gets a *CatalogVersionError.

func (*Store[T]) Close

func (s *Store[T]) Close() error

Close releases this process's hold. The owner stops serving, drains the spool for up to DrainOnClose, withdraws its state and unlocks; a reader stops waiting for the lock. Either then ingests any batch still in the spool if nobody else owns the store, so a producer that published while the owner let go is not left waiting for the next owner.

func (*Store[T]) OnRole

func (s *Store[T]) OnRole(fn func(Role))

OnRole calls fn whenever this process's role changes: when a reader takes the store over, fn is called with RoleOwner.

func (*Store[T]) Role

func (s *Store[T]) Role() Role

Role is what this process does with the store now.

func (*Store[T]) Status

func (s *Store[T]) Status() (State, error)

Status is the owner's state: this process's own as owner, or the one the owner published.

func (*Store[T]) Writer

func (s *Store[T]) Writer(schema recordstore.SchemaResolver, options spool.WriterOptions) (*Writer[T], error)

Writer gathers writes of the kinds schema resolves into batches cut at options' caps; options' Producer, Schema and Deliver are the store's.

type Target

type Target interface {
	recordstore.Backend
	recordstore.BatchAppender
	Promote(ctx context.Context) error
	ProducerSeq(ctx context.Context, instance string) (int64, error)
	SweepBatches(ctx context.Context, before time.Time, keep func(id string) bool) (int, error)
}

Target is a backend a store can be elected over: opened read-only it submits its writes, and Promote makes it write the file itself.

type Writer

type Writer[T Target] struct {
	// contains filtered or unexported fields
}

Writer gathers writes into batches for a store. It is not safe for concurrent use.

func (*Writer[T]) Append

func (w *Writer[T]) Append(ctx context.Context, stream, kind string, rows []recordstore.Row) error

Append gathers rows for stream.

func (*Writer[T]) Flush

func (w *Writer[T]) Flush(ctx context.Context, options FlushOptions) error

Flush delivers everything gathered. With options.Wait it waits for every published batch and returns the entries the owner refused, joined; the refusals of batches the owner applied directly are returned either way.

func (*Writer[T]) Seal

func (w *Writer[T]) Seal(ctx context.Context, stream string) error

Seal gathers a seal of stream.

Jump to

Keyboard shortcuts

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