[teamai] Push 87 resource(s) from XingfenD

This commit is contained in:
2026-09-10 16:10:45 +08:00
parent 425c9c078a
commit 65c04def51
1314 changed files with 211681 additions and 0 deletions
@@ -0,0 +1,324 @@
# Operators Guide
samber/ro provides 150+ operators organized by category. All operators are generic functions that take and return `Observable[T]`, designed for chaining via `Pipe`.
## Pipeline Construction
```go
// Typed pipes (Pipe1 through Pipe25) — compile-time type safety
result := ro.Pipe2(source, op1, op2)
result := ro.Pipe3(source, op1, op2, op3)
// Untyped pipe — uses any, loses type checking
result := ro.Pipe(source, op1, op2, op3)
// Reusable operator composition (curried)
transform := ro.PipeOp2(op1, op2) // returns func(Observable[A]) Observable[C]
result := transform(source)
// Synchronous collection
values, err := ro.Collect(observable)
```
## Creation Operators
Create observables from various sources. Typically the first argument to `Pipe`.
| Operator | Signature | Purpose |
| --- | --- | --- |
| `Just` | `Just[T](values ...T) Observable[T]` | Emit specific values then complete |
| `Of` | `Of[T](values ...T) Observable[T]` | Alias for Just |
| `FromSlice` | `FromSlice[T](items []T) Observable[T]` | Create from existing slice |
| `FromChannel` | `FromChannel[T](ch <-chan T) Observable[T]` | Wrap Go channel as observable |
| `Range` | `Range(start, end int) Observable[int]` | Emit integer sequence |
| `RangeWithStep` | `RangeWithStep(start, end, step int) Observable[int]` | Integer sequence with custom step |
| `RangeWithInterval` | `RangeWithInterval(start, end int, d time.Duration) Observable[int]` | Integers with delay between each |
| `RangeWithStepAndInterval` | `RangeWithStepAndInterval(start, end, step int, d time.Duration) Observable[int]` | Step + interval combined |
| `Interval` | `Interval(d time.Duration) Observable[int64]` | Emit sequential integers at intervals (infinite) |
| `IntervalWithInitial` | `IntervalWithInitial(initial, d time.Duration) Observable[int64]` | Interval with initial delay |
| `Timer` | `Timer(d time.Duration) Observable[int64]` | Single emission after delay |
| `Repeat` | `Repeat[T](item T, count int64) Observable[T]` | Repeat item N times |
| `RepeatWithInterval` | `RepeatWithInterval[T](item T, count int64, d time.Duration) Observable[T]` | Repeat with delay |
| `Defer` | `Defer[T](factory func() Observable[T]) Observable[T]` | Lazily create observable on subscription |
| `Future` | `Future[T](factory func() (T, error)) Observable[T]` | Single async value from function |
| `Start` | `Start(cb func() (any, error)) Observable[any]` | Execute callback and emit result |
| `Empty` | `Empty[T]() Observable[T]` | Complete immediately, no values |
| `Never` | `Never[T]() Observable[T]` | Never emit or complete |
| `Throw` | `Throw[T](err error) Observable[T]` | Immediately error |
```go
// Custom observable with direct control
obs := ro.NewObservable[int](func(ctx context.Context, observer ro.Observer[int]) error {
observer.Next(1)
observer.Next(2)
observer.Complete()
return nil
})
```
## Transformation Operators
Transform each value in the stream.
| Operator | Purpose |
| --- | --- |
| `Map[T, R](fn func(T) R)` | Transform each value T -> R |
| `MapI[T, R](fn func(T, int64) R)` | Map with index |
| `MapErr[T, R](fn func(T) (R, error))` | Map that can fail — error propagates |
| `MapWithContext[T, R](fn func(ctx, T) (ctx, R))` | Map with context access |
| `MapTo[T, R](output R)` | Replace all values with a constant |
| `FlatMap[T, R](fn func(T) Observable[R])` | Map each value to observable, flatten results |
| `MergeMap[T, R](fn func(T) Observable[R])` | Alias for FlatMap |
| `Scan[T, R](fn func(R, T) R, seed R)` | Running accumulation, emit each intermediate |
| `Reduce[T, R](fn func(R, T) R, seed R)` | Accumulate, emit only final result |
| `GroupBy[T, K](fn func(T) K)` | Partition into grouped observables by key |
| `Cast[T, U]()` | Type cast values |
| `Flatten[T]()` | Flatten `Observable[[]T]` to `Observable[T]` |
| `Materialize[T]()` | Wrap values in `Notification[T]` (next/error/complete) |
| `Dematerialize[T]()` | Unwrap `Notification[T]` back to values |
| `Timestamp[T]()` | Emit `TimestampValue[T]` with emission time |
| `TimeInterval[T]()` | Emit `IntervalValue[T]` with time since last emission |
```go
// FlatMap: for each user ID, fetch their orders (returns observable)
orders := ro.Pipe2(
userIDs,
ro.FlatMap(func(id int) ro.Observable[Order] {
return fetchOrders(id)
}),
ro.Filter(func(o Order) bool { return o.Status == "paid" }),
)
// Scan: running sum
ro.Pipe1(
ro.Just(1, 2, 3, 4, 5),
ro.Scan(func(acc, x int) int { return acc + x }, 0),
)
// Emits: 1, 3, 6, 10, 15
```
## Filtering Operators
Selectively emit values from the stream.
| Operator | Purpose |
| --- | --- |
| `Filter[T](fn func(T) bool)` | Emit only values matching predicate |
| `FilterI[T](fn func(T, int64) bool)` | Filter with index |
| `FilterWithContext[T](fn func(ctx, T) (ctx, bool))` | Filter with context access |
| `Distinct[T comparable]()` | Emit only unique values |
| `DistinctBy[T, K](fn func(T) K)` | Unique by key function |
| `Take[T](n int64)` | Emit first N values then complete |
| `TakeLast[T](n int)` | Emit last N values |
| `TakeWhile[T](fn func(T) bool)` | Emit while predicate true, then complete |
| `TakeUntil[T, S](signal Observable[S])` | Emit until signal observable emits |
| `Skip[T](n int64)` | Skip first N values |
| `SkipLast[T](n int)` | Skip last N values |
| `SkipWhile[T](fn func(T) bool)` | Skip while predicate true |
| `SkipUntil[T, S](signal Observable[S])` | Skip until signal emits |
| `Head[T]()` | First item only |
| `Tail[T]()` | All items except first |
| `ElementAt[T](n int)` | Item at specific index |
| `ElementAtOrDefault[T](n int64, fallback T)` | Item at index or default |
| `Find[T](fn func(T) bool)` | First value matching predicate |
| `First[T](fn func(T) bool)` | First matching value (alias-like) |
| `Last[T](fn func(T) bool)` | Last matching value |
| `Contains[T](fn func(T) bool)` | Emit bool: whether any value matches |
All filtering operators have `I` (indexed), `WithContext`, and `IWithContext` variants.
## Combining Operators
Merge multiple observables into one.
| Operator | Purpose |
| --- | --- |
| `Merge[T](sources ...Observable[T])` | Interleave emissions from all sources |
| `MergeWith[T](obs ...Observable[T])` | Chainable merge |
| `MergeAll[T]()` | Flatten `Observable[Observable[T]]` by merging |
| `Concat[T](obs ...Observable[T])` | Sequential: complete first, then start second |
| `ConcatWith[T](obs ...Observable[T])` | Chainable concat |
| `ConcatAll[T]()` | Flatten by concatenating sequentially |
| `Zip2[A, B](a, b)` | Pair values from 2 sources into `lo.Tuple2` |
| `Zip3` ... `Zip6` | Zip 3-6 sources into corresponding `lo.Tuple` |
| `ZipWith[A, B](obsB)` | Chainable zip |
| `CombineLatest2[A, B](a, b)` | Emit combined latest when either source emits |
| `CombineLatest3` ... `CombineLatest5` | Combine 3-5 sources |
| `CombineLatestWith[A, B](obsB)` | Chainable combine-latest |
| `Race[T](sources ...Observable[T])` | Emit from whichever source emits first |
| `Amb[T](sources ...)` | Alias for Race |
| `StartWith[T](prefixes ...T)` | Prepend values before source emissions |
| `EndWith[T](suffixes ...T)` | Append values after source completes |
```go
// Zip: pair user with their settings
ro.Zip2(userStream, settingsStream)
// Emits: lo.Tuple2[User, Settings]
// CombineLatest: re-emit whenever either changes
ro.CombineLatest2(priceStream, quantityStream)
// Emits latest (price, quantity) pair each time either updates
```
## Math and Aggregation
| Operator | Purpose |
| --- | --- |
| `Count[T]()` | Count of emitted values |
| `Sum[T Numeric]()` | Sum of all values |
| `Average[T Numeric]()` | Average as float64 |
| `Max[T Numeric]()` | Maximum value |
| `Min[T Numeric]()` | Minimum value |
| `Abs()` | Absolute value (float64) |
| `Ceil()` / `Floor()` / `Round()` / `Trunc()` | Rounding (float64) |
| `CeilWithPrecision(n)` / `FloorWithPrecision(n)` | Rounding with decimal places |
| `Clamp[T](lower, upper)` | Constrain values to range |
## Error Handling
| Operator | Purpose |
| --- | --- |
| `Catch[T](fn func(error) Observable[T])` | Catch error, switch to recovery observable |
| `OnErrorReturn[T](value T)` | Replace error with fallback value |
| `OnErrorResumeNextWith[T](obs ...Observable[T])` | Continue with fallback observables on error |
| `Retry[T]()` | Retry indefinitely on error |
| `RetryWithConfig[T](cfg RetryConfig)` | Retry with max attempts, delay, backoff |
| `ThrowIfEmpty[T](fn func() error)` | Error if stream completes empty |
```go
// RetryConfig for exponential backoff
ro.RetryWithConfig[Response](ro.RetryConfig{
Max: 3,
Delay: time.Second,
BackoffMultiplier: 2.0,
MaxDelay: 10 * time.Second,
})
```
## Timing and Buffering
| Operator | Purpose |
| --- | --- |
| `Delay[T](d time.Duration)` | Delay entire stream |
| `DelayEach[T](d time.Duration)` | Delay between each emission |
| `Timeout[T](d time.Duration)` | Error if no emission within duration |
| `ThrottleTime[T](d time.Duration)` | Ignore values within duration of last |
| `ThrottleWhen[T, t](tick Observable[t])` | Throttle by signal |
| `SampleTime[T](d time.Duration)` | Emit latest value at intervals |
| `SampleWhen[T, t](tick Observable[t])` | Sample by signal |
| `BufferWithCount[T](n int)` | Collect N items, emit as `[]T` |
| `BufferWithTime[T](d time.Duration)` | Collect within time window |
| `BufferWithTimeOrCount[T](n, d)` | Buffer with either condition |
| `BufferWhen[T, B](boundary Observable[B])` | Buffer until signal |
| `WindowWhen[T, B](boundary Observable[B])` | Window as nested observables |
| `Pairwise[T]()` | Emit consecutive pairs |
## Side Effects (Tap / Do)
Execute code without changing the stream. `Do` is an alias for `Tap`.
| Operator | Purpose |
| ------------------------------------- | --------------------------- |
| `Tap[T](onNext, onError, onComplete)` | Side effect on all events |
| `TapOnNext[T](fn func(T))` | Side effect on values |
| `TapOnError[T](fn func(error))` | Side effect on errors |
| `TapOnComplete[T](fn func())` | Side effect on completion |
| `TapOnSubscribe[T](fn func())` | Side effect on subscription |
| `TapOnFinalize[T](fn func())` | Side effect on teardown |
All have `WithContext` variants. `Do`, `DoOnNext`, `DoOnError`, `DoOnComplete`, `DoOnSubscribe`, `DoOnFinalize` are aliases.
## Connectable and Sharing
| Operator | Purpose |
| --- | --- |
| `Share[T]()` | Cold -> hot with reference counting |
| `ShareWithConfig[T](cfg ShareConfig[T])` | Share with reset options |
| `ShareReplay[T](bufferSize int)` | Share + replay last N values to late subscribers |
| `ShareReplayWithConfig[T](n, cfg)` | ShareReplay with reset options |
| `Serialize[T]()` | Queue emissions to ensure serial delivery |
```go
// ShareConfig controls lifecycle reset
ro.ShareConfig[T]{
ResetOnComplete: true, // re-subscribe on complete
ResetOnError: true, // re-subscribe on error
ResetOnReferenceCount: true, // re-subscribe when count drops to 0
}
```
## Context Operators
| Operator | Purpose |
| ---------------------------------------- | ------------------------------- |
| `ContextReset[T](ctx context.Context)` | Replace pipeline context |
| `ContextWithValue[T](key, value any)` | Add value to context |
| `ContextWithTimeout[T](d time.Duration)` | Add timeout |
| `ContextWithDeadline[T](t time.Time)` | Add deadline |
| `ContextMap[T](fn func(ctx) ctx)` | Transform context |
| `ThrowOnContextCancel[T]()` | Error when context is cancelled |
## Conditional Operators
| Operator | Purpose |
| --- | --- |
| `All[T](fn func(T) bool)` | Emit bool: whether all values match |
| `DefaultIfEmpty[T](value T)` | Emit default if source completes empty |
| `ThrowIfEmpty[T](fn func() error)` | Error if empty |
| `Iif[T](pred func() bool, a, b Observable[T])` | Choose between two observables |
| `While[T](cond func() bool)` | Repeat while condition true |
| `DoWhile[T](cond func() bool)` | Execute at least once, repeat while true |
| `SequenceEqual[T](obsB Observable[T])` | Emit bool: whether two streams match |
## Terminal Operators
| Operator | Purpose |
| --- | --- |
| `Collect[T](obs Observable[T]) ([]T, error)` | Block, return all values as slice |
| `CollectWithContext[T](ctx, obs) ([]T, ctx, error)` | Collect with context |
| `ToSlice[T]()` | Operator: emit `[]T` on complete |
| `ToChannel[T](size int)` | Convert to `<-chan Notification[T]` |
| `ToMap[T, K, V](fn func(T) (K, V))` | Collect into map |
## Observer and Subscription
```go
// Full observer (recommended)
observer := ro.NewObserver[T](onNext, onError, onComplete)
// Context-aware observer
observer := ro.NewObserverWithContext[T](onNext, onError, onComplete)
// Convenience shortcuts
observer := ro.OnNext[T](func(v T) { ... })
observer := ro.PrintObserver[T]() // debug: prints all events
observer := ro.NoopObserver[T]() // discard all events
// Subscription lifecycle
sub := observable.Subscribe(observer)
sub.Wait() // block until complete or error
sub.Unsubscribe() // cancel and cleanup
sub.IsActive() // check if still running
sub.GetError() // get terminal error
```
## Scheduling
| Operator | Purpose |
| -------------------------------- | ---------------------------------------- |
| `SubscribeOn[T](bufferSize int)` | Run subscription on async scheduler |
| `ObserveOn[T](bufferSize int)` | Deliver notifications on async scheduler |
## Concurrency Modes
```go
// Safe (default): synchronized emissions
ro.NewObservable[T](fn)
ro.NewSafeObservable[T](fn)
// Unsafe: no synchronization, caller must guarantee single-goroutine access
ro.NewUnsafeObservable[T](fn)
// Eventually safe: allows brief unsynchronized period, then synchronizes
ro.NewEventuallySafeObservable[T](fn)
```
@@ -0,0 +1,259 @@
# Reactive Patterns
Real-world patterns for building production reactive pipelines with samber/ro.
## Pattern 1: Remote Call with Retry and Timeout
Wrap a remote call (HTTP, gRPC, database) with automatic retry, exponential backoff, timeout, and fallback.
```go
result := ro.Pipe3(
fetchUser(userID), // ro.Observable[User] — wraps your remote call
ro.Timeout[User](5*time.Second),
ro.RetryWithConfig[User](ro.RetryConfig{
Max: 3,
Delay: 500 * time.Millisecond,
BackoffMultiplier: 2.0,
MaxDelay: 5 * time.Second,
}),
ro.Catch[User](func(err error) ro.Observable[User] {
log.Printf("remote call failed after retries: %v, using cache", err)
return getCachedUser(userID)
}),
)
user, err := ro.Collect(result)
```
**Why ro over plain calls:** declarative retry + timeout + fallback in 10 lines vs manual for-loops with sleep, context, and error tracking.
## Pattern 2: Continuous Event Stream (Hot Observable)
Share a single long-lived connection (WebSocket, SSE, message queue) across multiple consumers.
```go
// Cold observable wrapping any event stream source
eventStream := ro.NewObservable[TickerEvent](func(ctx context.Context, obs ro.Observer[TickerEvent]) error {
// connect to your stream source (WebSocket, NATS, Kafka, etc.)
for {
event, err := streamSource.Read(ctx)
if err != nil {
return err
}
obs.Next(event)
}
})
// Share: one connection, multiple consumers
shared := ro.Pipe1(eventStream, ro.Share[TickerEvent]())
// Consumer 1: update UI
shared.Subscribe(ro.OnNext(func(e TickerEvent) {
updateDashboard(e)
}))
// Consumer 2: record metrics
shared.Subscribe(ro.OnNext(func(e TickerEvent) {
metrics.RecordTick(e.Symbol, e.Price)
}))
// Consumer 3: alert on threshold
ro.Pipe1(shared, ro.Filter(func(e TickerEvent) bool {
return e.Price > alertThreshold
})).Subscribe(ro.OnNext(sendAlert))
```
The `rohttp` plugin provides WebSocket and HTTP streaming observables (see [Plugin Ecosystem](./plugin-ecosystem.md)).
## Pattern 3: Fan-In from Multiple Sources
Merge events from multiple independent sources, batch, and process.
```go
combined := ro.Pipe2(
ro.Merge(
apiStream,
pushStream,
cronScheduleStream,
),
ro.Distinct[Event](),
ro.BufferWithTimeOrCount[Event](100, 5*time.Second),
ro.Map(func(batch []Event) ProcessResult {
return processBatch(batch)
}),
)
```
**When to use Merge vs Concat vs Zip:**
| Operator | Behavior | Use when |
| --- | --- | --- |
| `Merge` | Interleave: emit from any source as it arrives | Independent streams, order doesn't matter |
| `Concat` | Sequential: finish first source, then start second | Ordered processing, fallback chains |
| `Zip` | Pair: wait for one value from each source | Correlated data (user + settings, request + response) |
| `CombineLatest` | Latest: re-emit combined whenever any source changes | Dependent state (price \* quantity, config + data) |
## Pattern 4: Dependent Data Combination
Combine data from multiple async sources that depend on each other.
```go
// Fetch user and their orders in parallel, combine
profile := ro.Pipe1(
ro.CombineLatest2(
fetchUser(userID),
fetchOrders(userID),
),
ro.Map(func(pair lo.Tuple2[User, []Order]) UserProfile {
return UserProfile{
User: pair.A,
Orders: pair.B,
}
}),
)
```
For independent data where you need exactly one value from each:
```go
// Zip: waits for one value from each, pairs them
configAndData := ro.Zip2(loadConfig(), loadData())
```
## Pattern 5: Running Aggregation with Scan
Maintain running state across stream values — useful for dashboards, analytics, monitoring.
```go
type Stats struct {
Count int
Sum float64
Avg float64
Max float64
}
statsStream := ro.Pipe2(
metricsStream,
ro.Scan(func(acc Stats, v float64) Stats {
acc.Count++
acc.Sum += v
acc.Avg = acc.Sum / float64(acc.Count)
if v > acc.Max {
acc.Max = v
}
return acc
}, Stats{}),
ro.SampleTime[Stats](5*time.Second), // emit stats every 5s
)
```
**Scan vs Reduce:** `Scan` emits every intermediate state (good for live dashboards). `Reduce` emits only the final accumulated value (good for batch summaries).
## Pattern 6: Error Recovery Cascade
Layer multiple error recovery strategies.
```go
resilient := ro.Pipe3(
primaryDataSource,
// Strategy 1: retry transient failures
ro.RetryWithConfig[Data](ro.RetryConfig{
Max: 2,
Delay: time.Second,
}),
// Strategy 2: fall back to secondary source
ro.Catch[Data](func(err error) ro.Observable[Data] {
log.Warn("primary failed, trying secondary", "err", err)
return secondaryDataSource
}),
// Strategy 3: return cached/default value
ro.OnErrorReturn[Data](cachedDefault),
)
```
**Order matters:** retry first (transient errors), then fallback source (persistent errors), then default value (total failure).
## Pattern 7: File System Watcher
React to file changes with debouncing.
```go
import rofsnotify "github.com/samber/ro/plugins/fsnotify"
watcher := ro.Pipe3(
rofsnotify.Watch("/etc/app/config/"),
ro.Filter(func(e fsnotify.Event) bool {
return e.Op&fsnotify.Write != 0
}),
ro.ThrottleTime[fsnotify.Event](2*time.Second), // debounce rapid saves
ro.Map(func(e fsnotify.Event) Config {
return reloadConfig(e.Name)
}),
)
watcher.Subscribe(ro.NewObserver(
func(cfg Config) { applyConfig(cfg) },
func(err error) { log.Error("config watch failed", "err", err) },
func() { log.Info("config watcher stopped") },
))
```
## Pattern 8: Graceful Shutdown
Use context or signal observable to cleanly terminate infinite streams.
```go
import rosignal "github.com/samber/ro/plugins/signal"
// Method 1: OS signal
shutdown := rosignal.Notify(syscall.SIGTERM, syscall.SIGINT)
sub := ro.Pipe1(
workStream,
ro.TakeUntil[Work, os.Signal](shutdown),
).Subscribe(ro.NewObserver(
processWork,
handleError,
func() { log.Info("gracefully stopped") },
))
sub.Wait() // blocks until SIGTERM/SIGINT
// Method 2: Context cancellation
ctx, cancel := context.WithCancel(context.Background())
sub := ro.Pipe2(
workStream,
ro.ContextReset[Work](ctx),
ro.ThrowOnContextCancel[Work](),
).Subscribe(worker)
// Later: cancel() triggers clean shutdown
```
## Pattern 9: Event-Driven Pipeline with Logging
Full production pipeline with observability at each stage.
```go
import roslog "github.com/samber/ro/plugins/observability/slog"
pipeline := ro.Pipe5(
eventSource,
ro.TapOnSubscribe[Event](func() {
slog.Info("pipeline started")
}),
ro.Filter(func(e Event) bool { return e.Valid() }),
roslog.TapOnNext[Event](logger, slog.LevelDebug), // log each event
ro.Map(enrichEvent),
ro.BufferWithTimeOrCount[EnrichedEvent](50, 10*time.Second),
ro.MapErr(func(batch []EnrichedEvent) (Result, error) {
return persistBatch(batch)
}),
ro.TapOnError[Result](func(err error) {
slog.Error("pipeline error", "err", err)
metrics.IncrCounter("pipeline.errors", 1)
}),
ro.RetryWithConfig[Result](ro.RetryConfig{Max: 3, Delay: time.Second}),
)
```
@@ -0,0 +1,152 @@
# Plugin Ecosystem
samber/ro ships 40+ plugins that extend the core library with domain-specific operators. Plugins are separate Go modules — install only what you need.
```bash
go get github.com/samber/ro/plugins/<category>/<name>
```
## Data Manipulation
| Plugin | Import | Purpose |
| --- | --- | --- |
| Bytes | `plugins/bytes` | Byte slice operations on streams |
| Strings | `plugins/strings` | String transformations (split, trim, join) |
| Sort | `plugins/sort` | Sorting operators for ordered streams |
| Strconv | `plugins/strconv` | Type conversion (string <-> numeric) |
| Iter | `plugins/iter` | Go 1.23+ iterator interop |
| SIMD (experimental) | `plugins/exp/simd` | SIMD-accelerated numeric transforms |
## Encoding and Serialization
| Plugin | Import | Purpose |
| ------ | ------------------------- | -------------------------------- |
| JSON | `plugins/encoding/json` | Marshal/unmarshal JSON in stream |
| CSV | `plugins/encoding/csv` | Parse/generate CSV rows |
| Base64 | `plugins/encoding/base64` | Encode/decode Base64 |
| Gob | `plugins/encoding/gob` | Go binary encoding |
```go
import rojson "github.com/samber/ro/plugins/encoding/json"
// Parse JSON stream: []byte -> MyStruct
parsed := ro.Pipe1(rawBytes, rojson.Unmarshal[MyStruct]())
```
## Scheduling
| Plugin | Import | Purpose |
| ------ | -------------- | ------------------------------------- |
| Cron | `plugins/cron` | Emit on cron expressions or intervals |
| ICS | `plugins/ics` | Parse iCal files into event streams |
```go
import rocron "github.com/samber/ro/plugins/cron"
// Emit every day at midnight
daily := rocron.Schedule("0 0 * * *")
```
## Network and I/O
| Plugin | Import | Purpose |
| -------- | ------------------ | ---------------------------------------- |
| HTTP | `plugins/http` | HTTP request operators (GET, POST, etc.) |
| I/O | `plugins/io` | File and stream reading/writing |
| FSNotify | `plugins/fsnotify` | File system change events |
```go
import rofsnotify "github.com/samber/ro/plugins/fsnotify"
// Watch directory for changes
events := rofsnotify.Watch("/var/log/app/")
ro.Pipe1(events, ro.Filter(func(e fsnotify.Event) bool {
return e.Op == fsnotify.Write
})).Subscribe(ro.OnNext(func(e fsnotify.Event) {
log.Println("Modified:", e.Name)
}))
```
## Observability and Logging
| Plugin | Import | Purpose |
| ------- | ------------------------------- | ------------------------------- |
| Log | `plugins/observability/log` | stdlib `log` integration |
| Zap | `plugins/observability/zap` | Uber Zap structured logging |
| Logrus | `plugins/observability/logrus` | Logrus logging |
| Slog | `plugins/observability/slog` | Go 1.21+ `log/slog` integration |
| Zerolog | `plugins/observability/zerolog` | Zerolog logging |
| Sentry | `plugins/observability/sentry` | Sentry error tracking |
| Oops | `plugins/samber/oops` | samber/oops structured errors |
```go
import roslog "github.com/samber/ro/plugins/observability/slog"
// Log all stream events via slog
ro.Pipe1(
dataStream,
roslog.Tap[Data](logger, slog.LevelInfo),
)
```
## Rate Limiting
| Plugin | Import | Purpose |
| ------ | -------------------------- | ---------------------------------- |
| Native | `plugins/ratelimit/native` | Built-in token bucket rate limiter |
| Ulule | `plugins/ratelimit/ulule` | ulule/limiter integration |
## Text Processing
| Plugin | Import | Purpose |
| -------- | ------------------ | ------------------------------------ |
| Regexp | `plugins/regexp` | Regex matching/extraction on streams |
| Template | `plugins/template` | Go text/html template rendering |
## System Integration
| Plugin | Import | Purpose |
| ------- | ---------------- | ------------------------------------------ |
| Process | `plugins/proc` | Process execution operators |
| Signal | `plugins/signal` | OS signal handling (SIGTERM, SIGINT, etc.) |
```go
import rosignal "github.com/samber/ro/plugins/signal"
// Observable that emits on SIGTERM/SIGINT
shutdown := rosignal.Notify(syscall.SIGTERM, syscall.SIGINT)
// Use as TakeUntil signal for graceful shutdown
ro.Pipe1(workStream, ro.TakeUntil[Work, os.Signal](shutdown))
```
## Validation
| Plugin | Import | Purpose |
| --- | --- | --- |
| Ozzo Validation | `plugins/ozzo/ozzo-validation` | Stream-level input validation |
## Testing
| Plugin | Import | Purpose |
| ------- | ----------------- | -------------------------------------- |
| Testify | `plugins/testify` | Test assertions for observable streams |
## Utilities
| Plugin | Import | Purpose |
| ----------- | --------------------- | -------------------------------------- |
| HyperLogLog | `plugins/hyperloglog` | Cardinality estimation on streams |
| Hot | `plugins/samber/hot` | In-memory caching integration |
| PSI | `plugins/samber/psi` | Starvation notifier / pressure metrics |
## Plugin Design Convention
All plugins follow the same pattern:
1. Import the plugin package
2. Use plugin-provided operators in `Pipe` chains
3. Plugin operators return `func(Observable[T]) Observable[R]` — standard operator signature
4. No global state — each operator instance is independent
Plugins are documented individually in their package directories. API details for each plugin are available at `pkg.go.dev/github.com/samber/ro/plugins/...`.
@@ -0,0 +1,183 @@
# Subjects Guide
Subjects are both Observable and Observer — they can receive values (via `Send`, `Error`, `Complete`) and be subscribed to. Subjects are natively **hot**: subscribers share a single execution, and late subscribers only see future emissions (unless replay is configured).
## When to Use Subjects vs Cold Observables
| Use case | Approach |
| --- | --- |
| Data pipeline from a known source (slice, channel, HTTP) | Cold observable (default) |
| Event bus where producers and consumers are decoupled | Subject |
| Multiple consumers need the same WebSocket/ticker stream | Cold observable + `Share()` or `ShareReplay()` |
| Imperatively push values from non-reactive code | Subject |
| Bridge between callback API and reactive pipeline | Subject (receive callbacks, emit to pipeline) |
## Subject Types
### PublishSubject
Standard multicast. Subscribers only see values emitted **after** they subscribe.
```go
subject := ro.NewPublishSubject[string]()
// Subscriber 1
subject.Subscribe(ro.OnNext(func(s string) {
fmt.Println("sub1:", s)
}))
subject.Send("hello") // sub1 sees this
// Subscriber 2 (late)
subject.Subscribe(ro.OnNext(func(s string) {
fmt.Println("sub2:", s)
}))
subject.Send("world") // both see this
subject.Complete()
```
**Use when:** broadcasting events where late subscribers don't need history — UI events, log streams, notifications.
### BehaviorSubject
Replays the **last emitted value** (or the initial value) to every new subscriber immediately on subscription.
```go
subject := ro.NewBehaviorSubject[int](0) // initial value = 0
// Subscriber 1 immediately receives 0
subject.Subscribe(ro.OnNext(func(v int) {
fmt.Println("sub1:", v) // 0, then 42
}))
subject.Send(42)
// Subscriber 2 immediately receives 42 (latest value)
subject.Subscribe(ro.OnNext(func(v int) {
fmt.Println("sub2:", v) // 42
}))
```
**Use when:** subscribers need the current state — config values, connection status, latest price.
### ReplaySubject
Buffers the last **N values** and replays them to every new subscriber.
```go
subject := ro.NewReplaySubject[string](3) // buffer size = 3
subject.Send("a")
subject.Send("b")
subject.Send("c")
subject.Send("d") // "a" evicted from buffer
// Late subscriber receives "b", "c", "d" (last 3)
subject.Subscribe(ro.OnNext(func(s string) {
fmt.Println(s)
}))
```
**Use when:** late subscribers need recent history — chat messages, recent logs, last N stock ticks.
### AsyncSubject
Emits **only the last value** and only when the subject completes. If the subject errors, no value is emitted.
```go
subject := ro.NewAsyncSubject[int]()
subject.Subscribe(ro.NewObserver(
func(v int) { fmt.Println(v) }, // receives 3 only
func(err error) { },
func() { fmt.Println("done") },
))
subject.Send(1)
subject.Send(2)
subject.Send(3)
subject.Complete() // triggers emission of 3, then "done"
```
**Use when:** only the final result matters — computation result, last response in a batch.
### UnicastSubject
Allows exactly **one subscriber**. Buffers values internally until that subscriber connects.
```go
subject := ro.NewUnicastSubject[int](100) // buffer size
subject.Send(1) // buffered
subject.Send(2) // buffered
// Single subscriber receives buffered + future values
subject.Subscribe(ro.OnNext(func(v int) {
fmt.Println(v) // 1, 2, then future values
}))
// Second subscribe would panic or error
```
**Use when:** single consumer with buffering — job queues, request pipelines where exactly one handler processes events.
## Cold to Hot Conversion
When you have a cold observable (e.g. an HTTP request) but need multiple subscribers to share it:
### Share
```go
// Each subscriber to `cold` would trigger a separate HTTP request
cold := httpPlugin.Get[Data](url)
// Share: single execution, multiple subscribers
hot := ro.Pipe1(cold, ro.Share[Data]())
hot.Subscribe(uiObserver) // shares one HTTP call
hot.Subscribe(metricsObserver) // same data, no extra request
```
`Share` uses reference counting: the source subscribes when the first subscriber arrives and unsubscribes when the last one leaves.
### ShareReplay
```go
// Late subscribers get the last N values + future values
hot := ro.Pipe1(cold, ro.ShareReplay[Data](1))
```
### Connectable Observable
For precise control over when the shared subscription starts:
```go
connectable := ro.Connectable[Data](cold)
// Set up subscribers first
connectable.Subscribe(observer1)
connectable.Subscribe(observer2)
// Start the shared execution explicitly
sub, err := connectable.Connect(ctx)
```
## Subject Decision Table
| Subject | Replay | Subscribers | Use case |
| --- | --- | --- | --- |
| `PublishSubject` | None | Many | Event bus, notifications |
| `BehaviorSubject` | Last 1 (+ initial) | Many | Current state, config |
| `ReplaySubject` | Last N | Many | Recent history, chat |
| `AsyncSubject` | Last 1 (on complete) | Many | Final computation result |
| `UnicastSubject` | Buffered (pre-subscribe) | Exactly 1 | Single-consumer queue |
## Common Subject Mistakes
| Mistake | Why | Fix |
| --- | --- | --- |
| Calling `Send()` after `Complete()` | Values are silently dropped — the subject is terminal | Track lifecycle, don't reuse completed subjects |
| Using PublishSubject when late subscribers need history | Late subscribers miss all prior events | Use BehaviorSubject (last 1) or ReplaySubject (last N) |
| Using ReplaySubject with unbounded buffer | Memory grows without limit | Set an explicit `bufferSize` |
| Multiple subscribers on UnicastSubject | Panics or undefined behavior | Use PublishSubject for multicast, UnicastSubject for single consumer |
| Not calling `Complete()` on subjects | Subscribers wait forever, goroutine leak | Always `Complete()` or `Error()` when the source is done |