Repository navigation
feat(tokenize): skip and record blobs whose tokenizer fails - #96
EllianCarlos wants to merge 13 commits into
Conversation
A timeout, a parser crash (exit 33) or empty output no longer stops the walk. blobExec drops the blob from the rewritten trees, as it drops a denylisted blob: no raw source, no empty tokenization. The walk continues, and the exit status is 0. Each dropped, denylisted or oversized blob gets one row in the file that --skipped-tsv names (sha, path, reason, detail, tokenizer). The runner passes <work>/tokenize-skipped.tsv and prints a summary with the count for each reason. A resumed run writes no duplicate row. A timed-out blob also gets a retry_blob row, and the next run tries it first. If it then tokenizes, the run empties commit_map, tree_map and ref_map and folds all of history again. --refold-marker tells the runner, and the runner removes the outputs of the old commits. --strict-tokenize (or CREGIT_STRICT_TOKENIZE=1) restores the old stop on exit 4 and 6. Sharded mode is always strict, because the shard merge does not carry the retry records. SkipFailedBlobSpec covers the three walkers, strict mode, the skip file and the re-fold. The README documents the skip file and strict mode.
--max-retries N (or CREGIT_MAX_RETRIES) gives a timed-out blob up to N more attempts in the same run. The default is 0: a timeout skips and records the blob at once, and the stall watchdog window stays at 3 x --blob-timeout (30 min). Each retry first waits while the 1-minute load average is above --load-limit (default 2 x the processor count, CREGIT_LOAD_LIMIT), for --load-wait-max seconds at most (default 600). Each retry has --timeout-retry-factor (default 3) times the budget. The first success ends the sequence. The skip detail lists all attempts, as in timeout=600s retries=3x1800s. Crashes get no retry. With retries, the defaulted stall window grows to the longest sequence for one blob plus one --blob-timeout (8420 s for 3 retries). An explicit window that is too small is refused. A recovered timeout from an earlier run folds history again, and the runner then deletes blame/ and blame-c100-incoming/: a full re-blame. So step 2 passes --no-retry-timed-out while either directory holds a .blame file, and prints how many blobs stay dropped and how many .blame files a retry would delete. --retry-skipped (or CREGIT_RETRY_SKIPPED=1) forces the retry. Tests cover the load gate, the retry sequence, the default of 0, an explicit 3, and the runner pass-through. The README documents the options, the defaults and the worst case.
EllianCarlos
left a comment
There was a problem hiding this comment.
Automated review (Claude Code, three axes: spec, standards, crop). I checked the main findings in the code before I posted it.
Verdict. The skip-and-record path does what the PR says in the three walkers, and the tests pass (I ran SkipFailedBlobSpec + MainOptionsSpec: 68 passed, 0 failed; test_tokenize_gate.sh in devenv shell: 90 passed, 0 failed). But the re-fold after a recovered timeout is not safe if the run stops part way: the recovered blob or the stale-blame cleanup can be lost for good. Strict mode also exits 0 while held blobs stay dropped. The changeset can lose about 120 lines with no loss of tested behaviour, and about 630 lines if the in-run retry moves to its own PR.
| # | Severity | Where | Finding | Lines saved |
|---|---|---|---|---|
| 1 | High | Walker.scala 1446-1463, Main.scala 864-869 | clearFold() and the refold marker come after the per-blob commits and after the walk. A stop in between omits the recovered blob for good, or keeps stale .blame files. |
0 |
| 2 | Medium | Walker.scala 1355-1366, 1225-1226 | Strict mode plus --no-retry-timed-out (passed whenever blame exists) drops blobs and exits 0. The doc says skippedKeys is empty under strict. |
0 |
| 3 | Low | run_pipeline_process.sh 615-617, Walker.scala 1256 | The warning says the next run retries a timeout. After step 7 it does not, unless --retry-skipped. |
0 |
| 4 | Crop | Walker.scala 947-959, 1104-1108, 1412-1417, 1243, 1227, 1276 | Dead AtomicBoolean/AtomicReference wrappers, stale "callback" comments, dead "unknown" fallbacks, a never-read map value. |
14 |
| 5 | Crop | Walker.scala 1254-1268 | One SKIPPED blob log line takes 15 source lines. |
8 |
| 6 | Crop | Main.scala 478-509 | Four copies of the same numeric-option parse block; !v.isNaN does nothing. |
13 |
| 7 | Standards / crop | LoadGate.scala 35-46 | onPoll is test-only, and its doc says the opposite of the caller. |
3 |
| 8 | Crop | SkipFailedBlobSpec.scala 386-422 | Two identical tests, one subset of a loop test, repeated resolveStallWindow asserts. |
27 |
| 9 | Low / crop | test_tokenize_gate.sh 376-393 | Label says factor 0 disables the retry (false). Two pairs of copy checks. | 6 |
| 10 | Crop | README.md 220-341, Main.scala 126 and 256-262, run_pipeline_process.sh 168-175 | The worst-case figures (7820 s, 8420 s, "2 h 20 min", "30 min") are written in 4 places, and the README tells the re-fold and hold story twice. Keep one copy in the README. | about 25 |
| 11 | Crop | Walker.scala 1286-1306, 1368-1387 | The scaladocs of tokenize (21 lines) and retryTimedOutBlobs (20 lines) repeat the code branch by branch. About 8 lines each is enough. |
about 25 |
| 12 | Split candidate | commit 4d3eecb | The in-run retry is off by default (the decision to keep the 30-min watchdog). It needs LoadGate.scala (69), the retry loop in tokenize, retryBudget/longestBlobSeconds, resolveStallWindow, 4 options with help, the runner pass-through and validation, case 19, the spec section at 331-528, and the README retry text. Move it to a follow-up PR. Keep --no-retry-timed-out, --retry-skipped and the blame hold here: they protect blame from the cross-run retry of commit 1. |
about 600 |
Total estimated lines saved: about 120 in this PR (items 4-11). With the split (item 12), about 630, because the split also removes most of items 6-10.
Other notes, not blocking:
- Sharded mode ignores
CREGIT_MAX_RETRIES,CREGIT_TIMEOUT_RETRY_FACTORandCREGIT_LOAD_LIMITwith no warning (run_pipeline_process.sh 1395-1398). The PR body lists this as not tested. test_tokenize_gate.shfails 48 of 90 checks outsidedevenv shell. In the shell, all pass. This is the environment, not the PR.
| if (recovered > 0) { | ||
| mapping.clearFold() | ||
| refolded = true | ||
| blobsRecovered.add(recovered) |
There was a problem hiding this comment.
[High, spec] The re-fold state is lost if the run stops before the walk ends.
Each recovered blob commits putBlob + deleteRetry at once (lines 1447-1450). clearFold() runs only after the loop (here), and Main writes --refold-marker only after the whole walk (Main.scala 864-869). Two failure windows follow:
- The run stops after the first recovery commit and before
clearFold(). The next loop rows can each take a full budget, so this window is long. The next run finds noretry_blobrow and a fullcommit_map. It does not visit the old commits again. The recovered blob is then omitted from the tokenized repository for good, with no record. - The run stops during the re-fold walk (the watchdog, OOM, or a kill). The next run finds no
retry_blobrow, sorefoldedis false and no marker is written. The runner keepsblame/andblame-c100-incoming/, which name commits that no longer exist. The runner comment at run_pipeline_process.sh 1468-1471 says that step 7 then keeps those stale.blamefiles silently.
Fix: reset the fold and write the marker before the first blob row. The runner acts on the marker only after step 2 exits 0 (tokenize_gate stops it first), so an early marker is safe.
@@ retryTimedOutBlobs, case other @@
inserter.flush()
+ if (recovered == 0) {
+ // Before the first blob row, so that a stop at any later point
+ // still re-folds and still drops the old outputs.
+ onRefold() // Main writes --refold-marker here
+ mapping.clearFold()
+ refolded = true
+ }
mapping.inTx {
@@ after the loop @@
if (recovered > 0) {
- mapping.clearFold()
- refolded = true
blobsRecovered.add(recovered)
@@ Main.scala 864-869 @@
- // Written after the walk, so that it exists only if the re-fold is complete.
- if (stats.refolded) refoldMarker.foreach { p => ... }
+ // Pass to the Walker instead: onRefold = () => refoldMarker.foreach(writeMarker)Add a test that stops the walk after retryTimedOutBlobs (for example a stall) and checks that the next run keeps the blob and that the marker exists.
| pending.foreach { key => | ||
| skippedKeys.put(key, BlobExec.Failure(BlobExec.Failure.Timeout, "retry pass off")) | ||
| pendingRetry.add(key) | ||
| } |
There was a problem hiding this comment.
[Medium, spec] Strict mode exits 0 while held blobs stay dropped.
The runner passes --no-retry-timed-out whenever blame output exists, also with --strict-tokenize (run_pipeline_process.sh 1420-1433). Here the held keys go into skippedKeys, and no counter changes. So a strict run after a default run drops those blobs from every new tree and exits 0. The PR says that strict mode restores the old stop on exit 4 and 6, and report_skipped_blobs tells the operator to use strict mode "to require a clean run". The doc of skippedKeys (lines 1225-1226) also says "Empty under strictTokenize". That is false here and in retryTimedOutBlobs (lines 1418, 1435).
Count the held blobs in strict mode (or make the runner refuse strict mode together with a hold), and correct the doc:
| pending.foreach { key => | |
| skippedKeys.put(key, BlobExec.Failure(BlobExec.Failure.Timeout, "retry pass off")) | |
| pendingRetry.add(key) | |
| } | |
| pending.foreach { key => | |
| skippedKeys.put(key, BlobExec.Failure(BlobExec.Failure.Timeout, "retry pass off")) | |
| pendingRetry.add(key) | |
| } | |
| // Strict mode must not exit 0 while it keeps blobs dropped. | |
| if (strictTokenize) blobsTimedOut.add(pending.size.toLong) |
| log " repository, not as raw source and not as empty tokenizations. A timeout" | ||
| log " is tried again by the next step-2 run. To require a clean run, use" | ||
| log " --strict-tokenize (or CREGIT_STRICT_TOKENIZE=1)." |
There was a problem hiding this comment.
[Low, spec] The warning says that the next run retries a timeout. That is false after step 7.
After the first blame, every later step-2 run gets --no-retry-timed-out (line 1433). Thus the next run does not retry the timeout unless the operator gives --retry-skipped. The same false claim is in the SKIPPED blob line of Walker.scala 1256.
| log " repository, not as raw source and not as empty tokenizations. A timeout" | |
| log " is tried again by the next step-2 run. To require a clean run, use" | |
| log " --strict-tokenize (or CREGIT_STRICT_TOKENIZE=1)." | |
| log " repository, not as raw source and not as empty tokenizations. The next" | |
| log " step-2 run tries a timeout again only if no blame output exists, or with" | |
| log " --retry-skipped. To require a clean run, use --strict-tokenize (or CREGIT_STRICT_TOKENIZE=1)." |
| // Set by either callback below: both mean "this blob's tokenization is | ||
| // unusable, so nothing about it may be persisted". | ||
| val unusable = new AtomicBoolean(false) | ||
| val outcome = BlobExec.run( | ||
| bytes = bytes, | ||
| origSha = task.origId.name, | ||
| filename = task.filename, | ||
| fullPath = task.fullPath, | ||
| command = command, | ||
| abortOnError = abortOnError, | ||
| inserter = workerInserter, | ||
| timeoutSeconds = blobTimeoutSeconds, | ||
| onTimeout = () => { blobsTimedOut.increment(); unusable.set(true) }, | ||
| onParserCrash = () => { blobsParserCrashed.increment(); unusable.set(true) } | ||
| ) | ||
| val (outcome, failed) = tokenize(bytes, task, workerInserter) | ||
| val unusable = new AtomicBoolean(failed.isDefined) | ||
| val failure = new AtomicReference[BlobExec.Failure](failed.orNull) | ||
| val res = outcome match { | ||
| case BlobExec.Outcome.Skip if unusable.get() && !strictTokenize => | ||
| // Dropped, exactly as a denylisted blob is: no id, so no tree entry, no | ||
| // blob_map row and no dataset row. The original bytes are NOT put into | ||
| // dst, because nothing refers to them. | ||
| dropFailedBlob(task, failure.get()) | ||
| BlobResult.Oversized(task.origId) | ||
| case BlobExec.Outcome.Skip if unusable.get() => |
There was a problem hiding this comment.
[Crop, about 14 lines] Dead wrappers and stale comments in 3 places.
tokenize now returns failed: Option[Failure]. The code puts it in a new AtomicBoolean and a new AtomicReference, and reads them back on the same thread. No callback sets them. The comment "Set by either callback below" is no longer true. The same pattern is at lines 1104-1108 and 1412-1413. Also, dropFailedBlob is called only when failed.isDefined, so the "unknown" fallbacks at lines 1243 and 1416-1417 are dead. skippedKeys stores a Failure that no code reads (only containsKey/putIfAbsent), so a key set like oversizedKeys is enough. In noteTokenized (1276) the contains check repeats the filter that forget does.
| // Set by either callback below: both mean "this blob's tokenization is | |
| // unusable, so nothing about it may be persisted". | |
| val unusable = new AtomicBoolean(false) | |
| val outcome = BlobExec.run( | |
| bytes = bytes, | |
| origSha = task.origId.name, | |
| filename = task.filename, | |
| fullPath = task.fullPath, | |
| command = command, | |
| abortOnError = abortOnError, | |
| inserter = workerInserter, | |
| timeoutSeconds = blobTimeoutSeconds, | |
| onTimeout = () => { blobsTimedOut.increment(); unusable.set(true) }, | |
| onParserCrash = () => { blobsParserCrashed.increment(); unusable.set(true) } | |
| ) | |
| val (outcome, failed) = tokenize(bytes, task, workerInserter) | |
| val unusable = new AtomicBoolean(failed.isDefined) | |
| val failure = new AtomicReference[BlobExec.Failure](failed.orNull) | |
| val res = outcome match { | |
| case BlobExec.Outcome.Skip if unusable.get() && !strictTokenize => | |
| // Dropped, exactly as a denylisted blob is: no id, so no tree entry, no | |
| // blob_map row and no dataset row. The original bytes are NOT put into | |
| // dst, because nothing refers to them. | |
| dropFailedBlob(task, failure.get()) | |
| BlobResult.Oversized(task.origId) | |
| case BlobExec.Outcome.Skip if unusable.get() => | |
| val (outcome, failed) = tokenize(bytes, task, workerInserter) | |
| val res = outcome match { | |
| case BlobExec.Outcome.Skip if failed.isDefined && !strictTokenize => | |
| // Dropped, exactly as a denylisted blob is: no id, so no tree entry, no | |
| // blob_map row and no dataset row. The original bytes are NOT put into | |
| // dst, because nothing refers to them. | |
| dropFailedBlob(task, failed.get) | |
| BlobResult.Oversized(task.origId) | |
| case BlobExec.Outcome.Skip if failed.isDefined => |
The same change in the batch path:
- // Set by either callback below: both mean "this blob's tokenization is
- // unusable, so nothing about it may be persisted".
val (outcome, failed) = tokenize(bytes, task, workerInserter)
- val unusable = new AtomicBoolean(failed.isDefined)
- val failure = new AtomicReference[BlobExec.Failure](failed.orNull)
- val dropped = unusable.get() && !strictTokenize
+ val dropped = failed.isDefined && !strictTokenize
- case _ if dropped => dropFailedBlob(task, failure.get())
+ case _ if dropped => dropFailedBlob(task, failed.get)
- if (!unusable.get()) noteTokenized(task)
+ if (failed.isEmpty) noteTokenized(task)
- if (dropped) None else Some((task, outcome, unusable.get()))
+ if (dropped) None else Some((task, outcome, failed.isDefined))| val next = | ||
| if (f.reason == BlobExec.Failure.Timeout) | ||
| "A timeout depends on load, so the next run over this memo tries this blob again." | ||
| else | ||
| "A parser crash or empty output is deterministic for this tokenizer, so a re-run " + | ||
| "does not try it again " + | ||
| "unless the tokenizer identity changes (--retokenize)." | ||
| System.err.println( | ||
| s"blobExec: SKIPPED blob: sha=$sha path=${task.fullPath} reason=${f.reason} " + | ||
| s"detail=${f.detail}. The tokenizer output is not usable, so the blob is dropped from " + | ||
| "the rewritten tree, as a denylisted blob is: it gives no blame and no dataset row, " + | ||
| "and the file is not in the tokenized repository as raw source. The walk continues. " + | ||
| next + | ||
| skipLog.path.map(p => s" Recorded in $p.").getOrElse("") | ||
| ) |
There was a problem hiding this comment.
[Crop, about 8 lines] One log line takes 15 source lines.
The README and --help already explain the drop. The line must give the facts and the SKIPPED blob prefix. This also removes the false claim that the next run retries the blob (see the comment on run_pipeline_process.sh 615).
| val next = | |
| if (f.reason == BlobExec.Failure.Timeout) | |
| "A timeout depends on load, so the next run over this memo tries this blob again." | |
| else | |
| "A parser crash or empty output is deterministic for this tokenizer, so a re-run " + | |
| "does not try it again " + | |
| "unless the tokenizer identity changes (--retokenize)." | |
| System.err.println( | |
| s"blobExec: SKIPPED blob: sha=$sha path=${task.fullPath} reason=${f.reason} " + | |
| s"detail=${f.detail}. The tokenizer output is not usable, so the blob is dropped from " + | |
| "the rewritten tree, as a denylisted blob is: it gives no blame and no dataset row, " + | |
| "and the file is not in the tokenized repository as raw source. The walk continues. " + | |
| next + | |
| skipLog.path.map(p => s" Recorded in $p.").getOrElse("") | |
| ) | |
| val next = | |
| if (f.reason == BlobExec.Failure.Timeout) " A later run can try it again (retry_blob)." | |
| else " A re-run tries it again only if the tokenizer identity changes (--retokenize)." | |
| System.err.println( | |
| s"blobExec: SKIPPED blob: sha=$sha path=${task.fullPath} reason=${f.reason} " + | |
| s"detail=${f.detail}. It is dropped from the rewritten tree, as a denylisted blob is. " + | |
| "The walk continues." + next + skipLog.path.map(p => s" Recorded in $p.").getOrElse("")) |
| case t if t.startsWith("--max-retries=") => | ||
| val spec = t.stripPrefix("--max-retries=") | ||
| spec.toIntOption.filter(n => n >= 0 && n <= 100) match { | ||
| case Some(n) => maxRetries = n | ||
| case None => | ||
| System.err.println(s"Error: --max-retries must be a whole number from 0 to 100 [$spec]") | ||
| sys.exit(1) | ||
| } | ||
| case t if t.startsWith("--timeout-retry-factor=") => | ||
| val spec = t.stripPrefix("--timeout-retry-factor=") | ||
| spec.toIntOption.filter(_ >= 0) match { | ||
| case Some(n) => timeoutRetryFactor = n | ||
| case None => | ||
| System.err.println(s"Error: --timeout-retry-factor must be a whole number >= 0 [$spec]") | ||
| sys.exit(1) | ||
| } | ||
| case t if t.startsWith("--load-limit=") => | ||
| val spec = t.stripPrefix("--load-limit=") | ||
| spec.toDoubleOption.filter(v => v >= 0 && !v.isNaN) match { | ||
| case Some(v) => loadLimit = v | ||
| case None => | ||
| System.err.println(s"Error: --load-limit must be a number >= 0 [$spec]") | ||
| sys.exit(1) | ||
| } | ||
| case t if t.startsWith("--load-wait-max=") => | ||
| val spec = t.stripPrefix("--load-wait-max=") | ||
| spec.toIntOption.filter(_ >= 0) match { | ||
| case Some(n) => loadWaitMax = n | ||
| case None => | ||
| System.err.println(s"Error: --load-wait-max must be a whole number of seconds >= 0 [$spec]") | ||
| sys.exit(1) | ||
| } |
There was a problem hiding this comment.
[Crop, about 13 lines] Four copies of the same parse block.
Each numeric option repeats stripPrefix, toIntOption.filter, match, println and sys.exit(1). The file already has parsePositiveSeconds (line 113) for the same job. Also, !v.isNaN at line 496 does nothing: NaN >= 0 is already false. One helper keeps the error texts the same:
+ private def intFlag(t: String, prefix: String, lo: Int, hi: Int, what: String): Int = {
+ val spec = t.stripPrefix(prefix)
+ spec.toIntOption.filter(n => n >= lo && n <= hi).getOrElse {
+ System.err.println(s"Error: ${prefix.dropRight(1)} must be $what [$spec]"); sys.exit(1)
+ }
+ }
...
case t if t.startsWith("--max-retries=") =>
- (8 lines)
+ maxRetries = intFlag(t, "--max-retries=", 0, 100, "a whole number from 0 to 100")
case t if t.startsWith("--timeout-retry-factor=") =>
- (8 lines)
+ timeoutRetryFactor = intFlag(t, "--timeout-retry-factor=", 0, Int.MaxValue, "a whole number >= 0")
case t if t.startsWith("--load-limit=") =>
val spec = t.stripPrefix("--load-limit=")
- (7 lines)
+ loadLimit = spec.toDoubleOption.filter(_ >= 0).getOrElse {
+ System.err.println(s"Error: --load-limit must be a number >= 0 [$spec]"); sys.exit(1)
+ }
case t if t.startsWith("--load-wait-max=") =>
- (8 lines)
+ loadWaitMax = intFlag(t, "--load-wait-max=", 0, Int.MaxValue, "a whole number of seconds >= 0")| /** Wait while the load is above the limit. `onPoll` is called at each poll with | ||
| * the load, so the caller can stamp progress (the stall watchdog must not count | ||
| * a deliberate wait as a stall). Returns the seconds waited. */ | ||
| def await(onPoll: Double => Unit): Long = { |
There was a problem hiding this comment.
[Standards, crop about 3 lines] onPoll is test-only, and its doc says the opposite of what the caller does.
The doc says the caller can stamp progress during the wait. The only shipped caller passes _ => () (Walker.scala 1328), and the tokenize doc says that no progress is stamped during the waits, by design (Main sizes the stall window for it). Only SkipFailedBlobSpec 489-497 uses the parameter. Remove it, delete onPoll(load) at line 46, and let the tests check the returned seconds.
| /** Wait while the load is above the limit. `onPoll` is called at each poll with | |
| * the load, so the caller can stamp progress (the stall watchdog must not count | |
| * a deliberate wait as a stall). Returns the seconds waited. */ | |
| def await(onPoll: Double => Unit): Long = { | |
| /** Wait while the load is above the limit. No progress is stamped: Main makes | |
| * the stall window larger than the longest wait. Returns the seconds waited. */ | |
| def await(): Long = { |
| test("--max-retries=0 gives no retry") { | ||
| val fx = fixture("timeout-1") | ||
| val stats = run(fx, "serial", maxRetries = 0) | ||
| stats.blobsTimeoutRetriedInRun shouldEqual 0L | ||
| stats.blobsSkipped shouldEqual 1L | ||
| fx.callsForB shouldEqual 1 | ||
| tsvLines(fx).last.split("\t")(3) shouldEqual "timeout=1s" | ||
| } | ||
|
|
||
| test("the default is no retry, and the stall window then stays at 3 x --blob-timeout") { | ||
| Main.DefaultMaxRetries shouldEqual 0 | ||
| Main.Usage should include("default 0: a timeout skips and records the blob") | ||
| // With the default, a defaulted window stays at 1800s (30 min) for 600s blobs. | ||
| Main.resolveStallWindow(600, Main.DefaultMaxRetries, Main.DefaultTimeoutRetryFactor, | ||
| 600, 1800, stallExplicit = false) shouldEqual Right(1800) | ||
| // With an explicit --max-retries=3 it grows to 8420s (about 2 h 20 min). | ||
| Main.resolveStallWindow(600, 3, Main.DefaultTimeoutRetryFactor, | ||
| 600, 1800, stallExplicit = false) shouldEqual Right(8420) | ||
| } | ||
|
|
||
| test("with the default, a timeout skips and records the blob at once") { | ||
| val fx = fixture("timeout-1") | ||
| val stats = run(fx, "serial", maxRetries = Main.DefaultMaxRetries) | ||
| stats.blobsTimeoutRetriedInRun shouldEqual 0L | ||
| stats.blobsSkipped shouldEqual 1L | ||
| fx.callsForB shouldEqual 1 | ||
| tsvLines(fx).last.split("\t")(3) shouldEqual "timeout=1s" | ||
| } | ||
|
|
||
| test("an explicit --max-retries=3 still retries") { | ||
| val fx = fixture("timeout-1") | ||
| val stats = run(fx, "serial", maxRetries = 3) | ||
| stats.blobsTimeoutRetriedInRun shouldEqual 1L | ||
| stats.blobsTimeoutRecoveredInRun shouldEqual 1L | ||
| stats.blobsSkipped shouldEqual 0L | ||
| fx.callsForB shouldEqual 2 | ||
| } |
There was a problem hiding this comment.
[Crop, about 27 lines] Four tests repeat other tests.
- 386-393 and 406-413 have the same body.
maxRetries = 0andMain.DefaultMaxRetriesare the same value (line 396 asserts it). - The
resolveStallWindowasserts at 399-403 repeat lines 511 and 521. - 415-422 is the serial
(1, 1)case of the loop at 335-352, which asserts more.
One test keeps the default, the help text and the behaviour:
| test("--max-retries=0 gives no retry") { | |
| val fx = fixture("timeout-1") | |
| val stats = run(fx, "serial", maxRetries = 0) | |
| stats.blobsTimeoutRetriedInRun shouldEqual 0L | |
| stats.blobsSkipped shouldEqual 1L | |
| fx.callsForB shouldEqual 1 | |
| tsvLines(fx).last.split("\t")(3) shouldEqual "timeout=1s" | |
| } | |
| test("the default is no retry, and the stall window then stays at 3 x --blob-timeout") { | |
| Main.DefaultMaxRetries shouldEqual 0 | |
| Main.Usage should include("default 0: a timeout skips and records the blob") | |
| // With the default, a defaulted window stays at 1800s (30 min) for 600s blobs. | |
| Main.resolveStallWindow(600, Main.DefaultMaxRetries, Main.DefaultTimeoutRetryFactor, | |
| 600, 1800, stallExplicit = false) shouldEqual Right(1800) | |
| // With an explicit --max-retries=3 it grows to 8420s (about 2 h 20 min). | |
| Main.resolveStallWindow(600, 3, Main.DefaultTimeoutRetryFactor, | |
| 600, 1800, stallExplicit = false) shouldEqual Right(8420) | |
| } | |
| test("with the default, a timeout skips and records the blob at once") { | |
| val fx = fixture("timeout-1") | |
| val stats = run(fx, "serial", maxRetries = Main.DefaultMaxRetries) | |
| stats.blobsTimeoutRetriedInRun shouldEqual 0L | |
| stats.blobsSkipped shouldEqual 1L | |
| fx.callsForB shouldEqual 1 | |
| tsvLines(fx).last.split("\t")(3) shouldEqual "timeout=1s" | |
| } | |
| test("an explicit --max-retries=3 still retries") { | |
| val fx = fixture("timeout-1") | |
| val stats = run(fx, "serial", maxRetries = 3) | |
| stats.blobsTimeoutRetriedInRun shouldEqual 1L | |
| stats.blobsTimeoutRecoveredInRun shouldEqual 1L | |
| stats.blobsSkipped shouldEqual 0L | |
| fx.callsForB shouldEqual 2 | |
| } | |
| test("the default is no retry: a timeout skips and records the blob at once") { | |
| Main.DefaultMaxRetries shouldEqual 0 | |
| Main.Usage should include("default 0: a timeout skips and records the blob") | |
| val fx = fixture("timeout-1") | |
| val stats = run(fx, "serial", maxRetries = Main.DefaultMaxRetries) | |
| stats.blobsTimeoutRetriedInRun shouldEqual 0L | |
| stats.blobsSkipped shouldEqual 1L | |
| fx.callsForB shouldEqual 1 | |
| tsvLines(fx).last.split("\t")(3) shouldEqual "timeout=1s" | |
| } |
| grep -q -- "--timeout-retry-factor=0" "$STUB_ARGV"; check "the factor is passed through (0 disables the retry)" $? | ||
| grep -q -- "--load-limit=8" "$STUB_ARGV"; check "the load limit is passed through" $? | ||
| rm -f "$STUB_ARGV" | ||
| OUT=$(STUB_RC=0 CREGIT_TIMEOUT_RETRY_FACTOR=5 run_step2 "$W") | ||
| grep -q -- "--timeout-retry-factor=5" "$STUB_ARGV"; check "CREGIT_TIMEOUT_RETRY_FACTOR is honoured" $? | ||
| rm -f "$STUB_ARGV" | ||
| OUT=$(STUB_RC=0 run_step2 "$W" --max-retries 0) | ||
| grep -q -- "--max-retries=0" "$STUB_ARGV"; check "--max-retries is passed through" $? | ||
| rm -f "$STUB_ARGV" | ||
| OUT=$(STUB_RC=0 run_step2 "$W" --max-retries 3) | ||
| grep -q -- "--max-retries=3" "$STUB_ARGV"; check "an explicit --max-retries 3 is passed through" $? | ||
| rm -f "$STUB_ARGV" | ||
| OUT=$(STUB_RC=0 CREGIT_MAX_RETRIES=3 run_step2 "$W") | ||
| grep -q -- "--max-retries=3" "$STUB_ARGV"; check "CREGIT_MAX_RETRIES=3 is passed through" $? | ||
| rm -f "$STUB_ARGV" | ||
| OUT=$(STUB_RC=0 CREGIT_MAX_RETRIES=5 run_step2 "$W") | ||
| grep -q -- "--max-retries=5" "$STUB_ARGV"; check "CREGIT_MAX_RETRIES is honoured" $? | ||
| rm -f "$STUB_ARGV" |
There was a problem hiding this comment.
[Low, crop about 6 lines] A wrong test label, and two pairs of copies.
The label at line 376 says "0 disables the retry". blobExec says the opposite: a factor of 0 or 1 gives the first budget (Walker.retryBudget, Main.scala 248, README 172). The retries still run. Also, 382-387 test the flag pass-through twice, and 388-393 test the CREGIT_MAX_RETRIES pass-through twice. Line 395 already covers the default of 0 (no flag is passed).
| grep -q -- "--timeout-retry-factor=0" "$STUB_ARGV"; check "the factor is passed through (0 disables the retry)" $? | |
| grep -q -- "--load-limit=8" "$STUB_ARGV"; check "the load limit is passed through" $? | |
| rm -f "$STUB_ARGV" | |
| OUT=$(STUB_RC=0 CREGIT_TIMEOUT_RETRY_FACTOR=5 run_step2 "$W") | |
| grep -q -- "--timeout-retry-factor=5" "$STUB_ARGV"; check "CREGIT_TIMEOUT_RETRY_FACTOR is honoured" $? | |
| rm -f "$STUB_ARGV" | |
| OUT=$(STUB_RC=0 run_step2 "$W" --max-retries 0) | |
| grep -q -- "--max-retries=0" "$STUB_ARGV"; check "--max-retries is passed through" $? | |
| rm -f "$STUB_ARGV" | |
| OUT=$(STUB_RC=0 run_step2 "$W" --max-retries 3) | |
| grep -q -- "--max-retries=3" "$STUB_ARGV"; check "an explicit --max-retries 3 is passed through" $? | |
| rm -f "$STUB_ARGV" | |
| OUT=$(STUB_RC=0 CREGIT_MAX_RETRIES=3 run_step2 "$W") | |
| grep -q -- "--max-retries=3" "$STUB_ARGV"; check "CREGIT_MAX_RETRIES=3 is passed through" $? | |
| rm -f "$STUB_ARGV" | |
| OUT=$(STUB_RC=0 CREGIT_MAX_RETRIES=5 run_step2 "$W") | |
| grep -q -- "--max-retries=5" "$STUB_ARGV"; check "CREGIT_MAX_RETRIES is honoured" $? | |
| rm -f "$STUB_ARGV" | |
| grep -q -- "--timeout-retry-factor=0" "$STUB_ARGV"; check "the factor is passed through (0 = the first budget)" $? | |
| grep -q -- "--load-limit=8" "$STUB_ARGV"; check "the load limit is passed through" $? | |
| rm -f "$STUB_ARGV" | |
| OUT=$(STUB_RC=0 CREGIT_TIMEOUT_RETRY_FACTOR=5 run_step2 "$W") | |
| grep -q -- "--timeout-retry-factor=5" "$STUB_ARGV"; check "CREGIT_TIMEOUT_RETRY_FACTOR is honoured" $? | |
| rm -f "$STUB_ARGV" | |
| OUT=$(STUB_RC=0 run_step2 "$W" --max-retries 3) | |
| grep -q -- "--max-retries=3" "$STUB_ARGV"; check "an explicit --max-retries 3 is passed through" $? | |
| rm -f "$STUB_ARGV" | |
| OUT=$(STUB_RC=0 CREGIT_MAX_RETRIES=5 run_step2 "$W") | |
| grep -q -- "--max-retries=5" "$STUB_ARGV"; check "CREGIT_MAX_RETRIES is honoured" $? | |
| rm -f "$STUB_ARGV" |
tokenize returns the failure, so the AtomicBoolean and AtomicReference that wrapped it on the same thread did nothing, and the "set by either callback" comment was stale. dropFailedBlob is called only with a failure, so its "unknown" fallbacks were dead. No code read the failure stored in skippedKeys, so it is a key set, as oversizedKeys is. And SkipLog.forget already checks for the row, so noteTokenized need not. Review finding 4 on #96.
The retry pass committed putBlob and deleteRetry for each recovered blob, but emptied commit_map, tree_map and ref_map only after the loop, and Main wrote --refold-marker only after the whole walk. A run that stopped in between lost the recovered blob for good (no retry row, a full commit_map), or kept blame files that name commits that no longer exist (no marker). Now the first recovery writes the marker (through a new onRefold callback), then empties the fold, and only then writes its blob row. A stop at any later point still re-folds on the next run, and the runner still drops the old outputs. The runner acts on the marker only after step 2 exits 0, so an early marker is safe. The retry of one row moves to retryOne, so that the loop reads as the three steps it is. Review finding 1 on #96.
The runner passes --no-retry-timed-out whenever blame output exists, also with --strict-tokenize. The held blobs then stayed dropped from every new tree, no counter changed, and a strict run exited 0. Strict mode promises a clean run or a non-zero exit. Now a strict run counts each held blob as a timeout, so it exits 4, and the hold line says that --retry-skipped tries them. The skippedKeys doc no longer claims the set is empty under strict mode. Review finding 2 on #96.
After the first blame, every later step-2 run gets --no-retry-timed-out, so "the next run tries it again" was false unless the operator passes --retry-skipped. The runner warning, its help text, the blobExec warning and the SKIPPED blob line now say when a retry happens. The SKIPPED blob line also shrinks from 15 source lines to 7. It keeps the facts and the prefix; the README and --help explain the drop. Review findings 3 and 5 on #96.
--max-retries, --timeout-retry-factor, --load-limit and --load-wait-max each repeated the same strip, parse, filter, print and exit block. flagValue does it once and keeps the error texts. The !v.isNaN check went too: NaN >= 0 is already false. Review finding 6 on #96.
The only shipped caller passed a no-op, and the doc said the caller could stamp progress during the wait, the opposite of what Walker does: Main sizes the stall window for the wait instead. The test now checks the seconds that await returns. Review finding 7 on #96.
SkipFailedBlobSpec had two tests with the same body (maxRetries 0 is the default), resolveStallWindow asserts that the stall-window test already makes, and a test that is the serial (1, 1) case of the recovery loop. One test now keeps the default, the help text and the behaviour. In test_tokenize_gate.sh the label said a factor of 0 disables the retry; it gives the first budget. Two pairs of pass-through checks tested the same thing twice. Review findings 8 and 9 on #96.
The worst-case figures (7820 s, 8420 s, about 2 h 20 min, 30 min) were written in the README, the blobExec help, a Main scaladoc and the runner help. Four copies drift. The README keeps them; the help texts say what blobExec does with the stall window and point there. The exit-status table also told the hold story a second time; the rules above it tell it once. Review finding 10 on #96.
Scaladocs that walked the code branch by branch (tokenize, the retry pass, resolveStallWindow, SkipLog, the WalkStats fields) are now at most three lines, each on a why, a trap or a reference. Comments that restated the next line went, in the code and in the specs. Three comments were stale: the Walker parameter note named Main defaults of 3 retries, the runner said the default load limit is the processor count, and BlobExec said the counters still use the onTimeout and onParserCrash callbacks. SkipLog.contains lost its last caller in the dead-wrapper cleanup, so it goes too. Review finding 11 on #96.
The five checks of CREGIT_STRICT_TOKENIZE, CREGIT_RETRY_SKIPPED, --max-retries, CREGIT_LOAD_LIMIT and --timeout-retry-factor each spelled out the same echo and exit, three of them inside an if around a case. refuse_value prints the same messages, and each check is now one flat case.
The pool took one timeout in its constructor, and Walker ran the pool before BlobExec.run and handed it the finished result. A caller that retries a timed-out blob with a larger budget (#96) could not use the pool: the budget was fixed, and a finished result cannot be re-run. The REQ header already carries a timeout per request. invoke now takes timeoutSeconds, and BlobExec.run takes the pool and calls it with its own timeoutSeconds, as it calls ChildRunner. Walker only passes the pool through. A new test checks that each request sends its own budget.
|
Replaced by a shorter version with no new flags. |
A srcML crash, an empty output or a timeout on one blob stopped step 2 for the whole project. With this PR, blobExec drops that blob, records it, and the run continues.
This PR targets
deploy/reblame-c100, the code the-C100re-blame runs use. It is not based onmaster.Commits
<work>/tokenize-skipped.tsv(sha path reason detail tokenizer), and step 2 prints a count for each reason. A timed-out blob gets aretry_blobrow, and the next run tries it first.--strict-tokenize(orCREGIT_STRICT_TOKENIZE=1) restores the old stop on exit 4 and 6. Sharded mode is always strict.--max-retries N(orCREGIT_MAX_RETRIES) gives a timed-out blob up to N more attempts in the same run. Each retry waits for the load to drop and gets 3x the time budget. The default is 0, so the stall watchdog window stays at 30 min. With 3 retries the window grows to about 2 h 20 min.Retry and blame
When a blob from an earlier run recovers, the run folds all of history again. All commit ids change, so the runner deletes
blame/andblame-c100-incoming/, and the project needs a full re-blame. To prevent that, step 2 does not retry earlier timeouts while a.blamefile exists. It prints how many blobs stay dropped and how many.blamefiles a retry would delete.--retry-skipped(orCREGIT_RETRY_SKIPPED=1) forces the retry.Tests
sbt -batch testtest_tokenize_gate.shtest_ensure_artifacts.shtest_retokenize_passthrough.shtest_reblame_passthrough.shtest_run_pipeline_workdir_guard.shprove(tokenize, tokenizeByBlobId, blameRepo, prettyPrint)pytest generate_datasetNot tested
ctp.py.