npm package discovery and stats viewer.

Discover Tips

  • General search

    [free text search, go nuts!]

  • Package details

    pkg:[package-name]

  • User packages

    @[username]

Sponsor

Optimize Toolset

I’ve always been into building performant and accessible sites, but lately I’ve been taking it extremely seriously. So much so that I’ve been building a tool to help me optimize and monitor the sites that I build to make sure that I’m making an attempt to offer the best experience to those who visit them. If you’re into performant, accessible and SEO friendly sites, you might like it too! You can check it out at Optimize Toolset.

About

Hi, 👋, I’m Ryan Hefner  and I built this site for me, and you! The goal of this site was to provide an easy way for me to check the stats on my npm packages, both for prioritizing issues and updates, and to give me a little kick in the pants to keep up on stuff.

As I was building it, I realized that I was actually using the tool to build the tool, and figured I might as well put this out there and hopefully others will find it to be a fast and useful way to search and browse npm packages as I have.

If you’re interested in other things I’m working on, follow me on Twitter or check out the open source projects I’ve been publishing on GitHub.

I am also working on a Twitter bot for this site to tweet the most popular, newest, random packages from npm. Please follow that account now and it will start sending out packages soon–ish.

Open Software & Tools

This site wouldn’t be possible without the immense generosity and tireless efforts from the people who make contributions to the world and share their work via open source initiatives. Thank you 🙏

© 2026 – Pkg Stats / Ryan Hefner

@inputless/dispatcher

v1.0.6

Published

Core signal dispatcher for routing behavioral signals to registered modules

Readme

@inputless/dispatcher

Core signal dispatcher for routing behavioral signals to registered modules.

Purpose

Acts as a central routing hub that captures behavioral signals from @inputless/tracker and dispatches them to registered modules (such as @inputless/ads, analytics processors, and other plugins).

Architecture

┌─────────────────────────────────────────┐
│   @inputless/tracker                   │
│   (Signal Capture)                     │
└──────────────┬─────────────────────────┘
               │
               │ Signals
               │
┌──────────────▼─────────────────────────┐
│   @inputless/dispatcher                │
│   ├── SignalReceiver                   │
│   ├── RoutingEngine                    │
│   ├── ModuleRegistry                  │
│   └── SignalQueue                     │
└──────────────┬─────────────────────────┘
               │
               │ Dispatched Signals
               │
    ┌──────────┴──────────┐
    │                     │
┌───▼─────────┐    ┌──────▼──────────┐
│ @inputless/ │    │ Other Modules   │
│    ads      │    │ (Analytics, etc)│
└─────────────┘    └─────────────────┘

Features

Signal Capture

  • Receives all behavioral signals from tracker
  • Buffers and queues signals for processing
  • Handles high-volume signal streams

Routing Engine

  • Routes signals to registered modules based on rules
  • Supports multiple routing strategies (broadcast, selective, filtered)
  • Configurable routing rules per module

Module Registry

  • Dynamic module registration/unregistration
  • Module priority management
  • Module health monitoring

Signal Queue

  • Async signal processing
  • Priority queues for time-sensitive signals
  • Backpressure handling

Installation

npm install @inputless/dispatcher

Usage

Basic Setup (Local Only - No API Endpoint)

import { SignalDispatcher } from '@inputless/dispatcher';
import { InputlessTracker } from '@inputless/tracker';
import type { BaseEvent } from '@inputless/events';

// Initialize dispatcher with default configuration (all processing is local)
const dispatcher = new SignalDispatcher({
  // No policy endpoint - all routing is local
  policy: {
    enabled: false, // Disable policy loading from backend
  },
  debug: true, // Enable debug logging
});

// Start the dispatcher
await dispatcher.start();

// Register a module handler
await dispatcher.registerModule('analytics', {
  onSignal: async (signal: BaseEvent) => {
    // Process signal locally (no API calls)
    console.log('Received signal:', {
      id: signal.id,
      type: signal.type,
      timestamp: new Date(signal.timestamp).toISOString(),
      sessionId: signal.sessionId,
    });
    
    // Return delivery result
    return {
      success: true,
      latencyMs: 10,
    };
  },
  onError: (error: Error, context) => {
    console.error('Delivery error:', error, context);
  },
});

// Initialize tracker WITHOUT API endpoint (local only)
const tracker = new InputlessTracker({
  enabledCategories: { ui: true, form: true },
  debug: true,
});

// Connect tracker to dispatcher
tracker.onEvent(async (event: BaseEvent) => {
  // Receive signals from tracker
  await dispatcher.receive(event);
});

tracker.start();

Example 1: Local Module Registration and Signal Processing

import { SignalDispatcher } from '@inputless/dispatcher';
import { InputlessTracker } from '@inputless/tracker';
import type { BaseEvent, DeliveryResult } from '@inputless/events';
import type { ModuleHandler, ModuleConfig } from '@inputless/dispatcher';

// Initialize dispatcher (local only)
const dispatcher = new SignalDispatcher({
  queue: {
    maxSize: 1000,
    overflowStrategy: 'drop-oldest',
  },
  workers: {
    count: 2,
    batchSize: 10,
    pollInterval: 100,
  },
  debug: true,
});

await dispatcher.start();

// Local event storage for analytics module
const analyticsEvents: BaseEvent[] = [];

// Register analytics module
const analyticsHandler: ModuleHandler = {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    // Store events locally
    analyticsEvents.push(signal);
    
    // Keep only last 1000 events
    if (analyticsEvents.length > 1000) {
      analyticsEvents.shift();
    }
    
    console.log(`Analytics: Stored event ${signal.type} (Total: ${analyticsEvents.length})`);
    
    // Return delivery result
    return {
      success: true,
      latencyMs: 5,
    };
  },
  onError: (error: Error, context) => {
    console.error('Analytics error:', error, context);
  },
};

// Module configuration
const analyticsConfig: ModuleConfig = {
  type: 'analytics',
  priority: 'normal',
  filters: {
    channel: ['analytics'], // Only accept signals with 'analytics' channel
  },
  enabled: true,
};

await dispatcher.registerModule('analytics', analyticsHandler, analyticsConfig);

// Initialize tracker (local only)
const tracker = new InputlessTracker({
  enabledCategories: { ui: true, form: true },
});

// Connect tracker to dispatcher
tracker.onEvent(async (event: BaseEvent) => {
await dispatcher.receive(event);
});

tracker.start();

// Get registered modules
const modules = dispatcher.getModules();
console.log('Registered modules:', modules.map(m => ({
  id: m.id,
  type: m.type,
  health: m.health.status,
})));

// Get module health
const health = dispatcher.getModuleHealth('analytics');
if (health) {
  console.log('Analytics module health:', {
    status: health.status, // 'healthy' | 'degraded' | 'unhealthy'
    lastSuccess: health.lastSuccess,
    lastFailure: health.lastFailure,
    consecutiveFailures: health.consecutiveFailures,
  });
}

Example 2: Advanced Local Configuration

import { SignalDispatcher, type DispatcherConfig } from '@inputless/dispatcher';
import { InputlessTracker } from '@inputless/tracker';
import type { BaseEvent } from '@inputless/events';

// Configure dispatcher with custom settings (all local)
const config: DispatcherConfig = {
  queue: {
    maxSize: 5000,              // Maximum queue size per module
    maxGlobalSize: 10000,       // Maximum queue size globally
    overflowStrategy: 'drop-oldest', // Strategy when queue is full
  },
  workers: {
    count: 4,                   // Number of delivery workers per module
    batchSize: 10,              // Signals per batch
    pollInterval: 100,          // Queue polling interval (ms)
  },
  retry: {
    maxAttempts: 5,             // Maximum retry attempts
    initialDelay: 1000,         // Initial retry delay (ms)
    maxDelay: 30000,            // Maximum retry delay (ms)
    backoffMultiplier: 2,       // Exponential backoff multiplier
    jitter: 0.1,                // Jitter factor (0-1)
  },
  circuitBreaker: {
    failureThreshold: 5,        // Failures before opening circuit
    successThreshold: 2,        // Successes before closing circuit
    timeout: 60000,             // Circuit open timeout (ms)
    windowMs: 60000,           // Rolling window for failures (ms)
  },
  policy: {
    enabled: false,             // Disable policy loading (local only)
  },
  metrics: {
    enabled: true,              // Enable metrics collection
    exportInterval: 60000,     // Metrics export interval (ms)
  },
  debug: true,                 // Enable debug logging
};

const dispatcher = new SignalDispatcher(config);
await dispatcher.start();

// Register multiple modules locally
await dispatcher.registerModule('analytics', {
  onSignal: async (signal: BaseEvent) => {
    console.log('Analytics:', signal.type);
    return { success: true, latencyMs: 10 };
  },
}, {
  type: 'analytics',
  priority: 'normal',
  filters: { channel: ['analytics'] },
});

await dispatcher.registerModule('ads', {
  onSignal: async (signal: BaseEvent) => {
    console.log('Ads:', signal.type);
    return { success: true, latencyMs: 5 };
  },
}, {
  type: 'ads',
  priority: 'high',
  filters: { channel: ['ads'] },
});

// Initialize tracker (local only)
const tracker = new InputlessTracker({
  enabledCategories: { ui: true, form: true },
});

tracker.onEvent(async (event: BaseEvent) => {
  await dispatcher.receive(event);
});

tracker.start();

Example 3: Multiple Local Modules with Different Priorities

import { SignalDispatcher } from '@inputless/dispatcher';
import { InputlessTracker } from '@inputless/tracker';
import type { BaseEvent, DeliveryResult } from '@inputless/events';
import type { ModuleHandler, ModuleConfig } from '@inputless/dispatcher';

const dispatcher = new SignalDispatcher({ debug: true });
await dispatcher.start();

// Local storage for each module
const analyticsData: BaseEvent[] = [];
const adsData: BaseEvent[] = [];
const securityData: BaseEvent[] = [];

// Register analytics module
const analyticsModule: ModuleHandler = {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    analyticsData.push(signal);
    console.log(`[Analytics] Processed: ${signal.type}`);
    return { success: true, latencyMs: 10 };
  },
  onError: (error: Error, context) => {
    console.error('[Analytics] Error:', error, context);
  },
  health: async () => {
    return {
      healthy: true,
      message: 'Analytics module healthy',
    };
  },
};

// Module configuration
const analyticsConfig: ModuleConfig = {
  type: 'analytics',
  priority: 'normal',           // Module priority
  filters: {
    channel: ['analytics'],
    type: ['ui.*', 'form.*'],    // Accept UI and form events
  },
  enabled: true,
  maxInFlight: 10,              // Maximum concurrent signals
};

await dispatcher.registerModule('analytics', analyticsModule, analyticsConfig);

// Register ads module with high priority
await dispatcher.registerModule('ads', {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    adsData.push(signal);
    console.log(`[Ads] Processed: ${signal.type}`);
    return { success: true, latencyMs: 5 };
  },
  onError: (error: Error) => {
    console.error('[Ads] Error:', error);
  },
}, {
  type: 'ads',
  priority: 'high',             // High priority for ads
  filters: {
    channel: ['ads'],
  },
  enabled: true,
});

// Register security module with high priority
await dispatcher.registerModule('security', {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    securityData.push(signal);
    console.log(`[Security] Processed: ${signal.type}`);
    
    // Check for suspicious patterns locally
    if (signal.type.startsWith('error.')) {
      console.warn('[Security] Error event detected:', signal);
    }
    
    return { success: true, latencyMs: 8 };
  },
  onError: (error: Error) => {
    console.error('[Security] Error:', error);
  },
}, {
  type: 'security',
  priority: 'high',             // High priority for security
  filters: {
    channel: ['security'],
    type: ['error.*'],           // Only error events
  },
  enabled: true,
});

// Initialize tracker (local only)
const tracker = new InputlessTracker({
  enabledCategories: { ui: true, form: true },
});

tracker.onEvent(async (event: BaseEvent) => {
  await dispatcher.receive(event);
});

tracker.start();

// Get all registered modules
setInterval(() => {
  const modules = dispatcher.getModules();
  console.log('Registered modules:', modules.map(m => ({
    id: m.id,
    type: m.type,
    priority: m.config.priority,
    health: m.health.status,
  })));
}, 10000);

Example 4: Local Custom Routing Rules

import { SignalDispatcher, type RoutingRule, type RoutingPolicy } from '@inputless/dispatcher';
import { InputlessTracker } from '@inputless/tracker';
import type { BaseEvent } from '@inputless/events';

const dispatcher = new SignalDispatcher({
  policy: { enabled: false }, // Local routing only
  debug: true,
});

await dispatcher.start();

// Define custom routing rules
const customRules: RoutingRule[] = [
  {
    id: 'route-performance-events',
    when: {
      type: 'perf.*',  // Wildcard: all performance events
    },
    action: 'route',
    targets: [
      {
        moduleId: 'analytics',
        priority: 'high',
      },
    ],
  },
  {
    id: 'route-error-events',
    when: {
      type: 'error.*',
    },
    action: 'route',
    targets: [
      {
        moduleId: 'security',
        priority: 'high',
        deadlineMs: 5000,  // Must be delivered within 5 seconds
      },
    ],
  },
  {
    id: 'route-ui-clicks-to-ads',
    when: {
      type: 'ui.click',
      channel: 'ads',
    },
    action: 'route',
    targets: [
      {
        moduleId: 'ads',
        priority: 'high',
      },
    ],
  },
  {
    id: 'deny-pii-signals',
    when: {
      pii: true,  // Block PII signals
    },
    action: 'deny',  // Block signal
  },
];

// Load routing policy
const policy: RoutingPolicy = {
  rules: customRules,
  version: '1.0',
  expires: Date.now() + 3600000, // 1 hour
};

await dispatcher.loadPolicy(policy);

// Register modules
await dispatcher.registerModule('analytics', {
  onSignal: async (signal: BaseEvent) => {
    console.log('[Analytics] Performance event:', signal.type);
    return { success: true, latencyMs: 10 };
  },
}, {
  type: 'analytics',
  filters: { type: ['perf.*'] },
});

await dispatcher.registerModule('security', {
  onSignal: async (signal: BaseEvent) => {
    console.warn('[Security] Error event:', signal.type);
    return { success: true, latencyMs: 5 };
  },
}, {
  type: 'security',
  filters: { type: ['error.*'] },
});

await dispatcher.registerModule('ads', {
  onSignal: async (signal: BaseEvent) => {
    console.log('[Ads] Click event:', signal.type);
    return { success: true, latencyMs: 3 };
  },
}, {
  type: 'ads',
  filters: { channel: ['ads'] },
});

// Initialize tracker (local only)
const tracker = new InputlessTracker({
  enabledCategories: { ui: true, performance: true },
});

tracker.onEvent(async (event: BaseEvent) => {
  await dispatcher.receive(event);
});

tracker.start();

Example 5: Local Module Health Monitoring

import { SignalDispatcher } from '@inputless/dispatcher';
import { InputlessTracker } from '@inputless/tracker';
import type { BaseEvent, DeliveryResult } from '@inputless/events';
import type { ModuleHealth } from '@inputless/dispatcher';

const dispatcher = new SignalDispatcher({ debug: true });
await dispatcher.start();

// Register module with health monitoring
await dispatcher.registerModule('analytics', {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    // Simulate occasional failures for health monitoring demo
    if (Math.random() < 0.1) {
      throw new Error('Simulated failure');
    }
    
    console.log('Processed:', signal.type);
    return { success: true, latencyMs: 10 };
  },
  onError: (error: Error) => {
    console.error('Error:', error);
  },
  health: async () => {
    // Custom health check
    return {
      healthy: true,
      message: 'Analytics module is healthy',
      details: {
        eventsProcessed: 1000,
        avgLatency: 10,
      },
    };
  },
}, {
  type: 'analytics',
  enabled: true,
});

// Initialize tracker (local only)
const tracker = new InputlessTracker({
  enabledCategories: { ui: true },
});

tracker.onEvent(async (event: BaseEvent) => {
  await dispatcher.receive(event);
});

tracker.start();

// Monitor module health
setInterval(() => {
  const health: ModuleHealth | null = dispatcher.getModuleHealth('analytics');
  
  if (health) {
    console.log('Module health:', {
      status: health.status,        // 'healthy' | 'degraded' | 'unhealthy'
      lastSuccess: health.lastSuccess,
      lastFailure: health.lastFailure,
      consecutiveFailures: health.consecutiveFailures,
      totalSignals: health.totalSignals,
      totalFailures: health.totalFailures,
      successRate: health.totalSignals > 0
        ? ((health.totalSignals - health.totalFailures) / health.totalSignals) * 100
        : 100,
    });
  }
  
  // Get all registered modules
  const modules = dispatcher.getModules();
  console.log('Registered modules:', modules.map(m => ({
    id: m.id,
    type: m.type,
    health: m.health.status,
    metrics: {
      totalSignals: m.health.totalSignals,
      totalFailures: m.health.totalFailures,
    },
  })));
}, 5000);

Example 6: Local Metrics and Observability

import { SignalDispatcher } from '@inputless/dispatcher';
import { InputlessTracker } from '@inputless/tracker';
import type { BaseEvent } from '@inputless/events';
import type { DispatcherMetrics } from '@inputless/dispatcher';

const dispatcher = new SignalDispatcher({
  metrics: {
    enabled: true,
    exportInterval: 60000,
  },
  debug: true,
});

await dispatcher.start();

// Register module
await dispatcher.registerModule('analytics', {
  onSignal: async (signal: BaseEvent) => {
    console.log('Processed:', signal.type);
    return { success: true, latencyMs: Math.random() * 50 };
  },
}, {
  type: 'analytics',
  enabled: true,
});

// Initialize tracker (local only)
const tracker = new InputlessTracker({
  enabledCategories: { ui: true, form: true },
});

tracker.onEvent(async (event: BaseEvent) => {
  await dispatcher.receive(event);
});

tracker.start();

// Monitor metrics
setInterval(() => {
  const metrics: DispatcherMetrics = dispatcher.getMetrics();
  
  console.log('📊 Dispatcher Metrics:', {
    counters: {
      signalsReceived: metrics.counters['signals.received'] || 0,
      signalsDelivered: metrics.counters['delivery.success'] || 0,
      signalsFailed: metrics.counters['delivery.failure'] || 0,
      signalsDropped: metrics.counters['signals.dropped'] || 0,
    },
    gauges: {
      queueDepth: metrics.gauges['queue.depth'] || 0,
      activeWorkers: metrics.gauges['workers.active'] || 0,
    },
    histograms: {
      deliveryLatency: metrics.histograms['delivery.latency'] ? {
        avg: metrics.histograms['delivery.latency'].avg,
        p50: metrics.histograms['delivery.latency'].p50,
        p95: metrics.histograms['delivery.latency'].p95,
        p99: metrics.histograms['delivery.latency'].p99,
        min: metrics.histograms['delivery.latency'].min,
        max: metrics.histograms['delivery.latency'].max,
      } : null,
    },
  });
}, 10000);

Example 7: Complete Local Integration - Tracker → Dispatcher → Modules

import { SignalDispatcher } from '@inputless/dispatcher';
import { InputlessTracker } from '@inputless/tracker';
import type { BaseEvent, DeliveryResult } from '@inputless/events';

// Initialize dispatcher (local only - no API endpoints)
const dispatcher = new SignalDispatcher({
  queue: {
    maxSize: 1000,
    overflowStrategy: 'drop-oldest',
  },
  workers: {
    count: 2,
    batchSize: 10,
    pollInterval: 100,
  },
  retry: {
    maxAttempts: 3,
    initialDelay: 100,
    maxDelay: 10000,
    backoffMultiplier: 2,
    jitter: 0.1,
  },
  circuitBreaker: {
    failureThreshold: 5,
    successThreshold: 2,
    timeout: 60000,
  },
  policy: {
    enabled: false, // Local routing only
  },
  metrics: {
    enabled: true,
  },
  debug: true,
});

await dispatcher.start();

// Local storage for each module
const moduleData = {
  analytics: [] as BaseEvent[],
  ads: [] as BaseEvent[],
  security: [] as BaseEvent[],
};

// Register analytics module
await dispatcher.registerModule('analytics', {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    moduleData.analytics.push(signal);
    console.log(`[Analytics] Stored: ${signal.type} (Total: ${moduleData.analytics.length})`);
    return { success: true, latencyMs: 10 };
  },
  onError: (error: Error) => {
    console.error('[Analytics] Error:', error);
  },
}, {
  type: 'analytics',
  priority: 'normal',
  filters: { channel: ['analytics'] },
});

// Register ads module
await dispatcher.registerModule('ads', {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    moduleData.ads.push(signal);
    console.log(`[Ads] Stored: ${signal.type} (Total: ${moduleData.ads.length})`);
    return { success: true, latencyMs: 5 };
  },
  onError: (error: Error) => {
    console.error('[Ads] Error:', error);
  },
}, {
  type: 'ads',
  priority: 'high',
  filters: { channel: ['ads'] },
});

// Register security module
await dispatcher.registerModule('security', {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    moduleData.security.push(signal);
    console.log(`[Security] Stored: ${signal.type} (Total: ${moduleData.security.length})`);
    
    // Local threat detection
    if (signal.type.startsWith('error.')) {
      console.warn('[Security] Error event detected:', signal.type);
    }
    
    return { success: true, latencyMs: 8 };
  },
  onError: (error: Error) => {
    console.error('[Security] Error:', error);
  },
}, {
  type: 'security',
  priority: 'high',
  filters: { channel: ['security'] },
});

// Initialize tracker WITHOUT API endpoint (local only)
const tracker = new InputlessTracker({
  enabledCategories: {
    ui: true,
    form: true,
    performance: true,
  },
  debug: true,
});

// Connect tracker to dispatcher (all processing is local)
tracker.onEvent(async (event: BaseEvent) => {
  // Receive signals from tracker
  await dispatcher.receive(event);
});

tracker.start();

// Display summary every 10 seconds
setInterval(() => {
  const metrics = dispatcher.getMetrics();
  const modules = dispatcher.getModules();
  
  console.log('📊 Complete Integration Summary:', {
    dispatcher: {
      signalsReceived: metrics.counters['signals.received'] || 0,
      signalsDelivered: metrics.counters['delivery.success'] || 0,
      signalsFailed: metrics.counters['delivery.failure'] || 0,
      queueDepth: metrics.gauges['queue.depth'] || 0,
    },
    modules: modules.map(m => ({
      id: m.id,
      type: m.type,
      health: m.health.status,
      storedEvents: moduleData[m.id as keyof typeof moduleData]?.length || 0,
    })),
    storage: {
      analytics: moduleData.analytics.length,
      ads: moduleData.ads.length,
      security: moduleData.security.length,
    },
  });
}, 10000);

Example 8: Local Integration with Cognitive SDK

import { SignalDispatcher } from '@inputless/dispatcher';
import { InputlessTracker } from '@inputless/tracker';
import { CognitiveSDK } from '@inputless/sdk-cognitive';
import type { BaseEvent, DeliveryResult } from '@inputless/events';

// Initialize dispatcher (local only)
const dispatcher = new SignalDispatcher({
  policy: { enabled: false },
  debug: true,
});

await dispatcher.start();

// Local storage for cognitive insights
const cognitiveInsights: Array<{
  event: BaseEvent;
  insights: any;
}> = [];

// Register cognitive module
await dispatcher.registerModule('cognitive', {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    // Store enriched events from cognitive SDK
    cognitiveInsights.push({
      event: signal,
      insights: (signal as any).insights || null,
    });
    
    console.log('[Cognitive] Stored enriched event:', signal.type);
    return { success: true, latencyMs: 15 };
  },
  onError: (error: Error) => {
    console.error('[Cognitive] Error:', error);
  },
}, {
  type: 'custom',
  priority: 'normal',
  filters: { channel: ['cognitive'] },
});

// Initialize tracker (local only)
const tracker = new InputlessTracker({
  enabledCategories: { ui: true, form: true },
});

// Initialize cognitive SDK (local only)
const cognitiveSDK = new CognitiveSDK({
  // No API endpoint - all processing is local
  debug: true,
});

// Connect tracker → cognitive SDK → dispatcher
tracker.onEvent(async (event: BaseEvent) => {
  // Process through cognitive SDK first
  const enriched = await cognitiveSDK.perceive(event);
  
  // Then send to dispatcher
  if (enriched) {
    await dispatcher.receive(enriched);
  }
});

tracker.start();
cognitiveSDK.start();

// Display cognitive insights
setInterval(() => {
  console.log('🧠 Cognitive Insights:', {
    totalInsights: cognitiveInsights.length,
    recentInsights: cognitiveInsights.slice(-10).map(i => ({
      eventType: i.event.type,
      hasInsights: !!i.insights,
    })),
  });
}, 10000);

Example 9: Local Error Handling and Circuit Breakers

import { SignalDispatcher } from '@inputless/dispatcher';
import { InputlessTracker } from '@inputless/tracker';
import type { BaseEvent, DeliveryResult } from '@inputless/events';

const dispatcher = new SignalDispatcher({
  circuitBreaker: {
    failureThreshold: 5,        // Failures before opening circuit
    successThreshold: 2,        // Successes before closing circuit
    timeout: 60000,             // Circuit open timeout (ms)
    windowMs: 60000,            // Rolling window for failures (ms)
  },
  retry: {
    maxAttempts: 3,
    initialDelay: 100,
    maxDelay: 10000,
    backoffMultiplier: 2,
    jitter: 0.1,
  },
  debug: true,
});

await dispatcher.start();

// Simulate a module that fails occasionally
let failureCount = 0;
const maxFailures = 3;

// Circuit breakers automatically protect against failing modules
await dispatcher.registerModule('unreliable-module', {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    // Simulate failures for first few signals
    failureCount++;
    if (failureCount <= maxFailures) {
      throw new Error(`Simulated failure ${failureCount}/${maxFailures}`);
    }
    
    // After failures, succeed
    console.log('[Unreliable] Processed:', signal.type);
    return { success: true, latencyMs: 10 };
  },
  onError: (error: Error, context) => {
    // Handle errors - circuit breaker tracks failures
    console.error('[Unreliable] Error:', error.message, {
      attempt: context.attempt,
      moduleId: context.moduleId,
      signalId: context.signal?.id,
    });
  },
}, {
  type: 'custom',
  priority: 'normal',
  enabled: true,
});

// Initialize tracker (local only)
const tracker = new InputlessTracker({
  enabledCategories: { ui: true },
});

tracker.onEvent(async (event: BaseEvent) => {
  await dispatcher.receive(event);
});

tracker.start();

// Monitor circuit breaker state
setInterval(() => {
  const module = dispatcher.getModules().find(m => m.id === 'unreliable-module');
  if (module) {
    console.log('Circuit Breaker State:', {
      moduleId: module.id,
      health: module.health.status,
      consecutiveFailures: module.health.consecutiveFailures,
      lastFailure: module.health.lastFailure,
      lastSuccess: module.health.lastSuccess,
    });
    
    // Circuit breaker state is managed automatically:
    // - After 5 failures → circuit opens (stops sending signals)
    // - After 60 seconds → circuit half-opens (tries one signal)
    // - After 2 successes → circuit closes (normal operation)
  }
}, 5000);

Example 10: Local Module Filtering and Selective Routing

import { SignalDispatcher } from '@inputless/dispatcher';
import { InputlessTracker } from '@inputless/tracker';
import type { BaseEvent, DeliveryResult } from '@inputless/events';

const dispatcher = new SignalDispatcher({ debug: true });
await dispatcher.start();

// Register module with specific filters
await dispatcher.registerModule('ui-only', {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    console.log('[UI-Only] Received:', signal.type);
    return { success: true, latencyMs: 5 };
  },
}, {
  type: 'analytics',
  filters: {
    type: ['ui.click', 'ui.hover', 'ui.scroll'], // Only UI events
    channel: ['analytics'],
  },
  enabled: true,
});

// Register module for form events only
await dispatcher.registerModule('form-only', {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    console.log('[Form-Only] Received:', signal.type);
    return { success: true, latencyMs: 5 };
  },
}, {
  type: 'analytics',
  filters: {
    type: ['form.*'], // All form events (wildcard)
    channel: ['analytics'],
  },
  enabled: true,
});

// Register module for performance events only
await dispatcher.registerModule('performance-only', {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    console.log('[Performance-Only] Received:', signal.type);
    return { success: true, latencyMs: 5 };
  },
}, {
  type: 'analytics',
  filters: {
    type: ['perf.*'], // All performance events (wildcard)
  },
  enabled: true,
});

// Initialize tracker (local only)
const tracker = new InputlessTracker({
  enabledCategories: {
    ui: true,
    form: true,
    performance: true,
  },
});

tracker.onEvent(async (event: BaseEvent) => {
  await dispatcher.receive(event);
});

tracker.start();

// Display routing statistics
setInterval(() => {
  const modules = dispatcher.getModules();
  console.log('Module Routing Statistics:', {
    modules: modules.map(m => ({
      id: m.id,
      filters: m.config.filters,
      signalsReceived: m.health.totalSignals,
    })),
  });
}, 10000);

Example 11: Graceful Shutdown and Cleanup

import { SignalDispatcher } from '@inputless/dispatcher';
import { InputlessTracker } from '@inputless/tracker';
import type { BaseEvent } from '@inputless/events';

const dispatcher = new SignalDispatcher({ debug: true });
await dispatcher.start();

await dispatcher.registerModule('analytics', {
  onSignal: async (signal: BaseEvent) => {
    console.log('Processed:', signal.type);
    return { success: true, latencyMs: 10 };
  },
});

const tracker = new InputlessTracker({
  enabledCategories: { ui: true },
});

tracker.onEvent(async (event: BaseEvent) => {
  await dispatcher.receive(event);
});

tracker.start();

// Graceful shutdown on page unload
window.addEventListener('beforeunload', async () => {
  // Stop dispatcher gracefully
  await dispatcher.stop();
  
  // This will:
  // - Stop all delivery workers
  // - Flush pending signals
  // - Stop policy refresh
  // - Clean up resources
  
  // Stop tracker
  await tracker.stop();
});

// Or manually stop when needed
// await dispatcher.stop();

Example 12: Complete Local Integration - All Components Working Together

import { SignalDispatcher } from '@inputless/dispatcher';
import { InputlessTracker } from '@inputless/tracker';
import { ContextEngine } from '@inputless/context';
import { CognitiveSDK } from '@inputless/sdk-cognitive';
import type { BaseEvent, DeliveryResult } from '@inputless/events';

// Initialize all components (all local - no API endpoints)
const dispatcher = new SignalDispatcher({
  queue: { maxSize: 1000 },
  workers: { count: 2, batchSize: 10 },
  retry: { maxAttempts: 3 },
  circuitBreaker: { failureThreshold: 5 },
  policy: { enabled: false }, // Local routing only
  metrics: { enabled: true },
  debug: true,
});

const tracker = new InputlessTracker({
  enabledCategories: { ui: true, form: true, performance: true },
  debug: true,
});

const contextEngine = new ContextEngine({
  patterns: { behavioral: true, sequence: true },
  anomaly: { enabled: true },
  debug: true,
});

const cognitiveSDK = new CognitiveSDK({
  debug: true,
});

// Start all components
await dispatcher.start();
contextEngine.start();
cognitiveSDK.start();

// Local storage
const storage = {
  analytics: [] as BaseEvent[],
  cognitive: [] as BaseEvent[],
  context: [] as any[],
};

// Register analytics module
await dispatcher.registerModule('analytics', {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    storage.analytics.push(signal);
    return { success: true, latencyMs: 10 };
  },
}, {
  type: 'analytics',
  filters: { channel: ['analytics'] },
});

// Register cognitive module
await dispatcher.registerModule('cognitive', {
  onSignal: async (signal: BaseEvent): Promise<DeliveryResult> => {
    storage.cognitive.push(signal);
    return { success: true, latencyMs: 15 };
  },
}, {
  type: 'custom',
  filters: { channel: ['cognitive'] },
});

// Complete pipeline: Tracker → Cognitive SDK → Context Engine → Dispatcher
tracker.onEvent(async (event: BaseEvent) => {
  // Step 1: Process through cognitive SDK
  const enriched = await cognitiveSDK.perceive(event);
  
  if (enriched) {
    // Step 2: Process through context engine
    const contextResult = await contextEngine.process([enriched]);
    storage.context.push(contextResult);
    
    // Step 3: Send to dispatcher
    await dispatcher.receive(enriched);
  }
});

tracker.start();

// Display complete integration dashboard
setInterval(() => {
  const metrics = dispatcher.getMetrics();
  const modules = dispatcher.getModules();
  const contextMetrics = contextEngine.getMetrics();
  
  console.log('🚀 Complete Integration Dashboard:', {
    tracker: {
      sessionId: tracker.getSessionId(),
      queueSize: tracker.getQueueSize(),
    },
    dispatcher: {
      signalsReceived: metrics.counters['signals.received'] || 0,
      signalsDelivered: metrics.counters['delivery.success'] || 0,
      queueDepth: metrics.gauges['queue.depth'] || 0,
      modules: modules.length,
    },
    context: {
      patternsDetected: contextMetrics.patternsDetected,
      anomaliesDetected: contextMetrics.anomaliesDetected,
      eventsProcessed: contextMetrics.eventsProcessed,
    },
    storage: {
      analytics: storage.analytics.length,
      cognitive: storage.cognitive.length,
      context: storage.context.length,
    },
  });
}, 10000);

API Reference

SignalDispatcher

Main Methods:

  • start(): Promise<void> - Start the dispatcher
  • stop(): Promise<void> - Stop the dispatcher
  • receive(signal: BaseEvent | unknown): Promise<void> - Receive a signal
  • registerModule(id: string, handler: ModuleHandler, config?: ModuleConfig): Promise<void> - Register a module
  • unregisterModule(id: string): Promise<void> - Unregister a module
  • getModules(): ModuleDescriptor[] - Get registered modules
  • getModuleHealth(moduleId: string): ModuleHealth | null - Get module health status
  • getMetrics(): DispatcherMetrics - Get dispatcher metrics
  • loadPolicy(policy: RoutingPolicy): Promise<void> - Load routing policy
  • updateConfig(config: Partial<DispatcherConfig>): void - Update configuration

ModuleHandler

interface ModuleHandler {
  onSignal(signal: BaseEvent): Promise<DeliveryResult>;
  onError?(error: Error, context?: ErrorContext): void;
  health?(): Promise<HealthResult>;
}

Note: onSignal must return a DeliveryResult, not void.

Configuration

See DispatcherConfig type for full configuration options including:

  • Queue settings: maxSize, maxGlobalSize, overflowStrategy
  • Worker settings: count, batchSize, pollInterval
  • Retry settings: maxAttempts, initialDelay, maxDelay, backoffMultiplier, jitter
  • Circuit breaker settings: failureThreshold, successThreshold, timeout, windowMs
  • Routing rules: rules (array of RoutingRule)
  • Policy settings: enabled, endpoint, refreshInterval
  • Metrics settings: enabled, exportInterval
  • Debug mode: debug

Module Structure

dispatcher/
├── src/
│   ├── SignalDispatcher.ts      // Main dispatcher class
│   ├── SignalReceiver.ts        // Receives signals from tracker
│   ├── RoutingEngine.ts         // Routes signals to modules
│   ├── ModuleRegistry.ts        // Manages registered modules
│   ├── SignalQueue.ts           // Queue management
│   ├── RoutingRules.ts          // Routing rule definitions
│   └── index.ts                 // Main exports
├── package.json
└── README.md

Dependencies

  • @inputless/tracker - Source of behavioral signals
  • @inputless/events - Event type definitions

See Also