From 74a53c8d273f7721cb251a21125bffecd89eb00a Mon Sep 17 00:00:00 2001 From: Simon Schrottner Date: Sun, 13 Sep 2026 00:31:49 +0200 Subject: [PATCH] feat(launchpad): make /start restore the baseline instead of restarting flagd A conformance client that wants per-scenario isolation has only /start to call -- the control API marks /reset optional (open-feature/spec#423) and this testbed does not implement it -- so it pays a full flagd stop-and-start before every scenario, for a configuration that has not changed. Post-#394 that is ~140-190ms each, over the testbed's ~325 executed scenario instances. The process is not what carries a scenario's leftovers; the flag definitions are. /start now reuses the running flagd when the requested configuration is the one already running, the process is alive, and no delayed restart is pending, and restores the baseline flag state in place instead. Restoring needs no pristine copy of the definitions, which is just as well, because the image does not carry one: /change overwrites its own source, rawflags/changing-flag.json, in the container's writable layer. But /change is a two-state toggle on a single flag, so the shipped baseline is recoverable by reading the current variant and flipping it back when it reads "bar". Reading it rather than remembering it is deliberate. A launchpad that restarts inside a container whose writable layer already holds "bar" would believe a remembered flag, and serve "bar" while reporting a restored baseline. That read is only sound because writes are now honest. /change used to wait for our own file watcher to regenerate the merged file and then return; measured against v0.16.0 in the built image, flagd serves the new variant about a second later, because it watches with the fileinfo watcher whose default poll interval is 1000ms. So /change reported success while flagd still served the old value, and a restore could not tell a file reading "foo" apart from a flagd that had not caught up with one. Both writes now wait for flagd to serve what was written, which is what makes the file content readable as the state flagd is in. Measured in the built image: | case | before | after | | --------------------------------------------- | ----------- | ------- | | /start, same config, nothing changed | ~140-190ms | ~9-10ms | | /start, same config, after a /change | ~140-190ms | ~1.0s | | /start, other config, or after /stop, or in a | | | | /restart window | ~140-190ms | ~125ms | | /change | ~0ms, wrong | ~1.0s | Two scenarios call /change, so a suite of a few hundred scenarios pays the second row about three times and the first row for everything else. Also fixes a leaked flagdLock: StartFlagd returned without unlocking when stopFlagDWithoutLock failed. Signed-off-by: Simon Schrottner --- launchpad/pkg/file.go | 130 +++++++++++++++++++----- launchpad/pkg/flagd.go | 193 ++++++++++++++++++++++++++++++++---- launchpad/pkg/flagd_test.go | 110 ++++++++++++++++++++ 3 files changed, 393 insertions(+), 40 deletions(-) create mode 100644 launchpad/pkg/flagd_test.go diff --git a/launchpad/pkg/file.go b/launchpad/pkg/file.go index cf5dbc9..81a7a3e 100644 --- a/launchpad/pkg/file.go +++ b/launchpad/pkg/file.go @@ -21,6 +21,20 @@ type FlagConfig struct { } `json:"flags"` } +const ( + // ChangingFlagFile is the only flag definition the launchpad ever writes to, + // and changing-flag is the only flag in it. Keeping that true is what lets + // the baseline be restored without a pristine copy of the definitions - see + // RestoreChangingFlag. + ChangingFlagFile = "rawflags/changing-flag.json" + // ChangingFlagKey is the flag /change toggles. + ChangingFlagKey = "changing-flag" + // BaselineChangingVariant is the defaultVariant changing-flag ships with. + BaselineChangingVariant = "foo" + // toggledChangingVariant is the other half of the toggle. + toggledChangingVariant = "bar" +) + var ( fileLock sync.Mutex // lock for file operations changeLock sync.Mutex // lock for change requests (so that multiple requests don't overlap) @@ -28,44 +42,116 @@ var ( changeFlagUpdateListeners []*sync.WaitGroup ) +// ToggleChangingFlag flips changing-flag to its other variant and reports the +// one now in force. func ToggleChangingFlag() (string, error) { changeLock.Lock() defer changeLock.Unlock() - // Path to the configuration file - configFile := "rawflags/changing-flag.json" + current, err := readChangingVariant() + if err != nil { + return "", err + } - // Read the existing file - data, err := os.ReadFile(configFile) + next := BaselineChangingVariant + if current == BaselineChangingVariant { + next = toggledChangingVariant + } + + return next, writeChangingVariant(next) +} + +// RestoreChangingFlag puts changing-flag back to the variant it ships with and +// reports whether it had to write anything. +// +// /change is the only endpoint that mutates a flag definition, and it mutates +// exactly one flag with exactly two states, so the shipped baseline is +// recoverable by flipping the toggle back rather than by restoring from a +// pristine copy of the definitions - which the image does not carry, because +// /change overwrites its own source in the container's writable layer. +// +// The current variant is read rather than remembered on purpose. A launchpad +// that restarts inside a container whose writable layer already holds "bar" +// would believe a remembered flag, and go on serving "bar" while reporting a +// restored baseline. Reading it is also what makes the common case free: most +// scenarios never call /change, so most restores write nothing at all. +// +// Reading it is only sound because every write waits for flagd to serve what it +// wrote - see writeChangingVariant. Without that, a file already reading "foo" +// could not be told apart from a flagd that has not caught up with it yet. +func RestoreChangingFlag() (bool, error) { + changeLock.Lock() + defer changeLock.Unlock() + + current, err := readChangingVariant() + if err != nil { + return false, err + } + if current == BaselineChangingVariant { + return false, nil + } + + return true, writeChangingVariant(BaselineChangingVariant) +} + +// readChangingVariant reports the defaultVariant currently written to +// changing-flag's definition. +func readChangingVariant() (string, error) { + data, err := os.ReadFile(ChangingFlagFile) if err != nil { return "", err } - // Parse the JSON into the FlagConfig struct var config FlagConfig if err := json.Unmarshal(data, &config); err != nil { return "", err } - // Find the "changing-flag" and toggle the default variant - flag, exists := config.Flags["changing-flag"] + flag, exists := config.Flags[ChangingFlagKey] if !exists { return "", errors.New("changing-flag not found in configuration") } + return flag.DefaultVariant, nil +} - // Toggle the defaultVariant between "foo" and "bar" - if flag.DefaultVariant == "foo" { - flag.DefaultVariant = "bar" - } else { - flag.DefaultVariant = "foo" +// writeChangingVariant sets changing-flag's defaultVariant and does not return +// until flagd serves it. +// +// Two watchers stand between the write and the served value: ours, which +// regenerates the merged flag file, and flagd's, which re-reads it. Waiting on +// ours alone used to be the whole of this function, and it is not enough - +// measured against v0.16.0, flagd serves the new variant around 500ms after +// /change has reported success. +// +// Waiting for both is what lets the current file content be read as the state +// flagd is in, which is the invariant RestoreChangingFlag depends on: if the +// file says "foo", flagd is serving "foo", so there is nothing to restore. +func writeChangingVariant(variant string) error { + // Read the existing file + data, err := os.ReadFile(ChangingFlagFile) + if err != nil { + return err } + // Parse the JSON into the FlagConfig struct + var config FlagConfig + if err := json.Unmarshal(data, &config); err != nil { + return err + } + + // Find the "changing-flag" and set the default variant + flag, exists := config.Flags[ChangingFlagKey] + if !exists { + return errors.New("changing-flag not found in configuration") + } + flag.DefaultVariant = variant + // Save the updated flag back to the configuration - config.Flags["changing-flag"] = flag + config.Flags[ChangingFlagKey] = flag // Serialize the updated configuration back to JSON updatedData, err := json.MarshalIndent(config, "", " ") if err != nil { - return "", err + return err } // the file watcher should be triggered instantly. If not, we add a timeout to prevent a hanging test @@ -86,20 +172,18 @@ func ToggleChangingFlag() (string, error) { // Write the updated JSON back to the file fileLock.Lock() - if err := atomicWriteFile(configFile, updatedData); err != nil { + if err := atomicWriteFile(ChangingFlagFile, updatedData); err != nil { fileLock.Unlock() - return "", err + return err } fileLock.Unlock() - select { - case <-ctx.Done(): - if errors.Is(ctx.Err(), context.DeadlineExceeded) { - return "", fmt.Errorf("Flags were not updated in time: %v", ctx.Err()) - } else { - return flag.DefaultVariant, nil - } + <-ctx.Done() + if errors.Is(ctx.Err(), context.DeadlineExceeded) { + return fmt.Errorf("flags were not updated in time: %v", ctx.Err()) } + + return awaitVariantInFlagd(variant) } func RestartFileWatcher() error { diff --git a/launchpad/pkg/flagd.go b/launchpad/pkg/flagd.go index 51da623..af4e4ef 100644 --- a/launchpad/pkg/flagd.go +++ b/launchpad/pkg/flagd.go @@ -11,6 +11,7 @@ import ( "net/url" "os" "os/exec" + "path/filepath" "sort" "strings" "sync" @@ -25,6 +26,10 @@ const ( readyPollInterval = 100 * time.Millisecond servedPollInterval = 10 * time.Millisecond + // probeTimeout bounds a single probe request. The shared startup deadline + // bounds the sequence of them. + probeTimeout = 500 * 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 @@ -91,13 +96,22 @@ func RestartFlagd(seconds int) { } func StartFlagd(config string) error { + flagdLock.Lock() + if config == "" { config = Config - } else { - Config = config } - flagdLock.Lock() + // The running process is only reusable if it is the one that would be + // started anyway: same configuration, still alive, and not inside a delayed + // restart's downtime window. A pending restart implies flagd is already + // stopped, so the process check covers it, but saying so is cheaper than + // relying on that. + reuse := flagdCmd != nil && flagdCmd.Process != nil && + restartCancelFunc == nil && config == Config + + Config = config + // Cancel any pending restart attempts if restartCancelFunc != nil { restartCancelFunc() @@ -105,14 +119,20 @@ func StartFlagd(config string) error { restartCancelFunc = nil } + configPath := flagdConfigPath(config) + + if reuse { + flagdLock.Unlock() + return resumeRunningFlagd() + } + if err := stopFlagDWithoutLock(); err != nil { + flagdLock.Unlock() return err } ensureStartConditions() - configPath := fmt.Sprintf("./configs/%s.json", config) - flagdCmd = exec.Command("./flagd", "start", "--config", configPath) flagdCmd.Stdout = os.Stdout flagdCmd.Stderr = os.Stderr @@ -123,7 +143,7 @@ func StartFlagd(config string) error { } flagdLock.Unlock() - client := &http.Client{Timeout: 500 * time.Millisecond} + client := &http.Client{Timeout: probeTimeout} // Every probe below runs against this context, so the budget bounds the // requests themselves and not just the gaps between them: a probe issued @@ -151,6 +171,130 @@ func StartFlagd(config string) error { return nil } +func flagdConfigPath(config string) string { + return fmt.Sprintf("./configs/%s.json", config) +} + +// resumeRunningFlagd restores the baseline flag state in the flagd that is +// already running, instead of restarting it. +// +// Restarting is how /start achieves isolation today, and for a client that has +// no /reset to call it is the only way to get it - which means a full flagd +// restart before every scenario, for a configuration that has not changed. The +// process is not what carries the scenario's leftovers; the flag definitions +// are. Putting those back is enough, and it is what the existing file watcher +// is already there to deliver. +func resumeRunningFlagd() error { + restored, err := RestoreChangingFlag() + if err != nil { + return fmt.Errorf("failed to restore the baseline flag state: %w", err) + } + + if !restored { + // Nothing has mutated a flag definition since the last restore, so the + // running flagd is already serving the baseline. Most scenarios never + // call /change, which makes this the common case and a free one. + fmt.Println("flagd reused; baseline flag state already in force.") + return nil + } + + // RestoreChangingFlag does not return until flagd serves the restored + // variant, so there is nothing left to wait for here. + fmt.Println("flagd reused; baseline flag state restored.") + return nil +} + +// awaitVariantInFlagd waits until the running flagd resolves changing-flag to +// the given variant. +// +// There is nothing to wait for when flagd is not running, or when the running +// configuration does not read the merged flag file - that file is the only +// source changing-flag reaches flagd through, so a configuration without it can +// never serve the flag and polling for it would burn the whole budget before +// failing. +func awaitVariantInFlagd(variant string) error { + if !flagdIsRunning() { + return nil + } + if !servesCombinedFlags(flagdConfigPath(Config)) { + return nil + } + + ctx, cancel := context.WithTimeout(context.Background(), startupBudget) + defer cancel() + + return awaitVariantServed(ctx, &http.Client{Timeout: probeTimeout}, ChangingFlagKey, variant) +} + +func flagdIsRunning() bool { + flagdLock.Lock() + defer flagdLock.Unlock() + + return flagdCmd != nil && flagdCmd.Process != nil +} + +// awaitVariantServed waits until flagd resolves key to the given variant. +func awaitVariantServed(ctx context.Context, client *http.Client, key, variant string) error { + ticker := time.NewTicker(servedPollInterval) + defer ticker.Stop() + + for { + if servedVariant(ctx, client, key) == variant { + return nil + } + select { + case <-ctx.Done(): + return fmt.Errorf("flagd did not serve %q as variant %q before the startup budget expired", key, variant) + case <-ticker.C: + } + } +} + +// servedVariant reports the variant flagd currently resolves key to, and an +// empty string for any answer that is not a successful evaluation. +func servedVariant(ctx context.Context, client *http.Client, key string) string { + req, err := http.NewRequestWithContext(ctx, http.MethodPost, ofrepEvaluateURL+url.PathEscape(key), strings.NewReader("{}")) + if err != nil { + return "" + } + req.Header.Set("Content-Type", "application/json") + + resp, err := client.Do(req) + if err != nil { + return "" + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return "" + } + + var evaluation struct { + Variant string `json:"variant"` + } + if err := json.NewDecoder(resp.Body).Decode(&evaluation); err != nil { + return "" + } + return evaluation.Variant +} + +// servesCombinedFlags reports whether the configuration reads the merged flag +// file that changing-flag reaches flagd through. +func servesCombinedFlags(configPath string) bool { + sources, err := fileSources(configPath) + if err != nil { + fmt.Printf("Cannot tell whether %s serves the merged flag file: %v\n", configPath, err) + return false + } + + for _, uri := range sources { + if filepath.Clean(uri) == filepath.Clean(OutputFile) { + return true + } + } + return false +} + // awaitReadyz waits for flagd's readiness probe to report that every sync // source has completed at least one successful data sync. func awaitReadyz(ctx context.Context, client *http.Client) error { @@ -259,6 +403,29 @@ func flagIsServed(ctx context.Context, client *http.Client, key string) bool { // covers each source separately because they are merged into the store // independently. func probeKeys(configPath string) ([]string, error) { + sources, err := fileSources(configPath) + if err != nil { + return nil, err + } + + var keys []string + for _, uri := range sources { + key, err := enabledFlagKey(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 +} + +// fileSources returns the URIs of every file source the given flagd +// configuration reads. +func fileSources(configPath string) ([]string, error) { content, err := os.ReadFile(configPath) if err != nil { return nil, fmt.Errorf("failed to read flagd config: %w", err) @@ -274,22 +441,14 @@ func probeKeys(configPath string) ([]string, error) { return nil, fmt.Errorf("failed to parse flagd config: %w", err) } - var keys []string + var uris []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") + uris = append(uris, source.URI) } - return keys, nil + return uris, nil } // enabledFlagKey picks a stable, enabled flag key from a flag definition file. diff --git a/launchpad/pkg/flagd_test.go b/launchpad/pkg/flagd_test.go new file mode 100644 index 0000000..b6e32e1 --- /dev/null +++ b/launchpad/pkg/flagd_test.go @@ -0,0 +1,110 @@ +package flagd + +import ( + "os" + "path/filepath" + "testing" +) + +// writeConfig writes a flagd configuration to a temporary file and returns its +// path. +func writeConfig(t *testing.T, body string) string { + t.Helper() + + path := filepath.Join(t.TempDir(), "config.json") + if err := os.WriteFile(path, []byte(body), 0o644); err != nil { + t.Fatalf("failed to write config: %v", err) + } + return path +} + +func TestFileSources(t *testing.T) { + tests := []struct { + name string + body string + want []string + }{ + { + name: "file sources are returned in order", + body: `{"sources":[{"uri":"flags/allFlags.json","provider":"file"},{"uri":"rawflags/selector-flags.json","provider":"file"}]}`, + want: []string{"flags/allFlags.json", "rawflags/selector-flags.json"}, + }, + { + name: "non-file providers are skipped", + body: `{"sources":[{"uri":"http://example.com","provider":"http"},{"uri":"flags/allFlags.json","provider":"file"}]}`, + want: []string{"flags/allFlags.json"}, + }, + { + name: "a configuration without sources yields none", + body: `{}`, + want: nil, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, err := fileSources(writeConfig(t, tt.body)) + if err != nil { + t.Fatalf("fileSources returned an error: %v", err) + } + if len(got) != len(tt.want) { + t.Fatalf("got %v, want %v", got, tt.want) + } + for i := range got { + if got[i] != tt.want[i] { + t.Fatalf("got %v, want %v", got, tt.want) + } + } + }) + } +} + +func TestFileSourcesRejectsUnreadableAndInvalidConfigs(t *testing.T) { + if _, err := fileSources(filepath.Join(t.TempDir(), "missing.json")); err == nil { + t.Fatal("expected an error for a missing configuration") + } + if _, err := fileSources(writeConfig(t, "not json")); err == nil { + t.Fatal("expected an error for an unparseable configuration") + } +} + +// TestServesCombinedFlags covers the distinction the restore path depends on: +// changing-flag only reaches flagd through the merged flag file, so a +// configuration that does not read it can never serve the flag and must not be +// waited on. +func TestServesCombinedFlags(t *testing.T) { + tests := []struct { + name string + body string + want bool + }{ + { + name: "the default shape reads the merged file", + body: `{"sources":[{"uri":"flags/allFlags.json","provider":"file"},{"uri":"rawflags/selector-flags.json","provider":"file"}]}`, + want: true, + }, + { + name: "the merged file is recognised through an unclean path", + body: `{"sources":[{"uri":"./flags/allFlags.json","provider":"file"}]}`, + want: true, + }, + { + name: "the metadata shape does not", + body: `{"sources":[{"uri":"rawflags/selector-flag-combined-metadata.json","provider":"file"}]}`, + want: false, + }, + { + name: "an unreadable configuration does not", + body: `not json`, + want: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got := servesCombinedFlags(writeConfig(t, tt.body)); got != tt.want { + t.Fatalf("got %v, want %v", got, tt.want) + } + }) + } +}