Vendor container-libs with PestoSocketPath interface method

- go.podman.io/storage@main
- go.podman.io/image/v5@main
- go.podman.io/common@main

Fixes: https://github.com/containers/podman/issues/29032

Signed-off-by: Jan Rodák <hony.com@seznam.cz>
This commit is contained in:
Jan Rodák
2026-06-29 11:40:44 +02:00
parent e4271bea59
commit 7ab215c0be
53 changed files with 4573 additions and 521 deletions

8
go.mod
View File

@@ -65,9 +65,9 @@ require (
github.com/vishvananda/netlink v1.3.1
go.etcd.io/bbolt v1.5.0
go.podman.io/buildah v1.44.0
go.podman.io/common v0.68.1-0.20260629144122-498ccd1bb53a
go.podman.io/image/v5 v5.40.1-0.20260629144122-498ccd1bb53a
go.podman.io/storage v1.63.1-0.20260629144122-498ccd1bb53a
go.podman.io/common v0.68.1-0.20260707152203-d126613c6575
go.podman.io/image/v5 v5.40.1-0.20260707152203-d126613c6575
go.podman.io/storage v1.63.1-0.20260707152203-d126613c6575
golang.org/x/crypto v0.53.0
golang.org/x/net v0.56.0
golang.org/x/sync v0.21.0
@@ -130,7 +130,7 @@ require (
github.com/hashicorp/go-cleanhttp v0.5.2 // indirect
github.com/hashicorp/go-retryablehttp v0.7.8 // indirect
github.com/inconshreveable/mousetrap v1.1.0 // indirect
github.com/klauspost/compress v1.18.6 // indirect
github.com/klauspost/compress v1.19.0 // indirect
github.com/kr/fs v0.1.0 // indirect
github.com/lufia/plan9stats v0.0.0-20240909124753-873cd0166683 // indirect
github.com/manifoldco/promptui v0.9.0 // indirect

16
go.sum
View File

@@ -209,8 +209,8 @@ github.com/kevinburke/ssh_config v1.5.0 h1:3cPZmE54xb5j3G5xQCjSvokqNwU2uW+3ry1+P
github.com/kevinburke/ssh_config v1.5.0/go.mod h1:q2RIzfka+BXARoNexmF9gkxEX7DmvbW9P4hIVx2Kg4M=
github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8=
github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck=
github.com/klauspost/compress v1.18.6 h1:2jupLlAwFm95+YDR+NwD2MEfFO9d4z4Prjl1XXDjuao=
github.com/klauspost/compress v1.18.6/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ=
github.com/klauspost/compress v1.19.0 h1:sXLILfc9jV2QYWkzFOPWStmcUVH2RHEB1JCdY2oVvCQ=
github.com/klauspost/compress v1.19.0/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ=
github.com/klauspost/pgzip v1.2.6 h1:8RXeL5crjEUFnR2/Sn6GJNWtSQ3Dk8pq4CL3jvdDyjU=
github.com/klauspost/pgzip v1.2.6/go.mod h1:Ch1tH69qFZu15pkjo5kYi6mth2Zzwzt50oCQKQE9RUs=
github.com/kr/fs v0.1.0 h1:Jskdu9ieNAYnjxsi0LbQp1ulIKZV1LAFgK1tWhpZgl8=
@@ -435,12 +435,12 @@ go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/
go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE=
go.podman.io/buildah v1.44.0 h1:hxQh/fZtm/r1Xuo2UnDoQH0PRafCRJ53oUeDSDCtKkE=
go.podman.io/buildah v1.44.0/go.mod h1:RUdF9lLmMmdD7K9Y6eO64ptYuUMdDZ84SkPJc4008Jc=
go.podman.io/common v0.68.1-0.20260629144122-498ccd1bb53a h1:uNlfd8m5TzyUdEEn/yK9k9pwpDRCVxMvVX+FsFSs9bs=
go.podman.io/common v0.68.1-0.20260629144122-498ccd1bb53a/go.mod h1:lbL8J2uE+YiFpRWhRotoxEm+68hIEgQU9XDGVe9DjZs=
go.podman.io/image/v5 v5.40.1-0.20260629144122-498ccd1bb53a h1:V07RRMju7VdvuHO/ZGfBJLdH18CkKPu8AdZf7NIbvAY=
go.podman.io/image/v5 v5.40.1-0.20260629144122-498ccd1bb53a/go.mod h1:9nU7fQ98VdpBpoTMIYZo+QsIW0uppqCTuJ3LjswmZXo=
go.podman.io/storage v1.63.1-0.20260629144122-498ccd1bb53a h1:LBNdyjEPpEzFuO/yYgNgAHXcCwP9yKEtLdIiY+o+IHE=
go.podman.io/storage v1.63.1-0.20260629144122-498ccd1bb53a/go.mod h1:a7ch5yo/9ygYk5oaO9IAz4LXJ2OPRe51w4NWnpgpRBU=
go.podman.io/common v0.68.1-0.20260707152203-d126613c6575 h1:TrtNYc0jj7bLzkwWybR2gZG/EhcpF8zc0brFMkP8omg=
go.podman.io/common v0.68.1-0.20260707152203-d126613c6575/go.mod h1:wwuevsG6ke58aALZp21G0wD6Jlm2WtBf5XJqRag6KqY=
go.podman.io/image/v5 v5.40.1-0.20260707152203-d126613c6575 h1:UWxDEQxYsiZaZPLwtciHpAnuvMS/q9f1xYuLc98Gc6c=
go.podman.io/image/v5 v5.40.1-0.20260707152203-d126613c6575/go.mod h1:DEGFX12gtx0ZxaElL7KYyJrvpGFYieNzvFoFE0Hc8Yk=
go.podman.io/storage v1.63.1-0.20260707152203-d126613c6575 h1:mVfzQLVJZDjsaK3yr2yI5dWVec4EQnLUNVHu8nEIId8=
go.podman.io/storage v1.63.1-0.20260707152203-d126613c6575/go.mod h1:LFeOixveujX5jTT0H8VuIOAED6KpPGFFwVMXMBxZPeI=
go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ=
go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ=
go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc=

View File

@@ -24,11 +24,7 @@ import (
)
func (r *Runtime) pestoSocketPath() string {
info, err := r.network.RootlessNetnsInfo()
if err != nil || info == nil {
return ""
}
return info.PestoSocketPath
return r.network.PestoSocketPath()
}
// setupRootlessPortMappingViaPesto adds this container's port forwarding

View File

@@ -28,9 +28,10 @@ type dictDecoder struct {
hist []byte // Sliding window history
// Invariant: 0 <= rdPos <= wrPos <= len(hist)
wrPos int // Current output position in buffer
rdPos int // Have emitted hist[:rdPos] already
full bool // Has a full window length been written yet?
wrPos int // Current output position in buffer
rdPos int // Have emitted hist[:rdPos] already
flushed int64 // Total bytes returned by readFlush since init
full bool // Has a full window length been written yet?
}
// init initializes dictDecoder to have a sliding window dictionary of the given
@@ -167,11 +168,22 @@ loop:
return dstPos - dstBase
}
// appendWindow appends the current sliding window (up to len(hist) most recent
// bytes, oldest first) to dst.
func (dd *dictDecoder) appendWindow(dst []byte) []byte {
if dd.full {
dst = append(dst, dd.hist[dd.wrPos:]...)
return append(dst, dd.hist[:dd.wrPos]...)
}
return append(dst, dd.hist[:dd.wrPos]...)
}
// readFlush returns a slice of the historical buffer that is ready to be
// emitted to the user. The data returned by readFlush must be fully consumed
// before calling any other dictDecoder methods.
func (dd *dictDecoder) readFlush() []byte {
toRead := dd.hist[dd.rdPos:dd.wrPos]
dd.flushed += int64(len(toRead))
dd.rdPos = dd.wrPos
if dd.wrPos == len(dd.hist) {
dd.wrPos, dd.rdPos = 0, 0
@@ -179,3 +191,9 @@ func (dd *dictDecoder) readFlush() []byte {
}
return toRead
}
// decoded reports the total number of bytes written into the dictionary since
// init (i.e. excluding any preset dict bytes).
func (dd *dictDecoder) decoded() int64 {
return dd.flushed + int64(dd.wrPos-dd.rdPos)
}

View File

@@ -342,6 +342,11 @@ type decompressor struct {
final bool
flushMode flushMode
cb func(InflateCheckpoint)
cp InflateCheckpoint
hasCP bool // WithResumeFrom was supplied
uncOffset int64 // baseline uncompressed offset (from a resume checkpoint)
cpBuf []byte
}
func (f *decompressor) nextBlock() {
@@ -676,6 +681,18 @@ func (f *decompressor) finishBlock() {
f.toRead = f.dict.readFlush()
}
if f.cb != nil {
bitPos := f.roffset*8 - int64(f.nb)
f.cpBuf = f.dict.appendWindow(f.cpBuf[:0])
f.cb(InflateCheckpoint{
UncompressedOffset: f.uncOffset + f.dict.decoded(),
CompressedOffset: bitPos / 8,
Final: f.final,
BitOffset: uint8(bitPos & 7),
Window: f.cpBuf,
})
}
f.step = nextBlock
}
@@ -806,6 +823,45 @@ func (f *decompressor) Reset(r io.Reader, dict []byte) error {
return nil
}
// ResetCP will adjust the input to the provided checkpoint.
// It is assumed the input stream is forwarded to cp.CompressedOffset.
func (f *decompressor) ResetCP(r io.Reader, cp InflateCheckpoint) error {
*f = decompressor{
r: makeReader(r),
bits: f.bits,
codebits: f.codebits,
h1: f.h1,
h2: f.h2,
dict: f.dict,
step: nextBlock,
cpBuf: f.cpBuf,
}
return f.applyCP(cp)
}
// applyCP seeds the decompressor state from a resume checkpoint:
// loads the sliding window, sets the absolute compressed/uncompressed
// offsets, and skips cp.BitOffset bits into the first input byte so
// the next decode aligns with the start of a deflate block.
func (f *decompressor) applyCP(cp InflateCheckpoint) error {
f.dict.init(maxMatchOffset, cp.Window)
f.roffset = cp.CompressedOffset
f.uncOffset = cp.UncompressedOffset
f.final = cp.Final
f.b = 0
f.nb = 0
if cp.BitOffset > 0 {
c, err := f.r.ReadByte()
if err != nil {
return noEOF(err)
}
f.roffset++
f.b = uint32(c) >> cp.BitOffset
f.nb = 8 - uint(cp.BitOffset)
}
return nil
}
type ReaderOpt func(*decompressor)
// WithPartialBlock tells decompressor to return after each block,
@@ -823,6 +879,36 @@ func WithDict(dict []byte) ReaderOpt {
}
}
// InflateCheckpoint provides a resumable checkpoint for inflate.
type InflateCheckpoint struct {
UncompressedOffset int64 // Byte offset in the decompressed stream
CompressedOffset int64 // Byte offset in the compressed stream
Final bool // True if this is the final block
BitOffset uint8 // 0-7 bits
Window []byte // 32KB sliding window dictionary
}
// WithEobCallback will call the provided function after each block
// with the current gzip checkpoint.
// After returning the provided window can no longer be referenced.
// The callback will not be triggered after a block is marked "final".
// The callback is not retained after Reset.
func WithEobCallback(cb func(InflateCheckpoint)) ReaderOpt {
return func(f *decompressor) {
f.cb = cb
}
}
// WithResumeFrom will adjust the input to the provided checkpoint.
// It is assumed the input stream is forwarded to the provided offset.
// The checkpoint is removed when Reset is called.
func WithResumeFrom(cp InflateCheckpoint) ReaderOpt {
return func(f *decompressor) {
f.cp = cp
f.hasCP = true
}
}
// NewReaderOpts returns new reader with provided options
func NewReaderOpts(r io.Reader, opts ...ReaderOpt) io.ReadCloser {
fixedHuffmanDecoderInit()
@@ -838,6 +924,12 @@ func NewReaderOpts(r io.Reader, opts ...ReaderOpt) io.ReadCloser {
opt(&f)
}
if f.hasCP {
if err := f.applyCP(f.cp); err != nil {
f.err = err
}
}
return &f
}

View File

@@ -0,0 +1,168 @@
package huff0
import "errors"
// BuildCTable builds a Huffman compression table from a precomputed symbol
// histogram and installs it as the previous (reuse) table on s.
//
// After this call:
// - EstimateSize/CanUseTable can probe the table against other histograms.
// - Compress1X/Compress4X with Reuse = ReusePolicyMust will encode without
// emitting a new table header.
// - TransferCTable can hand the table to a sibling Scratch.
//
// count[i] is the number of occurrences of symbol i. The histogram must have
// at least 2 distinct non-zero symbols; ErrUseRLE is returned for a single
// symbol and an error is returned for an empty histogram.
func (s *Scratch) BuildCTable(count *[256]uint32) error {
if s == nil {
return errors.New("huff0: BuildCTable on nil Scratch")
}
if count == nil {
return errors.New("huff0: nil count passed to BuildCTable")
}
var err error
s, err = s.prepare(nil)
if err != nil {
return err
}
s.count = *count
var total, maxCount int
var symLen uint16
for i, v := range s.count {
total += int(v)
if int(v) > maxCount {
maxCount = int(v)
}
if v != 0 {
symLen = uint16(i) + 1
}
}
if total == 0 {
return errors.New("huff0: empty histogram")
}
if symLen < 2 || maxCount == total {
return ErrUseRLE
}
// huff0's internal rank table assumes total ≤ BlockSizeMax (it uses
// highBit32(count+1) + 1 as a rank index into a fixed-size array).
// Histograms summed across multiple blocks can exceed that; scale the
// counts down preserving the distribution. Non-zero entries round up so
// rare symbols stay representable.
if total > BlockSizeMax {
shift := uint(0)
for total>>shift > BlockSizeMax {
shift++
}
round := uint32(1<<shift) - 1
var newTotal, newMax int
for i, v := range s.count {
if v == 0 {
continue
}
scaled := (v + round) >> shift
if scaled == 0 {
scaled = 1
}
s.count[i] = scaled
newTotal += int(scaled)
if int(scaled) > newMax {
newMax = int(scaled)
}
}
total = newTotal
maxCount = newMax
if maxCount == total {
return ErrUseRLE
}
}
s.symbolLen = symLen
s.maxCount = maxCount
s.srcLen = total
if err := s.buildCTable(); err != nil {
return err
}
if cap(s.prevTable) < len(s.cTable) {
s.prevTable = make(cTable, 0, maxSymbolValue+1)
}
s.prevTable = s.prevTable[:len(s.cTable)]
copy(s.prevTable, s.cTable)
s.prevTableLog = s.actualTableLog
// Force the next Compress* to recount from real input.
s.clearCount = true
s.maxCount = 0
return nil
}
// EstimateSize returns an estimated compressed payload size in bytes for the
// supplied histogram using the table currently stored in prevTable. It returns
// -1 when the table cannot encode every non-zero symbol of hist (i.e. when
// CanUseTable would return false). The estimate excludes the table header.
func (s *Scratch) EstimateSize(hist *[256]uint32) int {
if s == nil || hist == nil || len(s.prevTable) == 0 {
return -1
}
pt := s.prevTable
nbBits := uint32(7)
for i, v := range hist {
if v == 0 {
continue
}
if i >= len(pt) || pt[i].nBits == 0 {
return -1
}
nbBits += uint32(pt[i].nBits) * v
}
return int(nbBits >> 3)
}
// CanUseTable reports whether the table in prevTable can encode every
// non-zero symbol present in hist.
func (s *Scratch) CanUseTable(hist *[256]uint32) bool {
if s == nil || hist == nil || len(s.prevTable) == 0 {
return false
}
pt := s.prevTable
for i, v := range hist {
if v == 0 {
continue
}
if i >= len(pt) || pt[i].nBits == 0 {
return false
}
}
return true
}
// AppendTable serializes the table currently stored in prevTable (e.g. as
// installed by BuildCTable or carried over from a previous Compress call)
// into a self-delimiting zstd-style header and appends it to dst. The
// returned slice can be parsed back by ReadTable.
func (s *Scratch) AppendTable(dst []byte) ([]byte, error) {
if s == nil || len(s.prevTable) == 0 {
return dst, errors.New("huff0: AppendTable with empty table")
}
// cTable.write reads s.actualTableLog, s.symbolLen, s.huffWeight, s.fse
// and writes into s.Out. Save/restore Out so we don't disturb in-flight
// compression buffers.
saveOut := s.Out
saveTL := s.actualTableLog
saveSL := s.symbolLen
if s.fse == nil {
// Lazily init in case AppendTable is called on a fresh Scratch.
if _, err := s.prepare(nil); err != nil {
return dst, err
}
saveOut = s.Out
}
s.Out = s.Out[:0]
s.actualTableLog = s.prevTableLog
s.symbolLen = uint16(len(s.prevTable))
if err := s.prevTable.write(s); err != nil {
s.Out, s.actualTableLog, s.symbolLen = saveOut, saveTL, saveSL
return dst, err
}
dst = append(dst, s.Out...)
s.Out, s.actualTableLog, s.symbolLen = saveOut, saveTL, saveSL
return dst, nil
}

View File

@@ -31,7 +31,7 @@ func DecodedLen(src []byte) (int, error) {
// that the length header occupied.
func decodedLen(src []byte) (blockLen, headerLen int, err error) {
v, n := binary.Uvarint(src)
if n <= 0 || v > 0xffffffff {
if n <= 0 || n > 5 || v > 0xffffffff {
return 0, 0, ErrCorrupt
}

View File

@@ -75,14 +75,47 @@ The above is fine for big encodes. However, whenever possible try to *reuse* the
To reuse the encoder, you can use the `Reset(io.Writer)` function to change to another output.
This will allow the encoder to reuse all resources and avoid wasteful allocations.
Currently stream encoding has 'light' concurrency, meaning up to 2 goroutines can be working on part
of a stream. This is independent of the `WithEncoderConcurrency(n)`, but that is likely to change
By default, stream encoding has 'light' concurrency, meaning up to 2 goroutines can be working on part
of a stream. This is independent of the `WithEncoderConcurrency(n)`, but that is likely to change
in the future. So if you want to limit concurrency for future updates, specify the concurrency
you would like.
If you would like stream encoding to be done without spawning async goroutines, use `WithEncoderConcurrency(1)`
which will compress input as each block is completed, blocking on writes until each has completed.
#### Parallel Stream Compression
For maximum throughput on large streams, use `WithConcurrentBlocks(true)` together with
`WithEncoderConcurrency(n)` where n is the number of CPU cores you want to use.
This splits the input into large sections (jobs) that are compressed simultaneously by multiple goroutines,
similar to how the C zstd library does multithreaded compression.
```Go
enc, err := zstd.NewWriter(out,
zstd.WithEncoderLevel(zstd.SpeedDefault),
zstd.WithEncoderConcurrency(runtime.GOMAXPROCS(0)),
zstd.WithConcurrentBlocks(true),
)
```
Each non-first job receives an overlap prefix from the previous job for match context,
so compression ratio is only marginally affected. Output is flushed in order,
producing a valid single-frame zstd stream.
Benchmark on 1.8GB GOB stream (AMD Ryzen 9 9950X):
| Level | 1 thread | 4 threads | 16 threads | 1T ratio | 16T ratio |
|---------|:----------:|:------------------:|:-------------------:|:--------:|:---------:|
| fastest | 783 MB/s | 2950 MB/s (3.8×) | 6939 MB/s (8.9×) | 12.24% | 12.26% |
| default | 728 MB/s | 2533 MB/s (3.5×) | 5340 MB/s (7.3×) | 10.67% | 10.68% |
| better | 434 MB/s | 1105 MB/s (2.5×) | 2206 MB/s (5.1×) | 9.14% | 9.21% |
| best | 129 MB/s | 367 MB/s (2.8×) | 884 MB/s (6.8×) | 8.48% | 8.63% |
Notes:
* Not compatible with dictionary encoding.
* `Flush()` dispatches the current partial job, so latency-sensitive callers can force output.
* `EncodeAll` is unaffected — it uses its own concurrency via the encoder pool.
You can specify your desired compression level using `WithEncoderLevel()` option. Currently only pre-defined
compression settings can be specified.

View File

@@ -230,7 +230,7 @@ func BuildDict(o BuildDictOptions) ([]byte, error) {
}
block := blockEnc{lowMem: false}
block.init()
enc := encoder(&bestFastEncoder{fastBase: fastBase{maxMatchOff: int32(maxMatchLen), bufferReset: math.MaxInt32 - int32(maxMatchLen*2), lowMem: false}})
var enc encoder
if o.Level != 0 {
eOpts := encoderOptions{
level: o.Level,
@@ -242,6 +242,7 @@ func BuildDict(o BuildDictOptions) ([]byte, error) {
enc = eOpts.encoder()
} else {
o.Level = SpeedBestCompression
enc = encoder(&bestFastEncoder{fastBase: fastBase{maxMatchOff: int32(maxMatchLen), bufferReset: math.MaxInt32 - int32(maxMatchLen*2), lowMem: false}})
}
var (
remain [256]int

View File

@@ -128,6 +128,34 @@ func (e *fastBase) matchlen(s, t int32, src []byte) int32 {
return int32(matchLen(src[s:], src[t:]))
}
// resetBasePrefix resets the encoder state and loads prefix as initial history.
// This is used for parallel job encoding where non-first jobs need overlap context.
// Rep offsets are set to defaults [1,4,8] (invalidated, matching C behavior).
func (e *fastBase) resetBasePrefix(prefix []byte) {
if e.blk == nil {
e.blk = &blockEnc{lowMem: e.lowMem}
e.blk.init()
} else {
e.blk.reset(nil)
}
e.blk.initNewEncode()
if e.crc == nil {
e.crc = xxhash.New()
} else {
e.crc.Reset()
}
e.blk.dictLitEnc = nil
e.ensureHist(len(prefix) + maxCompressedBlockSize)
// Bump cur so old table entries fall outside the window.
// When cur >= bufferReset, leave it; the first Encode call
// will shift/clear tables, preserving valid prefix entries.
if e.cur < e.bufferReset {
e.cur += e.maxMatchOff + int32(len(e.hist))
}
e.hist = e.hist[:0]
e.hist = append(e.hist, prefix...)
}
// Reset the encoding table.
func (e *fastBase) resetBase(d *dict, singleBlock bool) {
if e.blk == nil {

View File

@@ -551,3 +551,18 @@ func (e *bestFastEncoder) Reset(d *dict, singleBlock bool) {
// Reset table to initial state
copy(e.table[:], e.dictTable)
}
func (e *bestFastEncoder) ResetPrefix(prefix []byte) {
e.resetBasePrefix(prefix)
if len(prefix) < 8 {
return
}
end := e.cur + int32(len(prefix)) - 8
for i := e.cur; i < end; i++ {
cv := load6432(prefix, i-e.cur)
h := hashLen(cv, bestLongTableBits, bestLongLen)
e.longTable[h] = prevEntry{offset: i, prev: e.longTable[h].offset}
h0 := hashLen(cv, bestShortTableBits, bestShortLen)
e.table[h0] = prevEntry{offset: i, prev: e.table[h0].offset}
}
}

View File

@@ -1096,6 +1096,20 @@ func (e *betterFastEncoder) Reset(d *dict, singleBlock bool) {
}
}
func (e *betterFastEncoder) ResetPrefix(prefix []byte) {
e.resetBasePrefix(prefix)
if len(prefix) < 8 {
return
}
end := e.cur + int32(len(prefix)) - 8
for i := e.cur; i < end; i += 2 {
cv := load6432(prefix, i-e.cur)
h := hashLen(cv, betterLongTableBits, betterLongLen)
e.longTable[h] = prevEntry{offset: i, prev: e.longTable[h].offset}
e.table[hashLen(cv>>8, betterShortTableBits, betterShortLen)] = tableEntry{val: uint32(cv >> 8), offset: i + 1}
}
}
// ResetDict will reset and set a dictionary if not nil
func (e *betterFastEncoderDict) Reset(d *dict, singleBlock bool) {
e.resetBase(d, singleBlock)
@@ -1229,6 +1243,10 @@ func (e *betterFastEncoderDict) Reset(d *dict, singleBlock bool) {
e.allDirty = false
}
func (e *betterFastEncoderDict) ResetPrefix([]byte) {
panic("ResetPrefix not supported for dict encoders")
}
func (e *betterFastEncoderDict) markLongShardDirty(entryNum uint32) {
e.longTableShardDirty[entryNum/betterLongTableShardSize] = true
}

View File

@@ -1037,6 +1037,18 @@ func (e *doubleFastEncoder) Reset(d *dict, singleBlock bool) {
}
}
func (e *doubleFastEncoder) ResetPrefix(prefix []byte) {
e.fastEncoder.ResetPrefix(prefix)
if len(prefix) < 8 {
return
}
end := e.cur + int32(len(prefix)) - 8
for i := e.cur + 1; i < end; i += 2 {
cv := load6432(prefix, i-e.cur)
e.longTable[hashLen(cv, dFastLongTableBits, dFastLongLen)] = tableEntry{val: uint32(cv), offset: i}
}
}
// ResetDict will reset and set a dictionary if not nil
func (e *doubleFastEncoderDict) Reset(d *dict, singleBlock bool) {
allDirty := e.allDirty
@@ -1102,6 +1114,10 @@ func (e *doubleFastEncoderDict) Reset(d *dict, singleBlock bool) {
}
}
func (e *doubleFastEncoderDict) ResetPrefix([]byte) {
panic("ResetPrefix not supported for dict encoders")
}
func (e *doubleFastEncoderDict) markLongShardDirty(entryNum uint32) {
e.longTableShardDirty[entryNum/dLongTableShardSize] = true
}

View File

@@ -797,6 +797,19 @@ func (e *fastEncoder) Reset(d *dict, singleBlock bool) {
}
}
func (e *fastEncoder) ResetPrefix(prefix []byte) {
e.resetBasePrefix(prefix)
if len(prefix) < 8 {
return
}
end := e.cur + int32(len(prefix)) - 8
// Index every 4th
for i := e.cur + 1; i < end; i += 4 {
cv := load6432(prefix, i-e.cur)
e.table[hashLen(cv, tableBits, tableFastHashLen)] = tableEntry{val: uint32(cv), offset: i}
}
}
// ResetDict will reset and set a dictionary if not nil
func (e *fastEncoderDict) Reset(d *dict, singleBlock bool) {
e.resetBase(d, singleBlock)
@@ -866,6 +879,10 @@ func (e *fastEncoderDict) Reset(d *dict, singleBlock bool) {
e.allDirty = false
}
func (e *fastEncoderDict) ResetPrefix([]byte) {
panic("ResetPrefix not supported for dict encoders")
}
func (e *fastEncoderDict) markAllShardsDirty() {
e.allDirty = true
}

352
vendor/github.com/klauspost/compress/zstd/enc_jobs.go generated vendored Normal file
View File

@@ -0,0 +1,352 @@
// Copyright 2019+ Klaus Post. All rights reserved.
// License information can be found in the LICENSE file.
// Based on work by Yann Collet, released under BSD License.
package zstd
import (
"fmt"
rdebug "runtime/debug"
"sync"
)
type encJob struct {
prefix []byte // overlap from previous job (nil for first)
input []byte // job's own input data (swapped from filling)
last bool // last block of last job gets last=true
output []byte // compressed blocks (filled by worker)
err error // encoding error
done chan struct{} // closed when complete
}
type jobState struct {
jobSize int
overlapSize int
filling []byte // accumulates input up to jobSize
nextPrefix []byte // overlap prefix prepared for the next dispatched job
jobSeq int // next job sequence number
jobCh chan *encJob // dispatch to workers
resultCh chan *encJob // ordered results to flusher
workerWg sync.WaitGroup
flusherWg sync.WaitGroup
mu sync.Mutex
flushedSeq int // last flushed sequence number
cond *sync.Cond
flusherErr error
started bool
inputPool sync.Pool // *[]byte buffers of jobSize cap
outputPool sync.Pool // *[]byte buffers for compressed output
overlapPool sync.Pool // *[]byte buffers for overlap prefixes
}
func (e *Encoder) startJobWorkers() {
js := &e.state.jobs
n := e.o.concurrent
js.jobCh = make(chan *encJob, n)
js.resultCh = make(chan *encJob, n)
js.flushedSeq = 0
js.cond = sync.NewCond(&js.mu)
// Workers borrow encoders from the shared e.encoders pool per-job.
// Ensure the pool is initialized before any worker tries to borrow.
e.init.Do(e.initialize)
for range n {
js.workerWg.Add(1)
go e.jobWorker()
}
js.flusherWg.Add(1)
go e.jobFlusher()
js.started = true
}
func (e *Encoder) jobWorker() {
js := &e.state.jobs
defer js.workerWg.Done()
for job := range js.jobCh {
enc := <-e.encoders
e.compressJob(enc, job)
e.encoders <- enc
close(job.done)
}
}
func (e *Encoder) compressJob(enc encoder, job *encJob) {
defer func() {
if r := recover(); r != nil {
job.err = fmt.Errorf("panic in parallel job: %v", r)
rdebug.PrintStack()
}
}()
if len(job.prefix) > 0 {
enc.ResetPrefix(job.prefix)
} else {
enc.Reset(nil, false)
}
data := job.input
if len(data) == 0 && job.last {
blk := enc.Block()
blk.reset(nil)
blk.last = true
blk.encodeRaw(nil)
job.output = append(job.output, blk.output...)
return
}
blk := enc.Block()
for len(data) > 0 {
todo := data
if len(todo) > e.o.blockSize {
todo = todo[:e.o.blockSize]
}
data = data[len(todo):]
blk.pushOffsets()
enc.Encode(blk, todo)
blk.last = len(data) == 0 && job.last
err := blk.encode(todo, e.o.noEntropy, !e.o.allLitEntropy)
if err != nil {
job.err = err
return
}
job.output = append(job.output, blk.output...)
blk.reset(nil)
}
}
func (js *jobState) getInputBuf(size int) []byte {
if v := js.inputPool.Get(); v != nil {
bp := v.(*[]byte)
b := *bp
if cap(b) >= size {
return b[:0]
}
}
return make([]byte, 0, size)
}
func (js *jobState) putInputBuf(b []byte) {
if cap(b) > 0 {
b = b[:0]
js.inputPool.Put(&b)
}
}
func (js *jobState) getOutputBuf(size int) []byte {
if v := js.outputPool.Get(); v != nil {
bp := v.(*[]byte)
b := *bp
if cap(b) >= size {
return b[:0]
}
}
return make([]byte, 0, size)
}
func (js *jobState) putOutputBuf(b []byte) {
if cap(b) > 0 {
b = b[:0]
js.outputPool.Put(&b)
}
}
func (js *jobState) getOverlapBuf(size int) []byte {
if v := js.overlapPool.Get(); v != nil {
bp := v.(*[]byte)
b := *bp
if cap(b) >= size {
return b[:size]
}
}
return make([]byte, size)
}
func (js *jobState) putOverlapBuf(b []byte) {
if cap(b) > 0 {
b = b[:0]
js.overlapPool.Put(&b)
}
}
func (e *Encoder) jobFlusher() {
js := &e.state.jobs
defer js.flusherWg.Done()
for job := range js.resultCh {
<-job.done
// Worker has fully exited compressJob, so the prefix is no longer
// in use. Return it to the pool regardless of outcome.
if job.prefix != nil {
js.putOverlapBuf(job.prefix)
job.prefix = nil
}
if job.err != nil {
js.mu.Lock()
js.flusherErr = job.err
js.cond.Broadcast()
js.mu.Unlock()
for range js.resultCh {
}
return
}
if len(job.output) > 0 {
_, err := e.state.w.Write(job.output)
if err != nil {
js.mu.Lock()
js.flusherErr = err
js.cond.Broadcast()
js.mu.Unlock()
for range js.resultCh {
}
return
}
e.state.nWritten += int64(len(job.output))
}
// Return buffers to pools.
js.putInputBuf(job.input)
js.putOutputBuf(job.output)
job.input = nil
job.output = nil
js.mu.Lock()
js.flushedSeq++
js.cond.Broadcast()
js.mu.Unlock()
}
}
func (e *Encoder) shutdownJobWorkers() {
js := &e.state.jobs
if !js.started {
return
}
close(js.jobCh)
js.workerWg.Wait()
close(js.resultCh)
js.flusherWg.Wait()
js.started = false
}
// waitAllJobs blocks until all dispatched jobs have been flushed.
func (e *Encoder) waitAllJobs() {
js := &e.state.jobs
if !js.started {
return
}
js.mu.Lock()
for js.flushedSeq < js.jobSeq && js.flusherErr == nil {
js.cond.Wait()
}
js.mu.Unlock()
}
func (e *Encoder) dispatchJob(final bool) error {
s := &e.state
js := &s.jobs
js.mu.Lock()
fErr := js.flusherErr
js.mu.Unlock()
if fErr != nil {
return fErr
}
if !s.headerWritten {
// Single-block optimization: fall through to encodeAll path.
if final && len(js.filling) > 0 && len(js.filling) <= e.o.blockSize {
s.current = e.encodeAll(s.encoder, js.filling, s.current[:0])
var n2 int
n2, s.err = s.w.Write(s.current)
if s.err != nil {
return s.err
}
s.nWritten += int64(n2)
s.nInput += int64(len(js.filling))
s.current = s.current[:0]
js.filling = js.filling[:0]
s.headerWritten = true
s.fullFrameWritten = true
s.eofWritten = true
return nil
}
if final && len(js.filling) == 0 && !e.o.fullZero {
s.headerWritten = true
s.fullFrameWritten = true
s.eofWritten = true
return nil
}
var tmp [maxHeaderSize]byte
fh := frameHeader{
ContentSize: uint64(s.frameContentSize),
WindowSize: uint32(s.encoder.WindowSize(s.frameContentSize)),
SingleSegment: false,
Checksum: e.o.crc,
DictID: 0,
}
dst := fh.appendTo(tmp[:0])
var n2 int
n2, s.err = s.w.Write(dst)
if s.err != nil {
return s.err
}
s.nWritten += int64(n2)
s.headerWritten = true
}
if len(js.filling) == 0 && !final {
return nil
}
if !js.started {
e.startJobWorkers()
}
// Estimate output size for pooled buffer.
outputEst := max(len(js.filling)/2, 512)
job := &encJob{
last: final,
done: make(chan struct{}),
output: js.getOutputBuf(outputEst),
}
// Each job owns its prefix slice; the flusher returns it to the pool
// after <-job.done, so workers and dispatch never share a buffer.
if js.nextPrefix != nil {
job.prefix = js.nextPrefix
js.nextPrefix = nil
}
// Build the next job's prefix from the tail of this job's input.
if !final && len(js.filling) > 0 {
overlapLen := min(js.overlapSize, len(js.filling))
np := js.getOverlapBuf(overlapLen)
copy(np, js.filling[len(js.filling)-overlapLen:])
js.nextPrefix = np
}
// Swap filling buffer into job — zero-copy for the input data.
job.input = js.filling
js.filling = js.getInputBuf(js.jobSize)
s.nInput += int64(len(job.input))
js.jobSeq++
if final {
s.eofWritten = true
}
js.resultCh <- job
js.jobCh <- job
return nil
}

View File

@@ -38,6 +38,7 @@ type encoder interface {
WindowSize(size int64) int32
UseBlock(*blockEnc)
Reset(d *dict, singleBlock bool)
ResetPrefix(prefix []byte)
}
type encoderState struct {
@@ -60,6 +61,9 @@ type encoderState struct {
wg sync.WaitGroup
// This waitgroup indicates we have a block encoding/writing.
wWg sync.WaitGroup
// Parallel job state (used when concurrentBlocks is enabled).
jobs jobState
}
// NewWriter will create a new Zstandard encoder.
@@ -74,6 +78,9 @@ func NewWriter(w io.Writer, opts ...EOption) (*Encoder, error) {
return nil, err
}
}
if e.o.concurrentBlocks && (e.o.dict != nil || e.o.concurrent <= 1) {
e.o.concurrentBlocks = false
}
if w != nil {
e.Reset(w)
}
@@ -95,12 +102,31 @@ func (e *Encoder) initialize() {
// as a new, independent stream.
func (e *Encoder) Reset(w io.Writer) {
s := &e.state
if e.o.concurrentBlocks {
e.shutdownJobWorkers()
js := &s.jobs
js.jobSize = e.o.jobSize()
js.overlapSize = e.o.overlapSize()
// js.filling is allocated lazily on first Write/ReadFrom so callers
// that only use EncodeAll don't pay the (up to ~32 MB) jobSize cost.
js.filling = js.filling[:0]
if js.nextPrefix != nil {
js.putOverlapBuf(js.nextPrefix)
js.nextPrefix = nil
}
js.jobSeq = 0
js.flushedSeq = 0
js.flusherErr = nil
js.started = false
}
s.wg.Wait()
s.wWg.Wait()
if cap(s.filling) == 0 {
s.filling = make([]byte, 0, e.o.blockSize)
}
if e.o.concurrent > 1 {
if e.o.concurrent > 1 && !e.o.concurrentBlocks {
if cap(s.current) == 0 {
s.current = make([]byte, 0, e.o.blockSize)
}
@@ -145,6 +171,9 @@ func (e *Encoder) ResetWithOptions(w io.Writer, opts ...EOption) error {
}
}
hasDict := e.o.dict != nil
if e.o.concurrentBlocks && hasDict {
e.o.concurrentBlocks = false
}
if hadDict != hasDict {
// Dict presence changed — encoder type must be recreated.
e.state.encoder = nil
@@ -176,6 +205,49 @@ func (e *Encoder) Write(p []byte) (n int, err error) {
if s.eofWritten {
return 0, ErrEncoderClosed
}
if e.o.concurrentBlocks {
return e.writeJobs(p)
}
return e.writeBlocks(p)
}
func (e *Encoder) writeJobs(p []byte) (n int, err error) {
s := &e.state
js := &s.jobs
jobSize := js.jobSize
if cap(js.filling) == 0 && len(p) > 0 {
js.filling = make([]byte, 0, jobSize)
}
for len(p) > 0 {
if len(p)+len(js.filling) < jobSize {
if e.o.crc {
_, _ = s.encoder.CRC().Write(p)
}
js.filling = append(js.filling, p...)
return n + len(p), nil
}
add := p
if len(p)+len(js.filling) > jobSize {
add = add[:jobSize-len(js.filling)]
}
if e.o.crc {
_, _ = s.encoder.CRC().Write(add)
}
js.filling = append(js.filling, add...)
p = p[len(add):]
n += len(add)
if len(js.filling) < jobSize {
return n, nil
}
if err := e.dispatchJob(false); err != nil {
return n, err
}
}
return n, nil
}
func (e *Encoder) writeBlocks(p []byte) (n int, err error) {
s := &e.state
for len(p) > 0 {
if len(p)+len(s.filling) < e.o.blockSize {
if e.o.crc {
@@ -374,6 +446,10 @@ func (e *Encoder) ReadFrom(r io.Reader) (n int64, err error) {
println("Using ReadFrom")
}
if e.o.concurrentBlocks {
return e.readFromJobs(r)
}
// Flush any current writes.
if len(e.state.filling) > 0 {
if err := e.nextBlock(false); err != nil {
@@ -387,7 +463,6 @@ func (e *Encoder) ReadFrom(r io.Reader) (n int64, err error) {
if e.o.crc {
_, _ = e.state.encoder.CRC().Write(src[:n2])
}
// src is now the unfilled part...
src = src[n2:]
n += int64(n2)
switch err {
@@ -420,15 +495,63 @@ func (e *Encoder) ReadFrom(r io.Reader) (n int64, err error) {
}
}
func (e *Encoder) readFromJobs(r io.Reader) (n int64, err error) {
js := &e.state.jobs
jobSize := js.jobSize
// Flush any current filling.
if len(js.filling) > 0 {
if err := e.dispatchJob(false); err != nil {
return 0, err
}
}
if cap(js.filling) < jobSize {
js.filling = make([]byte, 0, jobSize)
}
js.filling = js.filling[:jobSize]
src := js.filling
for {
n2, err := r.Read(src)
if e.o.crc {
_, _ = e.state.encoder.CRC().Write(src[:n2])
}
src = src[n2:]
n += int64(n2)
switch err {
case io.EOF:
js.filling = js.filling[:len(js.filling)-len(src)]
return n, nil
case nil:
default:
e.state.err = err
return n, err
}
if len(src) > 0 {
continue
}
if err = e.dispatchJob(false); err != nil {
return n, err
}
if cap(js.filling) < jobSize {
js.filling = make([]byte, 0, jobSize)
}
js.filling = js.filling[:jobSize]
src = js.filling
}
}
// Flush will send the currently written data to output
// and block until everything has been written.
// This should only be used on rare occasions where pushing the currently queued data is critical.
func (e *Encoder) Flush() error {
s := &e.state
if e.o.concurrentBlocks {
return e.flushJobs()
}
if len(s.filling) > 0 {
err := e.nextBlock(false)
if err != nil {
// Ignore Flush after Close.
if errors.Is(s.err, ErrEncoderClosed) {
return nil
}
@@ -438,7 +561,6 @@ func (e *Encoder) Flush() error {
s.wg.Wait()
s.wWg.Wait()
if s.err != nil {
// Ignore Flush after Close.
if errors.Is(s.err, ErrEncoderClosed) {
return nil
}
@@ -447,6 +569,20 @@ func (e *Encoder) Flush() error {
return s.writeErr
}
func (e *Encoder) flushJobs() error {
js := &e.state.jobs
if len(js.filling) > 0 {
if err := e.dispatchJob(false); err != nil {
return err
}
}
e.waitAllJobs()
js.mu.Lock()
fErr := js.flusherErr
js.mu.Unlock()
return fErr
}
// Close will flush the final output and close the stream.
// The function will block until everything has been written.
// The Encoder can still be re-used after calling this.
@@ -455,12 +591,16 @@ func (e *Encoder) Close() error {
if s.encoder == nil {
return nil
}
if e.o.concurrentBlocks {
return e.closeJobs()
}
if s.w == nil {
if len(s.filling) == 0 && !s.headerWritten && !s.eofWritten && s.nInput == 0 {
return nil
}
return errors.New("zstd: encoder has no writer")
}
err := e.nextBlock(true)
if err != nil {
if errors.Is(s.err, ErrEncoderClosed) {
@@ -511,6 +651,68 @@ func (e *Encoder) Close() error {
return s.err
}
func (e *Encoder) closeJobs() error {
s := &e.state
js := &s.jobs
if errors.Is(s.err, ErrEncoderClosed) {
return nil
}
if s.w == nil {
if len(js.filling) == 0 && !s.headerWritten && !s.eofWritten && s.nInput == 0 {
return nil
}
return errors.New("zstd: encoder has no writer")
}
if err := e.dispatchJob(true); err != nil {
e.shutdownJobWorkers()
if errors.Is(s.err, ErrEncoderClosed) {
return nil
}
return err
}
if s.frameContentSize > 0 && s.nInput != s.frameContentSize {
e.shutdownJobWorkers()
return fmt.Errorf("frame content size %d given, but %d bytes was written", s.frameContentSize, s.nInput)
}
if s.fullFrameWritten {
e.shutdownJobWorkers()
s.err = ErrEncoderClosed
return nil
}
e.shutdownJobWorkers()
if js.flusherErr != nil {
return js.flusherErr
}
// Write CRC
if e.o.crc {
var tmp [4]byte
_, s.err = s.w.Write(s.encoder.AppendCRC(tmp[:0]))
s.nWritten += 4
}
// Add padding
if s.err == nil && e.o.pad > 0 {
add := calcSkippableFrame(s.nWritten, int64(e.o.pad))
frame, err := skippableFrame(js.filling[:0], add, rand.Reader)
if err != nil {
return err
}
_, s.err = s.w.Write(frame)
}
if s.err == nil {
s.err = ErrEncoderClosed
return nil
}
return s.err
}
// EncodeAll will encode all input in src and append it to dst.
// This function can be called concurrently, but each call will only run on a single goroutine.
// If empty input is given, nothing is returned, unless WithZeroFrames is specified.

View File

@@ -14,22 +14,23 @@ type EOption func(*encoderOptions) error
// options retains accumulated state of multiple options.
type encoderOptions struct {
resetOpt bool
concurrent int
level EncoderLevel
single *bool
pad int
blockSize int
windowSize int
crc bool
fullZero bool
noEntropy bool
allLitEntropy bool
customWindow bool
customALEntropy bool
customBlockSize bool
lowMem bool
dict *dict
resetOpt bool
concurrent int
level EncoderLevel
single *bool
pad int
blockSize int
windowSize int
crc bool
fullZero bool
noEntropy bool
allLitEntropy bool
customWindow bool
customALEntropy bool
customBlockSize bool
lowMem bool
dict *dict
concurrentBlocks bool
}
func (o *encoderOptions) setDefault() {
@@ -333,6 +334,42 @@ func WithLowerEncoderMem(b bool) EOption {
}
}
// WithConcurrentBlocks enables job-based parallel compression for streams.
// When enabled and concurrent > 1, input is split into large sections (jobs)
// that are compressed simultaneously by multiple goroutines.
// Each non-first job receives an overlap prefix from the previous job for match context.
// Output is flushed in order, producing a valid single-frame zstd stream.
//
// Currently disabled when used with dictionary encoding.
// Cannot be changed with ResetWithOptions.
func WithConcurrentBlocks(b bool) EOption {
return func(o *encoderOptions) error {
if o.resetOpt && b != o.concurrentBlocks {
return errors.New("WithConcurrentBlocks cannot be changed on Reset")
}
o.concurrentBlocks = b
return nil
}
}
// jobSize returns the input section size per parallel job.
func (o *encoderOptions) jobSize() int {
s := max(o.windowSize*4, 512<<10)
return s
}
// overlapSize returns the overlap prefix size for parallel jobs.
func (o *encoderOptions) overlapSize() int {
switch o.level {
case SpeedBestCompression:
return o.windowSize / 2
case SpeedBetterCompression:
return o.windowSize / 4
default:
return o.windowSize / 8
}
}
// WithEncoderDict allows to register a dictionary that will be used for the encode.
//
// The slice dict must be in the [dictionary format] produced by

View File

@@ -1,4 +1,4 @@
// Code generated by command: go run gen_fse.go -out ../fse_decoder_amd64.s -pkg=zstd. DO NOT EDIT.
// Code generated by command: go run gen_fse.go -out ../fse_decoder.s -arch amd64,arm64 -pkg=zstd. DO NOT EDIT.
//go:build !appengine && !noasm && gc && !noasm

View File

@@ -0,0 +1,146 @@
// Code generated by command: go run gen_fse.go -out ../fse_decoder.s -arch amd64,arm64 -pkg=zstd. DO NOT EDIT.
// EXPERIMENTAL arm64 output lowered from an amd64 avo program.
//go:build arm64 && !appengine && !noasm && gc && !noasm
// func buildDtable_asm(s *fseDecoder, ctx *buildDtableAsmContext) int
TEXT ·buildDtable_asm(SB), $0-24
MOVD ctx+8(FP), R1
MOVD s+0(FP), R6
// Load values
MOVBU 4098(R6), R2
MOVD $0, R0
MOVD $1, R16
LSL R2, R16, R16
ORR R16, R0, R0
MOVD (R1), R3
MOVD 16(R1), R5
SUB $1, R0, R7
MOVD 8(R1), R1
MOVHU 4096(R6), R6
// End load values
// Init, lay down lowprob symbols
MOVD $0, R8
JMP init_main_loop_condition
init_main_loop:
ADD R8<<1, R1, R15
MOVH (R15), R9
CMN $1, R9
BNE do_not_update_high_threshold
ADD R7<<3, R5, R15
MOVB R8, 1(R15)
SUB $1, R7, R7
MOVD $0x0000000000000001, R9
do_not_update_high_threshold:
ADD R8<<1, R3, R15
MOVH R9, (R15)
ADD $1, R8, R8
init_main_loop_condition:
CMP R6, R8
BLT init_main_loop
// Spread symbols
// Calculate table step
MOVD R0, R8
LSR $0x01, R8, R8
MOVD R0, R9
LSR $0x03, R9, R9
ADD R9, R8, R8
ADD $3, R8, R8
// Fill add bits values
SUB $1, R0, R9
MOVD $0, R10
MOVD $0, R11
JMP spread_main_loop_condition
spread_main_loop:
MOVD $0, R12
ADD R11<<1, R1, R15
MOVH (R15), R13
JMP spread_inner_loop_condition
spread_inner_loop:
ADD R10<<3, R5, R15
MOVB R11, 1(R15)
adjust_position:
ADD R8, R10, R10
AND R9, R10, R10
CMP R7, R10
BGT adjust_position
ADD $1, R12, R12
spread_inner_loop_condition:
CMP R13, R12
BLT spread_inner_loop
ADD $1, R11, R11
spread_main_loop_condition:
CMP R6, R11
BLT spread_main_loop
TST R10, R10
BEQ spread_check_ok
MOVD ctx+8(FP), R0
MOVD R10, 24(R0)
MOVD $+1, R16
MOVD R16, ret+16(FP)
RET
spread_check_ok:
// Build Decoding table
MOVD $0, R6
build_table_main_table:
ADD R6<<3, R5, R15
MOVBU 1(R15), R1
ADD R1<<1, R3, R15
MOVHU (R15), R7
ADD $1, R7, R8
ADD R1<<1, R3, R15
MOVH R8, (R15)
MOVD R7, R8
CLZ R8, R16
MOVD $63, R8
SUB R16, R8, R8
MOVD R2, R1
SUB R8, R1, R1
LSL R1, R7, R7
SUB R0, R7, R7
ADD R6<<3, R5, R15
MOVB R1, (R15)
ADD R6<<3, R5, R15
MOVH R7, 2(R15)
CMP R0, R7
BLE build_table_check1_ok
MOVD ctx+8(FP), R1
MOVD R7, 24(R1)
MOVD R0, 32(R1)
MOVD $+2, R16
MOVD R16, ret+16(FP)
RET
build_table_check1_ok:
TST R1, R1
BNE build_table_check2_ok
CMP R6, R7
BNE build_table_check2_ok
MOVD ctx+8(FP), R0
MOVD R7, 24(R0)
MOVD R6, 32(R0)
MOVD $+3, R16
MOVD R16, ret+16(FP)
RET
build_table_check2_ok:
ADD $1, R6, R6
CMP R0, R6
BLT build_table_main_table
MOVD $+0, R16
MOVD R16, ret+16(FP)
RET

View File

@@ -1,4 +1,4 @@
//go:build amd64 && !appengine && !noasm && gc
//go:build (amd64 || arm64) && !appengine && !noasm && gc
package zstd
@@ -6,6 +6,10 @@ import (
"fmt"
)
// buildDtable_asm is generated by _generate/gen_fse.go and lowered to each
// architecture (amd64 by goasm, arm64 by the avo arm64 lowering printer). The
// Go side is identical across architectures, so it lives here.
type buildDtableAsmContext struct {
// inputs
stateTable *uint16
@@ -18,7 +22,7 @@ type buildDtableAsmContext struct {
errParam2 uint64
}
// buildDtable_asm is an x86 assembly implementation of fseDecoder.buildDtable.
// buildDtable_asm is an assembly implementation of fseDecoder.buildDtable.
// Function returns non-zero exit code on error.
//
//go:noescape

View File

@@ -1,4 +1,4 @@
//go:build !amd64 || appengine || !gc || noasm
//go:build (!amd64 && !arm64) || appengine || !gc || noasm
package zstd

View File

@@ -3,30 +3,47 @@
package zstd
import (
"fmt"
"io"
"github.com/klauspost/compress/internal/cpuinfo"
)
type decodeSyncAsmContext struct {
llTable []decSymbol
mlTable []decSymbol
ofTable []decSymbol
llState uint64
mlState uint64
ofState uint64
iteration int
litRemain int
out []byte
outPosition int
literals []byte
litPosition int
history []byte
windowSize int
ll int // set on error (not for all errors, please refer to _generate/gen.go)
ml int // set on error (not for all errors, please refer to _generate/gen.go)
mo int // set on error (not for all errors, please refer to _generate/gen.go)
// The shared decode/decodeSync/executeSimple wrappers and context structs live
// in seqdec_asm.go; this file only declares the amd64 asm routines and the
// dispatch helpers that pick the BMI2 / non-BMI2 (and 56-bit / safe) variant.
// sequenceDecs_decode implements the main loop of sequenceDecs in x86 asm.
//
// Please refer to seqdec_generic.go for the reference implementation.
//
//go:noescape
func sequenceDecs_decode_amd64(s *sequenceDecs, br *bitReader, ctx *decodeAsmContext) int
// sequenceDecs_decode_56_amd64 implements the main loop of sequenceDecs in x86 asm.
//
//go:noescape
func sequenceDecs_decode_56_amd64(s *sequenceDecs, br *bitReader, ctx *decodeAsmContext) int
// sequenceDecs_decode_bmi2 implements the main loop of sequenceDecs in x86 asm with BMI2 extensions.
//
//go:noescape
func sequenceDecs_decode_bmi2(s *sequenceDecs, br *bitReader, ctx *decodeAsmContext) int
// sequenceDecs_decode_56_bmi2 implements the main loop of sequenceDecs in x86 asm with BMI2 extensions.
//
//go:noescape
func sequenceDecs_decode_56_bmi2(s *sequenceDecs, br *bitReader, ctx *decodeAsmContext) int
// decodeAsm runs the sequenceDecs decode loop, choosing the BMI2 / 56-bit variant.
func decodeAsm(s *sequenceDecs, br *bitReader, ctx *decodeAsmContext, lte56bits bool) int {
if cpuinfo.HasBMI2() {
if lte56bits {
return sequenceDecs_decode_56_bmi2(s, br, ctx)
}
return sequenceDecs_decode_bmi2(s, br, ctx)
}
if lte56bits {
return sequenceDecs_decode_56_amd64(s, br, ctx)
}
return sequenceDecs_decode_amd64(s, br, ctx)
}
// sequenceDecs_decodeSync_amd64 implements the main loop of sequenceDecs.decodeSync in x86 asm.
@@ -51,273 +68,18 @@ func sequenceDecs_decodeSync_safe_amd64(s *sequenceDecs, br *bitReader, ctx *dec
//go:noescape
func sequenceDecs_decodeSync_safe_bmi2(s *sequenceDecs, br *bitReader, ctx *decodeSyncAsmContext) int
// decode sequences from the stream with the provided history but without a dictionary.
func (s *sequenceDecs) decodeSyncSimple(hist []byte) (bool, error) {
if len(s.dict) > 0 {
return false, nil
}
if s.maxSyncLen == 0 && cap(s.out)-len(s.out) < maxCompressedBlockSize {
return false, nil
}
// FIXME: Using unsafe memory copies leads to rare, random crashes
// with fuzz testing. It is therefore disabled for now.
const useSafe = true
/*
useSafe := false
if s.maxSyncLen == 0 && cap(s.out)-len(s.out) < maxCompressedBlockSizeAlloc {
useSafe = true
}
if s.maxSyncLen > 0 && cap(s.out)-len(s.out)-compressedBlockOverAlloc < int(s.maxSyncLen) {
useSafe = true
}
if cap(s.literals) < len(s.literals)+compressedBlockOverAlloc {
useSafe = true
}
*/
br := s.br
maxBlockSize := min(s.windowSize, maxCompressedBlockSize)
ctx := decodeSyncAsmContext{
llTable: s.litLengths.fse.dt[:maxTablesize],
mlTable: s.matchLengths.fse.dt[:maxTablesize],
ofTable: s.offsets.fse.dt[:maxTablesize],
llState: uint64(s.litLengths.state.state),
mlState: uint64(s.matchLengths.state.state),
ofState: uint64(s.offsets.state.state),
iteration: s.nSeqs - 1,
litRemain: len(s.literals),
out: s.out,
outPosition: len(s.out),
literals: s.literals,
windowSize: s.windowSize,
history: hist,
}
s.seqSize = 0
startSize := len(s.out)
var errCode int
// decodeSyncAsm runs the decodeSync loop, choosing the BMI2 / safe variant.
func decodeSyncAsm(s *sequenceDecs, br *bitReader, ctx *decodeSyncAsmContext, safe bool) int {
if cpuinfo.HasBMI2() {
if useSafe {
errCode = sequenceDecs_decodeSync_safe_bmi2(s, br, &ctx)
} else {
errCode = sequenceDecs_decodeSync_bmi2(s, br, &ctx)
}
} else {
if useSafe {
errCode = sequenceDecs_decodeSync_safe_amd64(s, br, &ctx)
} else {
errCode = sequenceDecs_decodeSync_amd64(s, br, &ctx)
if safe {
return sequenceDecs_decodeSync_safe_bmi2(s, br, ctx)
}
return sequenceDecs_decodeSync_bmi2(s, br, ctx)
}
switch errCode {
case noError:
break
case errorMatchLenOfsMismatch:
return true, fmt.Errorf("zero matchoff and matchlen (%d) > 0", ctx.ml)
case errorMatchLenTooBig:
return true, fmt.Errorf("match len (%d) bigger than max allowed length", ctx.ml)
case errorMatchOffTooBig:
return true, fmt.Errorf("match offset (%d) bigger than current history (%d)",
ctx.mo, ctx.outPosition+len(hist)-startSize)
case errorNotEnoughLiterals:
return true, fmt.Errorf("unexpected literal count, want %d bytes, but only %d is available",
ctx.ll, ctx.litRemain+ctx.ll)
case errorOverread:
return true, io.ErrUnexpectedEOF
case errorNotEnoughSpace:
size := ctx.outPosition + ctx.ll + ctx.ml
if debugDecoder {
println("msl:", s.maxSyncLen, "cap", cap(s.out), "bef:", startSize, "sz:", size-startSize, "mbs:", maxBlockSize, "outsz:", cap(s.out)-startSize)
}
return true, fmt.Errorf("output bigger than max block size (%d)", maxBlockSize)
default:
return true, fmt.Errorf("sequenceDecs_decode returned erroneous code %d", errCode)
if safe {
return sequenceDecs_decodeSync_safe_amd64(s, br, ctx)
}
s.seqSize += ctx.litRemain
if s.seqSize > maxBlockSize {
return true, fmt.Errorf("output bigger than max block size (%d)", maxBlockSize)
}
err := br.close()
if err != nil {
printf("Closing sequences: %v, %+v\n", err, *br)
return true, err
}
s.literals = s.literals[ctx.litPosition:]
t := ctx.outPosition
s.out = s.out[:t]
// Add final literals
s.out = append(s.out, s.literals...)
if debugDecoder {
t += len(s.literals)
if t != len(s.out) {
panic(fmt.Errorf("length mismatch, want %d, got %d", len(s.out), t))
}
}
return true, nil
}
// --------------------------------------------------------------------------------
type decodeAsmContext struct {
llTable []decSymbol
mlTable []decSymbol
ofTable []decSymbol
llState uint64
mlState uint64
ofState uint64
iteration int
seqs []seqVals
litRemain int
}
const noError = 0
// error reported when mo == 0 && ml > 0
const errorMatchLenOfsMismatch = 1
// error reported when ml > maxMatchLen
const errorMatchLenTooBig = 2
// error reported when mo > available history or mo > s.windowSize
const errorMatchOffTooBig = 3
// error reported when the sum of literal lengths exeeceds the literal buffer size
const errorNotEnoughLiterals = 4
// error reported when capacity of `out` is too small
const errorNotEnoughSpace = 5
// error reported when bits are overread.
const errorOverread = 6
// sequenceDecs_decode implements the main loop of sequenceDecs in x86 asm.
//
// Please refer to seqdec_generic.go for the reference implementation.
//
//go:noescape
func sequenceDecs_decode_amd64(s *sequenceDecs, br *bitReader, ctx *decodeAsmContext) int
// sequenceDecs_decode implements the main loop of sequenceDecs in x86 asm.
//
// Please refer to seqdec_generic.go for the reference implementation.
//
//go:noescape
func sequenceDecs_decode_56_amd64(s *sequenceDecs, br *bitReader, ctx *decodeAsmContext) int
// sequenceDecs_decode implements the main loop of sequenceDecs in x86 asm with BMI2 extensions.
//
//go:noescape
func sequenceDecs_decode_bmi2(s *sequenceDecs, br *bitReader, ctx *decodeAsmContext) int
// sequenceDecs_decode implements the main loop of sequenceDecs in x86 asm with BMI2 extensions.
//
//go:noescape
func sequenceDecs_decode_56_bmi2(s *sequenceDecs, br *bitReader, ctx *decodeAsmContext) int
// decode sequences from the stream without the provided history.
func (s *sequenceDecs) decode(seqs []seqVals) error {
br := s.br
maxBlockSize := min(s.windowSize, maxCompressedBlockSize)
ctx := decodeAsmContext{
llTable: s.litLengths.fse.dt[:maxTablesize],
mlTable: s.matchLengths.fse.dt[:maxTablesize],
ofTable: s.offsets.fse.dt[:maxTablesize],
llState: uint64(s.litLengths.state.state),
mlState: uint64(s.matchLengths.state.state),
ofState: uint64(s.offsets.state.state),
seqs: seqs,
iteration: len(seqs) - 1,
litRemain: len(s.literals),
}
if debugDecoder {
println("decode: decoding", len(seqs), "sequences", br.remain(), "bits remain on stream")
}
s.seqSize = 0
lte56bits := s.maxBits+s.offsets.fse.actualTableLog+s.matchLengths.fse.actualTableLog+s.litLengths.fse.actualTableLog <= 56
var errCode int
if cpuinfo.HasBMI2() {
if lte56bits {
errCode = sequenceDecs_decode_56_bmi2(s, br, &ctx)
} else {
errCode = sequenceDecs_decode_bmi2(s, br, &ctx)
}
} else {
if lte56bits {
errCode = sequenceDecs_decode_56_amd64(s, br, &ctx)
} else {
errCode = sequenceDecs_decode_amd64(s, br, &ctx)
}
}
if errCode != 0 {
i := len(seqs) - ctx.iteration - 1
switch errCode {
case errorMatchLenOfsMismatch:
ml := ctx.seqs[i].ml
return fmt.Errorf("zero matchoff and matchlen (%d) > 0", ml)
case errorMatchLenTooBig:
ml := ctx.seqs[i].ml
return fmt.Errorf("match len (%d) bigger than max allowed length", ml)
case errorNotEnoughLiterals:
ll := ctx.seqs[i].ll
return fmt.Errorf("unexpected literal count, want %d bytes, but only %d is available", ll, ctx.litRemain+ll)
case errorOverread:
return io.ErrUnexpectedEOF
}
return fmt.Errorf("sequenceDecs_decode_amd64 returned erroneous code %d", errCode)
}
if ctx.litRemain < 0 {
return fmt.Errorf("literal count is too big: total available %d, total requested %d",
len(s.literals), len(s.literals)-ctx.litRemain)
}
s.seqSize += ctx.litRemain
if s.seqSize > maxBlockSize {
return fmt.Errorf("output bigger than max block size (%d)", maxBlockSize)
}
if debugDecoder {
println("decode: ", br.remain(), "bits remain on stream. code:", errCode)
}
err := br.close()
if err != nil {
printf("Closing sequences: %v, %+v\n", err, *br)
}
return err
}
// --------------------------------------------------------------------------------
type executeAsmContext struct {
seqs []seqVals
seqIndex int
out []byte
history []byte
literals []byte
outPosition int
litPosition int
windowSize int
return sequenceDecs_decodeSync_amd64(s, br, ctx)
}
// sequenceDecs_executeSimple_amd64 implements the main loop of sequenceDecs.executeSimple in x86 asm.
@@ -334,54 +96,10 @@ func sequenceDecs_executeSimple_amd64(ctx *executeAsmContext) bool
//go:noescape
func sequenceDecs_executeSimple_safe_amd64(ctx *executeAsmContext) bool
// executeSimple handles cases when dictionary is not used.
func (s *sequenceDecs) executeSimple(seqs []seqVals, hist []byte) error {
// Ensure we have enough output size...
if len(s.out)+s.seqSize+compressedBlockOverAlloc > cap(s.out) {
addBytes := s.seqSize + len(s.out) + compressedBlockOverAlloc
s.out = append(s.out, make([]byte, addBytes)...)
s.out = s.out[:len(s.out)-addBytes]
// executeSimpleAsm runs the executeSimple loop, choosing the safe variant.
func executeSimpleAsm(ctx *executeAsmContext, safe bool) bool {
if safe {
return sequenceDecs_executeSimple_safe_amd64(ctx)
}
if debugDecoder {
printf("Execute %d seqs with literals: %d into %d bytes\n", len(seqs), len(s.literals), s.seqSize)
}
var t = len(s.out)
out := s.out[:t+s.seqSize]
ctx := executeAsmContext{
seqs: seqs,
seqIndex: 0,
out: out,
history: hist,
outPosition: t,
litPosition: 0,
literals: s.literals,
windowSize: s.windowSize,
}
var ok bool
if cap(s.literals) < len(s.literals)+compressedBlockOverAlloc {
ok = sequenceDecs_executeSimple_safe_amd64(&ctx)
} else {
ok = sequenceDecs_executeSimple_amd64(&ctx)
}
if !ok {
return fmt.Errorf("match offset (%d) bigger than current history (%d)",
seqs[ctx.seqIndex].mo, ctx.outPosition+len(hist))
}
s.literals = s.literals[ctx.litPosition:]
t = ctx.outPosition
// Add final literals
copy(out[t:], s.literals)
if debugDecoder {
t += len(s.literals)
if t != len(out) {
panic(fmt.Errorf("length mismatch, want %d, got %d, ss: %d", len(out), t, s.seqSize))
}
}
s.out = out
return nil
return sequenceDecs_executeSimple_amd64(ctx)
}

View File

@@ -1,4 +1,4 @@
// Code generated by command: go run gen.go -out ../seqdec_amd64.s -pkg=zstd. DO NOT EDIT.
// Code generated by command: go run gen.go -out ../seqdec.s -arch amd64,arm64 -pkg=zstd. DO NOT EDIT.
//go:build !appengine && !noasm && gc && !noasm

View File

@@ -0,0 +1,70 @@
//go:build arm64 && !appengine && !noasm && gc
package zstd
// The shared decode/decodeSync/executeSimple wrappers and context structs live
// in seqdec_asm.go; this file only declares the arm64 asm routines (generated
// by the avo arm64 lowering printer) and the dispatch helpers. arm64 has no
// BMI2, so each helper selects only between the 56-bit / safe variants.
// sequenceDecs_decode_arm64 implements the main loop of sequenceDecs in arm64 asm.
//
// Please refer to seqdec_generic.go for the reference implementation.
//
//go:noescape
func sequenceDecs_decode_arm64(s *sequenceDecs, br *bitReader, ctx *decodeAsmContext) int
// sequenceDecs_decode_56_arm64 implements the main loop of sequenceDecs in arm64 asm.
//
//go:noescape
func sequenceDecs_decode_56_arm64(s *sequenceDecs, br *bitReader, ctx *decodeAsmContext) int
// decodeAsm runs the sequenceDecs decode loop, choosing the 56-bit variant.
func decodeAsm(s *sequenceDecs, br *bitReader, ctx *decodeAsmContext, lte56bits bool) int {
if lte56bits {
return sequenceDecs_decode_56_arm64(s, br, ctx)
}
return sequenceDecs_decode_arm64(s, br, ctx)
}
// sequenceDecs_decodeSync_arm64 implements the main loop of sequenceDecs.decodeSync in arm64 asm.
//
// Please refer to seqdec_generic.go for the reference implementation.
//
//go:noescape
func sequenceDecs_decodeSync_arm64(s *sequenceDecs, br *bitReader, ctx *decodeSyncAsmContext) int
// sequenceDecs_decodeSync_safe_arm64 does the same as above, but does not write more than output buffer.
//
//go:noescape
func sequenceDecs_decodeSync_safe_arm64(s *sequenceDecs, br *bitReader, ctx *decodeSyncAsmContext) int
// decodeSyncAsm runs the decodeSync loop, choosing the safe variant.
func decodeSyncAsm(s *sequenceDecs, br *bitReader, ctx *decodeSyncAsmContext, safe bool) int {
if safe {
return sequenceDecs_decodeSync_safe_arm64(s, br, ctx)
}
return sequenceDecs_decodeSync_arm64(s, br, ctx)
}
// sequenceDecs_executeSimple_arm64 implements the main loop of sequenceDecs.executeSimple in arm64 asm.
//
// Returns false if a match offset is too big.
//
// Please refer to seqdec_generic.go for the reference implementation.
//
//go:noescape
func sequenceDecs_executeSimple_arm64(ctx *executeAsmContext) bool
// Same as above, but with safe memcopies
//
//go:noescape
func sequenceDecs_executeSimple_safe_arm64(ctx *executeAsmContext) bool
// executeSimpleAsm runs the executeSimple loop, choosing the safe variant.
func executeSimpleAsm(ctx *executeAsmContext, safe bool) bool {
if safe {
return sequenceDecs_executeSimple_safe_arm64(ctx)
}
return sequenceDecs_executeSimple_arm64(ctx)
}

2705
vendor/github.com/klauspost/compress/zstd/seqdec_arm64.s generated vendored Normal file
View File

File diff suppressed because it is too large Load Diff

289
vendor/github.com/klauspost/compress/zstd/seqdec_asm.go generated vendored Normal file
View File

@@ -0,0 +1,289 @@
//go:build (amd64 || arm64) && !appengine && !noasm && gc
package zstd
import (
"fmt"
"io"
)
// This file holds the parts of the assembly sequence decoder that are identical
// across architectures: the context structs exchanged with the asm, the error
// codes, and the decode/decodeSync/executeSimple wrappers. Each architecture
// supplies the small dispatch helpers (decodeAsm, decodeSyncAsm,
// executeSimpleAsm) that select the concrete asm routine — amd64 also chooses a
// BMI2 variant, arm64 has a single implementation.
type decodeSyncAsmContext struct {
llTable []decSymbol
mlTable []decSymbol
ofTable []decSymbol
llState uint64
mlState uint64
ofState uint64
iteration int
litRemain int
out []byte
outPosition int
literals []byte
litPosition int
history []byte
windowSize int
ll int // set on error (not for all errors, please refer to _generate/gen.go)
ml int // set on error (not for all errors, please refer to _generate/gen.go)
mo int // set on error (not for all errors, please refer to _generate/gen.go)
}
type decodeAsmContext struct {
llTable []decSymbol
mlTable []decSymbol
ofTable []decSymbol
llState uint64
mlState uint64
ofState uint64
iteration int
seqs []seqVals
litRemain int
}
type executeAsmContext struct {
seqs []seqVals
seqIndex int
out []byte
history []byte
literals []byte
outPosition int
litPosition int
windowSize int
}
const noError = 0
// error reported when mo == 0 && ml > 0
const errorMatchLenOfsMismatch = 1
// error reported when ml > maxMatchLen
const errorMatchLenTooBig = 2
// error reported when mo > available history or mo > s.windowSize
const errorMatchOffTooBig = 3
// error reported when the sum of literal lengths exeeceds the literal buffer size
const errorNotEnoughLiterals = 4
// error reported when capacity of `out` is too small
const errorNotEnoughSpace = 5
// error reported when bits are overread.
const errorOverread = 6
// decode sequences from the stream with the provided history but without a dictionary.
func (s *sequenceDecs) decodeSyncSimple(hist []byte) (bool, error) {
if len(s.dict) > 0 {
return false, nil
}
if s.maxSyncLen == 0 && cap(s.out)-len(s.out) < maxCompressedBlockSize {
return false, nil
}
// FIXME: Using unsafe memory copies leads to rare, random crashes
// with fuzz testing. It is therefore disabled for now.
const useSafe = true
br := s.br
maxBlockSize := min(s.windowSize, maxCompressedBlockSize)
ctx := decodeSyncAsmContext{
llTable: s.litLengths.fse.dt[:maxTablesize],
mlTable: s.matchLengths.fse.dt[:maxTablesize],
ofTable: s.offsets.fse.dt[:maxTablesize],
llState: uint64(s.litLengths.state.state),
mlState: uint64(s.matchLengths.state.state),
ofState: uint64(s.offsets.state.state),
iteration: s.nSeqs - 1,
litRemain: len(s.literals),
out: s.out,
outPosition: len(s.out),
literals: s.literals,
windowSize: s.windowSize,
history: hist,
}
s.seqSize = 0
startSize := len(s.out)
errCode := decodeSyncAsm(s, br, &ctx, useSafe)
switch errCode {
case noError:
break
case errorMatchLenOfsMismatch:
return true, fmt.Errorf("zero matchoff and matchlen (%d) > 0", ctx.ml)
case errorMatchLenTooBig:
return true, fmt.Errorf("match len (%d) bigger than max allowed length", ctx.ml)
case errorMatchOffTooBig:
return true, fmt.Errorf("match offset (%d) bigger than current history (%d)",
ctx.mo, ctx.outPosition+len(hist)-startSize)
case errorNotEnoughLiterals:
return true, fmt.Errorf("unexpected literal count, want %d bytes, but only %d is available",
ctx.ll, ctx.litRemain+ctx.ll)
case errorOverread:
return true, io.ErrUnexpectedEOF
case errorNotEnoughSpace:
size := ctx.outPosition + ctx.ll + ctx.ml
if debugDecoder {
println("msl:", s.maxSyncLen, "cap", cap(s.out), "bef:", startSize, "sz:", size-startSize, "mbs:", maxBlockSize, "outsz:", cap(s.out)-startSize)
}
return true, fmt.Errorf("output bigger than max block size (%d)", maxBlockSize)
default:
return true, fmt.Errorf("sequenceDecs_decode returned erroneous code %d", errCode)
}
s.seqSize += ctx.litRemain
if s.seqSize > maxBlockSize {
return true, fmt.Errorf("output bigger than max block size (%d)", maxBlockSize)
}
err := br.close()
if err != nil {
printf("Closing sequences: %v, %+v\n", err, *br)
return true, err
}
s.literals = s.literals[ctx.litPosition:]
t := ctx.outPosition
s.out = s.out[:t]
// Add final literals
s.out = append(s.out, s.literals...)
if debugDecoder {
t += len(s.literals)
if t != len(s.out) {
panic(fmt.Errorf("length mismatch, want %d, got %d", len(s.out), t))
}
}
return true, nil
}
// decode sequences from the stream without the provided history.
func (s *sequenceDecs) decode(seqs []seqVals) error {
br := s.br
maxBlockSize := min(s.windowSize, maxCompressedBlockSize)
ctx := decodeAsmContext{
llTable: s.litLengths.fse.dt[:maxTablesize],
mlTable: s.matchLengths.fse.dt[:maxTablesize],
ofTable: s.offsets.fse.dt[:maxTablesize],
llState: uint64(s.litLengths.state.state),
mlState: uint64(s.matchLengths.state.state),
ofState: uint64(s.offsets.state.state),
seqs: seqs,
iteration: len(seqs) - 1,
litRemain: len(s.literals),
}
if debugDecoder {
println("decode: decoding", len(seqs), "sequences", br.remain(), "bits remain on stream")
}
s.seqSize = 0
lte56bits := s.maxBits+s.offsets.fse.actualTableLog+s.matchLengths.fse.actualTableLog+s.litLengths.fse.actualTableLog <= 56
errCode := decodeAsm(s, br, &ctx, lte56bits)
if errCode != 0 {
i := len(seqs) - ctx.iteration - 1
switch errCode {
case errorMatchLenOfsMismatch:
ml := ctx.seqs[i].ml
return fmt.Errorf("zero matchoff and matchlen (%d) > 0", ml)
case errorMatchLenTooBig:
ml := ctx.seqs[i].ml
return fmt.Errorf("match len (%d) bigger than max allowed length", ml)
case errorNotEnoughLiterals:
ll := ctx.seqs[i].ll
return fmt.Errorf("unexpected literal count, want %d bytes, but only %d is available", ll, ctx.litRemain+ll)
case errorOverread:
return io.ErrUnexpectedEOF
}
return fmt.Errorf("sequenceDecs_decode_amd64 returned erroneous code %d", errCode)
}
if ctx.litRemain < 0 {
return fmt.Errorf("literal count is too big: total available %d, total requested %d",
len(s.literals), len(s.literals)-ctx.litRemain)
}
s.seqSize += ctx.litRemain
if s.seqSize > maxBlockSize {
return fmt.Errorf("output bigger than max block size (%d)", maxBlockSize)
}
if debugDecoder {
println("decode: ", br.remain(), "bits remain on stream. code:", errCode)
}
err := br.close()
if err != nil {
printf("Closing sequences: %v, %+v\n", err, *br)
}
return err
}
// executeSimple handles cases when dictionary is not used.
func (s *sequenceDecs) executeSimple(seqs []seqVals, hist []byte) error {
// Ensure we have enough output size...
if len(s.out)+s.seqSize+compressedBlockOverAlloc > cap(s.out) {
addBytes := s.seqSize + len(s.out) + compressedBlockOverAlloc
s.out = append(s.out, make([]byte, addBytes)...)
s.out = s.out[:len(s.out)-addBytes]
}
if debugDecoder {
printf("Execute %d seqs with literals: %d into %d bytes\n", len(seqs), len(s.literals), s.seqSize)
}
var t = len(s.out)
out := s.out[:t+s.seqSize]
ctx := executeAsmContext{
seqs: seqs,
seqIndex: 0,
out: out,
history: hist,
outPosition: t,
litPosition: 0,
literals: s.literals,
windowSize: s.windowSize,
}
// useSafe avoids overwriting the output buffer when the literals slice has
// not been allocated with the required over-allocation slack.
useSafe := cap(s.literals) < len(s.literals)+compressedBlockOverAlloc
ok := executeSimpleAsm(&ctx, useSafe)
if !ok {
return fmt.Errorf("match offset (%d) bigger than current history (%d)",
seqs[ctx.seqIndex].mo, ctx.outPosition+len(hist))
}
s.literals = s.literals[ctx.litPosition:]
t = ctx.outPosition
// Add final literals
copy(out[t:], s.literals)
if debugDecoder {
t += len(s.literals)
if t != len(out) {
panic(fmt.Errorf("length mismatch, want %d, got %d, ss: %d", len(out), t, s.seqSize))
}
}
s.out = out
return nil
}

View File

@@ -1,4 +1,4 @@
//go:build !amd64 || appengine || !gc || noasm
//go:build (!amd64 && !arm64) || appengine || !gc || noasm
package zstd

7
vendor/go.podman.io/common/libimage/errno_unix.go generated vendored Normal file
View File

@@ -0,0 +1,7 @@
//go:build !windows
package libimage
import "syscall"
const errNoSpace = syscall.ENOSPC

7
vendor/go.podman.io/common/libimage/errno_windows.go generated vendored Normal file
View File

@@ -0,0 +1,7 @@
//go:build windows
package libimage
import "golang.org/x/sys/windows"
const errNoSpace = windows.ERROR_DISK_FULL

View File

@@ -110,6 +110,12 @@ func (r *Runtime) Load(ctx context.Context, path string, options *LoadOptions) (
return loadedImages, err
}
logrus.Debugf("Error loading %s (%s): %v", path, transportName, err)
if errors.Is(err, errNoSpace) {
// %.0w makes err visible to error.Unwrap() without including any text
return nil, fmt.Errorf("no space left on device%.0w", err)
}
loadErrors = append(loadErrors, fmt.Errorf("%s: %v", transportName, err))
}

View File

@@ -24,6 +24,9 @@ const (
// Let's follow Firefox by limiting parallel downloads to 6. We do the
// same when pulling images in c/image.
searchMaxParallel = int64(6)
// searchResultOK represents true for Official and Automated.
searchResultOK = "[OK]"
)
// SearchResult is holding image-search related data.
@@ -44,6 +47,16 @@ type SearchResult struct {
Tag string
}
// IsAutomated returns true if the image was created by an automated build.
func (s *SearchResult) IsAutomated() bool {
return s.Automated == searchResultOK
}
// IsOfficial returns true if it's an official image.
func (s *SearchResult) IsOfficial() bool {
return s.Official == searchResultOK
}
// SearchOptions customize searching images.
type SearchOptions struct {
// Filter allows to filter the results.
@@ -214,11 +227,11 @@ func (r *Runtime) searchImageInRegistry(ctx context.Context, term, registry stri
}
official := ""
if results[i].IsOfficial {
official = "[OK]"
official = searchResultOK
}
automated := ""
if results[i].IsAutomated {
automated = "[OK]"
automated = searchResultOK
}
description := strings.ReplaceAll(results[i].Description, "\n", " ")
if len(description) > 44 && !options.NoTrunc {

View File

@@ -33,8 +33,7 @@ type HostContainersInternalOptions struct {
HostNetwork bool
}
// Lookup "host.containers.internal" dns name so we can add it to /etc/hosts when running inside podman machine.
var machineHostContainersInternalIP = sync.OnceValue(func() string {
func gvProxyHostIP() string {
var errMsg string
addrs, err := net.LookupIP(HostContainersInternal)
if err == nil {
@@ -47,15 +46,22 @@ var machineHostContainersInternalIP = sync.OnceValue(func() string {
}
logrus.Warnf("Failed to resolve %s for the host entry ip address: %s", HostContainersInternal, errMsg)
return ""
}
// Lookup "host.containers.internal" dns name so we can add it to /etc/hosts when running inside podman machine.
var machineHostContainersInternalIP = sync.OnceValue(func() string {
// If machine using gvproxy we let the gvproxy dns server handle resolve the name and then use that ip.
if machine.IsGvProxyBased() {
return gvProxyHostIP()
}
return wslHostIP()
})
// GetHostContainersInternalIP returns the host.containers.internal ip.
func GetHostContainersInternalIP(opts HostContainersInternalOptions) string {
switch opts.Conf.Containers.HostContainersInternalIP {
case "":
// If empty (default) we will automatically choose one below.
// If machine using gvproxy we let the gvproxy dns server handle resolve the name and then use that ip.
if machine.IsGvProxyBased() {
if machine.IsPodmanMachine() {
return machineHostContainersInternalIP()
}
case "none":

View File

@@ -0,0 +1,28 @@
package etchosts
import (
"github.com/sirupsen/logrus"
"github.com/vishvananda/netlink"
)
const defaultWSLRoute = "0.0.0.0/0"
// wslHostIP returns the Windows host's IP address. It only make
// sense to execute it when running in a WSL distribution.
// Instructions to retrieve the IP address are section "Identify IP address"
// (scenario 2) of the WSL networking documentation:
// https://learn.microsoft.com/en-us/windows/wsl/networking#identify-ip-address
func wslHostIP() string {
routes, err := netlink.RouteList(nil, netlink.FAMILY_V4)
if err != nil {
logrus.Warnf("Failed getting routes in the WSL machine: %v", err)
return ""
}
for _, r := range routes {
if r.Dst.String() == defaultWSLRoute && r.Gw != nil {
return r.Gw.String()
}
}
logrus.Warnf("No default route found in the WSL machine")
return ""
}

View File

@@ -0,0 +1,9 @@
//go:build !linux
package etchosts
// wslHostIP returns an empty string when running on OSes other than Linux
// (should never happen).
func wslHostIP() string {
return ""
}

View File

@@ -31,3 +31,7 @@ func (n *Netns) Run(lock *lockfile.LockFile, toRun func() error) error {
func (n *Netns) Info() *types.RootlessNetnsInfo {
return &types.RootlessNetnsInfo{}
}
func (n *Netns) PestoSocketPath() string {
return ""
}

View File

@@ -205,10 +205,8 @@ func (n *Netns) setupPasta(nsPath string) error {
extraOpts := []string{"--pid", pidPath}
var socketPath string
if n.config.Network.RootlessPortForwarder == config.RootlessPortForwarderPasta {
socketPath = n.getPath(pestoSocketFile)
extraOpts = append(extraOpts, "-c", socketPath)
extraOpts = append(extraOpts, "-c", n.getPath(pestoSocketFile))
}
pastaOpts := pasta.SetupOptions{
@@ -248,10 +246,9 @@ func (n *Netns) setupPasta(nsPath string) error {
}
n.info = &types.RootlessNetnsInfo{
IPAddresses: res.IPAddresses,
DnsForwardIps: res.DNSForwardIPs,
MapGuestIps: res.MapGuestAddrIPs,
PestoSocketPath: socketPath,
IPAddresses: res.IPAddresses,
DnsForwardIps: res.DNSForwardIPs,
MapGuestIps: res.MapGuestAddrIPs,
}
if err := n.serializeInfo(); err != nil {
return wrapError("serialize info", err)
@@ -632,6 +629,11 @@ func (n *Netns) Info() *types.RootlessNetnsInfo {
return n.info
}
// PestoSocketPath returns the path to the pesto control socket.
func (n *Netns) PestoSocketPath() string {
return n.getPath(pestoSocketFile)
}
func refCount(dir string, inc int) (int, error) {
file := filepath.Join(dir, refCountFile)
content, err := os.ReadFile(file)

View File

@@ -215,3 +215,11 @@ func (n *netavarkNetwork) RootlessNetnsInfo() (*types.RootlessNetnsInfo, error)
}
return n.rootlessNetns.Info(), nil
}
func (n *netavarkNetwork) PestoSocketPath() string {
if n.rootlessNetns == nil {
logrus.Debug("PestoSocketPath: rootlessNetns is nil")
return ""
}
return n.rootlessNetns.PestoSocketPath()
}

View File

@@ -36,6 +36,10 @@ type ContainerNetwork interface {
// Only used as rootless and should return an error as root.
RootlessNetnsInfo() (*RootlessNetnsInfo, error)
// PestoSocketPath returns the path to the pesto control socket
// for dynamic port forwarding. Empty when not available.
PestoSocketPath() string
// Drivers will return the list of supported network drivers
// for this interface.
Drivers() []string
@@ -377,9 +381,6 @@ type RootlessNetnsInfo struct {
DnsForwardIps []string
// MapGuestIps should be used for the host.containers.internal entry when set
MapGuestIps []string
// PestoSocketPath is the path to the pasta control socket for dynamic
// port forwarding via pesto. Empty when pasta was started without -c.
PestoSocketPath string
}
// FilterFunc can be passed to NetworkList to filter the networks.

View File

@@ -635,6 +635,10 @@ type NetworkConfig struct {
// bridge networks. Valid values are RootlessPortForwarderRootlessport
// (default, userspace TCP/UDP proxy) and RootlessPortForwarderPasta
// (experimental, pasta's kernel splice preserving the original source IP).
//
// Must only be changed when no containers are running. Switching while
// containers are active leads to leaked port-forwarding rules or cleanup
// failures.
RootlessPortForwarder string `toml:"rootless_port_forwarder,omitempty"`
}

View File

@@ -416,6 +416,11 @@ default_sysctls = [
# via kernel splice, which preserves the original source IP address inside the
# container. This option is experimental and subject to change.
#
# Important: This option must only be changed when no containers are running.
# Switching while containers are active leads to port-forwarding rules being
# leaked or cleanup failures because the running netns was created with the
# previous setting.
#
#rootless_port_forwarder = "rootlessport"
# Path to the directory where network configuration files are located.

View File

@@ -76,9 +76,12 @@ type Container struct {
// is set before using it.
Created time.Time `json:"created"`
// UIDMap and GIDMap are used for setting up a container's root
// filesystem for use inside of a user namespace where UID mapping is
// being used.
// UIDMap and GIDMap are the caller's requested UID/GID mapping for this
// container's user namespace. They always reflect what the caller
// asked for, regardless of whether the mapping was applied at layer
// creation time (chown) or is deferred to mount time (idmapped mounts).
// At mount time, these maps are passed to the graph driver so that the
// container sees the expected file ownership.
UIDMap []idtools.IDMap `json:"uidmap,omitempty"`
GIDMap []idtools.IDMap `json:"gidmap,omitempty"`

View File

@@ -95,9 +95,6 @@ type MountOpts struct {
// Volatile specifies whether the container storage can be optimized
// at the cost of not syncing all the dirty files in memory.
Volatile bool
// DisableShifting forces the driver to not do any ID shifting at runtime.
DisableShifting bool
}
// ApplyDiffOpts contains optional arguments for ApplyDiff methods.

View File

@@ -1504,7 +1504,7 @@ func (d *Driver) get(id string, disableShifting bool, options graphdriver.MountO
readWrite := !inAdditionalStore
if !d.SupportsShifting(options.UidMaps, options.GidMaps) || options.DisableShifting {
if !d.SupportsShifting(options.UidMaps, options.GidMaps) {
disableShifting = true
}

View File

@@ -167,8 +167,11 @@ type Layer struct {
// Flags is arbitrary data about the layer.
Flags map[string]any `json:"flags,omitempty"`
// UIDMap and GIDMap are used for setting up a layer's contents
// for use inside of a user namespace where UID mapping is being used.
// UIDMap and GIDMap are the on-disk ID mappings for this layer: the
// chown mapping that was applied to the layer's files at creation
// time. When the driver supports shifting (idmapped mounts), no
// chown occurs and these fields are empty. The caller's requested
// mapping is applied at mount time instead (see Container.UIDMap/GIDMap).
UIDMap []idtools.IDMap `json:"uidmap,omitempty"`
GIDMap []idtools.IDMap `json:"gidmap,omitempty"`
@@ -1600,8 +1603,8 @@ func (r *layerStore) create(id string, parentLayer *Layer, names []string, mount
UIDs: templateUIDs,
GIDs: templateGIDs,
Flags: newMapFrom(moreOptions.Flags),
UIDMap: copySlicePreferringNil(moreOptions.UIDMap),
GIDMap: copySlicePreferringNil(moreOptions.GIDMap),
UIDMap: copySlicePreferringNil(moreOptions.IDMappingOptions.UIDMap),
GIDMap: copySlicePreferringNil(moreOptions.IDMappingOptions.GIDMap),
BigDataNames: []string{},
location: r.pickStoreLocation(moreOptions.Volatile, writeable),
}
@@ -1645,7 +1648,10 @@ func (r *layerStore) create(id string, parentLayer *Layer, names []string, mount
}
}
idMappings := idtools.NewIDMappingsFromMaps(moreOptions.UIDMap, moreOptions.GIDMap)
idMappings := idtools.NewIDMappingsFromMaps(moreOptions.IDMappingOptions.UIDMap, moreOptions.IDMappingOptions.GIDMap)
if moreOptions.IDMappingOptions.HostUIDMapping && moreOptions.IDMappingOptions.HostGIDMapping {
idMappings = &idtools.IDMappings{}
}
opts := drivers.CreateOpts{
MountLabel: mountLabel,
StorageOpt: options,
@@ -2601,7 +2607,7 @@ func (r *layerStore) stageWithUnlockedStore(sl *maybeStagedLayerExtraction, pare
result, err := applyDiff(layerOptions, sl.diff, f, func(payload io.Reader) (int64, error) {
cleanup, stagedLayer, size, err := sl.staging.StartStagingDiffToApply(parent, drivers.ApplyDiffOpts{
Diff: payload,
Mappings: idtools.NewIDMappingsFromMaps(layerOptions.UIDMap, layerOptions.GIDMap),
Mappings: idtools.NewIDMappingsFromMaps(layerOptions.IDMappingOptions.UIDMap, layerOptions.IDMappingOptions.GIDMap),
// MountLabel is not supported for the unlocked extraction, see the comment in (*store).PutLayer()
MountLabel: "",
})

View File

@@ -10,6 +10,7 @@ import (
"io"
"io/fs"
"os"
"os/exec"
"path/filepath"
"runtime"
"strings"
@@ -29,6 +30,8 @@ func (dst *unpackDestination) Close() error {
return dst.root.Close()
}
const statusCodeENOSPC = 2
// tarOptionsDescriptor is passed as an extra file
const tarOptionsDescriptor = 3
@@ -82,6 +85,9 @@ func untar() {
}
if err := archive.Unpack(os.Stdin, dst, &options); err != nil {
if errors.Is(err, unix.ENOSPC) {
os.Exit(statusCodeENOSPC)
}
fatal(err)
}
// fully consume stdin in case it is zero padded
@@ -163,6 +169,14 @@ func invokeUnpack(decompressedArchive io.Reader, dest *unpackDestination, option
return fmt.Errorf("%w\nexhausting input failed (error: %w)", errorOut, err)
}
var exitErr *exec.ExitError
if errors.As(err, &exitErr) {
status := exitErr.ExitCode()
if status == statusCodeENOSPC {
return unix.ENOSPC
}
}
return errorOut
}
return nil

View File

@@ -1,25 +0,0 @@
package system
import (
"syscall"
"unsafe"
"golang.org/x/sys/unix"
)
// LUtimesNano is used to change access and modification time of the specified path.
// It's used for symbol link file because unix.UtimesNano doesn't support a NOFOLLOW flag atm.
func LUtimesNano(path string, ts []syscall.Timespec) error {
atFdCwd := unix.AT_FDCWD
var _path *byte
_path, err := unix.BytePtrFromString(path)
if err != nil {
return err
}
if _, _, err := unix.Syscall6(unix.SYS_UTIMENSAT, uintptr(atFdCwd), uintptr(unsafe.Pointer(_path)), uintptr(unsafe.Pointer(&ts[0])), unix.AT_SYMLINK_NOFOLLOW, 0, 0); err != 0 && err != unix.ENOSYS {
return err
}
return nil
}

View File

@@ -1,25 +0,0 @@
package system
import (
"syscall"
"unsafe"
"golang.org/x/sys/unix"
)
// LUtimesNano is used to change access and modification time of the specified path.
// It's used for symbol link file because unix.UtimesNano doesn't support a NOFOLLOW flag atm.
func LUtimesNano(path string, ts []syscall.Timespec) error {
atFdCwd := unix.AT_FDCWD
var _path *byte
_path, err := unix.BytePtrFromString(path)
if err != nil {
return err
}
if _, _, err := unix.Syscall6(unix.SYS_UTIMENSAT, uintptr(atFdCwd), uintptr(unsafe.Pointer(_path)), uintptr(unsafe.Pointer(&ts[0])), unix.AT_SYMLINK_NOFOLLOW, 0, 0); err != 0 && err != unix.ENOSYS {
return err
}
return nil
}

22
vendor/go.podman.io/storage/pkg/system/utimes_unix.go generated vendored Normal file
View File

@@ -0,0 +1,22 @@
//go:build unix
package system
import (
"syscall"
"golang.org/x/sys/unix"
)
// LUtimesNano is used to change access and modification time of the specified path.
func LUtimesNano(path string, ts []syscall.Timespec) error {
ts_ := [2]unix.Timespec{
{Sec: ts[0].Sec, Nsec: ts[0].Nsec},
{Sec: ts[1].Sec, Nsec: ts[1].Nsec},
}
err := unix.UtimesNanoAt(unix.AT_FDCWD, path, ts_[:], unix.AT_SYMLINK_NOFOLLOW)
if err == unix.ENOSYS {
err = nil
}
return err
}

View File

@@ -1,10 +1,10 @@
//go:build !linux && !freebsd
//go:build !unix
package system
import "syscall"
// LUtimesNano is only supported on linux and freebsd.
// LUtimesNano is only supported on Unix systems.
func LUtimesNano(path string, ts []syscall.Timespec) error {
return ErrNotSupportedPlatform
}

100
vendor/go.podman.io/storage/store.go generated vendored
View File

@@ -633,13 +633,39 @@ type AutoUserNsOptions = types.AutoUserNsOptions
type IDMappingOptions = types.IDMappingOptions
// LayerIDMappingOptions are the on-disk ID mappings for a layer.
//
// Unlike the caller-facing IDMappingOptions (which expresses what mapping
// the caller wants), these record how files are actually stored. The two
// may differ: when the graph driver supports shifting, no chown
// occurs so HostUIDMapping/HostGIDMapping are true and UIDMap/GIDMap
// are empty, even though the caller requested a non-trivial mapping.
// The caller's requested mapping is still honored at mount time via
// the Container's UIDMap/GIDMap.
type LayerIDMappingOptions struct {
// HostUIDMapping is true when files in this layer are stored with host
// UIDs.
HostUIDMapping bool
// HostGIDMapping is true when files in this layer are stored with host
// GIDs. See HostUIDMapping for details.
HostGIDMapping bool
// UIDMap is the on-disk UID mapping: it records the chown that was
// applied to the layer's files at creation time. Empty when
// HostUIDMapping is true.
UIDMap []idtools.IDMap
// GIDMap is the on-disk GID mapping: it records the chown that was
// applied to the layer's files at creation time. Empty when
// HostGIDMapping is true.
GIDMap []idtools.IDMap
}
// LayerOptions is used for passing options to a Store's CreateLayer() and PutLayer() methods.
type LayerOptions struct {
// IDMappingOptions specifies the type of ID mapping which should be
// used for this layer. If nothing is specified, the layer will
// inherit settings from its parent layer or, if it has no parent
// layer, the Store object.
types.IDMappingOptions
IDMappingOptions LayerIDMappingOptions
// TemplateLayer is the ID of a layer whose contents will be used to
// initialize this layer. If set, it should be a child of the layer
// which we want to use as the parent of the new layer.
@@ -708,10 +734,18 @@ type ImageBigDataOption struct {
// ContainerOptions is used for passing options to a Store's CreateContainer() method.
type ContainerOptions struct {
// IDMappingOptions specifies the type of ID mapping which should be
// used for this container's layer. If nothing is specified, the
// container's layer will inherit settings from the image's top layer
// or, if it is not being created based on an image, the Store object.
// IDMappingOptions specifies the caller's desired ID mapping for the
// container's user namespace.
//
// These express what the caller wants, not what ends up on disk.
// The store records them in the Container and uses them at mount
// time. How the layer's files are stored depends on whether the
// driver supports shifting: if it does, no chown occurs and the
// mapping is applied at mount time; otherwise files are chowned at
// layer creation time.
//
// If nothing is specified, mappings are inherited from the image's top
// layer or, if no image, from the Store's defaults.
types.IDMappingOptions
LabelOpts []string
// Flags is a set of named flags and their values to store with the container.
@@ -1518,14 +1552,14 @@ func populateLayerOptions(s *store, rlstore rwLayerStore, rlstores []roLayerStor
options.BigData = slices.Clone(lOptions.BigData)
options.Flags = copyMapPreferringNil(lOptions.Flags)
}
if options.HostUIDMapping {
options.UIDMap = nil
if options.IDMappingOptions.HostUIDMapping {
options.IDMappingOptions.UIDMap = nil
}
if options.HostGIDMapping {
options.GIDMap = nil
if options.IDMappingOptions.HostGIDMapping {
options.IDMappingOptions.GIDMap = nil
}
uidMap := options.UIDMap
gidMap := options.GIDMap
uidMap := options.IDMappingOptions.UIDMap
gidMap := options.IDMappingOptions.GIDMap
if parent != "" {
var err error
parentLayer, unlock, err = getParentLayer(rlstore, rlstores, parent)
@@ -1546,26 +1580,26 @@ func populateLayerOptions(s *store, rlstore rwLayerStore, rlstores []roLayerStor
return nil, nil, unlock, ErrParentIsContainer
}
}
if !options.HostUIDMapping && len(options.UIDMap) == 0 {
if !options.IDMappingOptions.HostUIDMapping && len(options.IDMappingOptions.UIDMap) == 0 {
uidMap = parentLayer.UIDMap
}
if !options.HostGIDMapping && len(options.GIDMap) == 0 {
if !options.IDMappingOptions.HostGIDMapping && len(options.IDMappingOptions.GIDMap) == 0 {
gidMap = parentLayer.GIDMap
}
} else {
if !options.HostUIDMapping && len(options.UIDMap) == 0 {
if !options.IDMappingOptions.HostUIDMapping && len(options.IDMappingOptions.UIDMap) == 0 {
uidMap = s.uidMap
}
if !options.HostGIDMapping && len(options.GIDMap) == 0 {
if !options.IDMappingOptions.HostGIDMapping && len(options.IDMappingOptions.GIDMap) == 0 {
gidMap = s.gidMap
}
}
if s.canUseShifting(uidMap, gidMap) {
options.IDMappingOptions = types.IDMappingOptions{HostUIDMapping: true, HostGIDMapping: true, UIDMap: nil, GIDMap: nil}
options.IDMappingOptions = LayerIDMappingOptions{HostUIDMapping: true, HostGIDMapping: true, UIDMap: nil, GIDMap: nil}
} else {
options.IDMappingOptions = types.IDMappingOptions{
HostUIDMapping: options.HostUIDMapping,
HostGIDMapping: options.HostGIDMapping,
options.IDMappingOptions = LayerIDMappingOptions{
HostUIDMapping: options.IDMappingOptions.HostUIDMapping,
HostGIDMapping: options.IDMappingOptions.HostGIDMapping,
UIDMap: copySlicePreferringNil(uidMap),
GIDMap: copySlicePreferringNil(gidMap),
}
@@ -1856,14 +1890,14 @@ func (s *store) imageTopLayerForMapping(image *Image, ristore roImageStore, rlst
// mappings, and register it as an alternate top layer in the image.
var layerOptions LayerOptions
if s.canUseShifting(options.UIDMap, options.GIDMap) {
layerOptions.IDMappingOptions = types.IDMappingOptions{
layerOptions.IDMappingOptions = LayerIDMappingOptions{
HostUIDMapping: true,
HostGIDMapping: true,
UIDMap: nil,
GIDMap: nil,
}
} else {
layerOptions.IDMappingOptions = types.IDMappingOptions{
layerOptions.IDMappingOptions = LayerIDMappingOptions{
HostUIDMapping: options.HostUIDMapping,
HostGIDMapping: options.HostGIDMapping,
UIDMap: copySlicePreferringNil(options.UIDMap),
@@ -2008,20 +2042,12 @@ func (s *store) CreateContainer(id string, names []string, image, layer, metadat
// But in transient store mode, all container layers are volatile.
Volatile: options.Volatile || s.transientStore,
}
if s.canUseShifting(uidMap, gidMap) {
layerOptions.IDMappingOptions = types.IDMappingOptions{
HostUIDMapping: true,
HostGIDMapping: true,
UIDMap: nil,
GIDMap: nil,
}
} else {
layerOptions.IDMappingOptions = types.IDMappingOptions{
HostUIDMapping: idMappingsOptions.HostUIDMapping,
HostGIDMapping: idMappingsOptions.HostGIDMapping,
UIDMap: copySlicePreferringNil(uidMap),
GIDMap: copySlicePreferringNil(gidMap),
}
useHostMapping := idMappingsOptions.HostUIDMapping || s.canUseShifting(uidMap, gidMap)
layerOptions.IDMappingOptions = LayerIDMappingOptions{
HostUIDMapping: useHostMapping,
HostGIDMapping: useHostMapping,
UIDMap: copySlicePreferringNil(uidMap),
GIDMap: copySlicePreferringNil(gidMap),
}
if options.Flags == nil {
options.Flags = make(map[string]any)
@@ -3074,10 +3100,6 @@ func (s *store) Mount(id, mountLabel string) (string, error) {
if err != nil {
return "", err
}
if options.UidMaps != nil || options.GidMaps != nil {
options.DisableShifting = !s.canUseShifting(options.UidMaps, options.GidMaps)
}
if err := rlstore.startWriting(); err != nil {
return "", err
}

View File

@@ -28,21 +28,34 @@ type AutoUserNsOptions struct {
AdditionalGIDMappings []idtools.IDMap
}
// IDMappingOptions are used for specifying how ID mapping should be set up for
// a layer or container.
// IDMappingOptions specifies the caller's desired UID/GID mapping for a
// layer or container.
//
// These options express what the caller wants, the mapping that the
// container's user namespace should use. They do not describe what is
// stored on disk. Depending on the graph driver, the store may apply
// the mapping at layer creation time (by chowning files) or defer it to
// mount time (using idmapped mounts or fuse-overlayfs options), but
// that distinction is transparent to the caller.
//
// The resolution order for the effective UID/GID maps is:
// 1. If HostUIDMapping/HostGIDMapping is true, no mapping is used (the
// corresponding UIDMap/GIDMap is ignored and treated as empty).
// 2. If UIDMap/GIDMap contain at least one entry, those mappings are used.
// 3. Otherwise, if the layer has a parent, the parent's mappings are inherited.
// 4. Otherwise, the Store-level default mappings are used.
type IDMappingOptions struct {
// UIDMap and GIDMap are used for setting up a layer's root filesystem
// for use inside of a user namespace where ID mapping is being used.
// If HostUIDMapping/HostGIDMapping is true, no mapping of the
// respective type will be used. Otherwise, if UIDMap and/or GIDMap
// contain at least one mapping, one or both will be used. By default,
// if neither of those conditions apply, if the layer has a parent
// layer, the parent layer's mapping will be used, and if it does not
// have a parent layer, the mapping which was passed to the Store
// object when it was initialized will be used.
// HostUIDMapping indicates that no UID mapping should be applied.
// When true, UIDMap is ignored and files are accessed with host UIDs.
HostUIDMapping bool
// HostGIDMapping indicates that no GID mapping should be applied.
// When true, GIDMap is ignored and files are accessed with host GIDs.
HostGIDMapping bool
UIDMap []idtools.IDMap
// UIDMap defines the UID mappings for the user namespace.
// Only used when HostUIDMapping is false.
UIDMap []idtools.IDMap
// GIDMap defines the GID mappings for the user namespace.
// Only used when HostGIDMapping is false.
GIDMap []idtools.IDMap
AutoUserNs bool
AutoUserNsOpts AutoUserNsOptions

View File

@@ -187,7 +187,7 @@ outer:
}
layerOptions := &LayerOptions{
IDMappingOptions: types.IDMappingOptions{
IDMappingOptions: LayerIDMappingOptions{
HostUIDMapping: true,
HostGIDMapping: true,
UIDMap: nil,

8
vendor/modules.txt vendored
View File

@@ -314,7 +314,7 @@ github.com/json-iterator/go
# github.com/kevinburke/ssh_config v1.5.0
## explicit; go 1.18
github.com/kevinburke/ssh_config
# github.com/klauspost/compress v1.18.6
# github.com/klauspost/compress v1.19.0
## explicit; go 1.24
github.com/klauspost/compress
github.com/klauspost/compress/flate
@@ -740,7 +740,7 @@ go.podman.io/buildah/pkg/sshagent
go.podman.io/buildah/pkg/util
go.podman.io/buildah/pkg/volumes
go.podman.io/buildah/util
# go.podman.io/common v0.68.1-0.20260629144122-498ccd1bb53a
# go.podman.io/common v0.68.1-0.20260707152203-d126613c6575
## explicit; go 1.25.7
go.podman.io/common/internal
go.podman.io/common/libimage
@@ -806,7 +806,7 @@ go.podman.io/common/pkg/umask
go.podman.io/common/pkg/util
go.podman.io/common/pkg/version
go.podman.io/common/version
# go.podman.io/image/v5 v5.40.1-0.20260629144122-498ccd1bb53a
# go.podman.io/image/v5 v5.40.1-0.20260707152203-d126613c6575
## explicit; go 1.25.7
go.podman.io/image/v5/copy
go.podman.io/image/v5/directory
@@ -883,7 +883,7 @@ go.podman.io/image/v5/transports
go.podman.io/image/v5/transports/alltransports
go.podman.io/image/v5/types
go.podman.io/image/v5/version
# go.podman.io/storage v1.63.1-0.20260629144122-498ccd1bb53a
# go.podman.io/storage v1.63.1-0.20260707152203-d126613c6575
## explicit; go 1.25.0
go.podman.io/storage
go.podman.io/storage/drivers