package link import ( "context" "encoding/json" "errors" "fmt" "time" "github.com/nats-io/nats.go" "github.com/novox/mesh-controller/internal/broker" ) // A hand-over asked of a machine's node-engine (novox/hq issue 356, issue 339). // // The operator hands a directory the node-engine uses as found to the mesh at the controller's terminal: // `nox node hand-over ` on the control-node (ADR 0272). The controller asks that machine's // engine on its own subject, a request on core NATS the engine answers once; the engine judges every // value and records the hand-over, or refuses and records nothing. The engine holds the same two shapes // in its own link code (mesh-host internal/link HandOverAsk, HandOverAnswer); a test on each side holds // the field names. // HandOverAsk is what the controller asks: the directory's absolute path as the engine states it, and who // asked, in the controller's words. type HandOverAsk struct { Path string `json:"path"` By string `json:"by"` } // HandOverAnswer is the engine's answer: what it recorded, or why it refused. type HandOverAnswer struct { Said string `json:"said,omitempty"` Refused string `json:"refused,omitempty"` } // HandOverWithin is how long the controller waits for the engine's answer: a file write, on a machine that is // up; a machine that is down is said as not answering. const HandOverWithin = 30 * time.Second // AskHandOver asks one machine's node-engine to hand a directory used as found to the mesh, and reads its // answer. An error is the ask not reaching an engine, or an answer that is not one; a refusal is the engine's // and comes back in the answer. func AskHandOver(ctx context.Context, conn *nats.Conn, node, path, by string, timeout time.Duration) (HandOverAnswer, error) { if conn == nil { return HandOverAnswer{}, errors.New("this controller is not on the bus") } body, err := json.Marshal(HandOverAsk{Path: path, By: by}) if err != nil { return HandOverAnswer{}, err } asking, cancel := context.WithTimeout(ctx, timeout) defer cancel() subject := broker.AskHandOverSubject(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 HandOverAnswer{}, fmt.Errorf("the bus refused the controller asking %s for a hand-over: %v", node, why) } switch { case errors.Is(err, nats.ErrNoResponders): return HandOverAnswer{}, fmt.Errorf("nothing on %s answers a hand-over: its node-engine is not running, is not "+ "on the bus, or is older than this ask (novox/hq issue 356); nothing was handed over", node) case errors.Is(err, context.DeadlineExceeded), errors.Is(err, nats.ErrTimeout): return HandOverAnswer{}, fmt.Errorf("%s did not answer the hand-over within %s; whether it was recorded is not "+ "known — the module's condition says whether the directory is still used as found", node, timeout) case err != nil: return HandOverAnswer{}, err } var answer HandOverAnswer if err := json.Unmarshal(reply.Data, &answer); err != nil { return HandOverAnswer{}, fmt.Errorf("%s answered the hand-over with something unreadable: %w", node, err) } if answer.Said == "" && answer.Refused == "" { return HandOverAnswer{}, fmt.Errorf("%s answered the hand-over with neither a record nor a refusal", node) } return answer, nil }