art_adk 1.0.3 copy "art_adk: ^1.0.3" to clipboard
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));
  }
}
4
likes
0
points
129
downloads

Documentation

Documentation

Publisher

verified publisherarealtimetech.com

Weekly Downloads

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.

Homepage

Topics

#websocket #realtime #messaging #pubsub #flutter

License

unknown (license)

Dependencies

flutter, http, pinenacl, web_socket_channel

More

Packages that depend on art_adk