@inputless/dispatcher
v1.0.6
Published
Core signal dispatcher for routing behavioral signals to registered modules
Maintainers
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/dispatcherUsage
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 dispatcherstop(): Promise<void>- Stop the dispatcherreceive(signal: BaseEvent | unknown): Promise<void>- Receive a signalregisterModule(id: string, handler: ModuleHandler, config?: ModuleConfig): Promise<void>- Register a moduleunregisterModule(id: string): Promise<void>- Unregister a modulegetModules(): ModuleDescriptor[]- Get registered modulesgetModuleHealth(moduleId: string): ModuleHealth | null- Get module health statusgetMetrics(): DispatcherMetrics- Get dispatcher metricsloadPolicy(policy: RoutingPolicy): Promise<void>- Load routing policyupdateConfig(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.mdDependencies
@inputless/tracker- Source of behavioral signals@inputless/events- Event type definitions
