You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
mergeResultFilesToParquet now holds one decoded input file at a time and writes to a caller-supplied parquet Writer, so merging a large run no longer needs every row plus the whole output buffer in memory.
What changed?
Breaking signature change. The only consumer is services/gcp-harbor-run/src/publish-native-run.ts in openrouter-web, which moves over in the subtree-pull PR.
It makes two passes over the readers, one file at a time:
Pass 1 collects row counts, light SampleScores (id, epoch, value), summed usage, the weighted primary score, and the first file's epochs, temperature and benchmark_config.
Pass 2 re-reads each non-empty file and checks it against pass 1: row count, each row's sample_id/epoch/score_value, and a fingerprint of the first row's run-level fields. Any mismatch throws. It then overwrites the run-level columns and writes rows through ParquetWriter in row groups of at most 100 rows.
The merge semantics and the full raw schema stay the same. Empty files are still skipped, extra_scores is still null, and format_version is still RESULT_FORMAT_VERSION. The schema and KV metadata match runResultToParquet.
The returned summary equals summarizeChunkRows(mergedRows), so callers don't need to decode the output again. summarizeChunkRows and the merge now share summarizeSampleScores.
rowsToSampleScores no longer copies answer, because it only feeds aggregates.
New export ./parquet-stream-writer. parquetStreamWriter(writeChunk, chunkBytes = 1_000_000) is a Writer that hands bytes to writeChunk in order on each row-group flush and on finish, instead of retaining the file. It's meant for piping into a GCS upload stream.
Why?
Providers need one GCS download link per run. The only existing merge built everything in memory: all decoded rows, the full output Buffer, and then the caller decoded it again. That will OOM on agentic runs with hundreds of trials.
How to test
Memory probe, run locally: 60 files × 20 rows, each row with a unique 200 KB explanation, 222 MB of parquet on disk. The old path mirrors publish-native-run (read all, merge to a Buffer, parse again). The new path uses file readers plus parquetStreamWriter appending to disk.
old (origin/main): max RSS 2,445,668 KB 1200 rows output 231,805,554 B
new (this branch): max RSS 285,208 KB 1200 rows output 232,037,709 B
The new output is about 0.1% larger because each row group carries its own dictionary and page index. These numbers are from the first commit (one row group per file). The unit tests check that the streamed bytes match an in-memory ByteWriter merge byte-for-byte, and that readers are called sequentially, twice each.
Reviewer focus
Row groups never span input files and hold at most 100 rows. Many tiny trial files mean many small row groups, and footer metadata grows linearly with the row-group count.
parquetStreamWriter overrides ensure, flush and finish on a ByteWriter, the same pattern as hyparquet-writer's own fileWriter. getBuffer and getBytes throw on purpose.
Checklist
Tests cover changed behavior
Public API or configuration changes are backward compatible, or the break is documented
Benchmark changes document dataset provenance and licensing (n/a)
No credentials, private results, or restricted dataset contents are included
The <latest_message> is the message that you should use to guide your goals + task for this session, and you should use the rest of the slack thread as context.
A [ts=...] marker on a Slack message is that message's timestamp. To act on a specific message with the slack tool (e.g. adding an emoji reaction via the reaction command), pass that value as timestamp along with the Channel ID — no extra lookup call is needed.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
TL;DR
mergeResultFilesToParquetnow holds one decoded input file at a time and writes to a caller-supplied parquetWriter, so merging a large run no longer needs every row plus the whole output buffer in memory.What changed?
services/gcp-harbor-run/src/publish-native-run.tsin openrouter-web, which moves over in the subtree-pull PR.SampleScores (id, epoch, value), summed usage, the weighted primary score, and the first file's epochs, temperature and benchmark_config.sample_id/epoch/score_value, and a fingerprint of the first row's run-level fields. Any mismatch throws. It then overwrites the run-level columns and writes rows throughParquetWriterin row groups of at most 100 rows.extra_scoresis still null, andformat_versionis stillRESULT_FORMAT_VERSION. The schema and KV metadata matchrunResultToParquet.summaryequalssummarizeChunkRows(mergedRows), so callers don't need to decode the output again.summarizeChunkRowsand the merge now sharesummarizeSampleScores.rowsToSampleScoresno longer copiesanswer, because it only feeds aggregates../parquet-stream-writer.parquetStreamWriter(writeChunk, chunkBytes = 1_000_000)is aWriterthat hands bytes towriteChunkin order on each row-group flush and on finish, instead of retaining the file. It's meant for piping into a GCS upload stream.Why?
Providers need one GCS download link per run. The only existing merge built everything in memory: all decoded rows, the full output
Buffer, and then the caller decoded it again. That will OOM on agentic runs with hundreds of trials.How to test
Memory probe, run locally: 60 files × 20 rows, each row with a unique 200 KB
explanation, 222 MB of parquet on disk. The old path mirrorspublish-native-run(read all, merge to a Buffer, parse again). The new path uses file readers plusparquetStreamWriterappending to disk.The new output is about 0.1% larger because each row group carries its own dictionary and page index. These numbers are from the first commit (one row group per file). The unit tests check that the streamed bytes match an in-memory
ByteWritermerge byte-for-byte, and that readers are called sequentially, twice each.Reviewer focus
parquetStreamWriteroverridesensure,flushandfinishon aByteWriter, the same pattern as hyparquet-writer's ownfileWriter.getBufferandgetBytesthrow on purpose.Checklist
Link to Devin session: https://openrouter.devinenterprise.com/sessions/cff9171c8dd8432ebb515e2a54374727
Open in Devin Desktop: https://openrouter.devinenterprise.com/desktop/session/cff9171c8dd8432ebb515e2a54374727?variant=devin
Requested by: @abhinav-pola