|
1 | 1 | package engine |
2 | 2 |
|
3 | 3 | import ( |
4 | | - "bufio" |
5 | 4 | "context" |
6 | 5 | "crypto/sha256" |
7 | 6 | "encoding/hex" |
8 | 7 | "errors" |
9 | 8 | "fmt" |
10 | | - "io" |
11 | 9 | "io/fs" |
12 | 10 | "os" |
13 | 11 | "path/filepath" |
@@ -162,11 +160,25 @@ type Options struct { |
162 | 160 | // OnProgress is called as the run advances. Nil means silence, which is |
163 | 161 | // what every caller that has nobody to show it to should pass. |
164 | 162 | // |
165 | | - // Called from the same goroutine doing the work, so there is no |
166 | | - // concurrency here to get wrong. Called often - once per write inside a |
167 | | - // file, not only once per finished file - so rate limiting what actually |
168 | | - // reaches a screen belongs to the caller. Without the writes inside a |
169 | | - // file, one 5 GB file would report once, at the end. |
| 163 | + // NEVER TWO AT ONCE, though not always from the same goroutine. Until |
| 164 | + // 2026-09-06 this promised the stronger thing - "from the same goroutine |
| 165 | + // doing the work" - and both callers were built on it: the command line bar |
| 166 | + // moves last and printed without a lock, and the window's throttle reads |
| 167 | + // and writes a timestamp without one. The files are written over several |
| 168 | + // goroutines now, so the engine serialises these calls instead. The lock |
| 169 | + // gives happens-before, so both callers stay correct unchanged. What a |
| 170 | + // caller may NOT do is assume the goroutine, which is why the sentence is |
| 171 | + // here rather than only in the commit that changed it. |
| 172 | + // |
| 173 | + // Called often - once per write inside a file, not only once per finished |
| 174 | + // file - so rate limiting what actually reaches a screen belongs to the |
| 175 | + // caller. Without the writes inside a file, one 5 GB file would report |
| 176 | + // once, at the end. |
| 177 | + // |
| 178 | + // With several files in flight the byte count is the whole run's, so it |
| 179 | + // moves while any writer moves rather than tracking one file. It still |
| 180 | + // falls back when a file fails, exactly as it did before, because a file |
| 181 | + // that failed counts for nothing. |
170 | 182 | OnProgress func(Progress) |
171 | 183 | } |
172 | 184 |
|
@@ -443,6 +455,12 @@ type Progress struct { |
443 | 455 | // |
444 | 456 | // A manifest is returned even when the run is cut short, otherwise cleanup |
445 | 457 | // has nothing to work with. |
| 458 | +// |
| 459 | +// The writing itself happens over several goroutines, and everything about |
| 460 | +// that lives in parallel.go - including why, and what it measured. What stays |
| 461 | +// here is everything a run does exactly once: the checks that decide whether |
| 462 | +// it may start at all, and the reading back of the answers in the order the |
| 463 | +// plan lists them. |
446 | 464 | func Run(ctx context.Context, files []PlannedFile, opt Options) (*Result, error) { |
447 | 465 | m := manifest.New( |
448 | 466 | "testing-files-generator", version.Version, |
@@ -516,133 +534,58 @@ func Run(ctx context.Context, files []PlannedFile, opt Options) (*Result, error) |
516 | 534 | } |
517 | 535 | }() |
518 | 536 |
|
519 | | - totalBytes := TotalBytes(files) |
520 | | - var bytesDone int64 |
521 | | - |
522 | | - for i, f := range files { |
523 | | - select { |
524 | | - case <-ctx.Done(): |
525 | | - // Stop starting new files. What is already finished stays, and |
526 | | - // the manifest describes exactly that. |
527 | | - m.Run.Complete = false |
528 | | - return res, ctx.Err() |
529 | | - default: |
530 | | - } |
| 537 | + // The files are written over several goroutines. Everything that runs |
| 538 | + // beside anything else lives in parallel.go, including the measurements |
| 539 | + // that put it there. |
| 540 | + written := writeAll(ctx, files, opt.OutDir, newProgressGate(files, opt.OnProgress)) |
531 | 541 |
|
532 | | - // Built per file rather than once, because it closes over how far the |
533 | | - // run had got before this file started. Left nil when nobody is |
534 | | - // listening, so a run without progress allocates nothing for it. |
535 | | - var report func(int64) |
536 | | - if opt.OnProgress != nil { |
537 | | - report = func(inFile int64) { |
538 | | - opt.OnProgress(Progress{ |
539 | | - FilesDone: i, FilesTotal: len(files), |
540 | | - BytesDone: bytesDone + inFile, BytesTotal: totalBytes, |
541 | | - }) |
542 | | - } |
543 | | - } |
544 | | - |
545 | | - sum, err := writeOne(ctx, f, opt.OutDir, report) |
546 | | - if err == nil { |
547 | | - // Only what reached the disk. Counting a file that failed would |
548 | | - // have the bar claim bytes nobody can find, and on a run where |
549 | | - // several fail the total would arrive before the files do. |
550 | | - bytesDone += f.Plan.Bytes |
551 | | - } |
552 | | - if opt.OnProgress != nil { |
553 | | - opt.OnProgress(Progress{ |
554 | | - FilesDone: i + 1, FilesTotal: len(files), |
555 | | - BytesDone: bytesDone, BytesTotal: totalBytes, |
556 | | - }) |
557 | | - } |
558 | | - if err != nil { |
559 | | - if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { |
560 | | - m.Run.Complete = false |
561 | | - return res, err |
| 542 | + // Read back in the order the plan lists, on this goroutine alone. Two |
| 543 | + // things rest on that and neither is tidiness: |
| 544 | + // |
| 545 | + // - the manifest keeps the order it has always had, which is the order |
| 546 | + // cleanup prints to a person before deleting from it, |
| 547 | + // - a run stopped part way names the LOWEST cancelled file rather than |
| 548 | + // whichever writer happened to notice first, so the same interruption |
| 549 | + // reports the same thing on every machine. |
| 550 | + var stopped error |
| 551 | + for i, r := range written { |
| 552 | + switch { |
| 553 | + case r.ok: |
| 554 | + m.Add(entryFor(files[i], r.sha, true, nil)) |
| 555 | + case r.err == nil: |
| 556 | + // Never started. A cancelled run leaves these behind and they are |
| 557 | + // neither a success nor a failure, so they get no entry - which is |
| 558 | + // what the sequential loop did by never reaching them. |
| 559 | + case errors.Is(r.err, context.Canceled) || errors.Is(r.err, context.DeadlineExceeded): |
| 560 | + // A writer stopped half way through wrote nothing that survived, |
| 561 | + // so there is nothing to record about it either. |
| 562 | + if stopped == nil { |
| 563 | + stopped = r.err |
562 | 564 | } |
| 565 | + default: |
563 | 566 | // One file failing does not end the run. Nine thousand good |
564 | 567 | // files are worth keeping, and the entry says what went wrong. |
565 | 568 | res.Failures++ |
566 | | - m.Add(entryFor(f, "", false, err)) |
567 | | - continue |
| 569 | + m.Add(entryFor(files[i], "", false, r.err)) |
568 | 570 | } |
569 | | - m.Add(entryFor(f, sum, true, nil)) |
570 | 571 | } |
571 | 572 |
|
572 | | - m.Run.Complete = true |
573 | | - return res, nil |
574 | | -} |
575 | | - |
576 | | -func writeOne(ctx context.Context, f PlannedFile, outDir string, report func(int64)) (string, error) { |
577 | | - final := filepath.Join(outDir, f.Name) |
578 | | - // The process id is in the name because two runs writing into one directory |
579 | | - // used to meet on it. Measured on 2026-08-03: two runs of the same target |
580 | | - // collided on the temporary file, one of them reported two files it could |
581 | | - // not produce, and the bytes of the other had already gone through the same |
582 | | - // handle. The name never survives the run, so nothing about it has to be |
583 | | - // repeatable - and the file it becomes is settled by the plan, not by this. |
584 | | - tmp := tempPathFor(outDir, f.Name) |
585 | | - |
586 | | - // os.Create, and O_EXCL was tried here and taken back out on 2026-08-25. |
587 | | - // |
588 | | - // The idea was sound: the check in preflight answers "this name is free" |
589 | | - // a few hundred lines before the write, and O_EXCL would have the |
590 | | - // filesystem answer it at the moment of writing instead. What it costs on |
591 | | - // Windows is not sound. Measured with a probe, a file created in a |
592 | | - // directory reached through a symbolic link: |
593 | | - // |
594 | | - // os.Create works |
595 | | - // O_CREATE|O_EXCL|O_WRONLY fails with "The file exists" |
596 | | - // |
597 | | - // about a file that does not exist. Go asks for the reparse point rather |
598 | | - // than what it points at when O_EXCL is set, so every file of a run whose |
599 | | - // output directory is a link fails - and this tool supports exactly that |
600 | | - // on purpose, because people keep fixtures on a mounted workspace or a |
601 | | - // scratch disk. Two guards said so within a minute of the change. |
602 | | - // |
603 | | - // The window O_EXCL would have closed is a real one and it is small: |
604 | | - // preflight refuses every name that is taken before the run starts, so |
605 | | - // what is left is somebody else creating our temporary name, with our |
606 | | - // process id in it, during the run. Trading a supported way of pointing |
607 | | - // the tool at a directory for that is the wrong way round. |
608 | | - fh, err := os.Create(tmp) |
609 | | - if err != nil { |
610 | | - return "", err |
611 | | - } |
612 | | - |
613 | | - h := sha256.New() |
614 | | - buffered := bufio.NewWriterSize(fh, 64<<10) |
615 | | - counter := &countingWriter{w: io.MultiWriter(buffered, h), report: report} |
616 | | - |
617 | | - writeErr := writeWithoutCrashing(ctx, f, counter) |
618 | | - if writeErr == nil { |
619 | | - writeErr = buffered.Flush() |
620 | | - } |
621 | | - closeErr := fh.Close() |
622 | | - |
623 | | - if writeErr != nil { |
624 | | - _ = os.Remove(tmp) |
625 | | - return "", writeErr |
626 | | - } |
627 | | - if closeErr != nil { |
628 | | - _ = os.Remove(tmp) |
629 | | - return "", closeErr |
| 573 | + // A stopped run keeps every file that FINISHED, which may leave a hole |
| 574 | + // where a writer was cut off. The sequential loop could only ever leave a |
| 575 | + // contiguous prefix, so this is the one thing a person can observe that |
| 576 | + // changed - decided by the owner on 2026-09-06, and the alternative is |
| 577 | + // worse in a way untouchable rule 7 names: a finished file with no entry |
| 578 | + // in the manifest is a file no command of this tool can remove. |
| 579 | + if stopped == nil && ctx.Err() != nil { |
| 580 | + stopped = ctx.Err() |
630 | 581 | } |
631 | | - |
632 | | - // The size is the promise. A generator that missed it by a byte is a bug |
633 | | - // worth catching here rather than in someone's test suite, so the file |
634 | | - // never reaches its final name. |
635 | | - if counter.n != f.Plan.Bytes { |
636 | | - _ = os.Remove(tmp) |
637 | | - return "", fmt.Errorf("generator for %s produced %d B where the plan said %d B", |
638 | | - f.Desc.ID, counter.n, f.Plan.Bytes) |
| 582 | + if stopped != nil { |
| 583 | + m.Run.Complete = false |
| 584 | + return res, stopped |
639 | 585 | } |
640 | 586 |
|
641 | | - if err := os.Rename(tmp, final); err != nil { |
642 | | - _ = os.Remove(tmp) |
643 | | - return "", err |
644 | | - } |
645 | | - return hex.EncodeToString(h.Sum(nil)), nil |
| 587 | + m.Run.Complete = true |
| 588 | + return res, nil |
646 | 589 | } |
647 | 590 |
|
648 | 591 | func entryFor(f PlannedFile, sha string, materialized bool, failure error) manifest.File { |
@@ -927,22 +870,3 @@ func runID(seed int64) string { |
927 | 870 | h := sha256.Sum256([]byte(fmt.Sprintf("run:%d", seed))) |
928 | 871 | return "run_" + hex.EncodeToString(h[:5]) |
929 | 872 | } |
930 | | - |
931 | | -type countingWriter struct { |
932 | | - w io.Writer |
933 | | - n int64 |
934 | | - // report, when set, is called with the running total for this file. It is |
935 | | - // what gives progress inside a single large file rather than only between |
936 | | - // files - the case where silence is worst, because one 5 GB file is one |
937 | | - // callback if you only count finished files. |
938 | | - report func(int64) |
939 | | -} |
940 | | - |
941 | | -func (c *countingWriter) Write(p []byte) (int, error) { |
942 | | - n, err := c.w.Write(p) |
943 | | - c.n += int64(n) |
944 | | - if c.report != nil { |
945 | | - c.report(c.n) |
946 | | - } |
947 | | - return n, err |
948 | | -} |
0 commit comments