A reading now keeps the path it was measured at, and the peak is read only from readings at the item's path now. Before, an agent's home that moved to its own account was compared with the operator's home it left, and data-shrank was raised for data that was never lost. Existing readings take their item's path now, so a shrink that is real stays raised.
326 lines
13 KiB
Go
326 lines
13 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. A reading keeps the item's path
|
|
// as it stands after the measurement — the holder's, or the last one known when it named none — and an
|
|
// item at a path with no recent reading is read at once (novox/hq issue 368).
|
|
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, path)
|
|
select $1, $2, $3, $4, $5, $6, d.path
|
|
from data_item d
|
|
where d.machine = $1 and d.module = $2 and d.item = $3
|
|
and not exists (select 1 from data_reading
|
|
where machine = $1 and module = $2 and item = $3 and path = d.path
|
|
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 — read only from
|
|
// readings at the item's path now (novox/hq issue 368): a directory is never compared with another one
|
|
// that once held the same item, so a moved item starts its size history again at its new path.
|
|
func (i *Inventory) DataPeaks(ctx context.Context, since time.Time) (map[string]int64, error) {
|
|
rows, err := i.store.Pool().Query(ctx, `select r.machine, r.module, r.item, max(r.size_bytes)
|
|
from data_reading r
|
|
join data_item d on d.machine = r.machine and d.module = r.module and d.item = r.item and d.path = r.path
|
|
where r.at >= $1 group by r.machine, r.module, r.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
|
|
}
|