Files
mesh-controller/internal/link/retirement.go
T
jochen 52af210e47 Derive data protection from a module's declared data (hq ADR 0233)
A module's data section says what it keeps and how precious it is; the backup holder's lines,
binding stickiness, retirement on unassign and D13's conditions follow from it, so issue 273's
empty replacement is said and an unassigned module's data is remembered, not forgotten.
2026-10-06 16:47:49 +02:00

252 lines
9.4 KiB
Go

package link
import (
"context"
"encoding/json"
"errors"
"fmt"
"strings"
"time"
"github.com/nats-io/nats.go"
"github.com/novox/mesh-controller/internal/broker"
)
// What becomes of a consumer the mesh stopped asking for (novox/hq ADR 0230).
//
// **A consumer is retired, not withdrawn, and deleted only by a person.** A provider that has seen the
// same consumers go unasked for in five passes disables their access and keeps their data, marked in
// its backend with when and why; asked for again, it enables them as they were. When more would go
// than its bound allows — more than three, or more than half of those it holds — it retires nothing and
// waits for `retire approve` or `retire reject`. A retired consumer is deleted only by `cleanup delete`,
// which the provider executes. Each of these is the provider's `provisioner.retirement` event, its
// `change` saying which; the controller keeps the ones that need a person as conditions.
// The changes a provider says.
const (
RetireWaiting = "waiting"
RetireSettled = "settled"
RetireApproved = "approved"
RetireRejected = "rejected"
RetireRetired = "retired"
RetireReenabled = "reenabled"
RetireDeleted = "deleted"
RetireAdopted = "adopted"
)
// ViaController is what the controller's verbs say they came through, so an act a provider was asked
// for some other way is told apart and recorded by hand.
const ViaController = "mesh-controller"
// RetiredConsumer is one consumer in a retirement word, or in a provider's answer.
type RetiredConsumer struct {
Consumer string `json:"consumer"`
Node string `json:"node,omitempty"`
RetiredAt string `json:"retired_at,omitempty"`
Why string `json:"why,omitempty"`
// SizeBytes is what it keeps on the backend; -1 or absent where the backend cannot say.
SizeBytes *int64 `json:"size_bytes,omitempty"`
// Kind is "consumer", or a backend's own word for something set aside ("set-aside-database").
Kind string `json:"kind,omitempty"`
// Access is "disabled" when the provider disabled it, "kept" for a mark-only provider whose consumer
// keeps its access until a person deletes it (ADR 0230); absent is disabled.
Access string `json:"access,omitempty"`
}
// Retirement is one provider's word about consumers the mesh stopped asking for.
type Retirement struct {
// Module is the emitter, read from the subject the bus let it publish on — never from the body.
Module string `json:"-"`
Provider string `json:"provider"`
ProviderNode string `json:"provider-node"`
Change string `json:"change"`
Consumers []RetiredConsumer `json:"consumers"`
Held int `json:"held"`
Bound string `json:"bound,omitempty"`
Why string `json:"why,omitempty"`
By string `json:"by,omitempty"`
Via string `json:"via,omitempty"`
At time.Time `json:"at"`
Since time.Time `json:"since,omitempty"`
}
// Names are the consumers it names, in its order.
func (r Retirement) Names() []string {
out := make([]string, 0, len(r.Consumers))
for _, c := range r.Consumers {
out = append(out, c.Consumer)
}
return out
}
// Retirements keeps what providers say about consumers they retire.
type Retirements interface {
// Retired records one word. An error the store is away for is held and asked again, like a report.
Retired(ctx context.Context, r Retirement) error
}
// KeepsRetirements says where retirement words are kept, and asks for them to be delivered.
func (s *Server) KeepsRetirements(r Retirements) error {
if err := s.inbound.Also(KindProvisioner); err != nil {
return err
}
s.retirements = r
return nil
}
// IsRetirement says a subject is a provider's retirement word.
func IsRetirement(subject string) bool {
_, ok := ProvisionerEmitter(subject)
return ok && strings.HasSuffix(subject, ".event."+broker.ProvisionerRetirement)
}
// ReadRetirement is one retirement word as the controller understands it.
func ReadRetirement(subject string, body []byte) (Retirement, error) {
module, ok := ProvisionerEmitter(subject)
if !ok || !IsRetirement(subject) {
return Retirement{}, fmt.Errorf("%s is not a provider's retirement word", subject)
}
var r Retirement
if err := json.Unmarshal(body, &r); err != nil {
return Retirement{}, fmt.Errorf("%s's retirement word could not be read: %w", module, err)
}
switch r.Change {
case RetireWaiting, RetireSettled, RetireApproved, RetireRejected, RetireRetired, RetireReenabled,
RetireDeleted, RetireAdopted:
default:
return Retirement{}, fmt.Errorf("%s said a retirement change the mesh has no name for: %q", module, r.Change)
}
if r.ProviderNode == "" {
return Retirement{}, fmt.Errorf("%s's retirement word named no machine", module)
}
r.Module = module
return r, nil
}
// retirement acts on one retirement word. Like a recovery, a word that clears a condition is said
// once, so a store that is away holds the message rather than dropping it.
func (s *Server) retirement(ctx context.Context, m Control) {
if s.retirements == nil {
_ = m.Took()
return
}
r, err := ReadRetirement(m.Subject(), m.Body())
if err != nil {
s.log.Printf("%v; ignored", err)
_ = m.Took()
return
}
err = s.retirements.Retired(ctx, r)
what := fmt.Sprintf("%s's retirement word (%s)", r.Module, r.Change)
switch s.decide(ctx, m, what, "", "", err) {
case Hold:
return
case Stale, GiveUp:
_ = m.Took()
return
}
if err != nil {
s.log.Printf("%s could not be kept: %v", what, err)
} else {
s.log.Printf("%s on %s %s %s", r.Module, r.ProviderNode, r.Change, strings.Join(r.Names(), ", "))
}
_ = m.Took()
}
// The tools every provider serves about its retired consumers (ADR 0230), asked of one machine.
const (
ToolRetirement = "provisioner_retirement"
ToolRetireApprove = "provisioner_retire_approve"
ToolRetireReject = "provisioner_retire_reject"
ToolRetiredDelete = "provisioner_delete"
)
// ErrNothingServes is a tool nothing on that machine serves: the module is not running there, or
// is older than the tool.
var ErrNothingServes = errors.New("nothing serves it")
// ModuleToolOn is where one machine's instance of a module answers a tool: the module's tool subject
// with the machine as its last token — the subject the node tools bind beside the shared one.
func ModuleToolOn(module, tool, node string) string { return ToolSubject(module, tool) + "." + node }
// AskModuleToolOn asks one machine's instance of a module one of its tools and reads its answer.
func AskModuleToolOn(ctx context.Context, conn *nats.Conn, module, tool, node string, args any,
timeout time.Duration) (Answer, error) {
body, err := json.Marshal(args)
if err != nil {
return Answer{}, err
}
asking, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
subject := ModuleToolOn(module, tool, node)
refused, stop := refusalsOf(conn, subject)
defer stop()
type replied struct {
msg *nats.Msg
err error
}
done := make(chan replied, 1)
go func() {
msg, err := conn.RequestWithContext(asking, subject, body)
done <- replied{msg, err}
}()
var reply *nats.Msg
select {
case r := <-done:
reply, err = r.msg, r.err
case why := <-refused:
cancel()
return Answer{}, fmt.Errorf("the bus refused the controller asking %s.%s on %s: %v", module, tool, node, why)
}
switch {
case errors.Is(err, nats.ErrNoResponders):
return Answer{}, fmt.Errorf("%w: %s on %s does not answer %s — it is not running there, or is "+
"older than the tool", ErrNothingServes, module, node, tool)
case errors.Is(err, context.DeadlineExceeded), errors.Is(err, nats.ErrTimeout):
return Answer{}, fmt.Errorf("%s on %s did not answer %s within %s", module, node, tool, timeout)
case err != nil:
return Answer{}, err
}
var answer Answer
if err := json.Unmarshal(reply.Data, &answer); err != nil {
return Answer{}, fmt.Errorf("%s on %s answered %s with something unreadable: %w", module, node, tool, err)
}
return answer, nil
}
// RetirementWaiting is the set a provider waits with for a person.
type RetirementWaiting struct {
Consumers []RetiredConsumer `json:"consumers"`
Since string `json:"since"`
Held int `json:"held"`
}
// RetirementRejected is the set a person's rejection keeps active.
type RetirementRejected struct {
Consumers []RetiredConsumer `json:"consumers"`
By string `json:"by"`
Why string `json:"why"`
At string `json:"at"`
}
// RetirementState is a provider's answer to provisioner_retirement.
type RetirementState struct {
Resource string `json:"resource"`
Node string `json:"node"`
Held []string `json:"held"`
// HeldSizes is what each held consumer keeps on the backend, in bytes, where the backend can say
// (novox/hq ADR 0233): what an empty replacement of a consumer's data is told by.
HeldSizes map[string]int64 `json:"held_sizes,omitempty"`
StablePasses int `json:"stable_passes"`
// StableForSeconds is how long the same unasked set must hold as well (ADR 0230: ten minutes).
StableForSeconds int `json:"stable_for_seconds,omitempty"`
// RetiresBy is "disable", or "mark-only" for a provider that cannot disable a consumer.
RetiresBy string `json:"retires_by,omitempty"`
Bound string `json:"bound"`
Waiting *RetirementWaiting `json:"waiting"`
Rejected *RetirementRejected `json:"rejected"`
Retired []RetiredConsumer `json:"retired"`
}