Files
mesh-controller/internal/inventory/data.go
T
jochen abf9125689 Measure each item its own way; never compare a partial size (hq ADR 0233)
A walk over a large library every hour loads the array that protects it. An item now says how it
is measured — a bounded daily walk, a dataset's counters, or its top level only — and a size that
is a lower bound is kept as such and never read as a shrink.
2026-10-06 17:00:22 +02:00

317 lines
12 KiB
Go

package inventory
import (
"context"
"errors"
"fmt"
"strings"
"time"
"github.com/jackc/pgx/v5"
"github.com/novox/mesh-controller/internal/catalogue"
)
// The data a module declares, as the mesh found it on each machine (novox/hq ADR 0233).
//
// The self-check composes what every machine declares and asks its backup holder what it measured;
// this keeps both. An irreplaceable item a machine no longer declares — its module unassigned, or no
// longer pulled in — is retired here, not forgotten: kept, with when and why, until a person deletes
// it through `cleanup delete`.
// DeclaredData is one item a machine's composition declares now.
type DeclaredData struct {
Module, Item, Class string
// Owned is whether it is in the module's own directory: only that is the mesh's to retire.
Owned bool
// Protection is backup, redundancy, both ("backup+redundancy") or none.
Protection string
}
// Redundancy is the redundant storage an item is on, as its machine's holder read it.
type Redundancy struct {
Kind string `json:"kind"`
Where string `json:"where"`
Healthy *bool `json:"healthy"`
Said string `json:"said"`
}
// Measurement is what a machine's backup holder said of one item. Nil fields were not measured.
type Measurement struct {
Path string
Size *int64
LastWrite *time.Time
MeasuredAt *time.Time
LastBackup *time.Time
// Error is what the holder could not measure, when it could not.
Error string
// Precision is what the size is, as the holder said it: "exact", "dataset …", "partial: …", "no size: …".
Precision string
Redundancy *Redundancy
}
// DataRecord is one item as the mesh keeps it.
type DataRecord struct {
Machine, Module, Item, Class, Path string
Owned bool
Protection string
FirstSeen, DeclaredAt time.Time
MeasureError string
Precision string
Redundancy *Redundancy
Size *int64
LastWrite, MeasuredAt, LastBackup *time.Time
RetiredAt *time.Time
RetiredWhy string
DeletedAt *time.Time
DeletedBy, DeletedWhy string
}
// Comparable is whether a size is fit to compare with another: exact, or a dataset's own counters.
func Comparable(precision string) bool {
return precision == "" || precision == "exact" || strings.HasPrefix(precision, "dataset ")
}
// Dataset is the ZFS dataset a size was read from, when it was: what several items on one dataset share.
func Dataset(precision string) string {
rest, ok := strings.CutPrefix(precision, "dataset ")
if !ok {
return ""
}
name, _, _ := strings.Cut(rest, ":")
return name
}
// Retired is whether the item is retired and not deleted.
func (r DataRecord) Retired() bool { return r.RetiredAt != nil && r.DeletedAt == nil }
// Key names it in one string, the way conditions and verbs do: <machine>/<module>/<item>.
func (r DataRecord) Key() string { return r.Machine + "/" + r.Module + "/" + r.Item }
// DataChange is what recording a machine's data changed: what it retired and what came back.
type DataChange struct {
Retired, Reenabled []DataRecord
}
// readingEvery is how often a measurement is kept as a reading: the shrink is read over days, and a
// row every five minutes would say the same thing sixty times an hour.
const readingEvery = 55 * time.Minute
// readingsKept is how long readings are kept.
const readingsKept = 90 * 24 * time.Hour
// dataKey is one item's identity on one machine.
type dataKey struct{ module, item string }
// RecordData keeps what one machine declares now and what its holder measured, at now.
//
// **Only from a composition that worked.** The caller passes the declared set of a machine whose
// composition succeeded; an item missing from it is then really no longer declared. An irreplaceable
// one is retired — never deleted, never forgotten — and anything else is forgotten, because what is
// rebuildable or a cache is not the mesh's to keep track of once its module is gone. An item declared
// again comes back from retirement as it was.
func (i *Inventory) RecordData(ctx context.Context, machine string, declared []DeclaredData,
measured map[string]map[string]Measurement, why string, now time.Time) (DataChange, error) {
var change DataChange
tx, err := i.store.Pool().Begin(ctx)
if err != nil {
return change, err
}
defer tx.Rollback(ctx) //nolint:errcheck
existing := map[dataKey]DataRecord{}
rows, err := tx.Query(ctx, dataSelect+` where machine = $1`, machine)
if err != nil {
return change, err
}
for rows.Next() {
r, err := scanData(rows)
if err != nil {
rows.Close()
return change, err
}
existing[dataKey{r.Module, r.Item}] = r
}
rows.Close()
if err := rows.Err(); err != nil {
return change, err
}
seen := map[dataKey]bool{}
for _, d := range declared {
k := dataKey{d.Module, d.Item}
seen[k] = true
m := measured[d.Module][d.Item]
var red struct {
Kind, Where, Said *string
Healthy *bool
}
if r := m.Redundancy; r != nil {
red.Kind, red.Where, red.Said, red.Healthy = &r.Kind, &r.Where, &r.Said, r.Healthy
}
was, had := existing[k]
if had && was.Retired() {
change.Reenabled = append(change.Reenabled, was)
}
// Deleted and declared again is new data: it starts over.
fresh := !had || was.DeletedAt != nil
if _, err := tx.Exec(ctx, `
insert into data_item (machine, module, item, class, path, first_seen, declared_at,
size_bytes, last_write, measured_at, last_backup, owned, protection,
measure_error, redundancy_kind, redundancy_where, redundancy_healthy, redundancy_said,
precision)
values ($1, $2, $3, $4, $5, $6, $6, $7, $8, $9, $10, $12, $13, $14, $15, $16, $17, $18, $19)
on conflict (machine, module, item) do update set
class = excluded.class,
owned = excluded.owned,
protection = excluded.protection,
precision = case when excluded.measured_at is not null then excluded.precision else data_item.precision end,
size_bytes = case when excluded.measured_at is not null then excluded.size_bytes else data_item.size_bytes end,
measure_error = case when excluded.measured_at is not null then excluded.measure_error else data_item.measure_error end,
redundancy_kind = case when excluded.measured_at is not null then excluded.redundancy_kind else data_item.redundancy_kind end,
redundancy_where = case when excluded.measured_at is not null then excluded.redundancy_where else data_item.redundancy_where end,
redundancy_healthy = case when excluded.measured_at is not null then excluded.redundancy_healthy else data_item.redundancy_healthy end,
redundancy_said = case when excluded.measured_at is not null then excluded.redundancy_said else data_item.redundancy_said end,
path = case when excluded.path <> '' then excluded.path else data_item.path end,
first_seen = case when $11::boolean then excluded.first_seen else data_item.first_seen end,
declared_at = excluded.declared_at,
last_write = coalesce(excluded.last_write, data_item.last_write),
measured_at = coalesce(excluded.measured_at, data_item.measured_at),
last_backup = coalesce(excluded.last_backup, data_item.last_backup),
retired_at = null, retired_why = null,
deleted_at = null, deleted_by = null, deleted_why = null`,
machine, d.Module, d.Item, d.Class, m.Path, now, m.Size, m.LastWrite, m.MeasuredAt, m.LastBackup,
fresh, d.Owned, d.Protection, nullable(m.Error), red.Kind, red.Where, red.Healthy, red.Said,
nullable(m.Precision)); err != nil {
return change, err
}
if m.Size != nil && m.MeasuredAt != nil && Comparable(m.Precision) {
if _, err := tx.Exec(ctx, `
insert into data_reading (machine, module, item, at, size_bytes, last_write)
select $1, $2, $3, $4, $5, $6
where not exists (select 1 from data_reading
where machine = $1 and module = $2 and item = $3 and at > $4::timestamptz - $7::interval)`,
machine, d.Module, d.Item, *m.MeasuredAt, *m.Size, m.LastWrite,
fmt.Sprintf("%d seconds", int(readingEvery.Seconds()))); err != nil {
return change, err
}
}
}
for k, r := range existing {
if seen[k] || r.DeletedAt != nil || r.Retired() {
continue
}
if !catalogue.Retires(r.Class) || !r.Owned {
if _, err := tx.Exec(ctx, `delete from data_item where machine = $1 and module = $2 and item = $3`,
machine, k.module, k.item); err != nil {
return change, err
}
continue
}
if _, err := tx.Exec(ctx, `update data_item set retired_at = $4, retired_why = $5
where machine = $1 and module = $2 and item = $3`, machine, k.module, k.item, now, why); err != nil {
return change, err
}
at := now
r.RetiredAt, r.RetiredWhy = &at, why
change.Retired = append(change.Retired, r)
}
if _, err := tx.Exec(ctx, `delete from data_reading where at < $1`, now.Add(-readingsKept)); err != nil {
return change, err
}
return change, tx.Commit(ctx)
}
const dataSelect = `select machine, module, item, class, path, first_seen, declared_at, size_bytes, last_write,
measured_at, last_backup, retired_at, coalesce(retired_why, ''), deleted_at, coalesce(deleted_by, ''),
coalesce(deleted_why, ''), owned, protection, coalesce(measure_error, ''), redundancy_kind, redundancy_where,
redundancy_healthy, redundancy_said, coalesce(precision, '') from data_item`
func scanData(rows pgx.Row) (DataRecord, error) {
var r DataRecord
var kind, where, said *string
var healthy *bool
err := rows.Scan(&r.Machine, &r.Module, &r.Item, &r.Class, &r.Path, &r.FirstSeen, &r.DeclaredAt, &r.Size,
&r.LastWrite, &r.MeasuredAt, &r.LastBackup, &r.RetiredAt, &r.RetiredWhy, &r.DeletedAt, &r.DeletedBy,
&r.DeletedWhy, &r.Owned, &r.Protection, &r.MeasureError, &kind, &where, &healthy, &said, &r.Precision)
if err == nil && kind != nil {
r.Redundancy = &Redundancy{Kind: *kind, Healthy: healthy}
if where != nil {
r.Redundancy.Where = *where
}
if said != nil {
r.Redundancy.Said = *said
}
}
return r, err
}
func nullable(s string) *string {
if s == "" {
return nil
}
return &s
}
// Data is every item the mesh keeps, by machine, module and item.
func (i *Inventory) Data(ctx context.Context) ([]DataRecord, error) {
rows, err := i.store.Pool().Query(ctx, dataSelect+` order by machine, module, item`)
if err != nil {
return nil, err
}
defer rows.Close()
var out []DataRecord
for rows.Next() {
r, err := scanData(rows)
if err != nil {
return nil, err
}
out = append(out, r)
}
return out, rows.Err()
}
// DataOf is one item.
func (i *Inventory) DataOf(ctx context.Context, machine, module, item string) (DataRecord, error) {
r, err := scanData(i.store.Pool().QueryRow(ctx, dataSelect+` where machine = $1 and module = $2 and item = $3`,
machine, module, item))
if errors.Is(err, pgx.ErrNoRows) {
return r, fmt.Errorf("the mesh knows no data %s of %s on %s", item, module, machine)
}
return r, err
}
// DataPeaks is each item's largest reading since a moment, keyed by DataRecord.Key.
func (i *Inventory) DataPeaks(ctx context.Context, since time.Time) (map[string]int64, error) {
rows, err := i.store.Pool().Query(ctx, `select machine, module, item, max(size_bytes) from data_reading
where at >= $1 group by machine, module, item`, since)
if err != nil {
return nil, err
}
defer rows.Close()
out := map[string]int64{}
for rows.Next() {
var machine, module, item string
var peak int64
if err := rows.Scan(&machine, &module, &item, &peak); err != nil {
return nil, err
}
out[machine+"/"+module+"/"+item] = peak
}
return out, rows.Err()
}
// MarkDataDeleted records that a person deleted a retired item. Refused for an item not retired.
func (i *Inventory) MarkDataDeleted(ctx context.Context, machine, module, item, by, why string, at time.Time) error {
tag, err := i.store.Pool().Exec(ctx, `update data_item set deleted_at = $4, deleted_by = $5, deleted_why = $6
where machine = $1 and module = $2 and item = $3 and retired_at is not null and deleted_at is null`,
machine, module, item, at, by, why)
if err != nil {
return err
}
if tag.RowsAffected() == 0 {
return fmt.Errorf("%s of %s on %s is not retired: only retired data is deleted", item, module, machine)
}
return nil
}