This release is intended to be the last one before rill's API and semantics stabilize in v1.0. Feedback is especially welcome now, while things can still change.
It includes new features, bug fixes and some breaking changes. A review of public GitHub usage suggests the breaking changes have little to no impact on existing users. The migration guide at the end lists what to check.
Rill's semantics are now documented in detail and pinned by a much stronger test suite. The library now requires Go 1.25 or newer (#81).
Highlights
First-class context support (#113)
By default, pipelines have the same lifecycle as in previous versions. Context support is opt-in through functional options accepted by ForEach, ToSlice, and other sinks. The new WithContext function gives pipelines errgroup-style cancellation and waiting: cancel on the first error, wait until nothing is running anymore.
// ctx is captured by the callbacks below.
// scope is a functional option passed to ForEach (or another sink).
ctx, scope := rill.WithContext(ctx)
// Source and other pipeline stages go here.
users := rill.Map(ids, 10, func(id int) (*User, error) {
return getUser(ctx, id) // context-aware
})
err := rill.ForEach(users, 5, func(u *User) error {
return process(ctx, u) // context-aware
}, scope)
// Nothing is running anymore; the context is canceled.
// Handle the error (if any).
fmt.Println("Error:", err)Non-commutative streaming reduction (#98)
Previously, Reduce combined values in whatever order its workers picked them up, so reducers had to be both associative and commutative. Now it computes
in[0] ⊕ in[1] ⊕ in[2] ⊕ ... ⊕ in[N-1]
where ⊕ is the reducer. Values stay in stream order, but the parenthesization is unspecified. This still requires associativity, but no longer requires commutativity. As a result, a much wider class of reducers is now supported. MapReduce does the same for the values of each key.
Reducers that were valid before still work. The only breaking change is a relaxed guarantee for n = 1. The documentation used to promise sequential processing, which implied a left-to-right fold. At n = 1, Reduce still folds left to right, but that is no longer guaranteed, and the parenthesization may change in a future release.
Errors are batch boundaries (#97)
Previously, Batch preserved the order of values and errors separately, but could reorder errors relative to values. Suppose we have a stream and want to split it into batches of 5:
1, 2, 3, err1, 4, 5, 6, 7, 8, err2, 9, 10
The old version produced something like:
err1, [1, 2, 3, 4, 5], err2, [6, 7, 8, 9, 10]
The new version deterministically produces:
[1, 2, 3], err1, [4, 5, 6, 7, 8], err2, [9, 10]
Stream conversion and nil handling (#93)
These changes were motivated by two goals: to improve composability with third-party code, and to stop giving special meaning to nil channels. In Go, a nil channel never emits and never closes, and rill now treats it that way.
- FromSlice(slice, err): emits all values followed by err, preserving partial results returned by APIs such as
os.ReadDir(previously the values were dropped in case of an error). Also, this makesToSlice(FromSlice(s, err))a round trip. - FromChan(values, err): treats err as a construction error: a non-nil error is emitted alone and values are ignored (previously the channel was consumed after the error). Matches APIs such as RabbitMQ's
Channel.Consume, which returns(nil, err)when consumer setup fails, enablingrill.FromChan(ch.Consume(...)). - FromChans(values, errs): treats both channels as independent inputs and closes only after both are fully consumed. If one or both of the inputs is nil, the output stream is never closed (previously nil inputs were ignored).
- FromSeq(seq, err): treats err as a construction error: a non-nil error is emitted alone and seq is ignored. When there's no error, a nil sequence is treated as misuse and panics (previously returned a nil stream).
- FromSeq2(seq): a nil sequence panics as above (previously returned a nil stream).
- Merge(ins...): with no arguments returns an empty closed channel (previously returned nil).
Stronger tests and documentation
The package and function documentation now specifies the pipeline model and lifecycle: when sinks return, what keeps running after an early return, and the rules custom stages must follow (#110). The test suite is much stronger and is now based on synctest (#85). Compared with v0.8.1, the suite is:
- Broader. More edge cases, more assertions, and the key documented contracts pinned.
- Precise. Sleeps happen in virtual time, where they are free, so tests construct the goroutine interleavings they need (still randomized) instead of hoping the scheduler produces them. Timing assertions are exact.
- Lifecycle-aware. Tests pin when sinks return, when the context is canceled, how much extra work runs after an early exit, and when the rest of the input is drained (#113, #115).
- Leak-checked. Every synctest bubble also verifies that no goroutine is left blocked, so leak detection comes for free.
- Faster. A race-enabled run takes 3 seconds instead of a minute, so CI now runs each test ten times, shuffled and under the race detector, at
GOMAXPROCS1 and 8 (#101).
Bug fixes
- All returned
truetogether with an error. It now returns(false, err)(#86). - First returned
found = truewhen the first item was an error, so callers checkingfoundbeforeerrcould use a value that didn't exist. It now returns(zero, false, err)(#90, #103). - Generate emitted a zero-valued item when
sendErrorwas called with a nil error. Such calls are now ignored (#104).
Other changes
- Wrap and FromSeq2 discard the accompanying value when an error is non-nil, matching
Try's value-or-error semantics (#103). - Tee was generalized and can now take any channel, not only streams. Existing calls are unaffected if they rely on type inference (#111).
sizepassed to Buffer now means the capacity of its output channel. Previously, it used a capacity ofsize-1to compensate for the item held by the forwarding goroutine (#96).- Invalid concurrency levels, nil callbacks, and invalid batch or buffer sizes now panic synchronously at the public call site (#95).
Performance improvements
- The cost of preserving order dropped by about half, and by more at higher concurrency, measured with no-op callbacks (#100).
- Batch allocates lazily, avoiding preallocated buffers while idle (#106).
Deprecations
- Split2 and OrderedSplit2 will be removed in v1.0. The same behavior can be composed from Tee and Filter, and the composition supports more than two branches. See the doc comments of these functions for migration examples (#89).
- DrainNB will be removed in v1.0 in favor of Discard.
Migration guide
Rill now requires Go 1.25 or newer. Most code needs no changes, though the cases below are worth checking.
- FromChans with a nil input. A nil channel is no longer ignored: it never closes, so the output never closes either, and the pipeline hangs. If you only have a value channel, use
FromChan(values, nil)instead. - FromSlice, FromChan, FromSeq, FromSeq2. Check call sites that can receive a nil input or a non-nil error. Their semantics changed as described in the conversions section above.
- Batch without a timeout. This only matters if your code relies on every batch except the last being full. Handle partial batches, since an error in the stream now becomes a batch boundary.
- All and First. Check call sites that use their other return values when the error is non-nil.
- Wrap, FromSeq2. These functions no longer construct items that hold a value and an error at the same time. If you have custom stages that rely on that, send the value and the error as separate items.
- Compile errors. Sinks got variadic options and Tee now accepts any channel, so only code that stores a sink as a function value or passes explicit type arguments to Tee is affected.
- Reduce at
n = 1. Reduce still folds left to right atn = 1, but that may change. If your code relies on that, use ForEach withn = 1and an accumulator variable. - Split2, OrderedSplit2 and DrainNB were deprecated. See Deprecations above.
New Contributors
- @maxtaran2010 made their first contribution in #87
All pull requests
🟢 New
- Make batching stricter: errors become batch boundaries by @destel in #97
- New ordered reduction engine: commutativity is no longer required by @destel in #98
- Add pipeline settlement signal by @destel in #108
- Add Scope API: context and settlement support by @destel in #109
- Add WithContext: context and structured concurrency support by @destel in #113
🔴 Fixes
- Fix
Allreturning true on error by @destel in #86 - Fix
Firstreturning found=true on error by @destel in #90 - Ignore nil errors in Generate by @destel in #104
🟣 Documentation
- fix: correct typos in doc comments (batch.go, reduce.go) by @maxtaran2010 in #87
- FromSlice: document that the slice must not be modified while it is read by @destel in #99
- Add a runnable example for Tee by @destel in #102
- Rewrite package and function docs by @destel in #110
- Update README and improve examples by @destel in #114
🟠 Dependencies
- Bump softprops/action-gh-release from 2 to 3 by @dependabot[bot] in #77
- Bump codecov/codecov-action from 5.5.2 to 7.0.0 by @dependabot[bot] in #79
- Bump actions/checkout from 6 to 7 by @dependabot[bot] in #82
- Bump actions/setup-go from 6 to 7 by @dependabot[bot] in #92
- Bump codecov/codecov-action from 7.0.0 to 7.1.1 by @dependabot[bot] in #112
🔵 Other
- Remove dead internal code by @destel in #80
- Raise minimum Go version to 1.25 by @destel in #81
- Add golangci-lint by @destel in #83
- Drop Coveralls, keep Codecov by @destel in #84
- Add pre-commit config with golangci-lint by @destel in #88
- Strengthen concurrency tests with synctest: exact assertions without scheduler luck by @destel in #85
- Deprecate Split2 and OrderedSplit2 by @destel in #89
- Rework blocking sinks: one internal pattern, stricter contract tests by @destel in #91
- Rework stream conversions: adapters follow the nature of their inputs by @destel in #93
- Create GitHub releases as drafts by @destel in #94
- Validate arguments at the public API boundary by @destel in #95
- Simplify Buffer: size is plain channel capacity by @destel in #96
- Optimize ordered transforms: ~50% less overhead by @destel in #100
- Run CI tests at different GOMAXPROCS settings by @destel in #101
- Try container is a value or an error, not both by @destel in #103
- mockapi: 32-bit compatibility by @destel in #105
- Allocate lazily in Batch by @destel in #106
- Ignore nil channels in Discard by @destel in #107
- Generalize Tee to any channel by @destel in #111
- Add more tests for core loops by @destel in #115
Full Changelog: v0.8.1...v0.9.0