liblsl_coordinator 0.4.1
liblsl_coordinator: ^0.4.1 copied to clipboard
A performance-focused Dart (and Flutter) LSL-based device coordination library.
example/liblsl_coordinator_example.dart
// packages/liblsl_coordinator/example/liblsl_coordinator_example.dart
import 'dart:async';
import 'dart:io';
import 'dart:math';
import 'package:liblsl_coordinator/liblsl_coordinator.dart';
// import 'package:liblsl_coordinator/transports/lsl.dart';
import 'package:logging/logging.dart';
import 'package:peer_coordinator/in_memory.dart';
/// Multi-node coordination test to verify election, joining, messaging, and stream management
Future<void> main(List<String> args) async {
// Configure logging
Logger.root.level = Level.INFO;
Logger.root.onRecord.listen(Log.defaultPrinter);
// Parse command line arguments
final nodeCount = args.isNotEmpty ? int.tryParse(args[0]) ?? 3 : 3;
/// Maxnodes includes the coordinator
final maxNodes = args.length > 1
? int.tryParse(args[1]) ?? 2
: 3; // Set low to test rejection
final testDuration = args.length > 2
? int.tryParse(args[2]) ?? 10
: 10; // seconds
logger.info('π Starting Multi-Node Coordination Test');
logger.info(' Nodes to start: $nodeCount');
logger.info(' Max nodes allowed: $maxNodes');
logger.info(' Test duration: ${testDuration}s');
logger.info(
' Expected: ${min(nodeCount, maxNodes)} nodes should join successfully',
);
logger.info('');
final sessionSuffix = Random().nextInt(10000);
final bus = InMemoryBus();
// Start nodes concurrently with slight delays to test election
final futures = <Future<void>>[];
final appId = "TestApp_${Random().nextInt(1000)}";
for (int i = 0; i < nodeCount; i++) {
futures.add(
_runNode(
sessionName: 'TestSession_$sessionSuffix',
nodeId: 'Node_$i',
appId: appId,
maxNodes: maxNodes,
delay: Duration(milliseconds: i * 500), // Stagger starts
testDuration: testDuration,
bus: bus,
),
);
}
// Wait for all nodes to complete
try {
await Future.wait(futures);
} catch (e) {
logger.info('β Test completed with errors: $e');
}
logger.info('π Multi-node coordination test completed');
// Give some time for cleanup
await Future.delayed(Duration(seconds: 2));
exit(0);
}
/// Run a single node instance
Future<void> _runNode({
required String sessionName,
required String nodeId,
required String appId,
required int maxNodes,
required Duration delay,
required int testDuration,
InMemoryBus? bus,
}) async {
// Stagger the start times to test election
if (delay > Duration.zero) {
logger.info(
'β³ $nodeId: Waiting ${delay.inMilliseconds}ms before starting...',
);
await Future.delayed(delay);
}
logger.info('π― $nodeId: Initializing...');
try {
// Create session configuration with the specified maxNodes
final sessionConfig = CoordinationSessionConfig(
name: sessionName,
heartbeatInterval: Duration(seconds: 1),
discoveryInterval: Duration(seconds: 5),
nodeTimeout: Duration(seconds: 10),
maxNodes: maxNodes, // This may cause some nodes to be rejected
consumeCoordinationStreamAsCoordinator: false,
);
// Create coordination configuration
// final coordinationConfig = CoordinationConfig(
// name: appId,
// sessionConfig: sessionConfig,
// topologyConfig: HierarchicalTopologyConfig(
// promotionStrategy: PromotionStrategyRandom(),
// maxNodes: maxNodes,
// ),
// streamConfig: CoordinationStreamConfig(
// name: 'coordination',
// sampleRate: 50.0,
// ),
// transportConfig: LSLTransportConfig(
// // This LSL API config specifically restricts to IPv4 and local machine
// // these wont go over the network
// lslApiConfig: LSLApiConfig(
// ipv6: IPv6Mode.disable,
// portRange: 128,
// logLevel: -2, // -2 Error only -> 9 is the most verbose
// resolveScope: ResolveScope.link,
// listenAddress: '127.0.0.1', // Use loopback for testing
// addressesOverride: ['224.0.0.183'],
// knownPeers: ['127.0.0.1'],
// ),
// coordinationFrequency: 50.0,
// ),
// );
// Create coordination configuration
final coordinationConfig = CoordinationConfig(
name: appId,
sessionConfig: sessionConfig,
topologyConfig: HierarchicalTopologyConfig(
promotionStrategy: PromotionStrategyRandom(),
maxNodes: maxNodes,
),
streamConfig: CoordinationStreamConfig(
name: 'coordination',
sampleRate: 50.0,
),
transportConfig: InMemoryTransportConfig(bus: bus ?? InMemoryBus()),
);
final dataStreamConfig = DataStreamConfig(
name: 'TestData',
channels: 3, // timestamp, node_id, sample_count
sampleRate: 10.0,
dataType: StreamDataType.double64,
// sendParticipantsReceiveCoordinator means that the participant nodes
// will not consume data sent either by other participants, or the coordinator
participationMode:
StreamParticipationMode.sendParticipantsReceiveCoordinator,
);
// Create session using the new simplified API
// final session = LSLCoordinationSession(
// coordinationConfig,
// thisNodeConfig: NodeConfigFactory().defaultConfig().copyWith(
// name: nodeId,
// ),
// );
final session = PeerSession.create(
coordinationConfig,
thisNodeConfig: NodeConfigFactory().defaultConfig().copyWith(
name: nodeId,
),
);
// Set up event listeners
_setupEventListeners(session, nodeId);
// Track test state
bool joinSuccessful = false;
bool testCompleted = false;
logger.info('π‘ $nodeId: Initializing session...');
await session.initialize();
logger.info('π $nodeId: Attempting to join coordination network...');
try {
await session.join();
joinSuccessful = true;
final role = session.isCoordinator ? 'Coordinator' : 'Participant';
logger.info('β
$nodeId: Successfully joined as $role');
logger.info(' Connected nodes: ${session.connectedNodes.length}');
// Run role-specific logic
if (session.isCoordinator) {
await _runCoordinatorTestLogic(
session,
nodeId,
testDuration,
maxNodes,
dataStreamConfig,
);
} else {
await _runParticipantTestLogic(
session,
nodeId,
testDuration,
dataStreamConfig,
);
}
testCompleted = true;
} catch (e) {
if (e.toString().contains('rejected')) {
logger.info('π« $nodeId: Join rejected (expected if > maxNodes): $e');
joinSuccessful = false;
} else {
logger.info('β $nodeId: Join failed unexpectedly: $e');
rethrow;
}
}
// Wait for test duration if joined successfully
if (joinSuccessful && !testCompleted) {
logger.info('β±οΈ $nodeId: Running test for ${testDuration}s...');
await Future.delayed(Duration(seconds: testDuration));
}
// Cleanup
if (joinSuccessful) {
logger.info('π§Ή $nodeId: Cleaning up...');
await session.leave();
await session.dispose();
}
logger.info('π $nodeId: Test completed successfully');
} catch (e, stack) {
logger.info('β $nodeId: Error during test: $e');
logger.info('π $nodeId: Stack trace: $stack');
rethrow;
}
}
/// Set up event listeners for testing
void _setupEventListeners(PeerSession session, String nodeId) {
// Phase changes
session.events.phaseChanges.listen((event) {
logger.info('π $nodeId: Phase changed to ${event.phase}');
});
// Node topology changes
session.events.nodeJoined.listen((event) {
logger.info(
'β $nodeId: Node joined: ${event.node.name} (${event.node.id})',
);
logger.info(' Total nodes: ${session.connectedNodes.length}');
});
session.events.nodeLeft.listen((event) {
logger.info('β $nodeId: Node left: ${event.node.name} (${event.node.id})');
logger.info(' Total nodes: ${session.connectedNodes.length}');
});
// User messages (coordination commands)
session.events.userCoordinationMessages.listen((event) {
logger.info(
'π¬ $nodeId: User Message: ${event.messageType} (${event.messageId}) - ${event.description}',
);
if (event.payload.isNotEmpty) {
logger.info(' Payload: ${event.payload}');
}
});
// Configuration updates
session.events.configUpdates.listen((event) {
logger.info('βοΈ $nodeId: Config Update: ${event.config}');
});
// Stream commands
session.events.streamStart.listen((event) {
logger.info('βΆοΈ $nodeId: Stream START command: ${event.streamName}');
if (event.startAt != null) {
logger.info(' Scheduled for: ${event.startAt}');
}
});
session.events.streamStop.listen((event) {
logger.info('βΉοΈ $nodeId: Stream STOP command: ${event.streamName}');
});
}
/// Coordinator test logic - manages the test sequence
Future<void> _runCoordinatorTestLogic(
PeerSession session,
String nodeId,
int testDuration,
int maxNodes,
DataStreamConfig streamConfig,
) async {
logger.info('π $nodeId: Running COORDINATOR test logic');
// Test sequence timeline
final testSteps = [
{'delay': 2, 'action': 'wait_for_nodes'},
{'delay': 3, 'action': 'pause_accepting'},
{'delay': 2, 'action': 'send_config'},
{'delay': 3, 'action': 'create_streams'},
{'delay': 2, 'action': 'start_data_collection'},
{'delay': testDuration, 'action': 'start_test_phase_1'},
{'delay': 2, 'action': 'pause_between_phases'},
{'delay': 1, 'action': 'resume_for_phase_2'},
{'delay': testDuration, 'action': 'start_test_phase_2'},
{'delay': 2, 'action': 'pause_before_stop'},
{'delay': 1, 'action': 'flush_and_resume'},
{'delay': 3, 'action': 'stop_data_collection'},
{'delay': 2, 'action': 'end_test'},
];
var elapsedTime = 0;
int messageCount = 0;
StreamSubscription? inboxSubscription;
DataStream? dataStream;
for (final step in testSteps) {
final delay = step['delay'] as int;
final action = step['action'] as String;
await Future.delayed(Duration(seconds: delay));
elapsedTime += delay;
try {
switch (action) {
case 'wait_for_nodes':
logger.info(
'β³ $nodeId: Waiting for [${maxNodes - 1}] participant nodes...',
);
try {
await session.waitForMinNodes(
maxNodes - 1,
timeout: Duration(seconds: 20),
);
logger.info(
'β
$nodeId: Participants joined (${session.connectedNodes.length} nodes)',
);
} catch (e) {
logger.info(
'β οΈ $nodeId: Timeout waiting for participants, continuing...',
);
}
break;
case 'pause_accepting':
logger.info('π $nodeId: Pausing acceptance of new nodes');
await session.pauseAcceptingNodes();
logger.info(' Is accepting nodes: ${session.isAcceptingNodes}');
break;
case 'send_config':
logger.info('π $nodeId: Broadcasting test configuration...');
await session.updateConfig({
'test_type': 'multi_node_coordination',
'test_duration': testDuration,
'data_rate': 10.0,
'channels': 3,
});
break;
case 'create_streams':
logger.info('π $nodeId: Creating test data stream...');
dataStream = await session.createDataStream(streamConfig);
inboxSubscription = dataStream.inbox.listen((data) {
// logger.info('π₯ $nodeId: Received data: $data');
messageCount++;
});
break;
case 'start_test_phase_1':
logger.info('π― $nodeId: Starting test phase 1...');
logger.info(' $nodeId: Current message count: $messageCount');
await session.sendUserMessage(
'start_test_phase',
'Starting coordinated test - Phase 1',
{
'phase': 1,
'intensity': 'low',
'start_at': DateTime.now()
.add(Duration(seconds: 5))
.toIso8601String(),
},
);
break;
case 'start_data_collection':
logger.info('π $nodeId: Starting data collection...');
logger.info(' $nodeId: Current message count: $messageCount');
await session.startStream('TestData');
break;
case 'pause_between_phases':
logger.info('βΈοΈ $nodeId: Pausing data streams between phases...');
logger.info(' $nodeId: Current message count: $messageCount');
await session.pauseStream('TestData');
logger.info(
' $nodeId: All nodes have paused busy-wait polling - system resources freed',
);
break;
case 'resume_for_phase_2':
logger.info('βΆοΈ $nodeId: Resuming data streams for Phase 2...');
await session.resumeStream('TestData', flushBeforeResume: true);
logger.info(' $nodeId: All nodes resumed with fresh buffers');
break;
case 'start_test_phase_2':
logger.info('π― $nodeId: Starting test phase 2...');
logger.info(' $nodeId: Current message count: $messageCount');
await session.sendUserMessage(
'start_test_phase',
'Starting coordinated test - Phase 2',
{
'phase': 2,
'intensity': 'high',
'start_at': DateTime.now()
.add(Duration(seconds: 1))
.toIso8601String(),
},
);
break;
case 'pause_before_stop':
logger.info('βΈοΈ $nodeId: Pausing before final data collection...');
logger.info(' $nodeId: Current message count: $messageCount');
await session.pauseStream('TestData');
break;
case 'flush_and_resume':
logger.info(
'πΏ $nodeId: Flushing and resuming for final collection...',
);
await session.flushStream('TestData');
await session.resumeStream('TestData', flushBeforeResume: false);
logger.info(
' $nodeId: Streams flushed manually and resumed without auto-flush',
);
break;
case 'stop_data_collection':
logger.info('π $nodeId: Stopping data collection...');
logger.info(' $nodeId: Current message count: $messageCount');
await session.stopStream('TestData');
break;
case 'end_test':
logger.info('π $nodeId: Ending test...');
await session.sendUserMessage('end_test', 'Test sequence completed', {
'total_duration': elapsedTime,
'final_node_count': session.connectedNodes.length,
});
logger.info(' $nodeId: FINAL message count: $messageCount');
inboxSubscription?.cancel();
break;
}
} catch (e, stack) {
logger.info('β $nodeId: Error in step $action: $e');
logger.info('π $nodeId: Stack trace: $stack');
}
}
logger.info('β
$nodeId: Coordinator test sequence completed');
}
/// Participant test logic - responds to coordinator commands
Future<void> _runParticipantTestLogic(
PeerSession session,
String nodeId,
int testDuration,
DataStreamConfig streamConfig,
) async {
logger.info('π€ $nodeId: Running PARTICIPANT test logic');
int currentPhase = 0;
String intensity = 'low';
Timer? dataTimer;
final nodeIdHash = nodeId.hashCode.abs() % 1000; // Unique ID for this node
int messageReceivedCount = 0;
int messageSentCount = 0;
StreamSubscription? inboxSubscription;
// Wait for the test duration or until test completes
bool testCompleted = false;
final testCompleter = Completer<void>();
// Listen for test commands
final ucSub = session.events.userCoordinationMessages.listen((event) async {
switch (event.messageType) {
case 'start_test_phase':
final phase = event.payload['phase'] as int;
intensity = event.payload['intensity'] as String;
final startAtStr = event.payload['start_at'] as String;
final startAt = DateTime.parse(startAtStr);
logger.info(
'π― $nodeId: Received phase command: Phase $phase ($intensity)',
);
logger.info(' Starting at: $startAt');
currentPhase = phase;
// Wait until the specified start time for synchronization
final delay = startAt.difference(DateTime.now());
if (delay.isNegative) {
logger.info('β‘ $nodeId: Starting immediately (past scheduled time)');
} else {
logger.info(
'β±οΈ $nodeId: Waiting ${delay.inMilliseconds}ms for synchronized start',
);
await Future.delayed(delay);
}
logger.info('βΆοΈ $nodeId: Phase $phase ($intensity) started!');
/// show current message count
logger.info(' $nodeId: Current message count: $messageReceivedCount');
break;
case 'end_test':
if (!testCompleted) {
testCompleted = true;
testCompleter.complete();
}
logger.info('π $nodeId: Test ended by coordinator');
logger.info(
' $nodeId: FINAL received message count '
'(expected to be 0 when '
'StreamParticipationMode.sendParticipantsReceiveCoordinator'
'): $messageReceivedCount',
);
logger.info(' $nodeId: FINAL sent message count: $messageSentCount');
dataTimer?.cancel();
inboxSubscription?.cancel();
final duration = event.payload['total_duration'];
logger.info(' Total test duration: ${duration}s');
break;
default:
logger.warning(
'π¬ $nodeId: Received unknown command: ${event.messageType} (${event.messageId}) - ${event.description}',
);
if (event.payload.isNotEmpty) {
logger.info(' Payload: ${event.payload}');
}
break;
}
});
// Listen for stream commands and generate data accordingly
final ssSub = session.events.streamStart.listen((event) async {
final testStream = await session.getDataStream(event.streamName);
inboxSubscription?.cancel(); // Cancel any previous subscription
inboxSubscription = testStream.inbox.listen((data) {
// logger.info('π₯ $nodeId: Received data: $data');
messageReceivedCount++;
});
logger.info('π $nodeId: Data stream started: ${event.streamName}');
dataTimer = _startDataGeneration(
nodeId,
nodeIdHash,
testStream,
() => currentPhase,
() => intensity,
() => messageSentCount++,
);
});
final srSub = session.events.streamReady.listen((event) async {
logger.info('β
$nodeId: Stream ready acknowledged: ${event.streamName}');
});
final stSub = session.events.streamStop.listen((event) {
logger.info('π $nodeId: Data stream stopped: ${event.streamName}');
dataTimer?.cancel();
});
// Listen for pause/resume commands to show coordination working
final spSub = session.events.streamPause.listen((event) {
logger.info(
'βΈοΈ $nodeId: PARTICIPANT received pause command for ${event.streamName}',
);
logger.info(
' $nodeId: Busy-wait polling paused - freeing system resources',
);
});
final rsSub = session.events.streamResume.listen((event) {
logger.info(
'βΆοΈ $nodeId: PARTICIPANT received resume command for ${event.streamName}',
);
logger.info(' $nodeId: Flush before resume: ${event.flushBeforeResume}');
logger.info(' $nodeId: Busy-wait polling resumed');
});
final sfSub = session.events.streamFlush.listen((event) {
logger.info(
'πΏ $nodeId: PARTICIPANT received flush command for ${event.streamName}',
);
logger.info(' $nodeId: Stream buffers cleared');
});
final sdSub = session.events.streamDestroy.listen((event) {
logger.info(
'π₯ $nodeId: PARTICIPANT received destroy command for ${event.streamName}',
);
logger.info(' $nodeId: Stream resources completely destroyed');
});
logger.info(
'β³ $nodeId: Participant ready, waiting for coordinator commands...',
);
// Wait for test completion or timeout
try {
await testCompleter.future.timeout(Duration(seconds: testDuration + 120));
} on TimeoutException {
logger.info('β° $nodeId: Participant test timeout, completing...');
}
// cancel subscriptions
await ucSub.cancel();
await ssSub.cancel();
await srSub.cancel();
await stSub.cancel();
await spSub.cancel();
await rsSub.cancel();
await sfSub.cancel();
await sdSub.cancel();
}
/// Generate test data for participants
Timer _startDataGeneration(
String nodeId,
int nodeIdHash,
DataStream testStream,
int Function() getCurrentPhase,
String Function() getIntensity,
void Function() onSampleSent,
) {
final random = Random();
var sampleCount = 0;
final timer = Timer.periodic(Duration(milliseconds: 100), (timer) {
// 10 Hz
if (!testStream.started) {
logger.info('π $nodeId: Data stream stopped, ending data generation');
timer.cancel();
return;
}
final now = DateTime.now().millisecondsSinceEpoch.toDouble();
final phase = getCurrentPhase();
final intensity = getIntensity();
// Generate data based on current phase and intensity
var dataValue = random.nextDouble() * 100;
if (intensity == 'high') {
dataValue *= 2; // Higher amplitude for high intensity
}
// Add phase-specific patterns
if (phase == 2) {
dataValue += sin(sampleCount * 0.1) * 20; // Add sine wave in phase 2
}
final data = [
now, // timestamp
nodeIdHash.toDouble(), // node identifier
dataValue, // sample value
];
testStream.sendData(data);
onSampleSent();
sampleCount++;
// logger.info status every 2 seconds
if (sampleCount % 20 == 0) {
logger.info(
'π $nodeId: Generated $sampleCount samples (Phase: $phase, Intensity: $intensity)',
);
}
});
logger.info('π΅ $nodeId: Data generation started (10 Hz)');
return timer;
}