spool

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: 21 Imported by: 0

Documentation

Overview

Codecs encode a batch entry's rows into its data file: ndjson and gzipped ndjson here, and others registered by the packages that implement them.

Package spool hands batches of stream writes from any process to the one that owns a store, as directories of data files committed by rename.

A spool directory holds tmp/ (batches being written), incoming/ (published batches, oldest first by name), failed/ (batches that could not be ingested, with an error.json) and trash/ (ingested batches awaiting collection). A batch is published by writing it under tmp/ and renaming it into incoming/: the rename is its only commit point, so a reader never sees half a batch.

Loading a published batch back: the manifest is checked, every data file is verified against its digest and row count, and the rows are decoded.

The manifest: a published batch's description, with explicit JSON names and string enums so a build reading another build's batch reads it the same.

Row values as a sqlite kind table stores them, so a spooled row and the same row appended directly store identical values.

Publishing a batch: its data files and manifest are written and made durable under tmp/, then the batch directory is renamed into incoming/.

The bulk Writer: a producer's writes gathered into batches cut at a row and byte cap, each delivered whole and in order.

Index

Constants

View Source
const (
	FormatNDJSON     = "ndjson"
	FormatNDJSONGzip = "ndjson.gz"
)
View Source
const ManifestFormat = 1

ManifestFormat is the manifest format this build writes and reads.

Variables

View Source
var ErrManifestFormat = errors.New("spool manifest format unsupported")

ErrManifestFormat reports a batch published in a manifest format this build does not read.

Functions

func Digest

func Digest(data []byte) string

Digest is the SHA256 a manifest records for a data file's bytes.

func Formats

func Formats() []string

Formats are the formats this build reads and writes, sorted.

func Normalize

func Normalize(schema recordstore.KindSchema, row recordstore.Row) (recordstore.Row, error)

Normalize converts row into the values sqlitetable.Value stores for schema's columns: a structured column's value as its JSON text, a time as sqlitetable.TimeLayout, and any other value the driver cannot store as its JSON text. A key schema does not declare is an error, but for a dynamic kind, whose key gets the value its inferred column stores — a null none, an object or an array itself, for the owner to infer its column from.

func Register

func Register(format string, codec Codec)

Register makes codec the one data files of format are written and read with. A format is registered once; registering it again panics, as a program wiring two codecs to one format is broken.

Types

type Codec

type Codec interface {
	Encode(w io.Writer, rows []recordstore.Row) error
	Decode(r io.Reader, fn func(recordstore.Row) error) error
}

Codec writes and reads the rows of one data file. Decode hands numbers over as json.Number, so an int64 survives the trip whole.

type CollectOptions

type CollectOptions struct {
	// TmpAge is how long a batch may stay unpublished in tmp, by its last
	// change, before it is taken as abandoned. Zero is 24 hours.
	TmpAge time.Duration

	// FailedAge is how long a failed batch is kept for inspection. Zero is 30
	// days.
	FailedAge time.Duration
}

CollectOptions say how long abandoned and failed batches are kept.

type Dir

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

Dir is a spool directory.

func OpenDir

func OpenDir(path string) (Dir, error)

OpenDir opens the spool directory at path, creating it and its subdirectories private to the user.

func (Dir) Collect

func (d Dir) Collect(now time.Time, options CollectOptions) error

Collect removes every trashed batch, the failed batches older than FailedAge and the tmp entries older than TmpAge, as of now.

func (Dir) Fail

func (d Dir) Fail(name string, reason error) error

Fail moves a batch that cannot be ingested to failed, with reason recorded beside it in error.json.

func (Dir) Incoming

func (d Dir) Incoming() ([]string, error)

Incoming lists the published batches, oldest first.

func (Dir) Load

func (d Dir) Load(name string) (recordstore.Batch, error)

Load reads the published batch name. An error means the batch as a whole cannot be ingested — ErrManifestFormat for a manifest format this build does not read — and belongs in failed/.

func (Dir) Path

func (d Dir) Path() string

Path is the spool directory.

func (Dir) Publish

func (d Dir) Publish(batch recordstore.Batch, format string) (string, error)

Publish writes batch into the spool, its rows in format, and returns the name it was published under in incoming/. Every kind the batch appends to must have its schema in batch.Schemas: the rows are normalized against it, and the owner ingests them under it. Nothing is published unless all of it is.

func (Dir) Publisher

func (d Dir) Publisher(format string) func(context.Context, recordstore.Batch) error

Publisher delivers a batch by publishing it in format, for a Writer.

func (Dir) Trash

func (d Dir) Trash(name string) error

Trash moves an ingested batch out of incoming, for Collect to remove. A batch trashed before — ingested again after a crash between its commit and its move — replaces the earlier copy.

type Producer

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

Producer numbers one process's batches: every batch it identifies gets the next seq, so an owner can ingest a producer's batches in order.

func NewProducer

func NewProducer(instance, build string) *Producer

NewProducer identifies the batches of this process as instance, made by build.

func (*Producer) Next

func (p *Producer) Next() recordstore.Producer

Next is the identity of the producer's next batch.

type Writer

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

Writer gathers writes into batches. A write it has taken is delivered by a later Append that fills a batch or by Flush; a bulk write spanning batches is not atomic. It is not safe for concurrent use.

func NewWriter

func NewWriter(options WriterOptions) (*Writer, error)

NewWriter returns a Writer delivering through options.Deliver.

func (*Writer) Append

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

Append gathers rows for stream, cutting and delivering a batch whenever one fills.

func (*Writer) Flush

func (w *Writer) Flush(ctx context.Context) error

Flush delivers everything gathered.

func (*Writer) Seal

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

Seal gathers a seal of stream, after the writes gathered before it.

type WriterOptions

type WriterOptions struct {
	// Producer identifies every batch the writer delivers. Required.
	Producer *Producer

	// Schema resolves the kinds appended to; a batch carries the schema of
	// every kind it appends to. Required.
	Schema recordstore.SchemaResolver

	// Deliver hands a batch on: Dir.Publisher, or a BatchAppender's
	// AppendBatch. A batch it fails is delivered again, with the same id and
	// seq, before any later one. Required.
	Deliver func(context.Context, recordstore.Batch) error

	// MaxRows and MaxBytes cap a batch's rows and their JSON size; zero is
	// 50,000 rows and 64MiB. A batch always takes at least one row.
	MaxRows  int
	MaxBytes int
}

WriterOptions configure a Writer.

Directories

Path Synopsis
Package parquet is the spool's parquet codec: importing it registers the "parquet" format, for producers exporting large batches in columnar form.
Package parquet is the spool's parquet codec: importing it registers the "parquet" format, for producers exporting large batches in columnar form.

Jump to

Keyboard shortcuts

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