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) + } + }) + } +}