Sitelet https://pkg.go.dev/github.com/reactivego/rx#Connectable.AutoConnect

rx

package module
v0.3.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 20, 2026 License: MIT Imports: 12 Imported by: 0

README

rx

import "github.com/reactivego/rx"

Go Reference

Package rx provides Reactive Extensions for Go: an API for asynchronous programming built around Observables — typed streams of data — and Operators that transform, filter, and combine them into pipelines.

Prerequisites

You'll need Go 1.23 or later; the implementation depends on generics and on iterator (range-over-func) support.

Quick Start

Create an Observable from some values, print each one, and wait for completion:

package main

import "github.com/reactivego/rx"

func main() {
    rx.From(1, 2, 3).Println().Wait()
}
1
2
3

The element type is inferred; use an explicit type parameter for mixed-type streams:

rx.From[any](1, "hi", 2.3).Println().Wait()
1
hi
2.3

Core Concepts

Observable and Observer

An Observable[T] is a stream of values that are pushed to an Observer[T] over time. Both are plain function types, so any function with the right signature works:

type Observable[T any] func(Observer[T], Scheduler, Subscriber)

type Observer[T any] func(next T, err error, done bool)

An Observable:

  • is a stream of events.
  • assumes zero to many values over time.
  • pushes values to its observers.
  • can take any amount of time to complete (or may never).
  • is cancellable.
  • is lazy: it does nothing until you subscribe.

The observer receives each value with done == false. The stream terminates with exactly one final call where done == true: with err == nil for normal completion, or a non-nil err for failure.

sub := rx.From("a", "b").Subscribe(rx.GoroutineContext(), func(next string, err error, done bool) {
    switch {
    case !done:
        fmt.Println(next)
    case err != nil:
        fmt.Println("error:", err)
    default:
        fmt.Println("complete")
    }
})
sub.Wait()
a
b
complete
Subscriptions

Subscribe returns a Subscription for observing and controlling the stream's lifecycle:

  • Subscribed() reports whether the subscription is still active.
  • Unsubscribe() cancels it.
  • Done() returns a channel that closes on termination, for use in select.
  • Wait() blocks until termination and returns the outcome.
  • Err() returns nil after normal completion, the stream's error after failure, ErrSubscriptionCanceled after Unsubscribe, or ErrSubscriptionActive while still running.

Every error produced by this package wraps the base error Err, so errors.Is(err, rx.Err) identifies package errors.

Context and Schedulers

Subscribe takes a context.Context that serves two purposes:

  1. Cancellation. Canceling the context unsubscribes the subscription; conversely the context is available to source operators through the subscriber, so sources can stop producing when it is canceled.
  2. Scheduler selection. A Scheduler controls where and when the stream runs, and is carried by the context:
// Concurrent: every task runs on its own goroutine.
ctx := rx.GoroutineContext()

// Serial: tasks run sequentially on the calling goroutine (the default).
ctx := rx.SchedulerContextWith(context.Background(), rx.NewScheduler())

GoroutineContext wraps an optional parent context, so cancellation composes with scheduler selection:

ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()

err := rx.Interval[int](100 * time.Millisecond).Println().Wait(rx.GoroutineContext(ctx))
// after one second: err is rx.ErrSubscriptionCanceled

Convenience entry points choose sensible defaults when no scheduler is attached: Wait, First, Last, and Slice use a new serial scheduler, while Go uses the concurrent Goroutine scheduler.

Iterators and Channels

Observables interoperate with Go's native streaming constructs in both directions:

  • Values and All expose an Observable as an iterator usable in a range loop; All pairs each value with its index.
  • Pull and Pull2 adapt an iter.Seq/iter.Seq2 into an Observable.
  • Recv turns a receive channel into an Observable; Send forwards emissions into a channel.
for value, err := range rx.From(1, 2, 3).Values(rx.GoroutineContext()) {
    if err != nil {
        break // terminal error; value is the zero value
    }
    fmt.Println(value)
}
Hot vs Cold Observables
  • Hot Observables emit values regardless of subscription status. Like a live broadcast, values emitted while nobody is subscribed are missed. Examples include system events, mouse movements, or real-time data feeds.

  • Cold Observables begin emission only when subscribed to, so each subscriber receives the complete sequence from the beginning. Examples include file contents, database queries, or HTTP requests executed on demand.

Multicasting

By default a cold Observable restarts for every subscriber. Multicasting shares one upstream subscription among many subscribers:

  • Publish returns a Connectable[T]: subscribers register first, and the source is subscribed only when Connect is called.
  • Behavior is like Publish but replays the latest value (starting from a seed) to each new subscriber, with conflating, non-backpressuring semantics suited to distributing current state.
  • RefCount turns a Connectable back into an ordinary Observable that connects when the first subscriber arrives and disconnects when the last one leaves.
  • AutoConnect connects once a given number of subscribers have subscribed, and reconnects after an error or cancellation when the threshold is crossed again.
  • Share is shorthand for Publish().RefCount().
  • Multicast and Subject provide a raw Observer/Observable pair for pushing values imperatively, with configurable buffering and replay.
state := source.Behavior(initial).RefCount()
// every subscriber immediately observes the current value, then tracks updates

Operators

Operators form a language for expressing programs with Observables. They transform, filter, and combine one or more Observables into new Observables, enabling complex asynchronous workflows through method chaining. Most operators exist both as a method on Observable[T] and as a standalone function returning a Pipe segment for use with Pipe; operators that change the element type (like Map between arbitrary types) are standalone functions.

Index

All converts the Observable into an iterator that yields each value paired with its zero-based index.

Append appends each emitted value to a provided slice while forwarding all emissions.

AsObservable converts between a typed Observable[T] and an Observable[any].

AsObserver converts between a typed Observer[T] and an Observer[any], emitting ErrTypecastFailed when a conversion fails.

Assign stores each emitted value into a provided pointer variable while forwarding all emissions.

AutoConnect connects a Connectable to its source once the specified number of subscribers have subscribed.

AutoUnsubscribe unsubscribes from the source automatically as soon as the stream terminates.

Behavior returns a replay-latest, conflating multicast Connectable seeded with an initial value, suited to distributing current state.

BufferCount collects values into slices of a given size, starting a new buffer every startBufferEvery values.

Catch recovers from an error by continuing with a replacement Observable instead of emitting the error.

CatchError recovers from an error via a selector function that maps the error to a replacement Observable.

CombineAll collects the Observables emitted by a higher-order Observable and combines their latest values into slices.

CombineLatest emits a slice of the latest values of all inputs whenever any input emits; CombineLatest2–CombineLatest5 return strongly typed tuples.

Concat emits all values from each Observable in turn, never overlapping emissions.

ConcatAll flattens a higher-order Observable by subscribing to each inner Observable only after the previous one completes.

ConcatMap projects each value to an inner Observable and concatenates their emissions in order.

ConcatWith appends other Observables, each starting only after the previous completes.

Connect subscribes a Connectable to its source and returns a Subscription; the scheduler can be provided via context.

Connectable is an Observable with delayed connection to its source: observers subscribe first, emission starts on Connect.

Connector is the Connect half of a Connectable.

Count emits a single value: the number of items the source emitted before completing.

Create constructs an Observable from a Creator function called with an incrementing index, bridging imperative code and the reactive pattern.

Creator is the function type used by Create: it returns the next value, an error, and a done flag.

Defer calls a factory to create a fresh Observable for every subscription.

Delay time-shifts each emission by a duration.

DistinctUntilChanged suppresses values that equal the previously emitted value.

Do calls a function for each value passing through the Observable.

ElementAt emits only the nth value.

Empty emits no values and terminates normally.

EndWith emits additional values after the source completes.

Equal returns an equality function for comparable types, for use with DistinctUntilChanged.

Err is the base error joined into every error this package produces; test with errors.Is.

ExhaustAll flattens a higher-order Observable, ignoring inner Observables that arrive while one is still active.

ExhaustMap projects each value to an inner Observable, ignoring source values while the current inner Observable is active.

Filter emits only the values that pass a predicate.

First blocks and returns the first emitted value.

Fprint, Fprintf, Fprintln write each value to an io.Writer while forwarding it.

From creates an Observable from the values passed in.

Go subscribes on the Goroutine scheduler, discards emissions, and returns a Subscription — useful for side effects.

GoroutineContext returns a context carrying the concurrent Goroutine scheduler, optionally wrapping a parent context.

Ignore creates an Observer that discards all emissions.

Interval emits an increasing sequence of numbers spaced by a time interval.

Last blocks and returns the last emitted value.

Map, MapE transform each value by applying a function; the E variant's function may also return an error.

Marshal marshals each value to []byte using a provided function.

Merge combines multiple Observables into one by interleaving their emissions.

MergeAll flattens a higher-order Observable by merging all inner Observables concurrently.

MergeMap projects each value to an inner Observable and merges their emissions.

MergeWith merges other Observables into the source.

Multicast returns an Observer/Observable pair that multicasts pushed values to every subscriber, with blocking or dropping backpressure depending on buffer-size sign.

Must returns its value argument, panicking if the accompanying error is non-nil.

Never emits no values and never terminates.

Of emits a single value and completes.

OnComplete, OnDone, OnError, OnNext register callbacks for the corresponding stream events while forwarding all emissions.

Passthrough forwards all emissions unchanged.

Pipe applies a chain of standalone operator segments to the Observable.

Print, Printf, Println print each value to stdout while forwarding it.

Publish returns a multicasting Connectable for the source Observable.

Pull, Pull2 adapt a Go iterator (iter.Seq, iter.Seq2) into an Observable.

Race mirrors the first Observable to emit or terminate, unsubscribing the others.

RaceWith races the source against other Observables.

Recv creates an Observable from a receive channel.

Reduce, ReduceE accumulate all values into a single final result using a seed and accumulator function.

RefCount converts a Connectable into an Observable that connects on the first subscription and disconnects when the last subscriber leaves.

Repeat resubscribes to the source each time it completes, indefinitely or a given number of times.

Retry, RetryTime resubscribe to the source when it errors, optionally limited and with a backoff schedule.

SampleTime emits the most recent value within periodic time intervals.

Scan, ScanE apply an accumulator to each value, emitting every intermediate result.

Send sends each value into a provided channel while forwarding it.

Share multicasts the source with reference counting: Publish().RefCount().

Skip suppresses the first n values.

Slice blocks and collects all values into a slice.

StartWith synchronously emits the given values on subscription before mirroring the source.

Subject returns an Observer/Observable pair with a replay buffer bounded by age and capacity.

Subscribe attaches an Observer with a context and returns a Subscription; the context carries the scheduler and cancels the subscription.

SubscribeOn specifies the concurrent scheduler the Observable should use when subscribed to.

Take emits only the first n values, then completes.

TakeWhile mirrors values until a predicate becomes false.

Tap observes all emissions with a secondary Observer without altering the stream.

Throw emits no values and terminates with the given error.

Ticker emits timestamps after an initial delay, following a schedule of intervals.

Timer emits an increasing sequence of numbers after an initial delay, following a schedule of intervals.

Tuple2–Tuple5 are the strongly typed tuples produced by the numbered combining operators.

Values converts the Observable into an iterator that yields each value, ending with the terminal error if the stream fails.

Wait subscribes and blocks until completion, returning nil, the stream's error, or a cancellation error.

WithLatestFrom emits a combination of each source value with the latest values from the other Observables; WithLatestFrom2–WithLatestFrom5 return strongly typed tuples.

WithLatestFromAll applies WithLatestFrom to the Observables emitted by a higher-order Observable.

Zip combines multiple Observables in lockstep, pairing the nth values of every input; Zip2–Zip5 return strongly typed tuples and accept a buffer-size option.

ZipAll zips the Observables emitted by a higher-order Observable.

Documentation

Overview

Package rx provides Reactive Extensions, a powerful API for asynchronous programming in Go, built around observables and operators to process streams of data seamlessly.

Example (All)
package main

import (
	"fmt"

	"github.com/reactivego/rx"
)

func main() {
	source := rx.From("ZERO", "ONE", "TWO")

	for next, err := range source.All() {
		if err != nil {
			fmt.Println("Unexpected error:", err)
		}
		fmt.Println(next.First, next.Second)
	}

	fmt.Println("OK")
}
Output:
0 ZERO
1 ONE
2 TWO
OK
Example (AutoConnect)
package main

import (
	"fmt"
	"time"

	"github.com/reactivego/rx"
)

func main() {
	// Create a multicaster hot observable that will emit every 100 milliseconds
	hot := rx.Interval[int](100 * time.Millisecond).Take(10).Publish()
	hotsub := hot.Connect(rx.GoroutineContext())
	defer hotsub.Unsubscribe()
	fmt.Println("Hot observable created and emitting 0,1,2,3,4,5,6 ...")

	// Publish the hot observable again but only Connect to it when 2
	// subscribers have connected.
	source := hot.Take(5).Publish().AutoConnect(2)

	// First subscriber
	sub1 := source.Printf("Subscriber 1: %d\n").Go()
	fmt.Println("First subscriber connected, waiting a bit...")

	// Wait a bit, nothing will emit yet
	time.Sleep(525 * time.Millisecond)

	fmt.Println("Second subscriber connecting, emissions begin!")
	// Second subscriber triggers the connection
	sub2 := source.Printf("Subscriber 2: %d\n").Go()

	// Wait for emissions to complete
	hotsub.Wait()
	sub1.Wait()
	sub2.Wait()

}
Output:
Hot observable created and emitting 0,1,2,3,4,5,6 ...
First subscriber connected, waiting a bit...
Second subscriber connecting, emissions begin!
Subscriber 1: 5
Subscriber 2: 5
Subscriber 1: 6
Subscriber 2: 6
Subscriber 1: 7
Subscriber 2: 7
Subscriber 1: 8
Subscriber 2: 8
Subscriber 1: 9
Subscriber 2: 9
Example (BufferCount)
package main

import (
	"fmt"

	"github.com/reactivego/rx"
)

func main() {
	source := rx.From(0, 1, 2, 3)

	fmt.Println("BufferCount(From(0, 1, 2, 3), 2, 1)")
	rx.BufferCount(source, 2, 1).Println().Wait()

	fmt.Println("BufferCount(From(0, 1, 2, 3), 2, 2)")
	rx.BufferCount(source, 2, 2).Println().Wait()

	fmt.Println("BufferCount(From(0, 1, 2, 3), 2, 3)")
	rx.BufferCount(source, 2, 3).Println().Wait()

	fmt.Println("BufferCount(From(0, 1, 2, 3), 3, 2)")
	rx.BufferCount(source, 3, 2).Println().Wait()

	fmt.Println("BufferCount(From(0, 1, 2, 3), 6, 6)")
	rx.BufferCount(source, 6, 6).Println().Wait()

	fmt.Println("BufferCount(From(0, 1, 2, 3), 2, 0)")
	rx.BufferCount(source, 2, 0).Println().Wait()
}
Output:
BufferCount(From(0, 1, 2, 3), 2, 1)
[0 1]
[1 2]
[2 3]
[3]
BufferCount(From(0, 1, 2, 3), 2, 2)
[0 1]
[2 3]
BufferCount(From(0, 1, 2, 3), 2, 3)
[0 1]
[3]
BufferCount(From(0, 1, 2, 3), 3, 2)
[0 1 2]
[2 3]
BufferCount(From(0, 1, 2, 3), 6, 6)
[0 1 2 3]
BufferCount(From(0, 1, 2, 3), 2, 0)
[0 1]
Example (CommandPattern)
package main

import (
	"time"

	"github.com/reactivego/rx"
)

func main() {
	type Data struct {
		Id  string
		Val int
	}

	// collect is an example of a long running observable producing data.
	collect := func(id string) rx.Observable[Data] {
		take5 := rx.Interval[int](100 * time.Millisecond).Take(5)
		// use rx.Defer to have a scope per subscription available.
		return rx.Defer(func() rx.Observable[Data] {
			// this scope is for storing data per subscription
			// for example say you want to support retries,
			// then this scope will be active per retry.
			return rx.Map(take5, func(val int) Data {
				return Data{Id: id, Val: val}
			})
		})
	}

	// setup a channel to push commands into, where a command is an rx.Observable[Data]
	commands := make(chan rx.Observable[Data], 2)

	// now user MergeAll to run all commands in parallel and merge their output
	sub := rx.MergeAll(rx.Recv(commands)).Println().Go()

	// launch a bunch of commands.
	commands <- collect("hello")
	commands <- collect("world")

	// We need to make sure the command channel is closed.
	close(commands)

	// then wait for everything to play out
	sub.Wait()

}
Output:
{hello 0}
{world 0}
{hello 1}
{world 1}
{hello 2}
{world 2}
{hello 3}
{world 3}
{hello 4}
{world 4}
Example (ConcatAll)
package main

import (
	"fmt"
	"time"

	"github.com/reactivego/rx"
)

func main() {
	source := rx.Empty[rx.Observable[string]]()
	rx.ConcatAll(source).Wait()

	source = rx.Of(rx.Empty[string]())
	rx.ConcatAll(source).Wait()

	req := func(request string, duration time.Duration) rx.Observable[string] {
		req := rx.From(request + " response")
		if duration == 0 {
			return req
		}
		return req.Delay(duration)
	}

	const ms = time.Millisecond

	req1 := req("first", 10*ms)
	req2 := req("second", 20*ms)
	req3 := req("third", 0*ms)
	req4 := req("fourth", 60*ms)

	source = rx.From(req1).ConcatWith(rx.From(req2, req3, req4).Delay(100 * ms))
	rx.ConcatAll(source).Println().Wait()

	fmt.Println("OK")
}
Output:
first response
second response
third response
fourth response
OK
Example (Count)
package main

import (
	"fmt"

	"github.com/reactivego/rx"
)

func main() {
	source := rx.From(1, 2, 3, 4, 5)

	count := source.Count()
	count.Println().Wait()

	emptySource := rx.Empty[int]()
	emptyCount := emptySource.Count()
	emptyCount.Println().Wait()

	fmt.Println("OK")
}
Output:
5
0
OK
Example (ElementAt)
package main

import (
	"github.com/reactivego/rx"
)

func main() {
	rx.From(0, 1, 2, 3, 4).ElementAt(2).Println().Wait()
}
Output:
2
Example (ExhaustAll)
package main

import (
	"fmt"
	"strconv"
	"time"

	"github.com/reactivego/rx"
)

func main() {
	const ms = time.Millisecond

	stream := func(name string, duration time.Duration, count int) rx.Observable[string] {
		return rx.Map(rx.Timer[int](0*ms, duration), func(next int) string {
			return name + "-" + strconv.Itoa(next)
		}).Take(count)
	}

	streams := []rx.Observable[string]{
		stream("a", 20*ms, 3),
		stream("b", 20*ms, 3),
		stream("c", 20*ms, 3),
		rx.Empty[string](),
	}

	streamofstreams := rx.Map(rx.Timer[int](20*ms, 30*ms, 250*ms, 100*ms).Take(4), func(next int) rx.Observable[string] {
		return streams[next]
	})

	err := rx.ExhaustAll(streamofstreams).Println().Wait()

	if err == nil {
		fmt.Println("success")
	}
}
Output:
a-0
a-1
a-2
c-0
c-1
c-2
success
Example (Marshal)
package main

import (
	"encoding/json"

	"github.com/reactivego/rx"
)

func main() {
	type R struct {
		A string `json:"a"`
		B string `json:"b"`
	}

	b2s := func(data []byte) string { return string(data) }

	rx.Map(rx.Of(R{"Hello", "World"}).Marshal(json.Marshal), b2s).Println().Wait()
}
Output:
{"a":"Hello","b":"World"}
Example (MergeMap)
package main

import (
	"context"
	"fmt"

	"github.com/reactivego/rx"
)

func main() {
	source := rx.From("https://reactivego.io", "https://github.com/reactivego")

	merged := rx.MergeMap(source, func(next string) rx.Observable[string] {
		fakeFetchData := rx.Of(fmt.Sprintf("content of %q", next))
		return fakeFetchData
	})

	merged.Println().Go(context.Background()).Wait()

}
Output:
content of "https://reactivego.io"
content of "https://github.com/reactivego"
Example (MergeMapSubject)
package main

import (
	"fmt"
	"time"

	"github.com/reactivego/rx"
)

func main() {
	source := rx.From("https://google.com", "https://reactivego.io", "https://github.com/reactivego")

	merged := rx.MergeMap(source, func(next string) rx.Observable[string] {
		fakeFetchData := rx.Of(fmt.Sprintf("content of %q", next))
		return fakeFetchData
	})

	// subject remembers last 2 emits by the observer for an hour.
	observer, subject := rx.Subject[string](time.Hour, 2)

	// First subscriber starts before merged completes, so sees all emits live
	wait := subject.Println().Go()
	merged.Tap(observer).Go().Wait()
	wait.Wait()

	// Sees only last 2 emits
	subject.Println().Go().Wait()

	// Sees only last 2 emits
	subject.Println().Go().Wait()

}
Output:
content of "https://google.com"
content of "https://reactivego.io"
content of "https://github.com/reactivego"
content of "https://reactivego.io"
content of "https://github.com/reactivego"
content of "https://reactivego.io"
content of "https://github.com/reactivego"
Example (Multicast)
package main

import (
	"context"
	"fmt"

	"github.com/reactivego/rx"
)

func main() {
	serial := rx.NewScheduler()
	ctx := rx.SchedulerContextWith(context.Background(), serial)

	in, out := rx.Multicast[int](1)

	// Ignore everything before any subscriptions, including the last!
	in.Next(-2)
	in.Next(-1)
	in.Next(0)
	in.Next(1)

	// Schedule the subsequent emits in a loop. This will be the first task to
	// run on the serial scheduler after the subscriptions have been added.
	serial.ScheduleLoop(2, func(index int, again func(next int)) {
		if index < 4 {
			in.Next(index)
			again(index + 1)
		} else {
			in.Done(rx.Err)
		}
	})

	// Add a couple of subscriptions
	sub1 := out.Println().Go(ctx)
	sub2 := out.Println().Go(ctx)

	// Let the scheduler run and wait for all of its scheduled tasks to finish.
	serial.Wait()
	fmt.Println(sub1.Wait())
	fmt.Println(sub2.Wait())
}
Output:
2
2
3
3
rx
rx
Example (MulticastDrop)
package main

import (
	"context"
	"fmt"

	"github.com/reactivego/rx"
)

func main() {
	serial := rx.NewScheduler()
	ctx := rx.SchedulerContextWith(context.Background(), serial)

	const onBackpressureDrop = -1

	// multicast with backpressure handling set to dropping incoming
	// items that don't fit in the buffer once it has filled up.
	in, out := rx.Multicast[int](1 * onBackpressureDrop)

	// ignore everything before any subscriptions, including the last!
	in.Next(-2)
	in.Next(-1)
	in.Next(0)
	in.Next(1)

	// add a couple of subscriptions
	sub1 := out.Println().Go(ctx)
	sub2 := out.Println().Go(ctx)

	in.Next(2)      // accepted: buffer not full
	in.Next(3)      // dropped: buffer full
	in.Done(rx.Err) // dropped: buffer full

	serial.Wait()
	fmt.Println(sub1.Wait())
	fmt.Println(sub2.Wait())
}
Output:
2
2
<nil>
<nil>
Example (Race)
package main

import (
	"errors"
	"fmt"
	"time"

	"github.com/reactivego/rx"
)

func main() {
	const ms = time.Millisecond

	req := func(request string, duration time.Duration) rx.Observable[string] {
		return rx.From(request + " response").Delay(duration)
	}

	req1 := req("first", 50*ms)
	req2 := req("second", 10*ms)
	req3 := req("third", 60*ms)

	rx.Race(req1, req2, req3).Println().Wait()

	err := func(text string, duration time.Duration) rx.Observable[int] {
		return rx.Throw[int](errors.New(text + " error")).Delay(duration)
	}

	err1 := err("first", 10*ms)
	err2 := err("second", 20*ms)
	err3 := err("third", 30*ms)

	fmt.Println(rx.Race(err1, err2, err3).Wait(rx.GoroutineContext()))
}
Output:
second response
first error
Example (Retry)
package main

import (
	"fmt"

	"github.com/reactivego/rx"
)

func main() {
	var first error = rx.Err
	a := rx.Create(func(index int) (next int, err error, done bool) {
		if index < 3 {
			return index, nil, false
		}
		err, first = first, nil
		return 0, err, true
	})
	err := a.Retry().Println().Wait()
	fmt.Println(first == nil)
	fmt.Println(err)
}
Output:
0
1
2
0
1
2
true
<nil>
Example (Share)
package main

import (
	"context"

	"github.com/reactivego/rx"
)

func main() {
	serial := rx.NewScheduler()

	shared := rx.From(1, 2, 3).Share()

	ctx := rx.SchedulerContextWith(context.Background(), serial)
	shared.Println().Go(ctx)
	shared.Println().Go(ctx)
	shared.Println().Go(ctx)

	serial.Wait()
}
Output:
1
1
1
2
2
2
3
3
3
Example (Skip)
package main

import (
	"github.com/reactivego/rx"
)

func main() {
	rx.From(1, 2, 3, 4, 5).Skip(2).Println().Wait()
}
Output:
3
4
5
Example (Subject)
package main

import (
	"context"
	"fmt"

	"github.com/reactivego/rx"
)

func main() {
	serial := rx.NewScheduler()
	ctx := rx.SchedulerContextWith(context.Background(), serial)

	// subject collects emits when there are no subscriptions active.
	in, out := rx.Subject[int](0, 1)

	// ignore everything before any subscriptions, except the last because buffer size is 1
	in.Next(-2)
	in.Next(-1)
	in.Next(0)
	in.Next(1)

	// add a couple of subscriptions
	sub1 := out.Println().Go(ctx)
	sub2 := out.Println().Go(ctx)

	// schedule the subsequent emits on the serial scheduler otherwise these calls
	// will block because the buffer is full.
	// subject will detect usage of scheduler on observable side and use it on the
	// observer side to keep the data flow through the subject going.
	serial.Schedule(func() {
		in.Next(2)
		in.Next(3)
		in.Done(rx.Err)
	})

	serial.Wait()
	fmt.Println(sub1.Wait())
	fmt.Println(sub2.Wait())
}
Output:
1
1
2
2
3
3
rx
rx
Example (SwitchAll)
package main

import (
	"fmt"
	"time"

	"github.com/reactivego/rx"
)

func main() {
	const ms = time.Millisecond

	// Emit 0,1,2,3 with 42ms in between
	interval42x4 := rx.Interval[int](42 * ms).Take(4)

	// Emit 0,1,2,3 with 16ms in between
	interval16x4 := rx.Interval[int](16 * ms).Take(4)

	overlapping := rx.Map(interval42x4, func(next int) rx.Observable[int] {
		return interval16x4
	})

	err := rx.SwitchAll(overlapping).Println().Wait(rx.GoroutineContext())

	if err == nil {
		fmt.Println("success")
	}
}
Output:
0
1
0
1
0
1
0
1
2
3
success
Example (SwitchMap)
package main

import (
	"fmt"
	"time"

	"github.com/reactivego/rx"
)

func main() {
	const ms = time.Millisecond

	webreq := func(request string, duration time.Duration) rx.Observable[string] {
		return rx.From(request + " result").Delay(duration)
	}

	first := webreq("first", 50*ms)
	second := webreq("second", 10*ms)
	latest := webreq("latest", 50*ms)

	switchmap := rx.SwitchMap(rx.Interval[int](20*ms).Take(3), func(i int) rx.Observable[string] {
		switch i {
		case 0:
			return first
		case 1:
			return second
		case 2:
			return latest
		default:
			return rx.Empty[string]()
		}
	})

	err := switchmap.Println().Wait()
	if err == nil {
		fmt.Println("success")
	}
}
Output:
second result
latest result
success
Example (Values)
package main

import (
	"fmt"

	"github.com/reactivego/rx"
)

func main() {
	source := rx.From(1, 3, 5)

	// Why choose the Goroutine concurrent scheduler?
	// An observable can actually be at the root of a tree
	// of separately running observables that have their
	// responses merged. The Goroutine scheduler allows
	// these observables to run concurrently.

	// run the observable on 1 or more goroutines
	for i := range source.Values(rx.GoroutineContext()) {
		// This is called from a newly created goroutine
		fmt.Println(i)
	}

	// run the observable on the current goroutine
	for i := range source.Values() {
		fmt.Println(i)
	}

	fmt.Println("OK")
}
Output:
1
3
5
1
3
5
OK

Index

Examples

Constants

This section is empty.

Variables

View Source
var Err = errors.New("rx")

Err declares the base error that is joined with every error returned by this package. It serves as the foundation for error types in the reactive extensions library, allowing for error type checking and custom error creation with consistent taxonomy.

View Source
var ErrInvalidCount = errors.Join(Err, errors.New("invalid count"))
View Source
var ErrOutOfSubjectSubscriptions = errors.Join(Err, errors.New("out of subject subscriptions"))
View Source
var ErrRepeatCountInvalid = errors.Join(Err, errors.New("repeat count invalid"))
View Source
var ErrSubscriptionActive = errors.Join(Err, errors.New("subscription active"))

ErrSubscriptionActive is the error returned by Err() when the subscription is still active and has not yet completed or been canceled.

View Source
var ErrSubscriptionCanceled = errors.Join(Err, errors.New("subscription canceled"))

ErrSubscriptionCanceled is the error returned by Wait() and Err() when the subscription was canceled by calling Unsubscribe() on the Subscription. This indicates the subscription was terminated by the subscriber rather than by the observable completing normally or with an error.

View Source
var ErrTypecastFailed = errors.Join(Err, errors.New("typecast failed"))

ErrTypecastFailed is returned when a type conversion fails during observer operations, typically when using AsObserver() to convert between generic and typed observers.

View Source
var ErrZipBufferOverflow = errors.Join(Err, errors.New("zip buffer overflow"))
View Source
var Goroutine = scheduler.Goroutine

Goroutine is a concurrent scheduler that runs every task on its own goroutine.

View Source
var NewScheduler = scheduler.New

NewScheduler returns a new serial scheduler that runs tasks sequentially on the calling goroutine.

View Source
var SchedulerContextWith = scheduler.ContextWith

SchedulerContextWith returns a context with the scheduler attached, replacing any scheduler already present.

View Source
var SchedulerFromContext = scheduler.FromContext

SchedulerFromContext returns the Scheduler attached to the context, or nil and false when none is attached.

Functions

func Equal added in v0.2.0

func Equal[T comparable]() func(T, T) bool

func GoroutineContext added in v0.3.0

func GoroutineContext(ctx ...context.Context) context.Context

GoroutineContext returns a context that is guaranteed to have a scheduler attached. If no context or a nil context is provided, then context.Background() is used. If the context already has a scheduler it is replaced with scheduler.Goroutine.

func Multicast added in v0.2.0

func Multicast[T any](size int) (Observer[T], Observable[T])

Multicast returns both an Observer and and Observable. The returned Observer is used to send items into the Multicast. The returned Observable is used to subscribe to the Multicast. The Multicast multicasts items send through the Observer to every Subscriber of the Observable.

size  size of the item buffer, number of items kept to replay to a new Subscriber.

Backpressure handling depends on the sign of the size argument. For positive size the multicast will block when one of the subscribers lets the buffer fill up. For negative size the multicast will drop items on the blocking subscriber, allowing the others to keep on receiving values. For hot observables dropping is preferred.

func Must added in v0.2.0

func Must[T any](t T, err error) T

func Subject

func Subject[T any](age time.Duration, capacity ...int) (Observer[T], Observable[T])

Subject returns both an Observer and and Observable. The returned Observer is used to send items into the Subject. The returned Observable is used to subscribe to the Subject. The Subject multicasts items send through the Observer to every Subscriber of the Observable.

age     max age to keep items in order to replay them to a new Subscriber (0 = no max age).
[size]  size of the item buffer, number of items kept to replay to a new Subscriber.
[cap]   capacity of the item buffer, number of items that can be observed before blocking.
[scap]  capacity of the subscription list, max number of simultaneous subscribers.

Types

type ConcurrentScheduler added in v0.2.0

type ConcurrentScheduler = scheduler.ConcurrentScheduler

ConcurrentScheduler is the interface for schedulers that run tasks concurrently.

type Connectable

type Connectable[T any] struct {
	Observable[T]
	Connector
}

Connectable[T] is an Observable[T] that provides delayed connection to its source. It combines both Observable and Connector interfaces:

  • Subscribe: Allows consumers to register for notifications from this Observable
  • Connect: Triggers the actual subscription to the underlying source Observable

The key feature of Connectable[T] is that it doesn't subscribe to its source until the Connect method is explicitly called, allowing multiple observers to subscribe before the source begins emitting items (multicast behavior).

func (Connectable[T]) AutoConnect added in v0.2.0

func (connectable Connectable[T]) AutoConnect(count int) Observable[T]

AutoConnect returns an Observable that automatically connects to the Connectable source when a specified number of subscribers subscribe to it.

When the specified number of subscribers (count) is reached, the Connectable source is connected, allowing it to start emitting items. The connection is shared among all subscribers. When all subscribers unsubscribe, the connection is terminated.

If count is less than 1, it returns an Observable that emits an ErrInvalidCount error.

func (Connectable[T]) RefCount added in v0.2.0

func (connectable Connectable[T]) RefCount() Observable[T]

RefCount converts a Connectable Observable into a standard Observable that automatically connects when the first subscriber subscribes and disconnects when the last subscriber unsubscribes.

When the first subscriber subscribes to the resulting Observable, it automatically calls Connect() on the source Connectable Observable. The connection is shared among all subscribers. When the last subscriber unsubscribes, the connection is automatically closed.

This is useful for efficiently sharing expensive resources (like network connections) among multiple subscribers.

type Connector added in v0.2.3

type Connector func(Scheduler, Subscriber)

Connector provides the Connect method for a Connectable[T].

func (Connector) Connect added in v0.2.3

func (connect Connector) Connect(ctx ...context.Context) Subscription

Connect instructs a Connectable[T] to subscribe to its source and begin emitting items to its subscribers. The scheduler can be provided via context using SchedulerContextWith, otherwise a new serial scheduler is created.

type Creator added in v0.2.0

type Creator[T any] func(index int) (Next T, Err error, Done bool)

Creator[T] is a function type that generates values for an Observable stream.

The Creator function receives a zero-based index for the current iteration and returns a tuple containing:

  • Next: The next value to emit (of type T)
  • Err: Any error that occurred during value generation
  • Done: Boolean flag indicating whether the sequence is complete

When Done is true, the Observable will complete after emitting any provided error. When Err is non-nil, the Observable will emit the error and then complete.

type Float added in v0.2.2

type Float interface {
	~float32 | ~float64
}

Float is a constraint that permits any floating-point type. If future releases of Go add new predeclared floating-point types, this constraint will be modified to include them.

type Integer added in v0.2.2

type Integer interface {
	Signed | Unsigned
}

Integer is a constraint that permits any integer type. If future releases of Go add new predeclared integer types, this constraint will be modified to include them.

type MaxBufferSizeOption added in v0.2.1

type MaxBufferSizeOption = func(*int)

MaxBufferSizeOption is a function type used for configuring the maximum buffer size of an observable stream.

func WithMaxBufferSize added in v0.2.1

func WithMaxBufferSize(n int) MaxBufferSizeOption

WithMaxBufferSize creates a MaxBufferSizeOption that sets the maximum buffer size to n. This option is typically used when creating new observables to control memory usage.

type Observable

type Observable[T any] func(Observer[T], Scheduler, Subscriber)

func AsObservable added in v0.2.0

func AsObservable[T any](observable Observable[any]) Observable[T]

func BufferCount added in v0.2.0

func BufferCount[T any](observable Observable[T], bufferSize, startBufferEvery int) Observable[[]T]

func CombineAll added in v0.2.0

func CombineAll[T any](observable Observable[Observable[T]]) Observable[[]T]

func CombineLatest

func CombineLatest[T any](observables ...Observable[T]) Observable[[]T]

func CombineLatest2 added in v0.2.1

func CombineLatest2[T, U any](first Observable[T], second Observable[U]) Observable[Tuple2[T, U]]

func CombineLatest3 added in v0.2.1

func CombineLatest3[T, U, V any](first Observable[T], second Observable[U], third Observable[V]) Observable[Tuple3[T, U, V]]

func CombineLatest4 added in v0.2.1

func CombineLatest4[T, U, V, W any](first Observable[T], second Observable[U], third Observable[V], fourth Observable[W]) Observable[Tuple4[T, U, V, W]]

func CombineLatest5 added in v0.2.1

func CombineLatest5[T, U, V, W, X any](first Observable[T], second Observable[U], third Observable[V], fourth Observable[W], fifth Observable[X]) Observable[Tuple5[T, U, V, W, X]]

func Concat

func Concat[T any](observables ...Observable[T]) Observable[T]

func ConcatAll added in v0.2.0

func ConcatAll[T any](observable Observable[Observable[T]]) Observable[T]

func ConcatMap added in v0.2.0

func ConcatMap[T, U any](observable Observable[T], project func(T) Observable[U]) Observable[U]

func Create

func Create[T any](create Creator[T]) Observable[T]

Create constructs a new Observable from a Creator function.

The Creator function is called repeatedly with an incrementing index value, and returns a tuple of (next value, error, done flag). The Observable will continue producing values until either:

  1. The Creator signals completion by returning done=true
  2. The Observer unsubscribes
  3. The Creator returns an error (which will be emitted with done=true)

This function provides a bridge between imperative code and the reactive Observable pattern.

func Defer

func Defer[T any](factory func() Observable[T]) Observable[T]

Defer creates an Observable that will use the provided factory function to create a new Observable every time it's subscribed to. This is useful for creating cold Observables or for delaying expensive Observable creation until subscription time.

func Empty

func Empty[T any]() Observable[T]

func ExhaustAll added in v0.2.0

func ExhaustAll[T any](observable Observable[Observable[T]]) Observable[T]

func ExhaustMap added in v0.2.0

func ExhaustMap[T, U any](observable Observable[T], project func(T) Observable[U]) Observable[U]

func From

func From[T any](slice ...T) Observable[T]

func Interval

func Interval[T Integer | Float](interval time.Duration) Observable[T]

func Map added in v0.2.0

func Map[T, U any](observable Observable[T], project func(T) U) Observable[U]

func MapE added in v0.2.0

func MapE[T, U any](observable Observable[T], project func(T) (U, error)) Observable[U]

func Merge

func Merge[T any](observables ...Observable[T]) Observable[T]

func MergeAll added in v0.2.0

func MergeAll[T any](observable Observable[Observable[T]]) Observable[T]

func MergeMap added in v0.2.0

func MergeMap[T, U any](observable Observable[T], project func(T) Observable[U]) Observable[U]

func Never

func Never[T any]() Observable[T]

func Of

func Of[T any](value T) Observable[T]

func Pull added in v0.2.2

func Pull[T any](seq iter.Seq[T]) Observable[T]

func Pull2 added in v0.2.2

func Pull2[T, U any](seq iter.Seq2[T, U]) Observable[Tuple2[T, U]]

func Race added in v0.2.0

func Race[T any](observables ...Observable[T]) Observable[T]

func Recv added in v0.2.0

func Recv[T any](ch <-chan T) Observable[T]

func Reduce added in v0.2.0

func Reduce[T, U any](observable Observable[T], seed U, accumulator func(acc U, next T) U) Observable[U]

func ReduceE added in v0.2.1

func ReduceE[T, U any](observable Observable[T], seed U, accumulator func(acc U, next T) (U, error)) Observable[U]

func Scan added in v0.2.0

func Scan[T, U any](observable Observable[T], seed U, accumulator func(acc U, next T) U) Observable[U]

func ScanE added in v0.2.0

func ScanE[T, U any](observable Observable[T], seed U, accumulator func(acc U, next T) (U, error)) Observable[U]

func SwitchAll added in v0.2.0

func SwitchAll[T any](observable Observable[Observable[T]]) Observable[T]

func SwitchMap added in v0.2.0

func SwitchMap[T, U any](o Observable[T], project func(T) Observable[U]) Observable[U]

func Throw

func Throw[T any](err error) Observable[T]

func Ticker

func Ticker(initialDelay time.Duration, intervals ...time.Duration) Observable[time.Time]

Ticker creates an ObservableTime that emits a sequence of timestamps after an initialDelay has passed. Subsequent timestamps are emitted using a schedule of intervals passed in. If only the initialDelay is given, Ticker will emit only once.

func Timer

func Timer[T Integer | Float](initialDelay time.Duration, intervals ...time.Duration) Observable[T]

func WithLatestFrom added in v0.2.0

func WithLatestFrom[T any](observables ...Observable[T]) Observable[[]T]

func WithLatestFrom2 added in v0.2.1

func WithLatestFrom2[T, U any](first Observable[T], second Observable[U]) Observable[Tuple2[T, U]]

func WithLatestFrom3 added in v0.2.1

func WithLatestFrom3[T, U, V any](first Observable[T], second Observable[U], third Observable[V]) Observable[Tuple3[T, U, V]]

func WithLatestFrom4 added in v0.2.1

func WithLatestFrom4[T, U, V, W any](first Observable[T], second Observable[U], third Observable[V], fourth Observable[W]) Observable[Tuple4[T, U, V, W]]

func WithLatestFrom5 added in v0.2.1

func WithLatestFrom5[T, U, V, W, X any](first Observable[T], second Observable[U], third Observable[V], fourth Observable[W], fifth Observable[X]) Observable[Tuple5[T, U, V, W, X]]

func WithLatestFromAll added in v0.2.0

func WithLatestFromAll[T any](observable Observable[Observable[T]]) Observable[[]T]

func Zip added in v0.2.1

func Zip[T any](observables ...Observable[T]) Observable[[]T]

func Zip2 added in v0.2.1

func Zip2[T, U any](first Observable[T], second Observable[U], options ...MaxBufferSizeOption) Observable[Tuple2[T, U]]

func Zip3 added in v0.2.1

func Zip3[T, U, V any](first Observable[T], second Observable[U], third Observable[V], options ...MaxBufferSizeOption) Observable[Tuple3[T, U, V]]

func Zip4 added in v0.2.1

func Zip4[T, U, V, W any](first Observable[T], second Observable[U], third Observable[V], fourth Observable[W], options ...MaxBufferSizeOption) Observable[Tuple4[T, U, V, W]]

func Zip5 added in v0.2.1

func Zip5[T, U, V, W, X any](first Observable[T], second Observable[U], third Observable[V], fourth Observable[W], fifth Observable[X], options ...MaxBufferSizeOption) Observable[Tuple5[T, U, V, W, X]]

func ZipAll added in v0.2.1

func ZipAll[T any](observable Observable[Observable[T]], options ...MaxBufferSizeOption) Observable[[]T]

func (Observable[T]) All

func (observable Observable[T]) All(ctx ...context.Context) iter.Seq2[Tuple2[int, T], error]

All converts an Observable stream into an iterator sequence that pairs each element with its index. It returns an iter.Seq2[Tuple2[int, T], error] which yields each element along with its position in the sequence as a Tuple2.

func (Observable[T]) Append added in v0.2.1

func (observable Observable[T]) Append(slice *[]T) Observable[T]

Append is a method variant of the Append function that appends each emitted value to the provided slice while forwarding all emissions to downstream operators. This is a convenience method that calls the standalone Append function.

func (Observable[T]) AsObservable

func (observable Observable[T]) AsObservable() Observable[any]

func (Observable[T]) Assign added in v0.2.0

func (observable Observable[T]) Assign(value *T) Observable[T]

Assign is a method version of the Assign function. It assigns every next value that is not done to the provided variable.

func (Observable[T]) AutoUnsubscribe

func (observable Observable[T]) AutoUnsubscribe() Observable[T]

func (Observable[T]) Behavior added in v0.2.8

func (observable Observable[T]) Behavior(seed T) Connectable[T]

Behavior returns a multicasting Observable[T] for an underlying Observable[T] as a Connectable[T] type, like Observable.Publish, but with replay-latest, conflating, non-backpressuring semantics suited to distributing current state:

  • Replay-latest on subscribe. A new Subscriber immediately observes the current value — the construction seed until the source emits, then the most recent source value.
  • Conflation for live Subscribers. A Subscriber that falls behind skips intermediate values and converges on the latest; it is never stranded on a stale value (as a dropping FIFO buffer would be) and never blocks the source (as a blocking FIFO buffer would, the way Multicast does).
  • No backpressure to the source. The source is never blocked by a slow Subscriber, so the source's own nature governs what Subscribers see: a fast cold source (e.g. From over many values) on a serial scheduler may run to completion before any Subscriber is scheduled, so Subscribers observe only the final value and completion; a source that emits over time lets Subscribers track it, modulo conflation when one lags.

Composes as a Connectable, so the connection lifecycle is managed by the standard combinators:

source.Behavior(seed).RefCount()                 // count-free; replays latest
source.Behavior(seed).AutoConnect(n)             // connect after n subscribers
source.StartWith(seed).Behavior(seed).RefCount() // StartWith composes too

Unlike Observable.Publish, whose Multicast core does not replay, Behavior's core delivers the current value to a late Subscriber, which is what makes Connectable.RefCount over Behavior count-free: a Subscriber that attaches after Connect still receives the latest value rather than nothing.

Behavior is scheduler-orthogonal. The conflating value cell and terminal latch are lock-light and shared; only the per-Subscriber wait differs by scheduler kind — a concurrent Subscriber blocks on a wake signal on its own goroutine, while a serial (trampoline) Subscriber polls by rescheduling, since a serial scheduler has no park/unpark and a blocking receiver would freeze the single goroutine it shares with the source.

func (Observable[T]) Catch

func (observable Observable[T]) Catch(other Observable[T]) Observable[T]

func (Observable[T]) CatchError

func (observable Observable[T]) CatchError(selector func(err error, caught Observable[T]) Observable[T]) Observable[T]

func (Observable[T]) ConcatWith

func (observable Observable[T]) ConcatWith(others ...Observable[T]) Observable[T]

func (Observable[T]) Count

func (observable Observable[T]) Count() Observable[int]

func (Observable[T]) Delay

func (observable Observable[T]) Delay(duration time.Duration) Observable[T]

func (Observable[T]) DistinctUntilChanged

func (observable Observable[T]) DistinctUntilChanged(equal func(T, T) bool) Observable[T]

func (Observable[T]) Do

func (observable Observable[T]) Do(f func(T)) Observable[T]

func (Observable[T]) ElementAt

func (observable Observable[T]) ElementAt(n int) Observable[T]

func (Observable[T]) EndWith added in v0.2.2

func (observable Observable[T]) EndWith(values ...T) Observable[T]

func (Observable[T]) Filter

func (observable Observable[T]) Filter(predicate func(T) bool) Observable[T]

func (Observable[T]) First

func (observable Observable[T]) First(ctx ...context.Context) (value T, err error)

func (Observable[T]) Fprint added in v0.2.0

func (observable Observable[T]) Fprint(out io.Writer) Observable[T]

func (Observable[T]) Fprintf added in v0.2.0

func (observable Observable[T]) Fprintf(out io.Writer, format string) Observable[T]

func (Observable[T]) Fprintln added in v0.2.0

func (observable Observable[T]) Fprintln(out io.Writer) Observable[T]

func (Observable[T]) Go added in v0.2.0

func (observable Observable[T]) Go(ctx ...context.Context) Subscription

Go subscribes to the observable and starts execution on a separate goroutine. It ignores all emissions from the observable sequence, making it useful when you only care about side effects and not the actual values. By default, it uses the Goroutine scheduler, but a scheduler can be provided via context using SchedulerContextWith. Returns a Subscription that can be used to cancel the subscription when no longer needed.

func (Observable[T]) Last

func (observable Observable[T]) Last(ctx ...context.Context) (value T, err error)

func (Observable[T]) Map

func (observable Observable[T]) Map(project func(T) any) Observable[any]

func (Observable[T]) MapE added in v0.2.1

func (observable Observable[T]) MapE(project func(T) (any, error)) Observable[any]

func (Observable[T]) Marshal added in v0.2.0

func (observable Observable[T]) Marshal(marshal func(any) ([]byte, error)) Observable[[]byte]

func (Observable[T]) MergeWith

func (observable Observable[T]) MergeWith(others ...Observable[T]) Observable[T]

func (Observable[T]) OnComplete added in v0.2.2

func (observable Observable[T]) OnComplete(f func()) Observable[T]

func (Observable[T]) OnDone added in v0.2.2

func (observable Observable[T]) OnDone(f func(error)) Observable[T]

func (Observable[T]) OnError added in v0.2.2

func (observable Observable[T]) OnError(f func(error)) Observable[T]

func (Observable[T]) OnNext added in v0.2.2

func (observable Observable[T]) OnNext(f func(T)) Observable[T]

func (Observable[T]) Passthrough added in v0.2.0

func (observable Observable[T]) Passthrough() Observable[T]

func (Observable[T]) Pipe added in v0.2.0

func (observable Observable[T]) Pipe(segments ...Pipe[T]) Observable[T]

func (Observable[T]) Print added in v0.2.0

func (observable Observable[T]) Print() Observable[T]

func (Observable[T]) Printf added in v0.2.0

func (observable Observable[T]) Printf(format string) Observable[T]

func (Observable[T]) Println

func (observable Observable[T]) Println() Observable[T]

func (Observable[T]) Publish

func (observable Observable[T]) Publish() Connectable[T]

Publish returns a multicasting Observable[T] for an underlying Observable[T] as a Connectable[T] type.

func (Observable[T]) RaceWith added in v0.2.0

func (observable Observable[T]) RaceWith(others ...Observable[T]) Observable[T]

func (Observable[T]) Repeat

func (observable Observable[T]) Repeat(count ...int) Observable[T]

Repeat emits the items emitted by the source Observable repeatedly.

Parameters:

  • count: Optional. The number of repetitions:
  • If omitted: The source Observable is repeated indefinitely
  • If 0: Returns an empty Observable
  • If negative: Returns an Observable that emits an error
  • If multiple count values: Returns an Observable that emits an error

The resulting Observable will subscribe to the source Observable repeatedly each time the source completes, up to the specified count.

func (Observable[T]) Retry

func (observable Observable[T]) Retry(limit ...int) Observable[T]

func (Observable[T]) RetryTime added in v0.2.0

func (observable Observable[T]) RetryTime(backoff func(int) time.Duration, limit ...int) Observable[T]

func (Observable[T]) SampleTime

func (observable Observable[T]) SampleTime(window time.Duration) Observable[T]

SampleTime emits the most recent item emitted by an Observable within periodic time intervals.

func (Observable[T]) Send added in v0.2.0

func (observable Observable[T]) Send(ch chan<- T) Observable[T]

func (Observable[T]) Share added in v0.2.0

func (observable Observable[T]) Share() Observable[T]

Share returns a new Observable that multicasts (shares) the original Observable. As long as there is at least one Subscriber this Observable will be subscribed and emitting data. When all subscribers have unsubscribed it will unsubscribe from the source Observable. Because the Observable is multicasting it makes the stream hot.

This method is useful when you have an Observable that is expensive to create or has side-effects, but you want to share the results of that Observable with multiple subscribers. By using `Share`, you can avoid creating multiple instances of the Observable and ensure that all subscribers receive the same data.

func (Observable[T]) Skip

func (observable Observable[T]) Skip(n int) Observable[T]

func (Observable[T]) Slice added in v0.2.0

func (observable Observable[T]) Slice(ctx ...context.Context) (slice []T, err error)

func (Observable[T]) StartWith

func (observable Observable[T]) StartWith(values ...T) Observable[T]

func (Observable[T]) Subscribe

func (observable Observable[T]) Subscribe(ctx context.Context, observe Observer[T]) Subscription

Subscribe subscribes to the observable with the given context, observer, and scheduler. The context will be available throughout the observable chain via subscriber.Context(). Source operators (like Create, Recv) are responsible for checking context cancellation.

func (Observable[T]) SubscribeOn

func (observable Observable[T]) SubscribeOn(scheduler ConcurrentScheduler) Observable[T]

func (Observable[T]) Take

func (observable Observable[T]) Take(n int) Observable[T]

Take returns an Observable that emits only the first count values emitted by the source Observable. If the source emits fewer than count values then all of its values are emitted. After that, it completes, regardless if the source completes.

func (Observable[T]) TakeWhile

func (observable Observable[T]) TakeWhile(condition func(T) bool) Observable[T]

func (Observable[T]) Tap added in v0.2.0

func (observable Observable[T]) Tap(tap Observer[T]) Observable[T]

func (Observable[T]) Values added in v0.2.0

func (observable Observable[T]) Values(ctx ...context.Context) iter.Seq2[T, error]

func (Observable[T]) Wait

func (observable Observable[T]) Wait(ctx ...context.Context) error

type Observer

type Observer[T any] func(next T, err error, done bool)

Observer[T] represents a consumer of values delivered by an Observable. It is implemented as a function that takes three parameters: - next: the next value emitted by the Observable - err: any error that occurred during emission (nil if no error) - done: a boolean indicating whether the Observable has completed

Observers follow the reactive pattern by receiving a stream of events (values, errors, or completion signals) and reacting to them accordingly.

func AsObserver added in v0.2.0

func AsObserver[T any](observe Observer[any]) Observer[T]

AsObserver converts an Observer of type `any` to an Observer of a specific type T. This allows adapting a generic Observer to a more specific type context.

func Ignore added in v0.2.0

func Ignore[T any]() Observer[T]

Ignore creates an Observer that simply discards any emissions from an Observable. It is useful when you need to create an Observer but don't care about its values.

func (Observer[T]) AsObserver

func (observe Observer[T]) AsObserver() Observer[any]

AsObserver converts a typed Observer[T] to a generic Observer[any]. It handles type conversion from 'any' back to T, and will emit an ErrTypecastFailed error when conversion fails.

func (Observer[T]) Done added in v0.2.2

func (observe Observer[T]) Done(err error)

Done signals that the Observable has completed emitting values, optionally with an error. If err is nil, it indicates normal completion. If err is non-nil, it indicates that the Observable terminated with an error.

After Done is called, the Observable will not emit any more values, regardless of whether the completion was successful or due to an error.

func (Observer[T]) Next

func (observe Observer[T]) Next(next T)

Next sends a new value to the Observer. This is a convenience method that handles the common case of emitting a new value without errors or completion signals.

type Pipe added in v0.2.0

type Pipe[T any] func(Observable[T]) Observable[T]

func Append added in v0.2.1

func Append[T any](slice *[]T) Pipe[T]

Append creates a pipe that appends each emitted value to the provided slice. It passes each value through to the next observer after appending it. This allows collecting all emitted values in a slice while still forwarding them. Only values emitted before completion (done=false) are appended.

func Assign added in v0.2.0

func Assign[T any](value *T) Pipe[T]

Assign creates a pipe that assigns every next value that is not done to the provided variable. The pipe will forward all events (next, err, done) to the next observer.

func AutoUnsubscribe added in v0.2.2

func AutoUnsubscribe[T any]() Pipe[T]

func Catch added in v0.2.0

func Catch[T any](other Observable[T]) Pipe[T]

func CatchError added in v0.2.0

func CatchError[T any](selector func(err error, caught Observable[T]) Observable[T]) Pipe[T]

func ConcatWith added in v0.2.1

func ConcatWith[T any](others ...Observable[T]) Pipe[T]

func Delay added in v0.2.1

func Delay[T any](duration time.Duration) Pipe[T]

func DistinctUntilChanged added in v0.2.0

func DistinctUntilChanged[T any](equal func(T, T) bool) Pipe[T]

func Do added in v0.2.0

func Do[T any](do func(T)) Pipe[T]

func ElementAt added in v0.2.1

func ElementAt[T any](n int) Pipe[T]

func EndWith added in v0.2.2

func EndWith[T any](values ...T) Pipe[T]

func Filter added in v0.2.0

func Filter[T any](predicate func(T) bool) Pipe[T]

func Fprint added in v0.2.0

func Fprint[T any](out io.Writer) Pipe[T]

func Fprintf added in v0.2.0

func Fprintf[T any](out io.Writer, format string) Pipe[T]

func Fprintln added in v0.2.0

func Fprintln[T any](out io.Writer) Pipe[T]

func MergeWith added in v0.2.1

func MergeWith[T any](others ...Observable[T]) Pipe[T]

func OnComplete added in v0.2.1

func OnComplete[T any](onComplete func()) Pipe[T]

func OnDone added in v0.2.1

func OnDone[T any](onDone func(error)) Pipe[T]

func OnError added in v0.2.1

func OnError[T any](onError func(error)) Pipe[T]

func OnNext added in v0.2.1

func OnNext[T any](onNext func(T)) Pipe[T]

func Passthrough added in v0.2.0

func Passthrough[T any]() Pipe[T]

func Print added in v0.2.0

func Print[T any]() Pipe[T]

func Printf added in v0.2.0

func Printf[T any](format string) Pipe[T]

func Println

func Println[T any]() Pipe[T]

func RaceWith added in v0.2.1

func RaceWith[T any](others ...Observable[T]) Pipe[T]

func Repeat added in v0.2.2

func Repeat[T any](count ...int) Pipe[T]

Repeat creates an Observable that emits the entire source sequence multiple times.

Parameters:

  • count: Optional. The number of repetitions:
  • If omitted: The source Observable is repeated indefinitely
  • If 0: Returns an empty Observable
  • If negative: Returns an Observable that emits an error
  • If multiple count values: Returns an Observable that emits an error

The resulting Observable will subscribe to the source Observable repeatedly each time the source completes, up to the specified count.

func Retry added in v0.2.2

func Retry[T any](limit ...int) Pipe[T]

func Send added in v0.2.0

func Send[T any](ch chan<- T) Pipe[T]

func Skip added in v0.2.0

func Skip[T any](n int) Pipe[T]

func StartWith added in v0.2.2

func StartWith[T any](values ...T) Pipe[T]

func Take added in v0.2.0

func Take[T any](n int) Pipe[T]

Take returns an Observable that emits only the first count values emitted by the source Observable. If the source emits fewer than count values then all of its values are emitted. After that, it completes, regardless if the source completes.

func TakeWhile added in v0.2.0

func TakeWhile[T any](condition func(T) bool) Pipe[T]

func Tap added in v0.2.0

func Tap[T any](tap Observer[T]) Pipe[T]

type Scheduler

type Scheduler = scheduler.Scheduler

Scheduler is the interface for scheduling tasks, either serially or concurrently.

type Signed added in v0.2.2

type Signed interface {
	~int | ~int8 | ~int16 | ~int32 | ~int64
}

Signed is a constraint that permits any signed integer type. If future releases of Go add new predeclared signed integer types, this constraint will be modified to include them.

type Subscriber

type Subscriber interface {
	// Subscribed returns true if the subscriber is in a subscribed state.
	// Returns false once Unsubscribe has been called.
	Subscribed() bool

	// Unsubscribe changes the state to unsubscribed and executes all registered
	// callback functions. Does nothing if already unsubscribed.
	Unsubscribe()

	// Add creates and returns a new child Subscriber.
	// If the parent is already unsubscribed, the child will be created in an
	// unsubscribed state. Otherwise, the child will be unsubscribed when the parent
	// is unsubscribed.
	Add() Subscriber

	// OnUnsubscribe registers a callback function to be executed when Unsubscribe is called.
	// If the subscriber is already unsubscribed, the callback is executed immediately.
	// If callback is nil, this method does nothing.
	OnUnsubscribe(callback func())

	// Context returns the context associated with this subscriber.
	// If no context was provided, returns context.Background().
	Context() context.Context
}

Subscriber is a subscribable entity that allows construction of a Subscriber tree.

type Subscription

type Subscription interface {
	// Subscribed returns true until Unsubscribe is called.
	Subscribed() bool

	// Unsubscribe will change the state to unsubscribed.
	Unsubscribe()

	// Done returns a channel that is closed when the subscription state changes to unsubscribed.
	// This channel can be used with select statements to react to subscription termination events.
	// If the scheduler is not concurrent, it will spawn a goroutine to wait for the scheduler.
	Done() <-chan struct{}

	// Err returns the subscription's terminal state:
	// - nil if the observable completed successfully
	// - the observable's error if it terminated with an error
	// - SubscriptionCanceled if the subscription was manually unsubscribed
	// - SubscriptionActive if the subscription is still active
	Err() error

	// Wait blocks until the subscription state becomes unsubscribed.
	// If the subscription is already unsubscribed, it returns immediately.
	// If the scheduler is not concurrent, it will wait for the scheduler to complete.
	// Returns:
	// - nil if the observable completed successfully
	// - the observable's error if it terminated with an error
	// - SubscriptionCanceled if the subscription was manually unsubscribed
	Wait() error
}

Subscription is an interface that allows monitoring and controlling a subscription. It provides methods for tracking the subscription's lifecycle.

type Tuple2 added in v0.2.1

type Tuple2[T, U any] struct {
	First  T
	Second U
}

type Tuple3 added in v0.2.1

type Tuple3[T, U, V any] struct {
	First  T
	Second U
	Third  V
}

type Tuple4 added in v0.2.1

type Tuple4[T, U, V, W any] struct {
	First  T
	Second U
	Third  V
	Fourth W
}

type Tuple5 added in v0.2.1

type Tuple5[T, U, V, W, X any] struct {
	First  T
	Second U
	Third  V
	Fourth W
	Fifth  X
}

type Unsigned added in v0.2.2

type Unsigned interface {
	~uint | ~uint8 | ~uint16 | ~uint32 | ~uint64 | ~uintptr
}

Unsigned is a constraint that permits any unsigned integer type. If future releases of Go add new predeclared unsigned integer types, this constraint will be modified to include them.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL