iop/apps/client/test/support/client_test_harness.dart

467 lines
14 KiB
Dart

import 'dart:async';
import 'dart:io';
import 'package:iop_client/iop_wire/client_wire_client.dart';
import 'package:iop_client/gen/proto/iop/control.pb.dart';
import 'package:iop_client/control_plane_status_client.dart';
import 'package:protobuf/protobuf.dart';
import 'package:fixnum/fixnum.dart';
import 'package:iop_client/src/integrations/nexo/nexo_notification_client.dart';
class FakeWebSocket implements WebSocket {
final _controller = StreamController<dynamic>();
@override
StreamSubscription<dynamic> listen(
void Function(dynamic event)? onData, {
Function? onError,
void Function()? onDone,
bool? cancelOnError,
}) {
return _controller.stream.listen(
onData,
onError: onError,
onDone: onDone,
cancelOnError: cancelOnError,
);
}
@override
dynamic noSuchMethod(Invocation invocation) {
return null;
}
}
class FakeClientWireClient extends ClientWireClient {
final bool shouldSuccess;
final String mockMessage;
FakeClientWireClient({
this.shouldSuccess = true,
this.mockMessage = 'Welcome to Toki CP!',
}) : super(FakeWebSocket(), heartbeatIntervalTime: 0);
@override
bool get isAlive => true;
@override
Future<Res> sendRequest<
Req extends GeneratedMessage,
Res extends GeneratedMessage
>(Req data, {Duration timeout = const Duration(seconds: 30)}) async {
if (data is ClientHelloRequest) {
final response = ClientHelloResponse()
..ready = shouldSuccess
..protocol = 'iop-wire-v1'
..serverTimeUnixNano = Int64(1716584400)
..message = mockMessage;
return response as Res;
}
throw UnimplementedError();
}
}
class FakeControlPlaneStatusRepository implements ControlPlaneStatusRepository {
String nextCommandStatus = 'accepted';
String nextCommandError = '';
String nextCommandSummary = 'Command accepted by control plane';
String? lastCommandOperation;
String? lastCommandTargetSelector;
Map<String, String>? lastCommandParameters;
String? lastCommandEdgeId;
Completer<EdgeOperationsResponseView>? edgeBOperationsCompleter;
String? emptyOperationsEdgeId;
Exception? fetchOperationsError;
Completer<EdgeCommandResponseView>? commandResponseCompleter;
List<EdgeRegistryView> mockEdges = [
EdgeRegistryView(
edgeId: 'edge-a',
edgeName: 'Edge Alpha',
version: '1.2.0',
capabilities: ['Serving', 'Automation'],
protocol: 'iop-wire-v1',
connected: true,
lastSeen: DateTime.now(),
),
EdgeRegistryView(
edgeId: 'edge-b',
edgeName: 'Edge Beta',
version: '1.2.1',
capabilities: ['Serving'],
protocol: 'iop-wire-v1',
connected: false,
lastSeen: DateTime.now().subtract(const Duration(minutes: 5)),
),
];
@override
Future<List<EdgeRegistryView>> fetchEdges() async {
return mockEdges;
}
@override
Future<EdgeStatusResponseView> fetchEdgeStatus(String edgeId) async {
return EdgeStatusResponseView(
requestId: 'req-123',
edgeId: edgeId,
edgeName: edgeId == 'edge-a' ? 'Edge Alpha' : 'Edge Beta',
observedTimeUnixNano: DateTime.now().microsecondsSinceEpoch * 1000,
nodes: [
EdgeNodeSnapshotView(
nodeId: 'node-1',
alias: 'Node One',
label: 'GPU-T4',
connected: true,
config: NodeConfigSummaryView(
adapters: [
AdapterSummaryView(type: 'ollama', enabled: true),
AdapterSummaryView(type: 'custom', enabled: false),
],
concurrency: 4,
),
providerSnapshots: [
ProviderSnapshotView(
id: 'provider-ollama',
adapter: 'ollama',
type: 'llm',
category: 'inference',
status: 'active',
health: 'healthy',
capacity: 10,
inFlight: 3,
queued: 1,
loadRatio: 0.3,
servedModels: ['llama-3.1', 'mistral'],
lifecycleCapabilities: ['start', 'stop', 'restart'],
),
ProviderSnapshotView(
id: 'provider-vllm',
adapter: 'vllm',
type: 'llm',
category: 'inference',
status: 'active',
health: 'healthy',
capacity: 8,
inFlight: 2,
queued: 0,
loadRatio: 0.25,
servedModels: ['vicuna'],
lifecycleCapabilities: ['start', 'stop'],
),
],
),
EdgeNodeSnapshotView(
nodeId: 'node-2',
alias: 'Node Two',
label: 'CPU-only',
connected: false,
config: NodeConfigSummaryView(
adapters: [AdapterSummaryView(type: 'python-cli', enabled: true)],
concurrency: 2,
),
providerSnapshots: [
ProviderSnapshotView(
id: 'provider-openai',
adapter: 'openai',
type: 'llm',
category: 'external',
status: 'active',
health: 'degraded',
capacity: 100,
inFlight: 50,
queued: 5,
loadRatio: 0.5,
servedModels: ['gpt-4', 'gpt-3.5-turbo'],
lifecycleCapabilities: [],
),
],
),
],
capabilities: [
EdgeCapabilitySummaryView(
kind: 'Serving',
available: true,
status: 'ready',
summary: 'Active',
),
EdgeCapabilitySummaryView(
kind: 'Automation',
available: true,
status: 'ready',
summary: 'Active',
),
],
domainAgents: edgeId == 'edge-a'
? [
EdgeDomainAgentSummaryView(
agentKind: 'deployer',
available: true,
lifecycleState: 'ready',
activeCommandId: 'cmd-alpha-active',
summary: 'Active',
),
EdgeDomainAgentSummaryView(
agentKind: 'build-deploy',
available: true,
lifecycleState: 'ready',
activeCommandId: '',
summary: 'Idle',
),
]
: [
EdgeDomainAgentSummaryView(
agentKind: 'deployer',
available: true,
lifecycleState: 'busy',
activeCommandId: 'cmd-beta-active',
summary: 'Processing',
),
EdgeDomainAgentSummaryView(
agentKind: 'build-deploy',
available: true,
lifecycleState: 'ready',
activeCommandId: '',
summary: 'Idle',
),
],
metadata: {'region': 'us-west'},
error: '',
);
}
@override
Future<List<EdgeNodeEventView>> fetchEdgeEvents(String edgeId) async {
return [
EdgeNodeEventView(
edgeId: edgeId,
eventId: 'evt-1',
type: 'online',
source: 'node-register',
nodeId: 'node-1',
alias: 'Node One',
reason: 'Node connected and registered successfully',
timestamp: DateTime.now().subtract(const Duration(minutes: 2)),
receivedAt: DateTime.now().subtract(const Duration(minutes: 2)),
metadata: {},
),
EdgeNodeEventView(
edgeId: edgeId,
eventId: 'evt-2',
type: 'warn',
source: 'runtime-dispatch',
nodeId: 'node-2',
alias: 'Node Two',
reason: 'Execution concurrency limit reached',
timestamp: DateTime.now().subtract(const Duration(minutes: 1)),
receivedAt: DateTime.now().subtract(const Duration(minutes: 1)),
metadata: {},
),
];
}
@override
Future<FleetStatusResponseView> fetchFleetStatus() async {
return FleetStatusResponseView(
generatedAt: DateTime.now(),
edges: mockEdges
.map(
(e) => FleetEdgeView(
edgeId: e.edgeId,
edgeName: e.edgeName,
version: e.version,
protocol: e.protocol,
connected: e.connected,
health: e.connected ? 'healthy' : 'unhealthy',
lastSeen: e.lastSeen,
nodeCount: 2,
capabilities: e.capabilities
.map(
(c) => EdgeCapabilitySummaryView(
kind: c,
available: true,
status: 'ready',
summary: '',
),
)
.toList(),
domainAgents: e.edgeId == 'edge-a'
? [
EdgeDomainAgentSummaryView(
agentKind: 'deployer',
available: true,
lifecycleState: 'ready',
activeCommandId: 'cmd-alpha-active',
summary: 'Active',
),
EdgeDomainAgentSummaryView(
agentKind: 'build-deploy',
available: true,
lifecycleState: 'ready',
activeCommandId: '',
summary: 'Idle',
),
]
: [
EdgeDomainAgentSummaryView(
agentKind: 'deployer',
available: true,
lifecycleState: 'busy',
activeCommandId: 'cmd-beta-active',
summary: 'Processing',
),
EdgeDomainAgentSummaryView(
agentKind: 'build-deploy',
available: true,
lifecycleState: 'ready',
activeCommandId: '',
summary: 'Idle',
),
],
error: '',
),
)
.toList(),
);
}
@override
Future<EdgeOperationsResponseView> fetchEdgeOperations(String edgeId) async {
if (fetchOperationsError != null) {
throw fetchOperationsError!;
}
if (emptyOperationsEdgeId == edgeId) {
return EdgeOperationsResponseView(
edgeId: edgeId,
operations: [],
);
}
if (edgeId == 'edge-a') {
return EdgeOperationsResponseView(
edgeId: edgeId,
operations: [
EdgeCommandRecordView(
edgeId: edgeId,
commandId: 'cmd-alpha-active',
operation: 'deployer.sync',
targetSelector: '',
status: 'success',
summary: 'Synced 5 models',
error: '',
),
],
);
} else {
if (edgeBOperationsCompleter != null) {
return edgeBOperationsCompleter!.future;
}
return EdgeOperationsResponseView(
edgeId: edgeId,
operations: [
EdgeCommandRecordView(
edgeId: edgeId,
commandId: 'cmd-beta-active',
operation: 'build-deploy.deploy',
targetSelector: '',
status: 'success',
summary: 'Deploy queued for edge beta',
error: '',
),
],
);
}
}
@override
Future<EdgeCommandResponseView> sendEdgeCommand(
String edgeId,
String operation, {
String? targetSelector,
Map<String, String>? parameters,
}) async {
lastCommandOperation = operation;
lastCommandTargetSelector = targetSelector;
lastCommandParameters = parameters;
lastCommandEdgeId = edgeId;
if (operation == 'agent.status' &&
(targetSelector == null || targetSelector.isEmpty)) {
return EdgeCommandResponseView(
requestId: 'req-111',
commandId: 'cmd-failed',
edgeId: edgeId,
status: 'error',
summary: 'Failed: agent.status requires target_selector',
error: 'target_selector is required',
);
}
if (operation == 'agent.command') {
final cmd = parameters?['command'];
if (targetSelector == null ||
targetSelector.isEmpty ||
cmd == null ||
cmd.isEmpty) {
return EdgeCommandResponseView(
requestId: 'req-111',
commandId: 'cmd-failed',
edgeId: edgeId,
status: 'error',
summary:
'Failed: agent.command requires target_selector and parameters.command',
error: 'target_selector and parameters.command are required',
);
}
}
if (commandResponseCompleter != null) {
return commandResponseCompleter!.future;
}
return EdgeCommandResponseView(
requestId: 'req-111',
commandId: 'cmd-new',
edgeId: edgeId,
status: nextCommandStatus,
summary: nextCommandSummary,
error: nextCommandError,
);
}
}
class FakeNexoClientsForUiTest implements NexoNotificationClient {
final StreamController<Map<String, dynamic>> controller;
FakeNexoClientsForUiTest(this.controller);
@override
Stream<Map<String, dynamic>> get onNotification => controller.stream;
@override
Future<void> initialize() async {}
@override
Future<String?> getDeviceToken() async => null;
@override
Future<void> setAuthToken(
String serverUrl,
String token, {
String? identifier,
}) async {}
@override
Future<void> setSigningKey(String serverUrl, String signingKey) async {}
@override
set onDeviceTokenReady(Future<void> Function(String token)? callback) {}
@override
set onNavigateToChannel(
void Function(String serverUrl, String channelId)? callback,
) {}
@override
set onNavigateToThread(
void Function(String serverUrl, String rootId)? callback,
) {}
}