kubernetes
| Input | Unit | Covered | Total | Percent |
| Go | statements | 968 | 983 | 98.5% |
Go
968 of 983 statements, 98.5%.
apiclient/client.go 98.2%
1// Package apiclient sends an operator's reads and writes to the2// Kubernetes API server.3//4// The Kubernetes API is HTTPS that serves JSON, and each read or write5// is one request, so this client uses only net/http and encoding/json.6// It imports nothing from k8s.io. A program that must not link7// client-go, such as a pod build of an operator or a command-line tool,8// can use it. The watches run on client-go's reflector, in the informer9// package, because upstream maintains and tests that loop.10//11// Every pod starts with what it needs to reach the API server.12// Kubernetes injects two environment variables that name the server's13// in-cluster address, and the kubelet mounts a CA certificate and a14// ServiceAccount token at a known path. Those five values are the whole15// of what client-go's rest.InClusterConfig() reads.16package apiclient1718import (19 "bytes"20 "context"21 "crypto/tls"22 "crypto/x509"23 "encoding/json"24 "errors"25 "fmt"26 "io"27 "net"28 "net/http"29 "os"30 "strings"31 "time"32 "unicode"33)3435// ServiceAccountDir is the directory where the kubelet mounts each36// container's API credentials.37const ServiceAccountDir = "/var/run/secrets/kubernetes.io/serviceaccount"3839// ErrNotFound separates "this object does not exist" from a real40// failure. An absent object is a normal state, and the caller handles41// it by creating the object. ErrConflict separates "something else42// wrote this object first" in the same way. That is a normal state43// under optimistic concurrency, and the caller handles it by reading44// the object again.45var (46 ErrNotFound = errors.New("not found")47 ErrConflict = errors.New("conflict: something else wrote this object first")48)4950// ErrThrottled is a 429 that lasted longer than the client waits for it51// (maxThrottleWait), or past the end of the client's context. A caller52// that must not give up, such as an operator's first read at start,53// asks again until its own context ends.54var ErrThrottled = errors.New("throttled: the API server asked the client to wait")5556// Stale reports whether a write failed because the copy it was made57// from is not the API server's copy: another writer changed the58// object, or deleted it.59func Stale(err error) bool {60 return errors.Is(err, ErrConflict) || errors.Is(err, ErrNotFound)61}6263// Client sends requests to one API server.64type Client struct {65 base string66 http *http.Client67 credentials string6869 // ctx is the context every request carries (WithContext). Nil means70 // no context ends a request.71 ctx context.Context7273 // waits is the context that ends the wait after a 42974 // (WithWaitContext). Nil means ctx ends it.75 waits context.Context7677 // writeGuard, when it is set, runs before each send of a request78 // that is not a GET (WithWriteGuard).79 writeGuard func() error8081 // observer, when it is set, hears the final answer to each request82 // (WithObserver).83 observer func(Outcome)84}8586// Outcome is the final answer to one request, for an observer.87type Outcome struct {88 Method string89 Path string9091 // Status is the HTTP status of the last answer, or zero when no92 // answer came: the connection failed, the context ended, or a93 // write guard refused the request.94 Status int9596 // Err is the error the caller received, or nil.97 Err error98}99100// New builds a client from its three parts. InCluster reads them from101// the pod's environment, and a test takes them from its fake API102// server. An empty credentials directory sends no token.103func New(base string, httpClient *http.Client, credentials string) *Client {104 return &Client{base: base, http: httpClient, credentials: credentials}105}106107// InClusterOptions are the optional parts of an in-cluster client.108type InClusterOptions struct {109 // ServiceAccountDir is the directory that holds the CA and the110 // token. Empty means the directory the kubelet mounts.111 ServiceAccountDir string112113 // Server is the API server's address, such as114 // https://127.0.0.1:6443. Empty means the address that the115 // environment names: the Service's virtual IP, where iptables pins116 // each new connection to one API server that the client cannot117 // choose. A pod on the host's network, on a machine that runs an118 // API server or k3s's local load balancer, has a better address:119 // its own loopback, where a remote server that died cannot hold a120 // connection. The CA and the token are the same at both addresses.121 Server string122123 // Timeout limits one request, from the dial to the last byte of124 // the answer. Zero means 30 seconds. A program that must abandon a125 // write before a deadline of its own, such as a leader whose Lease126 // another copy can take, sets a shorter limit.127 Timeout time.Duration128}129130// defaultTimeout is the limit of InClusterOptions.Timeout when it is131// zero.132const defaultTimeout = 30 * time.Second133134// InCluster builds a client from the pod's environment and its135// ServiceAccount.136func InCluster(options InClusterOptions) (*Client, error) {137 server := options.Server138 if server == "" {139 host, port := os.Getenv("KUBERNETES_SERVICE_HOST"), os.Getenv("KUBERNETES_SERVICE_PORT")140 if host == "" || port == "" {141 return nil, fmt.Errorf("not running in a cluster: KUBERNETES_SERVICE_HOST unset")142 }143 server = "https://" + host + ":" + port144 }145 timeout := options.Timeout146 if timeout == 0 {147 timeout = defaultTimeout148 }149 dir := options.ServiceAccountDir150 if dir == "" {151 dir = ServiceAccountDir152 }153154 // The mounted CA is the cluster's own. The client trusts only that155 // CA, not the system trust store, so it accepts the cluster's API156 // server and refuses any other server that answers on the address.157 caPEM, err := os.ReadFile(dir + "/ca.crt")158 if err != nil {159 return nil, fmt.Errorf("reading service account CA: %w", err)160 }161 roots := x509.NewCertPool()162 if !roots.AppendCertsFromPEM(caPEM) {163 return nil, fmt.Errorf("service account CA contains no certificates")164 }165166 return New(server, &http.Client{167 Transport: &http.Transport{168 TLSClientConfig: &tls.Config{RootCAs: roots},169 // Each timeout limits the same failure: a server that stops170 // answering without sending anything. A machine that fails171 // sends no FIN and no RST, so a connection to it goes silent172 // and every wait on it would otherwise have no limit.173 DialContext: (&net.Dialer{174 Timeout: 5 * time.Second,175 KeepAlive: 10 * time.Second,176 }).DialContext,177 ResponseHeaderTimeout: 10 * time.Second,178 IdleConnTimeout: 30 * time.Second,179 },180 // No request of this client streams, because the informer181 // package runs every watch. So the whole request, body182 // included, has a limit too, and a body that stops part way183 // cannot hold a pass.184 Timeout: timeout,185 }, dir), nil186}187188// WithContext answers a client whose requests carry ctx. A request189// ends when ctx ends, and so does the wait after a 429, so a caller190// that stops its work does not wait on the API server. The client it191// answers shares the connections of c.192func (c *Client) WithContext(ctx context.Context) *Client {193 bound := *c194 bound.ctx, bound.waits = ctx, nil195 return &bound196}197198// WithWaitContext answers a client whose wait after a 429 ends when ctx199// ends, and which then answers the 429. A request already sent runs to200// its end. A request that the client cut off may still land on the API201// server, so a writer that must know whether its write landed uses this202// in place of WithContext: an operator that releases its Lease after203// its last write waits for each write in flight, and not for the API204// server's advice to wait.205func (c *Client) WithWaitContext(ctx context.Context) *Client {206 bound := *c207 bound.waits = ctx208 return &bound209}210211// WithWriteGuard answers a client that asks guard before it sends each212// request that is not a GET, and each time it sends one again after a213// 429. An error from guard refuses the request, and the client sends214// nothing. A leader uses it to stop writing once it cannot show that it215// still holds its Lease. The guard sits in the client, so it covers216// every write, including those that memo and informer send. The client217// it answers shares the connections of c.218func (c *Client) WithWriteGuard(guard func() error) *Client {219 guarded := *c220 guarded.writeGuard = guard221 return &guarded222}223224// WithObserver answers a client that calls observe with the final225// answer to each request, after any wait for a 429. A caller that logs226// a failure and carries on still leaves a record there, so a program227// can learn from one place that something it sent did not land, and228// send it again sooner than its next scheduled try. The client it229// answers shares the connections of c.230func (c *Client) WithObserver(observe func(Outcome)) *Client {231 observed := *c232 observed.observer = observe233 return &observed234}235236// waitContext answers the context that ends the wait after a 429.237func (c *Client) waitContext() context.Context {238 if c.waits != nil {239 return c.waits240 }241 return c.context()242}243244// context answers the context of each request.245func (c *Client) context() context.Context {246 if c.ctx == nil {247 return context.Background()248 }249 return c.ctx250}251252// RequestJSON sends one request with a JSON body, and decodes the JSON253// answer into out. A nil out discards the answer.254func (c *Client) RequestJSON(method, path string, body []byte, out any) error {255 return c.Request(method, path, "application/json", body, out)256}257258// Request is RequestJSON with the body's content type stated. A PATCH259// needs it: the API server reads the patch's dialect from the header260// alone, and the same bytes mean different things as a merge patch, a261// JSON patch, and an apply.262//263// A 404 answers ErrNotFound, a 409 answers ErrConflict, and any other264// status that is not 2xx answers an error that holds the server's own265// message. A 429 is sent again after the wait the API server asks for,266// as maxThrottleWait describes.267func (c *Client) Request(method, path, contentType string, body []byte, out any) error {268 status, err := c.request(method, path, contentType, body, out)269 if c.observer != nil {270 c.observer(Outcome{Method: method, Path: path, Status: status, Err: err})271 }272 return err273}274275// request sends a request until it gets an answer that is not a 429 it276// can wait out, and answers the last status with the error.277func (c *Client) request(method, path, contentType string, body []byte, out any) (int, error) {278 var waited time.Duration279 for {280 status, err := c.send(method, path, contentType, body, out)281 var throttled *throttledError282 if !errors.As(err, &throttled) || waited+throttled.wait > maxThrottleWait*time.Second {283 return status, err284 }285 timer := time.NewTimer(throttled.wait)286 select {287 case <-c.waitContext().Done():288 timer.Stop()289 return status, err290 case <-timer.C:291 }292 waited += throttled.wait293 }294}295296// maxThrottleWait is the longest total wait, in units of the wait a297// 429 asks for, before a request answers the 429 to its caller. While298// the API server starts the storage of a CRD it just received, it299// answers 429 with the reason "storage is (re)initializing" for a300// second or two. A rollout that changes a CRD meets that answer, and a301// pass that sent the request again after the wait gets the object. A302// server that answers 429 for longer is overloaded, and the caller's303// own retry, which waits longer, handles it.304//305// The wait ends when the context of a client from WithContext or306// WithWaitContext ends.307// Otherwise the limit bounds it: a request holds its caller for at most308// ten seconds of waits, plus the request timeout of each send. A caller309// that holds a lock across the request, such as memo.Versions.Send,310// holds it that long too, and a shutdown that waits on the caller waits311// that long.312const maxThrottleWait = 10313314// send sends one request once, and answers the HTTP status of the315// answer, or zero when no answer came.316func (c *Client) send(method, path, contentType string, body []byte, out any) (int, error) {317 if c.writeGuard != nil && method != http.MethodGet {318 if err := c.writeGuard(); err != nil {319 return 0, fmt.Errorf("%s %s not sent: %w", method, path, err)320 }321 }322 var reader io.Reader323 if body != nil {324 reader = bytes.NewReader(body)325 }326 req, err := http.NewRequestWithContext(c.context(), method, c.base+path, reader)327 if err != nil {328 return 0, err329 }330 if err := c.authorize(req); err != nil {331 return 0, err332 }333 req.Header.Set("Accept", "application/json")334 if body != nil {335 req.Header.Set("Content-Type", contentType)336 }337338 resp, err := c.http.Do(req)339 if err != nil {340 return 0, err341 }342 defer drain(resp.Body)343344 status := resp.StatusCode345 switch {346 case status == http.StatusNotFound:347 return status, ErrNotFound348 case status == http.StatusConflict:349 return status, ErrConflict350 case status < 200 || status > 299:351 message := responseText(resp.Body)352 err := fmt.Errorf("%s %s: %s: %s", method, path, resp.Status, message)353 if status == http.StatusTooManyRequests {354 seconds := retryAfter(resp.Header.Get("Retry-After"), message)355 return status, &throttledError{err: err, wait: time.Duration(max(seconds, 1)) * time.Second, seconds: seconds}356 }357 return status, err358 case out == nil:359 return status, nil360 }361 return status, json.NewDecoder(resp.Body).Decode(out)362}363364// authorize puts the ServiceAccount token on a request.365//366// The client reads its token from disk on every request. The tokens are367// short-lived, and the kubelet refreshes the mounted file as each one368// nears its expiry, so a client that holds a token in memory eventually369// gets 401 answers.370func (c *Client) authorize(req *http.Request) error {371 if c.credentials == "" {372 return nil373 }374 token, err := os.ReadFile(c.credentials + "/token")375 if err != nil {376 return fmt.Errorf("reading service account token: %w", err)377 }378 req.Header.Set("Authorization", "Bearer "+string(token))379 return nil380}381382// throttledError is a 429 from the API server, which asks the client to383// wait and send the request again. seconds is the wait the answer384// stated, and zero when it stated none; wait is what the client waits385// before it sends again, one second when the answer stated none.386type throttledError struct {387 err error388 wait time.Duration389 seconds int390}391392func (e *throttledError) Error() string { return e.err.Error() }393394// Is makes a throttledError match ErrThrottled.395func (e *throttledError) Is(target error) bool { return target == ErrThrottled }396397// retryAfter reads how many seconds a 429 asks the client to wait: the398// Retry-After header, or the retryAfterSeconds of the Status body, or399// zero when the answer states neither.400func retryAfter(header, body string) int {401 var seconds int402 if _, err := fmt.Sscan(header, &seconds); err == nil && seconds > 0 {403 return seconds404 }405 var status struct {406 Details struct {407 RetryAfterSeconds int `json:"retryAfterSeconds"`408 } `json:"details"`409 }410 if json.Unmarshal([]byte(body), &status) == nil && status.Details.RetryAfterSeconds > 0 {411 return status.Details.RetryAfterSeconds412 }413 return 0414}415416// RetryAfterSeconds answers the seconds that the 429 an error holds417// asked the caller to wait, and zero for a 429 that stated no wait and418// for an error that holds no 429. A caller that asks again after419// ErrThrottled waits at least that long, because the API server asked420// for it. A 429 with no stated wait is often a refusal that time alone421// does not end, such as an eviction that a PodDisruptionBudget refuses,422// and the caller decides when to ask again.423func RetryAfterSeconds(err error) int {424 var throttled *throttledError425 if !errors.As(err, &throttled) {426 return 0427 }428 return throttled.seconds429}430431// responseText is the start of an answer's body, for an error that432// holds the server's own text. The API server ends its Status body with433// a newline, and the text leaves it out, so the error and the log line434// that prints it stay on one line.435func responseText(body io.Reader) string {436 message, _ := io.ReadAll(io.LimitReader(body, 2048))437 return strings.TrimRightFunc(string(message), unicode.IsSpace)438}439440// maxDrain limits the read below. The largest answer an operator asks441// for is one list, which the caller decodes into memory anyway, so442// reading the tail costs nothing new. Past this size, the connection is443// the cheaper thing to lose.444const maxDrain = 4 << 20445446// drain reads what the caller left in the response body, then closes447// it. Go returns a connection to its pool only when the body reaches448// EOF, so a body closed early costs a new TCP connection and TLS449// handshake on the next request. The early close also reaches the API450// server as a hang-up on a request it already answered. A 404, a 409,451// and a status write each answer with a body that no caller reads.452func drain(body io.ReadCloser) {453 _, _ = io.Copy(io.Discard, io.LimitReader(body, maxDrain))454 _ = body.Close()455}
apiclient/objects.go 91.7%
1package apiclient23import (4 "encoding/json"5 "net/http"6)78// Get reads one object.9func Get[T any](c *Client, path string) (*T, error) {10 out := new(T)11 if err := c.RequestJSON(http.MethodGet, path, nil, out); err != nil {12 return nil, err13 }14 return out, nil15}1617// ReplaceStatus writes one object's status to its status subresource,18// so a spec that a person edited between the read and the write stays19// as the person wrote it. The write sends the whole object, because20// that is what the subresource takes, and the API server keeps only the21// status half of it. The caller states the object's apiVersion and22// kind: an object read out of a list does not always hold them, and the23// API server refuses a write that holds neither.24//25// The API server's copy replaces the caller's, because the write makes26// a new resourceVersion and the next write must state it. The answer is27// decoded into a new value and not over the caller's object: a decode28// over the old object keeps each field that the answer leaves out, such29// as a status field the CRD schema drops, and the caller then holds a30// copy the API server does not. After an error the caller's object is31// as it was, with the status the API server did not take.32func ReplaceStatus[T any](c *Client, path string, object *T) error {33 body, err := json.Marshal(object)34 if err != nil {35 return err36 }37 stored := new(T)38 if err := c.RequestJSON(http.MethodPut, path+"/status", body, stored); err != nil {39 return err40 }41 *object = *stored42 return nil43}
apiservertest/errors.go 100.0%
1package apiservertest23// client-go reports each failed watch through apimachinery's4// utilruntime.HandleError, and the default handlers log the error and5// then run a global rate limit. The rate limit records the time of the6// last error in a package variable, first at program start, and sleeps7// until one millisecond after that time. Each bubble's clock starts at8// midnight UTC 2000-01-01, before any time the variable records outside9// a bubble or in an earlier bubble. So the first failed watch in a10// bubble sleeps for years of fake time, the reflector stops, and the11// bubble deadlocks when the test ends. A program that imports this12// package is a test, and in a test the rate limit protects nothing, so13// the package keeps the log and removes the rate limit.1415import (16 "context"1718 utilruntime "k8s.io/apimachinery/pkg/util/runtime"19 "k8s.io/klog/v2"20)2122func init() {23 utilruntime.ErrorHandlers = []utilruntime.ErrorHandler{logUnhandled}24}2526// logUnhandled logs the error the way apimachinery's own first handler27// does.28func logUnhandled(ctx context.Context, err error, message string, keysAndValues ...any) {29 klog.LoggerWithName(klog.FromContext(ctx), "UnhandledError").Error(err, message, keysAndValues...)30}
apiservertest/server.go 96.9%
1// Package apiservertest serves a test's fake API server over in-memory2// connections, so a test that talks to it can run on the fake clock of3// testing/synctest.4//5// A synctest bubble advances its clock only while every goroutine in6// it is durably blocked. A goroutine that reads a socket is not durably7// blocked, because an event outside the bubble can wake it. So a test8// whose client reads from an httptest.Server never advances the clock,9// and client-go's reflector backoff, a lease, or a retry wait takes its10// full duration in real time. This package replaces the socket with11// net.Pipe connections, whose reads block on channels. The server, its12// connections, and the client's connections start in the test's own13// goroutine tree, so in a bubble they are part of it, and a wait of a14// minute takes no real time.15//16// A test serves its handler with Start, and gives client-go the17// configuration that Config answers:18//19// synctest.Test(t, func(t *testing.T) {20// server := apiservertest.Start(t, handler)21// client, err := dynamic.NewForConfig(server.Config())22// ...23// informer.Start(t.Context(), client, source, options)24// time.Sleep(time.Minute)25// synctest.Wait()26// })27//28// synctest.Test waits for every goroutine in the bubble to exit, and29// fails the test when they block forever. The server closes itself and30// each of its connections when the test ends, so a client's connections31// end too. A goroutine that the test started, such as an informer, a32// watch, or a retry loop, runs until the test stops it. Start each one33// with t.Context(), which ends before the cleanups run, or stop it in a34// cleanup. A loop that keeps its own context sends to the closed server35// again after each backoff, and the bubble never ends.36//37// client-go's reflector has one wait that does not end with its context:38// after a streaming list meets a refused connection or a 429, it waits39// out its backoff, which is less than a minute. A test that takes the40// server down while a reflector runs sleeps for a minute in a cleanup,41// so that wait ends before the bubble does.42//43// The same server works outside a bubble, so a test can move to it44// before it moves to synctest.45//46// A program that imports the package also loses apimachinery's global47// rate limit on unhandled errors, because in a bubble that rate limit48// stops the reflector (errors.go).49package apiservertest5051import (52 "net"53 "net/http"54 "os"55 "sync"56 "syscall"57 "testing"58)5960// Server serves one handler over in-memory connections. It is up after61// Start, and SetDown takes it down and brings it back.62type Server struct {63 handler http.Handler6465 mu sync.Mutex66 serving *http.Server67 listener *pipeListener68}6970// Start serves the handler until the test ends.71func Start(t testing.TB, handler http.Handler) *Server {72 s := &Server{handler: handler}73 s.SetDown(false)74 t.Cleanup(func() { s.SetDown(true) })75 return s76}7778// SetDown takes the server down, or brings it up again. A server that79// goes down closes each connection it holds, the way a stopped process80// does, so a watch ends. Each request then fails with ECONNREFUSED, the81// error the kernel gives a client of an API server that restarts, until82// the server comes up. client-go's reflector checks for that error, and83// retries such a watch without a new list.84func (s *Server) SetDown(down bool) {85 s.mu.Lock()86 defer s.mu.Unlock()87 if down == (s.serving == nil) {88 return89 }90 if down {91 _ = s.serving.Close()92 s.serving, s.listener = nil, nil93 return94 }95 listener, serving := newPipeListener(), &http.Server{Handler: s.handler}96 s.listener, s.serving = listener, serving97 go func() { _ = serving.Serve(listener) }()98}99100// dial opens a pipe to the server and answers the client's end, once101// the server has accepted the other end. A server that is down refuses102// the connection. The lock keeps the server up until it accepts, and103// its accept loop takes each pipe at once.104func (s *Server) dial() (net.Conn, error) {105 s.mu.Lock()106 defer s.mu.Unlock()107 if s.listener == nil {108 return nil, &net.OpError{Op: "dial", Net: "tcp", Err: os.NewSyscallError("connect", syscall.ECONNREFUSED)}109 }110 client, server := net.Pipe()111 s.listener.accepted <- server112 return client, nil113}114115// pipeListener hands the server one end of each pipe that a dial opens.116type pipeListener struct {117 accepted chan net.Conn118 closed chan struct{}119 once sync.Once120}121122func newPipeListener() *pipeListener {123 return &pipeListener{accepted: make(chan net.Conn), closed: make(chan struct{})}124}125126func (l *pipeListener) Accept() (net.Conn, error) {127 select {128 case conn := <-l.accepted:129 return conn, nil130 case <-l.closed:131 return nil, net.ErrClosed132 }133}134135func (l *pipeListener) Close() error {136 l.once.Do(func() { close(l.closed) })137 return nil138}139140func (l *pipeListener) Addr() net.Addr { return &net.UnixAddr{Name: "pipe", Net: "pipe"} }
apiservertest/transport.go 100.0%
1package apiservertest23// The client side sends each request on a connection of its own, and4// closes the connection when the caller closes the answer's body.5//6// net/http's Transport does not work in a bubble. When a caller closes7// a streaming body early, such as a watch that stops, the Transport8// reads what is left of the body for up to 50 milliseconds, so it can9// keep the connection. The watch's own reader holds the body's lock10// while it waits for the next event, so the drain waits on that lock. A11// goroutine that waits on a lock is not durably blocked, the clock never12// reaches the 50 milliseconds, and the test never ends. Here, a closed13// body closes its connection first, which ends the reader's wait.1415import (16 "bufio"17 "context"18 "io"19 "net"20 "net/http"2122 "k8s.io/client-go/rest"23)2425// Host is the address that Config and Client send each request to. The26// server answers every address, because each request reaches it through27// a pipe and no name is resolved.28const Host = "http://apiserver.test"2930// Config answers a client-go configuration whose requests reach the31// server. Each call answers a new configuration, which a test can32// change, for example to wrap its Transport.33func (s *Server) Config() *rest.Config {34 return &rest.Config{Host: Host, Transport: s}35}3637// Client answers an HTTP client whose requests reach the server, for a38// client that is not client-go.39func (s *Server) Client() *http.Client {40 return &http.Client{Transport: s}41}4243// RoundTrip sends one request to the server on a new connection, and44// answers the server's response. The connection closes when the45// request's context ends or the caller closes the body.46func (s *Server) RoundTrip(req *http.Request) (*http.Response, error) {47 // A request whose context ended is never sent.48 err := context.Cause(req.Context())49 var conn net.Conn50 if err == nil {51 conn, err = s.dial()52 }53 if err != nil {54 if req.Body != nil {55 _ = req.Body.Close()56 }57 return nil, err58 }59 stop := context.AfterFunc(req.Context(), func() { _ = conn.Close() })60 // The server can answer before it reads the whole request, so the61 // request is written while the response is read.62 go func() { _ = req.Write(conn) }()63 resp, err := http.ReadResponse(bufio.NewReader(conn), req)64 if err != nil {65 stop()66 _ = conn.Close()67 // A request whose context ended fails with the context's error,68 // as it does through net/http's Transport, not with the error of69 // the connection that the end of the context closed.70 if ended := context.Cause(req.Context()); ended != nil {71 return nil, ended72 }73 return nil, err74 }75 resp.Body = &connBody{ReadCloser: resp.Body, conn: conn, stop: stop}76 return resp, nil77}7879// connBody is a response body that closes its connection.80type connBody struct {81 io.ReadCloser82 conn net.Conn83 stop func() bool84}8586func (b *connBody) Close() error {87 b.stop()88 _ = b.conn.Close()89 return b.ReadCloser.Close()90}
conditions/conditions.go 100.0%
1// Package conditions holds the condition type that the resources of2// every liken component report in their status, and the setter that3// changes a list of them.4//5// A condition answers "what is true now" for automation:6// `kubectl wait --for=condition=Ready` and a controller read it. An7// Event answers "what just happened" for a person, and the events8// package posts one Event for each condition transition9// (events.Recorder.SetCondition). The type is in its own package so10// that an API type package can declare it in its resources and link11// only the time package, not the HTTP client that the events package12// needs.13//14// A component adopts the type with an alias, so its CRDs and its Go15// code keep their names:16//17// type Condition = conditions.Condition18// type ConditionStatus = conditions.Status19package conditions2021import "time"2223// Status is a condition's verdict. Unknown is a third state: the24// component cannot tell yet.25type Status string2627const (28 True Status = "True"29 False Status = "False"30 Unknown Status = "Unknown"31)3233// Condition has the JSON shape of metav1.Condition, so34// `kubectl describe` and `kubectl wait` read it the way they read a35// Pod's conditions. ObservedGeneration records which36// metadata.generation the condition judged. LastTransitionTime is when37// the status last changed, not when the condition was last written.38type Condition struct {39 Type string `json:"type"`40 Status Status `json:"status"`41 ObservedGeneration int64 `json:"observedGeneration,omitempty"`42 Reason string `json:"reason"`43 Message string `json:"message"`44 LastTransitionTime time.Time `json:"lastTransitionTime"`45}4647// Set puts next in the list in place of the condition of its type, or48// appends it when the list holds none.49//50// Set keeps the Kubernetes rule that makes lastTransitionTime51// meaningful: the time moves only when the status changes. So52// `kubectl get` answers "how long has this been Ready?" and not "when53// did the operator last write?". When the status changes, the54// condition takes next's LastTransitionTime, or the present time when55// next states none. The present time is cut to the second, because the56// API server keeps a date-time to the second, and a status composed57// again must compare equal to the stored one.58//59// Set answers whether the condition transitioned: it is new, or its60// status or its reason changed. A new message alone is no transition,61// because a message often carries a value that changes on each pass,62// such as a count or a temperature. The list still takes the message.63func Set(list *[]Condition, next Condition) (transitioned bool) {64 for i := range *list {65 held := &(*list)[i]66 if held.Type != next.Type {67 continue68 }69 statusChanged := held.Status != next.Status70 transitioned = statusChanged || held.Reason != next.Reason71 if statusChanged {72 next.LastTransitionTime = transitionTime(next)73 } else {74 next.LastTransitionTime = held.LastTransitionTime75 }76 *held = next77 return transitioned78 }79 next.LastTransitionTime = transitionTime(next)80 *list = append(*list, next)81 return true82}8384// transitionTime answers the time a condition that changed its status85// records.86func transitionTime(next Condition) time.Time {87 if !next.LastTransitionTime.IsZero() {88 return next.LastTransitionTime89 }90 return time.Now().UTC().Truncate(time.Second)91}9293// Find answers the condition of one type, and false when the list94// holds none.95func Find(list []Condition, conditionType string) (Condition, bool) {96 for _, c := range list {97 if c.Type == conditionType {98 return c, true99 }100 }101 return Condition{}, false102}
election/election.go 96.2%
1// Package election keeps one acting copy of a singleton operator.2//3// A singleton operator creates and deletes objects and writes status4// from one view of the cluster. Two copies that act at once race: both5// create a pod for one object, or one deletes the pod the other just6// made, or two status writes overwrite each other. A rolling update, a7// node partition, and a replica count above one each run a second copy,8// so every singleton operator elects one copy with a coordination.k8s.io9// Lease, through this package.10//11// The election is client-go's leaderelection package with a LeaseLock.12// Only the copy that holds the Lease acts. Every other copy waits and13// reads the Lease, so a rolling update starts the new pod beside the old14// one, and the new pod takes over when the old one releases the Lease.15//16// This package adds four things to client-go's election:17//18// - A second release after client-go's, which survives a renewal that19// lands late (Lead.end).20// - The exit of the process when the Lease is lost.21// - The take of an expired Lease that an earlier process of the same22// pod held (renewalClock.Get).23// - A write guard for the operator's client (Lead.MayWrite).24//25// The election is not fencing. A leader that pauses, for example on a26// stalled node, can resume after its Lease expired and finish a request27// it had already sent. A leader that runs stops sending writes, and28// exits, before another copy can take the Lease, because of the timings29// below and the write guard. A write it sent before that can still land30// after another copy took the Lease (RequestTimeout says why). So a31// write that must not land late names a precondition that the next32// leader moves, such as the resourceVersion of the object it changes.33//34// The package links client-go's typed clientset for the coordination35// group. A program that must stay small, such as the pod build of36// media-operator, never imports it.37package election3839import (40 "context"41 "crypto/rand"42 "encoding/hex"43 "errors"44 "fmt"45 "os"46 "sync/atomic"47 "time"4849 coordination "k8s.io/api/coordination/v1"50 apierrors "k8s.io/apimachinery/pkg/api/errors"51 metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"52 coordinationv1 "k8s.io/client-go/kubernetes/typed/coordination/v1"53 "k8s.io/client-go/rest"54 "k8s.io/client-go/tools/leaderelection"55 "k8s.io/client-go/tools/leaderelection/resourcelock"56)5758// The election's three durations set the Lease's duration to 3059// seconds, twice client-go's default, with client-go's 10-second60// renewal deadline and a 5-second retry period in place of its 261// seconds.62//63// The retry period sets the load. The leader renews once per retry64// period, and a waiting copy reads the Lease once every 5 to 1165// seconds, because client-go adds up to 1.2 retry periods of jitter. A66// cluster can already have several clients that renew a Lease every one67// or two seconds, and an operator does not need faster failover.68//69// The duration sets the safety margin. After its last renewal at T, the70// leader tries again at T+5s, and gives up at T+15s when the renewal71// deadline passes. client-go then releases the Lease, bounded by one72// more renewal deadline, before it calls OnStoppedLeading, so the73// leader exits by T+25s. A waiting copy takes the Lease 30 seconds74// after it saw the last renewal, which is after T+30s. With client-go's75// 15-second duration, the two times would meet.76//77// The write guard is stricter than the elector. It refuses a write once78// the last renewal is one renewal deadline old, at T+10s. The client79// gives up on a write it allowed by T+25s, through RequestTimeout,80// before a new leader can start at T+30s.81//82// The cost is failover time. A waiting copy takes a released Lease on83// its next read, within about 11 seconds, and an abandoned Lease 30 to84// 41 seconds after the last renewal.85const (86 Duration = 30 * time.Second87 RenewDeadline = 10 * time.Second88 RetryPeriod = 5 * time.Second89)9091// RequestTimeout is the longest request the client of an elected92// operator may send: pass it as apiclient.InClusterOptions.Timeout. The93// client gives up on a write that the guard allowed just before T+10s94// by T+25s, one retry period before a new leader can start at T+30s.95// apiclient's default of 30 seconds is too long.96//97// The timeout ends the client's wait, not the write. kube-apiserver98// keeps running a request's handler after the request timed out, and99// counts each such handler in apiserver_request_post_timeout_total, so100// the API server can commit a write after the client gave up on it and101// after a new leader took the Lease. The timeout bounds how long a102// leader waits on a write. It does not keep a late write out.103const RequestTimeout = Duration - RenewDeadline - RetryPeriod104105// Options name the Lease and the process that competes for it.106type Options struct {107 // Name and Namespace name the Lease.108 Name string109 Namespace string110 // Pod is the name of the pod this process runs in. The identity the111 // Lease names is the pod's name and a random suffix, so `kubectl get112 // lease` shows which pod leads.113 Pod string114115 // Report receives each line the election logs. The operator adds116 // its own prefix and writes the line where its log goes.117 Report func(line string)118 // Exit ends the process after a lost Lease, with code 1. It is119 // os.Exit unless a test sets it.120 Exit func(code int)121}122123// Lead is this process's part in the election.124type Lead struct {125 options Options126 elector *leaderelection.LeaderElector127 identity string128 lock *renewalClock129130 // leases reaches the Lease itself, for the release that follows131 // client-go's own. end says why.132 leases coordinationv1.LeasesGetter133134 // started closes when this process takes the Lease, and done closes135 // when the election ends.136 started chan struct{}137 done chan struct{}138 cancel context.CancelFunc139140 // stepping marks an end this process chose, a shutdown, so the end141 // of the election is not a loss.142 stepping atomic.Bool143}144145// New builds the election. The identity is the pod's name and a random146// suffix, so a restarted container is a new candidate.147func New(config *rest.Config, options Options) (*Lead, error) {148 if options.Name == "" || options.Namespace == "" || options.Pod == "" {149 return nil, errors.New("the election needs the Lease's name and namespace and the pod's name")150 }151 if options.Report == nil {152 options.Report = func(string) {}153 }154 if options.Exit == nil {155 options.Exit = os.Exit156 }157 leases, err := coordinationv1.NewForConfig(config)158 if err != nil {159 return nil, err160 }161 suffix := make([]byte, 4)162 _, _ = rand.Read(suffix)163 l := &Lead{164 options: options,165 identity: options.Pod + "_" + hex.EncodeToString(suffix),166 leases: leases,167 started: make(chan struct{}),168 done: make(chan struct{}),169 }170 l.lock = &renewalClock{171 Interface: &resourcelock.LeaseLock{172 LeaseMeta: metav1.ObjectMeta{Name: options.Name, Namespace: options.Namespace},173 Client: leases,174 LockConfig: resourcelock.ResourceLockConfig{Identity: l.identity},175 },176 pod: options.Pod + "_",177 }178 report := options.Report179 subject := l.subject()180 l.elector, err = leaderelection.NewLeaderElector(leaderelection.LeaderElectionConfig{181 Lock: l.lock,182 LeaseDuration: Duration,183 RenewDeadline: RenewDeadline,184 RetryPeriod: RetryPeriod,185 // The release writes the Lease with no holder when the election's186 // context ends, so a waiting copy takes it on its next retry187 // instead of after the Lease's duration. end follows it with a188 // release of its own, for the case where client-go's fails.189 ReleaseOnCancel: true,190 Name: options.Name,191 Callbacks: leaderelection.LeaderCallbacks{192 OnStartedLeading: func(context.Context) {193 report(fmt.Sprintf("holding %s as %s", subject, l.identity))194 close(l.started)195 },196 OnNewLeader: func(holder string) {197 if holder != "" && holder != l.identity {198 report(fmt.Sprintf("waiting for %s, held by %s", subject, holder))199 }200 },201 // client-go calls this whenever the election ends. A loss ends202 // the process at once, because a pass in flight must not keep203 // writing after another copy takes the Lease. The kubelet204 // restarts the container, and it competes as a new candidate.205 OnStoppedLeading: func() {206 if l.stepping.Load() {207 return208 }209 report(fmt.Sprintf("lost %s", subject))210 l.options.Exit(1)211 },212 },213 })214 if err != nil {215 return nil, err216 }217 return l, nil218}219220// ErrStopped says that stop ended before this process could act. A copy221// that never led has nothing to release, so the operator exits cleanly.222var ErrStopped = errors.New("the stop came before the Lease")223224// Start builds the election, runs it, and blocks until this process may225// act. It answers ErrStopped when stop ends first. Every singleton226// operator starts its election this way, with the configuration of its227// pod's service account.228func Start(stop context.Context, config *rest.Config, options Options) (*Lead, error) {229 l, err := New(config, options)230 if err != nil {231 return nil, err232 }233 l.Run()234 if !l.Await(stop) {235 return nil, ErrStopped236 }237 return l, nil238}239240// Identity is the name this process holds the Lease under.241func (l *Lead) Identity() string { return l.identity }242243func (l *Lead) subject() string {244 return "Lease " + l.options.Namespace + "/" + l.options.Name245}246247// Run starts the election in the background.248func (l *Lead) Run() {249 ctx, cancel := context.WithCancel(context.Background())250 l.cancel = cancel251 go func() {252 l.elector.Run(ctx)253 close(l.done)254 }()255}256257// Await returns true when this process holds the Lease and may act. It258// returns false when stop ends first. After a false return the election259// has ended and this process holds nothing.260//261// A copy acts only while it holds the Lease. An error on a request of262// the Lease, a 403 Forbidden among them, is an ordinary failure of the263// election: a waiting copy keeps waiting, and a leader stops writing at264// the write guard and exits when the elector gives up.265func (l *Lead) Await(stop context.Context) bool {266 select {267 case <-l.started:268 case <-stop.Done():269 }270 if stop.Err() == nil {271 return true272 }273 l.end()274 return false275}276277// Leads answers whether the elector has started this process's lead.278func (l *Lead) Leads() bool {279 select {280 case <-l.started:281 return true282 default:283 return false284 }285}286287// ErrNotLeading refuses a write while this process cannot show that it288// holds the Lease.289var ErrNotLeading = errors.New("this copy does not hold a current leader Lease")290291// MayWrite is the client's write guard: pass it to292// apiclient.Client.WithWriteGuard, on a client whose requests end293// within RequestTimeout. It allows a write only while this294// process leads and its last renewal is younger than one renewal295// deadline. The elector itself gives up only after the renewal deadline296// passes with no success, which is up to one retry period later, so the297// guard stops writes first.298func (l *Lead) MayWrite() error {299 if !l.Leads() || l.stepping.Load() {300 return ErrNotLeading301 }302 if age := time.Since(l.lock.lastRenewal()); age >= RenewDeadline {303 return fmt.Errorf("%w: the last renewal was %s ago", ErrNotLeading, age.Round(time.Second))304 }305 return nil306}307308// StepDown ends this process's lead on a shutdown, and releases the309// Lease. The caller calls it after its last write has returned, so310// nothing writes once the Lease is free. A caller that cannot show311// that, because a writer did not stop in time, does not call StepDown:312// the process exits holding the Lease, and a waiting copy takes it when313// it expires, after any write still in flight can land.314func (l *Lead) StepDown() {315 l.end()316}317318// end stops the election and waits for the release. The kubelet kills319// the container at the end of its grace period. client-go bounds its320// release by the renewal deadline, and clearIfHeld is bounded by half of321// it, so the wait ends well before that.322//323// client-go's release can fail on a normal shutdown. The cancel ends324// the renewal loop, but a renewal that was already sent can still reach325// the API server after the release read the Lease. The release's update326// then carries a stale resourceVersion, the API server refuses it with a327// conflict, and a waiting copy waits out the whole duration. So once the328// election has ended, end reads the Lease again and clears it itself if329// it still names this process. The late renewal can land after this330// read too, so a conflict reads the Lease again. The renewal loop sends331// one renewal at a time, so at most one late write is in flight.332func (l *Lead) end() {333 l.stepping.Store(true)334 l.cancel()335 select {336 case <-l.done:337 case <-time.After(RenewDeadline + time.Second):338 return339 }340 ctx, cancel := context.WithTimeout(context.Background(), RenewDeadline/2)341 defer cancel()342 l.clearIfHeld(ctx)343}344345// clearIfHeld writes the Lease with no holder when it still names this346// process, or an earlier process of this pod whose Lease has expired347// (renewalClock.Get says why that one is safe to clear). A stop can348// arrive before the elector's next read takes such a Lease, and without349// this clear a copy in another pod waits out the whole duration. It350// reads the Lease again after a conflict, because the write it lost to351// can be this process's own late renewal. A Lease that names another352// process, or no process, is left alone.353func (l *Lead) clearIfHeld(ctx context.Context) {354 leases := l.leases.Leases(l.options.Namespace)355 subject := l.subject()356 for range 3 {357 lease, err := leases.Get(ctx, l.options.Name, metav1.GetOptions{})358 if err != nil {359 if !apierrors.IsNotFound(err) {360 l.options.Report(fmt.Sprintf("reading %s to release it: %v", subject, err))361 }362 return363 }364 if !l.mayClear(lease) {365 return366 }367 none, released, second := "", metav1.NewMicroTime(time.Now()), int32(1)368 lease.Spec.HolderIdentity = &none369 lease.Spec.LeaseDurationSeconds = &second370 lease.Spec.RenewTime = &released371 _, err = leases.Update(ctx, lease, metav1.UpdateOptions{})372 if err == nil {373 l.options.Report(fmt.Sprintf("released %s", subject))374 return375 }376 if !apierrors.IsConflict(err) {377 l.options.Report(fmt.Sprintf("releasing %s: %v", subject, err))378 return379 }380 }381}382383// mayClear answers whether clearIfHeld may write the Lease with no384// holder.385func (l *Lead) mayClear(lease *coordination.Lease) bool {386 if lease.Spec.HolderIdentity == nil || *lease.Spec.HolderIdentity == "" {387 return false388 }389 holder := *lease.Spec.HolderIdentity390 if holder == l.identity {391 return true392 }393 if lease.Spec.RenewTime == nil || lease.Spec.LeaseDurationSeconds == nil {394 return false395 }396 return l.lock.abandonedByThisPod(holder, lease.Spec.RenewTime.Time,397 int(*lease.Spec.LeaseDurationSeconds), time.Now())398}
election/electiontest/server.go 98.6%
1// Package electiontest is a fake API server that holds one Lease, for2// the tests of an election. It answers the way the API server does: a3// create conflicts with a Lease that exists, and an update from a stale4// resourceVersion conflicts with a newer write. A test serves it through5// apiservertest, so the election runs in a synctest bubble on the fake6// clock.7package electiontest89import (10 "encoding/json"11 "net/http"12 "strconv"13 "strings"14 "sync"15 "testing"16 "time"1718 "k8s.io/client-go/rest"1920 "github.com/liken-sh/liken/kubernetes/apiservertest"21 "github.com/liken-sh/liken/kubernetes/election"22)2324// Server holds the Lease named Name in Namespace.25type Server struct {26 Name string27 Namespace string2829 mu sync.Mutex30 current map[string]any31 version int32 // landAfterRead writes the Lease again, at a new version, right after33 // the next read is answered. It is the state a renewal leaves when34 // its client gave up on it and the API server applied it anyway.35 landAfterRead bool36 refuse map[string]int37 requests map[string]int38}3940// New is a server that holds no Lease yet.41func New(name, namespace string) *Server {42 return &Server{Name: name, Namespace: namespace, refuse: map[string]int{}, requests: map[string]int{}}43}4445// Config is a client configuration for the server, the way46// rest.InClusterConfig is one for the real API server. The typed client47// sends protobuf to the real API server, and JSON here, because the fake48// decodes JSON alone.49func (s *Server) Config(t *testing.T) *rest.Config {50 t.Helper()51 config := apiservertest.Start(t, s).Config()52 config.ContentType = "application/json"53 return config54}5556// store writes a Lease at the next version. Callers hold mu.57func (s *Server) store(object map[string]any) map[string]any {58 s.version++59 metadata, _ := object["metadata"].(map[string]any)60 metadata["resourceVersion"] = strconv.Itoa(s.version)61 s.current = object62 return object63}6465// HoldAs writes the Lease as another process would, with a renewal that66// is fresh now.67func (s *Server) HoldAs(holder string) {68 s.HoldAsRenewedAt(holder, time.Now())69}7071// HoldAsRenewedAt writes the Lease as a process that last renewed it at72// renewed, such as a process that has since exited.73func (s *Server) HoldAsRenewedAt(holder string, renewed time.Time) {74 s.mu.Lock()75 defer s.mu.Unlock()76 now := renewed.UTC().Format("2006-01-02T15:04:05.000000Z07:00")77 s.store(map[string]any{78 "apiVersion": "coordination.k8s.io/v1",79 "kind": "Lease",80 "metadata": map[string]any{"name": s.Name, "namespace": s.Namespace},81 "spec": map[string]any{82 "holderIdentity": holder,83 "leaseDurationSeconds": int(election.Duration / time.Second),84 "acquireTime": now,85 "renewTime": now,86 },87 })88}8990// Holder is the identity the Lease names, and empty when no Lease91// exists or it names no holder.92func (s *Server) Holder() string {93 s.mu.Lock()94 defer s.mu.Unlock()95 if s.current == nil {96 return ""97 }98 spec, _ := s.current["spec"].(map[string]any)99 holder, _ := spec["holderIdentity"].(string)100 return holder101}102103// LandAWriteAfterTheNextRead makes the next read of the Lease followed104// by a write that the reader did not see.105func (s *Server) LandAWriteAfterTheNextRead() {106 s.mu.Lock()107 defer s.mu.Unlock()108 s.landAfterRead = true109}110111// Refuse answers every request of the method with the status code,112// until Allow.113func (s *Server) Refuse(method string, status int) {114 s.mu.Lock()115 defer s.mu.Unlock()116 s.refuse[method] = status117}118119// Allow ends the refusal of the method.120func (s *Server) Allow(method string) {121 s.mu.Lock()122 defer s.mu.Unlock()123 delete(s.refuse, method)124}125126// Count is how many requests of the method the server received.127func (s *Server) Count(method string) int {128 s.mu.Lock()129 defer s.mu.Unlock()130 return s.requests[method]131}132133func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {134 s.mu.Lock()135 defer s.mu.Unlock()136 s.requests[r.Method]++137 w.Header().Set("Content-Type", "application/json")138 if code, refused := s.refuse[r.Method]; refused {139 // The API server names the reason and the code in the Status140 // it answers, and client-go reads both.141 reasons := map[int]string{http.StatusForbidden: "Forbidden", http.StatusNotFound: "NotFound"}142 reason, known := reasons[code]143 if !known {144 reason = "InternalError"145 }146 status(w, code, reason)147 return148 }149 body := map[string]any{}150 if r.Method == http.MethodPost || r.Method == http.MethodPut {151 _ = json.NewDecoder(r.Body).Decode(&body)152 }153 switch {154 case r.Method == http.MethodGet && strings.HasSuffix(r.URL.Path, "/leases/"+s.Name):155 if s.current == nil {156 status(w, http.StatusNotFound, "NotFound")157 return158 }159 _ = json.NewEncoder(w).Encode(s.current)160 if s.landAfterRead {161 s.landAfterRead = false162 s.store(s.current)163 }164 case r.Method == http.MethodPost:165 if s.current != nil {166 status(w, http.StatusConflict, "AlreadyExists")167 return168 }169 w.WriteHeader(http.StatusCreated)170 _ = json.NewEncoder(w).Encode(s.store(body))171 case r.Method == http.MethodPut:172 if s.current == nil {173 status(w, http.StatusNotFound, "NotFound")174 return175 }176 sent, _ := body["metadata"].(map[string]any)177 stored, _ := s.current["metadata"].(map[string]any)178 if sent["resourceVersion"] != stored["resourceVersion"] {179 status(w, http.StatusConflict, "Conflict")180 return181 }182 _ = json.NewEncoder(w).Encode(s.store(body))183 default:184 status(w, http.StatusNotFound, "NotFound")185 }186}187188// status answers the Status object the API server sends with an error,189// which is what client-go reads the reason from.190func status(w http.ResponseWriter, code int, reason string) {191 w.WriteHeader(code)192 _ = json.NewEncoder(w).Encode(map[string]any{193 "kind": "Status", "apiVersion": "v1", "status": "Failure", "reason": reason, "code": code,194 })195}
election/lock.go 100.0%
1package election23import (4 "context"5 "strings"6 "sync"7 "time"89 "k8s.io/client-go/tools/leaderelection/resourcelock"10)1112// renewalClock is the Lease lock with two additions. It records when13// this process last sent a write of the Lease that the API server14// accepted with this process as the holder, for the write guard. The15// time is taken before the request is sent, so it is never later than16// the time a waiting copy first reads the renewal and measures the17// duration from. And it frees the expired Lease of an earlier process18// of this pod (Get).19type renewalClock struct {20 resourcelock.Interface2122 mu sync.Mutex23 renewed time.Time2425 // pod is the prefix that every identity of this pod starts with.26 // A pod's name is a DNS subdomain, which has no underscore, so the27 // prefix names this pod and no other.28 pod string29}3031func (r *renewalClock) Create(ctx context.Context, record resourcelock.LeaderElectionRecord) error {32 sent := time.Now()33 err := r.Interface.Create(ctx, record)34 r.mark(sent, record, err)35 return err36}3738func (r *renewalClock) Update(ctx context.Context, record resourcelock.LeaderElectionRecord) error {39 sent := time.Now()40 err := r.Interface.Update(ctx, record)41 r.mark(sent, record, err)42 return err43}4445// Get reads the Lease, and shows client-go a Lease with no holder when46// an earlier process of this pod held it and its last renewal is one47// Lease duration old.48//49// client-go does not compare renewTime with its own clock, because the50// holder can run on another node with another clock. It measures the51// duration from the time it first read the current record. A process52// that restarts, or that starts after the API server comes back, first53// reads its earlier process's Lease late, and waits a whole duration54// from that read. When the only API server reboots, the leader cannot55// renew, exits, and restarts, and then nothing acts for a whole Lease56// duration after the API server returns. A rolling update in that time57// moves the wait to the new pod, which cannot clear a Lease that names58// another process.59//60// An earlier process of this pod is not a paused leader that can61// resume. The kubelet starts a container's process again only after62// the one before it ended, so that process sends no more writes. It63// wrote renewTime from the clock of this node, so the comparison with64// this process's clock is sound.65//66// A copy in another pod still waits the whole duration, and this67// package accepts that wait. client-go documents it:68// LeaderElectionConfig.LeaseDuration is how long a candidate that does69// not lead waits before it forces the take, measured from the last70// renewal it observed. Another pod cannot tell a holder that ended from71// one that paused. A missing pod does not prove that the holder ended:72// a force delete, or the out-of-service taint on its node, removes the73// pod from the API while its process can still run. The wait happens74// at most once for each restart of the only API server of a fleet.75//76// The Lease that client-go then updates keeps the resourceVersion of77// this read, so the take is a conditional write, and a conflict with78// any other writer ends it.79func (r *renewalClock) Get(ctx context.Context) (*resourcelock.LeaderElectionRecord, []byte, error) {80 record, raw, err := r.Interface.Get(ctx)81 if err == nil && r.abandonedByThisPod(record.HolderIdentity, record.RenewTime.Time,82 record.LeaseDurationSeconds, time.Now()) {83 free := *record84 free.HolderIdentity = ""85 return &free, raw, nil86 }87 return record, raw, err88}8990// abandonedByThisPod answers whether holder is an earlier process of91// this pod, whose last renewal at renewed is at least the Lease's92// duration old at now.93func (r *renewalClock) abandonedByThisPod(holder string, renewed time.Time, seconds int, now time.Time) bool {94 if holder == r.Identity() || !strings.HasPrefix(holder, r.pod) {95 return false96 }97 return !now.Before(renewed.Add(time.Duration(seconds) * time.Second))98}99100func (r *renewalClock) mark(sent time.Time, record resourcelock.LeaderElectionRecord, err error) {101 if err != nil || record.HolderIdentity != r.Identity() {102 return103 }104 r.mu.Lock()105 defer r.mu.Unlock()106 r.renewed = sent107}108109func (r *renewalClock) lastRenewal() time.Time {110 r.mu.Lock()111 defer r.mu.Unlock()112 return r.renewed113}
events/event.go 100.0%
1package events23import (4 "time"5 "unicode/utf8"6)78// ObjectReference names the object an Event is about. It is the9// involvedObject of a core/v1 Event, and `kubectl describe` finds the10// object's Events by its kind, name, and UID. An empty Namespace11// names a cluster-scoped object.12type ObjectReference struct {13 APIVersion string `json:"apiVersion,omitempty"`14 Kind string `json:"kind"`15 Namespace string `json:"namespace,omitempty"`16 Name string `json:"name"`17 UID string `json:"uid,omitempty"`18}1920// The two types of Event. Warning means a person may need to act.21// Normal means an expected transition or an action the component took.22const (23 TypeNormal = "Normal"24 TypeWarning = "Warning"25)2627// Event is a core/v1 Event as the recorder writes it. It states no28// eventTime: the API server applies the strict checks of29// events.k8s.io/v1 only to an Event that states one, and those checks30// refuse a note longer than 1024 bytes.31type Event struct {32 APIVersion string `json:"apiVersion"`33 Kind string `json:"kind"`34 Metadata Metadata `json:"metadata"`35 InvolvedObject ObjectReference `json:"involvedObject"`36 Reason string `json:"reason"`37 Message string `json:"message"`38 Type string `json:"type"`39 Source Source `json:"source"`40 FirstTimestamp time.Time `json:"firstTimestamp"`41 LastTimestamp time.Time `json:"lastTimestamp"`42 Count int32 `json:"count"`43 ReportingComponent string `json:"reportingComponent"`44 ReportingInstance string `json:"reportingInstance"`45}4647// Metadata is the part of an Event's metadata that the recorder writes48// and reads back. The recorder sends generateName, and the API server49// answers the name it chose, which a repeat patches.50type Metadata struct {51 Name string `json:"name,omitempty"`52 GenerateName string `json:"generateName,omitempty"`53 Namespace string `json:"namespace,omitempty"`54 ResourceVersion string `json:"resourceVersion,omitempty"`55}5657// Source is the older form of reportingComponent and58// reportingInstance. `kubectl describe` prints it in the From column.59type Source struct {60 Component string `json:"component,omitempty"`61 Host string `json:"host,omitempty"`62}6364// The limits that the API server's validation states for an Event's65// reason and note. The recorder cuts a longer reason or message to fit,66// so the API server never refuses an Event for its length.67const (68 maxReason = 12869 maxMessage = 102470)7172// clusterNamespace holds the Events about cluster-scoped objects. The73// API server accepts such an Event only in default or kube-system.74const clusterNamespace = "default"7576// namespaceOf answers the namespace that holds the Events about an77// object.78func namespaceOf(object ObjectReference) string {79 if object.Namespace == "" {80 return clusterNamespace81 }82 return object.Namespace83}8485// cut shortens text to at most limit bytes, at a character boundary,86// and marks a cut with an ellipsis.87func cut(text string, limit int) string {88 if len(text) <= limit {89 return text90 }91 const ellipsis = "…"92 end := limit - len(ellipsis)93 for end > 0 && !utf8.RuneStart(text[end]) {94 end--95 }96 return text[:end] + ellipsis97}
events/eventstest/events.go 100.0%
1// Package eventstest is a fake of the core/v1 events collection, for2// a test's fake API server. Every component's tests read the Events3// that events.Recorder wrote from it, so no component keeps its own4// copy of an Event store. A test serves it, with the rest of its fake,5// through apiservertest, so the test runs in a synctest bubble.6package eventstest78import (9 "encoding/json"10 "fmt"11 "net/http"12 "regexp"13 "slices"14 "strconv"15 "sync"1617 "github.com/liken-sh/liken/kubernetes/events"18)1920// eventPath matches the path of the events collection of a namespace,21// and of one Event in it.22var eventPath = regexp.MustCompile(`^/api/v1/namespaces/[^/]+/events(/[^/]+)?$`)2324// Events holds the Events a test's fake API server receives. It25// answers a create and a merge patch the way the API server does: a26// create names the Event from its generateName, a patch of an Event27// the server does not hold answers 404, and an Event about a28// cluster-scoped object is refused outside default and kube-system.29// The zero value holds no Events and is ready to serve.30type Events struct {31 mu sync.Mutex32 held []events.Event33 created int34 refusals int35 mux *http.ServeMux36}3738// Around answers a handler that serves each Event request itself and39// hands every other request to next, the test's own fake:40//41// recorded := &eventstest.Events{}42// server := apiservertest.Start(t, recorded.Around(fake))43func (e *Events) Around(next http.Handler) http.Handler {44 return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {45 if eventPath.MatchString(r.URL.Path) {46 e.ServeHTTP(w, r)47 return48 }49 next.ServeHTTP(w, r)50 })51}5253// ServeHTTP answers one request to the events collection.54func (e *Events) ServeHTTP(w http.ResponseWriter, r *http.Request) {55 e.mu.Lock()56 if e.mux == nil {57 e.mux = http.NewServeMux()58 e.mux.HandleFunc("POST /api/v1/namespaces/{namespace}/events", e.create)59 e.mux.HandleFunc("PATCH /api/v1/namespaces/{namespace}/events/{name}", e.patch)60 }61 mux := e.mux62 e.mu.Unlock()63 mux.ServeHTTP(w, r)64}6566// Refuse makes the server answer the next count writes with 503, the67// way an API server that restarts does.68func (e *Events) Refuse(count int) {69 e.mu.Lock()70 defer e.mu.Unlock()71 e.refusals = count72}7374// Expire deletes every Event, the way the API server's TTL does an75// hour after an Event's last write.76func (e *Events) Expire() {77 e.mu.Lock()78 defer e.mu.Unlock()79 e.held = nil80}8182// List answers every Event the server holds, in the order they were83// created.84func (e *Events) List() []events.Event {85 e.mu.Lock()86 defer e.mu.Unlock()87 return slices.Clone(e.held)88}8990// About answers the Events about the object of one kind and name, in91// the order they were created.92func (e *Events) About(kind, name string) []events.Event {93 var out []events.Event94 for _, event := range e.List() {95 if event.InvolvedObject.Kind == kind && event.InvolvedObject.Name == name {96 out = append(out, event)97 }98 }99 return out100}101102// refused answers whether the server refuses this write. The caller103// holds e.mu.104func (e *Events) refused(w http.ResponseWriter) bool {105 if e.refusals == 0 {106 return false107 }108 e.refusals--109 http.Error(w, "the API server is restarting", http.StatusServiceUnavailable)110 return true111}112113func (e *Events) create(w http.ResponseWriter, r *http.Request) {114 var event events.Event115 if err := json.NewDecoder(r.Body).Decode(&event); err != nil {116 http.Error(w, err.Error(), http.StatusBadRequest)117 return118 }119 namespace := r.PathValue("namespace")120 involved := event.InvolvedObject.Namespace121 switch {122 case event.Metadata.Namespace != namespace:123 http.Error(w, "the namespace of the Event does not match the namespace of the request", http.StatusBadRequest)124 return125 case involved == "" && namespace != "default" && namespace != "kube-system":126 http.Error(w, "an Event about a cluster-scoped object must be in default or kube-system", http.StatusUnprocessableEntity)127 return128 case involved != "" && involved != namespace:129 http.Error(w, "the involved object's namespace does not match the Event's", http.StatusUnprocessableEntity)130 return131 }132 e.mu.Lock()133 defer e.mu.Unlock()134 if e.refused(w) {135 return136 }137 e.created++138 if event.Metadata.Name == "" {139 event.Metadata.Name = event.Metadata.GenerateName + fmt.Sprintf("%05x", e.created)140 }141 event.Metadata.ResourceVersion = strconv.Itoa(e.created)142 e.held = append(e.held, event)143 w.WriteHeader(http.StatusCreated)144 _ = json.NewEncoder(w).Encode(event)145}146147func (e *Events) patch(w http.ResponseWriter, r *http.Request) {148 var patch map[string]any149 if err := json.NewDecoder(r.Body).Decode(&patch); err != nil || r.Header.Get("Content-Type") != "application/merge-patch+json" {150 http.Error(w, "the fake answers a JSON merge patch", http.StatusUnsupportedMediaType)151 return152 }153 e.mu.Lock()154 defer e.mu.Unlock()155 if e.refused(w) {156 return157 }158 at := slices.IndexFunc(e.held, func(held events.Event) bool {159 return held.Metadata.Namespace == r.PathValue("namespace") && held.Metadata.Name == r.PathValue("name")160 })161 if at < 0 {162 http.Error(w, "the Event does not exist", http.StatusNotFound)163 return164 }165 var document map[string]any166 held, _ := json.Marshal(e.held[at])167 _ = json.Unmarshal(held, &document)168 merged, _ := json.Marshal(mergePatch(document, patch))169 var event events.Event170 if err := json.Unmarshal(merged, &event); err != nil {171 http.Error(w, err.Error(), http.StatusUnprocessableEntity)172 return173 }174 e.held[at] = event175 _ = json.NewEncoder(w).Encode(event)176}177178// mergePatch applies a JSON merge patch (RFC 7386) to a document: a179// null removes a field, an object merges into the object it replaces,180// and any other value replaces the field.181func mergePatch(document, patch map[string]any) map[string]any {182 for key, value := range patch {183 inner, isObject := value.(map[string]any)184 held, heldObject := document[key].(map[string]any)185 switch {186 case value == nil:187 delete(document, key)188 case isObject && heldObject:189 document[key] = mergePatch(held, inner)190 default:191 document[key] = value192 }193 }194 return document195}
events/recorder.go 100.0%
1// Package events posts Kubernetes Events about the objects a liken2// component manages, so `kubectl describe` shows a person what3// happened to an object in the last hour.4//5// A status condition answers "what is true now", for automation. An6// Event answers "what just happened", for a person. A log line holds7// every attempt and detail. So a component posts one Event for each8// condition transition (SetCondition, or Transition after a status9// write) and one for each action it takes that changes no condition,10// such as a pod created again or a reboot requested. It never posts a11// reading, a key press, or each retry of a retry loop.12//13// A component builds one Recorder in main and passes it down:14//15// recorder := events.New(ctx, client, "observatory-operator", events.Options{})16// recorder.Warning(object, "StepFailed", "Connect failed: the mount did not answer")17// recorder.SetCondition(object, &status.Conditions, ready, conditions.False)18//19// The recorder writes core/v1 Events through apiclient, and imports20// nothing from k8s.io, so a build that must not link client-go can21// link it. It leaves out eventTime, because the API server applies the22// strict checks of events.k8s.io/v1, which refuse a note longer than23// 1024 bytes, only to an Event that states one. `kubectl describe` and24// `kubectl events` read core/v1.25//26// A write is best effort. Normal and Warning put the Event on a27// bounded queue and return at once, and one goroutine writes the28// queue, so a reconcile pass never waits on an Event and never fails29// because of one. The goroutine sends a failed write again (attempts),30// and logs one line when the last attempt fails. A full queue drops31// the Event and counts it (Dropped).32//33// A repeat of the same Event on the same object within the window34// patches count and lastTimestamp on the Event the recorder wrote, in35// place of a new Event, so `kubectl describe` prints one line such as36// "(x37 over 1h)" in place of 37 lines.37//38// The API server deletes an Event one hour after its last write39// (kube-apiserver's --event-ttl, which k3s does not change). So a40// condition or a status field, not an Event, holds a fact that must41// last longer.42//43// The Events about a cluster-scoped object go in the namespace44// default, because the API server accepts them only in default or45// kube-system. `kubectl describe` finds them there. `kubectl events46// --for` finds them only with -n default or -A.47//48// A component needs the RBAC verbs create and patch on events in each49// namespace it posts to.50//51// The recorder's goroutine ends with the context New takes, and a52// request in flight then ends at once. An Event still queued is lost,53// and a request that the client cut off can still land on the API54// server after the goroutine returns. A component that must know when55// its last write has returned, such as an elected operator that56// releases its Lease after its last write, passes New a context that57// outlives its last pass, and calls Shutdown after that pass:58//59// ctx, cancel := context.WithTimeout(context.Background(), events.ShutdownTimeout)60// defer cancel()61// if recorder.Shutdown(ctx) == nil {62// lead.StepDown()63// }64//65// A test serves eventstest.Events in front of its fake API server,66// through apiservertest, and reads what the recorder wrote from it. A67// test passes New t.Context().68package events6970import (71 "context"72 "encoding/json"73 "errors"74 "fmt"75 "io"76 "net/http"77 "os"78 "sync"79 "sync/atomic"80 "time"8182 "github.com/liken-sh/liken/kubernetes/apiclient"83)8485// queueLength bounds the Events that wait to be written. A component86// posts a few Events for each object in an hour, so a full queue87// means the API server stopped answering, and the recorder drops new88// Events instead of holding memory for them.89const queueLength = 2569091// attempts is how many times the recorder sends one write, and92// retryWait is the clock between two sends. An API server that93// restarts refuses connections for some seconds, and three sends 1094// seconds apart cover that without holding the queue for long. This95// is the interval of client-go's own recorder, which sends 12 times.96const (97 attempts = 398 retryWait = 10 * time.Second99)100101// Recorder writes the Events of one component.102type Recorder struct {103 client *apiclient.Client104 component string105 instance string106 log io.Writer107108 queue chan Event109 dropped atomic.Uint64110111 // mu guards stopping, so that no post queues an Event after112 // Shutdown began and the goroutine may have read the queue for the113 // last time. stop closes when Shutdown begins, abandon when the114 // context of Shutdown ends, and done when the goroutine returns.115 mu sync.Mutex116 stopping bool117 stop chan struct{}118 abandon chan struct{}119 abandonOnce sync.Once120 done chan struct{}121122 // series is read and written only by the goroutine that writes123 // the queue.124 series *series125}126127// Options are the optional parts of a Recorder.128type Options struct {129 // Instance is the reportingInstance of each Event: the name of the130 // pod, or of the node for a DaemonSet that runs on the host's131 // network. Empty means the host name, which is the pod's name in a132 // pod.133 Instance string134135 // Log receives one line for each Event the recorder could not136 // write, and for the first Event it dropped. Nil means standard137 // error.138 Log io.Writer139}140141// New starts a recorder that writes until ctx ends, or until Shutdown142// has written the queue. component is the143// reportingComponent of each Event, such as machine-operator.144func New(ctx context.Context, client *apiclient.Client, component string, options Options) *Recorder {145 instance := options.Instance146 if instance == "" {147 instance, _ = os.Hostname()148 }149 log := options.Log150 if log == nil {151 log = os.Stderr152 }153 r := &Recorder{154 client: client.WithContext(ctx),155 component: component,156 instance: instance,157 log: log,158 queue: make(chan Event, queueLength),159 series: newSeries(),160 stop: make(chan struct{}),161 abandon: make(chan struct{}),162 done: make(chan struct{}),163 }164 go r.run(ctx)165 return r166}167168// Normal posts an Event about an expected transition or an action the169// component took. A nil recorder posts nothing.170func (r *Recorder) Normal(object ObjectReference, reason, message string) {171 r.post(object, TypeNormal, reason, message)172}173174// Warning posts an Event about something a person may need to act on.175// A nil recorder posts nothing.176func (r *Recorder) Warning(object ObjectReference, reason, message string) {177 r.post(object, TypeWarning, reason, message)178}179180// Dropped answers how many Events the recorder dropped because its181// queue was full.182func (r *Recorder) Dropped() uint64 {183 if r == nil {184 return 0185 }186 return r.dropped.Load()187}188189// post builds an Event at the time of the call, and queues it.190func (r *Recorder) post(object ObjectReference, eventType, reason, message string) {191 if r == nil {192 return193 }194 now := time.Now().UTC().Truncate(time.Second)195 event := Event{196 APIVersion: "v1",197 Kind: "Event",198 Metadata: Metadata{199 GenerateName: object.Name + ".",200 Namespace: namespaceOf(object),201 },202 InvolvedObject: object,203 Reason: cut(reason, maxReason),204 Message: cut(message, maxMessage),205 Type: eventType,206 Source: Source{Component: r.component, Host: r.instance},207 FirstTimestamp: now,208 LastTimestamp: now,209 Count: 1,210 ReportingComponent: r.component,211 ReportingInstance: r.instance,212 }213 r.mu.Lock()214 defer r.mu.Unlock()215 if r.stopping {216 r.logf("dropped the %s Event %s on the %s %s: the recorder has shut down", eventType, reason, object.Kind, object.Name)217 return218 }219 select {220 case r.queue <- event:221 default:222 if r.dropped.Add(1) == 1 {223 r.logf("dropped the %s Event %s on the %s %s: %d Events wait to be written", eventType, reason, object.Kind, object.Name, queueLength)224 }225 }226}227228// ShutdownTimeout is the time an operator gives Shutdown. A queue of a229// few Events writes in milliseconds, and an API server that does not230// answer in 5 seconds will not answer before the pod's grace period231// ends. The rest of the default grace period of 30 seconds covers the232// release of a Lease, which kubernetes/election bounds at 16 seconds.233const ShutdownTimeout = 5 * time.Second234235// Shutdown stops the recorder and writes the Events it holds. It236// follows client-go's event broadcaster: a later Normal or Warning237// drops its Event, each Event still queued is sent once, and a write238// that fails is logged and not sent again. An Event that waits to be239// sent again is given up at once. Shutdown answers nil when the queue240// is empty and the last request has returned, so the caller's next241// step, such as the release of a Lease, comes after every write of the242// recorder.243//244// Unlike client-go's Shutdown, it waits for that last request. When ctx245// ends first, Shutdown answers ctx's error at once, and a request can246// still be in flight. The recorder then sends nothing after that247// request returns, and logs how many Events it did not write.248func (r *Recorder) Shutdown(ctx context.Context) error {249 if r == nil {250 return nil251 }252 r.mu.Lock()253 if !r.stopping {254 r.stopping = true255 close(r.stop)256 }257 r.mu.Unlock()258 select {259 case <-r.done:260 return nil261 case <-ctx.Done():262 }263 r.abandonOnce.Do(func() { close(r.abandon) })264 return fmt.Errorf("the Event recorder's last write had not returned: %w", ctx.Err())265}266267// run writes the queue until ctx ends, or until Shutdown has emptied268// it.269func (r *Recorder) run(ctx context.Context) {270 defer close(r.done)271 for {272 // A Shutdown takes the queue to drain, which stops when the273 // context of Shutdown ends. The select below picks at random274 // when the queue and stop are both ready, so stop is read275 // first.276 select {277 case <-r.stop:278 r.drain(ctx)279 return280 default:281 }282 select {283 case <-ctx.Done():284 return285 case <-r.stop:286 r.drain(ctx)287 return288 case event := <-r.queue:289 r.write(ctx, event)290 }291 }292}293294// drain sends each Event left in the queue once, after Shutdown. No295// post adds to the queue once Shutdown began, so an empty queue is the296// end.297func (r *Recorder) drain(ctx context.Context) {298 for {299 select {300 case <-ctx.Done():301 return302 case <-r.abandon:303 if left := len(r.queue); left > 0 {304 r.logf("the shutdown ran out of time, so these queued Events were not written: %d", left)305 }306 return307 default:308 }309 select {310 case event := <-r.queue:311 if err := r.writeOnce(event); err != nil {312 r.logFailure(event, err)313 }314 default:315 return316 }317 }318}319320// write sends one Event, and sends it again after a failure, up to321// attempts times.322func (r *Recorder) write(ctx context.Context, event Event) {323 for attempt := 1; ; attempt++ {324 err := r.writeOnce(event)325 if err == nil {326 return327 }328 if attempt == attempts {329 r.logFailure(event, err)330 return331 }332 timer := time.NewTimer(retryWait)333 select {334 case <-ctx.Done():335 timer.Stop()336 return337 case <-r.stop:338 timer.Stop()339 r.logFailure(event, err)340 return341 case <-timer.C:342 }343 }344}345346// writeOnce patches the Event of a series that is still open, or347// creates a new Event.348func (r *Recorder) writeOnce(event Event) error {349 key := keyOf(event)350 collection := "/api/v1/namespaces/" + event.Metadata.Namespace + "/events"351 if held, ok := r.series.get(key); ok && event.LastTimestamp.Sub(held.last) < window {352 count := held.count + 1353 // Strings, counts, and times of this century always encode,354 // so neither Marshal below can fail.355 patch, _ := json.Marshal(map[string]any{"count": count, "lastTimestamp": event.LastTimestamp})356 err := r.client.Request(http.MethodPatch, collection+"/"+held.name, "application/merge-patch+json", patch, nil)357 if err == nil {358 r.series.put(key, entry{name: held.name, count: count, last: event.LastTimestamp})359 return nil360 }361 if !errors.Is(err, apiclient.ErrNotFound) {362 return err363 }364 // The TTL deleted the Event, so the series starts again.365 r.series.remove(key)366 }367 body, _ := json.Marshal(event)368 var created Event369 if err := r.client.RequestJSON(http.MethodPost, collection, body, &created); err != nil {370 return err371 }372 r.series.put(key, entry{name: created.Metadata.Name, count: 1, last: event.LastTimestamp})373 return nil374}375376// logFailure logs the last failure of an Event that the recorder gives377// up on.378func (r *Recorder) logFailure(event Event, err error) {379 r.logf("writing the %s Event %s on the %s %s: %v", event.Type, event.Reason, event.InvolvedObject.Kind, event.InvolvedObject.Name, err)380}381382func (r *Recorder) logf(format string, args ...any) {383 fmt.Fprintf(r.log, r.component+": "+format+"\n", args...)384}
events/series.go 100.0%
1package events23import (4 "container/list"5 "time"6)78// window is how long after the last Event of a series a repeat still9// patches that Event. A repeat after a longer gap is a new occurrence,10// and a new Event shows where it began. client-go's recorder uses the11// same 10 minutes.12const window = 10 * time.Minute1314// seriesLimit bounds the series the recorder remembers. A component15// with more open series than this forgets the oldest, and its next16// repeat starts a new Event. client-go's recorder uses the same 4096.17const seriesLimit = 40961819// entry is the Event that a series patches: its name, the count it20// holds, and the time of its last repeat.21type entry struct {22 name string23 count int3224 last time.Time25}2627// series remembers the Event of each series, by key, and forgets the28// one used longest ago when it holds seriesLimit.29type series struct {30 order *list.List31 byKey map[string]*list.Element32}3334type seriesItem struct {35 key string36 entry entry37}3839func newSeries() *series {40 return &series{order: list.New(), byKey: map[string]*list.Element{}}41}4243// keyOf answers the key of an Event's series: the object, the type,44// the reason, and the message. The UID tells apart two objects of one45// name, one deleted and one created after it.46func keyOf(event Event) string {47 object := event.InvolvedObject48 return object.Kind + "\x00" + object.Namespace + "\x00" + object.Name + "\x00" + object.UID + "\x00" +49 event.Type + "\x00" + event.Reason + "\x00" + event.Message50}5152func (s *series) get(key string) (entry, bool) {53 element, ok := s.byKey[key]54 if !ok {55 return entry{}, false56 }57 s.order.MoveToFront(element)58 return element.Value.(*seriesItem).entry, true59}6061func (s *series) put(key string, e entry) {62 if element, ok := s.byKey[key]; ok {63 element.Value.(*seriesItem).entry = e64 s.order.MoveToFront(element)65 return66 }67 s.byKey[key] = s.order.PushFront(&seriesItem{key: key, entry: e})68 if s.order.Len() > seriesLimit {69 oldest := s.order.Back()70 s.order.Remove(oldest)71 delete(s.byKey, oldest.Value.(*seriesItem).key)72 }73}7475func (s *series) remove(key string) {76 if element, ok := s.byKey[key]; ok {77 s.order.Remove(element)78 delete(s.byKey, key)79 }80}
events/transition.go 100.0%
1package events23import "github.com/liken-sh/liken/kubernetes/conditions"45// SetCondition puts next in the list with conditions.Set, and posts one6// Event when the condition transitioned: it is new, or its status or7// its reason changed. It answers whether it transitioned.8//9// A caller that writes the status later, and must post only when the10// write lands, calls conditions.Set while it composes and Transition11// after the write.12func (r *Recorder) SetCondition(object ObjectReference, list *[]conditions.Condition, next conditions.Condition, bad conditions.Status) bool {13 transitioned := conditions.Set(list, next)14 if transitioned {15 r.Transition(object, next, bad)16 }17 return transitioned18}1920// Transition posts the Event of one condition transition, with the21// condition's reason and message, so `kubectl get -o yaml` and22// `kubectl describe` show the same words.23//24// bad is the status that needs a person: a transition to it is a25// Warning, and any other is Normal. Pass conditions.False for a26// condition such as Ready, conditions.True for one such as Stalled,27// and "" for a condition no status of which needs a person. A caller28// whose condition is bad only for some reasons, such as a Ready that is29// False both while a device starts and when it fails, passes bad only30// with the failing reasons.31func (r *Recorder) Transition(object ObjectReference, c conditions.Condition, bad conditions.Status) {32 if bad != "" && c.Status == bad {33 r.Warning(object, c.Reason, c.Message)34 return35 }36 r.Normal(object, c.Reason, c.Message)37}
informer/absent.go 100.0%
1package informer23// A kind can be defined by another operator, which a cluster may not4// run. The watch of such a kind reads the API server's 404 as an empty5// collection, so the operator starts and runs on a cluster without the6// definition, and a pass reads the kind from its copy with no request.7//8// Nothing the operator watches reports that a definition arrived. So9// the watch of an absent collection is a quiet stream that sends no10// event until the recheck is due, and then a 410 Gone. The reflector11// then reads the collection again, which finds it once it exists. The12// 410 makes that read a list, which answers a resourceVersion to watch13// from. Without the quiet stream, the reflector would back off, list14// again, and log a failure about every 30 seconds for as long as the15// definition is missing.16//17// The quiet stream asks the API server nothing. The API server answers18// every list with a resourceVersion, so a plain watch from the empty19// version follows only the empty list that stands in for an absent20// collection. A watch from no version starts at the present, and the21// API server sends it no object deleted before it opened. If the22// definition arrived between the list and the watch, such a watch would23// be accepted, and the reflector would resume it from no version, so an24// object deleted while no watch was open would stay in the copy.2526import (27 "context"28 "net/http"29 "time"3031 metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"32 "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"33 "k8s.io/apimachinery/pkg/runtime"34 "k8s.io/apimachinery/pkg/watch"35)3637// defaultAbsentRecheck is the wait of Options.AbsentRecheck when it is38// zero. A definition arrives when a person installs an operator, so a39// wait of minutes costs one request and delays nothing a person waits40// on.41const defaultAbsentRecheck = 5 * time.Minute4243// absence is how one watch treats an absent collection. The zero value44// treats every refusal as a failure.45type absence struct {46 refused func(error) bool47 recheck time.Duration48}4950// absence reads the options once, before the reflector starts, so the51// reflector's goroutines read no shared setting.52func (o Options) absence() absence {53 recheck := o.AbsentRecheck54 if recheck == 0 {55 recheck = defaultAbsentRecheck56 }57 return absence{refused: o.Absent, recheck: recheck}58}5960func (a absence) is(err error) bool {61 return err != nil && a.refused != nil && a.refused(err)62}6364// list answers an empty collection for a list the API server refused65// because the collection is absent.66func (a absence) list(items runtime.Object, err error) (runtime.Object, error) {67 if a.is(err) {68 return &unstructured.UnstructuredList{}, nil69 }70 return items, err71}7273// watch answers the quiet stream in place of the watch that follows74// the empty list of an absent collection: a plain watch from no75// version. The reflector's first read is a streaming list, which is a76// watch that asks for the initial events. It goes to the API server,77// and its refusal makes the reflector fall back to the list above.78func (a absence) watch(ctx context.Context, request metav1.ListOptions) (watch.Interface, bool) {79 if a.refused == nil || request.ResourceVersion != "" || request.SendInitialEvents != nil {80 return nil, false81 }82 return quietWatch(ctx, a.recheck), true83}8485// quietWatch sends no event until the recheck is due, and then a 41086// Gone. It ends when the context ends.87func quietWatch(ctx context.Context, recheck time.Duration) watch.Interface {88 stream := watch.NewRaceFreeFake()89 go func() {90 timer := time.NewTimer(recheck)91 defer timer.Stop()92 select {93 case <-ctx.Done():94 stream.Stop()95 case <-timer.C:96 // The reflector stops the stream once it reads the error. A97 // stream it stopped already takes no event.98 stream.Error(&metav1.Status{99 Status: metav1.StatusFailure, Code: http.StatusGone, Reason: metav1.StatusReasonExpired,100 Message: "the collection was absent; read it again",101 })102 }103 }()104 return stream105}
informer/cache.go 100.0%
1package informer23// A pass reads the objects a watch holds from the watch's store, not4// from the API server, so a settled pass sends the API server no read5// of a watched kind.6//7// A store answers only while it is ready: it holds the whole first8// read, and the API server accepted a watch and forbade none since, so9// a watch keeps the store current. A store that holds part of the first read would10// leave objects out of a list, and a store whose watch the API server11// forbids holds no change made since its last list. A store whose12// watch fails while the API server is down stays ready (Synced in13// informer.go). While a store is not ready, every read goes to the API14// server. An object that a ready store does not hold is read from the15// API server too, because it can be an object the operator created a16// moment ago.17//18// A store's copy can be older than the operator's own last write. The19// memo package records the version of each copy the API server20// answered, and a store's copy at another version is read from the API21// server once. When the watch delivers the write, the store answers22// again. A list from the API server while the store is not ready notes23// the version of each object too (List).24//25// The memo package holds the requests whose answers the memo notes,26// because a program that must not link client-go sends them too. This27// package names the read and the status write as well, so a pass that28// reads a store and writes needs one import.2930import (31 "errors"32 "slices"33 "strings"3435 "k8s.io/client-go/tools/cache"3637 "github.com/liken-sh/liken/kubernetes/apiclient"38 "github.com/liken-sh/liken/kubernetes/memo"39)4041// View is one watch's store, and whether the store holds the whole42// first read. The zero View holds nothing, and every read goes to the43// API server.44type View struct {45 Store cache.Store46 Synced func() bool4748 // Whole says the watch takes every object of the kind, with no49 // selector. A list from such a store also holds each object the50 // operator created or wrote that the store does not hold yet. Leave51 // it false for a store with a selector, because an object the52 // operator wrote can be outside the selection.53 Whole bool54}5556// Ready reports whether a list can come from the store.57func (v View) Ready() bool {58 return v.Store != nil && v.Synced != nil && v.Synced()59}6061// Held is one kind's store, and the memo of the copies the operator62// wrote or read.63type Held struct {64 View View65 Versions *memo.Versions66}6768// Meta is the part of an object's metadata that the cache reads.69type Meta = memo.Meta7071// Object is a pointer to an operator's struct for one kind, which72// answers the object's metadata.73type Object[T any] = memo.Object[T]7475// Key is the key a store holds an object under: namespace/name for a76// namespaced object and the name for a cluster-scoped one.77func Key(meta Meta) string { return memo.Key(meta) }7879// Cached answers the store's copy of one object, by its key, while the80// store is ready. A copy that does not convert is reported and not81// answered. The caller reads the API server when Cached answers82// nothing.83func Cached[T any](view View, key string) (*T, bool) {84 if !view.Ready() {85 return nil, false86 }87 object, held, err := view.Store.GetByKey(key)88 if err != nil || !held {89 return nil, false90 }91 item, err := Convert[T](object)92 if err != nil {93 Report("the cached "+key, err)94 return nil, false95 }96 return &item, true97}9899// CachedList answers every copy in the store, in the order of their100// keys, which is the order of a list from the API server. A copy that101// does not convert is reported and left out, the same as a watch event102// that does not convert.103func CachedList[T any](view View) []T {104 objects := view.Store.List()105 slices.SortFunc(objects, func(a, b any) int { return strings.Compare(objectKey(a), objectKey(b)) })106 var items []T107 for _, object := range objects {108 item, err := Convert[T](object)109 if err != nil {110 Report("the cached objects", err)111 continue112 }113 items = append(items, item)114 }115 return items116}117118// objectKey is the key a store holds an object under.119func objectKey(object any) string {120 key, _ := cache.MetaNamespaceKeyFunc(object)121 return key122}123124// ReadOne answers one object: the store's copy when the store is ready125// and holds a current copy, and the API server's copy otherwise.126func ReadOne[T any, P Object[T]](c *apiclient.Client, held Held, key, path string) (*T, error) {127 if copied, ok := Cached[T](held.View, key); ok && held.Versions.Current(key, P(copied).GetObjectMeta().GetResourceVersion()) {128 return copied, nil129 }130 return ReadFresh[T, P](c, held.Versions, key, path)131}132133// ReadFresh reads one object from the API server and notes its version.134func ReadFresh[T any, P Object[T]](c *apiclient.Client, versions *memo.Versions, key, path string) (*T, error) {135 return memo.ReadFresh[T, P](c, versions, key, path)136}137138// CurrentList answers the store's copies, in the order of their keys. A139// copy that is older than the operator's own last write, or that does140// not convert, is replaced by the API server's copy, and an object the141// API server no longer holds is left out. A store that holds the whole142// collection also answers each object the memo noted and the store does143// not hold yet, such as one the operator created a moment ago, so a144// pass does not create it again. path names the object of each key.145func CurrentList[T any, P Object[T]](c *apiclient.Client, held Held, path func(key string) string) ([]T, error) {146 // The informer can take a created object between the two reads, so147 // the memo's keys are read first. In the other order, the store's148 // keys miss the object, the memo skips it because the store now149 // holds it, and the list leaves it out, so a pass creates it again.150 // In this order the key is in one read or in both.151 var keys []string152 if held.View.Whole {153 keys = held.Versions.Unheld(held.View.Store)154 }155 keys = append(keys, held.View.Store.ListKeys()...)156 slices.Sort(keys)157 keys = slices.Compact(keys)158 current := make([]T, 0, len(keys))159 for _, key := range keys {160 if copied, ok := Cached[T](held.View, key); ok && held.Versions.Current(key, P(copied).GetObjectMeta().GetResourceVersion()) {161 current = append(current, *copied)162 continue163 }164 fresh, err := ReadFresh[T, P](c, held.Versions, key, path(key))165 if errors.Is(err, apiclient.ErrNotFound) {166 continue167 }168 if err != nil {169 return nil, err170 }171 current = append(current, *fresh)172 }173 return current, nil174}175176// List answers every object of one kind: the store's copies through177// CurrentList while the store is ready, and a list from the API server178// at listPath while it is not. The list notes the version of each object179// it answers (memo.Versions.SendList), so a later list from a store that180// became ready at an older version reads that object from the API server181// again. listPath carries the selector of the watch, so both answers182// hold the same objects. It carries no resourceVersion, so etcd answers183// the list. path names the object of each key.184//185// A pass that lists a kind it reads from a store calls List, and never186// lists the collection from the API server itself. A list outside List187// notes no version, and the next pass can act on an older copy.188func List[T any, P Object[T]](c *apiclient.Client, held Held, listPath string, path func(key string) string) ([]T, error) {189 if held.View.Ready() {190 return CurrentList[T, P](c, held, path)191 }192 return memo.ReadList[T, P](c, held.Versions, listPath)193}194195// SettleStatus writes the status that apply sets on a copy of an196// object, and settles on the API server's copy after a 409 or a 404197// (memo.SettleStatus).198func SettleStatus[T any, P Object[T]](c *apiclient.Client, versions *memo.Versions, path string, held *T, apply func(*T) bool) (bool, error) {199 return memo.SettleStatus[T, P](c, versions, path, held, apply)200}
informer/convert.go 100.0%
1package informer23import (4 "encoding/json"5 "fmt"6 "os"78 "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"9 "k8s.io/apimachinery/pkg/runtime"10 "k8s.io/client-go/tools/cache"11)1213// Convert decodes one object from a watch into the operator's own14// struct. The informer hands a handler an *unstructured.Unstructured,15// or, for an object that was deleted while the watch was down, a16// tombstone that holds the last copy the informer had. A tombstone can17// hold no copy at all, and then the error names the tombstone's key.18//19// An object that does not convert has a field whose type differs from20// the operator's struct, so the CRD schema and the struct disagree.21// The error names the object, and the caller reports it, because an22// object that is dropped with no word leaves nobody a way to find out23// why the operator ignored an edit.24func Convert[T any](object any) (T, error) {25 var out T26 if tombstone, ok := object.(cache.DeletedFinalStateUnknown); ok {27 if tombstone.Obj == nil {28 return out, fmt.Errorf("the tombstone for %s holds no copy of the object", tombstone.Key)29 }30 object = tombstone.Obj31 }32 item, ok := object.(*unstructured.Unstructured)33 if !ok {34 return out, fmt.Errorf("the watch delivered a %T, not an object", object)35 }36 err := runtime.DefaultUnstructuredConverter.FromUnstructured(item.Object, &out)37 if err == nil {38 return out, nil39 }40 // The converter refuses a number for a field whose type is a string41 // with its own UnmarshalJSON, and does not call that method. A42 // Service's targetPort is a number or a name, and an operator holds43 // it as text that its UnmarshalJSON reads from either, so the44 // converter refuses a Service that a read from the API server45 // decodes. An object the converter refuses is decoded again from its46 // JSON, the same way a read from the API server decodes it.47 if body, marshalErr := json.Marshal(item.Object); marshalErr == nil {48 var decoded T49 if json.Unmarshal(body, &decoded) == nil {50 return decoded, nil51 }52 }53 return out, fmt.Errorf("%s %s does not convert: %w", item.GetKind(), objectName(item), err)54}5556// objectName is namespace/name for a namespaced object and name for a57// cluster-scoped one.58func objectName(item *unstructured.Unstructured) string {59 if item.GetNamespace() == "" {60 return item.GetName()61 }62 return item.GetNamespace() + "/" + item.GetName()63}6465// Report logs an object that Convert refused. what names the watch or66// the read, such as "the PairingRequests".67func Report(what string, err error) {68 fmt.Fprintf(os.Stderr, "watching %s: %v\n", what, err)69}
informer/informer.go 96.2%
1// Package informer keeps an operator's copy of a Kubernetes collection2// current, and answers a pass's reads from that copy.3//4// The API server sends each change to a collection as it happens, on a5// watch, and a watch with no change to send costs nothing. client-go's6// reflector runs each watch. It reads the whole collection first, as a7// streaming list or as a plain list, and then watches from the version8// that read returned, so it receives every change made after the read.9// It resumes a watch that the API server closed from the last version10// it delivered, reads the collection again after a 410 Gone, and backs11// off while the API server fails. Upstream maintains and tests that12// loop, so no operator keeps one of its own.13//14// The package imports only three parts of client-go: the reflector and15// the informer in tools/cache, the dynamic client that lists and16// watches any kind with no generated code, and rest for the in-cluster17// configuration. The typed clientset and the informer factories link a18// client and an informer for every built-in kind, which doubles the19// size of a binary. The apiclient package sends every write, and every20// read that the copy does not answer (cache.go).21package informer2223import (24 "context"25 "sync"26 "sync/atomic"27 "time"2829 apierrors "k8s.io/apimachinery/pkg/api/errors"30 metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"31 "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"32 "k8s.io/apimachinery/pkg/runtime"33 "k8s.io/apimachinery/pkg/runtime/schema"34 "k8s.io/apimachinery/pkg/watch"35 "k8s.io/client-go/dynamic"36 "k8s.io/client-go/rest"37 "k8s.io/client-go/tools/cache"38)3940// InCluster builds the dynamic client for the watches from the pod's41// ServiceAccount, the same credentials apiclient.InCluster reads.42func InCluster() (dynamic.Interface, error) {43 return InClusterAt("")44}4546// InClusterAt is InCluster with the API server's address in place of47// the one the environment names, for the reason that48// apiclient.InClusterOptions.Server gives. An empty server means the49// environment's address.50func InClusterAt(server string) (dynamic.Interface, error) {51 config, err := rest.InClusterConfig()52 if err != nil {53 return nil, err54 }55 if server != "" {56 config.Host = server57 }58 return dynamic.NewForConfig(config)59}6061// Source names one collection and the part of it to watch. A watch62// covers only the objects the operator reads, so the API server filters63// the rest and never sends them. An empty Namespace watches a64// cluster-scoped kind, or every namespace of a namespaced kind.65type Source struct {66 Resource schema.GroupVersionResource67 Namespace string68 LabelSelector string69 FieldSelector string70}7172// String names the source in a log line.73func (s Source) String() string {74 name := s.Resource.Resource75 if s.Namespace != "" {76 name = s.Namespace + "/" + name77 }78 if s.FieldSelector != "" {79 name += " (" + s.FieldSelector + ")"80 }81 if s.LabelSelector != "" {82 name += " [" + s.LabelSelector + "]"83 }84 return name85}8687// Options are the optional parts of a watch.88type Options struct {89 // Handler receives each change after the copy holds it. A nil90 // Handler receives nothing, for a watch that only keeps a copy for91 // the pass to read.92 Handler cache.ResourceEventHandler9394 // Synced runs once, after the first read of the collection is in95 // the copy and the handler has taken each object of it. A pass that96 // started before then read the API server, and an edit made between97 // that read and the watch's first read is in no event. A wake here98 // makes the pass read again.99 Synced func()100101 // Reopened runs each time the API server accepts a watch after the102 // first one it accepted. The API server ends a watch on its own103 // schedule, so a low rate is normal, and a high rate says the stream104 // breaks faster than the watch can use it. A refused watch opened105 // nothing, so it does not count, and neither does the quiet stream106 // of an absent collection.107 Reopened func()108109 // Recovered runs when the API server accepts a watch after the last110 // attempt to watch the collection failed, such as after an outage of111 // the API server. A routine reopen follows no failure and does not112 // count. The copy answers again from that moment, before the113 // reflector delivers the events missed during the outage, so a114 // caller that wakes on it reads the collection from the API server115 // once, not from the copy.116 Recovered func()117118 // Absent, when it is not nil, names the refusals that mean the119 // collection is not there to watch, such as the 404 of a kind whose120 // definition another operator installs. The copy then holds an empty121 // collection, and the watch finds the collection when it arrives122 // (absent.go).123 Absent func(error) bool124125 // AbsentRecheck is how long the watch of an absent collection waits126 // before it asks the API server again. Zero means five minutes.127 AbsentRecheck time.Duration128129 // ListFailed, when it is not nil, runs each time a list of the whole130 // collection fails, with the error. An operator that waits for its131 // first read uses it to stop waiting for a read that fails, such as132 // the list of a kind the cluster does not serve.133 ListFailed func(error)134135 // Transform, when it is not nil, trims each object before the copy136 // holds it, after the managedFields are gone. A field that no pass137 // reads costs memory for every object of the kind, such as the list138 // of images a Node's runtime holds.139 Transform func(*unstructured.Unstructured)140141 // UnreadyOnWatchError, when it is true, makes the copy stop142 // answering after any failed watch, not only after a 401 or a 403,143 // until the API server accepts a watch again. While the API server144 // restarts, the reflector backs off for up to a minute, and a write145 // that lands on the API server in that time is in no copy. An146 // operator whose pass must never act on such a copy sets it, and147 // its pass reads the API server instead. liken's operators set it:148 // a machine that acted on a copy older than a withdrawn rollout149 // could stage the old target and reboot into it.150 UnreadyOnWatchError bool151152 // Indexers name the indexes the copy keeps, for a pass that reads153 // the objects of one label value without a scan of the whole copy.154 // The store of a watch with indexers is a cache.Indexer.155 Indexers cache.Indexers156}157158// Collection is the copy of one watched collection.159//160// A nil *Collection is valid and never syncs. Its View answers nothing,161// and each read goes to the API server.162type Collection struct {163 store cache.Store164 controller cache.Controller165 done chan struct{}166167 // watching is true after the API server accepts a watch, and false168 // again after it forbids one, or after any failed watch when strict169 // is true. Synced says why it matters.170 watching atomic.Bool171172 // strict is Options.UnreadyOnWatchError.173 strict bool174}175176// noteWatch records whether the API server grants the watch. An177// accepted watch grants it, and a 401 or a 403 refuses it. Any other178// failure, such as a refused connection while the API server restarts,179// or a 5xx, says nothing about the permission, so the flag keeps its180// last value, unless the copy is strict (Options.UnreadyOnWatchError).181// The dynamic client answers each refusal as a *errors.StatusError,182// which carries the HTTP status.183func (c *Collection) noteWatch(err error) {184 switch {185 case err == nil:186 c.watching.Store(true)187 case c.strict, apierrors.IsForbidden(err), apierrors.IsUnauthorized(err):188 c.watching.Store(false)189 }190}191192// Start opens the watch and returns at once. The watch runs until the193// context ends. Done closes after the informer has stopped and the last194// handler call and the Synced call have returned, so a caller that195// closes a channel after Done never sends on a closed channel.196func Start(ctx context.Context, client dynamic.Interface, source Source, options Options) *Collection {197 var collection dynamic.ResourceInterface = client.Resource(source.Resource)198 if source.Namespace != "" {199 collection = client.Resource(source.Resource).Namespace(source.Namespace)200 }201 c := &Collection{done: make(chan struct{}), strict: options.UnreadyOnWatchError}202 scope := func(list *metav1.ListOptions) {203 list.LabelSelector = source.LabelSelector204 list.FieldSelector = source.FieldSelector205 }206 var opened atomic.Int64207 // failed records that the last attempt to watch failed, so the next208 // accepted watch is a recovery (Options.Recovered).209 var failed atomic.Bool210 accepted := func() {211 if failed.Swap(false) && options.Recovered != nil {212 options.Recovered()213 }214 }215 absent := options.absence()216 lister := &cache.ListWatch{217 ListWithContextFunc: func(ctx context.Context, list metav1.ListOptions) (runtime.Object, error) {218 scope(&list)219 items, err := absent.list(collection.List(ctx, list))220 if err != nil && options.ListFailed != nil {221 options.ListFailed(err)222 }223 return items, err224 },225 WatchFuncWithContext: func(ctx context.Context, list metav1.ListOptions) (watch.Interface, error) {226 scope(&list)227 if quiet, ok := absent.watch(ctx, list); ok {228 // The quiet stream stands in for an accepted watch of an229 // empty collection, so the empty copy answers a pass. It230 // is no watch the API server accepted, so Reopened does231 // not count it.232 c.noteWatch(nil)233 accepted()234 return quiet, nil235 }236 stream, err := collection.Watch(ctx, list)237 c.noteWatch(err)238 if err != nil {239 failed.Store(true)240 } else {241 accepted()242 }243 // The first accepted watch is the streaming list of the first244 // read, or the watch after a plain list.245 if err == nil && opened.Add(1) > 1 && options.Reopened != nil {246 options.Reopened()247 }248 return stream, err249 },250 }251 handler := options.Handler252 if handler == nil {253 handler = cache.ResourceEventHandlerFuncs{}254 }255 c.store, c.controller = cache.NewInformerWithOptions(cache.InformerOptions{256 ListerWatcher: lister,257 ObjectType: &unstructured.Unstructured{},258 Handler: handler,259 Transform: trim(options.Transform),260 Indexers: options.Indexers,261 })262 go func() {263 defer close(c.done)264 var group sync.WaitGroup265 if options.Synced != nil {266 group.Go(func() {267 select {268 case <-c.controller.HasSyncedChecker().Done():269 options.Synced()270 case <-ctx.Done():271 }272 })273 }274 c.controller.RunWithContext(ctx)275 group.Wait()276 }()277 return c278}279280// Done closes when the watch has stopped.281func (c *Collection) Done() <-chan struct{} { return c.done }282283// Synced answers whether the copy holds the whole collection and the284// watch keeps it current. Until the first read is done, a read of the285// copy could miss an object that exists.286//287// A copy whose watch the API server forbids is not current, even after288// a list. The case is a release skew: a new binary under the previous289// release's RBAC, which grants list and not watch. The reflector then290// lists the collection again after each backoff, which grows to between291// thirty and sixty seconds, and in between the copy holds no change at292// all. So after a 401293// or a 403 on a watch, the copy does not answer and the pass reads the294// API server, until the API server accepts a watch again.295//296// A watch that fails for another reason does not stop the copy, unless297// the copy is strict (Options.UnreadyOnWatchError). While the API298// server is down, a read of the API server fails too, and the copy lets299// an operator keep its local work going. The reflector resumes the300// watch from the copy's last version when the API server returns, or301// lists again, so the copy then receives each change it missed.302func (c *Collection) Synced() bool {303 return c != nil && c.controller.HasSynced() && c.watching.Load()304}305306// View answers the copy for a pass to read (cache.go).307func (c *Collection) View() View {308 if c == nil {309 return View{}310 }311 return View{Store: c.store, Synced: c.Synced}312}313314// trim answers the transform that removes metadata.managedFields from315// each object before the informer stores it, and then runs the316// operator's own. The field records which client set each field of the317// object. No operator reads it, and without the transform the copy318// holds it for every object.319func trim(operators func(*unstructured.Unstructured)) cache.TransformFunc {320 return func(object any) (any, error) {321 if item, ok := object.(*unstructured.Unstructured); ok {322 item.SetManagedFields(nil)323 if operators != nil {324 operators(item)325 }326 }327 return object, nil328 }329}
informer/one.go 100.0%
1package informer23// The watch of one object by name, such as a CA ConfigMap or a serving4// Secret. An API pod follows each of them and hands every version of5// the object to its owner, so a rotated certificate reaches the next6// TLS handshake when the API server writes it.7//8// The watch opens the collection with the field selector9// metadata.name=<name>. The API server authorizes that request against10// a Role's resourceNames, because it reads the name from the selector,11// so a Role that names one object grants its list and its watch. A GET12// on the object's own path with watch=true is not a watch: the API13// server answers it as a plain get and closes the stream.1415import (16 "context"17 "sync"1819 "k8s.io/apimachinery/pkg/runtime/schema"20 "k8s.io/client-go/dynamic"21 "k8s.io/client-go/tools/cache"22)2324// One names the object a watch follows.25type One struct {26 Resource schema.GroupVersionResource27 Namespace string28 Name string2930 // What names the object in the line that reports a copy that does31 // not convert, such as "the Secret display-api-tls".32 What string33}3435// WatchOne calls seen with the object each time it changes, and with36// nil when the object does not exist. It returns at once, and the watch37// runs until the context ends. synced, when it is not nil, runs once38// after seen has taken the first read, whether that read found the39// object or not. Done on the answer closes after the last call to seen40// and to synced has returned.41//42// A copy that does not convert is reported, and seen does not run, so43// the owner keeps what it holds: the next version of the object can be44// valid. An owner that reads the object itself takes T as45// unstructured.Unstructured, which every copy converts to.46func WatchOne[T any](ctx context.Context, client dynamic.Interface, one One, seen func(held *T), synced func()) *Collection {47 report := &oneReport[T]{what: one.What, seen: seen}48 // Synced can run before Start returns the collection, so it takes49 // the store from a channel that fills after Start returns.50 started := make(chan cache.Store, 1)51 source := Source{Resource: one.Resource, Namespace: one.Namespace, FieldSelector: "metadata.name=" + one.Name}52 collection := Start(ctx, client, source, Options{53 Handler: cache.ResourceEventHandlerFuncs{54 AddFunc: report.take,55 UpdateFunc: func(_, object any) { report.take(object) },56 DeleteFunc: func(any) { report.gone() },57 },58 Synced: func() {59 store := <-started60 report.firstRead(func() bool { return len(store.ListKeys()) == 0 })61 if synced != nil {62 synced()63 }64 },65 })66 started <- collection.store67 return collection68}6970// oneReport turns the informer's calls into the owner's reports.71//72// An object absent from the first read has no event, so the watch73// reports it as nil once that read is done. The informer calls the74// handler on one goroutine and the end of the first read on another,75// and the lock runs one report at a time. The report of absence waits76// for both conditions:77//78// - No event reached the handler. An object that does not convert79// still exists, and an object that the read found and a later event80// deleted was reported already, so neither gets a report of its81// own.82// - The store is empty. The informer adds an object to its store83// before it calls the handler, so an object that arrived just after84// the read is in the store before its event reaches the owner. A85// report of absence would say the object does not exist while it86// does.87//88// An object that arrives after an empty first read is reported last in89// either order: its event runs after the report of absence, or the90// report finds the object in the store and reports nothing.91type oneReport[T any] struct {92 what string93 seen func(held *T)9495 mu sync.Mutex96 arrived bool97}9899func (r *oneReport[T]) take(object any) {100 r.mu.Lock()101 defer r.mu.Unlock()102 r.arrived = true103 held, err := Convert[T](object)104 if err != nil {105 Report(r.what, err)106 return107 }108 r.seen(&held)109}110111func (r *oneReport[T]) gone() {112 r.mu.Lock()113 defer r.mu.Unlock()114 r.arrived = true115 r.seen(nil)116}117118// firstRead reports absence at the end of the first read. empty reads119// the store under the lock, so no event runs between that read and the120// report.121func (r *oneReport[T]) firstRead(empty func() bool) {122 r.mu.Lock()123 defer r.mu.Unlock()124 if !r.arrived && empty() {125 r.seen(nil)126 }127}
memo/memo.go 100.0%
1// Package memo records the resourceVersion of the newest copy of each2// object that an operator wrote or read from the API server.3//4// An operator reads the objects that a watch holds from the watch's5// store, not from the API server. A store's copy can be older than the6// operator's own last write, because the watch delivers the write a7// moment after the API server answers it, and later still while the8// watch is down. A pass that acts on that copy acts again on a change it9// already made: it opens a pairing window again, sends a receiver a10// setting again, or skips a status write that the older copy hides. So11// the operator records the version of each copy the API server answered,12// and reads an object from the API server when the store's copy has13// another version. When the watch delivers the write, the versions14// match and the store answers again.15//16// The package imports nothing from k8s.io, so a program that must not17// link client-go can record its writes too.18package memo1920import "sync"2122// Versions records, for each object by its store key, the23// resourceVersion of the newest copy the operator wrote or read from24// the API server. It compares versions only for equality, because the25// API server gives them no order. A nil *Versions records nothing, and26// every copy in a store is current to it.27type Versions struct {28 mu sync.Mutex29 seen map[string]string3031 // notes counts every note, and notedAt holds the count at each32 // key's last note. A list compares the two to find the keys that a33 // request noted while the list was in flight (SendList).34 notes uint6435 notedAt map[string]uint643637 // requests holds, for each object, one request at a time with the38 // note of its answer, so the notes follow the order in which the API39 // server answered. Two goroutines that write one object could40 // otherwise note the older answer last, and a store's copy at that41 // older version would then count as current. Requests about other42 // objects do not wait.43 requests map[string]*sync.Mutex44}4546// New answers an empty memo.47func New() *Versions {48 return &Versions{seen: map[string]string{}, notedAt: map[string]uint64{}, requests: map[string]*sync.Mutex{}}49}5051// requestsOf answers the lock of one object's requests.52func (m *Versions) requestsOf(key string) *sync.Mutex {53 m.mu.Lock()54 defer m.mu.Unlock()55 held, ok := m.requests[key]56 if !ok {57 held = &sync.Mutex{}58 m.requests[key] = held59 }60 return held61}6263// Current reports whether a store's copy at this version is at least as64// new as every copy the operator wrote or read. An object the memo has65// not noted is current at any version.66func (m *Versions) Current(key, version string) bool {67 if m == nil {68 return true69 }70 m.mu.Lock()71 defer m.mu.Unlock()72 seen, noted := m.seen[key]73 return !noted || seen == version74}7576// Note records the version of a copy the API server answered. An empty77// version records an object the API server no longer holds, or one78// that another writer changed, and no copy in a store matches it.79func (m *Versions) Note(key, version string) {80 if m == nil {81 return82 }83 m.mu.Lock()84 defer m.mu.Unlock()85 m.note(key, version)86}8788// note records one version. The caller holds mu.89func (m *Versions) note(key, version string) {90 m.notes++91 m.seen[key] = version92 m.notedAt[key] = m.notes93}9495// KeyGetter reads one object from a store by its key. client-go's96// cache.Store satisfies it.97type KeyGetter interface {98 GetByKey(key string) (item any, exists bool, err error)99}100101// Unheld answers each key the memo noted at a version that the store102// does not hold. These are objects the operator created or wrote a103// moment ago, which the watch has not delivered yet.104func (m *Versions) Unheld(store KeyGetter) []string {105 if m == nil {106 return nil107 }108 m.mu.Lock()109 defer m.mu.Unlock()110 var keys []string111 for key, version := range m.seen {112 if _, held, _ := store.GetByKey(key); !held && version != "" {113 keys = append(keys, key)114 }115 }116 return keys117}118119// Send runs one request about one object, and notes the version of the120// copy the API server answered. A failed request notes the empty121// version: after a 404 or a 409 the operator holds no copy of the API122// server's, and a write whose answer was lost may have landed. The next123// read of the object then goes to the API server.124func (m *Versions) Send(key string, request func() (version string, err error)) error {125 if m == nil {126 _, err := request()127 return err128 }129 requests := m.requestsOf(key)130 requests.Lock()131 defer requests.Unlock()132 version, err := request()133 if err != nil {134 version = ""135 }136 m.Note(key, version)137 return err138}139140// SendList runs one list request, and notes the version of each object141// the API server answered, by its key. A pass that lists before a142// watch's store is ready must note these versions. The store can then143// become ready at an older version, because the reflector's first read144// can come from the API server's watch cache, which runs behind etcd.145// Without the notes, the next pass reads that older copy as current and146// acts on an object older than one it already read.147//148// An object that a request noted while the list was in flight keeps the149// request's version. The list can hold a copy from before the request's150// answer, and a note of that copy would make a store's copy at the older151// version current. When the request's version is the older one instead,152// the store's copy at the list's version is not current, and the next153// read of the object goes to the API server once. That costs one154// request and never answers an older copy.155//156// A list replaces the version of each object that the memo noted before157// the list was sent. That is correct only for a list that etcd answers,158// which holds every write the API server answered before the list. A159// list with resourceVersion=0 comes from the watch cache, which can be160// older than such a write, so the caller must not send one.161//162// A list that fails notes nothing, because it names no object.163func (m *Versions) SendList(request func() (versions map[string]string, err error)) error {164 if m == nil {165 _, err := request()166 return err167 }168 m.mu.Lock()169 start := m.notes170 m.mu.Unlock()171 listed, err := request()172 if err != nil {173 return err174 }175 m.mu.Lock()176 defer m.mu.Unlock()177 for key, version := range listed {178 if m.notedAt[key] <= start {179 m.note(key, version)180 }181 }182 return nil183}184185// ForgetGone drops the record of each object that a list did not186// answer and the store does not hold. The API server answered 404 for187// such an object, and its watch delivered the delete. An operator whose188// objects come and go, such as a Play for each film or a Job for each189// run, calls it after each list, or the memo would keep a record of190// every object for the life of the process. A key the store still holds191// keeps its record, so a copy of an object the API server already192// deleted is not current again before the watch removes it.193func (m *Versions) ForgetGone(store KeyGetter, listed map[string]bool) {194 if m == nil {195 return196 }197 m.mu.Lock()198 defer m.mu.Unlock()199 for key := range m.seen {200 if listed[key] {201 continue202 }203 if _, held, _ := store.GetByKey(key); held {204 continue205 }206 delete(m.seen, key)207 delete(m.notedAt, key)208 delete(m.requests, key)209 }210}211212// Forget drops the record of one object. An operator that reads an213// object by name from a store forgets it once the API server answers214// 404, so the next read of the name answers from the store and sends215// nothing.216//217// Forget and ForgetAt also drop the object's request lock (Send). A218// Send that holds the lock at that moment runs to its end and notes its219// answer, but a Send that starts after the forget takes a new lock, so220// the two can run at once and note their answers out of order. An221// operator that sends requests about one object from more than one222// goroutine forgets the object only while none of them runs.223func (m *Versions) Forget(key string) {224 if m == nil {225 return226 }227 m.mu.Lock()228 defer m.mu.Unlock()229 delete(m.seen, key)230 delete(m.notedAt, key)231 delete(m.requests, key)232}233234// Noted reports whether the memo holds a record of the key at any235// version, the empty one included. A store that holds a whole selection236// and not the key answers that no such object exists, unless the memo237// noted it: the operator created it a moment ago, or a request about it238// failed.239func (m *Versions) Noted(key string) bool {240 if m == nil {241 return false242 }243 m.mu.Lock()244 defer m.mu.Unlock()245 _, noted := m.seen[key]246 return noted247}248249// ForgetAt drops the record of one object when the record holds this250// version, and keeps any other record. A caller passes the version of251// the copy a store holds: once the store holds the version the memo252// noted, every later copy the store holds is newer, so the record253// guards nothing. Without it, each later write from another writer254// makes the store's copy differ from the record, and costs one read255// from the API server. The check and the delete are one step, so a256// write that notes a newer version in the meantime keeps its record.257func (m *Versions) ForgetAt(key, version string) {258 if m == nil || version == "" {259 return260 }261 m.mu.Lock()262 defer m.mu.Unlock()263 if seen, noted := m.seen[key]; noted && seen == version {264 delete(m.seen, key)265 delete(m.notedAt, key)266 delete(m.requests, key)267 }268}
memo/requests.go 94.0%
1package memo23// The requests whose answers the memo notes: a read of one object from4// the API server, a list of a collection, a write that answers the5// stored copy, and a status write that settles on the API server's6// copy. None of them reads a7// watch's store, so they import nothing from k8s.io, and a program that8// must not link client-go can link them. library-operator's pod build9// is such a program: it runs no pass, but it compiles the pass's code,10// which writes status through these functions, so it links them. The11// informer package names the read and the status write as well, for a12// pass that also reads a store.13//14// A write from a copy that another writer changed since carries an15// older resourceVersion, and the API server answers 409 Conflict.16// SettleStatus then reads the object from the API server and writes17// once more, if the fresh copy still needs the write. A copy of an18// object somebody deleted answers 404, and each caller handles that the19// way it handles an object that is absent.2021import (22 "strings"2324 "github.com/liken-sh/liken/kubernetes/apiclient"25)2627// Meta is the part of an object's metadata that the memo and a store28// read. It is an alias of an interface literal, not a type of its own,29// so a package can answer it without an import of this module. liken's30// machine package does that: init links it, and init must not link31// net/http, which this module's client brings.32type Meta = interface {33 GetNamespace() string34 GetName() string35 GetResourceVersion() string36}3738// Object is a pointer to an operator's struct for one kind, which39// answers the object's metadata.40type Object[T any] interface {41 *T42 GetObjectMeta() Meta43}4445// Key is the key a store holds an object under: namespace/name for a46// namespaced object and the name for a cluster-scoped one. The memo47// records each object under the same key.48func Key(meta Meta) string {49 if meta.GetNamespace() == "" {50 return meta.GetName()51 }52 return meta.GetNamespace() + "/" + meta.GetName()53}5455// NamespacedPath answers the API path of the object a key names, from56// the path of one object by its namespace and name. A list that reads57// an object again from the API server has only the key. A key with no58// namespace names a cluster-scoped object, and path takes an empty59// namespace for it.60func NamespacedPath(path func(namespace, name string) string) func(key string) string {61 return func(key string) string {62 namespace, name, namespaced := strings.Cut(key, "/")63 if !namespaced {64 return path("", key)65 }66 return path(namespace, name)67 }68}6970// ReadFresh reads one object from the API server and notes its version.71func ReadFresh[T any, P Object[T]](c *apiclient.Client, versions *Versions, key, path string) (*T, error) {72 return Written[T, P](versions, key, func() (*T, error) { return apiclient.Get[T](c, path) })73}7475// ReadList reads a collection from the API server, and notes the76// version of each object it answers (Versions.SendList). listPath is the77// path of the collection, with the selector of the watch whose store78// the list stands in for, and no resourceVersion (Versions.SendList).79func ReadList[T any, P Object[T]](c *apiclient.Client, versions *Versions, listPath string) ([]T, error) {80 var items []T81 err := versions.SendList(func() (map[string]string, error) {82 list, err := apiclient.Get[struct {83 Items []T `json:"items"`84 }](c, listPath)85 if err != nil {86 return nil, err87 }88 items = list.Items89 listed := make(map[string]string, len(items))90 for index := range items {91 meta := P(&items[index]).GetObjectMeta()92 listed[Key(meta)] = meta.GetResourceVersion()93 }94 return listed, nil95 })96 return items, err97}9899// Written sends one request that answers an object, such as a create,100// an apply, or an update, and notes the version of the copy the API101// server answered. A request that fails answers no copy, and notes that102// the operator holds no current copy.103func Written[T any, P Object[T]](versions *Versions, key string, request func() (*T, error)) (*T, error) {104 var answer *T105 err := versions.Send(key, func() (string, error) {106 copied, err := request()107 if err != nil {108 return "", err109 }110 answer = copied111 return P(copied).GetObjectMeta().GetResourceVersion(), nil112 })113 return answer, err114}115116// SettleStatus writes the status that apply sets on a copy of an117// object. apply composes the status from the copy it is given, sets it,118// and reports whether the copy needs the write. A write refused because119// the copy is older than the API server's, or because the object is120// gone, reads the object again, applies again to the fresh copy, and121// writes once more. SettleStatus reports whether a write landed, and an122// object that is gone answers apiclient.ErrNotFound. Each copy the API123// server answers is noted in versions. After an error, held holds the124// status that apply set, which the API server did not take.125func SettleStatus[T any, P Object[T]](c *apiclient.Client, versions *Versions, path string, held *T, apply func(*T) bool) (bool, error) {126 key := Key(P(held).GetObjectMeta())127 write := func() error {128 return versions.Send(key, func() (string, error) {129 if err := apiclient.ReplaceStatus(c, path, held); err != nil {130 return "", err131 }132 return P(held).GetObjectMeta().GetResourceVersion(), nil133 })134 }135 if !apply(held) {136 return false, nil137 }138 err := write()139 if !apiclient.Stale(err) {140 return err == nil, err141 }142 current, err := ReadFresh[T, P](c, versions, key, path)143 if err != nil {144 return false, err145 }146 *held = *current147 if !apply(held) {148 return false, nil149 }150 if err := write(); err != nil {151 return false, err152 }153 return true, nil154}