Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
56 changes: 56 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -612,6 +612,24 @@ These change each value into something else.
| `Switch()` | Another name for `SwitchTo`. | `Switch` |
| `Timestamp()`, `Timestamp(sequencer)` | Attaches the clock's current time to each value as a `Moment<T>`. | `Timestamp` |
| `TimeInterval()`, `TimeInterval(sequencer)` | Attaches the gap since the value before it. | `TimeInterval` |
| `GroupBy(keySelector)` | Splits the stream into one `GroupedSignal<TKey, T>` per key. A group arrives when its first value does. | `GroupBy` |
| `GroupBy(keySelector, comparer)`, `GroupBy(keySelector, capacity)`, `GroupBy(keySelector, capacity, comparer)` | Same, with a key comparer, an initial table capacity, or both. | `GroupBy` |
| `GroupBy(keySelector, elementSelector)` and the same comparer and capacity forms | Same, and runs `elementSelector` on each value before it enters its group. | `GroupBy` |
| `GroupByUntil(keySelector, durationSelector)` | Splits the stream into groups that end when their duration signal emits or completes. The next value with that key opens a new group. | `GroupByUntil` |
| `GroupByUntil(keySelector, durationSelector, comparer)`, `GroupByUntil(keySelector, durationSelector, capacity)`, `GroupByUntil(keySelector, durationSelector, capacity, comparer)` | Same, with a key comparer, an initial table capacity, or both. | `GroupByUntil` |
| `GroupByUntil(keySelector, elementSelector, durationSelector)` and the same comparer and capacity forms | Same, and runs `elementSelector` on each value before it enters its group. | `GroupByUntil` |
| `Slice(count)` | Hands out one signal per window of `count` values. The windows do not overlap. | `Window` |
| `Slice(count, skip)` | Opens a window every `skip` values, so windows can overlap or leave gaps. | `Window` |
| `Slice(timeSpan)`, `Slice(timeSpan, sequencer)` | Hands out one signal per time window. | `Window` |
| `Slice(timeSpan, timeShift)`, `Slice(timeSpan, timeShift, sequencer)` | Opens a window of `timeSpan` every `timeShift`. | `Window` |
| `Slice(timeSpan, count)`, `Slice(timeSpan, count, sequencer)` | Ends each window after `timeSpan` or `count` values, whichever comes first. | `Window` |
| `Slice(boundaries)` | Starts the next window each time the boundary signal emits. | `Window` |
| `Slice(closingSelector)` | Ends each window with the signal your selector returns for it. | `Window` |
| `Slice(openings, closingSelector)` | Opens a window for each opening value and ends it with the signal your selector returns. | `Window` |
| `Window(...)` | Another name for every `Slice` overload. | `Window` |

`Buffer` hands you a list per window. `Slice` hands you a signal per window, so you can subscribe to it and work on the
values as they arrive.

`Map` changes every value.

Expand Down Expand Up @@ -658,6 +676,29 @@ pages.SwitchTo().Subscribe(value => Console.Write(value + " "));
// prints: 1 2 10 11
```

`GroupBy` gives each key its own signal. The `Key` property tells you which group you hold.

```csharp
Signal.Sequence(1, 6)
.GroupBy(number => number % 2)
.Subscribe(group => group.Subscribe(number => Console.Write($"{group.Key}:{number} ")));
// prints: 1:1 0:2 1:3 0:4 1:5 0:6
```

Subscribe to a group as soon as it arrives. A group signal sends only the values that come after you subscribe.
The source stays subscribed until you dispose the outer subscription and every group subscription.

`Slice` gives each window its own signal.

```csharp
Signal.Sequence(1, 5)
.Slice(2)
.Subscribe(window => window.Subscribe(
number => Console.Write(number + " "),
() => Console.Write("| ")));
// prints: 1 2 | 3 4 | 5 |
```

### Filtering

These decide which values get through.
Expand Down Expand Up @@ -2718,6 +2759,9 @@ dotnet add xyz.Reactive/xyz.Reactive.csproj package ReactiveUI.Primitives.Maui.R
| `DelaySubscription` | `DelayStart` | Delay source subscription. |
| `Timeout` | `Expire` | Error on missing value before due time. |
| `Buffer(count)` | `Buffer(count)` | Fixed-size buffers. |
| `GroupBy` | `GroupBy` | Eight overloads, each returning `GroupedSignal<TKey, T>` groups. |
| `GroupByUntil` | `GroupByUntil` | Eight overloads; a group ends when its duration signal emits or completes. |
| `Window` | `Slice` or Rx-name `Window` | Eleven overloads; each window is a signal. |
| `SubscribeOn` | `SubscribeOn` | Schedule source subscription. |
| `ToList` / `ToArray` | `ToList` / `ToArray` or `CollectList` / `CollectArray` | Signal results. |
| `FirstAsync` / `LastAsync` | `FirstAsync` / `LastAsync` | Task result. |
Expand Down Expand Up @@ -2963,6 +3007,14 @@ the same thing.
| `RepeatSourceSignal` | `Repeat` over a source |
| `BufferSignal` | `Buffer` |
| `CollectSignal` | `Collect` |
| `GroupBySignal` | `GroupBy` |
| `GroupByUntilSignal` | `GroupByUntil` |
| `SliceCountSignal` | `Slice(count)`, `Slice(count, skip)` |
| `SliceTimeSignal` | `Slice(timeSpan)`, `Slice(timeSpan, timeShift)` |
| `SliceTimeCountSignal` | `Slice(timeSpan, count)` |
| `SliceBoundarySignal` | `Slice(boundaries)` |
| `SliceClosingSignal` | `Slice(closingSelector)` |
| `SliceOpeningSignal` | `Slice(openings, closingSelector)` |
| `EmitIfQuietSignal` | `EmitIfQuiet` |
| `SerializeSignal` | `Serialize` |
| `SynchronizeSignal` | `Synchronize` |
Expand Down Expand Up @@ -3032,6 +3084,10 @@ These sit beside the operators. You call them directly rather than through a cha
| `PublishingOption` | `ReactiveUI.Primitives.Async.Signals` | Chooses how an async signal publishes to its subscribers. |
| `SparkKind` | `ReactiveUI.Primitives.Core` | Which notification a `Spark` carries. |
| `Broadcaster<T>` | `ReactiveUI.Primitives.Signals` | The struct behind a signal's subscriber list. `Add`, `Remove`, `Next`, `Error`, `Completed` and `HasObservers`. Hold it as a field. |
| `GroupedSignal<TKey, T>` | `ReactiveUI.Primitives.Signals` | The signal of one group from `GroupBy` or `GroupByUntil`. It has a `Key` and sends the group's values. |
| `SliceWindow<T>` | `ReactiveUI.Primitives.Advanced` | One window from `Slice`. Publish to it and every current subscriber receives the value. A subscriber that arrives after it ends is completed at once. |
| `SliceRouter<TOuter, T>` | `ReactiveUI.Primitives.Advanced` | Posts window and group notifications while you hold your lock, then delivers them in order after you release it. |
| `SharedSubscription` | `ReactiveUI.Primitives.Advanced` | Keeps a source subscribed until its own handle and every lease taken from it are disposed. |
| `CopyOnWriteList<T>` | `ReactiveUI.Primitives.Advanced` | An immutable list that returns a new instance from `Add` and `Remove`. `Empty` starts one. |
| `SinkTerminal`, `SinkSubscription`, `SinkDelivery` | `ReactiveUI.Primitives.Advanced` | The static helpers a synchronous sink calls: forward a terminal notification, assign or dispose the upstream subscription, and forward a value while disposing the sink if the observer throws. |
| `AsyncContext`, `AsyncContextExtensions` | `ReactiveUI.Primitives.Async` | The ambient context an async signal flows through its operators. |
Expand Down
10 changes: 5 additions & 5 deletions src/Primitives.Extensions.Shared/ReactiveExtensions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -207,9 +207,9 @@ public IObservable<T> Conflate(TimeSpan minimumUpdatePeriod, ISequencer schedule
public IObservable<Heartbeat<T>> Heartbeat(TimeSpan heartbeatPeriod, ISequencer scheduler) =>
new HeartbeatObservable<T>(source, heartbeatPeriod, scheduler);

/// <summary>Emit the latest value or a default if none exists.</summary>
/// <param name="defaultValue">The default value.</param>
/// <returns>A sequence that emits the latest value or the default.</returns>
/// <summary>Emits <paramref name="defaultValue"/> on subscribe, then each source value that differs from the last value emitted.</summary>
/// <param name="defaultValue">The value emitted on subscribe before the source is subscribed; source values equal to the last emitted value, including this one, are dropped.</param>
/// <returns>A sequence that starts with the default value, then forwards changed values, errors and completion. Equality uses the default comparer.</returns>
[MethodImpl(MethodImplOptions.AggressiveInlining)]
public IObservable<T> LatestOrDefault(T defaultValue) =>
new LatestOrDefaultObservable<T>(source, defaultValue);
Expand Down Expand Up @@ -768,11 +768,11 @@ public IObservable<T> GetMin(params IObservable<T>[] sources)
}
}

/// <summary>Null-skipping operators for an observable source sequence of reference types.</summary>
/// <summary>Null-skipping operators for an observable source sequence of reference types, nullable or not.</summary>
/// <typeparam name="T">The reference type of the source elements.</typeparam>
/// <param name="source">The source observable.</param>
extension<T>(IObservable<T> source)
where T : class
where T : class?
{
/// <summary>Skip null values until the first non-null appears.</summary>
/// <returns>An IObservable of T.</returns>
Expand Down
55 changes: 55 additions & 0 deletions src/Primitives.Shared/Advanced/SliceTimeCountSignal{T}.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
// Copyright (c) 2019-2026 ReactiveUI Association Incorporated. All rights reserved.
// ReactiveUI Association Incorporated licenses this file to you under the MIT license.
// See the LICENSE file in the project root for full license information.

#if REACTIVE_SHIM
namespace ReactiveUI.Primitives.Reactive.Advanced;
#else
namespace ReactiveUI.Primitives.Advanced;
#endif

/// <summary>Cold signal that splits a source into windows that end after a duration or a number of values, whichever comes first.</summary>
/// <typeparam name="T">The value type.</typeparam>
[System.Diagnostics.DebuggerDisplay("SliceTimeCountSignal: TimeSpan = {_timeSpan}, Count = {_count}, Source = {_source}")]
public sealed class SliceTimeCountSignal<T> : IObservable<IObservable<T>>
{
/// <summary>The source observable.</summary>
private readonly IObservable<T> _source;

/// <summary>The maximum duration of each window.</summary>
private readonly TimeSpan _timeSpan;

/// <summary>The maximum number of values in each window.</summary>
private readonly int _count;

/// <summary>The sequencer that schedules the timer.</summary>
private readonly ISequencer _sequencer;

/// <summary>Initializes a new instance of the <see cref="SliceTimeCountSignal{T}"/> class.</summary>
/// <param name="source">The source observable.</param>
/// <param name="timeSpan">The maximum duration of each window.</param>
/// <param name="count">The maximum number of values in each window.</param>
/// <param name="sequencer">The sequencer that schedules the timer.</param>
/// <exception cref="ArgumentNullException"><paramref name="source"/> or <paramref name="sequencer"/> is <see langword="null"/>.</exception>
/// <exception cref="ArgumentOutOfRangeException"><paramref name="timeSpan"/> or <paramref name="count"/> is zero or negative.</exception>
public SliceTimeCountSignal(IObservable<T> source, TimeSpan timeSpan, int count, ISequencer sequencer)
{
SliceTimeGuard.ThrowIfNotPositive(timeSpan);
ArgumentOutOfRangeExceptionHelper.ThrowIfNegativeOrZero(count);
_source = source ?? throw new ArgumentNullException(nameof(source));
_timeSpan = timeSpan;
_count = count;
_sequencer = sequencer ?? throw new ArgumentNullException(nameof(sequencer));
}

/// <inheritdoc/>
public IDisposable Subscribe(IObserver<IObservable<T>> observer)
{
ArgumentExceptionHelper.ThrowIfNull(observer);

SliceTimeCountWitness<T> sink = new(observer, _timeSpan, _count, _sequencer);
sink.Start();
sink.SetSubscription(_source.Subscribe(sink));
return sink.Subscription;
}
}
Loading
Loading