Skip to content

getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write #2058

Description

@rpinckne

Summary

Calling getWritable() more than once within a single step body produces non-deterministic chunk ordering on Vercel prod. The API surface gives no hint that re-acquisition is unsafe.

Where the race lives

@workflow/core/dist/step/writable-stream.js:16-39 (workflow@5.0.0-beta.6):

exportfunctiongetWritable(options={}){// ...constserialize=getSerializeStream(getExternalReducers(...),ctx.encryptionKey);constserverWritable=newWorkflowServerWritableStream(runId,name);conststate=createFlushableState();ctx.ops.push(state.promise);flushablePipe(serialize.readable,serverWritable,state).catch(()=>{});pollWritableLock(serialize.writable,state);returnserialize.writable;}

Every invocation constructs a freshTransformStream + WorkflowServerWritableStream + an independent background flushablePipe(...). The (runId, name) is shared, but the pipes are independent client-side and race to flush into the same downstream queue.

Repro

Sequential per-chunk acquisition in a step body, no parallelism:

'use step';exportasyncfunctionemitStream(parts: UIMessageChunk[]){for(constpartofparts){constwriter=getWritable<UIMessageChunk>().getWriter();try{awaitwriter.write(part);}finally{writer.releaseLock();}}}

Emit ~6 small chunks (e.g. text-deltas "nov", "o", " e", "2", "e", " ok"). Read the resulting stream via workflowRun.getReadable<UIMessageChunk>(...).

  • Local dev (world-local): deltas arrive in order. Stream reads "novo e2e ok".
  • Vercel prod (world-vercel): deltas reorder deterministically. Stream reads e.g. "novo2e e ok", with text-delta arriving before text-start.

We verified the model output is correct upstream — direct streamText() and direct engine consumption both produce the right deltas in the right order. The reorder happens entirely between the per-chunk getWritable() calls and the downstream reader.

Mechanism

Documented in your own docs/changelog/eager-processing.mdx:330-331:

On local (world-local), stream writes go to the filesystem — effectively instant. On Vercel (world-vercel), writes go through HTTP to workflow-server → S3, adding 50-100ms latency.

That latency turns the per-pipe race window from microseconds (locally invisible) into tens of milliseconds (prod-observable). Small/fast deltas — exactly the shape of reasoning-model streaming output — surface the bug most reliably.

Canonical pattern (works)

docs/api-reference/workflow-ai/durable-agent.mdx:46 shows the right shape:

constwritable=getWritable<UIMessageChunk>();constresult=awaitagent.stream({ messages, writable });

One call. pipeTo(writable) (or the equivalent inside DurableAgent) holds the writer lock for the pipe's lifetime, so there's one TransformStream + one flushablePipe and the race goes away. We fixed our consumer by switching to this pattern + the AI SDK's toUIMessageStream().pipeTo(writable, { preventClose: true, preventAbort: true }).

Suggestion

The API surface gives no hint that re-acquisition is unsafe. Two options that would have caught this at dev time without needing prod telemetry:

  1. Dev-mode console.warn when getWritable() is called more than once per step context. contextStorage already tracks per-step state — a simple counter on the ctx would flag the misuse the second time it happens.

  2. Memoize per (runId, namespace): repeat calls within the same step ctx return the same handle (same underlying TransformStream, same flushablePipe). Idempotent, no behavior change for the canonical pattern, makes the unsafe pattern simply correct instead of subtly broken.

Either would have caught the bug in our first local test run instead of letting it ship to prod.

Context

We're building Novo Agents (hosted multi-agent service) on Workflow. The full fix on our side is in rpinckne/harness#50. Happy to help with a repro repo or PR if useful.

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions

      , 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
       blocks
      (function() {
      function addCopyButtons() {
      document.querySelectorAll('pre code').forEach(function(codeBlock) {
      if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
      codeBlock.parentElement.setAttribute('data-copy-added', 'true');
      var btn = document.createElement('button');
      btn.textContent = 'Copy';
      btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
      btn.onmouseover = function() { this.style.opacity = '1'; };
      btn.onmouseout = function() { this.style.opacity = '0.7'; };
      btn.onclick = function() {
      navigator.clipboard.writeText(codeBlock.textContent).then(function() {
      btn.textContent = 'Copied!';
      setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
      });
      };
      codeBlock.parentElement.style.position = 'relative';
      codeBlock.parentElement.appendChild(btn);
      });
      }
      addCopyButtons();
      // Re-run on dynamic content
      var observer = new MutationObserver(addCopyButtons);
      observer.observe(document.body, { childList: true, subtree: true });
      })();
      }
      } catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
      })();
      (function(){
      try {
      var __m = "github.com";
      var __re = new RegExp('^' + "github\\.com" + '
      getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write · Issue #2058 · vercel/workflow · GitHub
      Skip to content

      getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write #2058

      Description

      @rpinckne

      Summary

      Calling getWritable() more than once within a single step body produces non-deterministic chunk ordering on Vercel prod. The API surface gives no hint that re-acquisition is unsafe.

      Where the race lives

      @workflow/core/dist/step/writable-stream.js:16-39 (workflow@5.0.0-beta.6):

      exportfunctiongetWritable(options={}){// ...constserialize=getSerializeStream(getExternalReducers(...),ctx.encryptionKey);constserverWritable=newWorkflowServerWritableStream(runId,name);conststate=createFlushableState();ctx.ops.push(state.promise);flushablePipe(serialize.readable,serverWritable,state).catch(()=>{});pollWritableLock(serialize.writable,state);returnserialize.writable;}

      Every invocation constructs a freshTransformStream + WorkflowServerWritableStream + an independent background flushablePipe(...). The (runId, name) is shared, but the pipes are independent client-side and race to flush into the same downstream queue.

      Repro

      Sequential per-chunk acquisition in a step body, no parallelism:

      'use step';exportasyncfunctionemitStream(parts: UIMessageChunk[]){for(constpartofparts){constwriter=getWritable<UIMessageChunk>().getWriter();try{awaitwriter.write(part);}finally{writer.releaseLock();}}}

      Emit ~6 small chunks (e.g. text-deltas "nov", "o", " e", "2", "e", " ok"). Read the resulting stream via workflowRun.getReadable<UIMessageChunk>(...).

      • Local dev (world-local): deltas arrive in order. Stream reads "novo e2e ok".
      • Vercel prod (world-vercel): deltas reorder deterministically. Stream reads e.g. "novo2e e ok", with text-delta arriving before text-start.

      We verified the model output is correct upstream — direct streamText() and direct engine consumption both produce the right deltas in the right order. The reorder happens entirely between the per-chunk getWritable() calls and the downstream reader.

      Mechanism

      Documented in your own docs/changelog/eager-processing.mdx:330-331:

      On local (world-local), stream writes go to the filesystem — effectively instant. On Vercel (world-vercel), writes go through HTTP to workflow-server → S3, adding 50-100ms latency.

      That latency turns the per-pipe race window from microseconds (locally invisible) into tens of milliseconds (prod-observable). Small/fast deltas — exactly the shape of reasoning-model streaming output — surface the bug most reliably.

      Canonical pattern (works)

      docs/api-reference/workflow-ai/durable-agent.mdx:46 shows the right shape:

      constwritable=getWritable<UIMessageChunk>();constresult=awaitagent.stream({ messages, writable });

      One call. pipeTo(writable) (or the equivalent inside DurableAgent) holds the writer lock for the pipe's lifetime, so there's one TransformStream + one flushablePipe and the race goes away. We fixed our consumer by switching to this pattern + the AI SDK's toUIMessageStream().pipeTo(writable, { preventClose: true, preventAbort: true }).

      Suggestion

      The API surface gives no hint that re-acquisition is unsafe. Two options that would have caught this at dev time without needing prod telemetry:

      1. Dev-mode console.warn when getWritable() is called more than once per step context. contextStorage already tracks per-step state — a simple counter on the ctx would flag the misuse the second time it happens.

      2. Memoize per (runId, namespace): repeat calls within the same step ctx return the same handle (same underlying TransformStream, same flushablePipe). Idempotent, no behavior change for the canonical pattern, makes the unsafe pattern simply correct instead of subtly broken.

      Either would have caught the bug in our first local test run instead of letting it ship to prod.

      Context

      We're building Novo Agents (hosted multi-agent service) on Workflow. The full fix on our side is in rpinckne/harness#50. Happy to help with a repro repo or PR if useful.

      Metadata

      Metadata

      Assignees

      No one assigned

        Labels

        No labels
        No labels

        Type

        No type

        Projects

        No projects

          Milestone

          No milestone

          Relationships

          None yet

          Development

          No branches or pull requests

          Issue actions

          , 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write · Issue #2058 · vercel/workflow · GitHub
          Skip to content

          getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write #2058

          Description

          @rpinckne

          Summary

          Calling getWritable() more than once within a single step body produces non-deterministic chunk ordering on Vercel prod. The API surface gives no hint that re-acquisition is unsafe.

          Where the race lives

          @workflow/core/dist/step/writable-stream.js:16-39 (workflow@5.0.0-beta.6):

          exportfunctiongetWritable(options={}){// ...constserialize=getSerializeStream(getExternalReducers(...),ctx.encryptionKey);constserverWritable=newWorkflowServerWritableStream(runId,name);conststate=createFlushableState();ctx.ops.push(state.promise);flushablePipe(serialize.readable,serverWritable,state).catch(()=>{});pollWritableLock(serialize.writable,state);returnserialize.writable;}

          Every invocation constructs a freshTransformStream + WorkflowServerWritableStream + an independent background flushablePipe(...). The (runId, name) is shared, but the pipes are independent client-side and race to flush into the same downstream queue.

          Repro

          Sequential per-chunk acquisition in a step body, no parallelism:

          'use step';exportasyncfunctionemitStream(parts: UIMessageChunk[]){for(constpartofparts){constwriter=getWritable<UIMessageChunk>().getWriter();try{awaitwriter.write(part);}finally{writer.releaseLock();}}}

          Emit ~6 small chunks (e.g. text-deltas "nov", "o", " e", "2", "e", " ok"). Read the resulting stream via workflowRun.getReadable<UIMessageChunk>(...).

          • Local dev (world-local): deltas arrive in order. Stream reads "novo e2e ok".
          • Vercel prod (world-vercel): deltas reorder deterministically. Stream reads e.g. "novo2e e ok", with text-delta arriving before text-start.

          We verified the model output is correct upstream — direct streamText() and direct engine consumption both produce the right deltas in the right order. The reorder happens entirely between the per-chunk getWritable() calls and the downstream reader.

          Mechanism

          Documented in your own docs/changelog/eager-processing.mdx:330-331:

          On local (world-local), stream writes go to the filesystem — effectively instant. On Vercel (world-vercel), writes go through HTTP to workflow-server → S3, adding 50-100ms latency.

          That latency turns the per-pipe race window from microseconds (locally invisible) into tens of milliseconds (prod-observable). Small/fast deltas — exactly the shape of reasoning-model streaming output — surface the bug most reliably.

          Canonical pattern (works)

          docs/api-reference/workflow-ai/durable-agent.mdx:46 shows the right shape:

          constwritable=getWritable<UIMessageChunk>();constresult=awaitagent.stream({ messages, writable });

          One call. pipeTo(writable) (or the equivalent inside DurableAgent) holds the writer lock for the pipe's lifetime, so there's one TransformStream + one flushablePipe and the race goes away. We fixed our consumer by switching to this pattern + the AI SDK's toUIMessageStream().pipeTo(writable, { preventClose: true, preventAbort: true }).

          Suggestion

          The API surface gives no hint that re-acquisition is unsafe. Two options that would have caught this at dev time without needing prod telemetry:

          1. Dev-mode console.warn when getWritable() is called more than once per step context. contextStorage already tracks per-step state — a simple counter on the ctx would flag the misuse the second time it happens.

          2. Memoize per (runId, namespace): repeat calls within the same step ctx return the same handle (same underlying TransformStream, same flushablePipe). Idempotent, no behavior change for the canonical pattern, makes the unsafe pattern simply correct instead of subtly broken.

          Either would have caught the bug in our first local test run instead of letting it ship to prod.

          Context

          We're building Novo Agents (hosted multi-agent service) on Workflow. The full fix on our side is in rpinckne/harness#50. Happy to help with a repro repo or PR if useful.

          Metadata

          Metadata

          Assignees

          No one assigned

            Labels

            No labels
            No labels

            Type

            No type

            Projects

            No projects

              Milestone

              No milestone

              Relationships

              None yet

              Development

              No branches or pull requests

              Issue actions

              , 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write · Issue #2058 · vercel/workflow · GitHub
              Skip to content

              getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write #2058

              Description

              @rpinckne

              Summary

              Calling getWritable() more than once within a single step body produces non-deterministic chunk ordering on Vercel prod. The API surface gives no hint that re-acquisition is unsafe.

              Where the race lives

              @workflow/core/dist/step/writable-stream.js:16-39 (workflow@5.0.0-beta.6):

              exportfunctiongetWritable(options={}){// ...constserialize=getSerializeStream(getExternalReducers(...),ctx.encryptionKey);constserverWritable=newWorkflowServerWritableStream(runId,name);conststate=createFlushableState();ctx.ops.push(state.promise);flushablePipe(serialize.readable,serverWritable,state).catch(()=>{});pollWritableLock(serialize.writable,state);returnserialize.writable;}

              Every invocation constructs a freshTransformStream + WorkflowServerWritableStream + an independent background flushablePipe(...). The (runId, name) is shared, but the pipes are independent client-side and race to flush into the same downstream queue.

              Repro

              Sequential per-chunk acquisition in a step body, no parallelism:

              'use step';exportasyncfunctionemitStream(parts: UIMessageChunk[]){for(constpartofparts){constwriter=getWritable<UIMessageChunk>().getWriter();try{awaitwriter.write(part);}finally{writer.releaseLock();}}}

              Emit ~6 small chunks (e.g. text-deltas "nov", "o", " e", "2", "e", " ok"). Read the resulting stream via workflowRun.getReadable<UIMessageChunk>(...).

              • Local dev (world-local): deltas arrive in order. Stream reads "novo e2e ok".
              • Vercel prod (world-vercel): deltas reorder deterministically. Stream reads e.g. "novo2e e ok", with text-delta arriving before text-start.

              We verified the model output is correct upstream — direct streamText() and direct engine consumption both produce the right deltas in the right order. The reorder happens entirely between the per-chunk getWritable() calls and the downstream reader.

              Mechanism

              Documented in your own docs/changelog/eager-processing.mdx:330-331:

              On local (world-local), stream writes go to the filesystem — effectively instant. On Vercel (world-vercel), writes go through HTTP to workflow-server → S3, adding 50-100ms latency.

              That latency turns the per-pipe race window from microseconds (locally invisible) into tens of milliseconds (prod-observable). Small/fast deltas — exactly the shape of reasoning-model streaming output — surface the bug most reliably.

              Canonical pattern (works)

              docs/api-reference/workflow-ai/durable-agent.mdx:46 shows the right shape:

              constwritable=getWritable<UIMessageChunk>();constresult=awaitagent.stream({ messages, writable });

              One call. pipeTo(writable) (or the equivalent inside DurableAgent) holds the writer lock for the pipe's lifetime, so there's one TransformStream + one flushablePipe and the race goes away. We fixed our consumer by switching to this pattern + the AI SDK's toUIMessageStream().pipeTo(writable, { preventClose: true, preventAbort: true }).

              Suggestion

              The API surface gives no hint that re-acquisition is unsafe. Two options that would have caught this at dev time without needing prod telemetry:

              1. Dev-mode console.warn when getWritable() is called more than once per step context. contextStorage already tracks per-step state — a simple counter on the ctx would flag the misuse the second time it happens.

              2. Memoize per (runId, namespace): repeat calls within the same step ctx return the same handle (same underlying TransformStream, same flushablePipe). Idempotent, no behavior change for the canonical pattern, makes the unsafe pattern simply correct instead of subtly broken.

              Either would have caught the bug in our first local test run instead of letting it ship to prod.

              Context

              We're building Novo Agents (hosted multi-agent service) on Workflow. The full fix on our side is in rpinckne/harness#50. Happy to help with a repro repo or PR if useful.

              Metadata

              Metadata

              Assignees

              No one assigned

                Labels

                No labels
                No labels

                Type

                No type

                Projects

                No projects

                  Milestone

                  No milestone

                  Relationships

                  None yet

                  Development

                  No branches or pull requests

                  Issue actions

                  , 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write · Issue #2058 · vercel/workflow · GitHub
                  Skip to content

                  getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write #2058

                  Description

                  @rpinckne

                  Summary

                  Calling getWritable() more than once within a single step body produces non-deterministic chunk ordering on Vercel prod. The API surface gives no hint that re-acquisition is unsafe.

                  Where the race lives

                  @workflow/core/dist/step/writable-stream.js:16-39 (workflow@5.0.0-beta.6):

                  exportfunctiongetWritable(options={}){// ...constserialize=getSerializeStream(getExternalReducers(...),ctx.encryptionKey);constserverWritable=newWorkflowServerWritableStream(runId,name);conststate=createFlushableState();ctx.ops.push(state.promise);flushablePipe(serialize.readable,serverWritable,state).catch(()=>{});pollWritableLock(serialize.writable,state);returnserialize.writable;}

                  Every invocation constructs a freshTransformStream + WorkflowServerWritableStream + an independent background flushablePipe(...). The (runId, name) is shared, but the pipes are independent client-side and race to flush into the same downstream queue.

                  Repro

                  Sequential per-chunk acquisition in a step body, no parallelism:

                  'use step';exportasyncfunctionemitStream(parts: UIMessageChunk[]){for(constpartofparts){constwriter=getWritable<UIMessageChunk>().getWriter();try{awaitwriter.write(part);}finally{writer.releaseLock();}}}

                  Emit ~6 small chunks (e.g. text-deltas "nov", "o", " e", "2", "e", " ok"). Read the resulting stream via workflowRun.getReadable<UIMessageChunk>(...).

                  • Local dev (world-local): deltas arrive in order. Stream reads "novo e2e ok".
                  • Vercel prod (world-vercel): deltas reorder deterministically. Stream reads e.g. "novo2e e ok", with text-delta arriving before text-start.

                  We verified the model output is correct upstream — direct streamText() and direct engine consumption both produce the right deltas in the right order. The reorder happens entirely between the per-chunk getWritable() calls and the downstream reader.

                  Mechanism

                  Documented in your own docs/changelog/eager-processing.mdx:330-331:

                  On local (world-local), stream writes go to the filesystem — effectively instant. On Vercel (world-vercel), writes go through HTTP to workflow-server → S3, adding 50-100ms latency.

                  That latency turns the per-pipe race window from microseconds (locally invisible) into tens of milliseconds (prod-observable). Small/fast deltas — exactly the shape of reasoning-model streaming output — surface the bug most reliably.

                  Canonical pattern (works)

                  docs/api-reference/workflow-ai/durable-agent.mdx:46 shows the right shape:

                  constwritable=getWritable<UIMessageChunk>();constresult=awaitagent.stream({ messages, writable });

                  One call. pipeTo(writable) (or the equivalent inside DurableAgent) holds the writer lock for the pipe's lifetime, so there's one TransformStream + one flushablePipe and the race goes away. We fixed our consumer by switching to this pattern + the AI SDK's toUIMessageStream().pipeTo(writable, { preventClose: true, preventAbort: true }).

                  Suggestion

                  The API surface gives no hint that re-acquisition is unsafe. Two options that would have caught this at dev time without needing prod telemetry:

                  1. Dev-mode console.warn when getWritable() is called more than once per step context. contextStorage already tracks per-step state — a simple counter on the ctx would flag the misuse the second time it happens.

                  2. Memoize per (runId, namespace): repeat calls within the same step ctx return the same handle (same underlying TransformStream, same flushablePipe). Idempotent, no behavior change for the canonical pattern, makes the unsafe pattern simply correct instead of subtly broken.

                  Either would have caught the bug in our first local test run instead of letting it ship to prod.

                  Context

                  We're building Novo Agents (hosted multi-agent service) on Workflow. The full fix on our side is in rpinckne/harness#50. Happy to help with a repro repo or PR if useful.

                  Metadata

                  Metadata

                  Assignees

                  No one assigned

                    Labels

                    No labels
                    No labels

                    Type

                    No type

                    Projects

                    No projects

                      Milestone

                      No milestone

                      Relationships

                      None yet

                      Development

                      No branches or pull requests

                      Issue actions

                      , 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write · Issue #2058 · vercel/workflow · GitHub
                      Skip to content

                      getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write #2058

                      Description

                      @rpinckne

                      Summary

                      Calling getWritable() more than once within a single step body produces non-deterministic chunk ordering on Vercel prod. The API surface gives no hint that re-acquisition is unsafe.

                      Where the race lives

                      @workflow/core/dist/step/writable-stream.js:16-39 (workflow@5.0.0-beta.6):

                      exportfunctiongetWritable(options={}){// ...constserialize=getSerializeStream(getExternalReducers(...),ctx.encryptionKey);constserverWritable=newWorkflowServerWritableStream(runId,name);conststate=createFlushableState();ctx.ops.push(state.promise);flushablePipe(serialize.readable,serverWritable,state).catch(()=>{});pollWritableLock(serialize.writable,state);returnserialize.writable;}

                      Every invocation constructs a freshTransformStream + WorkflowServerWritableStream + an independent background flushablePipe(...). The (runId, name) is shared, but the pipes are independent client-side and race to flush into the same downstream queue.

                      Repro

                      Sequential per-chunk acquisition in a step body, no parallelism:

                      'use step';exportasyncfunctionemitStream(parts: UIMessageChunk[]){for(constpartofparts){constwriter=getWritable<UIMessageChunk>().getWriter();try{awaitwriter.write(part);}finally{writer.releaseLock();}}}

                      Emit ~6 small chunks (e.g. text-deltas "nov", "o", " e", "2", "e", " ok"). Read the resulting stream via workflowRun.getReadable<UIMessageChunk>(...).

                      • Local dev (world-local): deltas arrive in order. Stream reads "novo e2e ok".
                      • Vercel prod (world-vercel): deltas reorder deterministically. Stream reads e.g. "novo2e e ok", with text-delta arriving before text-start.

                      We verified the model output is correct upstream — direct streamText() and direct engine consumption both produce the right deltas in the right order. The reorder happens entirely between the per-chunk getWritable() calls and the downstream reader.

                      Mechanism

                      Documented in your own docs/changelog/eager-processing.mdx:330-331:

                      On local (world-local), stream writes go to the filesystem — effectively instant. On Vercel (world-vercel), writes go through HTTP to workflow-server → S3, adding 50-100ms latency.

                      That latency turns the per-pipe race window from microseconds (locally invisible) into tens of milliseconds (prod-observable). Small/fast deltas — exactly the shape of reasoning-model streaming output — surface the bug most reliably.

                      Canonical pattern (works)

                      docs/api-reference/workflow-ai/durable-agent.mdx:46 shows the right shape:

                      constwritable=getWritable<UIMessageChunk>();constresult=awaitagent.stream({ messages, writable });

                      One call. pipeTo(writable) (or the equivalent inside DurableAgent) holds the writer lock for the pipe's lifetime, so there's one TransformStream + one flushablePipe and the race goes away. We fixed our consumer by switching to this pattern + the AI SDK's toUIMessageStream().pipeTo(writable, { preventClose: true, preventAbort: true }).

                      Suggestion

                      The API surface gives no hint that re-acquisition is unsafe. Two options that would have caught this at dev time without needing prod telemetry:

                      1. Dev-mode console.warn when getWritable() is called more than once per step context. contextStorage already tracks per-step state — a simple counter on the ctx would flag the misuse the second time it happens.

                      2. Memoize per (runId, namespace): repeat calls within the same step ctx return the same handle (same underlying TransformStream, same flushablePipe). Idempotent, no behavior change for the canonical pattern, makes the unsafe pattern simply correct instead of subtly broken.

                      Either would have caught the bug in our first local test run instead of letting it ship to prod.

                      Context

                      We're building Novo Agents (hosted multi-agent service) on Workflow. The full fix on our side is in rpinckne/harness#50. Happy to help with a repro repo or PR if useful.

                      Metadata

                      Metadata

                      Assignees

                      No one assigned

                        Labels

                        No labels
                        No labels

                        Type

                        No type

                        Projects

                        No projects

                          Milestone

                          No milestone

                          Relationships

                          None yet

                          Development

                          No branches or pull requests

                          Issue actions

                          , 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write · Issue #2058 · vercel/workflow · GitHub
                          Skip to content

                          getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write #2058

                          Description

                          @rpinckne

                          Summary

                          Calling getWritable() more than once within a single step body produces non-deterministic chunk ordering on Vercel prod. The API surface gives no hint that re-acquisition is unsafe.

                          Where the race lives

                          @workflow/core/dist/step/writable-stream.js:16-39 (workflow@5.0.0-beta.6):

                          exportfunctiongetWritable(options={}){// ...constserialize=getSerializeStream(getExternalReducers(...),ctx.encryptionKey);constserverWritable=newWorkflowServerWritableStream(runId,name);conststate=createFlushableState();ctx.ops.push(state.promise);flushablePipe(serialize.readable,serverWritable,state).catch(()=>{});pollWritableLock(serialize.writable,state);returnserialize.writable;}

                          Every invocation constructs a freshTransformStream + WorkflowServerWritableStream + an independent background flushablePipe(...). The (runId, name) is shared, but the pipes are independent client-side and race to flush into the same downstream queue.

                          Repro

                          Sequential per-chunk acquisition in a step body, no parallelism:

                          'use step';exportasyncfunctionemitStream(parts: UIMessageChunk[]){for(constpartofparts){constwriter=getWritable<UIMessageChunk>().getWriter();try{awaitwriter.write(part);}finally{writer.releaseLock();}}}

                          Emit ~6 small chunks (e.g. text-deltas "nov", "o", " e", "2", "e", " ok"). Read the resulting stream via workflowRun.getReadable<UIMessageChunk>(...).

                          • Local dev (world-local): deltas arrive in order. Stream reads "novo e2e ok".
                          • Vercel prod (world-vercel): deltas reorder deterministically. Stream reads e.g. "novo2e e ok", with text-delta arriving before text-start.

                          We verified the model output is correct upstream — direct streamText() and direct engine consumption both produce the right deltas in the right order. The reorder happens entirely between the per-chunk getWritable() calls and the downstream reader.

                          Mechanism

                          Documented in your own docs/changelog/eager-processing.mdx:330-331:

                          On local (world-local), stream writes go to the filesystem — effectively instant. On Vercel (world-vercel), writes go through HTTP to workflow-server → S3, adding 50-100ms latency.

                          That latency turns the per-pipe race window from microseconds (locally invisible) into tens of milliseconds (prod-observable). Small/fast deltas — exactly the shape of reasoning-model streaming output — surface the bug most reliably.

                          Canonical pattern (works)

                          docs/api-reference/workflow-ai/durable-agent.mdx:46 shows the right shape:

                          constwritable=getWritable<UIMessageChunk>();constresult=awaitagent.stream({ messages, writable });

                          One call. pipeTo(writable) (or the equivalent inside DurableAgent) holds the writer lock for the pipe's lifetime, so there's one TransformStream + one flushablePipe and the race goes away. We fixed our consumer by switching to this pattern + the AI SDK's toUIMessageStream().pipeTo(writable, { preventClose: true, preventAbort: true }).

                          Suggestion

                          The API surface gives no hint that re-acquisition is unsafe. Two options that would have caught this at dev time without needing prod telemetry:

                          1. Dev-mode console.warn when getWritable() is called more than once per step context. contextStorage already tracks per-step state — a simple counter on the ctx would flag the misuse the second time it happens.

                          2. Memoize per (runId, namespace): repeat calls within the same step ctx return the same handle (same underlying TransformStream, same flushablePipe). Idempotent, no behavior change for the canonical pattern, makes the unsafe pattern simply correct instead of subtly broken.

                          Either would have caught the bug in our first local test run instead of letting it ship to prod.

                          Context

                          We're building Novo Agents (hosted multi-agent service) on Workflow. The full fix on our side is in rpinckne/harness#50. Happy to help with a repro repo or PR if useful.

                          Metadata

                          Metadata

                          Assignees

                          No one assigned

                            Labels

                            No labels
                            No labels

                            Type

                            No type

                            Projects

                            No projects

                              Milestone

                              No milestone

                              Relationships

                              None yet

                              Development

                              No branches or pull requests

                              Issue actions

                              , 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write · Issue #2058 · vercel/workflow · GitHub
                              Skip to content

                              getWritable() returns a new TransformStream per call — racing pipes reorder chunks when callers acquire per-write #2058

                              Description

                              @rpinckne

                              Summary

                              Calling getWritable() more than once within a single step body produces non-deterministic chunk ordering on Vercel prod. The API surface gives no hint that re-acquisition is unsafe.

                              Where the race lives

                              @workflow/core/dist/step/writable-stream.js:16-39 (workflow@5.0.0-beta.6):

                              exportfunctiongetWritable(options={}){// ...constserialize=getSerializeStream(getExternalReducers(...),ctx.encryptionKey);constserverWritable=newWorkflowServerWritableStream(runId,name);conststate=createFlushableState();ctx.ops.push(state.promise);flushablePipe(serialize.readable,serverWritable,state).catch(()=>{});pollWritableLock(serialize.writable,state);returnserialize.writable;}

                              Every invocation constructs a freshTransformStream + WorkflowServerWritableStream + an independent background flushablePipe(...). The (runId, name) is shared, but the pipes are independent client-side and race to flush into the same downstream queue.

                              Repro

                              Sequential per-chunk acquisition in a step body, no parallelism:

                              'use step';exportasyncfunctionemitStream(parts: UIMessageChunk[]){for(constpartofparts){constwriter=getWritable<UIMessageChunk>().getWriter();try{awaitwriter.write(part);}finally{writer.releaseLock();}}}

                              Emit ~6 small chunks (e.g. text-deltas "nov", "o", " e", "2", "e", " ok"). Read the resulting stream via workflowRun.getReadable<UIMessageChunk>(...).

                              • Local dev (world-local): deltas arrive in order. Stream reads "novo e2e ok".
                              • Vercel prod (world-vercel): deltas reorder deterministically. Stream reads e.g. "novo2e e ok", with text-delta arriving before text-start.

                              We verified the model output is correct upstream — direct streamText() and direct engine consumption both produce the right deltas in the right order. The reorder happens entirely between the per-chunk getWritable() calls and the downstream reader.

                              Mechanism

                              Documented in your own docs/changelog/eager-processing.mdx:330-331:

                              On local (world-local), stream writes go to the filesystem — effectively instant. On Vercel (world-vercel), writes go through HTTP to workflow-server → S3, adding 50-100ms latency.

                              That latency turns the per-pipe race window from microseconds (locally invisible) into tens of milliseconds (prod-observable). Small/fast deltas — exactly the shape of reasoning-model streaming output — surface the bug most reliably.

                              Canonical pattern (works)

                              docs/api-reference/workflow-ai/durable-agent.mdx:46 shows the right shape:

                              constwritable=getWritable<UIMessageChunk>();constresult=awaitagent.stream({ messages, writable });

                              One call. pipeTo(writable) (or the equivalent inside DurableAgent) holds the writer lock for the pipe's lifetime, so there's one TransformStream + one flushablePipe and the race goes away. We fixed our consumer by switching to this pattern + the AI SDK's toUIMessageStream().pipeTo(writable, { preventClose: true, preventAbort: true }).

                              Suggestion

                              The API surface gives no hint that re-acquisition is unsafe. Two options that would have caught this at dev time without needing prod telemetry:

                              1. Dev-mode console.warn when getWritable() is called more than once per step context. contextStorage already tracks per-step state — a simple counter on the ctx would flag the misuse the second time it happens.

                              2. Memoize per (runId, namespace): repeat calls within the same step ctx return the same handle (same underlying TransformStream, same flushablePipe). Idempotent, no behavior change for the canonical pattern, makes the unsafe pattern simply correct instead of subtly broken.

                              Either would have caught the bug in our first local test run instead of letting it ship to prod.

                              Context

                              We're building Novo Agents (hosted multi-agent service) on Workflow. The full fix on our side is in rpinckne/harness#50. Happy to help with a repro repo or PR if useful.

                              Metadata

                              Metadata

                              Assignees

                              No one assigned

                                Labels

                                No labels
                                No labels

                                Type

                                No type

                                Projects

                                No projects

                                  Milestone

                                  No milestone

                                  Relationships

                                  None yet

                                  Development

                                  No branches or pull requests

                                  Issue actions