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"` }