Documentation
¶
Overview ¶
Package db is an embedded ClassAd log: a persistent key->ClassAd store with optimistic multi-writer transactions, mirroring HTCondor's ClassAdLog (src/condor_utils/classad_log.h). It is the Go core that the cgo layer (package capi) exposes as C symbols for a C++ interface to sit on top of, and that the client/server module serves over CEDAR.
It maps directly onto the collections store: the key->ClassAd table is a Collection, and each transaction is a collections.Txn (snapshot isolation, write-write conflicts, per-ad commit). Unlike classad_log.h -- which allows only one active transaction -- this supports any number of independent concurrent transactions, each a distinct *Txn.
Index ¶
- Constants
- Variables
- func AggFilterAttrs(filter string) []string
- func AggProjection(groupCols []GroupCol, aggs []AggSpec) (attrs []string, groupCol, aggCol []int)
- func CompactInternPaths() (wirePath, astPath int64)
- func ConstraintRefs(expr string) (refs []string, dynamic bool)
- func FallbackReasons() (removal, bound, noBase, ineligible int64)
- func IsMatchAll(constraint string) bool
- func IsMatchNone(constraint string) bool
- func IsSystemKey(key string) bool
- func LastDeltaDecodeFailure() (stage, msg string)
- func LastNoBaseDetail() (versions int, chainBroken bool, sealedSkipped int)
- func MatchSignature(ad *classad.ClassAd, significantAttrs []string) uint64
- func PrivateConstraintRef(expr string) (attr string, dynamic bool)
- func SealWalkStats() (examined, deltas int64)
- func SealedProbesSkipped() int64
- func SpliceStats() (spliced, spliceRefused, objectPath int64)
- func SystemKey(name string) string
- func UnreadableBaseReasonNames() []string
- func UnreadableBaseReasons() map[string]int64
- func UnreadableBaseRefusals() int64
- func ValidTableName(name string) bool
- func ValueText(v classad.Value) string
- type AggFunc
- type AggRow
- func AggregateValues(seq iter.Seq[[]classad.Value], attrs []string, groupCols []GroupCol, ...) ([]AggRow, error)
- func ColumnarAggregate(stats func(attr string) (NumStats, bool), groupCols []GroupCol, aggs []AggSpec) ([]AggRow, bool)
- func GroupedFromColumns(src GroupStatsSource, constraint string, groupCols []GroupCol, aggs []AggSpec) ([]AggRow, bool)
- type AggSpec
- type ArchiveConfig
- type ArchiveTable
- func (t *ArchiveTable) AddIndex(categorical, value []string) bool
- func (t *ArchiveTable) Aggregate(constraint string, groupBy []string, aggs []AggSpec) ([]AggRow, error)
- func (t *ArchiveTable) AggregateCols(constraint string, groupCols []GroupCol, aggs []AggSpec) ([]AggRow, error)
- func (t *ArchiveTable) AggregateColsStats(constraint string, groupCols []GroupCol, aggs []AggSpec, ...) ([]AggRow, error)
- func (t *ArchiveTable) Append(ad *classad.ClassAd) error
- func (t *ArchiveTable) AppendOld(text string) error
- func (t *ArchiveTable) AutoIndexNames() []string
- func (t *ArchiveTable) AutoTune(opts AutoTuneOptions) AutoTuneResult
- func (t *ArchiveTable) BuildAndEnableSchemaScan(sampleMax, hotTopN int) bool
- func (t *ArchiveTable) CategoricalGroupCounts(attr string) (map[string]int64, bool)
- func (t *ArchiveTable) CategoricalGroupCountsBucketed(attr, bucketAttr string, width int64) (map[int64]map[string]int64, bool)
- func (t *ArchiveTable) CategoricalGroupCountsWhere(attr, constraint string) (map[string]int64, bool)
- func (t *ArchiveTable) Close() error
- func (t *ArchiveTable) CodecStats(sampleMax int) CodecStats
- func (t *ArchiveTable) Count() int
- func (t *ArchiveTable) CountConstraint(constraint string) (int, bool)
- func (t *ArchiveTable) DropIndex(names ...string) bool
- func (t *ArchiveTable) EncryptedAttrNames() []string
- func (t *ArchiveTable) EncryptionEnabled() bool
- func (t *ArchiveTable) Explain(constraint string) (QueryExplain, error)
- func (t *ArchiveTable) GCFloor() float64
- func (t *ArchiveTable) GroupCountAll(groupAttr string) ([]collections.GroupCount, bool)
- func (t *ArchiveTable) GroupCountConstraint(constraint, groupAttr string) ([]collections.GroupCount, bool)
- func (t *ArchiveTable) GroupSchemaAgreement(sampleMax, k int) GroupSchemaAgreement
- func (t *ArchiveTable) GroupSchemaChanges() []GroupSchemaChange
- func (t *ArchiveTable) GroupSchemaDrift() GroupSchemaDrift
- func (t *ArchiveTable) GroupSchemaLastAgreement() (GroupSchemaLastAgreement, bool)
- func (t *ArchiveTable) GroupSchemas(sampleMax, k int) GroupSchemaInfo
- func (t *ArchiveTable) GroupStatsAll(groupAttr string, aggAttrs []string) ([]collections.GroupStats, bool)
- func (t *ArchiveTable) GroupStatsConstraint(constraint, groupAttr string, aggAttrs []string) ([]collections.GroupStats, bool)
- func (t *ArchiveTable) HotAttrs() []string
- func (t *ArchiveTable) IndexSizes() IndexSizes
- func (t *ArchiveTable) IndexedAttrs() (categorical, value []string)
- func (t *ArchiveTable) MergePass(opts MergeOptions) int
- func (t *ArchiveTable) NumStats(constraint, attr string) (NumStats, bool)
- func (t *ArchiveTable) OpStats() OpStats
- func (t *ArchiveTable) Query(constraint string) (iter.Seq[*classad.ClassAd], error)
- func (t *ArchiveTable) QueryLimit(constraint string, limit int) (iter.Seq[*classad.ClassAd], error)
- func (t *ArchiveTable) QueryProject(constraint string, attrs []string) (iter.Seq[[]classad.Value], error)
- func (t *ArchiveTable) QueryProjectStats(constraint string, attrs []string, stats *collections.ScanStats) (iter.Seq[[]classad.Value], error)
- func (t *ArchiveTable) QueryRawProjected(constraint string, projection []string, redact bool) (iter.Seq[collections.RawAd], error)
- func (t *ArchiveTable) QueryRawProjectedRefs(constraint string, projection []string, redact bool) (iter.Seq[collections.RawAd], error)
- func (t *ArchiveTable) QueryRawProjectedRefsStats(constraint string, projection []string, redact bool, ...) (iter.Seq[collections.RawAd], error)
- func (t *ArchiveTable) QueryRawProjectedStats(constraint string, projection []string, redact bool, ...) (iter.Seq[collections.RawAd], error)
- func (t *ArchiveTable) Reindex()
- func (t *ArchiveTable) ReschemaScan(sampleMax, hotTopN int) bool
- func (t *ArchiveTable) Retention() collections.Retention
- func (t *ArchiveTable) RetrainDict(sampleMax int) (int, error)
- func (t *ArchiveTable) Rewrite() int
- func (t *ArchiveTable) Rotate(now float64) (int, error)
- func (t *ArchiveTable) RowGroupBytes() int
- func (t *ArchiveTable) SaveDemand()
- func (t *ArchiveTable) SchemaFit(sampleMax int) ([]SchemaFieldFit, int)
- func (t *ArchiveTable) SchemaScanInfo() SchemaScanInfo
- func (t *ArchiveTable) SetGCFloor(floor float64)
- func (t *ArchiveTable) SetRetention(r collections.Retention) error
- func (t *ArchiveTable) SetRowGroupBytes(n int) error
- func (t *ArchiveTable) SidecarSizes() SidecarSizes
- func (t *ArchiveTable) StaleIndexSegments() (stale, sealed int)
- func (t *ArchiveTable) StaleIndexSegmentsByPolicy() int
- func (t *ArchiveTable) Stats() Stats
- func (a *ArchiveTable) StopMaintenance()
- func (t *ArchiveTable) TopK(constraint string, attrs []string, orderAttr string, desc bool, k int) ([][]classad.Value, error)
- func (t *ArchiveTable) TopKStats(constraint string, attrs []string, orderAttr string, desc bool, k int) ([][]classad.Value, collections.ScanStats, error)
- func (t *ArchiveTable) Truncate()
- func (t *ArchiveTable) Watch(ctx context.Context, cursor []byte) (iter.Seq[WatchEvent], error)
- func (t *ArchiveTable) WatchCursor() ([]byte, error)
- func (t *ArchiveTable) ZoneAttrs() []string
- type AutoTuneOptions
- type AutoTuneResult
- type Catalog
- func (cat *Catalog) ArchiveTable(name string) (*ArchiveTable, bool)
- func (cat *Catalog) ArchiveTables() []string
- func (cat *Catalog) Close() error
- func (cat *Catalog) ConvertTableToMemory(name string) error
- func (cat *Catalog) CreateArchiveTable(name string, cfg ArchiveConfig) (*ArchiveTable, error)
- func (cat *Catalog) CreateExporter(def ExporterDef) error
- func (cat *Catalog) CreateTable(name string) (*DB, error)
- func (cat *Catalog) CreateTableInMemory(name string) (*DB, error)
- func (cat *Catalog) CreateTableOpts(name string, opts TableOptions) (*DB, error)
- func (cat *Catalog) CreateView(name string, spec ViewSpec) error
- func (cat *Catalog) DropArchiveTable(name string) error
- func (cat *Catalog) DropExporter(name string) error
- func (cat *Catalog) DropTable(name string) error
- func (cat *Catalog) DropView(name string) error
- func (cat *Catalog) EnsureTable(name string) (*DB, error)
- func (cat *Catalog) Exporter(name string) (ExporterDef, bool)
- func (cat *Catalog) Exporters() []ExporterDef
- func (cat *Catalog) LoadExporterState(name string) ([]byte, bool, error)
- func (cat *Catalog) Restore(r io.Reader) error
- func (cat *Catalog) SaveExporterState(name string, state []byte) error
- func (cat *Catalog) Snapshot(w io.Writer) error
- func (cat *Catalog) Table(name string) (*DB, bool)
- func (cat *Catalog) Tables() []string
- func (cat *Catalog) View(name string) (*View, bool)
- func (cat *Catalog) ViewBacking(name string) (*DB, bool)
- func (cat *Catalog) ViewSealed(name, constraint string) (seq iter.Seq[*classad.ClassAd], ok bool, err error)
- func (cat *Catalog) Views() []string
- type CatalogConfig
- type CodecStats
- type Config
- type ConflictError
- type Constraint
- type DB
- func (db *DB) AddHotAttrs(names ...string) []string
- func (db *DB) AddIndex(categorical, value []string) bool
- func (db *DB) BackupKey() []byte
- func (db *DB) Begin() *Txn
- func (db *DB) BeginRedacted() *Txn
- func (db *DB) Chained() bool
- func (db *DB) Close() error
- func (db *DB) CodecStats(sampleMax int) CodecStats
- func (db *DB) Compact() int
- func (db *DB) CountConstraint(constraint string) (int, bool)
- func (db *DB) Delete(key string) (bool, error)
- func (db *DB) DeleteWhere(constraint string) (int, error)
- func (db *DB) DeltaStats() (deltas, fulls int64)
- func (db *DB) DropIndex(names ...string) bool
- func (db *DB) EnableSchemaScan(sampleMax, hotTopN int) bool
- func (db *DB) EncryptedAttrNames() []string
- func (db *DB) EncryptionEnabled() bool
- func (db *DB) Explain(constraint string) (QueryExplain, error)
- func (db *DB) ExplainMatch(job *classad.ClassAd, targetConstraint string) MatchExplain
- func (db *DB) ForEach(fn func(ad *classad.ClassAd) bool)
- func (db *DB) ForEachSystemAd(fn func(key string, ad *classad.ClassAd) bool)
- func (db *DB) GroupSchemaAgreement(sampleMax, k int) GroupSchemaAgreement
- func (db *DB) GroupSchemaChanges() []GroupSchemaChange
- func (db *DB) GroupSchemaDrift() GroupSchemaDrift
- func (db *DB) GroupSchemaLastAgreement() (GroupSchemaLastAgreement, bool)
- func (db *DB) GroupSchemas(sampleMax, k int) GroupSchemaInfo
- func (db *DB) GroupStatsAll(groupAttr string, aggAttrs []string) ([]collections.GroupStats, bool)
- func (db *DB) GroupStatsConstraint(constraint, groupAttr string, aggAttrs []string) ([]collections.GroupStats, bool)
- func (db *DB) HotAttrs() []string
- func (db *DB) ID() string
- func (db *DB) InMemory() bool
- func (db *DB) IndexSizes() IndexSizes
- func (db *DB) IndexedAttrs() (categorical, value []string)
- func (db *DB) InstanceID() string
- func (db *DB) Keys() []string
- func (db *DB) KeysWhere(constraint string) (iter.Seq[string], error)
- func (db *DB) Len() int
- func (db *DB) LookupClassAd(key string) (*classad.ClassAd, bool)
- func (db *DB) LookupClassAdRedacted(key string) (*classad.ClassAd, bool)
- func (db *DB) Maintain(opts MaintainOptions)
- func (db *DB) Match(job *classad.ClassAd) iter.Seq[*classad.ClassAd]
- func (db *DB) MatchSorted(job *classad.ClassAd, limit int) []*classad.ClassAd
- func (db *DB) MatchSortedRanked(job *classad.ClassAd, limit int) []RankedMatch
- func (db *DB) MatchSortedRankedFiltered(job *classad.ClassAd, targetConstraint string, limit int) ([]RankedMatch, error)
- func (db *DB) NumStats(constraint, attr string) (NumStats, bool)
- func (db *DB) OpStats() OpStats
- func (db *DB) Ordered(index int, partition string, resume OrderCursor) iter.Seq[OrderedAd]
- func (db *DB) OrderedRedacted(index int, partition string, resume OrderCursor) iter.Seq[OrderedAd]
- func (db *DB) Put(key string, ad *classad.ClassAd) error
- func (db *DB) Query(constraint string) (iter.Seq[*classad.ClassAd], error)
- func (db *DB) QueryAsOf(constraint string, t time.Time) (iter.Seq[*classad.ClassAd], error)
- func (db *DB) QueryProject(constraint string, attrs []string) (iter.Seq[[]classad.Value], error)
- func (db *DB) QueryRaw(constraint string) (iter.Seq[collections.RawAd], error)
- func (db *DB) QueryRawProjected(constraint string, projection []string, redact bool) (iter.Seq[collections.RawAd], error)
- func (db *DB) QueryRawProjectedFromSeq(constraint string, projection []string, after SeqCursor, limit int) (iter.Seq[collections.RawAd], *SeqPage, error)
- func (db *DB) QueryRawProjectedRefs(constraint string, projection []string, redact bool) (iter.Seq[collections.RawAd], error)
- func (db *DB) QueryRawProjectedRefsStats(constraint string, projection []string, redact bool, ...) (iter.Seq[collections.RawAd], error)
- func (db *DB) QueryRawRedacted(constraint string) (iter.Seq[collections.RawAd], error)
- func (db *DB) QueryRawWire(constraint string, projection []string, redact bool) (iter.Seq[[]byte], error)
- func (db *DB) QueryRedacted(constraint string) (iter.Seq[*classad.ClassAd], error)
- func (db *DB) RecordDemand(constraint string)
- func (db *DB) RefreshHotSet(sampleMax, topN int) int
- func (db *DB) Reindex()
- func (db *DB) ReschemaScan(sampleMax, hotTopN int) bool
- func (db *DB) Restore(r io.Reader) error
- func (db *DB) RestoreWith(r io.Reader, keys SnapshotKeys) error
- func (db *DB) RestoreWithBackupKey(r io.Reader, backupKey []byte) error
- func (db *DB) RetrainDict(sampleMax int) (int, error)
- func (db *DB) Rewrite() int
- func (db *DB) SchemaFit(sampleMax int) ([]SchemaFieldFit, int)
- func (db *DB) SchemaScanInfo() SchemaScanInfo
- func (db *DB) SetEncryptedAttrs(attrs []string) error
- func (db *DB) SetTimeTravel(maxDistance, checkpoint time.Duration)
- func (db *DB) SidecarSizes() SidecarSizes
- func (db *DB) Snapshot(w io.Writer) error
- func (db *DB) SnapshotWithKey(w io.Writer) ([]byte, error)
- func (db *DB) StaleIndexSegments() (stale, sealed int)
- func (db *DB) StartMaintenance(interval time.Duration, opts MaintainOptions) (stop func())
- func (db *DB) Stats() Stats
- func (d *DB) StopMaintenance()
- func (db *DB) SuggestDrops(sampleMax int) []DropSuggestion
- func (db *DB) SuggestIndexes(sampleMax int) []collections.IndexSuggestion
- func (db *DB) TimeTravel() (maxDistance, checkpoint time.Duration, enabled bool)
- func (db *DB) TopK(constraint string, attrs []string, orderAttr string, desc bool, k int) ([][]classad.Value, error)
- func (db *DB) TopKStats(constraint string, attrs []string, orderAttr string, desc bool, k int) ([][]classad.Value, collections.ScanStats, error)
- func (db *DB) TrackedKeys() int
- func (db *DB) Truncate()
- func (db *DB) UpdateOld(key, text string) error
- func (db *DB) UpdateOldBatch(items []OldAdText) error
- func (db *DB) Watch(ctx context.Context, cursor []byte) (iter.Seq[WatchEvent], error)
- func (db *DB) WatchCursor() ([]byte, error)
- func (db *DB) WatchRedacted(ctx context.Context, cursor []byte) (iter.Seq[WatchEvent], error)
- type DropSuggestion
- type ExporterDef
- type GroupAgreementItem
- type GroupCol
- type GroupSchemaAgreement
- type GroupSchemaChange
- type GroupSchemaDelta
- type GroupSchemaDrift
- type GroupSchemaEntry
- type GroupSchemaInfo
- type GroupSchemaLastAgreement
- type GroupStatsSource
- type IndexSize
- type IndexSizes
- type IndexSuggestion
- type KEK
- type MaintainOptions
- type MatchExplain
- type MergeOptions
- type NumStats
- type OldAdText
- type OpStat
- type OpStats
- type OrderCursor
- type OrderSpec
- type OrderedAd
- type ProbeExplain
- type QueryExplain
- type RankedMatch
- type Retention
- type SchemaFieldFit
- type SchemaScanField
- type SchemaScanGroup
- type SchemaScanInfo
- type SeqCursor
- type SeqPage
- type SidecarSizes
- type SnapshotKeys
- type SortKey
- type Stats
- type TableOptions
- type Txn
- func (t *Txn) Abort()
- func (t *Txn) Commit() error
- func (t *Txn) CommitNondurable() error
- func (t *Txn) DeleteAttribute(key, name string)
- func (t *Txn) DestroyClassAd(key string)
- func (t *Txn) Has(key string) bool
- func (t *Txn) KeysWhere(constraint string) (iter.Seq[string], error)
- func (t *Txn) LookupAttr(key, name string) (string, bool)
- func (t *Txn) LookupClassAd(key string) (*classad.ClassAd, bool)
- func (t *Txn) NewClassAd(key string, ad *classad.ClassAd)
- func (t *Txn) NewClassAdOld(key, text string) bool
- func (t *Txn) Query(constraint string) (iter.Seq[*classad.ClassAd], error)
- func (t *Txn) SetAttribute(key, name, expr string) error
- type UnappliedError
- type View
- func (v *View) Backing() *DB
- func (v *View) Cursor() []byte
- func (v *View) LateDrops() int64
- func (v *View) Seal(now int64)
- func (v *View) SealedQuery(constraint string) (iter.Seq[*classad.ClassAd], error)
- func (v *View) SeriesCount() int
- func (v *View) Spec() ViewSpec
- func (v *View) State() (ViewState, error)
- type ViewAggFunc
- type ViewGroupCol
- type ViewMetric
- type ViewSpec
- type ViewState
- type WatchEvent
- type WatchKind
- type Watcher
Constants ¶
const DefaultDeltaMax = 16
Config opens a DB with indexing and ordered-index configuration. Dir empty is in-memory; a non-empty Dir is persistent. DefaultDeltaMax is the delta-record chain bound applied when Config.DeltaMax is left at zero, which is to say: delta records are ON by default. Storing only what a write changed is the cheaper way to hold a queue that is updated far more often than it is created -- measured at -42% ingest CPU and -32% allocation replaying a real schedd's job_queue.log -- and a feature that ships switched off is a feature that rots.
It is NOT reversible for data already written. A store that has written one delta record replays them for the rest of its life (the on-disk marker says so); setting DeltaMaxOff later stops new deltas and leaves the existing ones readable, it does not convert them back.
const DeltaMaxOff = -1
DeltaMaxOff disables delta records for a table. Any negative Config.DeltaMax means off; this spelling exists because zero means "use the default", so there has to be a way to say "none".
Variables ¶
var ErrRawWireUnsupported = errors.New("classad-db: table does not support wire-form rows")
ErrRawWireUnsupported reports that a table cannot serve the wire-form relay scan -- today, that it is in-memory rather than persistent. Callers fall back to a text row stream, which every table can serve.
Functions ¶
func AggFilterAttrs ¶ added in v0.23.1
AggFilterAttrs returns the attributes a per-aggregate filter reads. Callers use it to hold a filter to the same rules as the rest of the request -- notably the RPC layer's refusal to let an unprivileged reader touch a private attribute, which a filter could otherwise turn into an oracle (`COUNT(*) FILTER (WHERE Secret == "guess")` leaks by its count).
func AggProjection ¶ added in v0.16.8
AggProjection builds the deduplicated list of attributes the aggregation reads (group columns then non-"*" aggregate arguments) and the index of each group column / aggregate argument within that list. An aggregate whose argument is "*" (COUNT(*)) gets index -1. A caller projects a scan to attrs and feeds the resulting value rows, groupCol, and aggCol to AggregateValues.
func CompactInternPaths ¶ added in v0.30.0
func CompactInternPaths() (wirePath, astPath int64)
CompactInternPaths reports, process-wide, how compaction re-interned records: by transcoding the wire bytes, or by decoding into an ast and encoding it back. All-ast means the transcode is being attempted and refused on every record.
func ConstraintRefs ¶ added in v0.28.0
ConstraintRefs reports every attribute name an expression references, scoped or not, and whether the expression makes a reference this analysis cannot resolve statically (see dynamic below).
A constraint that does not parse returns no references and dynamic=false: it is rejected downstream where the parse error can be reported properly, and refusing it here would report the wrong reason.
func FallbackReasons ¶ added in v0.30.0
func FallbackReasons() (removal, bound, noBase, ineligible int64)
FallbackReasons reports, process-wide, why patch writes had to store a whole record rather than a delta: an attribute removal (which a delta cannot express), the chain reaching its bound, no whole record to chain to yet, or delta records not being in use. Each fallback costs a read of the stored ad, so this says which of those reads are worth attacking.
func IsMatchAll ¶ added in v0.21.3
IsMatchAll reports whether a constraint imposes no filter (an empty string or a literal "true"), so an aggregate over it covers every record. Shared with the mutable-table COUNT(*) fast path in dbrpc.
func IsMatchNone ¶ added in v0.27.1
IsMatchNone reports whether a constraint can never match: it references no attribute, so its value is the same for every record, and that value is not TRUE.
The dual of IsMatchAll, and missing until `select min(ProcId) from history where false` was measured at 3.5s -- a full scan to establish that nothing matches. FALSE, UNDEFINED and ERROR all match nothing, since a record matches only when the constraint is boolean TRUE, so this folds `where false`, `where 1 == 2`, `where undefined` and `where error` alike.
The no-reference check is the same load-bearing one IsMatchAll documents, in the other direction: a record-dependent expression must never fold even when it evaluates non-TRUE against an empty ad. `JobStatus == 2` is UNDEFINED with JobStatus absent and matches nothing HERE, but matches plenty of records, so anything referencing an attribute is left to the scan.
func IsSystemKey ¶ added in v0.9.0
IsSystemKey and SystemKey re-export the collections helpers so callers holding only a *db.DB (e.g. dbrpc building marker keys) can classify or construct a reserved system key without importing collections directly. A system key begins with a NUL byte and names a record hidden from client reads but retrievable by explicit LookupClassAd.
func LastDeltaDecodeFailure ¶ added in v0.30.5
func LastDeltaDecodeFailure() (stage, msg string)
LastDeltaDecodeFailure returns the most recently sampled decode error behind a refused delta chain merge, as (stage, message) where stage is "base" or "patch". The per-reason counts say how often the bytes would not decode; this says what the decoder objected to.
func LastNoBaseDetail ¶ added in v0.30.7
LastNoBaseDetail returns the shape of the most recent delta chain that had no whole record: versions found, whether the walk truncated at a dead link, and how many sealed segments could not be probed for want of a key index.
func MatchSignature ¶
MatchSignature is HTCondor's autocluster key: a 64-bit checksum over the given significant attributes' expression text in ad. Two ads with textually identical significant attributes (same Requirements, same RequestCpus literal, ...) hash equal, so a matchmaker can compute one candidate list per distinct signature and reuse it for every identical request.
func PrivateConstraintRef ¶ added in v0.28.0
PrivateConstraintRef reports the first private attribute an expression references, or "" if it references none. dynamic reports that the expression names attributes at runtime, which the caller must treat as unauthorized on an unprivileged connection: a reference it cannot see is one it cannot check.
What this CANNOT see, and what closes it: an ad may store a public attribute whose value is an expression over a private one (`Leak = ClaimId`). A constraint on `Leak` references nothing private and so passes here, while evaluation resolves ClaimId. Closing that needs the private value to be absent at EVALUATION time for an unprivileged reader, not a check on the constraint text -- the redacted decode walk already does this for rendering.
func SealWalkStats ¶ added in v0.30.0
func SealWalkStats() (examined, deltas int64)
SealWalkStats reports what the collapse walk examined and found. See collections.SealWalkStats.
func SealedProbesSkipped ¶ added in v0.30.7
func SealedProbesSkipped() int64
SealedProbesSkipped reports how many sealed-segment key probes were skipped because the segment's key index is not built yet.
func SpliceStats ¶ added in v0.30.0
func SpliceStats() (spliced, spliceRefused, objectPath int64)
SpliceStats reports, process-wide, how delta merges were served: by splicing attribute bytes, by decoding after the splice refused, and by decoding because the caller wanted an object rather than bytes. All three matter -- a merge that never attempts a splice and one that attempts and refuses look identical in a one-number report.
func SystemKey ¶ added in v0.9.0
SystemKey builds a reserved system key from name (prefixing the NUL sentinel).
func UnreadableBaseReasonNames ¶ added in v0.30.3
func UnreadableBaseReasonNames() []string
UnreadableBaseReasonNames lists every reason name UnreadableBaseReasons can report, including ones currently at zero, so a consumer can publish a stable set of counters.
func UnreadableBaseReasons ¶ added in v0.30.3
UnreadableBaseReasons breaks UnreadableBaseRefusals down by WHY the base could not be read, keyed by reason name ("not-visible", "segment-gone", "reassemble", "delta-chain", "decode"). The repair differs per reason -- a snapshot miss, a reaped segment, a lost columnar payload, an unresolvable delta chain and a decode failure have nothing in common -- so the total alone cannot direct an investigation. 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 in the store but its current record could not be read. Nonzero means a key the store holds failed to resolve; the write was reported to the caller as a conflict rather than stored against an empty ad, which is what used to turn such a miss into an identity-less row.
func ValidTableName ¶
ValidTableName reports whether name is usable as a table (and a directory): it must start with a letter or underscore and contain only letters, digits, underscores, and hyphens.
Types ¶
type AggFunc ¶ added in v0.16.8
type AggFunc uint8
AggFunc is a SQL aggregate function.
const ( AggCount AggFunc = iota // COUNT(*) or COUNT(col) AggSum AggAvg AggMin AggMax // AggCountDistinct is COUNT(DISTINCT col): the number of distinct defined values of // the argument in the group. It is EXACT, and therefore keeps one entry per distinct // value per group while the scan runs -- fine for the attributes people group and // count by (Owner, JobStatus, a host name), but not something to point at a // unique-per-row attribute over a large history. There is deliberately no silent // switch to a sketch: a query written COUNT(DISTINCT ...) gets a true count, and an // approximate one would have to be asked for by name. AggCountDistinct )
type AggRow ¶ added in v0.16.8
AggRow is one group's result: the group-by column values followed by the aggregate values, all rendered as strings (aligned with the request's group columns and aggregate specs).
func AggregateValues ¶ added in v0.16.8
func AggregateValues(seq iter.Seq[[]classad.Value], attrs []string, groupCols []GroupCol, aggs []AggSpec, groupCol, aggCol []int, stop func() bool) ([]AggRow, error)
AggregateValues is the shared GROUP BY core: it buckets a sequence of projected value rows (each aligned to the attribute list AggProjection returned) by the (possibly time-bucketed) group tuple and reduces each group with the COUNT/SUM/ AVG/MIN/MAX accumulators, returning one AggRow per group in first-seen order. groupCol[i]/aggCol[i] index into each value row (aggCol[i] < 0 means COUNT(*)). With no group columns it returns a single row aggregating the whole sequence, still yielding one row over an empty sequence (SQL semantics: COUNT is 0, others undefined). If stop is non-nil it is polled per row and, when it returns true, the scan halts early with the groups accumulated so far -- used to abandon a scan whose client has gone away.
func ColumnarAggregate ¶ added in v0.25.2
func ColumnarAggregate(stats func(attr string) (NumStats, bool), groupCols []GroupCol, aggs []AggSpec) ([]AggRow, bool)
ColumnarAggregate answers aggs over constraint from the columnar accelerator when it can: exactly one aggregate, no grouping, no per-aggregate FILTER, and a numeric argument the current schema carries. ok=false means nothing was computed and the caller must scan.
COUNT(attr), MIN, MAX, SUM and AVG are served. The result TYPE follows the reference exactly (see numAggValue): SUM accumulates integers in int64 and only becomes a real once a real value appears, AVG is always a real, and MIN/MAX keep their element's type. A formatting difference would be as wrong as a numeric one and far easier to ship unnoticed, so the test compares text.
Two refusals. A FILTER is refused rather than approximated: the columnar pass knows nothing about it, so answering would report an unfiltered aggregate as if it were filtered. And a column that turned up a BOOLEAN declines: the reference coerces booleans to 1/0 and then has a further quirk for a lone boolean element, which is not worth reproducing for data this pathological -- the scan gives the exact answer.
func GroupedFromColumns ¶ added in v0.28.0
func GroupedFromColumns(src GroupStatsSource, constraint string, groupCols []GroupCol, aggs []AggSpec) ([]AggRow, bool)
GroupedFromColumns answers a single-numeric-column GROUP BY from src's columns, or ok=false to scan.
Served: one group column with no bucket width, and aggregates drawn from COUNT(*), COUNT(attr), MIN, MAX, SUM and AVG over numeric attributes the schema carries. A per-aggregate FILTER declines -- the columnar pass knows nothing about it, so answering would report an unfiltered aggregate as a filtered one. COUNT(DISTINCT) declines: it needs the values, not their aggregate.
type AggSpec ¶ added in v0.16.8
AggSpec is one aggregate in a query: a function over an argument attribute. Arg "*" (only meaningful for plain COUNT) counts every row in the group; otherwise Arg is an attribute name evaluated per ad.
Filter, when non-empty, is a ClassAd expression restricting this aggregate -- and only this one -- to the rows of its group where the expression is true (SQL's `COUNT(*) FILTER (WHERE ...)`). It is what lets one pass over the data answer several differently-conditioned questions at once:
SELECT Owner, COUNT(*), COUNT(*) FILTER (WHERE JobStatus == 2) FROM jobs GROUP BY Owner
Without it that is one scan per condition. The filter narrows an aggregate, never the group: a group whose rows all fail every filter still appears, with COUNT 0 and the other functions undefined, exactly as SQL has it.
type ArchiveConfig ¶ added in v0.7.0
type ArchiveConfig struct {
// SegmentSize is the sealed-segment file size in bytes (default 8 MiB).
SegmentSize int
// RowGroupBytes is the uncompressed record-bytes budget for one columnar row group. 0 takes the
// default (see collections.Options.RowGroupBytes).
//
// Like the rest of this config it is read at CREATE time only -- a reopened archive uses the
// persisted value. Use SetRowGroupBytes to change it on a live archive: unlike the index set, this
// needs no reconciliation, because each block records its own layout, so a new budget governs
// segments sealed from then on and older ones keep theirs.
//
// It trades point-lookup cost against how far compression can see across records: a larger group
// stores the same data in fewer bytes and makes reading one record decompress more of it. The
// default is measured on real ads, but the right point depends on the read mix and on how much of
// the working set fits the block cache, which only a production-sized archive can answer.
RowGroupBytes int
// HotAttrs / CategoricalAttrs / ValueAttrs tune the per-segment hot header and
// indexes; ZoneAttrs names numeric attributes to keep per-segment min/max on for
// whole-segment query pruning (value-indexed attributes are included automatically).
HotAttrs []string
CategoricalAttrs, ValueAttrs []string
ZoneAttrs []string
// AutoAttrs names the subset of the above created by the auto-tuner rather than by a
// human. Provenance has to be persisted with the index set, or an auto index returns
// from a restart indistinguishable from one someone asked for -- exempt from trimming
// and from any future auto-drop, permanently, on the strength of nothing.
AutoAttrs []string
// IndexBackfillBytes bounds how far back an index configuration change is carried: only
// segments within the newest N bytes have their existing index rebuilt, older ones keep
// the one they have (see collections.Options.IndexBackfillBytes).
//
// Zero takes defaultIndexBackfillBytes. A NEGATIVE value means no bound: the change is
// carried across the whole archive, which for history means decompressing every record
// to add one index -- hours at scale, to speed up reads of the oldest data that a
// newest-first query with a limit may never reach. That is a reasonable thing to ask
// for deliberately and a bad thing to get by leaving a field unset, which is why it is
// not what zero means.
IndexBackfillBytes int64
// GroupSchemaCount is how many secondary columnar schemas to derive for attributes the base
// schema does not carry (see db.Config). 0 takes the default; a NEGATIVE value builds none.
// An archive's maintenance cadence is long, so the stability gate's several derivations span
// days rather than minutes -- the intended shape for a table whose structure changes on that
// scale, and the reason an archive builds none for its first few days.
GroupSchemaCount int
GroupStabilityRuns int
GroupMergeJaccard float64
GroupMaxPartialFrac float64
// Retention bounds what rotation keeps (max segments / bytes / age). Zero keeps all.
Retention collections.Retention
}
ArchiveConfig configures an archive table. Dir is set by the catalog.
type ArchiveTable ¶ added in v0.7.0
type ArchiveTable struct {
// contains filtered or unexported fields
}
An ArchiveTable is an append-only, rotated history store (the "condor history file" use case), exposed as a catalog table type alongside the mutable tables. Ads are appended once and never updated or deleted individually; old data is dropped in bulk by rotation. Queries are newest-first with an optional limit -- condor_history's "last K" -- with whole-segment pruning via zone maps. See collections.Archive.
func (*ArchiveTable) AddIndex ¶ added in v0.18.0
func (t *ArchiveTable) AddIndex(categorical, value []string) bool
AddIndex adds per-segment indexes on the named categorical and/or value attributes and rebuilds them over existing segments. Returns false if the index set was unchanged. A change is persisted (see saveIndexConfig).
func (*ArchiveTable) Aggregate ¶ added in v0.16.8
func (t *ArchiveTable) Aggregate(constraint string, groupBy []string, aggs []AggSpec) ([]AggRow, error)
Aggregate runs a server-side GROUP BY over the archive's matches: it applies the constraint (using the archive's zone-map pruning, so segments no matching record can fall in are never scanned), groups by the raw group columns, and reduces each group with COUNT/SUM/AVG/MIN/MAX. It shares the exact grouping/reduce engine (AggregateValues) the mutable-table aggregate uses, so an archive aggregate behaves identically to the same aggregate over a live table -- only the (small) grouped result is produced, not every matched ad. With no group columns it returns a single row over the whole match.
func (*ArchiveTable) AggregateCols ¶ added in v0.23.1
func (t *ArchiveTable) AggregateCols(constraint string, groupCols []GroupCol, aggs []AggSpec) ([]AggRow, error)
AggregateCols is Aggregate where a group column may carry a bucket width, so a numeric attribute can be grouped into fixed-width buckets (the "per day" dimension) server-side.
func (*ArchiveTable) AggregateColsStats ¶ added in v0.29.6
func (t *ArchiveTable) AggregateColsStats(constraint string, groupCols []GroupCol, aggs []AggSpec, stats *collections.ScanStats) ([]AggRow, error)
AggregateColsStats is AggregateCols that also fills stats (may be nil) with the scan's work breakdown for EXPLAIN ANALYZE. Only the fall-through projected scan -- the slow path taken when the aggregated attribute is not a numeric field the columnar accelerator covers -- fills it; the columnar/index fast paths do not reassemble records and leave stats zero (RecordsVisited == 0), which the caller reads as "served without a per-record scan". So a nonzero, large RecordsReassembled here is exactly the signature of the slow aggregate this trailer exists to surface.
func (*ArchiveTable) Append ¶ added in v0.7.0
func (t *ArchiveTable) Append(ad *classad.ClassAd) error
Append adds one ad to the archive (append-only; there is no update or per-key delete).
func (*ArchiveTable) AppendOld ¶ added in v0.7.0
func (t *ArchiveTable) AppendOld(text string) error
AppendOld appends an ad parsed from old-ClassAd text (the qmgmt/history line format).
func (*ArchiveTable) AutoIndexNames ¶ added in v0.25.0
func (t *ArchiveTable) AutoIndexNames() []string
AutoIndexNames returns the names of the archive's auto-created indexes.
func (*ArchiveTable) AutoTune ¶ added in v0.25.0
func (t *ArchiveTable) AutoTune(opts AutoTuneOptions) AutoTuneResult
AutoTune adds indexes the observed query demand justifies, persisting the result so the backfill it just paid for is not thrown away and re-paid on the next restart.
Adds only: opts.DropUnused is forced off. On a mutable table an auto-drop is reversible at roughly the cost of the index; on an archive, re-adding means reading history back through the decompressor, so the cheap direction and the expensive direction are the wrong way round for acting on an absence of evidence.
func (*ArchiveTable) BuildAndEnableSchemaScan ¶ added in v0.25.2
func (t *ArchiveTable) BuildAndEnableSchemaScan(sampleMax, hotTopN int) bool
BuildAndEnableSchemaScan builds or extends the archive's columnar accelerator (see collections.Archive.BuildAndEnableSchemaScan).
func (*ArchiveTable) CategoricalGroupCounts ¶ added in v0.23.1
func (t *ArchiveTable) CategoricalGroupCounts(attr string) (map[string]int64, bool)
CategoricalGroupCounts returns the exact per-value record counts for a categorically indexed attribute, read from the per-segment indexes without scanning any record. ok is false when the indexes cannot fully account for every record (see the collections-level doc), in which case the caller must scan. Aggregate uses this for GROUP BY COUNT(*); it is exported so a planner can ask whether the cheap path is available.
func (*ArchiveTable) CategoricalGroupCountsBucketed ¶ added in v0.23.1
func (t *ArchiveTable) CategoricalGroupCountsBucketed(attr, bucketAttr string, width int64) (map[int64]map[string]int64, bool)
CategoricalGroupCountsBucketed is CategoricalGroupCounts split by a numeric bucket (floor(bucketAttr/width)*width) -- the "per group per day" shape. bucketAttr must be zone-mapped; ok is false when the counts cannot be established from the indexes.
Not yet reachable from SQL: the archive aggregate crosses dbrpc as a []string group list, which has no room for a bucket width, so a bucketed GROUP BY over an archive still falls back to client-side reduction. Wiring it through is a protocol change.
func (*ArchiveTable) CategoricalGroupCountsWhere ¶ added in v0.23.1
func (t *ArchiveTable) CategoricalGroupCountsWhere(attr, constraint string) (map[string]int64, bool)
CategoricalGroupCountsWhere is CategoricalGroupCounts restricted to records matching constraint. ok is false unless the constraint is a pure conjunction of numeric comparisons on zone-mapped attributes and the indexes account for every record.
func (*ArchiveTable) Close ¶ added in v0.7.0
func (t *ArchiveTable) Close() error
Close flushes and closes the archive.
func (*ArchiveTable) CodecStats ¶ added in v0.19.0
func (t *ArchiveTable) CodecStats(sampleMax int) CodecStats
CodecStats reports the archive's compression (codec, dict size, last retrain, sampled ratio).
func (*ArchiveTable) Count ¶ added in v0.7.0
func (t *ArchiveTable) Count() int
Count is the number of records currently retained (reduced by rotation).
func (*ArchiveTable) CountConstraint ¶ added in v0.27.0
func (t *ArchiveTable) CountConstraint(constraint string) (int, bool)
CountConstraint counts the rows matching constraint via the columnar accelerator, or reports ok=false so the caller scans. See collections.Archive.CountConstraint.
func (*ArchiveTable) DropIndex ¶ added in v0.18.0
func (t *ArchiveTable) DropIndex(names ...string) bool
DropIndex removes the named per-segment indexes. Returns false if none matched. A change is persisted (see saveIndexConfig).
func (*ArchiveTable) EncryptedAttrNames ¶ added in v0.29.0
func (t *ArchiveTable) EncryptedAttrNames() []string
EncryptedAttrNames reports the attributes this archive seals, matching DB.EncryptedAttrNames.
func (*ArchiveTable) EncryptionEnabled ¶ added in v0.29.0
func (t *ArchiveTable) EncryptionEnabled() bool
EncryptionEnabled reports whether this archive's values are sealed. See collections.Archive.EncryptionEnabled: today this is false for every archive, because the archive open path passes no data key. It is reported rather than omitted so the difference from a mutable table is visible in .stats instead of implied by a missing line.
func (*ArchiveTable) Explain ¶ added in v0.25.2
func (t *ArchiveTable) Explain(constraint string) (QueryExplain, error)
MergePass merges cold segment runs when the archive has reached its trigger watermark, down to its target. Returns the number of merges performed.
An append-only table never compacts, so its segment count only grows, and every sealed segment costs a memory mapping at open. This is what bounds that; without it running, the count climbs until the daemon cannot start. Explain reports how a query against the archive would be executed (see collections.Archive.ExplainQuery), so `.explain` works on a history table rather than reporting it as a nonexistent table.
func (*ArchiveTable) GCFloor ¶ added in v0.21.4
func (t *ArchiveTable) GCFloor() float64
GCFloor returns the current runtime GC watermark (0 when unset).
func (*ArchiveTable) GroupCountAll ¶ added in v0.28.0
func (t *ArchiveTable) GroupCountAll(groupAttr string) ([]collections.GroupCount, bool)
GroupCountAll is GroupCountConstraint over every row: the whole column's histogram. AggregateCols uses it for a grouped COUNT(*) whose constraint is match-all, which the predicate analysis behind GroupCountConstraint cannot serve because there is no predicate to analyze. See collections.Archive.GroupCountAll.
func (*ArchiveTable) GroupCountConstraint ¶ added in v0.28.0
func (t *ArchiveTable) GroupCountConstraint(constraint, groupAttr string) ([]collections.GroupCount, bool)
GroupCountConstraint answers a per-value COUNT(*) of groupAttr over the rows matching constraint via the columnar accelerator, or reports ok=false so the caller scans. AggregateCols uses it for the single-numeric-column GROUP BY shape; it is exported so a caller (and a test) can tell which tier answered. See collections.Archive.GroupCountConstraint.
func (*ArchiveTable) GroupSchemaAgreement ¶ added in v0.28.0
func (t *ArchiveTable) GroupSchemaAgreement(sampleMax, k int) GroupSchemaAgreement
GroupSchemaAgreement reports how well per-segment derivations agree with the archive-wide one.
func (*ArchiveTable) GroupSchemaChanges ¶ added in v0.29.13
func (t *ArchiveTable) GroupSchemaChanges() []GroupSchemaChange
GroupSchemaChanges returns the archive's committed-group change log (see DB.GroupSchemaChanges).
func (*ArchiveTable) GroupSchemaDrift ¶ added in v0.28.0
func (t *ArchiveTable) GroupSchemaDrift() GroupSchemaDrift
GroupSchemaDrift reports how the archive's derived groups have moved.
func (*ArchiveTable) GroupSchemaLastAgreement ¶ added in v0.29.13
func (t *ArchiveTable) GroupSchemaLastAgreement() (GroupSchemaLastAgreement, bool)
GroupSchemaLastAgreement returns the archive's last persisted per-segment agreement, if any.
func (*ArchiveTable) GroupSchemas ¶ added in v0.28.0
func (t *ArchiveTable) GroupSchemas(sampleMax, k int) GroupSchemaInfo
GroupSchemas derives and reports candidate group schemas for the archive (see DB.GroupSchemas).
func (*ArchiveTable) GroupStatsAll ¶ added in v0.28.0
func (t *ArchiveTable) GroupStatsAll(groupAttr string, aggAttrs []string) ([]collections.GroupStats, bool)
GroupStatsAll is GroupStatsConstraint over every row, for a constraint this layer has established is match-all (the predicate analysis behind GroupStatsConstraint has no predicate to work with there).
func (*ArchiveTable) GroupStatsConstraint ¶ added in v0.28.0
func (t *ArchiveTable) GroupStatsConstraint(constraint, groupAttr string, aggAttrs []string) ([]collections.GroupStats, bool)
GroupStatsConstraint answers a per-group record count plus the aggregate inputs for each aggAttr over the rows matching constraint, via the columnar accelerator, or reports ok=false so the caller scans. AggregateCols uses it for the single-numeric-column GROUP BY shape; it is exported so a caller (and a test) can tell which tier answered. See collections.Archive.GroupStatsConstraint.
func (*ArchiveTable) HotAttrs ¶ added in v0.29.0
func (t *ArchiveTable) HotAttrs() []string
HotAttrs reports the archive's hot-header attributes, matching DB.HotAttrs for a mutable table.
func (*ArchiveTable) IndexSizes ¶ added in v0.19.0
func (t *ArchiveTable) IndexSizes() IndexSizes
IndexSizes reports the per-attribute index byte footprint.
func (*ArchiveTable) IndexedAttrs ¶ added in v0.19.0
func (t *ArchiveTable) IndexedAttrs() (categorical, value []string)
IndexedAttrs returns the archive's categorical and value index attributes.
func (*ArchiveTable) MergePass ¶ added in v0.24.1
func (t *ArchiveTable) MergePass(opts MergeOptions) int
func (*ArchiveTable) NumStats ¶ added in v0.25.2
func (t *ArchiveTable) NumStats(constraint, attr string) (NumStats, bool)
NumStats is DB.NumStats for an archive (history) table.
func (*ArchiveTable) OpStats ¶ added in v0.19.0
func (t *ArchiveTable) OpStats() OpStats
OpStats reports cumulative operational timings. An archive has no DB-level snapshot lock, so SnapshotLock is zero.
func (*ArchiveTable) Query ¶ added in v0.7.0
Query returns the archived ads matching constraint, newest first. QueryLimit caps the result at the newest limit matches (<= 0 = all) -- the scan stops after the newest satisfying segments, so "last K" is cheap.
func (*ArchiveTable) QueryLimit ¶ added in v0.7.0
func (*ArchiveTable) QueryProject ¶ added in v0.20.0
func (t *ArchiveTable) QueryProject(constraint string, attrs []string) (iter.Seq[[]classad.Value], error)
QueryProject scans the matching ads and yields each projected to just attrs' values, read wire-native where possible -- so an aggregate reads only the attributes it needs instead of fully decoding every record. Errors only on a malformed constraint.
func (*ArchiveTable) QueryProjectStats ¶ added in v0.29.6
func (t *ArchiveTable) QueryProjectStats(constraint string, attrs []string, stats *collections.ScanStats) (iter.Seq[[]classad.Value], error)
QueryProjectStats is QueryProject that also fills stats (may be nil) with the scan's work breakdown, so an EXPLAIN ANALYZE of an aggregate that falls to this projected scan can show the large RecordsReassembled that marks the slow per-record path.
func (*ArchiveTable) QueryRawProjected ¶ added in v0.20.1
func (t *ArchiveTable) QueryRawProjected(constraint string, projection []string, redact bool) (iter.Seq[collections.RawAd], error)
QueryRawProjected yields each matching ad as a raw projected subset (only the projection attributes, rendered from the stored representation), newest first — the archive-side of the server-side projection op. redact strips private attributes. It mirrors db.DB.QueryRawProjected so the same wire op serves archives and mutable tables uniformly. chaseRefs is false, matching HTCondor's projection protocol (exactly the requested attributes).
func (*ArchiveTable) QueryRawProjectedRefs ¶ added in v0.24.0
func (t *ArchiveTable) QueryRawProjectedRefs(constraint string, projection []string, redact bool) (iter.Seq[collections.RawAd], error)
QueryRawProjectedRefs is QueryRawProjected with the projected expressions' attribute references resolved too, so each yielded ad evaluates self-contained. See db.DB.QueryRawProjectedRefs for why that is a separate call rather than the default.
func (*ArchiveTable) QueryRawProjectedRefsStats ¶ added in v0.29.1
func (t *ArchiveTable) QueryRawProjectedRefsStats(constraint string, projection []string, redact bool, stats *collections.ScanStats) (iter.Seq[collections.RawAd], error)
QueryRawProjectedRefsStats is QueryRawProjectedRefs that also fills stats (may be nil) with the per-scan work breakdown for EXPLAIN ANALYZE.
func (*ArchiveTable) QueryRawProjectedStats ¶ added in v0.29.0
func (t *ArchiveTable) QueryRawProjectedStats(constraint string, projection []string, redact bool, stats *collections.ScanStats) (iter.Seq[collections.RawAd], error)
QueryRawProjectedStats is QueryRawProjected that also fills stats (may be nil) with the per-scan work breakdown for EXPLAIN ANALYZE (segments scanned/pruned, records decided from columns vs reassembled, rows matched).
func (*ArchiveTable) Reindex ¶ added in v0.18.0
func (t *ArchiveTable) Reindex()
Reindex rebuilds the per-segment indexes over all segments.
func (*ArchiveTable) ReschemaScan ¶ added in v0.25.3
func (t *ArchiveTable) ReschemaScan(sampleMax, hotTopN int) bool
ReschemaScan re-derives the archive's schema and rebuilds every sealed segment's columnar block (see DB.ReschemaScan).
func (*ArchiveTable) Retention ¶ added in v0.19.0
func (t *ArchiveTable) Retention() collections.Retention
Retention returns the archive's current retention bounds.
func (*ArchiveTable) RetrainDict ¶ added in v0.18.0
func (t *ArchiveTable) RetrainDict(sampleMax int) (int, error)
RetrainDict trains a fresh ZSTD dictionary from up to sampleMax records and recompresses every segment in place under it, returning the new dictionary's size in bytes.
func (*ArchiveTable) Rewrite ¶ added in v0.18.0
func (t *ArchiveTable) Rewrite() int
Rewrite recompresses and re-encodes every segment in place under the current codec and hot set, returning the number of records rewritten.
func (*ArchiveTable) Rotate ¶ added in v0.7.0
func (t *ArchiveTable) Rotate(now float64) (int, error)
Rotate drops sealed segments that fall outside the retention policy, given the current time (unix seconds, for age-based retention). Returns how many segments were dropped.
func (*ArchiveTable) RowGroupBytes ¶ added in v0.28.1
func (t *ArchiveTable) RowGroupBytes() int
RowGroupBytes reports the budget currently in effect (0 meaning the default).
func (*ArchiveTable) SaveDemand ¶ added in v0.25.0
func (t *ArchiveTable) SaveDemand()
SaveDemand checkpoints the archive's recorded query demand so index decisions survive a restart, ageing it by the time since the last checkpoint. See collections.SaveDemand.
func (*ArchiveTable) SchemaFit ¶ added in v0.25.3
func (t *ArchiveTable) SchemaFit(sampleMax int) ([]SchemaFieldFit, int)
SchemaFit measures the archive's derived schema against a fresh sample (see DB.SchemaFit).
func (*ArchiveTable) SchemaScanInfo ¶ added in v0.25.2
func (t *ArchiveTable) SchemaScanInfo() SchemaScanInfo
SchemaScanInfo reports the archive's columnar accelerator state.
func (*ArchiveTable) SetGCFloor ¶ added in v0.21.4
func (t *ArchiveTable) SetGCFloor(floor float64)
SetGCFloor installs a runtime GC watermark (in the archive's MinAgeAttr units) so the next Rotate may reclaim already-consumed records early: a change-feed source passes the feed's GC floor (min ack over live subscribers) here to drain records every subscriber has read. It only shortens retention -- it never keeps data past the configured Retention ceilings, and never drops anything younger than Retention.MinAge. Unlike SetRetention this is NOT persisted -- the caller re-asserts it from the current live floor each pass, so a stale value can never GC data across a restart. Passing floor <= 0 clears it.
func (*ArchiveTable) SetRetention ¶ added in v0.19.0
func (t *ArchiveTable) SetRetention(r collections.Retention) error
SetRetention updates the retention bounds and persists them (archiveconfig.json), so they take effect on the next Rotate and survive a restart.
func (*ArchiveTable) SetRowGroupBytes ¶ added in v0.28.1
func (t *ArchiveTable) SetRowGroupBytes(n int) error
SetRowGroupBytes changes the columnar row-group budget and persists it (archiveconfig.json), so a value tried on a live archive survives a restart. 0 restores the default.
Unlike the index set this needs no reconciliation pass: every block records the layout it was written with, so the new budget governs segments sealed from now on while everything already on disk keeps reading as before. That is what makes it worth exposing -- the right value depends on the read mix and on how much of the working set fits the block cache, which only a production-sized archive can answer, and finding it must not mean rebuilding the archive.
func (*ArchiveTable) SidecarSizes ¶ added in v0.19.0
func (t *ArchiveTable) SidecarSizes() SidecarSizes
SidecarSizes reports the sealed-segment sidecar index bytes (mmap-backed, evictable).
func (*ArchiveTable) StaleIndexSegments ¶ added in v0.21.2
func (t *ArchiveTable) StaleIndexSegments() (stale, sealed int)
StaleIndexSegments reports how many sealed segments still carry an index built under an older configuration, and how many are sealed in total (see collections.Archive.AddIndex: a sealed sidecar is immutable, so a runtime index change reaches old segments only via a Rewrite).
func (*ArchiveTable) StaleIndexSegmentsByPolicy ¶ added in v0.25.0
func (t *ArchiveTable) StaleIndexSegmentsByPolicy() int
StaleIndexSegmentsByPolicy reports segments deliberately left on an older index configuration because they fall outside the backfill horizon -- expected, and not a sign of anything failing. Reported apart from StaleIndexSegments so an alert on the latter stays meaningful once a horizon is configured.
func (*ArchiveTable) Stats ¶ added in v0.19.0
func (t *ArchiveTable) Stats() Stats
Stats reports storage accounting (records, segments, arena/used/dead bytes).
func (*ArchiveTable) StopMaintenance ¶ added in v0.30.4
func (a *ArchiveTable) StopMaintenance()
StopMaintenance asks any maintenance pass running on this archive to stop at its next safe boundary. See collections.Collection.StopMaintenance.
func (*ArchiveTable) TopK ¶ added in v0.29.1
func (t *ArchiveTable) TopK(constraint string, attrs []string, orderAttr string, desc bool, k int) ([][]classad.Value, error)
TopK returns the k rows matching constraint ordered by orderAttr -- the largest k when desc, the smallest k otherwise -- each projected to attrs, in sorted order. It is the server-side form of ORDER BY <orderAttr> {DESC|ASC} LIMIT k: only k rows are ever held or returned, not the whole match set, so a "newest few of a million" query does not ship a million rows to be sorted client-side.
orderAttr must resolve to a number (ClusterId, QDate, CompletionDate, ...); a row whose order value is not numeric has no place in the ordering and is skipped. If orderAttr is not among attrs it is fetched for the comparison and trimmed from the returned rows.
func (*ArchiveTable) TopKStats ¶ added in v0.29.4
func (t *ArchiveTable) TopKStats(constraint string, attrs []string, orderAttr string, desc bool, k int) ([][]classad.Value, collections.ScanStats, error)
TopKStats is TopK that also returns the cutoff scan's ScanStats, for EXPLAIN ANALYZE.
func (*ArchiveTable) Truncate ¶ added in v0.20.1
func (t *ArchiveTable) Truncate()
Truncate drops every record, resetting the archive to empty in place (see collections.Archive.Truncate). It is the destructive reset behind a from-scratch history re-sync: empty the table, then re-ingest from the source. The persisted config (indexes, zone maps, retention) is retained, so the archive keeps its shape.
func (*ArchiveTable) Watch ¶ added in v0.20.2
func (t *ArchiveTable) Watch(ctx context.Context, cursor []byte) (iter.Seq[WatchEvent], error)
Watch streams the archive's change events (append-only: upserts and the catch-up/live markers, no deletes), converting collections.WatchEvent to the db WatchEvent used by the mutable-table watch so the two are wire-identical. An empty cursor replays the retained history, then goes live. See collections.Archive.Watch.
func (*ArchiveTable) WatchCursor ¶ added in v0.29.12
func (t *ArchiveTable) WatchCursor() ([]byte, error)
WatchCursor returns an opaque cursor at the current head of the archive's change log. Watch(cursor) then streams only what is appended after it, which is how a client tails an archive without first replaying everything rotation still retains.
func (*ArchiveTable) ZoneAttrs ¶ added in v0.21.2
func (t *ArchiveTable) ZoneAttrs() []string
ZoneAttrs returns the attributes carrying per-segment [min,max] zone maps, on which a range query prunes whole segments rather than only postings.
type AutoTuneOptions ¶ added in v0.25.0
type AutoTuneOptions = collections.AutoTuneOptions
AutoTuneOptions and AutoTuneResult are re-exported so a caller can tune an archive's indexes without importing collections directly.
type AutoTuneResult ¶ added in v0.25.0
type AutoTuneResult = collections.AutoTuneResult
type Catalog ¶
type Catalog struct {
// contains filtered or unexported fields
}
A Catalog is a set of named tables, each an independent ClassAd store (its own keyspace, indexes, hot set, and persisted index config). Tables give the database more than one collection to work with -- e.g. separate machine and job ads -- without joins. A persistent catalog keeps each table in its own subdirectory (<dir>/tables/<name>), so tables are isolated on disk and the set is recovered by enumerating those subdirectories on open.
func OpenCatalog ¶
OpenCatalog opens the catalog rooted at dir with no encryption. See OpenCatalogConfig to enable encryption at rest.
func OpenCatalogConfig ¶ added in v0.7.0
func OpenCatalogConfig(cfg CatalogConfig) (*Catalog, error)
func (*Catalog) ArchiveTable ¶ added in v0.7.0
func (cat *Catalog) ArchiveTable(name string) (*ArchiveTable, bool)
ArchiveTable returns the archive (history) table named name.
func (*Catalog) ArchiveTables ¶ added in v0.7.0
ArchiveTables returns the archive table names, sorted.
func (*Catalog) ConvertTableToMemory ¶ added in v0.10.2
ConvertTableToMemory drops the on-disk backing of an existing persistent table, keeping its current contents live in RAM only. It copies the table through a consistent snapshot into a fresh in-memory table (preserving ads, index configuration, hot set, codec, and encryption policy), swaps that in, then closes and deletes the on-disk original.
It is a no-op when the catalog itself is in-memory or the table is already RAM-only. Like Rewrite, it takes a consistent snapshot but does not globally quiesce writers, so a write that races the swap can be lost -- run it during low write activity. DAEMON-level: it changes a table's durability, so callers gate it above ordinary WRITE.
func (*Catalog) CreateArchiveTable ¶ added in v0.7.0
func (cat *Catalog) CreateArchiveTable(name string, cfg ArchiveConfig) (*ArchiveTable, error)
CreateArchiveTable creates (or returns the existing) append-only archive table named name, persisted under <dir>/archives/<name>. cfg configures indexes / zone maps / retention on first creation and is ignored for an existing archive.
func (*Catalog) CreateExporter ¶ added in v0.14.0
func (cat *Catalog) CreateExporter(def ExporterDef) error
CreateExporter registers a new exporter definition. The name must be a valid identifier and must not already be registered. Exporters occupy their own namespace, so an exporter may share a name with a table (an exporter named "jobs" that mirrors the "jobs" table is expected). On a persistent catalog the definition is written to disk before it is returned; on an in-memory catalog it lives only in memory.
func (*Catalog) CreateTable ¶
CreateTable creates (or returns the existing) table named name. Its data persists under <dir>/tables/<name> for a persistent catalog.
func (*Catalog) CreateTableInMemory ¶ added in v0.10.2
CreateTableInMemory creates (or returns the existing) table named name as RAM-only, even in a persistent catalog -- shorthand for CreateTableOpts with InMemory set. If the table already exists its backing is unchanged (options are not re-applied); use ConvertTableToMemory to drop an existing persistent table's on-disk backing.
func (*Catalog) CreateTableOpts ¶ added in v0.10.1
func (cat *Catalog) CreateTableOpts(name string, opts TableOptions) (*DB, error)
CreateTableOpts is CreateTable with options; see TableOptions. If the table already exists it is returned as-is (options are not re-applied).
func (*Catalog) CreateView ¶ added in v0.12.0
CreateView creates a materialized view named name over spec.BaseTable, materializing it synchronously (so a cardinality-limit overflow or an aggregation error fails the create and leaves nothing behind) and then maintaining it live from the base table's change stream. The definition is persisted for a persistent catalog.
func (*Catalog) DropArchiveTable ¶ added in v0.7.0
DropArchiveTable closes and removes the archive named name, deleting its on-disk data.
func (*Catalog) DropExporter ¶ added in v0.14.0
DropExporter removes an exporter's definition and its resume state. It is not an error to drop an exporter that does not exist (idempotent teardown).
func (*Catalog) DropTable ¶
DropTable closes and removes the table named name, deleting its on-disk data.
func (*Catalog) DropView ¶ added in v0.12.0
DropView removes a view: it stops the updater, drops the in-memory data, and deletes the persisted definition.
func (*Catalog) EnsureTable ¶
EnsureTable creates the table if it does not exist, returning it.
func (*Catalog) Exporter ¶ added in v0.14.0
func (cat *Catalog) Exporter(name string) (ExporterDef, bool)
Exporter returns a single exporter definition.
func (*Catalog) Exporters ¶ added in v0.14.0
func (cat *Catalog) Exporters() []ExporterDef
Exporters returns all registered exporter definitions, sorted by name.
func (*Catalog) LoadExporterState ¶ added in v0.14.0
LoadExporterState returns an exporter's last checkpointed resume-state blob. The bool is false when the exporter has never checkpointed (a fresh exporter with no state yet), which callers treat as "start from the beginning", not as an error.
func (*Catalog) Restore ¶ added in v0.7.0
Restore replaces each table named in the catalog snapshot with its snapshotted contents (creating the table if absent), each under that table's DB-wide lock. Tables present in the catalog but ABSENT from the snapshot are left untouched. Encrypted sections are opened with the catalog's pool keys (the tables share them).
func (*Catalog) SaveExporterState ¶ added in v0.14.0
SaveExporterState durably records an exporter's opaque resume-state blob, replacing any prior state. The exporter calls this after its data has been accepted downstream, so that on restart it resumes just past the last delivered change (at-least-once). The exporter must already exist.
func (*Catalog) Snapshot ¶ added in v0.7.0
Snapshot writes a backup of every table in the catalog to w. Each table is captured under its own DB-wide lock (consistent per table); the set of tables is the snapshot's membership at the moment each name is read.
func (*Catalog) ViewBacking ¶ added in v0.12.0
ViewBacking returns the in-memory backing table of a view, so reads (SELECT/query) resolve a view name to its materialized rows.
func (*Catalog) ViewSealed ¶ added in v0.16.0
func (cat *Catalog) ViewSealed(name, constraint string) (seq iter.Seq[*classad.ClassAd], ok bool, err error)
ViewSealed streams a continuous aggregate's sealed (archived) rows matching the constraint, so a view read can union sealed history with the live backing. ok is false for a plain table or a gauge view (nothing sealed). The returned sequence is empty (not an error) for a continuous aggregate with no archive (in-memory catalog).
type CatalogConfig ¶ added in v0.7.0
type CatalogConfig struct {
// OnOpenStep, when set, is called once per table and archive opened by OpenCatalogConfig, with how long
// that one took. Catalog open is a loop over directories, so without it the whole thing is a single
// number -- which is how a 15s startup stayed unattributed: every other phase reported 0s and this one
// reported all of it, with nothing inside.
//
// kind is "table" or "archive". MAY BE CALLED CONCURRENTLY, from whichever goroutine opened that table --
// opens run in parallel -- so a handler must be safe for concurrent use and must not block. It is still
// called once per open, which is what keeps per-table attribution intact under parallelism: only the SUM
// stops matching the wall clock.
OnOpenStep func(kind, name string, d time.Duration)
Dir string
// DeltaMax sets the delta-record chain bound (see Config.DeltaMax) for every table this
// catalog opens, and DeltaMaxFor overrides it per table name. Zero takes DefaultDeltaMax --
// delta records are ON by default -- and DeltaMaxOff turns them off.
//
// Both belong on the CATALOG rather than on CreateTableOpts: OpenCatalogConfig opens every
// existing table directory at startup, so by the time a caller reaches CreateTable the table
// is already open and per-call options are not re-applied. A per-table option would
// therefore take effect the first time a table was created and be silently ignored on every
// restart after -- a knob that works once.
//
// Note that a store which has written even one delta record replays them for the rest of its
// life (the on-disk marker says so), so this is not reversible for data already written;
// DeltaMaxOff stops NEW deltas and leaves the old ones readable.
DeltaMax int
DeltaMaxFor map[string]int
// PoolKeys enables encryption at rest for every table (each table's master key is
// wrapped under these keys). EncryptedAttrs is the default explicit encrypted-attr
// set for each table (private attributes are always encrypted). See db/encrypt.go.
PoolKeys []KEK
EncryptedAttrs []string
// OnSealMigration, when set, is called for each table whose open rewrote segments to seal private
// attributes that predate sealing, with how many segments it rewrote. It carries the TABLE NAME
// because a catalog opens many, and "some table rewrote 400 segments" does not tell an operator
// which one to expect the cost from.
//
// The per-table Config hook exists too, but a catalog builds those configs itself, so without this
// the hook was unreachable for every table a daemon actually opens -- the only place the one-off
// startup cost needs explaining.
//
// Called from the goroutine that opened that table, so it may run concurrently (see OnOpenStep).
OnSealMigration func(table string, segments int)
// SealMigrationWorkers bounds the concurrent segment rewrites within one table's migration (see
// collections.MigrateSealedAttrs). 0 takes a per-table default. Tables are already opened in
// parallel, so this multiplies with that.
SealMigrationWorkers int
}
CatalogConfig configures a catalog, including encryption at rest applied to every table. Dir empty is in-memory.
type CodecStats ¶
type CodecStats = collections.CodecStats
Diagnostic and management types, re-exported from collections for callers that only import db.
type Config ¶
type Config struct {
Dir string
// Ordered configures maintained, filtered, sorted indexes -- e.g. the negotiator's
// resource-request lists (partition by Owner, sort by JobPrio then QDate), iterated
// in order via Ordered. Optional.
Ordered []OrderSpec
// HotAttrs / CategoricalAttrs / ValueAttrs / MatchClosureRoots tune storage and
// query/match push-down (see collections.Options). Optional.
HotAttrs []string
CategoricalAttrs, ValueAttrs []string
MatchClosureRoots []string
// DeltaMax bounds the delta-record chain length for this table; see DefaultDeltaMax and
// DeltaMaxOff for what 0 and a negative value mean.
//
// It turns on delta records for this table and bounds the chain length: a
// SetAttribute stores only the attributes it changed, until a key has accumulated
// DeltaMax of them and the next write stores the whole ad again. See
// collections.Options.DeltaMax for the trade.
DeltaMax int
// GroupSchemaCount is how many SECONDARY columnar schemas to derive: sets of attributes the
// base schema does not carry which are present or absent together, stored columnar for the
// ads that hold them without costing a slot in the ads that do not. 0 takes the default (4);
// a NEGATIVE value builds none.
//
// A group's blocks are built only once its members have kept showing up together across
// GroupStabilityRuns maintenance passes, so a freshly opened table builds none for a while --
// SchemaScanInfo.GroupSchemas reports how many actually exist. See collections.Options.
GroupSchemaCount int
GroupStabilityRuns int
GroupMergeJaccard float64
GroupMaxPartialFrac float64
// SealMigrationWorkers bounds the concurrent segment rewrites when an existing store is migrated
// to sealed private attributes at open (see collections.MigrateSealedAttrs). 0 takes a default
// derived from GOMAXPROCS.
SealMigrationWorkers int
// OnSealMigration, if set, is called with the number of segments rewritten when the open migrated a
// store to sealed private attributes. It exists so a daemon can log a one-off startup cost rather
// than leaving an operator to wonder why the first open after an upgrade took longer.
OnSealMigration func(segments int)
// SegmentSize overrides the arena segment size in bytes (see collections.Options).
// 0 uses the default (8 MiB). A smaller value seals segments sooner -- useful for
// tests and for tuning the sealed-segment accelerators (columnar scan, sealed indexes).
SegmentSize int
// MutatingBlockCacheBytes and ArchiveBlockCacheBytes set the PROCESS-GLOBAL shared
// decompressed-columnar-block cache budgets (see collections.Options): one budget shared by all
// mutating tables and a separate one shared by all archive tables. They are global ceilings, not
// per-DB or per-table, so a process opening many tables (as htcondordb does) no longer multiplies
// a fixed per-table cache by the table count. 0 keeps the collections default (512 MiB each). The
// same Config feeds both the mutating and archive opens, so setting them once configures the
// whole process; the last non-zero value wins and resizes the live cache.
MutatingBlockCacheBytes int64
ArchiveBlockCacheBytes int64
// PoolKeys enables encryption at rest: the DB master key is wrapped under each of
// these pool/signing keys (any one opens the DB; a rotated-in key is added on the
// next open). Empty ⇒ encryption disabled. See db/encrypt.go.
PoolKeys []KEK
// EncryptedAttrs names the attributes whose values are sealed at rest (case-
// insensitive). Only meaningful with PoolKeys set. An encrypted attribute may not
// also be indexed. See collections.Options.EncryptedAttrs.
EncryptedAttrs []string
}
type ConflictError ¶
type ConflictError struct{ Keys []string }
ConflictError reports the keys whose writes lost an optimistic write-write race at commit. The other writes in the transaction committed; the caller re-reads and retries the conflicted keys.
func (*ConflictError) Error ¶
func (e *ConflictError) Error() string
type Constraint ¶
type Constraint struct {
// contains filtered or unexported fields
}
Constraint is a compiled ClassAd boolean expression, for evaluating the same filter against many ads without re-parsing.
func ParseConstraint ¶
func ParseConstraint(expr string) (*Constraint, error)
ParseConstraint compiles a ClassAd boolean expression.
type DB ¶
type DB struct {
// contains filtered or unexported fields
}
DB is an embedded ClassAd log. Safe for concurrent use.
func Open ¶
Open opens a ClassAd log with default configuration. A non-empty dir makes it persistent (memory-mapped arenas under dir, recovered on reopen); an empty dir is in-memory. See OpenConfig for indexing / ordered-index configuration.
func OpenConfig ¶
OpenConfig opens a ClassAd log with the given configuration. Every DB is stamped with a stable DB id (persisted alongside a persistent store, so it survives reopen) and a fresh instance id for this open. Until high-availability DB servers exist they are effectively the same identity, but a follower/replica shares the DB id while carrying its own instance id.
func (*DB) AddHotAttrs ¶
AddHotAttrs pins the named attributes into the hot set and returns the resulting hot attributes. The hot set is persisted.
func (*DB) AddIndex ¶
AddIndex adds categorical and/or value indexes at runtime, returning whether the configuration changed. The spec is updated immediately; call Reindex to build the new index over existing ads. A change is persisted so it survives a restart.
func (*DB) BackupKey ¶ added in v0.7.0
BackupKey returns a copy of the backup key -- the key that unwraps a snapshot's encryption, so an operator can escrow it and decrypt/restore backups without the pool keys. It is DISTINCT from the live-data key (it cannot decrypt the running store) and from the master. Returns nil when encryption is disabled.
func (*DB) BeginRedacted ¶ added in v0.28.0
BeginRedacted is Begin for a caller not entitled to sealed values: reads through the transaction decode with no key. Writes are unaffected. See collections.Collection.BeginRedacted.
func (*DB) Chained ¶ added in v0.21.3
Chained reports whether the table has structural (parent-only) ads that Query hides, so a caller knows when Len equals the match-all row count (the COUNT(*) fast path).
func (*DB) CodecStats ¶
func (db *DB) CodecStats(sampleMax int) CodecStats
CodecStats reports the storage codec's state and effectiveness (name, dictionary size, last-retrain time, and a sampled compression ratio).
func (*DB) Compact ¶
Compact reclaims dead space in shards whose dead-byte ratio warrants it, returning the number of shards compacted.
func (*DB) CountConstraint ¶ added in v0.23.0
CountConstraint counts the rows matching constraint via the columnar schema scan when the constraint is columnar-eligible (Native, numeric comparisons on one int schema field) and schema-scan is enabled; ok=false ⇒ the caller should use the normal count path. See collections.Collection.CountConstraint.
func (*DB) Delete ¶ added in v0.8.0
Delete removes the ad at key in its own optimistic transaction, retrying on a write-write conflict. It reports whether an ad was present to remove. Equivalent to Begin + DestroyClassAd + Commit with retry.
func (*DB) DeltaStats ¶ added in v0.30.0
DeltaStats reports how many of this table's writes were stored as delta records versus in full. Both zero means delta records are off (or nothing has been written). It is the only way to tell a working delta store from one that has quietly fallen back to storing whole ads -- which reads back perfectly correctly and saves nothing.
func (*DB) DropIndex ¶
DropIndex removes the named attributes from the configured indexes, returning whether the configuration changed. A change is persisted.
func (*DB) EnableSchemaScan ¶ added in v0.23.0
EnableSchemaScan builds the per-segment adschema columnar accelerator over the table's sealed segments and enables it, choosing the uncompressed hot numeric tier as the top-hotTopN int/real fields by query demand (see collections.Collection.BuildAndEnableSchemaScan). sampleMax bounds the ad sample. Opt-in; a table that never calls this is unaffected.
func (*DB) EncryptedAttrNames ¶ added in v0.7.0
EncryptedAttrNames returns the explicit encrypted-attribute set (not the always-on private attributes). EncryptionEnabled reports whether encryption at rest is active.
func (*DB) EncryptionEnabled ¶ added in v0.7.0
EncryptionEnabled reports ENCRYPTION AT REST: the master key is wrapped under pool keys, so the data is protected from someone holding the disk. It is not "are values sealed" -- private attributes are always sealed, and without pool keys the master sits beside the data in the clear, which protects nothing and is reported here as false.
func (*DB) Explain ¶
func (db *DB) Explain(constraint string) (QueryExplain, error)
Explain reports how the store would execute a constraint query -- which conjuncts are index-usable and the resulting access path.
func (*DB) ExplainMatch ¶
func (db *DB) ExplainMatch(job *classad.ClassAd, targetConstraint string) MatchExplain
ExplainMatch reports how matchmaking job against this (resource) collection would execute: the job's Requirements rewritten over the slot with the job's attributes baked to constants, and which of the resulting probes prune via a configured index. targetConstraint, if non-empty, is the MATCH resource-side filter (WHERE TARGET / NOPREEMPT) melded into the explanation.
func (*DB) ForEach ¶
ForEach calls fn for every committed ad, in no particular order, until fn returns false. It reads a consistent snapshot (concurrent writers do not block it).
func (*DB) ForEachSystemAd ¶ added in v0.9.0
ForEachSystemAd calls fn for every internal system-keyed ad and its key -- the reaper enumeration path. System records are durable bookkeeping (e.g. idempotency markers) stored in the same table as data but hidden from every client scan/query/DeleteWhere; this is the only way to enumerate them. Reads a consistent snapshot. See SystemKey.
func (*DB) GroupSchemaAgreement ¶ added in v0.28.0
func (db *DB) GroupSchemaAgreement(sampleMax, k int) GroupSchemaAgreement
GroupSchemaAgreement reports how well per-segment derivations agree with the table-wide one.
func (*DB) GroupSchemaChanges ¶ added in v0.29.13
func (db *DB) GroupSchemaChanges() []GroupSchemaChange
GroupSchemaChanges returns the committed-group change log: the moments the accelerator adopted a different set of secondary schemas, with each change's diff and reason. Reads persisted state only (no sampling), so it is available at READ.
func (*DB) GroupSchemaDrift ¶ added in v0.28.0
func (db *DB) GroupSchemaDrift() GroupSchemaDrift
GroupSchemaDrift reports how the derived groups have moved across retained derivations.
func (*DB) GroupSchemaLastAgreement ¶ added in v0.29.13
func (db *DB) GroupSchemaLastAgreement() (GroupSchemaLastAgreement, bool)
GroupSchemaLastAgreement returns the last persisted per-segment agreement result, if any.
func (*DB) GroupSchemas ¶ added in v0.28.0
func (db *DB) GroupSchemas(sampleMax, k int) GroupSchemaInfo
GroupSchemas derives and reports candidate group schemas: sets of attributes the base schema does not carry which are present or absent together, and could therefore be stored columnar for the ads that have them without costing a slot in the ads that do not. Report-only.
func (*DB) GroupStatsAll ¶ added in v0.28.0
func (db *DB) GroupStatsAll(groupAttr string, aggAttrs []string) ([]collections.GroupStats, bool)
GroupStatsAll is GroupStatsConstraint over every row, for a constraint the caller has established is match-all. Declines on a chained table for the reason above.
func (*DB) GroupStatsConstraint ¶ added in v0.28.0
func (db *DB) GroupStatsConstraint(constraint, groupAttr string, aggAttrs []string) ([]collections.GroupStats, bool)
GroupStatsConstraint answers a per-group record count plus the aggregate inputs for each aggAttr over the rows matching constraint, via the columnar accelerator, or ok=false so the caller scans. This is the mutable-table half of what ArchiveTable already exposes; dbrpc's aggregate routes a grouped COUNT(*)/MIN/MAX/SUM/AVG through GroupedFromColumns for both.
A CHAINED table declines. The columnar paths read a record's own columns and know nothing about parentKeyFor/mergeParent, so on a chained table a group or aggregate attribute inherited from the parent would be missing from the column and the answer would differ from the scan's. (The same is true of the existing CountConstraint fast path, which does not guard it; chaining has no production consumer today -- only classad's own tests set IsStructural -- so that is a latent gap rather than a live bug.)
func (*DB) HotAttrs ¶
HotAttrs returns the current hot attributes (front-loaded in each ad's hot header for cheap access).
func (*DB) InMemory ¶ added in v0.10.2
InMemory reports whether the table's data lives only in RAM (no on-disk backing), i.e. it was opened with an empty Dir. Such a table is not recovered across restarts.
func (*DB) IndexSizes ¶
func (db *DB) IndexSizes() IndexSizes
IndexSizes reports each configured index's resident bytes (with human/auto provenance) against the live data bytes -- the memory cost of indexing.
func (*DB) IndexedAttrs ¶
IndexedAttrs returns the currently-indexed attribute names, split into categorical (string equality/membership) and value (numeric + range) indexes.
func (*DB) InstanceID ¶
InstanceID is this open's identity (fresh each Open). Equal in spirit to ID until HA replicas exist, when replicas of one DB share ID but differ by InstanceID.
func (*DB) Keys ¶
Keys returns every committed key at a consistent snapshot, in no particular order. Useful for administrative enumeration and for a replica that must clear its keyspace before a full re-sync (see the leader-follower replicator).
func (*DB) KeysWhere ¶ added in v0.20.4
KeysWhere returns an iterator over the storage keys of every row whose ad matches constraint. Unlike a query, it yields the DB keys themselves, so a caller can address matched rows for UPDATE/DELETE without depending on any self-reported key attribute in the ad. The constraint is parsed eagerly (a parse error returns before any scan); the scan is lazy and stops early if the caller's yield returns false. It matches against decoded ads (a full scan, like DeleteWhere).
func (*DB) Len ¶
Len returns the number of committed ads (including structural parent-only ads of a chained collection, which Query hides -- see Collection.Len / Chained).
func (*DB) LookupClassAd ¶
LookupClassAd returns the committed ad for key (the hash table, outside any transaction), or (nil, false).
func (*DB) LookupClassAdRedacted ¶ added in v0.28.0
LookupClassAdRedacted is LookupClassAd for a caller not entitled to sealed values; see QueryRedacted.
func (*DB) Maintain ¶
func (db *DB) Maintain(opts MaintainOptions)
Maintain runs one self-tuning pass: it auto-tunes indexes (adds demand-driven ones, trims auto indexes over the memory budget, never touches human-created ones), refreshes the hot-attribute set, and optionally retrains the compression dictionary. Index/hot changes are persisted. Synchronous; a server drives it on a schedule (see dbrpc.Server.StartMaintenance).
func (*DB) Match ¶
Match returns the ads that symmetrically match job (bilateral Requirements), pushed down to the store. For the negotiator's pick-best pattern, prefer MatchSorted.
func (*DB) MatchSorted ¶
MatchSorted returns job's matches ranked by the job's Rank, best first, at most limit (<=0 = all) -- the negotiator resource-request path, with the store's deferred materialization so only the returned top-N are built.
func (*DB) MatchSortedRanked ¶
func (db *DB) MatchSortedRanked(job *classad.ClassAd, limit int) []RankedMatch
MatchSortedRanked is MatchSorted that also returns each match's Rank. The job's Requirements prunes candidate slots via any covering index (the matchmaking pushdown) rather than bilaterally evaluating every slot.
func (*DB) MatchSortedRankedFiltered ¶
func (db *DB) MatchSortedRankedFiltered(job *classad.ClassAd, targetConstraint string, limit int) ([]RankedMatch, error)
MatchSortedRankedFiltered is MatchSortedRanked restricted to slots that also satisfy targetConstraint (a MATCH's WHERE TARGET / NOPREEMPT filter over the resource ad). The constraint's index probes narrow the candidate scan (pushdown) and it is re-checked on each matched slot. An empty constraint is exactly MatchSortedRanked.
func (*DB) NumStats ¶ added in v0.25.2
NumStats computes the aggregate inputs for attr over the records matching constraint, via the columnar scan. ok=false means the columnar path cannot serve it and the caller should scan.
func (*DB) OpStats ¶ added in v0.11.0
OpStats returns the store's operational timing counters (see the OpStats type).
func (*DB) Ordered ¶
Ordered iterates one partition of the index-th configured ordered index in sort order (Config.Ordered), yielding each member ad with a resume cursor and its cluster signature (for run-length folding into resource-request lists). partition selects the run (e.g. an Owner); it is ignored for an index with no Partition. A zero resume starts at the beginning. The snapshot is O(1) and stable under concurrent churn.
func (*DB) OrderedRedacted ¶ added in v0.28.0
OrderedRedacted is Ordered for a caller NOT entitled to sealed values; see QueryRedacted.
func (*DB) Put ¶ added in v0.8.0
Put inserts or replaces the ad at key in its own optimistic transaction, retrying on a write-write conflict (another committer touched key since this attempt's snapshot) until it lands or maxWriteAttempts is exhausted. A blind overwrite that loses the optimistic race would otherwise be silently dropped; Put is the retrying convenience for the common single-ad upsert so callers do not each reimplement the loop. Equivalent to Begin + NewClassAd + Commit.
func (*DB) Query ¶
Query returns the committed ads matching the constraint expression (an "old ClassAd" boolean expression over each ad, e.g. `JobStatus == 2 && Owner == "alice"`). The store pushes the filter down -- indexed constraints visit only candidates -- so this is far cheaper than ForEach + client-side filtering. Errors only on a malformed constraint.
func (*DB) QueryAsOf ¶ added in v0.12.0
QueryAsOf runs a point-in-time ("AS OF") query: it returns the ads matching the constraint as they were at time t. It errors on a malformed constraint, when time travel is not enabled on this table, or when t is older than the retained window.
func (*DB) QueryProject ¶
QueryProject returns, for each ad matching the constraint, just the named attributes' values (aligned with attrs), read wire-native where possible so an aggregate or projection does not pay the full-ad decode Query costs. The yielded slice is reused across iterations; copy any value to retain it past the next step. Errors only on a malformed constraint.
func (*DB) QueryRaw ¶ added in v0.8.0
QueryRaw yields each matching ad as a collections.RawAd -- the wire-form attribute strings decoded straight from the stored representation with no AST, for a persistent (inline) store as well as an in-memory one -- so a whole-ad result set can be relayed without materializing and re-encoding each ad. Errors only on a malformed constraint.
func (*DB) QueryRawProjected ¶ added in v0.16.3
func (db *DB) QueryRawProjected(constraint string, projection []string, redact bool) (iter.Seq[collections.RawAd], error)
QueryRawProjected is QueryRaw restricted to the projected attribute names, applied inside the collection's decode walk: a non-projected attribute is skipped before any name resolution or value rendering, and a hot-header- covered projection is served from the hot header alone (see collections.ScanRawProjected). redact additionally strips private attributes. An empty projection means no attribute filter.
func (*DB) QueryRawProjectedFromSeq ¶ added in v0.29.5
func (db *DB) QueryRawProjectedFromSeq(constraint string, projection []string, after SeqCursor, limit int) (iter.Seq[collections.RawAd], *SeqPage, error)
QueryRawProjectedFromSeq is QueryRawProjected a page at a time: it yields at most limit matching ads after the cursor and reports where to resume.
The order is the store's commit sequence, which is what makes the cursor stable — see collections.QueryRawFromSeq. A caller pages by handing back the previous page's Next until More is false. The zero cursor starts at the beginning.
This is for a mutable table. An archive is already ordered newest-first with its limit pushed down, so a "last K" read there needs no cursor at all.
func (*DB) QueryRawProjectedRefs ¶ added in v0.24.0
func (db *DB) QueryRawProjectedRefs(constraint string, projection []string, redact bool) (iter.Seq[collections.RawAd], error)
QueryRawProjectedRefs is QueryRawProjected that also carries the attributes the projected expressions reference, transitively, so each yielded ad EVALUATES self-contained (collections chaseRefs).
The distinction matters whenever a projected attribute holds an expression rather than a literal, which is the norm for HTCondor data -- Requirements, Rank and friends are expressions over sibling attributes. Projecting to exactly the requested names drops those siblings, so the expression evaluates to undefined at the far end: asking for Requirements alone answers undefined where asking for the whole ad answers true.
QueryRawProjected remains the right call for a relay that must reproduce HTCondor's query protocol, which specifies exactly the requested attributes and nothing more. Use this one when the recipient is going to EVALUATE what it receives.
func (*DB) QueryRawProjectedRefsStats ¶ added in v0.29.1
func (db *DB) QueryRawProjectedRefsStats(constraint string, projection []string, redact bool, stats *collections.ScanStats) (iter.Seq[collections.RawAd], error)
QueryRawProjectedRefsStats is QueryRawProjectedRefs that also fills stats (may be nil) with the per-scan work breakdown, for EXPLAIN ANALYZE. An empty/true constraint takes the no-WHERE scan path and leaves stats zero.
func (*DB) QueryRawRedacted ¶ added in v0.16.3
QueryRawRedacted is QueryRaw with private (secret) attributes stripped inside the collection's decode walk -- an unprivileged consumer's whole-ad query pays no per-attribute re-classification and never renders a private value (see collections.ScanRawRedacted).
func (*DB) QueryRawWire ¶ added in v0.16.4
func (db *DB) QueryRawWire(constraint string, projection []string, redact bool) (iter.Seq[[]byte], error)
QueryRawWire yields each matching ad as a self-contained WIRE-FORM ROW (an inline-names subset ad assembled by slice copies -- see collections.ScanRawWire): the relay form for shipping ads to a remote consumer with the old-ClassAd render deferred to that consumer's client edge. projection restricts the entries (empty = whole ad); redact strips private attributes at the source. At-rest-encrypted values are opened during assembly (the consumer holds no data key).
Only a persistent (inline) store can serve wire rows. An in-memory table returns ErrRawWireUnsupported rather than an empty sequence: a relay scan that yields nothing is indistinguishable from a query that matched nothing, so returning one would turn every RAM-table query into a silent empty result at the consumer.
func (*DB) QueryRedacted ¶ added in v0.28.0
QueryRedacted is Query for a caller NOT entitled to sealed values: the ads are decoded with no key, so a sealed attribute arrives undefined rather than opened. See collections.Collection.QueryRedacted.
Query decodes with the table's key and leaves it to the serializer to drop private attributes, which means an unprivileged reader's secret is decrypted in this process and then filtered on the way out.
func (*DB) RecordDemand ¶
RecordDemand notes the attributes constraint filters on (for index suggestions) without scanning. Use it for filters applied outside the normal Query path -- e.g. a MATCH's resource-side WHERE TARGET constraint -- so those attributes still surface in SuggestIndexes. A malformed constraint is ignored.
func (*DB) RefreshHotSet ¶
RefreshHotSet recomputes the hot set as the topN most frequent attributes from a sample of up to sampleMax live ads, returning how many were chosen. The resulting hot set is persisted.
func (*DB) Reindex ¶
func (db *DB) Reindex()
Reindex rebuilds all configured indexes from the live ads.
func (*DB) ReschemaScan ¶ added in v0.25.2
ReschemaScan derives a new schema from a fresh sample and rebuilds every sealed segment's columnar block against it, replacing the one pinned at first enable. Heavy -- a block per sealed segment is re-encoded and re-persisted -- and never done by routine maintenance, which deliberately keeps the schema stable. Returns false if there was nothing to sample or the accelerator cannot run here.
func (*DB) Restore ¶ added in v0.7.0
Restore replaces the entire DB with the contents of the snapshot in r. It holds the DB-wide lock exclusively: all writers are blocked and the truncate+reload is atomic. An encrypted snapshot is opened with this DB's pool keys against the snapshot's embedded master envelope, so a snapshot taken by any DB sharing a pool key can be restored.
func (*DB) RestoreWith ¶ added in v0.7.0
func (db *DB) RestoreWith(r io.Reader, keys SnapshotKeys) error
RestoreWith is Restore using any level of the key hierarchy (see SnapshotKeys).
func (*DB) RestoreWithBackupKey ¶ added in v0.7.0
RestoreWithBackupKey is Restore using an explicitly-provided backup key (from BackupKey / the backup.key command) to unwrap the snapshot key, instead of the DB's pool keys. This recovers an encrypted backup independently of the pool keys -- the escrow path for disaster recovery when the original pool keys are unavailable.
func (*DB) RetrainDict ¶
RetrainDict trains a fresh ZSTD dictionary from a sample of up to sampleMax live ads, switches new writes to it, and recompresses existing records under it. It returns the new dictionary's size in bytes. This is what turns on (or refreshes) compression for a collection that started with the identity codec. See collections.Collection.RetrainDict.
func (*DB) Rewrite ¶
Rewrite re-encodes every live ad with the current hot set (so a changed hot set applies to existing ads) and force-compacts, returning the number of ads rewritten. A maintenance operation -- see collections.Collection.Rewrite.
func (*DB) SchemaFit ¶ added in v0.25.2
func (db *DB) SchemaFit(sampleMax int) ([]SchemaFieldFit, int)
SchemaFit measures the current derived schema against a fresh sample, reporting each field's escape rate -- how often its value is not in the fixed slot, and how much of that is the attribute simply being absent. This is how you tell whether the schema still matches the data (see collections.Collection.SchemaFit). Returns nil when the accelerator is not enabled.
func (*DB) SchemaScanInfo ¶ added in v0.23.1
func (db *DB) SchemaScanInfo() SchemaScanInfo
SchemaScanInfo reports the columnar (adschema) accelerator's state: whether it is enabled (so a numeric COUNT(*) WHERE routes to the columnar fast path), its hot columns, and segment coverage.
func (*DB) SetEncryptedAttrs ¶ added in v0.7.0
SetEncryptedAttrs replaces the explicit set of attributes encrypted at rest (the human-toggled set; private attributes are always encrypted). It errors if encryption is disabled or a named attribute is indexed. The policy is persisted so it survives a restart and, in HA, lets a follower converge on reload. New writes seal the new set; existing records re-seal when next rewritten (compaction/Rewrite).
func (*DB) SetTimeTravel ¶ added in v0.12.0
SetTimeTravel enables (with a positive maxDistance), retunes, or disables (maxDistance <= 0) point-in-time queries for this table, and persists the setting so it survives a restart. A zero checkpoint interval uses the collections default (1 minute). Enabling is not retroactive -- only changes from now on are travelable.
func (*DB) SidecarSizes ¶ added in v0.29.0
func (db *DB) SidecarSizes() SidecarSizes
SidecarSizes reports the sealed-segment sidecar index bytes (mmap-backed, evictable), matching ArchiveTable.SidecarSizes. A mutable table has these sidecars too -- only the accessor was missing, so its .stats could not report the on-disk footprint an archive's could.
func (*DB) Snapshot ¶ added in v0.7.0
Snapshot writes a consistent backup of every ad to w. It holds the DB-wide lock shared, so it is consistent against writers without blocking other readers or snapshots.
func (*DB) SnapshotWithKey ¶ added in v0.7.0
SnapshotWithKey is Snapshot that also returns the per-backup (snapshot) key it minted -- the finest-grained escrow key, which decrypts THIS backup only (unlike the backup key, which decrypts all of them). Returns nil when encryption is disabled.
func (*DB) StaleIndexSegments ¶ added in v0.29.0
StaleIndexSegments reports how many sealed segments still carry an index built under an older configuration, and how many are sealed in total, matching ArchiveTable.StaleIndexSegments.
func (*DB) StartMaintenance ¶
func (db *DB) StartMaintenance(interval time.Duration, opts MaintainOptions) (stop func())
StartMaintenance starts a background goroutine that runs Maintain with the given options every interval, returning a stop function. Prefer the catalog-wide dbrpc.Server.StartMaintenance, which also covers tables created later.
func (*DB) Stats ¶
Stats returns a snapshot of the store's storage (ad count, segment/arena/dead bytes) for observability.
func (*DB) StopMaintenance ¶ added in v0.30.4
func (d *DB) StopMaintenance()
StopMaintenance asks any maintenance pass running on this table to stop at its next safe boundary, so a caller closing the server does not wait out a whole in-flight pass. See collections.Collection.StopMaintenance.
func (*DB) SuggestDrops ¶
func (db *DB) SuggestDrops(sampleMax int) []DropSuggestion
SuggestDrops recommends indexes to drop (unused or low-cardinality) from observed demand and a sample of up to sampleMax live ads.
func (*DB) SuggestIndexes ¶
func (db *DB) SuggestIndexes(sampleMax int) []collections.IndexSuggestion
SuggestIndexes samples the store and returns attributes that queries filter on but are not yet indexed (advisory; a server may log or auto-apply them).
func (*DB) TimeTravel ¶ added in v0.12.0
TimeTravel reports the table's current point-in-time settings and whether enabled.
func (*DB) TopK ¶ added in v0.29.1
func (db *DB) TopK(constraint string, attrs []string, orderAttr string, desc bool, k int) ([][]classad.Value, error)
TopK is ArchiveTable.TopK for a mutable table.
func (*DB) TopKStats ¶ added in v0.29.4
func (db *DB) TopKStats(constraint string, attrs []string, orderAttr string, desc bool, k int) ([][]classad.Value, collections.ScanStats, error)
TopKStats is DB.TopK that also returns the cutoff scan's ScanStats, for EXPLAIN ANALYZE.
func (*DB) TrackedKeys ¶ added in v0.30.0
TrackedKeys reports how many keys the delta tracker holds in memory (0 when delta records are off). It is the feature's one new in-memory structure, so it is the number to watch.
func (*DB) Truncate ¶ added in v0.7.0
func (db *DB) Truncate()
Truncate removes every ad from the DB, atomically against all writers (it takes the DB-wide lock exclusively). A DAEMON-level operation -- the first half of a restore, or a deliberate wipe.
func (*DB) UpdateOld ¶ added in v0.8.0
UpdateOld ingests an ad from old-ClassAd wire text under key, skipping the AST build -- the wire-native ingest path (as collections.UpdateOld). It writes through the same shard storage, change log, and watch feed as a committed Put, but bypasses the optimistic-concurrency layer (last-writer-wins), which suits high-rate single-key upserts where per-key write-write races do not occur (a collector re-advertising its own ad).
An encrypted store no longer takes the AST path wholesale. The wire-native encoder cannot seal, so the ingest defers to the sealing path for an ad that HAS something to seal and streams the rest (see collections.encodeOld) -- for job and history ads, nearly all of them.
func (*DB) UpdateOldBatch ¶ added in v0.8.0
UpdateOldBatch ingests many ads (key + old-ClassAd text) in one shard-commit batch -- the wire-native bulk ingest, so a burst of upserts costs one commit instead of one per ad. Bypasses the optimistic-concurrency layer (last-writer-wins) like UpdateOld. An encrypted store keeps the batch: sealing is decided per ad inside the ingest (see UpdateOld), so one commit still covers the whole batch.
func (*DB) WatchCursor ¶ added in v0.7.0
Watch streams changes committed after the given cursor (nil = from now). Cancel via ctx. Events arrive in commit order per key; see collections.Watch for the delivery and coalescing guarantees. WatchCursor returns an opaque cursor at the current head of the change log, so Watch(ctx, cursor) then streams only subsequent changes (no replay of current contents). Cheap; requires the store was opened with watch history.
func (*DB) WatchRedacted ¶ added in v0.28.0
WatchRedacted is Watch for a watcher NOT entitled to sealed values: each event's ad is decoded with no key, so a sealed attribute arrives undefined. A watch has the longest exposure of any read path -- it streams ads continuously. See collections.Collection.WatchRedacted.
type DropSuggestion ¶
type DropSuggestion = collections.DropSuggestion
Diagnostic and management types, re-exported from collections for callers that only import db.
type ExporterDef ¶ added in v0.14.0
type ExporterDef struct {
Name string `json:"name"`
// Kind selects the exporter implementation (e.g. "kafka"). The catalog does not
// enumerate valid kinds; an unknown kind is simply ignored by any exporter that does
// not recognize it.
Kind string `json:"kind"`
// Config is the exporter's own configuration, opaque to the catalog. It is stored and
// returned verbatim.
Config json.RawMessage `json:"config,omitempty"`
}
ExporterDef defines an external sink that mirrors this database's change stream elsewhere (for example, a Kafka change-data exporter that federates several instances through a broker). The catalog is a passive registry: it stores the definition and an opaque resume-state blob and never runs, interprets, or validates Config -- that is the job of the out-of-process exporter (which reads these over dbrpc). Keeping the DB free of any sink-specific dependency is deliberate.
type GroupAgreementItem ¶ added in v0.29.13
type GroupAgreementItem = collections.GroupAgreementItem
type GroupCol ¶ added in v0.16.8
GroupCol is one GROUP BY column for a (possibly bucketed) aggregate: the attribute Attr, optionally floored into fixed-width buckets. BucketWidth == 0 groups by the raw attribute value; BucketWidth > 0 (seconds) groups by the epoch-aligned bucket floor(number(Attr)/BucketWidth)*BucketWidth, and a row whose Attr is not a finite number drops out of the result. This is the shared shape a bucketed aggregate and a bucketed materialized view can both group by.
type GroupSchemaAgreement ¶ added in v0.28.0
type GroupSchemaAgreement = collections.GroupSchemaAgreement
type GroupSchemaChange ¶ added in v0.29.13
type GroupSchemaChange = collections.GroupSchemaChange
type GroupSchemaDelta ¶ added in v0.29.13
type GroupSchemaDelta = collections.GroupSchemaDelta
type GroupSchemaDrift ¶ added in v0.28.0
type GroupSchemaDrift = collections.GroupSchemaDrift
type GroupSchemaEntry ¶ added in v0.28.0
type GroupSchemaEntry = collections.GroupSchemaEntry
type GroupSchemaInfo ¶ added in v0.28.0
type GroupSchemaInfo = collections.GroupSchemaInfo
GroupSchemaInfo, GroupSchemaEntry, GroupSchemaDrift and GroupSchemaAgreement are re-exported so a caller can read the group-schema reports without importing collections.
type GroupSchemaLastAgreement ¶ added in v0.29.13
type GroupSchemaLastAgreement = collections.GroupSchemaLastAgreement
type GroupStatsSource ¶ added in v0.28.0
type GroupStatsSource interface {
// GroupStatsConstraint answers the grouped stats for rows matching constraint.
GroupStatsConstraint(constraint, groupAttr string, aggAttrs []string) ([]collections.GroupStats, bool)
// GroupStatsAll answers them over every row, for a constraint the caller has established is
// match-all (the predicate analysis has no predicate to work with there).
GroupStatsAll(groupAttr string, aggAttrs []string) ([]collections.GroupStats, bool)
}
GroupStatsSource is a table that can answer a grouped columnar aggregate: an archive or a mutable table. Both back onto the same collections primitive, so both get the same fast path and the same row-shaping rather than a copy per table type that could drift.
type IndexSize ¶
type IndexSize = collections.IndexSize
Diagnostic and management types, re-exported from collections for callers that only import db.
type IndexSizes ¶
type IndexSizes = collections.IndexSizes
Diagnostic and management types, re-exported from collections for callers that only import db.
type IndexSuggestion ¶
type IndexSuggestion = collections.IndexSuggestion
Diagnostic and management types, re-exported from collections for callers that only import db.
type KEK ¶ added in v0.7.0
KEK is a pool / signing key that can wrap the DB master key. Re-exported so db callers need not import crypt.
type MaintainOptions ¶
type MaintainOptions struct {
// SampleMax caps the ads sampled for index tuning, hot-set frequency, and dictionary
// training. Default 4096.
SampleMax int
// HotTopN refreshes the hot set to the topN most common attributes; 0 disables it.
HotTopN int
// Retrain retrains the ZSTD dictionary (recompacting + reindexing). Expensive on a
// large store, so a server may run it on a longer cadence than the rest.
Retrain bool
// CompactInterval is the cadence of a lightweight, standalone compaction pass
// (dbrpc.Server.StartMaintenance) that reclaims dead space independently of the
// expensive retrain-dominated Maintain pass. Compaction is cheap and self-limiting
// (it unlinks fully-dead segments for free and only recompacts shards past the
// dead-byte threshold), so it runs far more often than Maintain -- essential for a
// high-churn table (e.g. a collector re-advertising every few minutes) where dead
// space would otherwise grow unbounded between the throttled retrain passes. 0 uses
// a default cadence; a negative value disables the standalone compaction pass.
CompactInterval time.Duration
// MinIndexDemand is the minimum observed query count for the auto-tuner to add an
// index; 0 leaves index auto-tune off (no demand-driven adds).
MinIndexDemand int64
// IndexBudgetHighFrac / IndexBudgetLowFrac / IndexBudgetSlackBytes bound auto-created
// index memory as a fraction of the live data bytes (see collections.AutoTuneOptions);
// 0 high frac disables the budget (auto indexes grow unbounded).
IndexBudgetHighFrac float64
IndexBudgetLowFrac float64
IndexBudgetSlackBytes int64
// ArchiveMerge tunes the merge pass run over each archive table. The zero value uses
// the policy defaults, which target a segment count well clear of the process's mapping
// budget. ArchiveMergeDisabled turns the archive side of maintenance off entirely.
ArchiveMerge MergeOptions
ArchiveMergeDisabled bool
// ArchiveIndexMinDemand is the minimum observed query demand for the auto-tuner to add
// an index to an ARCHIVE table; 0 leaves archive index auto-tune off.
//
// It is separate from MinIndexDemand, and off by default, because the two costs are not
// comparable. Indexing a mutable table is bounded by its live size; indexing an archive
// is paid for by decompressing history, so the threshold wants to be set against
// ArchiveConfig.IndexBackfillBytes -- how far back a new index is actually carried --
// rather than inherited from the mutable side.
//
// Auto-DROP is deliberately not offered for archives at any setting. Dropping an index
// is nearly free and re-adding it costs a backfill, so a wrong drop is expensive and
// asymmetric in the direction that punishes acting on a weak signal.
ArchiveIndexMinDemand int64
// ArchiveSchemaScanHotTopN, when > 0, builds the per-segment columnar accelerator on
// ARCHIVE tables too, keeping the topN most query-read numeric fields uncompressed.
//
// Separate from SchemaScanHotTopN, and off by default, for the same reason
// ArchiveIndexMinDemand is: the first build reads every sealed record of the whole
// history once, so turning it on is a deliberate decision about an existing deployment
// rather than something an upgrade should start doing. Once built it is cheap -- an
// archive's segments are immutable, so a block is never invalidated, and later passes
// only cover newly-sealed segments.
//
// The payoff is confined to what the accelerator serves: single-int-field COUNT
// comparisons (25-65x on a 50k-record archive). Other aggregates -- MAX, SUM, GROUP BY
// -- take their existing paths and are unaffected.
ArchiveSchemaScanHotTopN int
// SchemaScanHotTopN, when > 0, builds/refreshes the per-segment adschema columnar
// accelerator (used by CountConstraint's fast path) keeping the topN most query-read
// numeric fields uncompressed. The first pass chooses a stable schema and hot set from
// the current sample and demand; later passes only extend coverage to newly-sealed
// segments, so a table opts into the columnar count path once maintenance runs and stays
// covered as it grows. 0 leaves it off.
SchemaScanHotTopN int
}
MaintainOptions configures one maintenance pass (DB.Maintain).
type MatchExplain ¶
type MatchExplain = collections.MatchExplain
Diagnostic and management types, re-exported from collections for callers that only import db.
type MergeOptions ¶ added in v0.24.1
type MergeOptions = collections.MergeOptions
MergeOptions tunes a merge pass over an archive; see collections.MergeOptions. The zero value is usable and fills in defaults.
type NumStats ¶ added in v0.25.2
type NumStats = collections.NumStats
NumStats is one columnar pass's numeric aggregate inputs (see collections.NumStats).
type OldAdText ¶ added in v0.8.0
OldAdText is one keyed ad in old-ClassAd wire text, for UpdateOldBatch.
type OpStat ¶ added in v0.11.0
type OpStat = collections.OpStat
OpStat is one operation's cumulative call count and total wall-nanoseconds.
type OpStats ¶ added in v0.11.0
type OpStats struct {
collections.OpStats
SnapshotLock OpStat `json:"snapshotLock"`
// SnapshotLockWait is the time committers spent BLOCKED on the shared snapshot lock (stalled
// behind an exclusive Truncate/Restore) -- the waiters' side of SnapshotLock's holder time.
SnapshotLockWait OpStat `json:"snapshotLockWait"`
}
OpStats is a snapshot of the store's operational timing counters -- where callers spent time blocked in, or holding, each stall point (see collections.OpStats) -- plus the DB-wide snapshot lock's exclusive-hold time (Truncate/Restore). Every value is a monotonic cumulative total; a scraper derives rate and mean latency from deltas.
type OrderCursor ¶
type OrderCursor = collections.OrderCursor
OrderSpec, SortKey, OrderedAd, OrderCursor configure and drive maintained ordered indexes (the schedd priority-queue / resource-request-list pattern). Re-exported from collections so callers of db need not import it.
type OrderSpec ¶
type OrderSpec = collections.OrderSpec
OrderSpec, SortKey, OrderedAd, OrderCursor configure and drive maintained ordered indexes (the schedd priority-queue / resource-request-list pattern). Re-exported from collections so callers of db need not import it.
type OrderedAd ¶
type OrderedAd = collections.OrderedAd
OrderSpec, SortKey, OrderedAd, OrderCursor configure and drive maintained ordered indexes (the schedd priority-queue / resource-request-list pattern). Re-exported from collections so callers of db need not import it.
type ProbeExplain ¶ added in v0.21.2
type ProbeExplain = collections.ProbeExplain
Diagnostic and management types, re-exported from collections for callers that only import db.
type QueryExplain ¶
type QueryExplain = collections.QueryExplain
Diagnostic and management types, re-exported from collections for callers that only import db.
type RankedMatch ¶
type RankedMatch = collections.RankedMatch
RankedMatch is a matched ad with the job's Rank of it.
type Retention ¶ added in v0.19.0
type Retention = collections.Retention
Diagnostic and management types, re-exported from collections for callers that only import db.
type SchemaFieldFit ¶ added in v0.25.2
type SchemaFieldFit = collections.SchemaFieldFit
Diagnostic and management types, re-exported from collections for callers that only import db.
type SchemaScanField ¶ added in v0.25.2
type SchemaScanField = collections.SchemaScanField
Diagnostic and management types, re-exported from collections for callers that only import db.
type SchemaScanGroup ¶ added in v0.29.13
type SchemaScanGroup = collections.SchemaScanGroup
Diagnostic and management types, re-exported from collections for callers that only import db.
type SchemaScanInfo ¶ added in v0.23.1
type SchemaScanInfo = collections.SchemaScanInfo
Diagnostic and management types, re-exported from collections for callers that only import db.
type SeqCursor ¶ added in v0.29.5
type SeqCursor = collections.SeqCursor
SeqCursor and SeqPage re-export the collections types so a caller of the DB API does not have to import the storage layer to paginate.
type SeqPage ¶ added in v0.29.5
type SeqPage = collections.SeqPage
SeqCursor and SeqPage re-export the collections types so a caller of the DB API does not have to import the storage layer to paginate.
type SidecarSizes ¶ added in v0.19.0
type SidecarSizes = collections.SidecarSizes
Diagnostic and management types, re-exported from collections for callers that only import db.
type SnapshotKeys ¶ added in v0.7.0
SnapshotKeys carries any ONE level of the key hierarchy sufficient to decrypt a snapshot. A snapshot is decryptable with, in decreasing specificity: the per-backup SnapshotKey (seals the body directly), the BackupKey (unwraps the snapshot key), the MasterKey (derives the backup key), or PoolKeys (open the embedded master envelope). Restore uses the most specific field provided; all are optional (plain Restore uses the DB's own pool keys). This lets an operator escrow whichever level fits their recovery model and decrypt a backup even without the original pool keys or a running daemon.
type SortKey ¶
type SortKey = collections.SortKey
OrderSpec, SortKey, OrderedAd, OrderCursor configure and drive maintained ordered indexes (the schedd priority-queue / resource-request-list pattern). Re-exported from collections so callers of db need not import it.
type Stats ¶
type Stats = collections.Stats
Diagnostic and management types, re-exported from collections for callers that only import db.
type TableOptions ¶ added in v0.10.1
type TableOptions struct {
// InMemory makes the table's data live only in RAM even in a persistent
// catalog: no <dir>/tables/<name> directory is created, so the table is not
// recovered across restarts (it reappears empty). Useful for high-churn,
// reconstructible data (e.g. frequently-replaced ads) that is not worth the
// disk I/O of persistence. In an already-in-memory catalog it is a no-op.
InMemory bool
}
TableOptions tunes how a table is created. The zero value creates a normal table (persistent when the catalog has a directory).
type Txn ¶
type Txn struct {
// contains filtered or unexported fields
}
Txn is an independent optimistic transaction. Operations are buffered and applied at Commit under snapshot-isolation OCC. A *Txn is not safe for concurrent use by multiple goroutines; independent transactions are.
func (*Txn) Abort ¶
func (t *Txn) Abort()
Abort discards the transaction's buffered operations. Nothing is written.
func (*Txn) Commit ¶
Commit applies the buffered operations. It returns a *ConflictError if any key was modified by another committer since this transaction's snapshot (the non-conflicted operations still committed), or nil on full success.
func (*Txn) CommitNondurable ¶
CommitNondurable is Commit that defers the disk durability sync (classad_log.h CommitNondurableTransaction): the writes are visible immediately but their flush is batched to a later durable commit. On an in-memory DB it is identical to Commit.
func (*Txn) DeleteAttribute ¶
DeleteAttribute removes one attribute of key (classad_log.h LogDeleteAttribute). A no-op if key or the attribute is absent.
func (*Txn) DestroyClassAd ¶
DestroyClassAd removes key (classad_log.h LogDestroyClassAd).
func (*Txn) Has ¶ added in v0.29.11
Has reports whether key exists as the transaction sees it (its own buffered writes over the snapshot, exactly as LookupClassAd does) without decoding the stored record. Use it wherever only presence matters: LookupClassAd's decode of a wide ad dominates the cost of an ingest that asks the question once per operation.
func (*Txn) KeysWhere ¶ added in v0.24.0
KeysWhere returns the storage keys of the rows matching the constraint as the transaction sees them (DB.KeysWhere with the transaction's writes overlaid). It is what lets an UPDATE or DELETE inside a transaction address a row the transaction itself created. Errors only on a malformed constraint.
func (*Txn) LookupAttr ¶
LookupAttr returns the unparsed expression of one attribute as the transaction sees it (classad_log.h LookupInTransaction), or ("", false).
func (*Txn) LookupClassAd ¶
LookupClassAd returns key's ad as the transaction sees it: its own buffered writes (read-your-writes) merged over the snapshot (classad_log.h Lookup + the LookupInTransaction overlay in one call).
func (*Txn) NewClassAd ¶
NewClassAd stores ad under key (classad_log.h LogNewClassAd). An existing ad at key is replaced.
func (*Txn) NewClassAdOld ¶ added in v0.25.1
NewClassAdOld stores an ad supplied as old-ClassAd text under key, encoding it straight to the stored wire form instead of parsing it into a ClassAd for the commit to re-encode. The transaction is unaffected: the write is buffered, conflict-checked and committed exactly as NewClassAd's is.
It reports whether the wire-native path was taken. False means the caller must parse the text and use NewClassAd -- an encrypted store seals its values and the streaming encoder does not seal, and a few ad shapes (a repeated attribute name, an escape the fast lexer would read differently) defer to the reference parser by design.
func (*Txn) Query ¶ added in v0.24.0
Query returns the ads matching the constraint as the transaction sees them: the committed rows with the transaction's own buffered writes overlaid, so a query inside a transaction observes work the transaction has not committed yet. DB.Query cannot -- it reads the committed store -- which is why a caller that has staged writes and then wants to read them back must go through the transaction.
Isolation is snapshot plus read-your-writes: the committed half is read at the transaction's own snapshot sequence, the same one Get reads at, so a scan and a point lookup in one transaction agree and a concurrent commit is invisible to both. The cost is a full scan -- reading at a past sequence and overlaying by key both need the per-record walk the indexed query path skips. See collections.Txn.Query. Errors only on a malformed constraint.
func (*Txn) SetAttribute ¶
SetAttribute sets one attribute of key to the expression parsed from expr (classad_log.h LogSetAttribute) -- a read-modify-write within the transaction, so it composes with the transaction's own earlier writes to key. The ad is created if absent.
type UnappliedError ¶ added in v0.30.2
type UnappliedError struct {
Keys []string
// Reasons is parallel to Keys: why each one could not be composed ("delta-no-base",
// "delta-chain-broken", ...). Carried so a caller can log what happened rather than only
// that something did.
Reasons []string
}
UnappliedError reports writes that could not be composed at all, as distinct from a conflict. Re-applying the identical write cannot succeed -- the key is present but the record behind it will not come back -- so a caller should record these and make progress rather than retry. Treating them as conflicts made a tailer rewind and re-apply forever, escalating to full log replays that changed nothing. See collections.CommitResult.Unapplied.
func (*UnappliedError) Error ¶ added in v0.30.2
func (e *UnappliedError) Error() string
type View ¶ added in v0.12.0
type View struct {
// contains filtered or unexported fields
}
View is a live materialized view: its spec, the per-key contribution store, per-group accumulators, and an in-memory backing DB holding one rendered ad per group (so the view is queried like a table). All mutable state is guarded by mu.
func (*View) Backing ¶ added in v0.12.0
Backing returns the in-memory table holding the materialized rows (read-only in intent).
func (*View) Cursor ¶ added in v0.12.0
Cursor returns the last durable resume cursor (for persistence).
func (*View) LateDrops ¶ added in v0.15.0
LateDrops reports how many base rows were dropped because their time bucket was already sealed (out-of-window late data).
func (*View) Seal ¶ added in v0.16.0
Seal appends and evicts every bucket whose window has closed as of now (unix seconds) into the archive, advancing the watermark. It is what the periodic tick invokes; exported so an operator (or a test) can force sealing at a chosen instant.
func (*View) SealedQuery ¶ added in v0.16.0
SealedQuery streams this continuous aggregate's sealed (archived) rows matching the constraint. It returns an empty sequence for a gauge view or an in-memory continuous aggregate (no archive). Reads union this with the live backing so a view read returns the full series, not just the unsealed buckets.
func (*View) SeriesCount ¶ added in v0.12.0
SeriesCount returns the current number of groups (Prometheus series).
type ViewAggFunc ¶ added in v0.12.0
type ViewAggFunc string
ViewAggFunc is a delta-maintainable aggregate.
const ( ViewCount ViewAggFunc = "count" ViewSum ViewAggFunc = "sum" ViewAvg ViewAggFunc = "avg" )
type ViewGroupCol ¶ added in v0.12.0
type ViewGroupCol struct {
Attr string `json:"attr"`
Alias string `json:"alias"`
BucketWidth int64 `json:"bucketWidth,omitempty"` // seconds; 0 = raw value
}
ViewGroupCol is one GROUP BY column: the base-table attribute and the alias it is stored under in the materialized rows (e.g. attr "Owner", alias "label_owner"). When BucketWidth > 0, the column is time-bucketed: the (numeric, unix-seconds) attribute is floored to epoch-aligned buckets of that width, turning the view into a time series -- each interval is a distinct, permanent group rather than an overwritten gauge. This is the same shape as dbrpc.GroupCol.
type ViewMetric ¶ added in v0.12.0
type ViewMetric struct {
Func ViewAggFunc `json:"func"`
Arg string `json:"arg"`
Alias string `json:"alias"`
}
ViewMetric is one aggregate output: the function, its argument attribute ("*" for COUNT(*)), and the alias it is stored under (e.g. "metric_jobs").
type ViewSpec ¶ added in v0.12.0
type ViewSpec struct {
BaseTable string `json:"baseTable"`
Groups []ViewGroupCol `json:"groups"`
Metrics []ViewMetric `json:"metrics"`
// Cardinality is the hard maximum number of groups (output series). Exceeding it fails
// the build (CreateView error) or, at runtime, moves the view to the failed state.
Cardinality int `json:"cardinality"`
// SelectText is the original SELECT for display (.views); not used for execution.
SelectText string `json:"selectText"`
// Grace is seconds to wait after a time bucket's window closes before sealing it
// (so late-arriving base rows still land). Only meaningful for a continuous aggregate
// (a spec with a time-bucketed group column). 0 seals as soon as the window closes.
Grace int64 `json:"grace,omitempty"`
// Retention is seconds of sealed history to keep in the continuous aggregate's archive;
// 0 keeps all.
Retention int64 `json:"retention,omitempty"`
}
ViewSpec is a materialized view's persisted definition.
func (ViewSpec) IsContinuous ¶ added in v0.15.0
IsContinuous reports whether the view is a continuous aggregate (has a time bucket).
type ViewState ¶ added in v0.12.0
type ViewState uint8
ViewState is a view's lifecycle state.
const ( // ViewBuilding: the initial catch-up is in progress. ViewBuilding ViewState = iota // ViewActive: built and maintaining live. ViewActive // ViewStale: definition loaded but the base table is absent; will bind on the next // activation once the base table exists. ViewStale // ViewFailed: a runtime error (e.g. cardinality exceeded); the updater has stopped. ViewFailed )
type WatchEvent ¶
type WatchEvent struct {
Kind WatchKind
Key string
Ad *classad.ClassAd // nil for a delete
Cursor []byte
}
WatchEvent is one change to the log. Cursor resumes a watch just after this event (survives reconnect / restart, so no change is missed).
type WatchKind ¶
type WatchKind uint8
WatchKind classifies a watch event.
const ( // WatchUpsert: key was added or updated; Ad holds its new value. WatchUpsert WatchKind = iota // WatchDelete: key was removed; Ad is nil. WatchDelete // WatchReset: discard prior state and rebuild from the upserts that follow (the // initial full replay, or a re-sync after the cursor fell out of retention). Key // and Ad are empty. WatchReset // WatchSynced: end of the initial catch-up/replay; the watcher is now live. Cursor is // a durable resume point (and, after a Reset, the point at which a shadow rebuild goes // live). Key and Ad are empty. WatchSynced // WatchResync: the live stream fell behind; reconnect with the last persisted cursor // (which re-enters catch-up). Key and Ad are empty; no state is implied. WatchResync )
type Watcher ¶
type Watcher struct {
// contains filtered or unexported fields
}
Watcher is a watch whose readiness is signalled on a file descriptor -- for the C (DaemonCore) side, which registers the fd in its poll loop and, on wakeup, drains events with Next. Go writes a single wakeup byte when the queue goes non-empty (coalesced: the reader drains fully per wakeup). The fd is a pipe or eventfd the caller created and owns; Watcher never closes it.
func NewWatcher ¶
NewWatcher starts watching db, signalling notifyFD when events are queued. cursor nil starts from now.
func (*Watcher) Next ¶
func (w *Watcher) Next() (WatchEvent, bool)
Next dequeues one event without blocking. ok is false when the queue is empty; the caller (having been woken via the fd) drains until ok is false.
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
Package main builds the embedded ClassAd log as a C archive: it exports C symbols (cadb_*) mirroring HTCondor's classad_log.h, so a C++ interface (built on libcondor_utils) can sit on top.
|
Package main builds the embedded ClassAd log as a C archive: it exports C symbols (cadb_*) mirroring HTCondor's classad_log.h, so a C++ interface (built on libcondor_utils) can sit on top. |
|
Package replicate is the transport-neutral core for replicating a db table/archive's change stream into another store: an in-memory Change (a db.WatchEvent enriched with replication metadata) and an idempotent Sink that applies Changes into a *db.DB or *db.ArchiveTable.
|
Package replicate is the transport-neutral core for replicating a db table/archive's change stream into another store: an in-memory Change (a db.WatchEvent enriched with replication metadata) and an idempotent Sink that applies Changes into a *db.DB or *db.ArchiveTable. |