Documentation
¶
Overview ¶
Package seekable writes and reads streams using the Zstandard seekable format.
A seekable stream is a valid Zstandard stream made from one or more compressed frames followed by a final skippable frame containing a seek table. Standard Zstandard decoders can read the stream from the beginning, while Reader uses the seek table to serve Read, ReadAt, and Seek calls by uncompressed byte offset and exposes the parsed metadata through Reader.SeekTable.
Writer and Encoder produce seekable streams by storing each non-empty input chunk as a separate Zstandard frame. Close or EndStream must be called to append or retrieve the final seek-table skippable frame. Close is idempotent; EndStream finalizes the Encoder and must be called once. Operations after Close or EndStream return ErrClosed.
The package accepts small encoder and decoder interfaces and is tested with github.com/klauspost/compress/zstd.
Index ¶
Examples ¶
Constants ¶
This section is empty.
Variables ¶
var ErrClosed = errors.New("seekable: closed")
ErrClosed is returned by operations after a Reader, Writer, or Encoder has been closed or finalized.
Functions ¶
This section is empty.
Types ¶
type Encoder ¶
type Encoder struct {
// contains filtered or unexported fields
}
Encoder is a byte-oriented API that is useful where wrapping io.Writer is not desirable.
Each non-empty Encode call returns one compressed Zstandard frame and appends one entry to the in-memory seek table. EndStream returns the final seek-table skippable frame, which must be appended after all encoded frames to form a complete seekable stream. EndStream finalizes the encoder. After EndStream, Encode and EndStream return ErrClosed.
func NewEncoder ¶
func NewEncoder(encoder ZSTDEncoder, opts ...WriterOption) (*Encoder, error)
NewEncoder returns a byte-oriented encoder that uses encoder for Zstandard compression.
The caller remains responsible for closing encoder, if it requires closing.
type FrameOffsetEntry ¶ added in v0.10.0
type FrameOffsetEntry struct {
// ID is the sequence number of the frame in the index.
ID int64
// CompressedOffset is the offset within the compressed stream.
CompressedOffset uint64
// DecompressedOffset is the offset within the decompressed stream.
DecompressedOffset uint64
// CompressedSize is the size of the compressed frame.
CompressedSize uint32
// DecompressedSize is the size of the frame's decompressed data.
DecompressedSize uint32
// Checksum is the lower 32 bits of the XXH64 hash of the uncompressed data.
// It is meaningful only when SeekTable.HasChecksums reports true.
Checksum uint32
}
FrameOffsetEntry is the post-processed view of a seek table entry suitable for indexing.
func (FrameOffsetEntry) LogValue ¶ added in v0.10.0
func (o FrameOffsetEntry) LogValue() slog.Value
LogValue implements slog.LogValuer.
type FrameSource ¶ added in v0.7.3
FrameSource returns one frame of data at a time.
When there are no more frames, it returns nil, nil. A non-nil error stops the write. Empty frames are ignored by Writer.WriteMany.
WriteMany may retain each returned slice until WriteMany returns. FrameSource implementations must not mutate or reuse returned slices before then.
type Reader ¶
type Reader struct {
// contains filtered or unexported fields
}
Reader provides sequential and random access to a seekable Zstandard stream and exposes its parsed seek-table metadata.
Offsets are expressed in the decompressed stream. Read and Seek use an internal current offset; ReadAt does not.
SeekTable may be called concurrently with any Reader method except Close. ReadAt may be called concurrently with other Reader methods when the decoder and read environment support concurrent use. Read and Seek share the current offset and must be serialized by the caller. Do not call Close concurrently with any Reader method.
func NewReader ¶
func NewReader(rs io.ReadSeeker, decoder ZSTDDecoder, opts ...ReaderOption) (*Reader, error)
NewReader returns a Zstandard stream reader that supports random access by decompressed offset.
The stream must end with a seek-table skippable frame. Unless WithReaderEnvironment supplies a read environment, rs must be non-nil. When NewReader uses rs directly and rs also implements io.ReaderAt, frame reads do not move rs's current offset.
The decoder must be non-nil. NewReader reads and validates the seek table during construction. Reader caches one decoded frame by default; use WithReaderFrameCache to change or disable caching.
The caller remains responsible for closing rs and decoder, if they require closing.
Example ¶
package main
import (
"bytes"
"fmt"
"io"
"log"
"github.com/klauspost/compress/zstd"
seekable "github.com/SaveTheRbtz/zstd-seekable-format-go/pkg"
)
func exampleFrames() [][]byte {
return [][]byte{[]byte("Hello"), []byte(" "), []byte("World!")}
}
func writeExampleFrames(w *seekable.Writer) {
for _, frame := range exampleFrames() {
if _, err := w.Write(frame); err != nil {
log.Fatal(err)
}
}
}
func exampleSeekableStream() []byte {
var buf bytes.Buffer
enc, err := zstd.NewWriter(nil, zstd.WithEncoderLevel(zstd.SpeedFastest))
if err != nil {
log.Fatal(err)
}
defer func() {
if err := enc.Close(); err != nil {
log.Fatal(err)
}
}()
w, err := seekable.NewWriter(&buf, enc)
if err != nil {
log.Fatal(err)
}
writeExampleFrames(w)
if err := w.Close(); err != nil {
log.Fatal(err)
}
return buf.Bytes()
}
func main() {
compressed := exampleSeekableStream()
dec, err := zstd.NewReader(nil)
if err != nil {
log.Fatal(err)
}
defer dec.Close()
r, err := seekable.NewReader(bytes.NewReader(compressed), dec)
if err != nil {
log.Fatal(err)
}
defer func() {
_ = r.Close()
}()
ello := make([]byte, 4)
if _, err := r.ReadAt(ello, 1); err != nil {
log.Fatal(err)
}
fmt.Println(string(ello))
world := make([]byte, 5)
if _, err := r.Seek(-6, io.SeekEnd); err != nil {
log.Fatal(err)
}
if _, err := r.Read(world); err != nil {
log.Fatal(err)
}
fmt.Println(string(world))
}
Output: ello World
func (*Reader) Close ¶
Close releases Reader-owned resources.
Close is idempotent. After Close, Read, ReadAt, Seek, and SeekTable return ErrClosed. Close does not close the io.ReadSeeker, decoder, or custom read environment passed to NewReader.
func (*Reader) Read ¶
Read reads decompressed bytes from the Reader's current offset and advances the current offset by the bytes read.
func (*Reader) ReadAt ¶
ReadAt reads len(p) decompressed bytes starting at off.
ReadAt does not move the Reader's current offset.
ReadAt follows io.ReaderAt EOF behavior. For non-empty p, it returns io.EOF when off is at or past the decompressed size, or with the bytes read when p extends past the end. Before Close, a zero-length p returns 0, nil.
Before Close, ReadAt may be called concurrently if the supplied decoder and read environment support concurrent use.
func (*Reader) SeekTable ¶ added in v0.10.0
SeekTable returns the parsed seek table for this Reader.
The returned SeekTable is immutable. SeekTable returns ErrClosed after Close. Before Close, SeekTable may be called concurrently with other Reader methods except Close.
Example ¶
package main
import (
"bytes"
"fmt"
"log"
"github.com/klauspost/compress/zstd"
seekable "github.com/SaveTheRbtz/zstd-seekable-format-go/pkg"
)
func exampleFrames() [][]byte {
return [][]byte{[]byte("Hello"), []byte(" "), []byte("World!")}
}
func writeExampleFrames(w *seekable.Writer) {
for _, frame := range exampleFrames() {
if _, err := w.Write(frame); err != nil {
log.Fatal(err)
}
}
}
func exampleSeekableStream() []byte {
var buf bytes.Buffer
enc, err := zstd.NewWriter(nil, zstd.WithEncoderLevel(zstd.SpeedFastest))
if err != nil {
log.Fatal(err)
}
defer func() {
if err := enc.Close(); err != nil {
log.Fatal(err)
}
}()
w, err := seekable.NewWriter(&buf, enc)
if err != nil {
log.Fatal(err)
}
writeExampleFrames(w)
if err := w.Close(); err != nil {
log.Fatal(err)
}
return buf.Bytes()
}
func main() {
compressed := exampleSeekableStream()
dec, err := zstd.NewReader(nil)
if err != nil {
log.Fatal(err)
}
defer dec.Close()
r, err := seekable.NewReader(bytes.NewReader(compressed), dec)
if err != nil {
log.Fatal(err)
}
defer func() {
_ = r.Close()
}()
table, err := r.SeekTable()
if err != nil {
log.Fatal(err)
}
entry, ok := table.EntryByDecompressedOffset(6)
if !ok {
log.Fatal("missing seek-table entry")
}
fmt.Printf("frames=%d size=%d checksums=%t\n", table.NumFrames(), table.Size(), table.HasChecksums())
fmt.Printf("offset 6 is in frame %d\n", entry.ID)
}
Output: frames=3 size=12 checksums=true offset 6 is in frame 2
type ReaderEnvironment ¶ added in v0.10.0
type ReaderEnvironment interface {
// GetFrameByIndex returns one complete compressed frame by seek-table entry.
//
// The returned slice must contain exactly index.CompressedSize bytes starting
// at index.CompressedOffset in the compressed stream.
//
// Reader may call GetFrameByIndex concurrently from concurrent ReadAt calls.
// Implementations used that way must support concurrent calls and must not
// mutate returned slices after returning them.
GetFrameByIndex(index FrameOffsetEntry) ([]byte, error)
ReadFooter() ([]byte, error)
// ReadSkipFrame returns the full final seek-table skippable frame.
//
// skippableFrameOffset is the number of bytes from the end of the stream
// back to the start of the final skippable frame. The returned bytes must
// include the Skippable_Magic_Number and Frame_Size fields.
ReadSkipFrame(skippableFrameOffset int64) ([]byte, error)
}
ReaderEnvironment is an advanced hook for custom storage implementations. It can be used to read frames and seek tables from somewhere other than a normal io.ReadSeeker.
type ReaderOption ¶ added in v0.10.0
ReaderOption configures NewReader.
func WithReaderEnvironment ¶ added in v0.10.0
func WithReaderEnvironment(e ReaderEnvironment) ReaderOption
WithReaderEnvironment sets a custom read environment for advanced storage implementations.
When this option is supplied, NewReader uses e instead of the io.ReadSeeker argument for all seek-table and frame reads.
func WithReaderFrameCache ¶ added in v0.10.0
func WithReaderFrameCache(cache framecache.Cache) ReaderOption
WithReaderFrameCache sets the Reader's decoded-frame cache.
A nil cache selects the default one-frame FIFO cache. To disable caching, use framecache.NewFIFO(framecache.Limits{MaxFrames: 0}).
If NewReader succeeds, the Reader owns the cache, clearing it before use and on Close. Callers must not use the cache directly or share it with another Reader.
func WithReaderLogger ¶ added in v0.10.0
func WithReaderLogger(l *slog.Logger) ReaderOption
WithReaderLogger sets the logger used by Reader internals.
Passing nil restores the default discard logger.
type SeekTable ¶ added in v0.10.0
type SeekTable struct {
// contains filtered or unexported fields
}
SeekTable is parsed random-access metadata from a Zstandard seek-table skippable frame.
Use NewSeekTable to construct a SeekTable from bytes written through WriterEnvironment.WriteSeekTable or returned by Encoder.EndStream. Lookup methods can be used concurrently.
func NewSeekTable ¶ added in v0.10.0
NewSeekTable parses a seek-table skippable frame.
buf must contain the final seek-table skippable frame itself, including the skippable-frame magic number and frame-size header. This is the byte sequence returned by Encoder.EndStream or passed to WriterEnvironment.WriteSeekTable, not the whole compressed stream.
Example ¶
package main
import (
"fmt"
"log"
"github.com/klauspost/compress/zstd"
seekable "github.com/SaveTheRbtz/zstd-seekable-format-go/pkg"
)
func main() {
enc, err := zstd.NewWriter(nil, zstd.WithEncoderLevel(zstd.SpeedFastest))
if err != nil {
log.Fatal(err)
}
defer func() {
if err := enc.Close(); err != nil {
log.Fatal(err)
}
}()
e, err := seekable.NewEncoder(enc)
if err != nil {
log.Fatal(err)
}
if _, err := e.Encode([]byte("Hello")); err != nil {
log.Fatal(err)
}
if _, err := e.Encode([]byte(" World!")); err != nil {
log.Fatal(err)
}
seekTableFrame, err := e.EndStream()
if err != nil {
log.Fatal(err)
}
table, err := seekable.NewSeekTable(seekTableFrame)
if err != nil {
log.Fatal(err)
}
entry, ok := table.EntryByDecompressedOffset(7)
if !ok {
log.Fatal("missing seek-table entry")
}
fmt.Printf("frames=%d size=%d checksums=%t\n", table.NumFrames(), table.Size(), table.HasChecksums())
fmt.Printf("offset 7 is in frame %d\n", entry.ID)
}
Output: frames=2 size=12 checksums=true offset 7 is in frame 1
func (SeekTable) EntryByDecompressedOffset ¶ added in v0.10.0
func (t SeekTable) EntryByDecompressedOffset(off uint64) (FrameOffsetEntry, bool)
EntryByDecompressedOffset returns the frame containing off in the decompressed stream. It returns false if off is greater than or equal to Size().
func (SeekTable) EntryByID ¶ added in v0.10.0
func (t SeekTable) EntryByID(id int64) (FrameOffsetEntry, bool)
EntryByID returns the frame with id. It returns false if id is greater than or equal to NumFrames() or less than 0.
func (SeekTable) HasChecksums ¶ added in v0.10.0
HasChecksums reports whether entries in the seek table include checksum fields.
type WriteManyOption ¶ added in v0.7.3
type WriteManyOption func(options *writeManyOptions) error
WriteManyOption configures Writer.WriteMany.
func WithConcurrency ¶ added in v0.7.3
func WithConcurrency(concurrency int) WriteManyOption
WithConcurrency sets the maximum number of concurrent frame encoding operations.
The default is runtime.GOMAXPROCS(0).
func WithWriteCallback ¶ added in v0.7.3
func WithWriteCallback(cb func(entry FrameOffsetEntry)) WriteManyOption
WithWriteCallback calls cb after each non-empty frame is written.
cb receives the seek-table entry for the frame that was just written. It is called in stream order from the WriteMany writer goroutine, after the frame has been written and added to the Writer's in-memory seek table.
cb must not call methods on the same Writer. To stop WriteMany from cb, cancel the WriteMany context or make the FrameSource return an error.
type Writer ¶
type Writer struct {
// contains filtered or unexported fields
}
Writer writes a seekable Zstandard stream.
Each non-empty Write call becomes one Zstandard frame in the output stream. Close must be called to write the final seek-table skippable frame; without it, Reader and NewSeekTable cannot find the random-access metadata. Close is idempotent. Write and WriteMany return ErrClosed after Close.
func NewWriter ¶
func NewWriter(w io.Writer, encoder ZSTDEncoder, opts ...WriterOption) (*Writer, error)
NewWriter wraps w and encoder into an indexed Zstandard stream.
w must be non-nil unless WithWriterEnvironment supplies a custom write environment. The caller remains responsible for closing w and encoder, if they require closing.
The resulting stream can be randomly accessed through Reader or NewSeekTable.
Example ¶
package main
import (
"bytes"
"fmt"
"io"
"log"
"github.com/klauspost/compress/zstd"
seekable "github.com/SaveTheRbtz/zstd-seekable-format-go/pkg"
)
func exampleFrames() [][]byte {
return [][]byte{[]byte("Hello"), []byte(" "), []byte("World!")}
}
func writeExampleFrames(w *seekable.Writer) {
for _, frame := range exampleFrames() {
if _, err := w.Write(frame); err != nil {
log.Fatal(err)
}
}
}
func exampleSeekableStream() []byte {
var buf bytes.Buffer
enc, err := zstd.NewWriter(nil, zstd.WithEncoderLevel(zstd.SpeedFastest))
if err != nil {
log.Fatal(err)
}
defer func() {
if err := enc.Close(); err != nil {
log.Fatal(err)
}
}()
w, err := seekable.NewWriter(&buf, enc)
if err != nil {
log.Fatal(err)
}
writeExampleFrames(w)
if err := w.Close(); err != nil {
log.Fatal(err)
}
return buf.Bytes()
}
func main() {
compressed := exampleSeekableStream()
dec, err := zstd.NewReader(bytes.NewReader(compressed))
if err != nil {
log.Fatal(err)
}
defer dec.Close()
all, err := io.ReadAll(dec)
if err != nil {
log.Fatal(err)
}
fmt.Println(string(all))
}
Output: Hello World!
func (*Writer) Close ¶
Close implements io.Closer. It writes the seek table, releases the in-memory frame index, and causes future Writer method calls to fail.
The caller is still responsible for closing the underlying writer.
func (*Writer) Write ¶
Write writes a chunk of data as a separate frame into the data stream.
Note that Write does not do any coalescing nor splitting of data, so each non-empty write will map to a separate Zstandard frame. Empty writes do not add seek-table entries.
If the underlying frame write fails or writes only part of the frame, the writer stops accepting more frames. Close may still be called to write the seek table for frames that were fully written. Bytes already accepted by the underlying writer are not rolled back.
func (*Writer) WriteMany ¶ added in v0.10.0
func (s *Writer) WriteMany(ctx context.Context, frameSource FrameSource, options ...WriteManyOption) error
WriteMany writes many frames concurrently.
It reads frames from frameSource sequentially, compresses up to the configured concurrency in parallel, and writes compressed frames in the same order returned by frameSource. Close must still be called after a successful WriteMany call to write the final seek table. Frame write failures have the same no-more-frames behavior as Writer.Write.
Example ¶
package main
import (
"bytes"
"context"
"fmt"
"log"
"github.com/klauspost/compress/zstd"
seekable "github.com/SaveTheRbtz/zstd-seekable-format-go/pkg"
)
func exampleFrames() [][]byte {
return [][]byte{[]byte("Hello"), []byte(" "), []byte("World!")}
}
func main() {
var buf bytes.Buffer
enc, err := zstd.NewWriter(nil, zstd.WithEncoderLevel(zstd.SpeedFastest))
if err != nil {
log.Fatal(err)
}
defer func() {
if err := enc.Close(); err != nil {
log.Fatal(err)
}
}()
w, err := seekable.NewWriter(&buf, enc)
if err != nil {
log.Fatal(err)
}
frames := exampleFrames()
next := func() ([]byte, error) {
if len(frames) == 0 {
return nil, nil
}
frame := frames[0]
frames = frames[1:]
return frame, nil
}
err = w.WriteMany(context.Background(), next,
seekable.WithConcurrency(2),
seekable.WithWriteCallback(func(entry seekable.FrameOffsetEntry) {
fmt.Printf("%d %d\n", entry.ID, entry.DecompressedSize)
}),
)
if err != nil {
log.Fatal(err)
}
if err := w.Close(); err != nil {
log.Fatal(err)
}
}
Output: 0 5 1 1 2 6
type WriterEnvironment ¶ added in v0.10.0
type WriterEnvironment interface {
// WriteFrame writes one complete compressed Zstandard frame.
//
// It is called once for each non-empty frame, in stream order. It should
// follow io.Writer conventions and return len(p), nil after a complete write.
WriteFrame(p []byte) (n int, err error)
// WriteSeekTable writes the final seek-table skippable frame.
//
// p includes the skippable-frame magic number, frame-size field, seek-table
// entries, and seek-table footer. It should follow io.Writer conventions and
// return len(p), nil after a complete write.
WriteSeekTable(p []byte) (n int, err error)
}
WriterEnvironment is an advanced hook for custom storage implementations. It can be used to write frames and seek tables somewhere other than a normal io.Writer.
type WriterOption ¶ added in v0.10.0
WriterOption configures NewWriter and NewEncoder. Options that configure output environments apply only to NewWriter.
func WithWriterEnvironment ¶ added in v0.10.0
func WithWriterEnvironment(e WriterEnvironment) WriterOption
WithWriterEnvironment sets a custom write environment for advanced storage implementations.
When this option is supplied to NewWriter, NewWriter uses e instead of the io.Writer argument for all frame and seek-table writes. NewEncoder returns compressed frames directly, so WithWriterEnvironment has no effect there.
func WithWriterLogger ¶ added in v0.10.0
func WithWriterLogger(l *slog.Logger) WriterOption
WithWriterLogger sets the logger used by Writer and Encoder internals.
Passing nil restores the default discard logger.
type ZSTDDecoder ¶
ZSTDDecoder is the decompressor.
It is compatible with the DecodeAll method provided by github.com/klauspost/compress/zstd.
Reader may call DecodeAll concurrently from concurrent ReadAt calls. Decoders used that way must support concurrent DecodeAll calls.
type ZSTDEncoder ¶
ZSTDEncoder is the compressor.
It is compatible with the EncodeAll method provided by github.com/klauspost/compress/zstd.
Writer.WriteMany may call EncodeAll concurrently. Encoders used with WriteMany must support concurrent EncodeAll calls.
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
Package framecache provides decoded-frame cache implementations for seekable readers.
|
Package framecache provides decoded-frame cache implementations for seekable readers. |