The problem
Need reactive data structures that notify observers of changes
The solution
class Observable {
constructor(value) {
this._value = value;
this._observers = new Set();
this._onChangeCallbacks = [];
}
get value() {
return this._value;
}
set value(newValue) {
const oldValue = this._value;
if (oldValue !== newValue) {
this._value = newValue;
this._notifyObservers(newValue, oldValue);
this._executeOnChangeCallbacks(newValue, oldValue);
}
}
subscribe(observer) {
this._observers.add(observer);
// Return unsubscribe function
return () => {
this._observers.delete(observer);
};
}
subscribeNext(callback) {
const observer = {
next: callback
};
return this.subscribe(observer);
}
onChange(callback) {
this._onChangeCallbacks.push(callback);
return () => {
const index = this._onChangeCallbacks.indexOf(callback);
if (index > -1) {
this._onChangeCallbacks.splice(index, 1);
}
};
}
map(transform) {
const mapped = new Observable(transform(this._value));
this.subscribe((newValue, oldValue) => {
mapped.value = transform(newValue);
});
return mapped;
}
filter(predicate) {
const filtered = new Observable(
predicate(this._value) ? this._value : undefined
);
this.subscribe((newValue) => {
if (predicate(newValue)) {
filtered.value = newValue;
}
});
return filtered;
}
debounceTime(delay) {
const debounced = new Observable(this._value);
let timeoutId;
this.subscribe((newValue) => {
clearTimeout(timeoutId);
timeoutId = setTimeout(() => {
debounced.value = newValue;
}, delay);
});
return debounced;
}
distinctUntilChanged(compareFn = (a, b) => a === b) {
const distinct = new Observable(this._value);
this.subscribe((newValue, oldValue) => {
if (!compareFn(newValue, oldValue)) {
distinct.value = newValue;
}
});
return distinct;
}
_notifyObservers(newValue, oldValue) {
for (const observer of this._observers) {
try {
if (typeof observer === 'function') {
observer(newValue, oldValue);
} else if (observer.next) {
observer.next(newValue, oldValue);
}
} catch (error) {
console.error('Error in observer:', error);
}
}
}
_executeOnChangeCallbacks(newValue, oldValue) {
for (const callback of this._onChangeCallbacks) {
try {
callback(newValue, oldValue);
} catch (error) {
console.error('Error in onChange callback:', error);
}
}
}
}
class ObservableObject {
constructor(obj = {}) {
this._data = new Proxy(obj, this._createHandler());
this._observers = new Map(); // property -> Set of observers
this._onAnyChangeCallbacks = [];
}
get data() {
return this._data;
}
subscribe(property, observer) {
if (!this._observers.has(property)) {
this._observers.set(property, new Set());
}
this._observers.get(property).add(observer);
return () => {
const observers = this._observers.get(property);
if (observers) {
observers.delete(observer);
if (observers.size === 0) {
this._observers.delete(property);
}
}
};
}
subscribeToAll(observer) {
this._onAnyChangeCallbacks.push(observer);
return () => {
const index = this._onAnyChangeCallbacks.indexOf(observer);
if (index > -1) {
this._onAnyChangeCallbacks.splice(index, 1);
}
};
}
set(property, value) {
this._data[property] = value;
}
get(property) {
return this._data[property];
}
_createHandler() {
return {
set: (target, property, value) => {
const oldValue = target[property];
if (oldValue !== value) {
target[property] = value;
this._notifyObservers(property, value, oldValue);
return true;
}
return true;
},
deleteProperty: (target, property) => {
const oldValue = target[property];
const deleted = delete target[property];
if (deleted) {
this._notifyObservers(property, undefined, oldValue);
}
return deleted;
}
};
}
_notifyObservers(property, newValue, oldValue) {
// Notify specific property observers
const propertyObservers = this._observers.get(property);
if (propertyObservers) {
for (const observer of propertyObservers) {
try {
if (typeof observer === 'function') {
observer(newValue, oldValue, property);
} else if (observer.next) {
observer.next(newValue, oldValue, property);
}
} catch (error) {
console.error(`Error in observer for property "${property}":`, error);
}
}
}
// Notify any-change observers
for (const observer of this._onAnyChangeCallbacks) {
try {
if (typeof observer === 'function') {
observer(property, newValue, oldValue);
} else if (observer.next) {
observer.next(property, newValue, oldValue);
}
} catch (error) {
console.error('Error in any-change observer:', error);
}
}
}
toJSON() {
return { ...this._data };
}
}Put it to work
// Simple observable
const count = new Observable(0);
// Subscribe to changes
const unsubscribe = count.subscribe((newValue, oldValue) => {
console.log(`Count changed from ${oldValue} to ${newValue}`);
});
// Change value (triggers observers)
count.value = 1;
count.value = 2;
count.value = 2; // No change, no notification
// Unsubscribe
unsubscribe();
count.value = 3; // No console log
// Observable with operators
const searchQuery = new Observable('');
const debouncedSearch = searchQuery.debounceTime(300);
const filteredSearch = debouncedSearch.filter(query => query.length >= 3);
filteredSearch.subscribe(query => {
console.log('Searching for:', query);
});
// Simulate user typing
searchQuery.value = 'h';
searchQuery.value = 'he';
searchQuery.value = 'hel';
searchQuery.value = 'hell';
searchQuery.value = 'hello';
// After 300ms: "Searching for: hello"
// Observable object
const user = new ObservableObject({
name: 'John',
age: 30,
email: '[email protected]'
});
// Subscribe to specific property
user.subscribe('age', (newAge, oldAge) => {
console.log(`Age changed from ${oldAge} to ${newAge}`);
});
// Subscribe to any property change
user.subscribeToAll((property, newValue, oldValue) => {
console.log(`${property} changed from ${oldValue} to ${newValue}`);
});
// Modify properties
user.data.age = 31; // Triggers both observers
user.data.name = 'Jane'; // Triggers any-change observer
user.data.email = '[email protected]'; // Triggers any-change observer
// Create computed observable
const price = new Observable(100);
const taxRate = new Observable(0.08);
const total = new Observable(price.value * (1 + taxRate.value));
price.subscribe(() => {
total.value = price.value * (1 + taxRate.value);
});
taxRate.subscribe(() => {
total.value = price.value * (1 + taxRate.value);
});
total.subscribe(value => {
console.log('Total updated:', value);
});
price.value = 150; // Triggers total update