JavaScripthardJavaScript

Observer Utilities

Observable pattern implementations for reactive programming

01

The problem

Need reactive data structures that notify observers of changes

02

The solution

JavaScript
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 };
  }
}

03

Put it to work

Example
// 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