collections

package module
v0.30.11 Latest Latest
Warning

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

Go to latest
Published: Sep 26, 2026 License: Apache-2.0 Imports: 36 Imported by: 0

Documentation

Overview

Package collections is an in-memory, sharded, memory-dense store of ClassAds with a compiled query engine. Ads are serialized to the compact wire form and packed into append-only arena segments; a per-shard directory maps a stable key hash to a record location. Updates are MVCC-stamped so table scans see every ad exactly once even while background compaction moves bytes.

Index

Constants

View Source
const DefaultDictSize = 112 * 1024

DefaultDictSize is the target size in bytes for a trained dictionary's content.

View Source
const DefaultRetrainInterval = 15 * time.Minute

DefaultRetrainInterval is the suggested cadence for background dictionary retraining. A dictionary trained on a representative sample stays effective as long as the ad population's shape is stable, so retraining need not be frequent.

Variables

View Source
var ErrTimeTravelDisabled = errors.New("collections: time travel is not enabled on this collection")

ErrTimeTravelDisabled is returned by the AS OF query paths when the collection has no time-travel configuration.

View Source
var OpenIndexDiagHook func(OpenIndexDiag)

OpenIndexDiagHook, when non-nil, is called once at the end of each Collection Open with a summary of sidecar adoption for that store. Set it (once, at process start) to log why a reopen rebuilds indexes -- the slow-startup path -- rather than adopting the persisted ones. Diagnostic only; the callback must not block (it runs on the Open path).

Functions

func ColNativeCRCFailures added in v0.28.0

func ColNativeCRCFailures() int64

ColNativeCRCFailures reports how many columnar payloads have been refused for a bad checksum since the process started.

func CollapseDeferrals added in v0.30.0

func CollapseDeferrals() int64

CollapseDeferrals reports how many delta-chain collapses have been postponed to a later pass.

func ColumnarDeltaDisposition added in v0.30.0

func ColumnarDeltaDisposition() (dropped, materialized, unresolved int64)

ColumnarDeltaDisposition reports what columnarization did with the delta records it met: dropped as garbage, materialized into whole ads, and left alone because the chain would not resolve.

func ColumnarDeltasDropped added in v0.30.0

func ColumnarDeltasDropped() int64

ColumnarDeltasDropped reports how many delta records columnarization has dropped rather than carried into a rewritten segment.

func ColumnarFallbacks added in v0.28.0

func ColumnarFallbacks() int64

ColumnarFallbacks reports how many records the columnar accelerator has deferred to the row path since the process started.

It is a correctness-adjacent metric rather than a performance one. Deferring is always right, so a rising count is the only outward sign that the accelerator is being bypassed -- for instance because a block is addressing attribute ids from a previous process's intern table. A durability test that asserts this does not grow across a restart catches that; one that only compares query answers cannot, because the answers stay correct.

func ColumnarizedSegments added in v0.28.0

func ColumnarizedSegments() (segments, bytesSaved int64)

ColumnarizedSegments reports how many segments have been rewritten into columnar form, and how many arena bytes that removed, since the process started.

func CompactDeferredLiveDelta added in v0.30.9

func CompactDeferredLiveDelta() int64

CompactDeferredLiveDelta reports how many shard compactions were deferred because the active segment still held a live delta. See compactDeferredLiveDelta.

func CompactInternPaths added in v0.30.0

func CompactInternPaths() (wirePath, astPath int64)

CompactInternPaths reports how many records compaction re-interned by transcoding the wire bytes versus by decoding into an ast and encoding it back.

func ConflictChecks added in v0.6.0

func ConflictChecks() int64

ConflictChecks returns the cumulative number of per-key conflict checks committed transactions have performed -- zero while a single writer runs (the fast path).

func CorruptChainLinks() int64

CorruptChainLinks reports how many times a bucket-chain walk found a link naming a segment that does not exist, across every collection in this process.

It is exported so an operator can see it without a debugger. Nonzero means some reader walked a segment whose mapping it did not hold alive, and the walk read whatever replaced it; the records themselves are not corrupt on disk.

func DeltaAnomalies added in v0.30.0

func DeltaAnomalies() (liveAtCompact, droppedFromHistory, flagMismatch int64)

DeltaAnomalies reports the three things that should all stay zero in a healthy delta store: live deltas met by compaction, superseded deltas dropped from the travel window, and records whose header and payload disagree about being a delta.

func FallbackReasons added in v0.30.0

func FallbackReasons() (removal, bound, noBase, ineligible int64)

FallbackReasons reports why patch writes stored whole records rather than deltas.

func IsSystemKey added in v0.9.0

func IsSystemKey(key string) bool

IsSystemKey reports whether key is a reserved internal system key (begins with a NUL byte). Exported so callers (and the db/dbrpc layers) can classify a key without duplicating the sentinel.

func LastDeltaDecodeFailure added in v0.30.5

func LastDeltaDecodeFailure() (stage, msg string)

LastDeltaDecodeFailure returns the most recently sampled decode error behind a refused chain merge, or ("", "") if none has happened.

func LastNoBaseDetail added in v0.30.7

func LastNoBaseDetail() (versions int, chainBroken bool, sealedSkipped int)

LastNoBaseDetail returns the shape of the most recent no-base chain: how many versions the walk did find, whether it truncated at a dead link, and how many sealed segments it could not probe.

func LatencyBucketBoundsNanos added in v0.12.1

func LatencyBucketBoundsNanos() []int64

LatencyBucketBoundsNanos returns a copy of the histogram's inclusive upper bounds (in nanoseconds). An OpStat's Buckets slice has one more entry than this: buckets[i] counts observations <= bound[i], and the final entry counts observations above the last bound.

func MergedBytes added in v0.24.1

func MergedBytes() int64

MergedBytes reports the total source bytes rewritten by merges since this process started.

func MigrateSkippedLiveDelta added in v0.30.10

func MigrateSkippedLiveDelta() int64

MigrateSkippedLiveDelta reports how many shards the sealed-attribute migration skipped because their active segment still held a live delta. See migrateSkippedLiveDelta.

func ProjectionChecksum added in v0.5.2

func ProjectionChecksum(ad *classad.ClassAd, attrs []string) uint64

ProjectionChecksum returns a 64-bit FNV-1a checksum over the given attributes' literal expression text in ad, taken in the order given and delimited so two distinct projections cannot alias. A missing attribute contributes a fixed marker, so "absent" is distinguished from every present value.

This is the standalone, caller-supplied-attrs form of an ordered index's cluster Signature (OrderSpec.Cluster). Unlike that index-time signature -- whose attribute set is fixed when the index is opened -- ProjectionChecksum takes the projection at call time, for callers whose significant attributes vary per query (e.g. a schedd whose negotiator sends the significant-attribute set in each negotiation header).

It hashes each attribute's *expression text* rather than its evaluated value, matching HTCondor autocluster semantics: two job ads whose significant attributes are textually identical (same Requirements expression, same RequestMemory literal, ...) hash equal, so a run-length fold over a priority-ordered stream folds them into one resource request with a stored-value compare. A 64-bit collision could merge two adjacent runs; over the projected value bytes that is negligible.

func ProvisionalIndexBuilds added in v0.30.11

func ProvisionalIndexBuilds() int64

ProvisionalIndexBuilds reports how many sealed segments were given a provisional key index at seal time. See segment.keyIdxMem.

func RenderRawAdInline added in v0.16.4

func RenderRawAdInline(w []byte, buf []byte, offs []int) (outBuf []byte, outOffs []int, myType, targetType string, ok bool)

RenderRawAdInline renders a wire-form row (a self-contained inline-names ad, e.g. one shipped by ScanRawWire over dbrpc) to old-ClassAd "Name = Value" expression text: each expression is buf[offs[i]:offs[i+1]], with MyType/TargetType lifted out exactly as the collection scans do -- the shape message.PutClassAdRawBytes consumes. This is the LAST-MINUTE conversion at the client edge; everything upstream of it stays in wire form. buf/offs are reset and reused (pass the previous call's returns to avoid allocation).

func SealWalkStats added in v0.30.0

func SealWalkStats() (examined, deltas int64)

SealWalkStats reports what the collapse walk has examined and found: live records looked at, and live deltas among them. The ratio is the walk's yield -- whether it is finding chains to collapse or scanning a segment to conclude there is nothing to do.

func SealedProbesSkipped added in v0.30.7

func SealedProbesSkipped() int64

SealedProbesSkipped reports how many sealed-segment probes were skipped for want of a key index. A rising count alongside delta-index-pending says reads are racing the reindex pass.

func SetArchiveBlockCacheBudget added in v0.29.6

func SetArchiveBlockCacheBudget(bytes int64)

func SetMutatingBlockCacheBudget added in v0.29.6

func SetMutatingBlockCacheBudget(bytes int64)

SetMutatingBlockCacheBudget and SetArchiveBlockCacheBudget set the process-global shared block-cache byte budget for each workload kind. They are the config knob behind Options' {Mutating,Archive}BlockCacheBytes; db/ (and above it htcondordb) should call them once at startup from configuration. Safe to call at any time and concurrently: if the kind's cache already exists it is resized in place. A value ≤ 0 is ignored (keeps the current or default budget).

func SpliceStats added in v0.30.0

func SpliceStats() (spliced, spliceRefused, objectPath int64)

SpliceStats reports how delta merges were served: by splicing attribute bytes, by decoding after the splice refused, and by decoding because the caller asked for an object rather than bytes (a point read, or a collapse in a collection with an ordered index or a watcher).

func StrandedSealedDeltas added in v0.30.8

func StrandedSealedDeltas() int64

StrandedSealedDeltas reports how many live delta records were found in sealed segments at open. Any value above zero means the seal-collapse invariant is broken on disk, and the keys behind those records are on their way to becoming unreadable.

func SystemKey added in v0.9.0

func SystemKey(name string) string

SystemKey builds a system key from a name by prefixing the reserved NUL byte. The name is otherwise opaque; callers namespace it however they like.

func TrainDict

func TrainDict(samples [][]byte) ([]byte, error)

TrainDict builds a ZSTD compression dictionary from sample records (the wire-encoded bytes of a representative set of ads; see CollectSamples). The resulting dictionary can be handed to NewZSTDCodec. It uses DefaultDictSize.

func TrainDictSize

func TrainDictSize(samples [][]byte, dictSize int) (dict []byte, err error)

TrainDictSize is TrainDict with an explicit dictionary content size.

The pure-Go zstd.BuildDict does not perform ZDICT-style content *selection* (the cover algorithm); the dictionary content is whatever we supply as the builder's History. We therefore assemble the content ourselves by concatenating distinct samples up to dictSize — for a pool of similar ClassAds, this captures the shared attribute names and values that later ads back- reference. BuildDict then trains the entropy tables from the full sample set.

func UnreadableBaseReasonNames added in v0.30.3

func UnreadableBaseReasonNames() []string

UnreadableBaseReasonNames lists every reason name UnreadableBaseReasons can report, including the ones currently at zero. A consumer publishing these as metrics needs the whole set up front: a reason that appears only once it is non-zero reads as "no such counter" at exactly the moment someone is checking whether it is the one firing.

func UnreadableBaseReasons added in v0.30.3

func UnreadableBaseReasons() map[string]int64

UnreadableBaseReasons reports why patch writes were refused for an unreadable base, keyed by reason name. The values sum to UnreadableBaseRefusals.

func UnreadableBaseRefusals added in v0.30.1

func UnreadableBaseRefusals() int64

UnreadableBaseRefusals reports, process-wide, how many patch writes were refused because the key was present but its current record could not be read. Nonzero means the store failed to resolve a key it holds -- the write was reported as a conflict rather than stored as a fragment.

func WatchFilter added in v0.4.0

func WatchFilter(seq iter.Seq[WatchEvent], match func(*classad.ClassAd) bool) iter.Seq[WatchEvent]

WatchFilter wraps a Watch event stream to deliver only events for keys whose ad satisfies match. It keeps a set of the keys currently matching so a filtered view stays correct as ads change:

  • an Upsert whose ad matches is delivered (the key is marked matching);
  • an Upsert whose ad no longer matches, for a key that was matching, is converted to a Delete so the client drops it from its filtered view;
  • a Delete is delivered for a key known to be matching, and -- during catch-up (before Synced), where the prior match state of a resumed key is unknown -- forwarded conservatively (the client no-ops an unknown key);
  • Reset clears the matched set; Synced and Resync pass through.

A nil match returns seq unchanged (no filtering). match is called on each Upsert's ad and must be safe for concurrent-free sequential use.

Types

type AdUpdate

type AdUpdate struct {
	Key []byte
	Ad  *classad.ClassAd
}

AdUpdate is one insert-or-update in a batch.

type Archive

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

Archive is an append-only, larger-than-RAM, rotated store of ClassAds — the "condor history file" use case. Ads are appended once and never updated; old data is dropped in bulk by rotation, with whole-segment pruning via zone maps and newest-first + LIMIT queries.

It is a thin facade over a Collection configured as an append log (Options.AppendOnly) with newest-first scans, retention rotation, and zone maps. All the archive-specific machinery — sealed segments, per-segment sidecar indexes, zone pruning, O(segments) recovery, retention, and the append-stream watch — now lives in the mainline Collection, so the archive is "a persistent store with a few rules and no compaction" rather than a fork. See docs/history-archive.md.

func CreateArchive

func CreateArchive(opts ArchiveOptions) (*Archive, error)

CreateArchive creates a new, empty Archive. The directory must not already hold one (use OpenArchive to reopen).

func OpenArchive

func OpenArchive(opts ArchiveOptions) (*Archive, error)

OpenArchive reopens an existing Archive (or creates one if the directory is empty), recovering its segments. The same options must be supplied as at creation (the Collection recovers each segment with the codec it names).

func (*Archive) AddIndex added in v0.18.0

func (a *Archive) AddIndex(categorical, value []string) bool

AddIndex adds per-segment indexes on the named categorical and/or value attributes and rebuilds every segment's index under the new set, so existing records are covered too. Returns false if the index set was unchanged.

Cost: the rebuild decompresses every record in the archive once. It does NOT rewrite any segment data -- a sealed segment's bytes are already correct and are left untouched; only the derived index sidecar beside them is rebuilt (see Collection.reindexSealed). That is the difference from Rewrite, which re-encodes and rewrites the whole store.

func (*Archive) Append

func (a *Archive) Append(ad *classad.ClassAd) error

Append adds one ad to the archive: it is appended to the log (never superseding any prior record) and, when the active segment fills, a new one is started and the old one sealed with its index and zone maps. Safe for one writer concurrent with queries.

func (*Archive) AutoIndexNames added in v0.25.0

func (a *Archive) AutoIndexNames() []string

AutoIndexNames returns the names of the archive's auto-created indexes, for persisting provenance alongside the index set itself.

func (*Archive) AutoTune added in v0.25.0

func (a *Archive) AutoTune(opts AutoTuneOptions) AutoTuneResult

AutoTune applies demand-driven index changes to the archive. See Collection.AutoTune.

func (*Archive) BuildAndEnableSchemaScan added in v0.25.2

func (a *Archive) BuildAndEnableSchemaScan(sampleMax, hotTopN int) bool

BuildAndEnableSchemaScan builds (or extends coverage of) the per-segment columnar accelerator over the archive's sealed segments, as Collection.BuildAndEnableSchemaScan does for a mutable table. An archive is the natural target for it -- the segments are immutable once sealed, so a block never needs invalidating -- and the payoff is large: a numeric COUNT over 50k history records drops from ~50ms of row scan to ~1ms.

The first call reads every sealed record once to sample and to build a block per segment, which over a long history is the same "paid for by reading history back" cost an archive index backfill has. Later calls only cover segments sealed since.

func (*Archive) CategoricalGroupCounts added in v0.23.1

func (a *Archive) CategoricalGroupCounts(attr string) (map[string]int64, bool)

CategoricalGroupCounts returns the exact number of records carrying each distinct value of a categorically indexed attribute, derived from the per-segment indexes instead of by scanning records. ok is false when the answer cannot be established from the indexes alone, in which case the caller must scan; see the package comment above for what makes it decline.

The counts are keyed by the value's exact spelling, matching what a scan would group by: ClassAd string comparison folds case, but two spellings of the same value are distinct groups, and the index keeps its exact-case run for precisely that reason.

func (*Archive) CategoricalGroupCountsBucketed added in v0.23.1

func (a *Archive) CategoricalGroupCountsBucketed(attr, bucketAttr string, width int64) (map[int64]map[string]int64, bool)

CategoricalGroupCountsBucketed is CategoricalGroupCounts split by a second, numeric dimension: it returns per-value counts keyed by the bucket floor(bucketAttr/width)*width. This is the "per group per day" shape, with width the bucket size in the attribute's own units (86400 for daily buckets over a unix-seconds attribute).

It is cheap for the same reason the ungrouped form is, plus one more: an archive is append-ordered, so a segment's values for a monotonically advancing attribute like a completion time span a narrow range. When a segment's zone map for bucketAttr falls wholly inside one bucket -- the common case -- every record in it belongs to that bucket and its per-value counts are attributed wholesale, with no record read. Only the segments straddling a bucket boundary are scanned, and with buckets much wider than a segment's time span those are a small minority.

bucketAttr must carry a zone map (ZoneAttrs at creation); without one there is no way to place a segment in a bucket without reading it, and the whole point is lost, so this declines. It also declines for the same completeness reasons as CategoricalGroupCounts.

func (*Archive) CategoricalGroupCountsWhere added in v0.23.1

func (a *Archive) CategoricalGroupCountsWhere(attr, constraint string) (map[string]int64, bool)

CategoricalGroupCountsWhere is CategoricalGroupCounts restricted to the records matching constraint. ok is false unless the constraint is a pure conjunction of numeric comparisons against zone-mapped attributes AND the indexes fully account for every record (the same completeness bar as the unconstrained form) -- see the file comment for why the constraint's syntax is re-checked rather than its probes trusted.

Segments are decided by their zone maps: one whose bounds imply every condition contributes its postings wholesale, one that cannot satisfy some condition contributes nothing, and only the genuinely straddling ones are read.

func (*Archive) Close

func (a *Archive) Close() error

Close flushes and unmaps the archive. It must not be used afterward.

func (*Archive) CodecStats added in v0.19.0

func (a *Archive) CodecStats(sampleMax int) CodecStats

CodecStats reports the archive's compression: codec, dictionary size, last retrain time, and the sampled compression ratio (up to sampleMax records).

func (*Archive) Count added in v0.7.0

func (a *Archive) Count() int

Count returns the number of records currently retained. Rotation reduces it in whole-segment steps.

func (*Archive) CountConstraint added in v0.27.0

func (a *Archive) CountConstraint(constraint string) (int, bool)

CountConstraint counts the records matching constraint via the columnar accelerator, returning ok=false when the columnar path cannot serve it (see Collection.CountConstraint) so the caller scans instead.

An archive is the big append-only scan target, and it is exactly where this was missing: a mutable table's COUNT(*) WHERE reaches the columnar count through DB.CountConstraint, while an archive's went to a wire-native scan of every record, even with an accelerator built over its segments.

func (*Archive) DropIndex added in v0.18.0

func (a *Archive) DropIndex(names ...string) bool

DropIndex removes the named per-segment indexes. Returns false if none matched.

func (*Archive) EncryptedAttrNames added in v0.29.0

func (a *Archive) EncryptedAttrNames() []string

EncryptedAttrNames reports the attributes this archive seals, mirroring Collection.EncryptedAttrNames.

func (*Archive) EncryptionEnabled added in v0.29.0

func (a *Archive) EncryptionEnabled() bool

EncryptionEnabled reports whether this archive's values are sealed (mirroring Collection.EncryptionEnabled). It is currently FALSE for every archive a catalog opens: the archive open path passes no data key, so an archive stores private attributes in the clear where a mutable table always seals them. Exposing it is what makes that asymmetry visible instead of absent from the diagnostics -- an operator reading "encryption at rest: on" for a table should not have to guess what its history table does.

func (*Archive) ExplainQuery added in v0.25.2

func (a *Archive) ExplainQuery(q *vm.Query) QueryExplain

ExplainQuery reports how the archive would execute q -- which conjuncts are index-usable and the resulting access path -- as Collection.ExplainQuery does for a mutable table.

Two archive-specific notes on reading the result. An archive's scan is newest-first over sealed segments, and whole segments are additionally skipped by zone map (see ZoneAttrs), which the plan string does not name: a "serial-scan" over an archive may still visit far fewer records than TotalAds. And TotalAds counts everything retained, so probe selectivity is against the whole history rather than a working set.

func (*Archive) Flush

func (a *Archive) Flush() error

Flush makes prior appends durable and queryable. The Collection keeps appended records queryable in place and mmap-durable across a process crash without an explicit seal, so this is a no-op retained for API compatibility.

func (*Archive) GCFloor added in v0.21.4

func (a *Archive) GCFloor() float64

GCFloor returns the current runtime GC watermark (0 when unset).

func (*Archive) GroupCountAll added in v0.28.0

func (a *Archive) GroupCountAll(groupAttr string) ([]GroupCount, bool)

GroupCountAll returns the per-value record counts of groupAttr over EVERY record, read out of the columnar blocks. The caller must have established that its query matches all records. ok=false ⇒ not columnar-eligible; group by scanning records.

func (*Archive) GroupCountConstraint added in v0.28.0

func (a *Archive) GroupCountConstraint(constraint, groupAttr string) ([]GroupCount, bool)

GroupCountConstraint returns the per-value record counts of groupAttr over the records matching constraint, read out of the columnar blocks. ok=false ⇒ the shape is not columnar-eligible and the caller must group by scanning records.

func (*Archive) GroupSchemaAgreement added in v0.28.0

func (a *Archive) GroupSchemaAgreement(sampleMax, k int) GroupSchemaAgreement

GroupSchemaAgreement reports how well per-segment derivations agree with the archive-wide one (see Collection.GroupSchemaAgreement).

func (*Archive) GroupSchemaChanges added in v0.29.13

func (a *Archive) GroupSchemaChanges() []GroupSchemaChange

GroupSchemaChanges returns the archive's committed-group change log (see Collection.GroupSchemaChanges).

func (*Archive) GroupSchemaDrift added in v0.28.0

func (a *Archive) GroupSchemaDrift() GroupSchemaDrift

GroupSchemaDrift reports how the archive's derived groups have moved (see Collection.GroupSchemaDrift).

func (*Archive) GroupSchemaLastAgreement added in v0.29.13

func (a *Archive) GroupSchemaLastAgreement() (GroupSchemaLastAgreement, bool)

GroupSchemaLastAgreement returns the archive's last persisted per-segment agreement (see Collection.GroupSchemaLastAgreement).

func (*Archive) GroupSchemas added in v0.28.0

func (a *Archive) GroupSchemas(sampleMax, k int) GroupSchemaInfo

GroupSchemas derives and reports candidate group schemas for the archive (see Collection.GroupSchemas). Report-only: nothing about storage or reads changes.

func (*Archive) GroupStatsAll added in v0.28.0

func (a *Archive) GroupStatsAll(groupAttr string, aggAttrs []string) ([]GroupStats, bool)

GroupStatsAll is GroupStatsConstraint over EVERY record. The caller must have established that its query matches all records.

func (*Archive) GroupStatsConstraint added in v0.28.0

func (a *Archive) GroupStatsConstraint(constraint, groupAttr string, aggAttrs []string) ([]GroupStats, bool)

GroupStatsConstraint returns the per-group record count plus a NumStats per aggAttr over the records matching constraint, read out of the columnar blocks. ok=false ⇒ not columnar-eligible; group by scanning records. See Collection.GroupStatsQuery.

func (*Archive) HotAttrNames added in v0.29.0

func (a *Archive) HotAttrNames() []string

HotAttrNames reports the attributes stored in the hot header, mirroring Collection.HotAttrNames. An archive has a hot set exactly as a mutable collection does; only the accessor was missing, which made an archive look like it had none.

func (*Archive) IndexSizes added in v0.19.0

func (a *Archive) IndexSizes() IndexSizes

IndexSizes reports the per-attribute index byte footprint (heap postings) versus data.

func (*Archive) IndexedAttrs added in v0.19.0

func (a *Archive) IndexedAttrs() (categorical, value []string)

IndexedAttrs returns the categorical and value attributes the archive indexes.

func (*Archive) MarkAutoIndexes added in v0.25.0

func (a *Archive) MarkAutoIndexes(names []string)

MarkAutoIndexes records existing indexes as auto-created, restoring the provenance a saved config carries without rebuilding anything. See Collection.MarkAutoIndexes.

func (*Archive) MergePass added in v0.24.1

func (a *Archive) MergePass(opts MergeOptions) int

MergePass merges cold segment runs when the archive has reached opts.TriggerSegments, continuing until it is at or below opts.TargetSegments, no eligible run remains, or opts.MaxMerges merges have been done. It returns the number of merges performed.

The two watermarks are the point: merging is rewriting, so it should happen in occasional batches with headroom either side, not continuously to pin the count at one value. Below the trigger a pass is a cheap no-op (a segment count and a comparison), so it is safe to call often.

Safe to call on a live archive: merges take the maintenance lock (so they never overlap a reseal or another pass), the writer's open segment is never a candidate, and an in-flight scan holding a source keeps its mapping alive through the usual pin/reap path.

func (*Archive) NumStatsQuery added in v0.25.2

func (a *Archive) NumStatsQuery(q *vm.Query, attr string) (NumStats, bool)

NumStatsQuery computes numeric aggregate inputs for attr over the archive's matching records via the columnar scan (see Collection.NumStatsQuery). Only meaningful once the archive carries a columnar accelerator; ok=false otherwise, and the caller scans.

func (*Archive) OpStats added in v0.19.0

func (a *Archive) OpStats() OpStats

OpStats reports cumulative operational timing counters (write wait/hold, sync, retrain, reindex, ...) for the archive's writes and maintenance.

func (*Archive) Query

func (a *Archive) Query(q *vm.Query) iter.Seq[*classad.ClassAd]

Query returns an iterator over the ads matching q, newest first.

func (*Archive) QueryLimit

func (a *Archive) QueryLimit(q *vm.Query, limit int) iter.Seq[*classad.ClassAd]

QueryLimit returns an iterator over the ads matching q, newest first, stopping after limit results (limit <= 0 ⇒ unlimited). Because the backing Collection scans newest-first, stopping early yields the most recent limit matches — a pushed-down LIMIT.

func (*Archive) QueryProject added in v0.20.0

func (a *Archive) QueryProject(q *vm.Query, attrs []string) iter.Seq[[]classad.Value]

func (*Archive) QueryProjectStats added in v0.29.6

func (a *Archive) QueryProjectStats(q *vm.Query, attrs []string, stats *ScanStats) iter.Seq[[]classad.Value]

QueryProjectStats is QueryProject that also fills stats (may be nil) with the scan's work breakdown, for an EXPLAIN ANALYZE of an aggregate that falls to the projected scan.

func (*Archive) QueryRawProjected added in v0.20.1

func (a *Archive) QueryRawProjected(q *vm.Query, projection []string, chaseRefs, redact bool) iter.Seq[RawAd]

QueryProject scans the ads matching q and yields each one projected to just attrs' values (aligned with attrs), read wire-native where possible -- so an aggregate does not pay the full-ad decode QueryLimit costs. The yielded slice is reused; copy to retain. QueryRawProjected yields each ad matching q as a raw projected subset (only the projection attributes, rendered from the stored representation), newest first — the archive-side of the server projection relay. It shares the backing Collection's projection walk, so the same newest-first ordering and pushed-down LIMIT (via early stop by the caller) apply. chaseRefs and redact are as in Collection.QueryRawProjected.

func (*Archive) QueryRawProjectedStats added in v0.29.0

func (a *Archive) QueryRawProjectedStats(q *vm.Query, projection []string, chaseRefs, redact bool, stats *ScanStats) iter.Seq[RawAd]

QueryRawProjectedStats is QueryRawProjected that also fills stats (may be nil) with the per-scan work breakdown for EXPLAIN ANALYZE.

func (*Archive) Reindex added in v0.18.0

func (a *Archive) Reindex()

Reindex rebuilds the per-segment indexes over all segments: those sealed since the last build, and those whose sidecar was built under a superseded index configuration. Segment data is never touched.

func (*Archive) ReschemaScan added in v0.25.3

func (a *Archive) ReschemaScan(sampleMax, hotTopN int) bool

ReschemaScan re-derives the archive's schema and rebuilds every sealed segment's columnar block against it (see Collection.ReschemaScan).

func (*Archive) Retention added in v0.19.0

func (a *Archive) Retention() Retention

Retention returns the archive's current retention bounds.

func (*Archive) RetrainDict added in v0.18.0

func (a *Archive) RetrainDict(sampleMax int) (int, error)

RetrainDict trains a fresh ZSTD dictionary from up to sampleMax records and recompresses every segment in place under it (an append-only reseal that preserves order), returning the new dictionary's size in bytes. This is how an archive's compression adapts to the data it has accumulated.

func (*Archive) Rewrite added in v0.18.0

func (a *Archive) Rewrite() int

Rewrite recompresses and re-encodes every segment in place under the current codec and hot set (e.g. after a hot-set change), preserving order, and returns the number of records rewritten.

func (*Archive) Rotate

func (a *Archive) Rotate(now float64) (int, error)

Rotate drops whole oldest segments until the archive is back within its Retention bounds, returning the number dropped. now is the caller's wall clock (unix seconds) for age-based retention.

func (*Archive) RowGroupBytes added in v0.28.1

func (a *Archive) RowGroupBytes() int

RowGroupBytes reports the budget currently in effect (0 meaning the default).

func (*Archive) SaveDemand added in v0.25.0

func (a *Archive) SaveDemand()

SaveDemand checkpoints recorded query demand to the archive's directory, ageing it by the time since the last checkpoint, so index decisions are not restarted from zero every time the process is. See Collection.SaveDemand.

func (*Archive) SchemaFit added in v0.25.3

func (a *Archive) SchemaFit(sampleMax int) ([]SchemaFieldFit, int)

SchemaFit measures the archive's derived schema against a fresh sample, reporting per-field escape rates (see Collection.SchemaFit).

func (*Archive) SchemaScanInfo added in v0.25.2

func (a *Archive) SchemaScanInfo() SchemaScanInfo

SchemaScanInfo reports the columnar accelerator's state for the archive.

func (*Archive) SetGCFloor added in v0.21.4

func (a *Archive) SetGCFloor(floor float64)

SetGCFloor installs a runtime GC watermark (in Retention.MinAgeAttr units) so Rotate may reclaim consumed records early -- above MinAge, below the configured ceilings (see Collection.SetGCFloor). Not persisted; re-assert it each maintenance pass. floor <= 0 clears it.

func (*Archive) SetRetention added in v0.19.0

func (a *Archive) SetRetention(r Retention)

SetRetention updates the retention bounds at runtime; the next Rotate enforces them.

func (*Archive) SetRowGroupBytes added in v0.28.1

func (a *Archive) SetRowGroupBytes(n int)

SetRowGroupBytes changes the archive's columnar row-group budget at runtime (see Collection.SetRowGroupBytes). It governs segments sealed from now on; nothing already written is rewritten or needs to be.

func (*Archive) SidecarSizes added in v0.8.0

func (a *Archive) SidecarSizes() SidecarSizes

SidecarSizes reports the archive's sealed-segment sidecar index bytes (mmap-backed, evictable page cache), broken out by structure. An operator diagnostic.

func (*Archive) StaleIndexSegments added in v0.21.2

func (a *Archive) StaleIndexSegments() (stale, sealed int)

StaleIndexSegments reports how many of the archive's sealed segments still carry an index built under an older configuration, and how many are sealed in total. Normally zero: AddIndex/DropIndex rebuild as they go. A non-zero count means a rebuild was interrupted or failed (it is best-effort per segment); a Reindex retries those segments.

func (*Archive) StaleIndexSegmentsByPolicy added in v0.25.0

func (a *Archive) StaleIndexSegmentsByPolicy() int

StaleIndexSegmentsByPolicy reports segments left on an older index configuration because they fall outside the backfill horizon: expected, and not actionable.

func (*Archive) Stats added in v0.19.0

func (a *Archive) Stats() Stats

Stats reports storage accounting: record count, segment count, and arena/used/dead bytes (dead is ~0 for an append log, which never supersedes). The same struct the mutable store reports, so an archive's storage is visible on the same terms.

func (*Archive) StopMaintenance added in v0.30.4

func (a *Archive) StopMaintenance()

StopMaintenance asks the archive's underlying collection to stop maintenance. See Collection.StopMaintenance.

func (*Archive) TopKOrderThreshold added in v0.29.4

func (a *Archive) TopKOrderThreshold(q *vm.Query, orderAttr string, desc bool, k int, stats *ScanStats) (float64, int, bool)

TopKOrderThreshold computes the ORDER BY orderAttr {DESC|ASC} LIMIT k cutoff over the archive's matching records via the columnar scan (see Collection.TopKOrderThreshold). stats (may be nil) records the cutoff scan's work for EXPLAIN ANALYZE.

func (*Archive) Truncate added in v0.20.1

func (a *Archive) Truncate()

Truncate drops every record, resetting the archive to empty in place: all segments are unmapped and their data + sidecar-index files unlinked (via the backing Collection's Truncate). Segment counters keep advancing, so a fresh Append starts a new segment. This is the destructive reset a from-scratch history re-sync uses -- empty the store, then re-ingest from the source. Retention/index/zone-map configuration is preserved. Safe against concurrent queries (a scan holding a pin reads its old data until it finishes); callers must serialize Truncate against the single Append writer.

func (*Archive) UpgradeCodecPass added in v0.25.0

func (a *Archive) UpgradeCodecPass(opts UpgradeOptions) int

UpgradeCodecPass re-encodes sealed segments that are still on an older dictionary, where measurement says the newer one compresses them meaningfully better. Returns the number of segments upgraded.

Segments measured as no better (or worse) are recorded and not measured again until the dictionary changes, so a pass over a fully-evaluated archive costs a directory read.

func (*Archive) Watch added in v0.4.0

func (a *Archive) Watch(ctx context.Context, cursor []byte) (iter.Seq[WatchEvent], error)

Watch streams the archive as change data: a full replay of retained records (oldest first) then live appends, resumable from an opaque cursor. A cursor older than what rotation still retains yields a WatchReset and resumes from the current floor. See Collection.Watch and docs/WATCH.md.

func (*Archive) WatchCursor added in v0.29.12

func (a *Archive) WatchCursor() ([]byte, error)

WatchCursor returns an opaque cursor at the current head of the archive's change log, so a following Watch streams only subsequent appends rather than replaying what is retained. See Collection.WatchCursor.

func (*Archive) ZoneAttrs added in v0.21.2

func (a *Archive) ZoneAttrs() []string

ZoneAttrs returns the attributes carrying per-segment [min,max] zone maps -- the ones a range query prunes whole segments on, not just postings.

type ArchiveOptions

type ArchiveOptions struct {
	// Dir is the directory holding the archive's files. Required.
	Dir string
	// SegmentSize is the segment (mmap file) size in bytes; a segment rolls over when
	// the next ad will not fit. Default 8 MiB.
	SegmentSize int
	// Codec compresses stored ad bytes. Default identity. For recovery the codec must
	// match what a segment was written under (recorded per segment in its file).
	Codec Codec
	// HotAttrs front-loads these attributes in each ad's hot header (see Collection).
	HotAttrs []string
	// CategoricalAttrs / ValueAttrs configure the per-segment indexes (see Collection).
	CategoricalAttrs []string
	ValueAttrs       []string
	// ZoneAttrs names numeric attributes to keep per-segment min/max on, so a query with
	// a range/equality constraint on one can skip whole segments. ValueAttrs are
	// automatically included.
	ZoneAttrs []string
	// Retention bounds what rotation keeps. Zero ⇒ keep everything.
	Retention Retention
	// IndexBackfillBytes bounds how far back an index configuration change is carried (see
	// Options.IndexBackfillBytes). 0 carries it across the whole archive.
	IndexBackfillBytes int64

	// RowGroupBytes is the uncompressed record-bytes budget for one columnar row group (see
	// Options.RowGroupBytes). 0 takes the default. It decides how much of a segment has to be
	// decompressed to read a single record, against how far compression can see across records, and
	// an archive is where that trade is worth tuning: its segments are immutable, so the choice is
	// made once per segment and lived with.
	RowGroupBytes int

	// GroupSchemaCount and its companions enable secondary columnar schemas (see Options).
	GroupSchemaCount    int
	GroupStabilityRuns  int
	GroupMergeJaccard   float64
	GroupMaxPartialFrac float64

	// InternAtSeal interns each segment as soon as it seals (see Options.InternAtSeal), so an
	// archive gets the interning density/decode win immediately instead of only after a later
	// RetrainDict/Rewrite. Optional; default off. Trades a per-seal transcode (off the write
	// lock) for the win landing eagerly.
	InternAtSeal bool

	// BlockCacheBytes sets the PROCESS-GLOBAL shared archive block-cache budget (see
	// Options.ArchiveBlockCacheBytes). It is shared by ALL archives in the process, not scoped to
	// this one. 0 keeps the current/default archive budget. Setting it here is equivalent to
	// calling SetArchiveBlockCacheBudget.
	BlockCacheBytes int64
}

ArchiveOptions configures an Archive.

type AutoTuneChange

type AutoTuneChange struct {
	Attr, Kind, Action, Reason string // Action: "add" | "drop"
}

AutoTuneChange records one add/drop AutoTune applied.

type AutoTuneOptions

type AutoTuneOptions struct {
	// SampleMax caps the ads profiled. Default maxDistinctSample.
	SampleMax int
	// MinDemand is the minimum equality+range probe count for an attribute to be
	// added as an index. Default 1 (any observed demand).
	MinDemand int64
	// DropUnused, if set, drops AUTO-created indexes that no query has filtered on
	// ("unused"). Off by default so AutoTune never removes an index a workload has
	// simply not exercised yet. Human-created indexes and low-cardinality indexes are
	// never auto-dropped (SuggestDrops reports them for manual review).
	//
	// Two guards keep this from thrashing, because rebuilding a dropped index is far more
	// expensive than carrying an idle one:
	//
	//   - It requires EVIDENCE. Demand counters are process-local and never decay, so
	//     right after a restart every attribute reads zero; dropping on that would discard
	//     the whole auto set on every restart and rebuild it as traffic returns. Nothing is
	//     dropped until the tracker has watched for DropUnusedMinWindow and seen at least
	//     DropUnusedMinQueries probes.
	//   - When a memory budget is configured it defers to it entirely. An unused index
	//     that fits costs little, and budgetTrim already evicts least-demanded first when
	//     space is actually short -- with hysteresis, which dropping on "unused" alone has
	//     none of. Reclaim when you need the room, not merely because you could.
	DropUnused bool
	// DropUnusedMinWindow and DropUnusedMinQueries are how much observation makes a zero
	// meaningful. Defaults: defaultDropWindow and defaultDropMinQueries. Negative values
	// waive the requirement (as with BudgetSlackBytes), which is for tests and for a caller
	// who has its own evidence that an index is unwanted.
	DropUnusedMinWindow  time.Duration
	DropUnusedMinQueries int64
	// BudgetHighFrac / BudgetLowFrac, when BudgetHighFrac > 0, bound index memory as a
	// fraction of the live data bytes: AutoTune stops adding demand-driven indexes once
	// index bytes reach the high mark, and trims the least-used AUTO indexes (never
	// human-created ones) until index bytes fall below BudgetLowFrac. This is a high/low-
	// watermark with hysteresis: grow until over the high mark, trim back to the low
	// mark. 0 disables the budget (unbounded). BudgetLowFrac defaults to 0.7*BudgetHighFrac.
	BudgetHighFrac float64
	BudgetLowFrac  float64
	// BudgetSlackBytes is absolute leeway on top of the high watermark: index bytes may
	// exceed BudgetHighFrac by up to this many bytes before any trim triggers, so a small
	// overage (or a small database, where the percentage is tiny in absolute terms) does
	// not cause churn. 0 uses defaultBudgetSlackBytes (10 MiB); set a negative value for
	// no slack (trim exactly at the fraction).
	BudgetSlackBytes int64
	// Reindex, if set, calls Reindex after applying changes so the new/removed
	// indexes take effect immediately instead of at the caller's next Reindex.
	Reindex bool
}

AutoTuneOptions configures AutoTune.

type AutoTuneResult

type AutoTuneResult struct {
	Changes []AutoTuneChange
	Changed bool
}

AutoTuneResult reports what AutoTune changed.

type Codec

type Codec interface {
	// Compress appends the compressed form of src to dst and returns it.
	Compress(dst, src []byte) []byte
	// Decompress appends the decompressed form of src to dst and returns it.
	Decompress(dst, src []byte) ([]byte, error)
	// Name identifies the codec (for diagnostics).
	Name() string
}

Codec compresses and decompresses encoded ad bytes. It is a seam so the store can run with no compression (identityCodec, the default) in tests and benchmarks, and with ZSTD — optionally with a pre-trained shared dictionary — in production.

A codec's dictionary is fixed for the life of a collection: stored records are opaque compressed bytes that compaction copies verbatim, so every record must be decodable by the same codec. Re-training the dictionary between compactions (which requires versioning dictionaries per record and recompressing) is a future extension; a dictionary trained once from a representative sample (see TrainDict) already captures the cross-ad redundancy that dominates a pool of similar ClassAds.

func NewZSTDCodec

func NewZSTDCodec(dict []byte) (Codec, error)

NewZSTDCodec returns a ZSTD codec. If dict is non-empty it is used as a shared compression dictionary (see TrainDict). Pass nil for dictionary-less ZSTD.

type CodecStats added in v0.6.0

type CodecStats struct {
	// Codec is the current codec name: "identity" (no compression), "zstd", or "zstd+dict".
	Codec string `json:"codec"`
	// DictBytes is the trained ZSTD dictionary size (0 if none). Set by RetrainDict.
	DictBytes int64 `json:"dictBytes"`
	// LastRetrain is when RetrainDict last succeeded in this process (zero if never; note
	// a persistent collection recovers its codec but not this in-process timestamp).
	LastRetrain time.Time `json:"lastRetrain,omitempty"`
	// SampleRecords is how many live records were sampled for the ratio below.
	SampleRecords int `json:"sampleRecords"`
	// CompressedBytes / UncompressedBytes are the sampled records' stored (compressed) and
	// decompressed sizes; Ratio is UncompressedBytes/CompressedBytes (1.0 = no gain).
	CompressedBytes   int64   `json:"compressedBytes"`
	UncompressedBytes int64   `json:"uncompressedBytes"`
	Ratio             float64 `json:"ratio"`
}

CodecStats reports the storage codec's state and effectiveness, for the diagnostic question "is compression on, when was the dictionary last (re)trained, and how well is it compressing?". A collection defaults to the identity codec (no compression) until Options.Codec or RetrainDict switches it, so a Ratio near 1.0 with Codec "identity" means compression was never enabled.

type Collection

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

Collection is an in-memory, sharded, memory-dense store of ClassAds. It is safe for concurrent use.

func New

func New(opts Options) *Collection

New creates an empty Collection.

func Open

func Open(opts Options) (*Collection, error)

Open opens a persistent collection under opts.Dir, whose arenas are memory-mapped files. Committed data is flushed to disk on Close (per-commit msync durability is added in a later milestone). If opts.Dir is empty, Open is equivalent to New (an in-memory collection). Persistence is unix-only.

NOTE (P2): this creates a fresh persistent collection; recovering an existing directory (rebuilding the directory + index from the segment files) is the next milestone.

func (*Collection) AddAutoIndex added in v0.6.0

func (c *Collection) AddAutoIndex(categorical, value []string) bool

AddAutoIndex adds indexes marked as auto-created (provenance auto), like AddIndex but eligible for automatic trimming by AutoTune's memory budget. Used by the persistence layer to restore auto provenance on restart.

func (*Collection) AddHotAttrs added in v0.6.0

func (c *Collection) AddHotAttrs(names ...string) []string

AddHotAttrs pins the named attributes into the hot set (front-loaded in future writes' hot headers), merging them with the current set, and returns the resulting hot attribute names. Unlike RefreshHotSet, which recomputes the set from sampled frequency, this forces specific attributes in regardless of how often they appear. Works in both RAM (interned) and persistent (inline) modes.

func (*Collection) AddIndex

func (c *Collection) AddIndex(categorical, value []string) bool

AddIndex adds categorical (string equality / membership) and value (numeric equality + range) indexes at runtime and returns whether the configuration changed. Names already indexed are ignored; a name given as both categorical and value, or already indexed as the other kind, is placed as categorical (equality on strings is the common case and categorical also serves it). The new indexes take effect for existing segments on the next Reindex; new segments pick them up when they are first indexed.

func (*Collection) AutoIndexNames added in v0.6.0

func (c *Collection) AutoIndexNames() []string

AutoIndexNames returns the names of auto-created indexes (provenance auto), so the human/auto distinction can be persisted and survive a restart.

func (*Collection) AutoTune

func (c *Collection) AutoTune(opts AutoTuneOptions) AutoTuneResult

AutoTune applies the advisory signals: it adds the demand-driven indexes from SuggestIndexes (those at or above MinDemand) and, if DropUnused is set, drops the unused ones from SuggestDrops. It is the "make it so" companion to the SuggestX methods — call it on whatever schedule you like (e.g. alongside Reindex). All suggestions are computed from one consistent snapshot before any change is applied.

func (*Collection) Begin added in v0.6.0

func (c *Collection) Begin() *Txn

Begin starts an optimistic transaction. Its snapshot for a shard is captured the first time the transaction reads or writes a key in that shard.

func (*Collection) BeginRedacted added in v0.28.0

func (c *Collection) BeginRedacted() *Txn

BeginRedacted is Begin for a caller NOT entitled to sealed values: reads through the transaction decode with no key, so a sealed attribute comes back undefined. Writes are unaffected -- a redacted transaction can still store an ad, and encodeAd seals what the collection's policy says to seal.

The entitlement lives on the transaction rather than on each read, so a new read method inherits it instead of having to remember it.

func (*Collection) BuildAndEnableSchemaScan added in v0.22.0

func (c *Collection) BuildAndEnableSchemaScan(sampleMax, hotTopN int) bool

BuildAndEnableSchemaScan samples the collection, builds an adschema, chooses the hot numeric tier as the top-hotTopN int/real fields by accumulated read demand (c.demand -- the same signal RefreshHotSet uses; the schema's hot tier IS the hot set), and enables the columnar scan over the sealed segments. Returns false if there is nothing to sample. Re-callable to pick up newly-sealed segments (existing blocks are kept).

func (*Collection) Chained added in v0.21.3

func (c *Collection) Chained() bool

Chained reports whether the collection has structural (parent-only) ads that Scan/Query hide (HTCondor cluster/proc chaining). When false, Len is exactly the match-all row count, which the COUNT(*) fast path relies on.

func (*Collection) Close

func (c *Collection) Close() error

Close flushes all committed data to disk and unmaps the collection's segment files. It is a no-op for an in-memory collection. The collection must not be used after Close.

func (*Collection) CodecStats added in v0.6.0

func (c *Collection) CodecStats(sampleMax int) CodecStats

CodecStats measures the storage codec's effectiveness over a sample of up to sampleMax live records: it decompresses each and compares the stored (compressed) size to the decompressed size. It takes each shard's read lock briefly, so it is safe alongside readers and writers.

func (*Collection) CollectRegionSamples added in v0.26.0

func (c *Collection) CollectRegionSamples(maxBytes, lookBackBytes int) [][]byte

CollectRegionSamples returns dictionary-training samples drawn from sealed segments' columnar block regions, in equal byte shares across the three kinds, up to about maxBytes in total.

For an APPEND-ONLY collection the draw is biased to recent segments: future writes resemble recent history far more than the oldest retained records, and an archive's oldest segments may predate schema or workload changes entirely. lookBackBytes caps how far back the draw reaches, measured in segment bytes from the newest; 0 means no limit.

Returns nil when no sealed segment carries a block -- there is nothing to sample from, and the caller should fall back to record samples (see CollectSamples).

func (*Collection) CollectSamples

func (c *Collection) CollectSamples(max int) [][]byte

CollectSamples returns the decompressed wire bytes of up to max ads, for training a ZSTD dictionary (see TrainDict). It samples across shards, decoding each record with the codec it was stored under.

func (*Collection) CollectSamplesRecent added in v0.26.0

func (c *Collection) CollectSamplesRecent(maxRecords, maxBytes, lookBackBytes int) [][]byte

CollectSamplesRecent returns record samples drawn at random from the newest lookBackBytes of segment data, honoring MVCC visibility, bounded by BOTH a record count and a byte budget -- whichever binds first.

Both bounds exist because they answer different questions and callers supply different ones. maxRecords is what every existing caller has (RetrainDict's sampleMax), and keeping it a count means switching to this sampler cannot change what an existing call collects by reinterpreting its units. maxBytes is what the work actually scales with, since a record's size varies by two orders of magnitude across tables.

maxRecords <= 0 means no count bound; maxBytes <= 0 uses defaultSampleBytes; lookBackBytes <= 0 uses defaultLookBackBytes. Samples are flattened to inline wire exactly as CollectSamples does, so the two are interchangeable as TrainDict input.

Returns nil when the collection holds no visible records.

func (*Collection) CollectSamplesRecentN added in v0.26.0

func (c *Collection) CollectSamplesRecentN(maxRecords int) [][]byte

CollectSamplesRecentN returns up to maxRecords records drawn at random from the recent window, with no meaningful byte bound -- the sampler for decisions that need a representative COUNT of records rather than a fixed volume of bytes. See statSampleBytes.

func (*Collection) ColumnarEvalCount added in v0.27.0

func (c *Collection) ColumnarEvalCount(q *vm.Query) (int, bool)

ColumnarEvalCount counts the records matching q by evaluating q itself against each record's columns, for any NATIVE query.

ok=false when it cannot serve the query at all: the accelerator is off, or the query is not native (a delegated subtree needs a real ClassAd scope). Individual records whose stored value is an expression are re-evaluated the ordinary way; that never fails the query.

Because it evaluates the query rather than a summary of it, it is exact by construction -- the class of bug that required ExactProbes cannot arise here.

func (*Collection) ColumnarizeSealed added in v0.28.0

func (c *Collection) ColumnarizeSealed() int

ColumnarizeSealed rewrites sealed segments into columnar-native form as a maintenance pass: each attribute the schema carries is stored once, in the segment's own columnar payload, and removed from the records. Returns how many segments were rewritten.

Idempotent and bounded. An already-columnarized segment is skipped and the active write target is never touched, so re-calling is cheap; one pass rewrites at most ColumnarSegmentBudget segments, so a large existing archive converges over several passes instead of stalling on one.

Requires a derived schema, which BuildAndEnableSchemaScan produces and also calls this from -- so a caller on the ordinary schema-review interval gets this without asking for it. Holds maintMu, so it serializes against Compact/RetrainDict/Rotate/Rewrite and against a merge.

func (*Collection) Compact

func (c *Collection) Compact() int

Compact reclaims space from superseded/deleted records. For each shard it first unlinks any fully-dead segment cheaply (no rewrite), then, if the shard's dead-byte ratio still warrants it, recompresses the live records into fresh segments. It is safe to call concurrently with reads and writes. Returns the number of shards where space was reclaimed (by either mechanism).

func (*Collection) CountConstraint added in v0.22.0

func (c *Collection) CountConstraint(constraint string) (int, bool)

CountConstraint parses a constraint string and, if it is columnar-eligible, counts via the columnar scan (see CountQuery). ok=false ⇒ use the normal count path.

func (*Collection) CountQuery added in v0.22.0

func (c *Collection) CountQuery(q *vm.Query) (int, bool)

CountQuery counts the records matching q using the columnar scan, when q is exactly (Native) one or more numeric comparisons on a single INT schema field -- the common count-where. It returns (count, true) on that fast path, or (0, false) to signal the caller to use the normal scan (schema-scan not enabled, or the predicate is not columnar-eligible). Correctness: the per-record numeric comparison matches the store's evaluation of `field OP number`; a value escaped out of the hot column is read from the cold tail, and Native() rules out any residual program the columnar path would miss.

func (*Collection) Delete

func (c *Collection) Delete(key []byte) bool

Delete removes key, returning whether it was present. The removal is an MVCC tombstone: scans already in progress still see the pre-delete version.

func (*Collection) DeltaStats added in v0.30.0

func (c *Collection) DeltaStats() (deltas, fulls int64)

DeltaStats reports how many records this collection has stored as deltas versus in full since it was opened. Both zero means delta mode is off (or nothing has been written).

Both count decisions made by the delta tracker, which is the write path -- NOT every whole record the store writes. A COLLAPSE does not appear here at all: collapseBatch writes its merged records with tx.putWireAd, which takes bytes that are already encoded and so never consults the tracker. So fulls does not move when a collapse rewrites thousands of chains, and a "did the collapse do anything" check built on it reads zero in every case -- which is exactly the wrong answer to take from a counter that looks like it should know. Use SealWalkStats for what the collapse examined and found.

func (*Collection) DropIndex

func (c *Collection) DropIndex(names ...string) bool

DropIndex removes the given attributes from the configured indexes (categorical or value, whichever they were) and returns whether the configuration changed. The attributes stop being used by queries immediately; their postings are reclaimed from existing segments at the next Reindex. Names that were never indexed are ignored.

func (*Collection) EnableSchemaScan added in v0.22.0

func (c *Collection) EnableSchemaScan(s *adSchema, hot []int)

EnableSchemaScan builds a columnar block over every currently-sealed segment (skipping the active append target) for the given schema and hot field set, publishes it on each segment, and records the state so CountQuery can auto-route matching queries. Additive and opt-in: with no state/block a query takes the normal row path, so a collection that never calls this is unaffected. Reads immutable sealed bytes only.

func (*Collection) EncryptedAttrNames added in v0.7.0

func (c *Collection) EncryptedAttrNames() []string

EncryptedAttrNames returns the explicit encrypted-attribute set (case-folded, sorted) -- the human-toggled attributes, not the always-on private ones. Used for persistence and diagnostics.

func (*Collection) EncryptionEnabled added in v0.7.0

func (c *Collection) EncryptionEnabled() bool

EncryptionEnabled reports whether this collection has a data key (encryption at rest is active). Without it, EncryptedAttrs is inert.

func (*Collection) ExplainMatch added in v0.6.0

func (c *Collection) ExplainMatch(job *classad.ClassAd, targetConstraint string) MatchExplain

ExplainMatch reports how matchmaking job against this collection would execute: it rewrites the job's Requirements over the slot (baking the job's attribute values to constants), extracts the index-satisfiable probes, and reports which are covered by a configured index and the resulting access path. targetConstraint, if non-empty, is the MATCH's resource-side filter (WHERE TARGET, including NOPREEMPT's `State =!= "Claimed"`) -- already slot-scoped -- which is melded into the shown predicate, probes, and evaluation order so the explanation reflects the full effective slot filter. No I/O beyond reading the spec.

func (*Collection) ExplainQuery added in v0.6.0

func (c *Collection) ExplainQuery(q *vm.Query) QueryExplain

ExplainQuery reports how the store would execute q: which of its conjuncts are index-usable, and the resulting access path (indexed / parallel scan / serial scan). It performs no I/O beyond reading the current index spec.

func (*Collection) Flatten added in v0.4.0

func (c *Collection) Flatten(key []byte) (*classad.ClassAd, bool)

Flatten returns a standalone, self-contained copy of the ad at key with every inherited parent attribute materialized into it (the child's own values winning, parent-private attributes excluded) and no residual parent link. It is the form to archive as a history record: like condor_history, which stores each completed job flattened rather than as a live cluster/proc pair, so the record stays readable after its structural parent is gone. For a collection without parent chaining this is identical to Get.

func (*Collection) ForEachAd added in v0.7.0

func (c *Collection) ForEachAd(fn func(key string, ad *classad.ClassAd) bool)

ForEachAd calls fn with every stored (non-system) ad and its key, including structural (parent-only) ads -- a client-facing image, unlike Scan/Keys which hide structural ads. It is the driver for db.DeleteWhere's key matching, so it must NOT expose internal system records (they are neither client-visible nor client-deletable); system records are enumerated separately via ForEachSystemAd. Each ad is fully decoded (decompressed and, for encrypted attributes, decrypted). Iteration stops early if fn returns false. Per-shard consistent snapshot, like Scan.

func (*Collection) ForEachAdAt added in v0.24.0

func (c *Collection) ForEachAdAt(snapOf func(shardIdx int) uint64, fn func(key string, ad *classad.ClassAd) bool)

ForEachAdAt is ForEachAd over the versions visible at a caller-chosen snapshot sequence per shard, instead of each shard's current commit sequence.

snapOf is called once per shard, with the shard's index, and returns the sequence to read that shard at. It is a callback rather than a map so a transaction can capture a shard's sequence lazily on first touch -- and, having captured it, stay pinned to it for every later read of that shard.

This is what lets a transaction scan the committed store at its own snapshot: the same versions its point lookups see, rather than whatever has committed since.

func (*Collection) ForEachSystemAd added in v0.9.0

func (c *Collection) ForEachSystemAd(fn func(key string, ad *classad.ClassAd) bool)

ForEachSystemAd calls fn with every internal system-keyed ad and its key -- the mirror image of ForEachAd, which hides them. It is the enumeration path a TTL reaper uses to find and expire durable bookkeeping records. Each ad is fully decoded; iteration stops early if fn returns false. Per-shard consistent snapshot, like Scan.

func (*Collection) GCFloor added in v0.21.4

func (c *Collection) GCFloor() float64

GCFloor returns the current runtime GC watermark (0 when unset).

func (*Collection) Get

func (c *Collection) Get(key []byte) (*classad.ClassAd, bool)

Get returns the current ad for key, decoded, or (nil, false).

func (*Collection) GetRedacted added in v0.28.0

func (c *Collection) GetRedacted(key []byte) (*classad.ClassAd, bool)

GetRedacted is Get for a reader NOT entitled to sealed values: the decode holds no key, so a sealed attribute comes back undefined rather than opened. See QueryRedacted.

func (*Collection) GroupCountAll added in v0.28.0

func (c *Collection) GroupCountAll(groupAttr string) ([]GroupCount, bool)

GroupCountAll is GroupCountQuery with no constraint: the histogram of the whole column.

func (*Collection) GroupCountConstraint added in v0.28.0

func (c *Collection) GroupCountConstraint(constraint, groupAttr string) ([]GroupCount, bool)

GroupCountConstraint parses a constraint string and, if the shape is columnar-eligible, answers the per-value counts of groupAttr over the matching records (see GroupCountQuery). ok=false ⇒ use the normal grouping path.

func (*Collection) GroupCountQuery added in v0.28.0

func (c *Collection) GroupCountQuery(q *vm.Query, groupAttr string) ([]GroupCount, bool)

GroupCountQuery answers `SELECT groupAttr, COUNT(*) ... WHERE q GROUP BY groupAttr` from the columnar blocks, returning the groups sorted by value. See GroupStatsQuery, of which this is the no-value- aggregate case.

func (*Collection) GroupSchemaAgreement added in v0.28.0

func (c *Collection) GroupSchemaAgreement(sampleMax, k int) GroupSchemaAgreement

GroupSchemaAgreement derives groups from each sealed segment independently and reports how often each table-level group reappears. Diagnostic only, and O(segments x sample) -- for an operator deciding whether to commit to a group, not for a maintenance pass.

func (*Collection) GroupSchemaChanges added in v0.29.13

func (c *Collection) GroupSchemaChanges() []GroupSchemaChange

GroupSchemaChanges returns the committed-group change log, oldest first. Reads the sidecar only -- no sampling -- so it is available to a READ-level client.

func (*Collection) GroupSchemaDrift added in v0.28.0

func (c *Collection) GroupSchemaDrift() GroupSchemaDrift

GroupSchemaDrift reads the checkpoint and reports what moved between the first and last retained derivations.

func (*Collection) GroupSchemaLastAgreement added in v0.29.13

func (c *Collection) GroupSchemaLastAgreement() (GroupSchemaLastAgreement, bool)

GroupSchemaLastAgreement returns the last persisted per-segment agreement, if any.

func (*Collection) GroupSchemas added in v0.28.0

func (c *Collection) GroupSchemas(sampleMax, k int) GroupSchemaInfo

GroupSchemas derives group schemas from a fresh sample and returns the report. It does not change how anything is stored or read; see deriveGroupSchemas.

k <= 0 uses defaultGroupSchemas. The report is also checkpointed to the collection's directory when it has one, so a later call can show what moved.

func (*Collection) GroupStatsAll added in v0.28.0

func (c *Collection) GroupStatsAll(groupAttr string, aggAttrs []string) ([]GroupStats, bool)

GroupStatsAll is GroupStatsQuery with no constraint: the whole column's histogram plus aggregates. The caller is asserting the query matches every record, which only it can know (see db.IsMatchAll) -- there is no predicate here to check that against.

This is not reachable through GroupStatsQuery: an unconstrained query has no probes to analyze, so the predicate analysis declines it. And the index path that answers other unconstrained grouped counts only covers CATEGORICALLY indexed attributes, so `GROUP BY <numeric>` with no WHERE fell to a record scan the same way the constrained form did.

func (*Collection) GroupStatsConstraint added in v0.28.0

func (c *Collection) GroupStatsConstraint(constraint, groupAttr string, aggAttrs []string) ([]GroupStats, bool)

GroupStatsConstraint parses a constraint string and, if the shape is columnar-eligible, answers the per-group record count plus a NumStats per aggAttr over the matching records (see GroupStatsQuery). ok=false ⇒ use the normal grouping path.

func (*Collection) GroupStatsQuery added in v0.28.0

func (c *Collection) GroupStatsQuery(q *vm.Query, groupAttr string, aggAttrs []string) ([]GroupStats, bool)

GroupStatsQuery answers `SELECT groupAttr, COUNT(*), <aggregates over aggAttrs> ... WHERE q GROUP BY groupAttr` from the columnar blocks, returning the groups sorted by value. Each group carries a NumStats per aggAttr, from which the caller renders MIN/MAX/SUM/AVG/COUNT(attr) -- the same inputs NumStatsQuery returns for the ungrouped case, so the rendering is shared and cannot diverge.

ok=false when the shape is not columnar-servable -- no accelerator, a non-native constraint, a predicate that is not a conjunction of scalar numeric comparisons, or a group/aggregate attribute the schema does not carry as a numeric field -- and the caller then scans as before. It also declines mid-flight if a matching record's GROUP value is not a scalar number, because such a record still forms a group: the row path emits a group per distinct rendered value (`undefined` for an absent attribute, the text for a string), and a column read cannot tell absent from present-but-non-numeric, let alone render one. Dropping those records instead would return a result silently missing groups, so the query goes back to the scan that can render them.

An absent AGGREGATE value is different and is simply skipped: MIN/SUM/COUNT(attr) over a set that does not include the record is what the reference computes for a record whose attribute is undefined. Only the group column decides which rows exist.

Segments whose own schema lacks a field fall back to a row walk that reads only the attributes involved, so a schema change costs those segments and not the query.

func (*Collection) HotAttrNames

func (c *Collection) HotAttrNames() []string

HotAttrNames returns the current hot attributes by name (for diagnostics).

func (*Collection) IndexSizes added in v0.6.0

func (c *Collection) IndexSizes() IndexSizes

IndexSizes measures each configured index's resident bytes across all live segments, tagged with provenance (human vs auto), against the live data bytes. It takes each shard's read lock briefly, so it is safe alongside readers and writers.

func (*Collection) IndexedAttrs

func (c *Collection) IndexedAttrs() (categorical, value []string)

IndexedAttrs returns the currently-indexed attribute names, split by kind, in the collection's canonical (interned) casing.

func (*Collection) InternSealed added in v0.23.0

func (c *Collection) InternSealed()

InternSealed transcodes every sealed, still-inline segment of a persistent append-only collection to interned form under the current codec, swapping each in and rebuilding its indexes. It is the eager complement to interning-at-compaction (resealAppendOnly): an archive that never retrains still gets the density/decode win as soon as segments seal. Idempotent -- already-interned segments and the active write target are skipped, so re-calling is cheap and it composes with a later RetrainDict. Append-only only: a mutable segment's records can be superseded in place after seal, which this off-lock transcode would race (there interning rides compaction, which reconciles supersedes in a final critical section). Holds maintMu, so it serializes against Compact/RetrainDict/Rotate/Rewrite. Called opt-in from the Archive.Append eager-seal hook (Options.InternAtSeal) and available as a manual maintenance pass.

func (*Collection) Keys added in v0.6.0

func (c *Collection) Keys() []string

Keys returns every visible key in the collection at a consistent per-shard snapshot, in no particular order. Structural (parent-only) keys of a chained collection are excluded, matching Scan's output. Each key is a fresh copy the caller owns. Useful for administrative enumeration and for a replica that must clear its keyspace before a full re-sync.

func (*Collection) Len

func (c *Collection) Len() int

Len returns the number of live keys across all shards. Note this includes structural (parent-only) ads of a chained collection, which Scan/Query hide -- so Len equals the number of rows a match-all query returns only when the collection is not Chained.

func (*Collection) MarkAutoIndexes added in v0.25.0

func (c *Collection) MarkAutoIndexes(names []string)

MarkAutoIndexes records the named EXISTING indexes as auto-created. It changes only provenance: no index is added, dropped, or rebuilt, and a name that is not currently indexed is ignored.

This is the restore side of AutoIndexNames for a store that rebuilds its whole index set at open (an archive does). Going through AddAutoIndex there would not work: the indexes already exist as human-created by then, and an auto add deliberately refuses to downgrade a human index. Adding them separately instead would restore provenance correctly but rebuild them, which for an archive means reading history back through the decompressor on every start.

func (*Collection) Match added in v0.5.2

func (c *Collection) Match(job *classad.ClassAd) iter.Seq[*classad.ClassAd]

Match returns every ad in the collection that symmetrically matches job, in no particular order. Use MatchSorted for a Rank-ordered result. job is not modified.

func (*Collection) MatchSorted added in v0.5.2

func (c *Collection) MatchSorted(job *classad.ClassAd, limit int) []*classad.ClassAd

MatchSorted returns the matching ads ranked by job's Rank expression, best (highest Rank) first. limit <= 0 returns all matches; limit > 0 returns at most the top limit. Ads whose Rank does not evaluate to a number sort after ranked ones. Ties in Rank are broken in an unspecified order (it depends on scan/fan-out order); a caller needing a deterministic tiebreak should apply its own. job is not modified.

func (*Collection) MatchSortedRanked added in v0.6.0

func (c *Collection) MatchSortedRanked(job *classad.ClassAd, limit int) []RankedMatch

MatchSortedRanked is MatchSorted that also returns each match's Rank, best (highest Rank) first, up to limit (<= 0 = all). Ads whose Rank is not numeric sort after ranked ones. Like MatchSorted, the job's Requirements is used to visit only index-candidate slots when an index covers it (the matchmaking pushdown), rather than bilaterally evaluating every slot.

func (*Collection) MatchSortedRankedFiltered added in v0.6.0

func (c *Collection) MatchSortedRankedFiltered(job *classad.ClassAd, targetConstraint string, limit int) ([]RankedMatch, error)

MatchSortedRankedFiltered is MatchSortedRanked restricted to slots that also satisfy targetConstraint (a constraint over the resource ad, e.g. `State =!= "Claimed"`). The constraint's index-satisfiable probes narrow the candidate scan (pushdown); the full constraint is then re-checked on each matched slot. An empty constraint is exactly MatchSortedRanked. Returns an error only if targetConstraint fails to parse.

func (*Collection) MigrateSealedAttrs added in v0.28.0

func (c *Collection) MigrateSealedAttrs(workers int) int

MigrateSealedAttrs rewrites every sealed segment that still holds a value this collection would seal, so a store written before sealing was turned on stops carrying private attributes in the clear. Returns the number of segments rewritten.

Safe to call on a store that needs nothing: the scan finds no candidate and it returns 0. Safe to call twice, and safe to interrupt (see the file comment). A collection with no sealer, or no persistence, has nothing to migrate.

workers bounds the concurrent rewrites; <= 0 takes a default derived from GOMAXPROCS. Rewriting is decompress + re-encode + compress per record, so it is CPU-bound and worth parallelizing on a store with many segments -- but each rewrite also holds a whole new segment in memory, which is why it is bounded rather than one goroutine per segment.

func (*Collection) NumStatsQuery added in v0.25.2

func (c *Collection) NumStatsQuery(q *vm.Query, attr string) (NumStats, bool)

NumStatsQuery computes the numeric aggregate inputs for attr over the records matching q, using the columnar scan. A nil q means every record.

It returns ok=false when the columnar path cannot serve the request -- the accelerator is off, attr is not a numeric field of the current schema, or q is not a conjunction of scalar numeric comparisons ON attr ITSELF. That last restriction is the same one CountQuery has: a predicate over a different field would need a second column read per record, which this pass does not do. A caller that gets false scans instead.

func (*Collection) OpStats added in v0.11.0

func (c *Collection) OpStats() OpStats

OpStats returns a snapshot of the collection's operational timing counters -- where its callers spent time blocked in, or holding, each of the store's stall points (shard write lock, segment allocation, durability sync, and the maintenance passes). Per-shard counters are summed across all shards; every value is a monotonic cumulative total, so a scraper derives rate and mean latency from deltas. It reads atomics locklessly, so it never contends with readers or writers.

func (*Collection) Ordered added in v0.5.2

func (c *Collection) Ordered(index int, partition classad.Value, resume OrderCursor) iter.Seq[OrderedAd]

Ordered iterates one partition of the index-th configured ordered index in sort order, yielding each member ad with its resume cursor and cluster signature. Iteration is over an O(1) snapshot taken at the call, so it is stable even as the index churns; each ad is fetched live by key, so a member deleted since the snapshot is skipped. For an index configured without a Partition, the partition argument is ignored (there is a single global run). resume's zero value starts at the beginning.

func (*Collection) OrderedRedacted added in v0.28.0

func (c *Collection) OrderedRedacted(index int, partition classad.Value, resume OrderCursor) iter.Seq[OrderedAd]

OrderedRedacted is Ordered for a caller NOT entitled to sealed values: each ad is materialized with no key, so a sealed attribute arrives undefined. See Collection.QueryRedacted.

This is served over the RPC to unprivileged sessions (dbrpc's ordered op), so it needs the entitlement like any other client-facing read -- it is not only the negotiator's internal list.

func (*Collection) PresenceCountQuery added in v0.27.0

func (c *Collection) PresenceCountQuery(q *vm.Query) (int, bool)

PresenceCountQuery counts the records matching a single `attr is undefined` / `attr isnt undefined` predicate over the columnar accelerator.

Returns ok=false when the columnar path cannot serve it -- the accelerator is off, the predicate is not a lone presence probe on a schema field, or some record's value for the attribute is an EXPRESSION, which needs evaluating. A caller that gets false scans instead.

func (*Collection) Put

func (c *Collection) Put(key []byte, ad *classad.ClassAd) error

func (*Collection) Query

func (c *Collection) Query(q *vm.Query) iter.Seq[*classad.ClassAd]

Query returns an iterator over the ads matching q, with the same scan-exactly-once guarantee as Scan.

Each ad is match-tested cheaply and only matching ads are fully decoded to be yielded. The match test uses, in order of preference: wire-native evaluation (the query reads scalar-literal attributes directly from the encoded ad, building no ClassAd); partial decode (decode only the attributes the query reads, transitively) when an attribute is a non-literal expression; or full decode for queries that read attributes by a runtime-computed name (eval()).

func (*Collection) QueryAsOf added in v0.12.0

func (c *Collection) QueryAsOf(q *vm.Query, t time.Time) (iter.Seq[*classad.ClassAd], error)

QueryAsOf runs q against the point-in-time snapshot at time t: it resolves t to a per-shard commit sequence and scans each shard at that sequence, so the result is exactly the ads that were current at t. It errors if time travel is disabled or t predates the retained window. v1 uses a serial full scan per shard (the index, parallel, and chained fast paths are current-time only); the query constraint is still applied to every candidate.

func (*Collection) QueryProject added in v0.6.0

func (c *Collection) QueryProject(q *vm.Query, attrs []string) iter.Seq[[]classad.Value]

QueryProject is Query specialized for aggregation and projection: for each ad matching q it yields just the named attributes' values, read wire-native (straight from the encoded ad, no *classad.ClassAd built) whenever they are scalar literals -- the common case for the numeric/string attributes an aggregate groups or sums by. Only an ad that stores one of the requested attributes as a non-literal expression falls back to a full decode, and only that ad.

This avoids the dominant cost of a GROUP BY / aggregate over Query, which fully decodes every matching ad (tens of allocations each) just to read one or two attributes.

The yielded slice is reused across iterations and aligned with attrs; the consumer must copy any value it needs to retain past the next step (the aggregate reads each value into its group state immediately).

func (*Collection) QueryProjectStats added in v0.29.6

func (c *Collection) QueryProjectStats(q *vm.Query, attrs []string, stats *ScanStats) iter.Seq[[]classad.Value]

QueryProjectStats is QueryProject that also fills stats (may be nil) with the scan's work breakdown -- how many segments were pruned by zone maps vs scanned, and how many records were answered from columns vs reassembled from the arena. It exists so an EXPLAIN ANALYZE of an aggregate that falls to this projected scan (an attribute the columnar accelerator does not cover numerically) can show WHERE the time went -- specifically a large RecordsReassembled, the slow per-record path. Stats accumulate as the returned sequence is consumed.

func (*Collection) QueryRaw

func (c *Collection) QueryRaw(q *vm.Query) iter.Seq[RawAd]

QueryRaw is Query, but yields each matching ad as RawAd (AST-free) instead of a *classad.ClassAd. It uses the same index/scan machinery -- an indexed lookup still visits only candidate ads -- so it does not regress selective queries. Inline-name (persistent) collections have no intern table, so QueryRaw yields nothing for them; callers that must support those should use Query.

func (*Collection) QueryRawFromSeq added in v0.29.5

func (c *Collection) QueryRawFromSeq(q *vm.Query, projection []string, after SeqCursor, limit int) (iter.Seq[RawAd], *SeqPage)

QueryRawFromSeq yields the records matching q that sort after the cursor, in (shard, seq, key) order, at most limit of them, projected to the named attributes. A nil query matches everything; an empty projection keeps the whole ad. It is the paginated form of QueryRaw: successive calls, each handed the previous page's Next cursor, visit every record exactly once.

The zero cursor starts at the beginning. Each shard is frozen at the sequence it had when its first page was taken, and is finished before the next shard is opened — so a record written during pagination is invisible if its shard had already started, and visible if it had not. That is the guarantee Scan already makes, which is likewise per shard.

limit <= 0 means no limit, in which case More is always false.

func (*Collection) QueryRawProjected added in v0.16.3

func (c *Collection) QueryRawProjected(q *vm.Query, projection []string, chaseRefs, redact bool) iter.Seq[RawAd]

QueryRawProjected is QueryRaw with the same in-walk projection as ScanRawProjected (see there for chaseRefs/redact).

func (*Collection) QueryRawProjectedStats added in v0.29.0

func (c *Collection) QueryRawProjectedStats(q *vm.Query, projection []string, chaseRefs, redact bool, stats *ScanStats) iter.Seq[RawAd]

QueryRawProjectedStats is QueryRawProjected that also fills stats (may be nil) with the per-scan work breakdown, for EXPLAIN ANALYZE. Iterate the returned sequence to completion (or stop early) before reading stats; a partial iteration leaves partial counts.

func (*Collection) QueryRawRedacted added in v0.16.3

func (c *Collection) QueryRawRedacted(q *vm.Query) iter.Seq[RawAd]

QueryRawRedacted is QueryRaw with private (secret) attributes stripped from every yielded ad, using the intern table's precomputed per-id private flags (see wire.NewInternTableWithPrivacy): the redaction costs one bool check per attribute and private nodes are never rendered at all. Use it to serve a public query; QueryRaw preserves private attributes for trusted consumers (e.g. the negotiator's private-ad query).

func (*Collection) QueryRawWire added in v0.16.4

func (c *Collection) QueryRawWire(q *vm.Query, projection []string, redact bool) iter.Seq[[]byte]

QueryRawWire is ScanRawWire restricted to ads matching q.

func (*Collection) QueryRedacted added in v0.28.0

func (c *Collection) QueryRedacted(q *vm.Query) iter.Seq[*classad.ClassAd]

QueryRedacted is Query for a reader NOT entitled to sealed values: the decode holds no key, so a sealed attribute comes back undefined rather than opened (see wire.DecodeResolveEncRedact).

This is the full-decode counterpart of QueryRawRedacted. Query decodes with the collection's key and leaves it to the SERIALIZER to drop private attributes -- so the secret was decrypted in this process for every unentitled read and then filtered on the way out. A caller that is not entitled to a value should not be able to obtain it, rather than obtain it and be trusted to drop it.

func (*Collection) RecordDemand added in v0.6.0

func (c *Collection) RecordDemand(probes []vm.Probe)

RecordDemand notes the attributes a constraint filters on (for SuggestIndexes) without running a scan. Query records demand automatically; this is for callers that filter outside the normal scan path -- e.g. cross-table MATCH applying a resource-side (WHERE TARGET) constraint to already-matched candidates -- so those attributes still surface as index suggestions.

func (*Collection) RefreshHotSet

func (c *Collection) RefreshHotSet(sampleMax, topN int) int

RefreshHotSet samples up to sampleMax live ads, tallies how often each attribute appears, and installs the topN most common attributes as the hot set used to front-load future writes' hot headers. It counts attribute ids directly from the wire form (no full decode). Existing ads keep the hot header they were written with; because daemons rewrite ads periodically, the population's hot headers converge on the refreshed set over time.

Returns the number of attributes chosen. A no-op (returns 0) when there are no ads yet.

func (*Collection) Reindex

func (c *Collection) Reindex()

Reindex (re)builds the per-segment value/categorical indexes for every live segment, covering all records written so far. It reads only immutable segment bytes, so it runs off the write path and does not block writers or compaction. Call it on whatever schedule you like: queries use whatever coverage exists and full-scan the rest, so Reindex only affects query speed, never results.

Reindex also reconciles segments with the current index configuration: a segment indexed before an AddIndex is rebuilt so the new attribute is backfilled, and one indexed before a DropIndex is rebuilt (or, if nothing is indexed anymore, its index is dropped) so the removed attribute's postings are reclaimed. A whole span of segments therefore evolves toward the current spec at whatever cadence the caller reindexes — no write-path or compaction coupling.

func (*Collection) ReschemaScan added in v0.25.2

func (c *Collection) ReschemaScan(sampleMax, hotTopN int) bool

ReschemaScan derives a NEW schema from a fresh sample and rebuilds every sealed segment's columnar block against it, replacing the schema that BuildAndEnableSchemaScan pinned at first enable. This is the deliberate, heavy operation that routine maintenance refuses to do: it re-encodes and re-persists a block per sealed segment, so its cost scales with the whole table, not with what changed.

Use it when SchemaFit shows the schema no longer matching the data -- escapes climbing on a queried column, or a field that has become common since. Returns false if the accelerator cannot run here (encryption at rest) or there was nothing to sample, leaving the existing schema in place.

Existing blocks are dropped first: EnableSchemaScan builds only where a segment has none, so without that it would keep every old block and the new schema would match none of them. A COLUMNARIZED segment is exempt -- its block holds attributes its records no longer carry, so it keeps both the block and the schema it was written with; adopting a new schema there means rewriting the segment. Queries during the rebuild take the row path, which is correct but slower.

func (*Collection) Retention added in v0.19.0

func (c *Collection) Retention() Retention

Retention returns the current retention bounds.

func (*Collection) RetrainDict

func (c *Collection) RetrainDict(sampleMax int) (int, error)

func (*Collection) Rewrite added in v0.6.0

func (c *Collection) Rewrite() int

Rewrite re-encodes every live ad with the current hot set (and match closure) so a changed hot set takes effect on existing ads, not just future writes, then force-compacts every shard to reclaim the superseded pre-rewrite records. Returns the number of ads rewritten.

It re-Puts ads on the normal write path, so it is a maintenance operation: an update to a key that races the rewrite may be overwritten by the pre-rewrite value. Run it during low write activity (or, in an HA deployment, on the sole writer).

func (*Collection) Rotate added in v0.17.0

func (c *Collection) Rotate(now float64) (int, error)

Rotate drops whole oldest segments until the append-only collection is back within its Retention bounds, returning the number of segments dropped. now is the caller's wall clock (unix seconds) for age-based retention -- the store keeps no clock of its own -- and is ignored unless Retention.MaxAgeAttr/MaxAge are set. Rotate is a no-op (returns 0) on a non-append-only collection or when Retention is the zero value.

The segment still being appended to is never dropped, so Rotate can bound a growing archive without racing the writer. A dropped segment's file is unmapped and unlinked once no in-flight scan still pins it (a scan that started earlier keeps reading a consistent snapshot until it finishes; the reclaim is deferred to its last unpin).

Retention honors MaxSegments, MaxBytes, and MaxAgeAttr/MaxAge (drop a segment whose newest value of the age attribute is older than now-MaxAge); MaxAge requires the age attribute to be a configured ZoneAttr so its per-segment max is available.

func (*Collection) RowGroupBytes added in v0.28.1

func (c *Collection) RowGroupBytes() int

RowGroupBytes reports the budget currently in effect (0 meaning the default).

func (*Collection) SaveDemand added in v0.25.0

func (c *Collection) SaveDemand()

SaveDemand checkpoints query demand to the collection's directory, decaying it first by the time since the last checkpoint. Call it from maintenance; it is a no-op on an in-memory collection.

Best-effort by design, as for the other sidecar metadata: a failed write costs the interval's demand, not correctness. Nothing reads this file except a later restore, and a restore that finds it missing or unparseable simply starts from what it has.

func (*Collection) Scan

func (c *Collection) Scan() iter.Seq[*classad.ClassAd]

Scan returns an iterator over every ad in the collection. It is scan-exactly- once: each key present at the moment a shard's scan begins is yielded exactly once (never duplicated, never skipped), even while concurrent updates and compaction run. An ad updated mid-scan is seen at whichever version was current when that shard's scan started.

func (*Collection) ScanRaw

func (c *Collection) ScanRaw() iter.Seq[RawAd]

ScanRaw is Scan, yielding every ad as RawAd (AST-free).

func (*Collection) ScanRawProjected added in v0.16.3

func (c *Collection) ScanRawProjected(projection []string, chaseRefs, redact bool) iter.Seq[RawAd]

ScanRawProjected is ScanRaw restricted to the projected attribute names, applied INSIDE the wire walk: the projection is resolved to interned ids once per scan, and each ad is walked as raw (id, node) pairs -- a non-projected attribute costs one id comparison and a TLV length hop, and is never name- resolved or rendered. For a typical monitoring projection (10-20 attributes of a several-hundred-attribute ad) this skips the overwhelming majority of the per-ad decode work that a decode-everything-then-filter projection pays.

chaseRefs additionally resolves each emitted expression's attribute references against the same ad (transitively, per ad -- an "elevator" of linear passes that repeats only while new references surface), so a projected ad evaluates self-contained. HTCondor's query protocol sends exactly the requested attributes, so a protocol-compatible server passes false.

It applies to both representations: an interned collection resolves references by id, a persistent (inline-name) one by name.

redact strips private attributes exactly as ScanRawRedacted does. An empty projection means no attribute filter (the whole ad, matching QueryRawProject semantics upstream). Inline-name collections yield nothing, as with ScanRaw.

func (*Collection) ScanRawRedacted added in v0.16.3

func (c *Collection) ScanRawRedacted() iter.Seq[RawAd]

ScanRawRedacted is ScanRaw with private attributes stripped (see QueryRawRedacted).

func (*Collection) ScanRawWire added in v0.16.4

func (c *Collection) ScanRawWire(projection []string, redact bool) iter.Seq[[]byte]

ScanRawWire yields each ad as a WIRE-FORM ROW: a self-contained inline-names ad holding only the selected entries, assembled by slice copies from the stored bytes (see wire.AppendAdSubsetInline) -- nothing is decoded or rendered here. It is the relay scan for shipping a persistent table's ads across a process boundary (dbrpc) with the old-ClassAd render deferred to the far edge (RenderRawAdInline); the consumer needs no intern table and no data key (at-rest-encrypted values are opened during assembly).

projection restricts the entries to the named attributes (case-insensitive; MyType/TargetType always ship so the ad stays typed); empty means every entry. redact drops private attributes -- pre-pruned from a projection, or tested per entry with a byte-gated predicate for the whole-ad case. The yielded row aliases a buffer reused across the iteration. Non-inline (RAM) collections yield nothing: their ads are not self-contained, and an in-process consumer reads them through the collection directly.

func (*Collection) SchemaFit added in v0.25.2

func (c *Collection) SchemaFit(sampleMax int) ([]SchemaFieldFit, int)

SchemaFit measures the current schema against a fresh sample of up to sampleMax records, reporting per-field escape rates and the number of records actually sampled.

It reads samples and re-runs the encoder's storability test on each; it encodes nothing and writes nothing. Returns nil when the accelerator is not enabled (there is no schema to judge) or there is nothing to sample.

func (*Collection) SchemaScanInfo added in v0.23.1

func (c *Collection) SchemaScanInfo() SchemaScanInfo

SchemaScanInfo returns the columnar accelerator's current state (see SchemaScanInfo).

func (*Collection) SetEncryptedAttrs added in v0.7.0

func (c *Collection) SetEncryptedAttrs(attrs []string) error

SetEncryptedAttrs replaces the explicit encrypted-attribute set at runtime (the toggle meta-command). It is a no-op-with-error if encryption is disabled, and errors if any named attribute is currently indexed (an encrypted value is opaque, so it cannot be indexed). Private attributes are always encrypted regardless of this set. New records use the new set immediately; existing records keep their prior form until rewritten (compaction/Rewrite re-encodes them under the current policy).

func (*Collection) SetGCFloor added in v0.21.4

func (c *Collection) SetGCFloor(floor float64)

SetGCFloor installs a runtime GC watermark, in Retention.MinAgeAttr units, that lets Rotate reclaim already-consumed records EARLY: a segment whose newest MinAgeAttr value is below floor may be dropped before it reaches MaxAge. It is how a change-feed source drains records every live subscriber has acknowledged (floor is the feed's GC floor, i.e. the min ack over live subscribers). It only ever shortens retention -- it can never keep data past the configured ceilings (a slow or absent subscriber holds a low floor, so its data simply ages out under MaxAge), and it never drops anything younger than Retention.MinAge. Requires MinAgeAttr set (and zone-mapped). Passing floor <= 0 clears it. Not persisted -- callers re-assert it each pass from the current live floor; a stale saved value must never GC data across a restart.

func (*Collection) SetRetention added in v0.19.0

func (c *Collection) SetRetention(r Retention)

SetRetention updates the retention bounds at runtime (an append-only collection). The next Rotate enforces them. Serialized against Rotate/maintenance via maintMu, so c.ret is only ever read or written under that lock.

func (*Collection) SetRowGroupBytes added in v0.28.1

func (c *Collection) SetRowGroupBytes(n int)

SetRowGroupBytes changes the uncompressed record-bytes budget for a columnar row group. 0 restores the default.

Safe on a live collection, and it needs no reconciliation pass: every block records the layout it was written with, so a new budget governs row groups sealed from now on while everything already on disk keeps reading exactly as before. That is what makes this tunable against a real archive rather than only at creation.

func (*Collection) SetTimeTravel added in v0.12.0

func (c *Collection) SetTimeTravel(o *TimeTravelOptions)

SetTimeTravel enables, retunes, or (with a nil/zero option) disables point-in-time queries at runtime. Enabling starts recording checkpoints and retaining superseded versions from now on (it is not retroactive); disabling lets the next compaction reclaim the retained history. Safe to call concurrently with reads and writes.

func (*Collection) SidecarSizes added in v0.8.0

func (c *Collection) SidecarSizes() SidecarSizes

SidecarSizes reports the live Collection's sealed-segment sidecar bytes: the mmap-backed index each sealed segment holds after the flip from the in-RAM segIndex -- a file mapping for a persistent collection (page-cache resident, evictable to disk) or an anonymous mapping for an in-memory one (off-heap, MADV_FREE-able to swap). Either way these bytes are NOT Go-heap memory, so they are reported apart from IndexSizes (which now measures only the active, still-in-RAM segments' postings). Together the two give the operator the full picture: heap postings on the hot active segment plus off-heap sidecar bytes on the sealed tail. It reads each segment's already-mapped bytes under the shard read lock -- no re-mapping -- so it is cheap enough for periodic sampling.

func (*Collection) StaleIndexSegments added in v0.21.2

func (c *Collection) StaleIndexSegments() (stale, sealed int)

StaleIndexSegments reports sealed segments whose index was built under a superseded configuration and is still due a rebuild, out of the sealed segments total.

"Due a rebuild" is the operative part. With Options.IndexBackfillBytes set, a configuration change is deliberately carried only to recent segments; older ones keep the index they have, permanently and by design. Counting those would leave the number permanently non-zero and say nothing about whether anything is wrong -- so they are excluded here and reported separately by StaleIndexSegmentsByPolicy.

A persistent non-zero value from this therefore still means what it always meant: a reindex is failing, or never running.

func (*Collection) StaleIndexSegmentsByPolicy added in v0.25.0

func (c *Collection) StaleIndexSegmentsByPolicy() int

StaleIndexSegmentsByPolicy reports sealed segments left on a superseded index configuration because they fall outside the backfill horizon. Expected, not actionable: it is the size of the deliberate mixture, and it grows as the store does.

func (*Collection) StartAutoRetrain

func (c *Collection) StartAutoRetrain(interval time.Duration, sampleMax, hotTopN int) (stop func())

StartAutoRetrain runs periodic maintenance in a background goroutine every interval, until the returned stop function is called (stop blocks until the goroutine exits). Each tick:

  • retrains the ZSTD dictionary from a sample of up to sampleMax ads (RetrainDict), and
  • if hotTopN > 0, refreshes the hot-attribute set to the topN most common attributes (RefreshHotSet), so the store self-tunes which attributes it front-loads for fast queries.

Retrain errors (e.g. the pure-Go BuildDict declining a small/homogeneous corpus) are ignored — the previous codec keeps working.

Cost note: retraining recompacts every shard. Compaction is concurrent (the recompression runs without the shard lock), so writers and scanners are not blocked for its duration; only two brief per-shard critical sections take the lock. The 15-minute default still keeps the work rare.

func (*Collection) Stats

func (c *Collection) Stats() Stats

Stats returns a snapshot of the collection's storage. It takes each shard's read lock briefly, so it is safe to call concurrently with readers and writers.

func (*Collection) StopMaintenance added in v0.30.4

func (c *Collection) StopMaintenance()

StopMaintenance asks any maintenance pass running on this collection to stop at its next safe boundary. It does not wait, and it does not interrupt work already in progress on a single segment -- a pass abandoned mid-rewrite would be the one thing a caller must never do, since the shutdown that follows unmaps the segments a rewrite is still reading.

It exists because a maintenance pass is bounded but not short. The columnar budget caps a pass at 64 segment rewrites so a schema change converges over several passes instead of stalling on one, which keeps a pass unnoticeable in steady state -- but a caller closing the server waits for the pass in flight to finish, and on an archive at history scale that is tens of seconds of apparently-hung process after the logs say it stopped.

The abandoned work is not lost: every segment is rewritten at most once, so the backlog only shrinks, and the next pass after reopen picks up exactly where this one left off.

Once stopped a collection stays stopped; this is a shutdown signal, not a pause.

func (*Collection) SuggestDrops

func (c *Collection) SuggestDrops(sampleMax int) []DropSuggestion

SuggestDrops recommends configured indexes to remove, from observed query demand and a sample of up to sampleMax live ads. It flags two cases: an index no query has ever filtered on ("unused"), and one whose sampled values are effectively constant ("low-cardinality", ≤1 distinct value — the index cannot prune). It is advisory: apply via DropIndex. Demand is cumulative since the collection was created, so give a workload time to run before trusting "unused".

func (*Collection) SuggestIndexes

func (c *Collection) SuggestIndexes(sampleMax int) []IndexSuggestion

SuggestIndexes recommends value/categorical indexes from observed query demand and a sample of up to sampleMax live ads. It is advisory: apply the returned CategoricalAttrs/ValueAttrs via New (or a future dynamic Reindex). Attributes already indexed are not re-suggested. Results are ordered by demand, most first.

func (*Collection) SupportsRawWire added in v0.25.0

func (c *Collection) SupportsRawWire() bool

SupportsRawWire reports whether ScanRawWire/QueryRawWire can serve this collection. Only an inline (persistent) collection can: a RAM collection's ads are not self-contained, so the relay scans yield nothing for one. A caller MUST check this rather than treat an empty relay scan as an empty result -- they look identical.

func (*Collection) TimeTravelConfig added in v0.12.0

func (c *Collection) TimeTravelConfig() (opts TimeTravelOptions, enabled bool)

TimeTravelConfig reports the collection's current time-travel settings and whether it is enabled (for persisting the runtime toggle; see db.saveIndexConfig).

func (*Collection) TopKOrderThreshold added in v0.29.4

func (c *Collection) TopKOrderThreshold(q *vm.Query, orderAttr string, desc bool, k int, stats *ScanStats) (threshold float64, seen int, ok bool)

TopKOrderThreshold computes the cutoff order value for ORDER BY orderAttr {DESC|ASC} LIMIT k over the records matching q: the k-th largest (desc) or k-th smallest (asc). It returns (threshold, seen, ok): seen is the number of matching records with a numeric order value, and ok reports whether the columnar path could serve the request at all (same gate as NumStatsQuery: accelerator on, orderAttr a numeric schema field, q a conjunction of scalar numeric comparisons).

The threshold is meaningful only when seen > k (the keep filled); at seen <= k there are k or fewer rows and the caller should just fetch them directly.

stats (may be nil) records the cutoff scan's work for EXPLAIN ANALYZE: which records were read from columns vs reassembled by the active-segment fallback, and how many matched -- the same ScanStats a projected scan reports, so the top-K path is no longer invisible to the diagnostic.

func (*Collection) TrackedKeys added in v0.30.0

func (c *Collection) TrackedKeys() int

TrackedKeys reports how many keys the delta tracker holds. The tracker is the one new in-memory structure delta records add -- one small entry per key that has been written -- so this is what an operator checks when asking what the feature costs in RAM, and what a test checks to confirm the bound is doing something.

func (*Collection) Truncate added in v0.7.0

func (c *Collection) Truncate()

Truncate removes every ad from the collection, leaving it empty (an in-place reset used by a DB restore before it reloads a snapshot). It resets each shard's directory and count and retires its segments -- RAM segments are dropped for the GC; persistent segments are munmap'd + unlinked once no in-flight scan references them (the compaction pin/reap protocol). A concurrent scan that already captured its window reads the old data safely until it finishes; scans and writes that start after Truncate see an empty collection. Ordered indexes are cleared in place. Callers needing atomicity against concurrent writers must serialize Truncate with them (the db layer holds its DB-wide lock); at the shard level Truncate is itself consistent.

func (*Collection) Update

func (c *Collection) Update(batch []AdUpdate) error

Update applies a batch of inserts/updates and returns only once every ad is committed and visible to new scans. Ads are encoded (and compressed) outside the shard locks; each affected shard then applies its updates under one lock acquisition at a single new commit sequence.

func (*Collection) UpdateOld

func (c *Collection) UpdateOld(batch []OldAdUpdate) error

UpdateOld applies a batch of inserts/updates whose ads arrive in old-ClassAd form. It encodes each ad directly to the wire form, attribute by attribute, without materializing an intermediate ast.ClassAd: scalar-literal attributes (the common case) are written straight to wire, and only genuinely computed values are parsed with the expression parser. This is the efficient path for ads read from a socket. Commit semantics match Update.

func (*Collection) VectorEvalCount added in v0.27.0

func (c *Collection) VectorEvalCount(q *vm.Query) (int, bool)

VectorEvalCount counts records satisfying q, evaluating it a column at a time.

ok=false only when there is no columnar state at all; individual blocks that cannot be vectorized are served by colScope, and segments without a block by a row walk, so the answer is always complete.

func (*Collection) Watch added in v0.4.0

func (c *Collection) Watch(ctx context.Context, cursor []byte) (iter.Seq[WatchEvent], error)

Watch replays everything that may have changed since cursor (nil ⇒ a full replay from empty), then streams live changes until ctx is cancelled or the consumer stops. Requires Options.WatchHistory > 0. See docs/WATCH.md and WatchEvent.

func (*Collection) WatchCursor added in v0.7.0

func (c *Collection) WatchCursor() ([]byte, error)

WatchCursor returns an opaque cursor pointing at the current head of the change log, so a subsequent Watch(ctx, cursor) streams only changes from now on -- no initial replay of the current contents. Requires Options.WatchHistory > 0. It only snapshots each shard's commit sequence (no registration, no replay), so it is cheap.

func (*Collection) WatchRedacted added in v0.28.0

func (c *Collection) WatchRedacted(ctx context.Context, cursor []byte) (iter.Seq[WatchEvent], error)

WatchRedacted is Watch for a watcher NOT entitled to sealed values: every event's ad is decoded with no key, so a sealed attribute arrives undefined rather than opened. See Collection.QueryRedacted.

A watch streams ads continuously, so it is the read path with the longest exposure: decoding with the collection's key and trusting the serializer to drop private attributes means the secret is opened for every event of every unentitled watcher.

func (*Collection) ZoneAttrs added in v0.21.2

func (c *Collection) ZoneAttrs() []string

ZoneAttrs returns the attributes this collection keeps per-segment [min,max] zone maps on, in configuration order (explicit ZoneAttrs first, then the value-indexed attributes that are zoned automatically). Empty for a mutable collection, which has no zone maps. Diagnostic only -- it is how `.indexes` on an archive can show which range queries prune whole segments rather than only postings.

type CommitResult added in v0.6.0

type CommitResult struct {
	Committed int
	Conflicts [][]byte
	Unapplied [][]byte
	// UnappliedReasons is parallel to Unapplied: why each key could not be composed. A dropped
	// key reported without a reason is a bare fact with no action attached -- an operator reading
	// it had to go to the daemon ad's counters and correlate by hand.
	UnappliedReasons []string
}

CommitResult reports a transaction's outcome. Conflicts holds the keys whose write lost a write-write race and were not applied; the caller may re-read and retry just those. The other buffered writes committed.

Unapplied is the OTHER kind of not-applied, and the distinction is the point: those writes could not be composed at all, and re-applying the identical write cannot change that. A patch whose stored base is unreadable is the case that exists today -- the key is present but the record behind it will not come back (a delta chain that no longer materializes, a segment that went away), so there is nothing to merge the patch onto.

Reporting those as Conflicts made a caller that retries conflicts loop forever. On a production mirror each one rewound the tailer, re-applied, failed identically, and after three attempts escalated to a full 460-second replay of a 923 MB log that wrote nothing -- 285 times, which is why the mirror could never catch up. A caller must be able to tell "try again" from "this will never work"; nothing else about the commit changes.

func (CommitResult) Conflicted added in v0.6.0

func (r CommitResult) Conflicted() bool

Conflicted reports whether any buffered write lost a conflict.

func (CommitResult) HasUnapplied added in v0.30.2

func (r CommitResult) HasUnapplied() bool

HasUnapplied reports whether any buffered write could not be composed at all. Retrying it is futile; the caller should record it and make progress.

type ConjunctExplain added in v0.6.0

type ConjunctExplain struct {
	// Text is the conjunct rewritten over the slot (job refs baked, TARGET scope dropped).
	Text string `json:"text"`
	// Probed is true when the conjunct is covered by an active index probe -- it prunes
	// the candidate set before the bilateral re-verify.
	Probed bool `json:"probed"`
	// Indexed is true when the conjunct's selectivity is estimable from an index (a
	// superset of Probed: it also covers bare boolean flags, which are estimable via
	// `attr == true` but are not extracted as pushdown probes).
	Indexed bool `json:"indexed"`
	// HasSelectivity reports whether TrueFrac is populated.
	HasSelectivity bool `json:"hasSelectivity"`
	// TrueFrac is the estimated fraction of slots for which the conjunct holds.
	TrueFrac float64 `json:"trueFrac"`
	// ResourceSide is true for a conjunct that came from the MATCH's resource-side filter
	// (WHERE TARGET / NOPREEMPT) rather than the job's Requirements. It is applied to
	// matched candidates as a post-filter, so it re-checks rather than prunes today.
	ResourceSide bool `json:"resourceSide,omitempty"`
}

ConjunctExplain is one top-level && conjunct in the match's evaluation order.

type DropSuggestion

type DropSuggestion struct {
	Attr   string // canonical attribute name
	Kind   string // "categorical" or "value" (how it is currently indexed)
	Reason string // "unused" (never queried) or "low-cardinality" (no pruning power)

	QueriesEq      int64
	QueriesRange   int64
	SampledPresent int
	DistinctValues int
	Capped         bool
}

DropSuggestion recommends removing one configured index, with the rationale.

type GroupAgreementItem added in v0.29.13

type GroupAgreementItem struct {
	Attrs []string `json:"attrs"`
	Frac  float64  `json:"frac"`
}

GroupAgreementItem pairs a group's members with the fraction of segments that re-derived it.

type GroupCount added in v0.28.0

type GroupCount struct {
	Value classad.Value
	Count int
}

GroupCount is one group of a GROUP BY: the group column's value and how many records held it.

type GroupSchemaAgreement added in v0.28.0

type GroupSchemaAgreement struct {
	Segments int `json:"segments"`
	// PerGroup[i] is the fraction of segments whose own derivation produced group i's exact
	// member set, for the table-level groups in report order.
	PerGroup []float64 `json:"perGroup,omitempty"`
}

GroupSchemaAgreement reports how well per-segment derivations agree with a whole-table one -- the other half of "is this a property of the data". A group derived from the table but absent from most segments is a sampling artifact, and a schema pointer spent on it would buy coverage in some segments and nothing in others.

type GroupSchemaChange added in v0.29.13

type GroupSchemaChange struct {
	Unix    int64              `json:"unix"`
	Reason  string             `json:"reason"`
	Added   [][]string         `json:"added,omitempty"`
	Removed [][]string         `json:"removed,omitempty"`
	Changed []GroupSchemaDelta `json:"changed,omitempty"`
}

GroupSchemaChange is one committed-group-set change, for reporting.

type GroupSchemaDelta added in v0.29.13

type GroupSchemaDelta struct {
	Before  []string `json:"before"`
	After   []string `json:"after"`
	Added   []string `json:"added,omitempty"`
	Removed []string `json:"removed,omitempty"`
}

GroupSchemaDelta is one surviving group whose membership shifted across a change.

type GroupSchemaDrift added in v0.28.0

type GroupSchemaDrift struct {
	// Derivations is how many candidate SNAPSHOTS the sampler has retained (capped at
	// groupHistoryMax), NOT how many times the committed set changed -- most snapshots are
	// identical. CommittedChanges is the count that answers "how much real churn": entries in the
	// committed-set change log (see GroupSchemaChanges).
	Derivations      int   `json:"derivations"`
	CommittedChanges int   `json:"committedChanges"`
	FirstUnix        int64 `json:"firstUnix,omitempty"`
	LastUnix         int64 `json:"lastUnix,omitempty"`
	// Retained is how many of the FIRST derivation's groups still appear, by exact member
	// set, in the last -- the number that matters, because a group whose membership changed
	// is a different group and its block would have to be rebuilt.
	Retained int `json:"retained"`
	OfFirst  int `json:"ofFirst"`
	// MaxPartialFrac is the worst partial fraction across the last derivation's groups.
	// Nonzero means co-occurrence has decayed: some ad now holds part of a group.
	MaxPartialFrac float64 `json:"maxPartialFrac"`
}

GroupSchemaDrift compares the earliest and latest retained derivations, so an operator can see whether the groups are holding without having captured the earlier report themselves.

type GroupSchemaEntry added in v0.28.0

type GroupSchemaEntry struct {
	Attrs []string `json:"attrs"`
	// InFrac / NoneFrac / PartialFrac are the sample's split over this group. PartialFrac is
	// the one to watch: it is the fraction that would fall back to a row decode, and it is 0
	// by construction at derivation, so any growth is drift.
	InFrac      float64 `json:"inFrac"`
	NoneFrac    float64 `json:"noneFrac"`
	PartialFrac float64 `json:"partialFrac"`
	// Cells is the attribute occurrences this group would make columnar, and CellsFrac that
	// as a fraction of every occurrence in the sample.
	Cells     int     `json:"cells"`
	CellsFrac float64 `json:"cellsFrac"`
}

GroupSchemaEntry is one group in a report.

type GroupSchemaInfo added in v0.28.0

type GroupSchemaInfo struct {
	// Sampled is how many ads the derivation looked at, and BaseFields the size of the base
	// schema it derived against.
	Sampled    int `json:"sampled"`
	BaseFields int `json:"baseFields"`
	// BaseCells / TotalCells are attribute occurrences the base schema covers, out of all
	// occurrences in the sample -- the coverage the groups are measured against.
	BaseCells  int                `json:"baseCells"`
	TotalCells int                `json:"totalCells"`
	Groups     []GroupSchemaEntry `json:"groups,omitempty"`
}

GroupSchemaInfo reports the derived group schemas, for diagnostics. Phase 1 is report-only: these are what WOULD be stored columnar per group, not what is.

type GroupSchemaLastAgreement added in v0.29.13

type GroupSchemaLastAgreement struct {
	Unix     int64                `json:"unix"`
	Segments int                  `json:"segments"`
	Groups   []GroupAgreementItem `json:"groups,omitempty"`
}

GroupSchemaLastAgreement is the last per-segment agreement result, read from the sidecar so a READ-level client can see it without recomputing (see GroupSchemaAgreement).

type GroupStats added in v0.28.0

type GroupStats struct {
	Value classad.Value
	Count int
	Stats []NumStats
}

GroupStats is one group of a GROUP BY with value aggregates: the group column's value, how many records fell in the group, and one NumStats per requested aggregate attribute, in the order requested.

type Hasher

type Hasher interface {
	Hash(key []byte) uint64
}

Hasher maps a stable key to a 64-bit hash. The hash routes a key to a shard and indexes the per-shard directory; it is never used as the ad's identity (the full key is stored inline in each record and compared on lookup), so hash collisions are resolved exactly and never lose data.

type IndexSize added in v0.6.0

type IndexSize struct {
	Attr  string `json:"attr"`
	Kind  string `json:"kind"`  // "categorical" | "value"
	Bytes int64  `json:"bytes"` // resident posting bytes (roaring bitmaps + keys) across all live segments
	// SketchBytes is the resident memory of this attribute's per-segment sketches --
	// the categorical bloom filter and the HyperLogLog distinct-count registers --
	// reported apart from Bytes so it is visible rather than hidden in the posting total.
	SketchBytes int64   `json:"sketchBytes"`
	Auto        bool    `json:"auto"` // created by the auto-tuner (vs human/Options)
	Frac        float64 `json:"frac"` // Bytes as a fraction of the live data bytes
}

IndexSize is the measured memory footprint of one attribute's index.

type IndexSizes added in v0.6.0

type IndexSizes struct {
	PerIndex   []IndexSize `json:"perIndex"`
	TotalBytes int64       `json:"totalBytes"` // posting bytes (the auto-tuner's budget denominator)
	// TotalSketchBytes is the sum of every index's SketchBytes (bloom + HLL). It is
	// reported separately and is NOT folded into TotalBytes/Frac, so the watermark and
	// auto-tuner budget stay calibrated on posting bytes; sketch memory is bounded and
	// small (<=8 KiB bloom + 1 KiB HLL per categorical attr per segment).
	TotalSketchBytes int64   `json:"totalSketchBytes"`
	DataBytes        int64   `json:"dataBytes"` // live compressed record bytes
	Frac             float64 `json:"frac"`      // TotalBytes / DataBytes
}

IndexSizes is the collection's index memory, per attribute and in total, against the live data bytes -- the denominator for the "index is N% of data" watermark.

type IndexSuggestion

type IndexSuggestion struct {
	Attr string // canonical (first-seen) attribute name
	Kind string // "categorical" (string equality/membership) or "value" (numeric + range)

	QueriesEq    int64 // equality/membership probes observed
	QueriesRange int64 // range probes observed

	SampledPresent int     // sampled ads with this attribute present
	StringFrac     float64 // fraction of present values that are string literals
	NumericFrac    float64 // fraction that are numeric literals
	DistinctValues int     // distinct values in the sample (a lower bound if Capped)
	Capped         bool
}

IndexSuggestion recommends indexing one attribute, with the rationale (the two signals) so the decision is explainable.

type MatchExplain added in v0.6.0

type MatchExplain struct {
	// HasRequirements is false when the job has no Requirements (it then matches every
	// slot -- a full scan with no pruning).
	HasRequirements bool `json:"hasRequirements"`
	// SlotPredicate is the job's Requirements rewritten over the slot: TARGET.attr
	// becomes the slot's own attribute and every job reference is baked to its value,
	// e.g. `Memory >= 4096 && Arch == "X86_64"`. This is the predicate whose probes
	// drive candidate pruning (the bilateral match still re-verifies every candidate).
	SlotPredicate string `json:"slotPredicate"`
	// Probes are the rewritten predicate's index-satisfiable conjuncts and their index
	// status on the resource collection.
	Probes []ProbeExplain `json:"probes"`
	// IndexUsable is how many probes prune via an index.
	IndexUsable int `json:"indexUsable"`
	// Plan is the access path over the resource slots: "empty" (the request contradicts
	// itself, so no slot can match), "indexed" (visit only candidate slots),
	// "parallel-scan", or "serial-scan" (match every slot).
	Plan        string `json:"plan"`
	Parallelism int    `json:"parallelism"`
	Shards      int    `json:"shards"`
	// TotalResources is the resource (slot) count, the denominator for selectivity.
	TotalResources int `json:"totalResources"`
	// EvalOrder is the top-level && conjuncts in the order the bilateral match evaluates
	// them (after short-circuit reordering), each tagged with its role: an active pruning
	// probe, or a re-check whose selectivity (from the index, when available) drove where
	// it sorts. It makes the ordering transparent -- e.g. an always-true capability flag
	// appears last as a ~99%-true re-check, not a candidate filter.
	EvalOrder []ConjunctExplain `json:"evalOrder,omitempty"`
}

MatchExplain describes how matchmaking a specific job against this (resource) collection would execute: the job's Requirements rewritten over the slot (with the job's attributes baked to constants), the index-satisfiable probes that rewrite yields, and which of them prune candidates via a configured index.

type MergeOptions added in v0.24.1

type MergeOptions struct {
	// TargetSegments is the LOW watermark: once a pass starts it merges until the archive
	// is at or below this. Default defaultTargetSegments.
	TargetSegments int
	// TriggerSegments is the HIGH watermark: a pass does nothing until the archive reaches
	// it. The gap between the two is what stops the policy from merging a run on every pass
	// forever once it is sitting on the target -- rewriting data continuously to hold a line
	// that does not need holding to the segment. Default: TargetSegments plus a quarter.
	TriggerSegments int
	// MaxSegmentBytes caps a merged segment. Segment offsets are uint32 throughout the
	// record and sidecar formats, so 4 GiB is a hard structural ceiling; the default stays
	// well under it. A run stops before exceeding this.
	MaxSegmentBytes int64
	// MinMergeBytes is the smallest output worth producing. A run below it is still merged
	// if it is the best available -- reducing the count is the point -- but the policy
	// prefers to keep accumulating.
	MinMergeBytes int64
	// MaxRun bounds how many segments one merge consumes, so a single merge's work and
	// memory stay predictable however small the segments are.
	MaxRun int
	// KeepRecent leaves the newest sealed segments untouched, keeping the hot end of the
	// archive finely divided (and away from the open segment the writer is appending to).
	KeepRecent int
	// MaxMerges bounds one pass, so maintenance stays interruptible and a large backlog is
	// worked down over several passes instead of one long stall.
	MaxMerges int
	// MaxBytesPerPass bounds the SOURCE bytes one pass rewrites. MaxMerges alone does not
	// bound the work, because a merge's cost is its inputs' size, not its count -- one pass
	// of large runs can move orders of magnitude more than another of the same length.
	//
	// This exists for the catch-up case: an archive that grew while merging was off has a
	// backlog measured in the size of the archive, and without a byte bound the first pass
	// would try to rewrite all of it at once. Steady state does not need it -- a pass stops
	// at the low watermark long before this -- so the default is sized to cap one pass at a
	// tolerable stall (~25s at the ~160 MB/s merging sustains) rather than to throttle.
	//
	// Sustained RATE is the scheduler's job, not this: pace passes so the long-run average
	// stays within the byte budget the deployment allows.
	MaxBytesPerPass int64
}

MergeOptions tunes a merge pass. The zero value is usable: it fills in the defaults below.

type NumStats added in v0.25.2

type NumStats struct {
	N   int     `json:"n"`
	Sum float64 `json:"sum"`
	Min float64 `json:"min"`
	Max float64 `json:"max"`
	// IntSum accumulates the INTEGER contributions in int64, exactly as the reference SUM does
	// -- a float64 accumulator loses precision past 2^53, so an all-integer sum rendered from
	// Sum would disagree with the scanning aggregator on large values.
	IntSum int64 `json:"intSum,omitempty"`
	// AnyReal records whether any contributing value was a REAL rather than an integer, which is
	// exactly the reference's promotion rule: SUM stays an integer unless a real appears, and
	// MIN/MAX keep their element's type. Without it the same query would print differently
	// depending on whether the accelerator was on.
	AnyReal bool `json:"anyReal,omitempty"`
	// AnyBool records that a boolean turned up in the column (only possible as an escaped value,
	// since a bool is not a numeric schema field). The reference coerces booleans to 1/0 and has
	// a further quirk for a lone boolean element; rather than reproduce that, a caller declines
	// and lets the scan answer. Pathological data, exact answer.
	AnyBool bool `json:"anyBool,omitempty"`
}

NumStats is one columnar pass's worth of aggregate inputs for a single numeric attribute, from which MIN, MAX, SUM, AVG and COUNT(attr) all follow.

N counts the records whose value for the attribute is a number -- which is COUNT(attr)'s definition, and the divisor for AVG. Records where the attribute is absent or non-numeric contribute to none of these fields. Min/Max are meaningless when N == 0.

type OldAdUpdate

type OldAdUpdate struct {
	Key  []byte
	Text string
}

OldAdUpdate is one insert-or-update whose ad is supplied in "old ClassAd" serialization (newline-separated `Name = Value` lines, as sent over a TCP socket), rather than as a parsed *classad.ClassAd.

type OpStat added in v0.11.0

type OpStat struct {
	Count int64 `json:"count"`
	Nanos int64 `json:"nanos"`
	// MaxNanos is the longest single occurrence seen, and Buckets is the latency
	// histogram (see LatencyBucketBoundsNanos: buckets[i] counts occurrences <= bound[i],
	// the last entry counts the overflow above the final bound). Both surface the tail
	// the mean (Nanos/Count) hides. Omitted from JSON when empty so an older peer that
	// never populated them stays byte-compatible.
	MaxNanos int64   `json:"maxNanos,omitempty"`
	Buckets  []int64 `json:"buckets,omitempty"`
}

OpStat is one operation's cumulative call count and total wall-nanoseconds. Both are monotonic; a scraper derives rate and mean latency (Nanos/Count) from deltas.

type OpStats added in v0.11.0

type OpStats struct {
	ShardWriteWait OpStat `json:"shardWriteWait"`
	ShardWriteHold OpStat `json:"shardWriteHold"`
	SegmentAlloc   OpStat `json:"segmentAlloc"`
	Sync           OpStat `json:"sync"`
	// CommitSync is the per-commit durability wall time: how long a Txn.Commit blocked
	// flushing all of its shards to disk. Since those shard flushes now run in parallel,
	// this is the commit's critical path (≈ the slowest shard), which the per-shard Sync
	// total -- now the summed fsync WORK, not the latency -- no longer reveals. Count is
	// the number of durable commits; Nanos/Count is mean commit durability latency.
	CommitSync OpStat `json:"commitSync,omitempty"`
	Compact    OpStat `json:"compact"`
	Retrain    OpStat `json:"retrain"`
	Reindex    OpStat `json:"reindex"`
}

OpStats is a snapshot of a collection's operational timing counters -- the wall time its callers spent blocked in, or holding, each of the store's stall points. Shard counters are summed across all shards. See opMetrics for what each measures. OpStats reports, per stall point, the cumulative call count and total wall time. The collection-wide snapshot lock (Truncate/Restore) lives one layer up in the db package, which composes its own snapshot-lock counter onto this snapshot.

type OpenIndexDiag added in v0.29.0

type OpenIndexDiag struct {
	Dir              string // the collection's directory
	SealedSegments   int    // sealed (non-active) segments seen at open
	SidecarFiles     int    // of those, how many had a .idx sidecar file on disk
	KeyIndexAdopted  int    // key-index sidecar mmapped (the primary-key index; no key rebuild)
	AttrIndexAdopted int    // attribute-index sidecar mmapped -- THESE need no Reindex rebuild
	// Reasons counts, per segment whose attribute index was NOT adopted (and so will be rebuilt),
	// why it was not:
	//   "no-sidecar-file"              -- no .idx on disk (segment predates sidecar persistence)
	//   "sidecar-rejected"            -- .idx present but the container/key index was invalid
	//                                    (older on-disk format, truncation, or a version bump)
	//   "attr-section-missing-or-stale" -- key index adopted, but the attribute-index section was
	//                                    absent, CRC-bad, or did not cover the segment
	Reasons map[string]int

	// StrandedSealedDeltas is how many LIVE delta records were found in sealed segments, and
	// StrandedSealedKeys names up to openStrandedSampleMax of them. Both are zero on a healthy
	// store: the seal-collapse invariant says a live fragment only ever sits in an active
	// segment, and every collapse pass relies on that. A non-zero count means those keys are
	// invisible to the collapse and will lose their base to the next compaction.
	StrandedSealedDeltas int
	StrandedSealedKeys   [][]byte

	// Timing breaks a persistent Open into its phases so a slow reopen names the phase that
	// dominated, instead of surfacing only as one elapsed number. Durations accumulate across the
	// collection's shards; all zero for an in-memory Open. See OpenTiming. Purely observational.
	Timing OpenTiming
}

OpenIndexDiag summarizes, for a single Collection Open, how each sealed segment's persisted index sidecar was handled: adopted (mmapped, no work) vs. left for the post-Open Reindex to rebuild -- which decodes every record the index covers and is the dominant reopen cost (~95% of reopen time on a large archive). It is emitted once per Open via OpenIndexDiagHook.

It exists to answer, from a production log, "why is startup rebuilding the archive indexes instead of adopting the ones on disk?" -- distinguishing an archive whose segments predate sidecar persistence (no files) from one whose sidecars are present but rejected (format/version) or whose attribute index is absent/stale. Purely observational; it changes no behavior.

type OpenTiming added in v0.29.2

type OpenTiming struct {
	MapSegments    time.Duration // readdir + mmap each segment file + derive its write extent
	DirRestore     time.Duration // load the dir snapshot, or (on a miss) rebuild it by walking records
	DirRebuilt     bool          // true if the directory was rebuilt (snapshot missing/stale) not restored
	PublishDict    time.Duration // restore each interned segment's attribute dictionary
	PublishColumns time.Duration // publish each columnarized segment's columnar payload
	LoadSidecars   time.Duration // map each sealed segment's index sidecar (adopt zone map + indexes)
	ZoneRecompute  time.Duration // of LoadSidecars, time decoding records to rebuild a missing zone map
	Reindex        time.Duration // post-recovery rebuild of indexes not adopted from sidecars
	RebuildOrdered time.Duration // post-recovery rebuild of the maintained ordered indexes
	AdoptSchema    time.Duration // re-enable columnar schema-scan from persisted blocks
	LoadDicts      time.Duration // load on-disk dictionaries before segment recovery
	LoadDemand     time.Duration // load checkpointed query demand
}

OpenTiming records how long each phase of a persistent Open took, accumulated across shards. A reopen that adopts every sidecar (OpenIndexDiag reasons empty) but is still slow is explained here: the cost is in mapping segments, rebuilding the directory (a snapshot miss forces a full record walk), recomputing zone maps for segments that lack one, or the post-recovery index/order rebuilds -- and this says which.

type Options

type Options struct {
	// Shards is the number of independently-locked shards; rounded up to a power
	// of two. Default 16. Forced to 1 when AppendOnly (an append log is a single
	// ordered sequence: it needs one segment chain, not sharded key routing).
	Shards int
	// AppendOnly makes the collection a pure append log: every Put appends a new
	// record (no per-key supersession, no directory, no compaction, no deletes) --
	// the model a history archive needs. Records accumulate in commit order and are
	// reclaimed only by whole-segment retention, not compaction. Scan/Query see every
	// appended record; Get/Delete-by-key are not meaningful (there is no key index).
	AppendOnly bool
	// ReverseScan makes Scan/Query yield records newest-first (most recently
	// appended first) instead of in stored order. It walks segments from the last
	// (newest) backward and, within a segment, records back-to-front -- the natural
	// "last K" order a history archive presents (condor_history shows the most
	// recent jobs first). A consumer that stops early (yield/emit returns false)
	// after K records therefore gets the K newest, a pushed-down LIMIT. Reverse
	// scans always run serial (never fan out). Default false ⇒ stored order.
	// Most meaningful with AppendOnly, where stored order is commit order.
	ReverseScan bool
	// Retention bounds how much an AppendOnly collection keeps: Rotate(now) drops
	// whole oldest segments until the collection is back within these bounds (a
	// history archive's aging-out). The zero value keeps everything (Rotate is
	// then inert). Only meaningful with AppendOnly. See Rotate.
	Retention Retention
	// IndexBackfillBytes bounds how far back an index configuration change is carried. When
	// set, only segments within the newest IndexBackfillBytes have their existing index
	// rebuilt under a new configuration; older ones keep the index they already have.
	//
	// The point is that rebuilding is not free -- it decompresses every record it covers --
	// while an out-of-date index is still perfectly usable for the attributes it does hold
	// (recovery uses whatever a segment's index covers, deciding per probe). On an archive
	// queried newest-first with a limit, the oldest segments may never be reached at all, so
	// re-indexing them buys little and costs the whole archive.
	//
	// 0 carries a change across every segment, which is the historical behaviour.
	IndexBackfillBytes int64

	// InternAtSeal makes an AppendOnly collection intern each segment as soon as it
	// seals (when the active segment fills), rather than only when RetrainDict/Rewrite
	// later reseal it. Interning stores segment-local ids + a per-segment name dictionary
	// instead of inline names, halving raw mmap bytes and speeding decode (see the
	// interning design). This trades a per-seal transcode (off the write lock) for the
	// density/decode win landing immediately on an archive that never retrains. Only
	// meaningful with AppendOnly (a mutable segment's records can be superseded in place,
	// which an off-lock transcode would race -- there interning rides compaction instead).
	// Optional; default off preserves exactly today's behavior. See InternSealed.
	InternAtSeal bool
	// ZoneAttrs names numeric attributes to keep a per-segment [min,max] zone map on,
	// so a range/equality query on one (e.g. CompletionDate > T) skips whole sealed
	// segments whose span cannot match, and age-based Retention (MaxAgeAttr) can drop
	// a segment by its newest timestamp. Value-indexed attributes (ValueAttrs) are
	// zoned automatically. Only meaningful with AppendOnly. Optional. See zonemap.go.
	ZoneAttrs []string
	// SegmentSize is the arena segment size in bytes. Default 8 MiB.
	SegmentSize int
	// Hasher routes keys to shards / directory buckets. Default 64-bit FNV-1a.
	Hasher Hasher
	// Codec compresses stored ad bytes. Default identity (no compression).
	Codec Codec
	// DeltaMax turns on DELTA RECORDS and bounds the chain: it is IGNORED on an AppendOnly
	// collection, where there is no earlier version of a key to merge with and the tracker
	// would grow per record forever. A write that names the attributes
	// it changed stores only those, until a key has accumulated DeltaMax deltas, at which
	// point the next write stores the whole ad again. 0 (default) stores whole ads always,
	// which is the historical behavior.
	//
	// It trades write volume for read work. A live job ad averages ~7.3 KB while one
	// transaction's changes average ~150 B, so writes shrink by ~49x; a point read has to
	// merge up to DeltaMax fragments over a full record. Requires inline-name records (a
	// persistent store) and is not combined with columnarization, which assumes each record
	// holds a whole ad.
	DeltaMax int
	// HotAttrs names the "popular" attributes to front-load in each ad's hot
	// header, so a query filtering on them resolves each in O(1) instead of
	// scanning the ad body. Typically the attributes common queries filter on
	// (e.g. Cpus, Memory, Arch, OpSys, State). Optional.
	HotAttrs []string
	// MatchClosureRoots names attributes (typically "Requirements") whose transitive
	// self-reference closure is front-loaded into each ad's hot header at encode time,
	// so matching a wide ad reads only the match-relevant attributes (via the hot
	// header) instead of decoding the whole ad. Optional; enables the hot-closure match
	// fast path. The frequency/HotAttrs hot set is unioned in, so queries are unaffected.
	MatchClosureRoots []string
	// GroupSchemaCount is how many SECONDARY columnar schemas to derive and build alongside the
	// base one, for attributes the base schema does not carry which are present or absent
	// together. 0 uses defaultGroupSchemas; a NEGATIVE value builds none.
	//
	// On by default because the base schema's coverage is otherwise hostage to the mix of ads in
	// the table: measured on a production AP, a history table's base schema covers 90.3% of
	// attribute occurrences until jobs removed before ever running are mixed in, at which point
	// it falls to 51.3%, while base plus four groups holds 87-94% across the range. Enabling by
	// default still commits no storage on the strength of one sample -- nothing is built until a
	// group's members have kept recurring across GroupStabilityRuns maintenance passes.
	GroupSchemaCount int
	// ColumnarSegmentBudget is how many sealed segments one maintenance pass may rewrite into
	// COLUMNAR-NATIVE form -- each attribute the schema carries stored once, in the segment's own
	// columnar payload, and removed from the records (see colnative.go). 0 uses
	// defaultColumnarBudget; a NEGATIVE value disables the rewrite and keeps every segment
	// whole-record with a columnar copy beside it in the sidecar.
	//
	// On by default because the duplication is pure cost: without it every value the schema
	// carries is stored twice, row-form in the arena and again in the sidecar block. Measured on
	// 1500 real OSPool machine ads, moving them removed 43% of the table (4104 -> 2327 bytes per
	// record), and the sidecar megabyte it reclaimed was entirely the duplicate copy.
	//
	// Budgeted rather than unbounded because turning it on over an existing archive rewrites every
	// sealed segment once, and at history scale doing that in a single pass is hours of I/O and a
	// flushed page cache -- the same failure a retrain caused when it resealed a whole archive.
	// Each segment is rewritten at most once, so a bounded budget still converges.
	ColumnarSegmentBudget int
	// RowGroupBytes is the uncompressed record-bytes budget for one columnar row group. 0 uses
	// colGroupTargetBytes (128 KiB).
	//
	// This is the knob that decides how much of a segment has to be decompressed to read a single
	// record, against how far compression can see across records -- larger groups compress better and
	// cost more per point lookup. The default is measured (see colGroupTargetBytes); it is exposed
	// because the right point depends on the read mix and on how much of the working set fits the
	// block cache, which a large production table is the only honest place to find out.
	//
	// The row cap and the 8-row group alignment are deliberately not configurable.
	RowGroupBytes int
	// GroupStabilityRuns is how many consecutive derivations a group's members must have
	// co-occurred in before its blocks are built. 0 uses the default; 1 disables the gate.
	GroupStabilityRuns int
	// GroupMergeJaccard widens a group by absorbing attributes whose presence pattern is at least
	// this similar. 0 uses defaultGroupJaccard; a NEGATIVE value keeps exact co-occurrence only.
	// GroupMaxPartialFrac bounds the fraction of ads that may then hold only PART of a group and
	// take the slow path for it; 0 uses defaultGroupMaxPartial.
	//
	// Measured against a later snapshot of the same table, widening still recovered more than
	// exact grouping (25.14% of attribute occurrences against 23.53%), because a partial ad reads
	// from the cold tail -- where it would have read anyway -- rather than paying a penalty. The
	// in-sample partial rate does NOT predict the later one, though, and no widened member set
	// reproduced across snapshots, so the stability gate refuses to build one. See
	// mergeNearPatterns.
	GroupMergeJaccard   float64
	GroupMaxPartialFrac float64
	// DemandHalfLife is how quickly recorded query demand fades, so index decisions
	// track the current workload rather than everything the process has ever seen. 0
	// uses defaultDemandHalfLife; negative disables decay (counters accumulate for the
	// lifetime of the data, as they did before decay existed).
	DemandHalfLife time.Duration
	// CommitSync, if set, is called once per committed (possibly group-coalesced)
	// batch on a shard, after the writes are applied — the point at which a future
	// durable collection would fsync/serialize the batch. Group commit amortizes
	// this across all writers that commit together. It must be safe for concurrent
	// use across shards. Optional; default no-op (in-memory store).
	CommitSync func()
	// CategoricalAttrs names string-valued attributes to index for equality and
	// set-membership (`Attr == "x"`, `Attr == "x" || Attr == "y"`). ValueAttrs
	// names numeric attributes to index for equality and range (`Attr >= n`). A
	// query filtering on an indexed attribute visits only candidate ads instead of
	// scanning all of them. Optional.
	CategoricalAttrs []string
	ValueAttrs       []string
	// Ordered configures maintained, filtered, ordered indexes -- the schedd
	// priority-queue pattern (see docs/MATCH.md §3 and Collection.Ordered). Each
	// spec groups members into partitions and keeps them sorted on the write path,
	// so a negotiator can iterate a partition in order (and resume) without
	// re-sorting each cycle. A malformed expression in a spec panics New. Optional.
	Ordered []OrderSpec
	// Dir, if set, makes the collection persistent: arenas are memory-mapped files
	// under this directory and committed writes are flushed to disk (see Open).
	// Empty ⇒ in-memory (the default). Unix-only.
	Dir string
	// QueryParallelism controls cross-segment fan-out for large full-scan queries:
	//   0 (default) ⇒ auto — the library picks the policy (currently: up to 6 workers
	//                 per query, clamped to GOMAXPROCS). The meaning of auto may change
	//                 across releases.
	//   1           ⇒ serial (never fan out).
	//   N ≥ 2       ⇒ fan out with up to N workers per query.
	// In every non-serial case fan-out is still gated by a work-size threshold, never
	// exceeds the segment count, and draws from a machine-wide worker budget shared
	// across concurrent queries (so it degrades to serial under load rather than
	// oversubscribing). Small scans and indexed queries run serial regardless.
	// See parallel_scan.go / docs/PARALLEL_QUERY.md.
	QueryParallelism int
	// WatchHistory enables the Watch subscription verb (see docs/WATCH.md). It is the
	// per-shard capacity of the delete journal — how many recent deletes are retained
	// so a resuming watcher can be told precisely which keys were removed. 0 (default)
	// disables Watch entirely (no journal, no live notification, zero write-path cost).
	// A larger value lets clients resume incrementally from further back before
	// falling to a full replay.
	WatchHistory int
	// WatchBuffer is the per-watcher live event-channel capacity (default 1024). A
	// watcher that overflows it is told to resync rather than stalling writers.
	WatchBuffer int
	// WatchCoalesce, if > 0, batches live Watch events into windows of this
	// duration and emits only the newest event per key in each window (default 0 =
	// off, one event per change). It smooths bursty churn -- e.g. a freshly
	// submitted job that takes many SetAttribute updates in quick succession is
	// delivered as a single Upsert of its settled state. Only the last event of a
	// flushed window carries a cursor, so a mid-window consumer crash resumes from
	// the prior window and re-delivers (at-least-once is preserved). Catch-up is
	// never coalesced.
	WatchCoalesce time.Duration

	// ParentKeyFor, if set, gives an ad a parent: it maps a child's key to its
	// parent's key (return nil for a top-level ad with no parent). A child then
	// resolves attribute references it lacks by falling through to its parent --
	// the primitive "join" the job queue needs, where a proc ad chains to its
	// cluster ad so a query like `DAGManJobId == 42` sees the cluster's attribute.
	// A family (a root and all its descendants) is co-located in one shard -- the
	// collection routes every key to the shard of its ultimate root -- so chained
	// evaluation is consistent and lock-free. Parent bindings are immutable: a
	// key's parent must not change across updates. Default nil disables chaining
	// (every ad is standalone) and keeps the plain single-key routing/scan.
	ParentKeyFor func(key []byte) []byte

	// IsStructural, if set, marks a key as a structural (parent-only) ad that
	// exists solely to be chained to -- e.g. a job cluster ad. Structural ads are
	// stored and used as parents but are hidden from Query/Scan/Watch results by
	// default (like condor_q omitting cluster ads). Default nil treats every ad as
	// a normal, visible ad. Only meaningful together with ParentKeyFor.
	IsStructural func(key []byte) bool

	// ParentPrivateAttrs are parent attributes children do NOT inherit and that do
	// NOT trigger a watch fan-out when they change -- e.g. a job cluster ad's
	// factory bookkeeping (JobMaterializeNextProcId, ...), which mutates per proc
	// materialized but is meaningless to a proc. They are excluded from flattening
	// and from the parent-change diff, so high-frequency parent-internal churn does
	// not fan out to every child. Case-insensitive. Only meaningful with ParentKeyFor.
	ParentPrivateAttrs []string

	// DataKey enables encryption at rest: the named attributes' values are sealed with
	// AES-256-GCM under this key before being written to a segment. It must be the DB
	// master key's DataInfo subkey (crypt.Subkey(master, crypt.DataInfo)) -- distinct
	// from the master itself -- so a stolen database is useless without a pool key.
	// Empty ⇒ encryption disabled (EncryptedAttrs is then inert). See encrypt.go.
	DataKey []byte
	// EncryptedAttrs names the attributes to encrypt at rest (case-insensitive). Only
	// meaningful with DataKey set. An encrypted attribute may NOT also be indexed
	// (CategoricalAttrs/ValueAttrs) -- New panics on overlap -- since its value is
	// opaque at rest. The set can be changed at runtime via the toggle meta-command;
	// existing records keep their prior form until rewritten.
	EncryptedAttrs []string

	// TimeTravel, if non-nil, enables point-in-time ("AS OF") queries: superseded
	// record versions committed within MaxDistance of now are retained (not reclaimed
	// by compaction), and a sparse time->seq checkpoint is recorded so a wall-clock
	// time resolves to the commit sequence that was current then. Nil (default)
	// disables it with zero write-path cost -- retention collapses to the current seq,
	// exactly as before. See timeseq.go. It can be toggled at runtime.
	TimeTravel *TimeTravelOptions

	// MutatingBlockCacheBytes and ArchiveBlockCacheBytes set the PROCESS-GLOBAL shared
	// decompressed-columnar-block cache budgets, in bytes, for the two workload kinds: one budget
	// shared by ALL mutating (sharded, key-superseding) collections, and a separate one shared by ALL
	// append-only ARCHIVE collections. This replaces the old fixed 256 MB-per-collection cache, whose
	// footprint grew as 256 MB × (number of columnarized tables) with no knob. Both are global
	// ceilings, not per-collection budgets. A value ≤ 0 leaves the current/default budget
	// (defaultMutatingCacheBytes / defaultArchiveCacheBytes, 512 MiB each). Setting either here is
	// equivalent to calling Set{Mutating,Archive}BlockCacheBudget; the last non-zero setting wins and
	// resizes the live shared cache. Because the budget is global, prefer setting it once (e.g. from
	// db/ at startup) rather than per table. An archive collection uses the archive budget; every
	// other collection uses the mutating budget.
	MutatingBlockCacheBytes int64
	ArchiveBlockCacheBytes  int64
}

type OrderCursor added in v0.5.2

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

OrderCursor marks a position in an ordered scan. Its zero value starts at the beginning of a partition; the cursor yielded alongside an ad resumes strictly after that ad, so a negotiator can iterate a partition across several calls.

type OrderSpec added in v0.5.2

type OrderSpec struct {
	Partition string    // attribute partitioning the index; empty = one global partition
	Where     string    // membership predicate; empty = every ad is a member
	Keys      []SortKey // sort order within a partition
	// Cluster names the attributes whose combined value forms a member's cluster
	// signature -- a hash surfaced alongside each ad by Ordered. The schedd's RRL
	// fold groups the ordered stream into runs of equal signature (a stored-value
	// compare instead of re-hashing the requirement attributes each cycle). Computed
	// once on the write path. Empty = no signature (Ordered reports 0). Optional.
	Cluster []string
}

OrderSpec configures a maintained, filtered, ordered index over a collection -- the schedd's priority-queue pattern (§3 of docs/MATCH.md). Members are grouped into independent runs by Partition (e.g. per Owner); within a partition they are ordered by Keys, with insertion order as the final tiebreaker ("the order the job entered the queue"). Where restricts membership to the ads that matter (e.g. idle jobs), so the index tracks only the churn-prone subset.

type OrderedAd added in v0.5.2

type OrderedAd struct {
	Ad        *classad.ClassAd
	Cursor    OrderCursor
	Signature uint64
}

OrderedAd is one step of an ordered scan: the member ad, a cursor that resumes right after it, and the cluster signature (0 unless OrderSpec.Cluster is set). The signature lets an app run-length-fold the stream into RRLs with a stored-value compare instead of re-hashing each ad's clustering attributes.

type PrefixDecompressor added in v0.16.4

type PrefixDecompressor interface {
	DecompressPrefix(dst, src []byte, want int) ([]byte, error)
}

PrefixDecompressor is an optional Codec capability: decompress only (approximately) the first want bytes of a record. With the hot region encoded as a physical prefix (see wire.EncodeInlineWithHotEnc), a hot-covered projected read decompresses a couple of KB instead of the whole record. The result may be shorter than want (the record ends first -- then it is the complete record) and is otherwise a TRUNCATED record: the caller must be prepared for parses to run off the end and fall back to a full Decompress.

type ProbeExplain added in v0.6.0

type ProbeExplain struct {
	Attr    string `json:"attr"`
	Op      string `json:"op"`
	Indexed bool   `json:"indexed"`        // this probe can use an index to prune
	Kind    string `json:"kind,omitempty"` // "categorical" | "value" | "" (attr not indexed)

	// Selectivity and EstCandidates are the planner's estimate of how much this
	// probe prunes, present only when Indexed: EstCandidates is the estimated
	// number of ads the index would visit for this probe (summed over the segment
	// indexes' selectivity stats), and Selectivity is that as a fraction of the
	// total (lower is more selective). HasSelectivity distinguishes a real 0 from
	// "not estimated".
	HasSelectivity bool    `json:"hasSelectivity"`
	Selectivity    float64 `json:"selectivity,omitempty"`
	EstCandidates  int64   `json:"estCandidates,omitempty"`

	// Coalesced marks a conjunct the planner folded into another probe on the same value
	// attribute: the opposite bound (`Memory > 1024 && Memory < 4096` becomes one two-sided
	// probe), a tighter same-side bound that subsumes it (`Memory > 1024 && Memory > 512`
	// probes only `> 1024`), or a range an equality absorbed (`Memory == 2048 && Memory >
	// 1024` probes only the equality). Selectivity/EstCandidates then describe the folded
	// probe, not this conjunct's own reach; the conjunct that IS the probe is left unmarked.
	Coalesced bool `json:"coalesced,omitempty"`
}

planIndex matches the query's probes against the configured indexes. Empty means no index-usable constraint (the store full-scans). ProbeExplain describes how one index-satisfiable conjunct of a query relates to the configured indexes.

type QueryExplain added in v0.6.0

type QueryExplain struct {
	// Native reports wire-native evaluation: the query reads scalar-literal
	// attributes directly from the encoded ads, building no ClassAd per ad.
	Native bool `json:"native"`
	// Probes are the query's index-satisfiable conjuncts and their index status.
	Probes []ProbeExplain `json:"probes"`
	// IndexUsable is how many probes can prune via an index.
	IndexUsable int `json:"indexUsable"`
	// Plan is the chosen access path: "indexed" (visit index candidates),
	// "parallel-scan", "serial-scan" (full scan), or "empty" (a contradictory conjunct
	// makes the query unsatisfiable, so no records are visited at all).
	Plan string `json:"plan"`
	// Parallelism is the configured per-query worker cap; Shards is the shard count.
	Parallelism int `json:"parallelism"`
	Shards      int `json:"shards"`
	// TotalAds is the live ad count, the denominator for probe selectivity.
	TotalAds int `json:"totalAds"`
}

QueryExplain is a description of how the store would execute a query, for the diagnostic ".explain" command -- what the planner sees and the access path it would choose.

type RankedMatch added in v0.6.0

type RankedMatch struct {
	Ad      *classad.ClassAd
	Rank    float64
	HasRank bool
}

RankedMatch is a matched ad and the job's Rank of it (HasRank is false when the job's Rank does not evaluate to a number for that ad).

type RawAd

type RawAd struct {
	Exprs      [][]byte
	MyType     string
	TargetType string
}

RawAd is a query result rendered as old-ClassAd wire parts -- the "Name = Value" expression byte slices plus MyType/TargetType -- decoded straight from the stored form with no ast/classad ClassAd. A collector can hand Exprs to message.PutClassAdRawBytes to stream a result set without ever building an AST.

Exprs (and their backing bytes) alias a buffer reused across the iteration, so a consumer must finish with one RawAd -- e.g. write it to the wire -- before advancing the iterator to the next.

type Retention

type Retention struct {
	MaxSegments int   // keep at most this many sealed segments
	MaxBytes    int64 // keep at most this many bytes of sealed segment files
	// MaxAgeAttr names the numeric attribute (e.g. "CompletionDate", unix seconds) that
	// MaxAge is measured against. It must be a ZoneAttr so its per-segment max is available.
	// Now is supplied per Rotate call so the store needs no clock of its own.
	MaxAgeAttr string
	// MaxAge is a hard ceiling: drop a segment whose newest MaxAgeAttr value is older than
	// Now-MaxAge, regardless of anything else (a slow or absent consumer never keeps data
	// past it). Zero ⇒ no age ceiling.
	MaxAge float64
	// MinAgeAttr names the numeric attribute the minimum-retention floor and the GC floor are
	// measured against (a ZoneAttr). It is independent of MaxAgeAttr -- MinAge works with no
	// MaxAge configured -- though a change-feed queue typically points both at the same
	// entry-time attribute.
	MinAgeAttr string
	// MinAge is the minimum-retention floor, for using the store as a short-lived queue. The
	// GC floor (SetGCFloor) may reclaim already-consumed records early, but never any younger
	// than Now-MinAge -- so a consumer always has at least MinAge to drain a record. MinAge
	// only bounds the GC floor; it never overrides the hard ceilings, so MaxBytes / MaxSegments
	// / MaxAge still evict even younger-than-MinAge data under pressure. Zero ⇒ consumed
	// records may be GC'd as soon as the floor passes them.
	MinAge float64
}

Retention bounds how much history an append-only store keeps. Rotate drops the oldest whole segments until every set bound is met. A zero field means "no bound on that axis". Shared with Collection (Options.Retention / Collection.Rotate).

type ScanStats added in v0.29.0

type ScanStats struct {
	SegmentsTotal   int // segments in the snapshot
	SegmentsPruned  int // dropped by zone maps before any record was read
	SegmentsScanned int // segments actually walked

	RecordsVisited       int // records the scan iterated
	RecordsColumnDecided int // WHERE answered from the columnar prefilter (no reassembly to test it)
	RecordsReassembled   int // record rebuilt from the arena (the ~O(ad width) path)
	RowsMatched          int // records that satisfied the WHERE and were emitted
}

ScanStats accumulates, for one query scan, how the work broke down -- so a diagnostic (EXPLAIN ANALYZE) can show WHERE a slow scan spent itself rather than leaving it to be guessed. All counts are for a single scan; pass a fresh &ScanStats to a *Stats query method.

The distinction that matters is RecordsReassembled vs RecordsColumnDecided: a scan that decides the WHERE from the columnar prefilter and emits only matches touches the record arena RowsMatched times; a scan that reassembles most of what it visits is the slow shape (the record is rebuilt from the arena -- ~O(ad width) each -- even for non-matches).

type SchemaFieldFit added in v0.25.2

type SchemaFieldFit struct {
	Name string `json:"name"`
	Kind string `json:"kind"`
	// Width is the fixed slot's size in bytes (0 for a bit-packed bool).
	Width int  `json:"width,omitempty"`
	Hot   bool `json:"hot,omitempty"`
	// Escaped is the fraction of sampled records whose value was not in the fixed slot, for
	// any reason. This is the number that matters for speed.
	Escaped float64 `json:"escaped"`
	// Missing is the fraction absent from the record entirely -- an escape that no width or
	// kind change would fix. Escaped-Missing is the part a re-schema could recover.
	Missing float64 `json:"missing"`
}

SchemaFieldFit reports how well one schema field still fits the data, measured against a fresh sample rather than the one the schema was built from.

A field escapes when its value is not in the fixed slot -- either the attribute is absent (Missing) or it is present but unstorable in the slot: a different kind, or an int too wide for the width chosen (Escaped - Missing). Escapes are what the schema exists to avoid: a queried column that escapes forces the block's cold-stream decompression and a cold-tail walk, which is the slow path the accelerator was meant to skip.

Rates are fractions of Sampled (see SchemaFit), in [0,1].

type SchemaScanField added in v0.25.2

type SchemaScanField struct {
	Name     string `json:"name"`
	Kind     string `json:"kind"` // bool, int, real, string
	Width    int    `json:"width,omitempty"`
	Unsigned bool   `json:"unsigned,omitempty"`
	Hot      bool   `json:"hot,omitempty"`
}

SchemaScanField is one attribute in the derived schema: the name it was recovered under, the storable kind chosen for it, and the fixed width its slot occupies. Hot marks the numeric columns kept uncompressed for O(1) scan.

Width is bytes for an int (1/2/4/6/8, the narrowest that fits the sampled values), 8 for a real or a string slot, and 0 for a bool (bit-packed).

type SchemaScanGroup added in v0.29.13

type SchemaScanGroup struct {
	Fields []SchemaScanField `json:"fields"`
}

SchemaScanGroup is one committed secondary (group) schema: the co-occurring attributes stored columnar together, field by field in layout order. Its members are what `.schema groups` shows at READ; the candidate report from a fresh sample is a separate, DAEMON-level derivation.

type SchemaScanInfo added in v0.23.1

type SchemaScanInfo struct {
	Enabled         bool     `json:"enabled"`
	HotFields       []string `json:"hotFields,omitempty"`
	SchemaFields    int      `json:"schemaFields,omitempty"`
	SealedSegments  int      `json:"sealedSegments,omitempty"`
	CoveredSegments int      `json:"coveredSegments,omitempty"`
	// GroupSchemas is how many SECONDARY schemas are built alongside the base one, and
	// GroupSchemaFields their total field count. Zero when the feature is off, and also when it
	// is on but no group has yet kept its members together long enough to be committed to
	// storage -- which is the normal state for the first few maintenance passes, and the one
	// worth being able to tell apart from "not configured".
	GroupSchemas      int `json:"groupSchemas,omitempty"`
	GroupSchemaFields int `json:"groupSchemaFields,omitempty"`
	// Schema is the derived schema itself, field by field, in layout order. The counts above
	// say how much of the table the accelerator covers; this says what it decided the ads
	// look like -- which is what you need to judge whether the sampling recovered the shape
	// you expected, or picked up something odd.
	Schema []SchemaScanField `json:"schema,omitempty"`
	// Groups is the committed secondary schemas, field by field (the members that GroupSchemas /
	// GroupSchemaFields count). It is read from the live committed state -- no sampling -- so it is
	// available to a READ-level client, unlike GroupSchemas(sampleMax, k), which derives fresh
	// CANDIDATE groups from a new sample. Empty when no group is committed.
	Groups []SchemaScanGroup `json:"groups,omitempty"`
}

SchemaScanInfo reports the state of the per-segment columnar (adschema) accelerator, for diagnostics. Enabled means a numeric COUNT(*) WHERE can take the columnar fast path; HotFields are the demand-hot numeric columns kept uncompressed for O(1) scan; SchemaFields is the schema's field count; and CoveredSegments of SealedSegments carry a columnar block (the rest row-fall-back until a rewrite/reindex reaches them).

type SeqCursor added in v0.29.5

type SeqCursor struct {
	Shard    uint32
	Snapshot uint64
	Seq      uint64
	Key      string
}

SeqCursor marks a position in a sequence-ordered scan: the last record the caller consumed. The zero value starts at the beginning.

It is per shard, because the commit sequence is: each shard stamps records from its own counter (`seq := sh.commitSeq + 1`), so sequences from different shards are unrelated numbers and only order records within one shard. A scan therefore finishes a shard before moving to the next, and the cursor names which shard it is in — the same shape as Scan's guarantee, which is likewise per shard ("each key present at the moment A SHARD's scan begins").

Snapshot is that shard's sequence when its first page was taken; carrying it forward is what keeps later pages from seeing writes that landed meanwhile. It is only meaningful together with Shard.

Key breaks ties within a commit, where several records share a sequence. It is compared bytewise, and only against records with the same sequence, so it imposes no ordering of its own beyond making the position unambiguous.

type SeqPage added in v0.29.5

type SeqPage struct {
	// Next is where a following page resumes: hand it back unchanged. It
	// carries the shard and that shard's snapshot as well as the position.
	// Meaningful only when More is true.
	Next SeqCursor
	// More reports whether the scan stopped at the page limit with records
	// still to come, as opposed to having reached the end of the collection.
	More bool
}

SeqPage is one page of a sequence-ordered scan.

type SidecarSizes added in v0.8.0

type SidecarSizes struct {
	Segments    int   `json:"segments"`    // sealed segments with a sidecar
	MappedBytes int64 `json:"mappedBytes"` // total sidecar bytes (mmap-backed, evictable)
	MPHBytes    int64 `json:"mphBytes"`    // of MappedBytes, minimal-perfect-hash structures
	BloomBytes  int64 `json:"bloomBytes"`  // of MappedBytes, bloom filters
}

SidecarSizes reports the on-disk index bytes of an Archive's sealed sidecars, broken out so the minimal-perfect-hash and bloom overhead is visible. These are distinct from the live Collection's IndexSizes: those measure HEAP-resident postings + sketches, whereas sidecar bytes live in the page cache -- demand-paged and evictable under memory pressure. They are reported as a separate budget, never folded into a heap figure, so the "index is N% of data" watermark stays an honest measure of resident memory.

type SortKey added in v0.5.2

type SortKey struct {
	Expr string // attribute name or expression, e.g. "JobPrio" or "QDate"
	Desc bool   // sort this component descending
}

SortKey is one component of an ordered index's sort order: an attribute name or expression evaluated against each member ad, ascending by default.

type Stats

type Stats struct {
	// Ads is the number of live keys (== Len).
	Ads int
	// Segments is the number of arena segments across all shards.
	Segments int
	// ArenaBytes is the total capacity of those segments -- the memory the
	// collection has reserved for record storage, and the dominant term in its
	// resident footprint. (RAM segments back this on the Go heap; persistent
	// segments back it with an mmap.)
	ArenaBytes int64
	// UsedBytes is the bytes written into segments: live records plus superseded
	// (dead) records not yet reclaimed by compaction.
	UsedBytes int64
	// DeadBytes is the bytes belonging to superseded records -- reclaimable by
	// Compact. LiveBytes = UsedBytes - DeadBytes.
	DeadBytes int64
}

Stats is a point-in-time summary of a collection's storage, for observability and capacity planning (e.g. Prometheus metrics). All byte counts are of the dictionary-compressed on-arena form, not the ads' decompressed text.

func (Stats) LiveBytes

func (s Stats) LiveBytes() int64

LiveBytes is the compressed size of the live records (UsedBytes - DeadBytes).

type TimeTravelOptions added in v0.12.0

type TimeTravelOptions struct {
	// MaxDistance is how far back queries may travel: versions superseded more than
	// this long ago are eligible for reclamation. Larger keeps more history (more
	// retained dead bytes). Must be > 0.
	MaxDistance time.Duration
	// CheckpointInterval is the granularity of the time->seq map: a checkpoint is
	// recorded at most once per interval, on the first commit that crosses the
	// boundary, so an AS OF time resolves to within one interval. Default 1 minute;
	// lower for finer resolution at proportionally more checkpoints.
	CheckpointInterval time.Duration
}

TimeTravelOptions configures point-in-time queries for a collection (see the TimeTravel option and timeseq.go).

type Txn added in v0.6.0

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

Txn is an optimistic, snapshot-isolation transaction over a Collection. Not safe for concurrent use by multiple goroutines; each goroutine uses its own Txn.

func (*Txn) Commit added in v0.6.0

func (tx *Txn) Commit() CommitResult

Commit applies the buffered writes, each independently: a write whose key is unchanged since the transaction's snapshot commits; one whose key was modified by another committer is reported in CommitResult.Conflicts and not applied (the successful writes are not rolled back). The transaction must not be used after Commit.

func (*Txn) Delete added in v0.6.0

func (tx *Txn) Delete(key []byte)

Delete buffers a delete of key. Nothing is written until Commit.

func (*Txn) Get added in v0.6.0

func (tx *Txn) Get(key []byte) (*classad.ClassAd, bool)

Get returns the ad for key as the transaction sees it: its own buffered write if any (read-your-writes), else the version live at the transaction's snapshot. On a chained (parent/child) collection it resolves inherited attributes by merging the parent as of the same snapshot -- mirroring Collection.Get, transactionally.

func (*Txn) Has added in v0.29.11

func (tx *Txn) Has(key []byte) bool

Has reports whether key exists as the transaction sees it, without reading the stored record: it resolves the key to a live record location and stops there. Get is the wrong tool for a presence question -- it copies the record bytes out from under the shard lock and decodes them, and for a wide job ad that decode is the most expensive thing in an ingest -- so a caller that only needs presence (a diagnostic counter, an upsert-vs-insert branch) should ask here.

func (*Txn) KeysWhere added in v0.24.0

func (tx *Txn) KeysWhere(q *vm.Query) iter.Seq[string]

KeysWhere returns the storage keys of the rows matching q as the transaction sees them, with the same overlay and the same caveats as Query. It is what lets an UPDATE or DELETE inside a transaction address rows the transaction itself created.

func (*Txn) PatchAttrs added in v0.30.0

func (tx *Txn) PatchAttrs(key []byte, patch *classad.ClassAd, removed []string)

PatchAttrs buffers a change to key expressed ONLY as the attributes it changes -- the whole ad is never read, built, or buffered. That is the point: a mirror applying a schedd's attribute updates spends most of its time reading an ad back just to hand a copy of it to the encoder, and the update itself does not need it.

Commit stores this as a delta record. When it cannot -- the key has no full record yet, the chain has reached its bound, or the write removed an attribute -- Commit reads the stored ad and merges there, so the read happens once per re-materialization instead of once per write.

removed names attributes the write deletes. Reads through this transaction still see the merged view (Get applies the patch over the stored ad), so read-your-writes is unaffected.

func (*Txn) Put added in v0.6.0

func (tx *Txn) Put(key []byte, ad *classad.ClassAd)

Put buffers an insert or update of key. Nothing is written until Commit.

func (*Txn) PutOld added in v0.25.1

func (tx *Txn) PutOld(key []byte, text string) bool

PutOld buffers an insert or update of key whose ad arrives as old-ClassAd text, encoding it straight to the stored wire form here rather than building an ast.ClassAd for Commit to encode. Every transactional guarantee is unchanged -- the write is buffered, conflict-checked against the same snapshot, and committed identically; only the encoding path differs.

It reports whether the fast path was taken. False means the caller should parse the text and use Put.

An encrypted collection no longer refuses this path wholesale. The streaming encoder cannot seal, so encodeOld defers to the sealing path for an ad that HAS something to seal, and streams the rest -- which for job and history ads is nearly all of them. Refusing wholesale made encryption cost every ingest.

func (*Txn) PutPatch added in v0.30.0

func (tx *Txn) PutPatch(key []byte, ad *classad.ClassAd, changed []string, removed bool)

PutPatch is Put for a caller that knows which attributes it changed. The full ad is still buffered -- read-your-writes and the value lookup both need it -- but Commit may store only the named attributes as a delta record, which for a wide ad is the difference between encoding and compressing ~7.3 KB and ~150 B. removed reports that the write also deleted an attribute, which a delta cannot express, so the whole ad is stored instead.

It is always safe to call Put instead; a caller that does loses only the optimization.

func (*Txn) Query added in v0.24.0

func (tx *Txn) Query(q *vm.Query) iter.Seq[*classad.ClassAd]

Query returns the ads matching q as the transaction sees them: the committed rows, with the transaction's own buffered writes overlaid (read-your-writes). A row the transaction deleted is absent, a row it rewrote is matched and yielded in its rewritten form, and a row it created is included if it matches -- so a query inside a transaction observes the transaction's own work, which Collection.Query cannot.

Consistency: snapshot isolation plus read-your-writes. The committed half is read at the transaction's own snapshot sequence per shard -- the same sequence Get reads at -- so a scan and a point lookup in one transaction agree, and a row another transaction commits mid-scan is invisible to both. Scanning a shard captures its snapshot if the transaction has not touched it yet, which pins every later read of that shard too.

Cost: a full scan of the committed half, because reading at a past sequence and overlaying by key both need the per-record walk that the indexed query path (which yields ads without keys, at the current sequence) does not do. A transaction that has bought nothing from either -- no buffered writes, and content to read the live store -- is better served by Collection.Query, which is what Txn with an untouched write buffer used to fall back to; that fast path is gone now that the snapshot is the point.

func (*Txn) SetDurable added in v0.6.0

func (tx *Txn) SetDurable(d bool)

SetDurable controls whether Commit runs the durability sync (default true). A nondurable commit is visible immediately (readers and watchers see it) but its disk flush is deferred to a later durable commit or flush -- the classad_log.h CommitNondurableTransaction batching. No effect on an in-memory collection, whose sync is already a no-op.

type UpgradeOptions added in v0.25.0

type UpgradeOptions struct {
	// MinGainFrac is the fraction of a segment's bytes a re-encode must be projected to
	// save before it is worth doing. Default defaultMinGainFrac.
	MinGainFrac float64
	// SampleRecords is how many records are decompressed and re-compressed to project the
	// gain. The estimate only has to separate "clearly worth it" from "not", so this stays
	// small; the cost of being wrong is one segment, not the archive.
	SampleRecords int
	// MaxBytesPerPass bounds the source bytes one pass re-encodes, as for merging. Zero
	// uses the default.
	MaxBytesPerPass int64
	// MaxSegments bounds how many segments one pass upgrades.
	MaxSegments int
}

UpgradeOptions tunes a codec-upgrade pass. The zero value is usable.

type WatchEvent added in v0.4.0

type WatchEvent struct {
	Kind   WatchKind
	Key    []byte
	Ad     *classad.ClassAd
	Cursor []byte
}

WatchEvent is one item in a Watch stream. Cursor, when non-nil, is an opaque token the client persists after processing the event and passes back to Watch on resume. Catch-up data events carry no cursor (persist at WatchSynced); live events do.

type WatchKind added in v0.4.0

type WatchKind uint8

WatchKind is the type of a WatchEvent.

const (
	// WatchUpsert carries the full ad for an added or updated key (Ad is set).
	WatchUpsert WatchKind = iota
	// WatchDelete signals a key was removed (Ad is nil).
	WatchDelete
	// WatchReset tells the client to discard its state (build into a shadow): an
	// authoritative full snapshot of Upserts follows, ending at WatchSynced. Emitted
	// when a precise incremental resume is impossible (first subscribe, cursor older
	// than the delete-retention window, or a different store generation).
	WatchReset
	// WatchSynced marks the end of the initial catch-up/snapshot: the client is now
	// live. Its Cursor is a durable resume point (and, after a Reset, the point to
	// swap the shadow state live).
	WatchSynced
	// WatchResync tells the client the live stream fell behind and it must reconnect
	// with its last persisted cursor (which re-enters catch-up). No state is implied.
	WatchResync
)

Directories

Path Synopsis
Package crypt is the encryption-at-rest key hierarchy for a ClassAd database.
Package crypt is the encryption-at-rest key hierarchy for a ClassAd database.
Package vm compiles ClassAd expressions to a linear instruction stream and interprets them against a scope, replacing per-query AST walks.
Package vm compiles ClassAd expressions to a linear instruction stream and interprets them against a scope, replacing per-query AST walks.
Package wire defines a compact TLV binary form of a ClassAd with attribute-key interning and a hot-attribute header, plus the encoder/decoder that convert to and from the fully-public ast.ClassAd representation.
Package wire defines a compact TLV binary form of a ClassAd with attribute-key interning and a hot-attribute header, plus the encoder/decoder that convert to and from the fully-public ast.ClassAd representation.

Jump to

Keyboard shortcuts

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