Skip to content

Add a build/write-prefetch unparse path driven by a dedicated Builder tree - #1736

Open
olabusayoT wants to merge 3 commits into
apache:mainfrom
olabusayoT:daf-3065-build-cursor
Open

olabusayoT wants to merge 3 commits into
apache:mainfrom
olabusayoT:daf-3065-build-cursor

Conversation

@olabusayoT

@olabusayoT olabusayoT commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

Adds an off-by-default (useBuildWritePrefetch) two-pass unparse path in
which build races ahead of write over the same infoset tree, resolving
OVCs directly against the tree instead of through Suspensions where
possible.

Build is driven by a Builder tree paralleling the Unparser tree and wired
through the grammar combinators, with per-parse BuildFrames on an
explicit BuildCursor stack so building can stop after any step and resume
later. Write pulls build forward through awaitChild and
childExistsOrFinal until unparsePrefetchWindowNodes are built ahead or
pending suspensions exceed unparsePendingSuspensionTripLimit, all on the
calling thread.

The Builder tree is only constructed when the tunable is on and
hasAnyPrefetchBeneficialOVC holds, which is scoped per compiling root
and requires a non-constant OVC.

writeContent and unparse share their setup, dispatch and teardown through
per-combinator run methods, WriteUnparser.dispatchBody and
runElementContent, so SpecifiedLengthPrefixedUnparser now also resolves
its prefix length when the body throws during unparse.

DaffodilTunables.withTunable is generated as withTunablePart0..N to stay
under the JVM method size limit.

DAFFODIL-3065

Design: https://cwiki.apache.org/confluence/spaces/DAFFODIL/pages/451975094/Build+Write-Prefetch+Unparse+Path

… tree

Adds an off-by-default (useBuildWritePrefetch) two-pass unparse path in
which build races ahead of write over the same infoset tree, resolving
OVCs directly against the tree instead of through Suspensions where
possible.

Build is driven by a Builder tree paralleling the Unparser tree and wired
through the grammar combinators, with per-parse BuildFrames on an
explicit BuildCursor stack so building can stop after any step and resume
later. Write pulls build forward through awaitChild and
childExistsOrFinal until unparsePrefetchWindowNodes are built ahead or
pending suspensions exceed unparsePendingSuspensionTripLimit, all on the
calling thread.

The Builder tree is only constructed when the tunable is on and
hasAnyPrefetchBeneficialOVC holds, which is scoped per compiling root
and requires a non-constant OVC.

writeContent and unparse share their setup, dispatch and teardown through
per-combinator run methods, WriteUnparser.dispatchBody and
runElementContent, so SpecifiedLengthPrefixedUnparser now also resolves
its prefix length when the body throws during unparse.

DaffodilTunables.withTunable is generated as withTunablePart0..N to stay
under the JVM method size limit.

DAFFODIL-3065
Cache each element's ChoiceBranchStartEvent/EndEvent on its
ElementRuntimeData instead of looking them up through a shared interning
cache, whose read-write lock was contended with many threads unparsing.

Pass the ElementUnparserBase to ElementBuilder instead of two closures per
element, create BuildState's no-op output stream lazily, and drop the
child index stack from the build side.

Make the TypedEquality operators, MStack's element operations, and
withRetryIfBlocking inline so each call site uses static types and typed
arrays and allocates no closure. Test Maybes directly in element content
dispatch instead of converting them to Options.

Build/write prefetch is the default for every schema; the tunable
description and tests are updated to match.

DAFFODIL-3065

@stevedlawrence stevedlawrence left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

checkpoint comments

* explicit stack, so building can stop after any step and continue later
* without holding a thread or a JVM call stack.
*/
trait Builder extends Serializable {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we rename these classes to make it more clear they are about building our internal infoset? Otherwise being in the unparser namespace one might assume they are about building unparsers or something. For example, maybe this wants to be InfosetBuilder.

Might additionally want to move it to the infoset package as well so it lives along side the InfosetInputter stuff sine the two are very closely related.

* the Unparser tree. Most Grams create or select no infoset content and
* inherit this Nope default; only those that do override it.
*/
def builder: Maybe[Builder] = Maybe.Nope

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should we follow the patter of Parsers and Unparsers and change this to return a Builder instead of a Maybe[Builder]. Primitives that do not return a builder can return a NadaBuilder (similar to NadaParser and NadaUnparser, which should get optimized out?

I'm honestly not sure if there is a reason for NadaParser/Unaprsers (maybe to avoid a bunch of One/Some's, but maybe there's value in consistency?

It might be worth considering if there's value in removing NadaParser/Unprsers and changing it all to Maybe.s It will then at least force primitives to consider how to handle parsers. Right now they don't necessarily have to do that and we just throw an abort it a NadaParser/Unparse ends up not being optimized out.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added DAFFODIL-3103 to track replacement of Nada* objects with Maybe[Object], for now will replace Builder with NadaBuilder

eValue.builder
}
}
private lazy val eReptypeBuilder: Maybe[Builder] = repTypeElementGram.builder

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggest we name this eRepTypeBuilder (with a capital T) to match the reptype parser/unparser converntion.


// Shares the memoized unparser above for unparseBegin/unparseEnd, so
// build and write see identical node-creation behavior.
override lazy val builder: Maybe[Builder] = unparser match {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This might read a bit better if we change the logic in val unparser to something like

private lazy val isSomethingAboutSpecifiedLength = ...

val unparser = 
  if (context.isOutputValueCalc) {
    new ElementOCVSpecifiedLengthUnparser(...)
  } else if (isSomethingAboutSpecifiedLength) {
    new ElementSpecifiedLengthUnparser(...)
  } else {
    subComb.unparser
  }

And then this builder can have very similar same conditionals just with Builders instead of Unparsers. eg.

val builder =
  if (context.isOUtputValueCalc || isSomethignAboutSpecifiedLength) {
    ElementBuidler(...)
  } else {
    subComb.builder
  }

That way it's very clear that the logic is basically the same rather than relying on type matching on parsers.

// Shares the memoized unparser above for unparseBegin/unparseEnd, so
// build and write see identical nilled/OVC/IVC node-creation behavior.
override lazy val builder: Maybe[Builder] = {
val eu = unparser.asInstanceOf[ElementUnparserBase]

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Same here, I would rather see the same logic rather than type matching. It's much easier to visually verify the conditionalsare the same rather than figuring out which Unparsers are created.

writeState.getDataOutputStream.setFinished(writeState)

(singlePassBytes, walkerOut.toByteArray)
}

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

These tests seem like they are using too much of our internals for a generic TestUtil function. This feels much more like unit tests specific to testing that prefetching is working as expected. I'd prefer if we dont' invent another test infrasture for this (I found AI really likes to create it's own test frameworks) so if there's an alternative (like TDML tests) I'd prefer that. But if we really do need this, suggest it becomes it's own file specific to testing prefetch.

* Off by default: growing past initialSize isn't itself wrong, so paying
* this bookkeeping cost on every push isn't worth it normally.
*/
final val trackMaxSizeReached: Boolean = false

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can you split mstack changes to a new PR?

Note that this is a prime example to avoid combine unrelated changes in PR, especially when they are complex. I believe I looked at the mstack changes as part of another PR and thought it looked reasonable, but we ended up closing that PR and now it needs to be reviewed again.


/**
* True if a component reachable from this element has a
* dfdl:outputValueCalc resolvable without writing; gates

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Again, let's avoid the term "writing", it's not a term we've traditionally used except for actually writing bits and could cause confusion.

// Not forced eagerly: onPath only references this when
// tunable.useBuildWritePrefetch is on (the tunable is fixed at compile
// time), so schemas that never enable it never pay to construct the
// Builder tree.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Doe sthis increase compilation time very much? If not, I might suggest we always create the builder, regardless ofthe tunable. That way the tunable could be dynamically changed at unparse time without needing to recompile. I'm not sure that would actually be used in practice, but it could be useful for performance testing. E.g. compile once, run once with tunable true, run once with tunable false, and then you have a good idea if a certain value is significantly better. We can still compile a DataProcessor with a specific value, but it means it can be overridden.


final def freeChildIfNoLongerNeeded(index: Int, doFree: Boolean): Unit = {
val node = _contents(index)
// A null slot means write already freed it before build's redundant

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

s/write/unparse/

Also, I'm not sure I understand this comment or why this change is needed. Seems like the builder should never call this function since all it's doing is creating nodes, and skimming the code it's not immediately obvious that it ever does call this. Seems only the unparse() function should free these nodes after it has unparsed them. Same with the below change in this file.

@stevedlawrence stevedlawrence left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

checkpoint comments

* against a `BuildState`, never a write-side `UState`.
*/
final class BuildCursor(root: Builder, val state: UState, ctx: UnparseSharedContext) {
private var stack = new Array[BuildFrame](32)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Does this want to be an MStack? Seems like this does things very similar to what an MStack is used for and some of the things like depth/push/pop/etc you get for free.

final class BuildCursor(root: Builder, val state: UState, ctx: UnparseSharedContext) {
private var stack = new Array[BuildFrame](32)
private var depth = 0
private var failure: Throwable = null

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggest we drop this for simplicity? I assume once we throw a BuildAbortedException we'll immediately turn that into an UnparseError or something and terminate the unparse so there shouldn't be a concern about calling back into this.

*/
def advance(): Unit = {
if (failure != null) {
throw new BuildAbortedException(failure)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If we do keep the failure var, does this want to be an assert? It feels like calling the builder after it's thrown a exception should be considered a usage error.

* The lead only rises as build adds nodes and only falls as write
* finishes them, so this is the peak over a whole run.
*/
def peakLead: Long = peakLead_

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This doesn't seem to be used, is this for diagnostics?

lead = newLead
if (
ctx.leadExceedsPrefetchLimit ||
ctx.suspensionTracker.pendingCount > ctx.pendingSuspensionTripLimit

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Doesn't this mean that if we have a lot of suspensions (more than the trip limit) that we stop trying to build the infoset, and from that point on we'll only build elements one at a time until the suspensions go down? Could switching to one at a time building going to cause performance issues?

I feel like we actually want the opposite behavior: if we have a bunch of suspensions, then maybe we consider increasing the prefetch limit hoping to read some more elements that unblock some of those existing suspensions. I'm not sure we really need that since unparse() will continue to make progress and eventually get there. Also the default value of 500 feels pretty high, I wonder if we're never hitting hitting that limit and so aren't finding out what happens when this is hit?

If we aren't actually hitting this, I might suggest we remove the suspensionTripLimit. A format with lots of suspensions could always increase the prefetch limit which should have a similar effect.

private final val NextChild = 0
private final val AfterScalar = 1
private final val InArray = 2
private final val AfterOccurrence = 3

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Might put this in an object SequenceBuilder to avoid allocating these vals for every new frame. It's only 4 vals so probably not a big deal, but the less memory we allocate the better.

* `getDataOutputStream` is NOT stubbed: generic `UState` utility methods
* (toString, currentLocation, bitPos0b) call into it, so `BuildState`
* lazily constructs an actual DOS wrapping a no-op sink purely to satisfy
* that.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would think if something is calling getDataOutputStream on a BuildState that's a usage error and it should assert. Things like bitPosotion/currentLocation/etc are pretty meaningless in a BuildState so if someting has access to a BuildState and queries those it's probably a bug.

final class BuildState(
private val inputter: InfosetInputter,
sharedCtx: UnparseSharedContext,
diagnosticsArg: Seq[api.Diagnostic],

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Our usages always pass Nil to the diagnosticsArg. I woudl assume any diagnostics that a Builder/BuildState create would bubbble up to whatever called build and add the diagnostics to the real Ustate?

sharedCtx.tunable,
areDebugging
)
with SuspensionCapableUState

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is this really suspension capable? I wouldn't think the builder could create suspensios since it doesn't ever call unparse() which I think is the only thing that ever creates suspensions?

}

final override def documentElement: DIDocument = inputter.documentElement
}

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Feels like we might want to rethink this BuildState. Wonder if it wants to be it's own thing that UState mixes in, and the functions that work on both a BuildState and a UState just accept a BuildState.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants