fix(02): write fixtures exclusively and only on successful flush (WR-01)

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
Jakub Zych
2026-09-17 14:41:18 +02:00
parent 5994e671b3
commit 7e70afe819
5 changed files with 92 additions and 9 deletions

View File

@@ -49,6 +49,15 @@ func ParseFlow(raw []byte) (Flow, error) {
// SaveFlow writes a validated flow atomically, syncing before rename.
func SaveFlow(path string, flow Flow) error {
return saveFlow(path, flow, false)
}
// SaveFlowExclusive writes a validated flow and fails if path already exists.
func SaveFlowExclusive(path string, flow Flow) error {
return saveFlow(path, flow, true)
}
func saveFlow(path string, flow Flow, exclusive bool) error {
if err := validateFlow(flow); err != nil {
return err
}
@@ -62,6 +71,30 @@ func SaveFlow(path string, flow Flow) error {
return fmt.Errorf("tide: create fixture dir: %w", err)
}
}
if exclusive {
f, err := os.OpenFile(path, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0o644)
if err != nil {
return fmt.Errorf("tide: create fixture: %w", err)
}
ok := false
defer func() {
_ = f.Close()
if !ok {
_ = os.Remove(path)
}
}()
if _, err := f.Write(raw); err != nil {
return fmt.Errorf("tide: write fixture: %w", err)
}
if err := f.Sync(); err != nil {
return fmt.Errorf("tide: sync fixture: %w", err)
}
if err := f.Close(); err != nil {
return fmt.Errorf("tide: close fixture: %w", err)
}
ok = true
return nil
}
tmp, err := os.CreateTemp(dir, ".tide-*.tmp")
if err != nil {
return fmt.Errorf("tide: create temp fixture: %w", err)

View File

@@ -336,7 +336,10 @@ func recordSeed(ctx context.Context, seed *Seed, cfg ManifestConfig) error {
if err != nil {
return err
}
return SaveFlow(dest, flow)
if _, err := os.Stat(dest); err == nil {
return SaveFlow(dest, flow)
}
return SaveFlowExclusive(dest, flow)
}
func seedAlreadyRecorded(path string) (bool, error) {
@@ -404,7 +407,13 @@ func recordRoute(ctx context.Context, route Route, cfg ManifestConfig) error {
if err != nil {
return err
}
if err := SaveFlow(dest, flow); err != nil {
if cfg.Update {
if err := SaveFlow(dest, flow); err != nil {
return err
}
continue
}
if err := SaveFlowExclusive(dest, flow); err != nil {
return err
}
}

View File

@@ -48,6 +48,7 @@ type ProxyConfig struct {
VarsPath string
Fixtures string
MaxBody int64
Update bool
}
// Proxy is a fixed-upstream recording reverse proxy.
@@ -282,7 +283,7 @@ func (p *Proxy) recordStep(state *captureState, resp *http.Response, respBody []
return err
}
buf.steps = append(buf.steps, step)
return p.writeSessionLocked(buf)
return nil
}
func (p *Proxy) sessionLocked(name string) (*sessionBuf, error) {
@@ -290,8 +291,10 @@ func (p *Proxy) sessionLocked(name string) (*sessionBuf, error) {
return buf, nil
}
dest := p.sessionPath(name)
if _, err := os.Stat(dest); err == nil {
return nil, fmt.Errorf("tide: duplicate session %q", name)
if !p.cfg.Update {
if _, err := os.Stat(dest); err == nil {
return nil, fmt.Errorf("tide: duplicate session %q", name)
}
}
buf := &sessionBuf{name: name}
p.sessions[name] = buf
@@ -310,7 +313,11 @@ func (p *Proxy) writeSessionLocked(buf *sessionBuf) error {
Name: buf.name,
Steps: append([]Step(nil), buf.steps...),
}
return SaveFlow(p.sessionPath(buf.name), flow)
dest := p.sessionPath(buf.name)
if p.cfg.Update {
return SaveFlow(dest, flow)
}
return SaveFlowExclusive(dest, flow)
}
func (p *Proxy) sessionPath(name string) string {
@@ -323,9 +330,7 @@ func (p *Proxy) failSession(name string, err error) {
if _, ok := p.failed[name]; !ok {
p.failed[name] = err
}
if buf := p.sessions[name]; buf != nil && len(buf.steps) == 0 {
_ = os.Remove(p.sessionPath(name))
}
_ = os.Remove(p.sessionPath(name))
}
func validateSessionName(name string) error {

View File

@@ -305,6 +305,40 @@ func TestProxyDuplicateSessionName(t *testing.T) {
}
}
func TestProxyFailedSessionLeavesNoFixture(t *testing.T) {
upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
if r.URL.Path == "/ok" {
_, _ = w.Write([]byte(`{"ok":true}`))
return
}
_, _ = w.Write([]byte(`{"token":"` + testJWT + `"}`))
}))
t.Cleanup(upstream.Close)
fixtures := t.TempDir()
proxy := newTestProxy(t, upstream.URL, fixtures, testRulesYAML())
srv := httptest.NewServer(proxy.Handler())
t.Cleanup(srv.Close)
doProxy(t, srv.URL, "partial", "/ok", "")
req, err := http.NewRequest(http.MethodGet, srv.URL+"/leak", nil)
if err != nil {
t.Fatal(err)
}
req.Header.Set(SessionHeader, "partial")
resp, err := http.DefaultClient.Do(req)
if err != nil {
t.Fatal(err)
}
resp.Body.Close()
if err := proxy.Flush(); err == nil {
t.Fatal("failed session flush must surface the capture error")
}
if _, err := os.Stat(filepath.Join(fixtures, "nuxt", "partial.yaml")); !os.IsNotExist(err) {
t.Fatalf("partial session fixture committed: %v", err)
}
}
func newTestProxy(t *testing.T, upstream, fixtures, rulesYAML string) *Proxy {
t.Helper()
p, err := NewProxy(ProxyConfig{