Core 16 Operators topic
The "Core 16" Essential Cell Operators
A Practical Guide for First-Time Developers
Table of Contents
- Introduction
- Learning Path
- Phase 1: Get Data In
- Phase 2: Hold State
- Phase 3: React in UI
- Phase 4: Shape Streams
- Phase 5: Go Async
- Phase 6: Combine Sources
- Quick Reference Card
- Common Patterns
- Next Steps
Introduction
The Cell Framework provides 16 essential operators that cover most of your daily reactive programming needs. They are carefully ordered to guide you through a natural learning path:
get data in → hold state → react in the UI → shape streams → go async → combine sources
Each operator works with zero knowledge of the underlying framework mechanics (Receptor, TestCell, Context, Synapses) — all optional, all defaulted. They represent the Standard Entry Point for building reactive systems.
Learning Path
Read left to right, then down. Numbers match the operator index in this guide.
get data in → hold state → react in UI → shape streams → go async → combine → isolate
┌──────────────┐ ┌──────────────┐ ┌──────────────┐
│ GET DATA IN │ │ HOLD STATE │ │ REACT IN UI │
│ │ │ │ │ │
│ 2 ingress │ │ 1 state │ │ 3 observe │
│ │ │ 4 derive │ │ │
└──────┬───────┘ └──────┬───────┘ └──────┬───────┘
│ │ │
└────────────┬─────┴────────┬─────────┘
▼ ▼
┌────────────────┐ ┌────────────────┐
│ SHAPE STREAMS │ │ GO ASYNC │
│ │ │ │
│ 5 debounce │ │ 8 asyncMap │
│ 6 distinct │ │ 9 fromFuture │
│ 7 throttle │ │ 10 fromStream │
└───────┬────────┘ └───────┬────────┘
│ │
└─────────┬──────────┘
▼
┌─────────────────────┐
│ COMBINE SOURCES │
│ │
│ 11 synthesis │
│ 12 hub │
│ 13 switchMap │
│ 14 sanitized │
└──────────┬──────────┘
▼
┌─────────────────────┐
│ ISOLATE WRITES │
│ │
│ 15 transaction │
│ 16 txApply │
└─────────────────────┘
Cell.valve and Cell.open are used in /example but are not part of this 16.
Phase 1: Get Data In
1. Cell.ingress — How intent/events enter the graph
Purpose: Create a manual entry point for external events (UI clicks, network messages, hardware sensors).
When to use:
- User interactions (button clicks, form submissions)
- External events (WebSocket messages, sensor readings)
- Testing and simulation
How it works:
- Creates a cell that accepts raw data
- Wraps data in a Pulse automatically
- Broadcasts to all observers
Basic Usage:
// Create an ingress for string events
final events = Cell.ingress<String>();
// Emit events
events.emit('click');
events.emit('submit');
// Observe events
Cell.observe(
source: events.cell,
effect: (pulse) => print('Event: ${pulse.payload}'),
);
With Refinement:
final events = Cell.ingress<String>(
refine: (host, input) {
// Transform or filter incoming events
if (input.payload.isEmpty) return null;
return Pulse(input.payload.toUpperCase());
},
);
events.emit('hello'); // Observers receive "HELLO"
events.emit(''); // Filtered out (null)
Advanced Usage:
final events = Cell.ingress<Map<String, dynamic>>(
refine: (host, input) {
final data = input.payload;
if (!data.containsKey('type')) return null;
return Pulse(data, type: data['type'] as String?);
},
context: Context.module('event-processor'),
testRule: TestCell.allowAll,
);
// Emit with type
events.emit({'type': 'user.login', 'user': 'alice'});
// Observers receive a Pulse with type 'user.login'
2. Cell.state — Retained app state you read and update
Purpose: Create a persistent, mutable state atom that maintains its value across reactive cycles.
When to use:
- Application state (counter, settings, user profile)
- UI state (form inputs, toggle states)
- Domain models (entities from database)
How it works:
- Initialized with a starting value
- Updates are processed through
evolvefunction - Changes are validated and committed atomically
- Observers are notified
Basic Usage:
final counter = Cell.state<int>(
initial: 0,
evolve: (host, input) {
final delta = input.payload as int? ?? 1;
return Pulse(host.value + delta);
},
);
// Update
counter.update(1); // value becomes 1
counter.update(5); // value becomes 6
// Read
print(counter.cell.value); // 6
With Validation:
final isPositive = TestCell<int>((value, {host, ...}) => value > 0);
final counter = Cell.state<int>(
initial: 0,
evolve: (host, input) {
final delta = input.payload as int? ?? 1;
return Pulse(host.value + delta);
},
testRule: isPositive,
);
counter.update(5); // ✓ Allowed
counter.update(-1); // ✗ Blocked (value stays 5)
Complex State:
class User {
final String name;
final int age;
User(this.name, this.age);
}
final user = Cell.state<User>(
initial: User('Alice', 25),
evolve: (host, input) {
final updates = input.payload as Map<String, dynamic>?;
if (updates == null) return null;
final current = host.value;
final newUser = User(
updates['name'] as String? ?? current.name,
updates['age'] as int? ?? current.age,
);
return Pulse(newUser);
},
);
user.update({'name': 'Bob'}); // Name changes to Bob
user.update({'age': 30}); // Age changes to 30
user.update({'name': 'Charlie', 'age': 35}); // Both change
Phase 2: Hold State
4. Cell.derive — Pure view-models / projections from state
Purpose: Create a derived value that transforms one source into another.
When to use:
- View models from state
- Data transformations
- Formatting and sanitization
How it works:
- Listens to a source cell
- Applies a pure transformation
- Emits the transformed result
- Preserves causal provenance
Basic Usage:
final counter = Cell.state<int>(
initial: 0,
evolve: (host, input) {
final delta = input.payload as int? ?? 1;
return Pulse(host.value + delta);
},
);
// Derive a display string
final display = Cell.derive<int, String>(
source: counter.cell,
project: (pulse) => Pulse('Count: ${pulse.payload}'),
);
print(display.value); // "Count: 0"
With Type Conversion:
// Convert User object to display name
final user = Cell.state<User>(
initial: User('Alice', 25),
evolve: (host, input) => Pulse(input.payload as User),
);
final displayName = Cell.derive<User, String>(
source: user.cell,
project: (pulse) => Pulse('${pulse.payload.name} (${pulse.payload.age})'),
);
Filtering with Derive:
// Only emit when the value is positive
final positiveOnly = Cell.derive<int, int>(
source: source.cell,
project: (pulse) {
final value = pulse.payload;
return value > 0 ? Pulse(value) : null; // Filter out non-positive
},
);
5. Cell.debounce — Everyday input: search, validation, autosave
Purpose: Wait for a period of stability before emitting the latest value.
When to use:
- Search-as-you-type
- Form validation
- Autosave
- Window resize
How it works:
- Each new pulse resets a timer
- Only emits when silence period passes
- Leading edge option for immediate first response
Basic Usage:
final searchInput = Cell.ingress<String>();
final debouncedSearch = Cell.debounce(
searchInput.cell,
Duration(milliseconds: 300),
);
// Only emits after user stops typing for 300ms
Cell.observe(
source: debouncedSearch,
effect: (pulse) => performSearch(pulse.payload),
);
With Leading Edge:
final debouncedSearch = Cell.debounce(
searchInput.cell,
Duration(milliseconds: 300),
leading: true, // First pulse emits immediately
);
// Emits: first character immediately, then after silence
Real-World Example:
class SearchWidget extends StatefulWidget {
@override
_SearchWidgetState createState() => _SearchWidgetState();
}
class _SearchWidgetState extends State<SearchWidget> {
final searchInput = Cell.ingress<String>();
late Cell debounced;
late EgressHandle observer;
@override
void initState() {
super.initState();
debounced = Cell.debounce(
searchInput.cell,
Duration(milliseconds: 300),
);
observer = Cell.observe(
source: debounced,
effect: (pulse) => performSearch(pulse.payload),
);
}
void onSearchChanged(String query) {
searchInput.emit(query);
}
}
8. Cell.throttle — Frequency-based rate limiting
Purpose: Limit the frequency of updates to a predictable, constant rate.
When to use:
- UI performance (scroll, mouse-move)
- API rate limiting
- Sensor data sampling
How it works:
- Leading edge: first pulse emits immediately
- Silent window: subsequent pulses are suppressed
- Trailing edge: last pulse emits when window ends
Basic Usage:
final throttled = Cell.throttle(
source.cell,
Duration(milliseconds: 100), // At most one per 100ms
leading: true, // First pulse immediate
trailing: false, // No trailing emission
);
// Use case: scroll events
final scroll = Cell.ingress<int>();
final throttledScroll = Cell.throttle(
scroll.cell,
Duration(milliseconds: 16), // ~60fps
);
With Trailing:
final throttled = Cell.throttle(
source.cell,
Duration(milliseconds: 200),
leading: true, // First immediate
trailing: true, // Last at end of window
);
// Pattern: First → (silence) → Last
Throttle vs Debounce:
// Throttle: "At most one per N ms"
final throttled = Cell.throttle(source, Duration(milliseconds: 300));
// Debounce: "Wait for N ms of silence"
final debounced = Cell.debounce(source, Duration(milliseconds: 300));
// Throttle is better for: scroll events, mouse movement
// Debounce is better for: search input, form validation
9. Cell.distinct — Skip redundant updates and rebuilds
Purpose: Suppress consecutive duplicate payloads.
When to use:
- Prevent unnecessary UI rebuilds
- Reduce network requests
- Filter sensor noise
How it works:
- Compares current payload with previous
- Skips if equal (using
==or custom comparator) - First emission always passes
Basic Usage:
final unique = Cell.distinct(
source: source.cell,
equals: (a, b) => a == b, // Optional custom equality
);
// Sequence: 1, 1, 2, 2, 3, 1
// Emits: 1, 2, 3, 1
Custom Equality:
final unique = Cell.distinct(
source: source.cell,
equals: (a, b) {
// Case-insensitive comparison
if (a is String && b is String) {
return a.toLowerCase() == b.toLowerCase();
}
return a == b;
},
);
Phase 5: Go Async
9. Cell.asyncMap — HTTP/DB work (latestOnly / exhaust)
Purpose: Run background tasks for each input, with configurable concurrency control.
When to use:
- API requests
- Database queries
- Heavy computations
- Data enrichment
How it works:
- Maps each input to an async task
- Controls concurrency (parallel, sequential, throttled)
- Emits results as they complete
Concurrency Modes:
// Parallel (default) - unlimited concurrency
final parallel = Cell.asyncMap<int, User>(
source,
(id) => api.fetchUser(id),
concurrency: 0, // Unlimited
);
// Sequential - one at a time
final sequential = Cell.asyncMap<int, User>(
source,
(id) => api.fetchUser(id),
concurrency: 1, // One at a time
);
// Limited - at most 3 concurrent
final limited = Cell.asyncMap<int, User>(
source,
(id) => api.fetchUser(id),
concurrency: 3,
);
SwitchMap (latest only):
// Only care about the most recent result
final latest = Cell.asyncMap<int, User>(
source,
(id) => api.fetchUser(id),
latestOnly: true,
);
// If new request starts before old completes, old result is dropped
ExhaustMap (ignore while busy):
// Ignore new inputs while processing
final exhaust = Cell.asyncMap<int, User>(
source,
(id) => api.fetchUser(id),
exhaust: true,
);
// If busy, new inputs are ignored
Real-World Example:
final userId = Cell.ingress<int>();
// Only care about the most recently selected user
final userProfile = Cell.asyncMap<int, UserProfile>(
userId.cell,
(id) => api.fetchProfile(id),
latestOnly: true,
);
Cell.observe(
source: userProfile,
effect: (pulse) => displayUser(pulse.payload),
);
userId.emit(1); // Fetches user 1
userId.emit(2); // Cancels user 1, fetches user 2
10. Cell.fromFuture — Bridge Futures and Streams into Cell
Purpose: Bridge a single asynchronous result into the reactive graph.
When to use:
- Loading initial configuration
- One-time data fetch
- Initialization tasks
Basic Usage:
final config = Cell.fromFuture(loadConfig());
Cell.observe(
source: config,
effect: (pulse) => applyConfig(pulse.payload),
);
With Error Handling:
Future<String> loadData() async {
try {
final response = await http.get(url);
return response.body;
} catch (e) {
return 'Error: $e';
}
}
final data = Cell.fromFuture(loadData());
11. Cell.fromStream — Bridge continuous streams
Purpose: Bridge a continuous Stream into the reactive graph.
When to use:
- WebSockets
- File watchers
- Timers and periodic events
- Hardware sensors
Basic Usage:
final ticks = Cell.fromStream(
Stream.periodic(Duration(seconds: 1), (i) => i),
);
Cell.observe(
source: ticks,
effect: (pulse) => print('Tick: ${pulse.payload}'),
);
With WebSocket:
final ws = WebSocket.connect('wss://example.com');
final messages = Cell.fromStream(ws.asBroadcastStream());
Cell.observe(
source: messages,
effect: (pulse) => handleMessage(pulse.payload),
);
Phase 6: Combine Sources
12. Cell.synthesis — Forms & dashboards: latest of several fields
Purpose: Merge multiple source cells into one, reacting when any source changes.
When to use:
- Form validation (combining multiple fields)
- Dashboards (multiple data sources)
- Calculations (price + tax + shipping)
How it works:
- Links to all source cells
- Reacts when any source changes
- Aggregator receives all sources and the triggering pulse
- Emits a new aggregated result
Basic Usage:
final total = Cell.synthesis<double>(
[price.cell, tax.cell, shipping.cell],
aggregator: (sources, pulse) {
final p = sources.elementAt(0).value as double? ?? 0;
final t = sources.elementAt(1).value as double? ?? 0;
final s = sources.elementAt(2).value as double? ?? 0;
return Pulse(p + t + s);
},
);
Form Validation Example:
final email = Cell.state<String>(
initial: '',
evolve: (host, input) => Pulse(input.payload as String? ?? ''),
);
final password = Cell.state<String>(
initial: '',
evolve: (host, input) => Pulse(input.payload as String? ?? ''),
);
final formValid = Cell.synthesis<bool>(
[email.cell, password.cell],
aggregator: (sources, pulse) {
final email = sources.elementAt(0).value as String? ?? '';
final password = sources.elementAt(1).value as String? ?? '';
final valid = email.contains('@') && password.length >= 8;
return Pulse(valid);
},
);
Cell.observe(
source: formValid,
effect: (pulse) {
submitButton.enabled = pulse.payload;
},
);
15. Cell.switchMap — Dynamic source switching
Purpose: Switch to a new reactive source dynamically based on a selection.
When to use:
- Tab switching
- User selection
- Language/locale changes
- Module activation
How it works:
- Listens to a source cell (the selector)
- Each new value triggers a mapper that returns a new cell
- Automatically unlinks old source, links new one
- Forwards pulses from the active source
Basic Usage:
final profile = Cell.switchMap<int, Profile>(
selectedId.cell,
(id) => getProfileCellFor(id),
);
// When selectedId changes, profile switches to a new source
Real-World Example:
// Tab navigation
final selectedTab = Cell.state<String>(
initial: 'home',
evolve: (host, input) => Pulse(input.payload as String? ?? 'home'),
);
final tabContent = Cell.switchMap<String, Widget>(
selectedTab.cell,
(tab) {
switch (tab) {
case 'home': return homeTab.cell;
case 'settings': return settingsTab.cell;
case 'profile': return profileTab.cell;
default: return emptyTab.cell;
}
},
);
16. Cell.transaction — Multi-cell atomic updates
Purpose: Group multiple cell updates into a single atomic unit.
When to use:
- Financial transfers
- Inventory adjustments
- Complex form submissions
- Any operation with invariants across cells
How it works:
- Begin: Register participants
- Update: Buffer writes (not applied yet)
- Read: Observe values according to isolation level
- Commit: Acquire locks, validate, apply atomically
- Rollback: Discard buffered changes
Basic Usage:
final tx = Cell.transaction();
await tx.begin([accountA, accountB]);
final fromBalance = tx.read(accountA) as int;
final toBalance = tx.read(accountB) as int;
tx.update(accountA, fromBalance - 50);
tx.update(accountB, toBalance + 50);
await tx.commit(); // All or nothing
With Savepoints:
final tx = Cell.transaction();
await tx.begin([cell1, cell2, cell3]);
tx.update(cell1, 10);
tx.update(cell2, 20);
// Create a checkpoint
final sp = tx.savepoint();
// Speculative updates
tx.update(cell2, 30);
tx.update(cell3, 40);
// Something went wrong - rollback to checkpoint
await tx.rollback(savepoint: sp);
// cell1 = 10, cell2 = 20, cell3 unchanged
await tx.commit();
With Isolation Levels:
final tx = Cell.transaction(TransactionOptions(
isolation: IsolationLevel.repeatableRead,
timeout: Duration(seconds: 5),
onEvent: (e) => print(e),
));
await tx.begin([accountA, accountB]);
// Reads return snapshot from begin
final a = tx.read(accountA) as int;
tx.update(accountA, a + 100);
await tx.commit();
17. Cell.txApply — Batch multiple apply() into a single commit
Purpose: Batch multiple apply() calls into a single atomic commit.
When to use:
- Multiple function calls that must be atomic
- Complex state mutations
- Performance optimization (reduce downstream churn)
Basic Usage:
final tx = Cell.txApply();
await tx.execute(
participants: [cell1, cell2],
body: (tx) {
cell1.apply(updateFunction1);
cell2.apply(updateFunction2);
},
);
With Compensation:
final tx = Cell.txApply(TxApplyOptions(
compensationErrorPolicy: CompensationErrorPolicy.bestEffort,
compensationMaxAttempts: 3,
));
await tx.execute(
participants: [cell1, cell2],
body: (tx) {
cell1.apply(
updateFunction,
compensate: rollbackFunction,
);
},
);
Quick Reference Card
Essential Operators Summary
| # | Operator | Category | Purpose |
|---|---|---|---|
| 1 | state |
Entry Point | Retained app state |
| 2 | ingress |
Entry Point | Event entry |
| 3 | observe |
Observation | Side effects |
| 4 | derive |
Transformation | Projections |
| 5 | debounce |
Flow Control | Stability-based |
| 6 | distinct |
Flow Control | Deduplication |
| 7 | throttle |
Flow Control | Rate limiting |
| 8 | asyncMap |
Transformation | Background tasks |
| 9 | fromFuture |
Bridge | One-time async |
| 10 | fromStream |
Bridge | Continuous async |
| 11 | synthesis |
Transformation | Multi-source aggregation |
| 12 | hub |
Routing | Signal routing |
| 13 | switchMap |
Transformation | Dynamic source selection |
| 14 | sanitized |
Transformation | Data redaction |
| 15 | transaction |
Orchestration | Atomic updates |
| 16 | txApply |
Orchestration | Batch apply() |
Common Patterns
1. Counter with UI
final counter = Cell.state<int>(
initial: 0,
evolve: (host, input) {
final delta = input.payload as int? ?? 1;
return Pulse(host.value + delta);
},
);
final display = Cell.derive<int, String>(
source: counter.cell,
project: (pulse) => Pulse('Count: ${pulse.payload}'),
);
2. Search with Debounce
final search = Cell.ingress<String>();
final debounced = Cell.debounce(search.cell, Duration(milliseconds: 300));
final results = Cell.asyncMap<String, List<Result>>(
debounced,
(query) => api.search(query),
latestOnly: true,
);
3. Form Validation
final email = Cell.state<String>(...);
final password = Cell.state<String>(...);
final isValid = Cell.synthesis<bool>([email.cell, password.cell], ...);
4. Event Router
final hub = Cell.hub(
routing: HubRouting.pattern,
registrations: [
(key: 'user.*', priority: 10, handler: userHandler),
(key: 'admin.*', priority: 20, handler: adminHandler),
],
fallback: 'unknown',
);
5. Atomic Transfer
final tx = Cell.transaction();
await tx.begin([accountA, accountB]);
final a = tx.read(accountA) as int;
final b = tx.read(accountB) as int;
tx.update(accountA, a - 50);
tx.update(accountB, b + 50);
await tx.commit();
6. Async Data Loading
final userId = Cell.ingress<int>();
final userProfile = Cell.asyncMap<int, UserProfile>(
userId.cell,
(id) => api.fetchProfile(id),
latestOnly: true,
);
7. Multi-Source Synthesis
final summary = Cell.synthesis([a.cell, b.cell, c.cell], aggregator: ...);
Next Steps
Continue Learning
| Resource | Description |
|---|---|
| HowTo-Instruction.md | Deep dive into Instructions |
| HowTo-Receptor.md | Building transformation pipelines |
| HowTo-TestCell.md | Validation and security |
| HowTo-Start.md | Getting started guide |
Run the Examples
# Instruction pipeline
dart run examples/instruction_pipeline_walkthrough.dart
# Hub demo
dart run examples/hub_demo.dart
# Throttle demo
dart run examples/throttle_demo.dart
Build Your First App
import 'package:cell/cell.dart';
void main() {
// State
final counter = Cell.state<int>(...);
// Derive
final display = Cell.derive<int, String>(...);
// Observe
final observer = Cell.observe<int>(
source: counter.cell,
effect: (pulse) => print(pulse.payload),
);
// Use
counter.update(1);
}
Summary
Key Takeaways
- 16 operators cover most daily reactive needs
- Learning path guides you from simple to complex
- Each operator has a specific purpose
- Composition creates powerful systems
- Zero boilerplate - default settings work
Quick Start Checklist
- Understand the learning path
- Start with
state,ingress,observe - Add
derivefor projections - Use
debouncefor user input - Use
asyncMapfor network calls - Use
synthesisfor multi-source - Use
transactionfor atomic updates - Use
hubfor routing
Happy coding with the Cell Framework! 🚀
Classes
- Cell Getting Started Core 16 Operators Demo:aircraft Demo:hotel Core
- A reactive node: holds state (or relays signals), validates every incoming change against a policy, and broadcasts accepted changes to whatever else is listening.
- OpenCell Core 16 Operators Core
- A specialized, interactive Cell that serves as a Reactive Bridge for external stimulus injection and dynamic topology management.