From 1328676daffd3188d5b298f5f8f7331a86bbde03 Mon Sep 17 00:00:00 2001 From: Simon Schrottner Date: Fri, 11 Sep 2026 16:10:39 +0200 Subject: [PATCH 1/2] fix(launchpad): make /start wait until flags are actually served /start returned as soon as flagd's /readyz reported 200, and treated that as "the seeded flag state is being served". It is not. flagd's file sync calls sendDataSync(), which pushes the payload onto a channel buffered to the number of sources, and only then setReady(true); the parse and store swap that make those flags evaluable happen afterwards, on the goroutine that drains that channel. Every source can therefore hand over its payload and flip ready without a single one having been applied, so /start could return while flagd still answered FLAG_NOT_FOUND for flags the configuration plainly defines. Measured against the v3.10.1 image over 40 starts, 14 of them (35%) left such a window, median 6ms and up to 32ms. A provider that blocks during its own initialisation absorbs the window and never sees it; a stateless provider evaluates the instant /start returns and races it, which reads as a catastrophically broken provider rather than as a racing testbed. The Go OFREP conformance suite went from 22 and 29 failures over two runs, with near disjoint failing sets, to the same 2 failures twice, both of them the known fixture gap that open-feature/flagd-testbed#392 fills. After /readyz reports 200, poll a real evaluation over OFREP until it resolves, within the existing 10s budget. The probe keys are derived from the flagd configuration being started, one enabled flag per file source, rather than hardcoded: the configurations do not all share a flag file, so a fixed key would not survive them, and the sources are merged into the store independently, so each needs its own probe. Disabled flags are skipped because flagd reports FLAG_NOT_FOUND for them too. No sleep and no blanket retry. A sleep would hide the window from every other language's adoption and turn a deterministic contract into a flake, and a retry in the client would hide a genuine backend defect that the conformance suite exists to surface. Signed-off-by: Simon Schrottner --- launchpad/pkg/flagd.go | 177 +++++++++++++++++++++++++++++++++++++++-- 1 file changed, 171 insertions(+), 6 deletions(-) diff --git a/launchpad/pkg/flagd.go b/launchpad/pkg/flagd.go index 99c13e5..631f71a 100644 --- a/launchpad/pkg/flagd.go +++ b/launchpad/pkg/flagd.go @@ -1,16 +1,37 @@ package flagd import ( + "bytes" "context" + "encoding/json" "errors" "fmt" + "io" "net/http" + "net/url" "os" "os/exec" + "sort" + "strings" "sync" "time" ) +const ( + // startupBudget bounds the whole of /start: waiting for the readiness + // probe and waiting for the flags to actually be served. + startupBudget = 10 * time.Second + + readyPollInterval = 100 * time.Millisecond + servedPollInterval = 10 * time.Millisecond + + readyzURL = "http://localhost:8014/readyz" + // flagd serves OFREP over plain HTTP on 8016 for every configuration we + // ship, including the one that configures server certificates (those apply + // to the flag evaluation port only), so one URL works for all of them. + ofrepEvaluateURL = "http://localhost:8016/ofrep/v1/evaluate/flags/" +) + var ( flagdCmd *exec.Cmd flagdLock sync.Mutex @@ -102,23 +123,45 @@ func StartFlagd(config string) error { } flagdLock.Unlock() - // Poll health endpoint until ready client := &http.Client{Timeout: 500 * time.Millisecond} - ticker := time.NewTicker(100 * time.Millisecond) + deadline := time.Now().Add(startupBudget) + + if err := awaitReadyz(client, deadline); err != nil { + _ = StopFlagd() + return err + } + + // /readyz reports ready as soon as every sync source has handed its payload + // over; flagd parses it and swaps its flag store afterwards, on a separate + // goroutine. In that window flagd answers FLAG_NOT_FOUND for flags the + // configuration plainly defines. A provider that blocks during its own + // initialisation absorbs the window, but a stateless one evaluates the + // instant /start returns and races it, so wait for a real evaluation. + if err := awaitFlagsServed(client, configPath, deadline); err != nil { + _ = StopFlagd() + return err + } + + fmt.Println("flagd started successfully.") + return nil +} + +// awaitReadyz waits for flagd's readiness probe to report that every sync +// source has completed at least one successful data sync. +func awaitReadyz(client *http.Client, deadline time.Time) error { + ticker := time.NewTicker(readyPollInterval) defer ticker.Stop() - timeout := time.After(10 * time.Second) + timeout := time.After(time.Until(deadline)) for { select { case <-timeout: - _ = StopFlagd() return fmt.Errorf("flagd health check timed out") case <-ticker.C: - resp, err := client.Get("http://localhost:8014/readyz") + resp, err := client.Get(readyzURL) if err == nil { resp.Body.Close() if resp.StatusCode == http.StatusOK { - fmt.Println("flagd started successfully.") return nil } } @@ -126,6 +169,128 @@ func StartFlagd(config string) error { } } +// awaitFlagsServed waits until flagd actually resolves a flag from every file +// source the configuration lists, so that a successful /start is a promise +// that the next evaluation resolves against the new baseline. +func awaitFlagsServed(client *http.Client, configPath string, deadline time.Time) error { + keys, err := probeKeys(configPath) + if err != nil { + // Without a probe key there is nothing to verify the store with. Fall + // back to the readiness probe alone rather than failing a start that + // would otherwise have worked. + fmt.Printf("Cannot verify that flags are served for %s, falling back to the readiness probe: %v\n", configPath, err) + return nil + } + + for _, key := range keys { + if err := awaitFlagServed(client, key, deadline); err != nil { + return err + } + } + return nil +} + +func awaitFlagServed(client *http.Client, key string, deadline time.Time) error { + for { + if flagIsServed(client, key) { + return nil + } + if time.Now().After(deadline) { + return fmt.Errorf("flagd did not serve flag %q before the startup budget expired", key) + } + time.Sleep(servedPollInterval) + } +} + +// flagIsServed reports whether flagd holds the given flag in its store. flagd +// answers FLAG_NOT_FOUND both while the store is still empty and for a key it +// genuinely does not hold; every other answer means the flag is being served. +func flagIsServed(client *http.Client, key string) bool { + resp, err := client.Post(ofrepEvaluateURL+url.PathEscape(key), "application/json", strings.NewReader("{}")) + if err != nil { + return false + } + defer resp.Body.Close() + + body, err := io.ReadAll(resp.Body) + if err != nil { + return false + } + return !bytes.Contains(body, []byte("FLAG_NOT_FOUND")) +} + +// probeKeys returns one enabled flag key per file source of the given flagd +// configuration. Deriving the keys from the configuration rather than hard +// coding one keeps the check working for every configuration we ship, and +// covers each source separately because they are merged into the store +// independently. +func probeKeys(configPath string) ([]string, error) { + content, err := os.ReadFile(configPath) + if err != nil { + return nil, fmt.Errorf("failed to read flagd config: %w", err) + } + + var cfg struct { + Sources []struct { + URI string `json:"uri"` + Provider string `json:"provider"` + } `json:"sources"` + } + if err := json.Unmarshal(content, &cfg); err != nil { + return nil, fmt.Errorf("failed to parse flagd config: %w", err) + } + + var keys []string + for _, source := range cfg.Sources { + if source.Provider != "file" { + continue + } + key, err := enabledFlagKey(source.URI) + if err != nil { + return nil, err + } + keys = append(keys, key) + } + + if len(keys) == 0 { + return nil, fmt.Errorf("no file source with an enabled flag found") + } + return keys, nil +} + +// enabledFlagKey picks a stable, enabled flag key from a flag definition file. +// Disabled flags are skipped because flagd reports FLAG_NOT_FOUND for them, +// which would be indistinguishable from the store not being populated yet. +func enabledFlagKey(path string) (string, error) { + content, err := os.ReadFile(path) + if err != nil { + return "", fmt.Errorf("failed to read flag source %s: %w", path, err) + } + + var definition struct { + Flags map[string]struct { + State string `json:"state"` + } `json:"flags"` + } + if err := json.Unmarshal(content, &definition); err != nil { + return "", fmt.Errorf("failed to parse flag source %s: %w", path, err) + } + + candidates := make([]string, 0, len(definition.Flags)) + for key, flag := range definition.Flags { + if flag.State == "DISABLED" { + continue + } + candidates = append(candidates, key) + } + if len(candidates) == 0 { + return "", fmt.Errorf("flag source %s defines no enabled flag", path) + } + + sort.Strings(candidates) + return candidates[0], nil +} + func StopFlagd() error { flagdLock.Lock() defer flagdLock.Unlock() From 55adfabd2128e3f8c76cc7a66b18a103cd5a30ae Mon Sep 17 00:00:00 2001 From: Simon Schrottner Date: Fri, 11 Sep 2026 17:41:17 +0200 Subject: [PATCH 2/2] fix(launchpad): bound the startup probes and ignore answers flagd could not give Two gaps in the readiness probing added by the previous commit, both raised in review. The probes ran off a deadline the loops only consulted between requests, so a probe issued just short of it could still block for the client's 500ms timeout and push /start past the budget it advertises. Derive a context from the budget instead and issue every probe through client.Do with it, so the bound covers the requests themselves and a probe cannot outlive it. flagIsServed treated any response without FLAG_NOT_FOUND in the body as proof that the store was populated, including a 500. A 500 is not an evaluation at all - flagd is saying it could not answer - so it carries no information about the store, and accepting it could let /start report success while flagd was unable to evaluate the probe flag. Poll again on 5xx instead. Non-5xx answers still count: flagd returns 400 with PARSE_ERROR or GENERAL for a flag it holds but cannot resolve from the empty context the probe sends, and that answer only exists once the flag is in the store, so rejecting everything but 200 would turn such a configuration into a 10s timeout and a failed start. Signed-off-by: Simon Schrottner --- launchpad/pkg/flagd.go | 76 ++++++++++++++++++++++++++++++------------ 1 file changed, 55 insertions(+), 21 deletions(-) diff --git a/launchpad/pkg/flagd.go b/launchpad/pkg/flagd.go index 631f71a..51da623 100644 --- a/launchpad/pkg/flagd.go +++ b/launchpad/pkg/flagd.go @@ -124,9 +124,14 @@ func StartFlagd(config string) error { flagdLock.Unlock() client := &http.Client{Timeout: 500 * time.Millisecond} - deadline := time.Now().Add(startupBudget) - if err := awaitReadyz(client, deadline); err != nil { + // Every probe below runs against this context, so the budget bounds the + // requests themselves and not just the gaps between them: a probe issued + // just short of the deadline cannot stretch /start by its own timeout. + ctx, cancel := context.WithTimeout(context.Background(), startupBudget) + defer cancel() + + if err := awaitReadyz(ctx, client); err != nil { _ = StopFlagd() return err } @@ -137,7 +142,7 @@ func StartFlagd(config string) error { // configuration plainly defines. A provider that blocks during its own // initialisation absorbs the window, but a stateless one evaluates the // instant /start returns and races it, so wait for a real evaluation. - if err := awaitFlagsServed(client, configPath, deadline); err != nil { + if err := awaitFlagsServed(ctx, client, configPath); err != nil { _ = StopFlagd() return err } @@ -148,31 +153,42 @@ func StartFlagd(config string) error { // awaitReadyz waits for flagd's readiness probe to report that every sync // source has completed at least one successful data sync. -func awaitReadyz(client *http.Client, deadline time.Time) error { +func awaitReadyz(ctx context.Context, client *http.Client) error { ticker := time.NewTicker(readyPollInterval) defer ticker.Stop() - timeout := time.After(time.Until(deadline)) for { select { - case <-timeout: + case <-ctx.Done(): return fmt.Errorf("flagd health check timed out") case <-ticker.C: - resp, err := client.Get(readyzURL) - if err == nil { - resp.Body.Close() - if resp.StatusCode == http.StatusOK { - return nil - } + if isReady(ctx, client) { + return nil } } } } +// isReady reports whether flagd's readiness probe answers 200. +func isReady(ctx context.Context, client *http.Client) bool { + req, err := http.NewRequestWithContext(ctx, http.MethodGet, readyzURL, nil) + if err != nil { + return false + } + + resp, err := client.Do(req) + if err != nil { + return false + } + defer resp.Body.Close() + + return resp.StatusCode == http.StatusOK +} + // awaitFlagsServed waits until flagd actually resolves a flag from every file // source the configuration lists, so that a successful /start is a promise // that the next evaluation resolves against the new baseline. -func awaitFlagsServed(client *http.Client, configPath string, deadline time.Time) error { +func awaitFlagsServed(ctx context.Context, client *http.Client, configPath string) error { keys, err := probeKeys(configPath) if err != nil { // Without a probe key there is nothing to verify the store with. Fall @@ -183,35 +199,53 @@ func awaitFlagsServed(client *http.Client, configPath string, deadline time.Time } for _, key := range keys { - if err := awaitFlagServed(client, key, deadline); err != nil { + if err := awaitFlagServed(ctx, client, key); err != nil { return err } } return nil } -func awaitFlagServed(client *http.Client, key string, deadline time.Time) error { +func awaitFlagServed(ctx context.Context, client *http.Client, key string) error { + ticker := time.NewTicker(servedPollInterval) + defer ticker.Stop() + for { - if flagIsServed(client, key) { + if flagIsServed(ctx, client, key) { return nil } - if time.Now().After(deadline) { + select { + case <-ctx.Done(): return fmt.Errorf("flagd did not serve flag %q before the startup budget expired", key) + case <-ticker.C: } - time.Sleep(servedPollInterval) } } // flagIsServed reports whether flagd holds the given flag in its store. flagd // answers FLAG_NOT_FOUND both while the store is still empty and for a key it -// genuinely does not hold; every other answer means the flag is being served. -func flagIsServed(client *http.Client, key string) bool { - resp, err := client.Post(ofrepEvaluateURL+url.PathEscape(key), "application/json", strings.NewReader("{}")) +// genuinely does not hold. A 5xx is not an evaluation at all - flagd is telling +// us it could not answer - so it says nothing about the store and is worth +// another poll. Every other answer means flagd resolved the key against a +// populated store, including a 4xx for a flag whose targeting cannot be +// satisfied from the empty context we send. +func flagIsServed(ctx context.Context, client *http.Client, key string) bool { + req, err := http.NewRequestWithContext(ctx, http.MethodPost, ofrepEvaluateURL+url.PathEscape(key), strings.NewReader("{}")) + if err != nil { + return false + } + req.Header.Set("Content-Type", "application/json") + + resp, err := client.Do(req) if err != nil { return false } defer resp.Body.Close() + if resp.StatusCode >= http.StatusInternalServerError { + return false + } + body, err := io.ReadAll(resp.Body) if err != nil { return false