kubernetes

InputUnitCoveredTotalPercent
Gostatements96898398.5%

Go

968 of 983 statements, 98.5%.

FileCovered statementsTotal statementsPercent
apiclient/client.go11211498.2%
apiclient/objects.go111291.7%
apiservertest/errors.go22100.0%
apiservertest/server.go313296.9%
apiservertest/transport.go2929100.0%
conditions/conditions.go2121100.0%
election/election.go10110596.2%
election/electiontest/server.go727398.6%
election/lock.go2525100.0%
events/event.go1010100.0%
events/eventstest/events.go8787100.0%
events/recorder.go104104100.0%
events/series.go2020100.0%
events/transition.go88100.0%
informer/absent.go1919100.0%
informer/cache.go5151100.0%
informer/convert.go2020100.0%
informer/informer.go757896.2%
informer/one.go3333100.0%
memo/memo.go9090100.0%
memo/requests.go475094.0%
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}