art_adk 1.0.3
art_adk: ^1.0.3 copied to clipboard
Flutter SDK for ART realtime communication with WebSocket channels, AI Agents, AI Orchestrators, presence tracking, end-to-end encrypted messaging, and CRDT-backed shared objects.
example/lib/main.dart
import 'dart:async';
import 'dart:convert';
import 'package:art_adk/art_adk.dart';
import 'package:flutter/material.dart';
import 'package:flutter/services.dart' show rootBundle;
import 'package:http/http.dart' as http;
/// ─────────────────────────────────────────────────────────────────────────────
/// 0. Configure these three values for your environment
/// ─────────────────────────────────────────────────────────────────────────────
const String kServerUri = 'YOUR_WEBSOCKET_URI';
const String kPasscodeEndpoint = 'PASSCODE_ENDPOINT';
const String kUsername = 'USER_NAME';
const String kChannel = 'YOUR_CHANNEL_NAME';
void main() => runApp(const AdkExampleApp());
class AdkExampleApp extends StatelessWidget {
const AdkExampleApp({super.key});
@override
Widget build(BuildContext context) {
return MaterialApp(
title: 'ART ADK Example',
theme: ThemeData(colorSchemeSeed: Colors.indigo, useMaterial3: true),
home: const AdkHomePage(),
);
}
}
/// ─────────────────────────────────────────────────────────────────────────────
/// Home page — walks through every major ADK feature via buttons.
/// ─────────────────────────────────────────────────────────────────────────────
class AdkHomePage extends StatefulWidget {
const AdkHomePage({super.key});
@override
State<AdkHomePage> createState() => _AdkHomePageState();
}
class _AdkHomePageState extends State<AdkHomePage> {
Adk? _adk;
BaseSubscription? _subscription;
final List<String> _log = <String>[];
String? _selectedUser;
//Agent
Agent? _agent;
AgentThread? _thread;
Run? _run;
//Orchestrator
Orchestrator? _orchestrator;
OrchestratorThread? _orchestratorThread;
LiveObjSubscription? get _liveObj {
final sub = _subscription;
return sub is LiveObjSubscription ? sub : null;
}
void _logEvent(String message) {
debugPrint(message);
if (mounted) {
setState(() {
_log.add(
'[${DateTime.now().toIso8601String().substring(11, 19)}] $message');
});
}
}
/// ───────────────────────────────────────────────────────────────────────────
/// 1. Connect
/// ───────────────────────────────────────────────────────────────────────────
Future<void> _connect() async {
try {
_logEvent('connecting...');
final credentials = await _loadCredentials();
final passcode = await _fetchPasscode(credentials);
final updatedCredentials = credentials.copyWith(accessToken: passcode);
final adk = Adk(
adkConfig: AdkConfig(
uri: kServerUri,
authToken: passcode,
getCredentials: () => updatedCredentials,
),
);
adk.on('connection', (dynamic data) {
if (data is ConnectionDetail) {
_logEvent('connected · ${data.connectionId}');
} else {
_logEvent('connected · $data');
}
});
adk.on('close', (dynamic reason) => _logEvent('closed · $reason'));
await adk.connect();
setState(() => _adk = adk);
} catch (e) {
_logEvent('connect failed · $e');
}
}
/// ───────────────────────────────────────────────────────────────────────────
/// 2. Subscribe to a channel (default, secure, or shared object)
/// ───────────────────────────────────────────────────────────────────────────
Future<void> _subscribe() async {
final adk = _adk;
if (adk == null) {
_logEvent('connect first');
return;
}
try {
final sub = await adk.subscribe(channel: kChannel);
setState(() => _subscription = sub);
_logEvent('subscribed to $kChannel (${sub.channelConfig.channelType})');
// Bind a named event across any channel type.
sub.emitter.on('message', (dynamic data) {
_logEvent('message · $data');
});
// On default channels, also stream every event with listen().
if (sub is Subscription) {
sub.listen((Map<String, dynamic> data) {
_logEvent("event=${data['event']}");
});
}
// Track presence (must be enabled on the channel in the dashboard).
unawaited(sub.fetchPresence(
callback: (users) {
final normalized =
users.map((e) => e.split(':').first).toSet().toList();
if (!mounted) return;
setState(() {
_selectedUser ??= normalized.firstWhere(
(u) => u != kUsername,
orElse: () => '',
);
});
_logEvent('presence · $normalized');
},
)
// .catchError((e) {
// _logEvent('Presence ACK timeout: $e');
// }),
);
// On CRDT channels, observe the document tree.
if (sub is LiveObjSubscription) {
await sub.query(path: 'document').listen((dynamic data) {
_logEvent('doc · $data');
});
}
} catch (e) {
_logEvent('subscribe failed · $e');
}
}
/// ───────────────────────────────────────────────────────────────────────────
/// 3. Push a message (optionally targeted)
/// ───────────────────────────────────────────────────────────────────────────
Future<void> _sendMessage() async {
final sub = _subscription;
if (sub == null) {
_logEvent('Subscribe first');
return;
}
if (_selectedUser == null || _selectedUser!.isEmpty) {
_logEvent('No recipient available');
return;
}
final messageId = '$kUsername-${DateTime.now().millisecondsSinceEpoch}';
try {
final refId = await sub.push(
event: 'message',
data: {
'message': 'Hello from Flutter ADK',
'from': kUsername,
'id': messageId,
},
options: PushConfig(
to: ["USER_NAME"],
),
);
_logEvent('Message sent');
_logEvent('Reference Id: $refId');
} catch (e) {
_logEvent('Send failed: $e');
}
}
/// ───────────────────────────────────────────────────────────────────────────
/// 4. Encryption — generate a keypair once per session
/// ───────────────────────────────────────────────────────────────────────────
Future<void> _generateKeyPair() async {
final adk = _adk;
if (adk == null) {
_logEvent('connect first');
return;
}
try {
final pair = await adk.generateKeyPair();
_logEvent('keypair ready · pub=${pair.publicKey.substring(0, 10)}…');
} catch (e) {
_logEvent('keygen failed · $e');
}
}
/// ───────────────────────────────────────────────────────────────────────────
/// 5. CRDT — set a document title
/// ───────────────────────────────────────────────────────────────────────────
Future<void> _setDocTitle() async {
final live = _liveObj;
if (live == null) {
_logEvent('subscribe to a shared-object channel first');
return;
}
live.state()['document']['title'].set(
'Title @ ${DateTime.now().millisecondsSinceEpoch}',
);
await live.flush();
_logEvent('title written');
}
/// ───────────────────────────────────────────────────────────────────────────
/// 6. CRDT — array push / pop
/// ───────────────────────────────────────────────────────────────────────────
Future<void> _pushToArray() async {
final live = _liveObj;
if (live == null) return;
final length = live.state()['items'].push('item ${DateTime.now()}');
await live.flush();
_logEvent('items · push (len=$length)');
}
Future<void> _popFromArray() async {
final live = _liveObj;
if (live == null) return;
final removed = live.state()['items'].pop();
await live.flush();
_logEvent('items · pop ($removed)');
}
/// ───────────────────────────────────────────────────────────────────────────
/// 7. Interceptor — log every message that passes through
/// ───────────────────────────────────────────────────────────────────────────
Future<void> _addInterceptor() async {
final adk = _adk;
if (adk == null) return;
try {
await adk.intercept(
interceptor: 'demo-logger',
fn: (
Map<String, dynamic> payload,
void Function(dynamic data) resolve,
void Function(String error) reject,
) {
_logEvent('intercepted · ${payload['event']}');
resolve(payload);
},
);
_logEvent('interceptor installed');
} catch (e) {
_logEvent('interceptor failed · $e');
}
}
/// ───────────────────────────────────────────────────────────────────────────
/// 8. Teardown
/// ───────────────────────────────────────────────────────────────────────────
Future<void> _disconnect() async {
try {
await _subscription?.unsubscribe();
await _adk?.disconnect();
} finally {
if (mounted) {
setState(() {
_subscription = null;
_adk = null;
});
}
_logEvent('disconnected');
}
}
/// ───────────────────────────────────────────────────────────────────────────
/// Helpers
/// ───────────────────────────────────────────────────────────────────────────
Future<CredentialStore> _loadCredentials() async {
try {
final raw = await rootBundle.loadString('assets/adk-services.json');
final json = jsonDecode(raw) as Map<String, dynamic>;
return CredentialStore(
environment: json['Environment'] as String? ?? '',
projectKey: json['ProjectKey'] as String? ?? '',
orgTitle: json['Org-Title'] as String? ?? '',
clientID: json['Client-ID'] as String? ?? '',
clientSecret: json['Client-Secret'] as String? ?? '',
);
} catch (e) {
throw Exception('Failed to load assets/adk-services.json: $e');
}
}
Future<String> _fetchPasscode(CredentialStore creds) async {
final response = await http.post(
Uri.parse(kPasscodeEndpoint),
headers: <String, String>{
'Client-Id': creds.clientID,
'Client-Secret': creds.clientSecret,
'X-Org': creds.orgTitle,
'Environment': creds.environment,
'ProjectKey': creds.projectKey,
'Content-Type': 'application/json',
},
body: jsonEncode(<String, dynamic>{
'username': kUsername,
'first_name': 'Alice',
'last_name': 'Example',
}),
);
if (response.statusCode < 200 || response.statusCode >= 300) {
throw Exception('passcode request failed (${response.statusCode})');
}
final decoded = jsonDecode(response.body) as Map<String, dynamic>;
final data = decoded['data'];
final passcode = data is Map<String, dynamic>
? data['passcode'] as String?
: decoded['passcode'] as String?;
if (passcode == null || passcode.isEmpty) {
throw Exception('passcode missing in response');
}
return passcode;
}
/// ───────────────────────────────────────────────────────────────────────────
/// Agent
/// ───────────────────────────────────────────────────────────────────────────
Future<void> _createAgent() async {
final adk = _adk;
if (adk == null) {
_logEvent('connect first');
return;
}
try {
_agent = adk.agent('AGENT_ID');
_logEvent('Agent created');
} catch (e) {
_logEvent('create agent failed · $e');
}
}
Future<void> _startThread() async {
final agent = _agent;
if (agent == null) {
_logEvent('create agent first');
return;
}
try {
final thread = agent.thread();
_thread = thread;
_logEvent('Thread started');
_logEvent('Thread Id : ${thread.threadId}');
} catch (e) {
_logEvent('thread failed · $e');
}
}
Future<void> _listenAgent() async {
final thread = _thread;
if (thread == null) {
_logEvent('start thread first');
return;
}
await thread.listen((AgentEventEnvelope envelope) {
switch (envelope.event) {
case 'agent_general_response':
final output = envelope.content as AgentOutput;
_logEvent('Agent : ${output.message}');
break;
case 'agent_error_response':
final error = envelope.content as AgentError;
_logEvent('Error : ${error.message}');
break;
case 'human_input_request':
final request = envelope.content as HumanInputRequest;
_logEvent('Human Input : ${request.prompt}');
break;
case 'agent_wait_response':
final wait = envelope.content as AgentWait;
_logEvent('Waiting : ${wait.waitingForAgentId}');
break;
default:
_logEvent(envelope.event);
}
});
await thread.listenTrace((frame) {
_logEvent('TRACE : $frame');
});
thread.feedbackRequest(
(HumanInputRequest request, Run run) async {
_logEvent('Question : ${request.prompt}');
// Demo response
await run.sendFeedback(
'Budget 50,000 and travelling in December',
);
},
);
_logEvent('Agent listeners registered');
}
Future<void> _sendPrompt() async {
final thread = _thread;
if (thread == null) {
_logEvent('start thread first');
return;
}
try {
_run = await thread.run(
'Plan a 3-day trip to Goa',
);
final output = await _run!.done();
_logEvent('Final Response');
_logEvent(output.message);
} on AgentError catch (e) {
_logEvent('${e.code} : ${e.message}');
} catch (e) {
_logEvent(e.toString());
}
}
Future<void> _resetAgent() async {
_run = null;
_thread = null;
_agent = null;
_logEvent('Agent cleared');
}
/// ───────────────────────────────────────────────────────────────────────────
/// Orchestrator
/// ───────────────────────────────────────────────────────────────────────────
Future<void> _createOrchestrator() async {
final adk = _adk;
if (adk == null) {
_logEvent('connect first');
return;
}
try {
_orchestrator = adk.orchestrator("ORCHESTRATOR_ID");
_logEvent('Orchestrator created');
} catch (e) {
_logEvent('create orchestrator failed · $e');
}
}
Future<void> _startOrchestratorThread() async {
final orchestrator = _orchestrator;
if (orchestrator == null) {
_logEvent('create orchestrator first');
return;
}
try {
final thread = await orchestrator.thread();
_orchestratorThread = thread;
_logEvent('Orchestrator thread started');
_logEvent('Thread Id : ${thread.threadId}');
} catch (e) {
_logEvent('thread failed · $e');
}
}
Future<void> _listenOrchestrator() async {
final thread = _orchestratorThread;
if (thread == null) {
_logEvent('start orchestrator thread first');
return;
}
thread.listen((Map<String, dynamic> data) {
final event = data['event'] as String? ?? '';
final content = data['content'];
if (content is Map && content['reply'] is Function) {
_logEvent('Question : ${content['prompt']}');
final reply = content['reply'] as Function;
reply(<String, dynamic>{
'user_input': 'Budget 50,000 and travelling in December',
});
return;
}
final type = content is Map ? (content['type'] as String? ?? '') : '';
switch (type.isNotEmpty ? type : event) {
case 'agent_general_response':
_logEvent('Answer : ${content['message']}');
break;
case 'agent_error_response':
_logEvent(
'Error : ${content['message']}',
);
break;
case 'human_input_request':
_logEvent(
'Human Input : ${content['prompt']}',
);
break;
case 'agent_wait_response':
_logEvent(
'Waiting : ${content['waiting_for_agent_id']}',
);
break;
case 'planner_correction_request':
_logEvent(
'Planner : ${content['reason']}',
);
break;
default:
_logEvent('[$event] $content');
}
});
thread.listenTrace((frame) {
_logEvent('TRACE : $frame');
});
_logEvent('Orchestrator listeners registered');
}
Future<void> _runWorkflow() async {
final thread = _orchestratorThread;
if (thread == null) {
_logEvent('start orchestrator thread first');
return;
}
try {
await thread.push(
event: 'user_input',
data: <String, dynamic>{
'message': 'Plan a 3-day trip to Goa',
},
);
_logEvent('Workflow started');
} catch (e) {
_logEvent('Workflow failed · $e');
}
}
Future<void> _resetOrchestrator() async {
_orchestratorThread?.dispose();
_orchestratorThread = null;
_orchestrator = null;
_logEvent('Orchestrator cleared');
}
/// ───────────────────────────────────────────────────────────────────────────
/// UI
/// ───────────────────────────────────────────────────────────────────────────
@override
Widget build(BuildContext context) {
return Scaffold(
appBar: AppBar(
title: const Text('ART ADK Example'),
actions: <Widget>[
Padding(
padding: const EdgeInsets.only(right: 12),
child: Center(
child: Text(
_adk?.getState() ?? 'stopped',
style: Theme.of(context).textTheme.labelMedium,
),
),
),
IconButton(
icon: Row(
children: [
Text("Clear Logs"),
const Icon(Icons.delete_outline),
],
),
onPressed: () => setState(() => _log.clear()),
tooltip: 'Clear Logs',
),
],
),
body: SingleChildScrollView(
child: Padding(
padding: const EdgeInsets.all(12),
child: Column(
children: [
header('Connection'),
_btn('Connect', _connect, primary: true),
_btn('Subscribe', _subscribe),
_btn('Generate KeyPair', _generateKeyPair),
_btn('Disconnect', _disconnect, destructive: true),
const Divider(height: 1),
header('Chat'),
_btn('Send Message', _sendMessage),
const Divider(height: 1),
header('CRDT Operations'),
_btn('CRDT · Set Title', _setDocTitle),
_btn('CRDT · Array Push', _pushToArray),
_btn('CRDT · Array Pop', _popFromArray),
_btn('Add Interceptor', _addInterceptor),
const Divider(height: 1),
header('Agent'),
_btn(
'Create Agent',
_createAgent,
),
_btn(
'Start Thread',
_startThread,
),
_btn(
'Register Listeners',
_listenAgent,
),
_btn(
'Run Prompt',
_sendPrompt,
),
_btn('Reset Agent', _resetAgent, primary: true),
const Divider(height: 1),
header('Orchestrator'),
_btn(
'Create Orchestrator',
_createOrchestrator,
),
_btn(
'Start Orchestrator Thread',
_startOrchestratorThread,
),
_btn(
'Register Orchestrator Listeners',
_listenOrchestrator,
),
_btn(
'Run Workflow',
_runWorkflow,
),
_btn('Reset Orchestrator', _resetOrchestrator, primary: true),
const Divider(height: 1),
Container(
color: Colors.black12.withValues(alpha: 0.03),
child: ListView.builder(
shrinkWrap: true,
padding: const EdgeInsets.all(12),
itemCount: _log.length,
itemBuilder: (_, int i) => SingleChildScrollView(
child: Text(
_log[i],
style: const TextStyle(fontFamily: 'Menlo', fontSize: 12),
),
),
),
),
],
),
),
),
);
}
Widget header(String text) =>
Align(alignment: Alignment.centerLeft, child: Text(text));
Widget _btn(
String label,
Future<void> Function() onTap, {
bool primary = false,
bool destructive = false,
}) {
final style = primary
? FilledButton.styleFrom()
: destructive
? FilledButton.styleFrom(
backgroundColor: Colors.red.shade400,
foregroundColor: Colors.white,
)
: null;
return primary || destructive
? FilledButton(style: style, onPressed: onTap, child: Text(label))
: OutlinedButton(onPressed: onTap, child: Text(label));
}
}