- Add oto_server_job_client.dart for direct server dispatch - Update runner.proto with new dispatch contract - Update agent_runner and edge_registration_client for control plane separation - Add Go protobuf files for runner service - Update HTTP server with new dispatch endpoints - Add registration and smoke tests - Archive completed task documents
579 lines
19 KiB
Dart
579 lines
19 KiB
Dart
import 'dart:async';
|
|
import 'dart:convert';
|
|
import 'dart:io';
|
|
|
|
import 'package:test/test.dart';
|
|
import 'package:oto/oto/agent/agent_config.dart';
|
|
import 'package:oto/oto/agent/agent_runner.dart';
|
|
import 'package:oto/oto/agent/edge_registration_client.dart';
|
|
import 'package:oto/oto/agent/oto_server_job_client.dart';
|
|
import 'package:oto/oto/core/build_result.dart';
|
|
import 'package:oto/oto/core/execution_context.dart';
|
|
|
|
class FakeEdgeAgentSession implements EdgeAgentSession {
|
|
@override
|
|
final RegistrationResult result;
|
|
bool closed = false;
|
|
|
|
FakeEdgeAgentSession(this.result);
|
|
|
|
@override
|
|
Future<void> close() async {
|
|
closed = true;
|
|
}
|
|
}
|
|
|
|
class FakeEdgeRegistrationClient extends RegistrationClient {
|
|
final RegistrationResult Function(AgentConfig config) _onRegister;
|
|
FakeEdgeAgentSession? lastSession;
|
|
|
|
FakeEdgeRegistrationClient(this._onRegister);
|
|
|
|
@override
|
|
Future<EdgeAgentSession> openSession(AgentConfig config) async {
|
|
final session = FakeEdgeAgentSession(_onRegister(config));
|
|
lastSession = session;
|
|
return session;
|
|
}
|
|
}
|
|
|
|
void main() {
|
|
const validBootstrapYaml = '''
|
|
agent:
|
|
id: "agent-123"
|
|
alias: "my-agent"
|
|
enrollment_token: "token-456"
|
|
edge:
|
|
url: "127.0.0.1:8080"
|
|
runtime:
|
|
install_dir: "/usr/bin"
|
|
workspace_root: "/var/oto"
|
|
log_dir: "/var/log/oto"
|
|
''';
|
|
|
|
late AgentConfig validConfig;
|
|
|
|
setUp(() {
|
|
validConfig = AgentConfig.fromYamlContent(validBootstrapYaml);
|
|
});
|
|
|
|
group('EdgeEndpoint', () {
|
|
test('EdgeEndpoint parses host and port from config edge url', () {
|
|
final endpoint1 = EdgeEndpoint.parse('127.0.0.1:8080');
|
|
expect(endpoint1.host, '127.0.0.1');
|
|
expect(endpoint1.port, 8080);
|
|
|
|
final endpoint2 = EdgeEndpoint.parse('http://edge-server:9000');
|
|
expect(endpoint2.host, 'edge-server');
|
|
expect(endpoint2.port, 9000);
|
|
|
|
final endpoint3 = EdgeEndpoint.parse('https://secure-edge');
|
|
expect(endpoint3.host, 'secure-edge');
|
|
expect(endpoint3.port, 443);
|
|
|
|
final endpoint4 = EdgeEndpoint.parse('edge-no-port');
|
|
expect(endpoint4.host, 'edge-no-port');
|
|
expect(endpoint4.port, 80);
|
|
});
|
|
});
|
|
|
|
group('EdgeRegistrationClient.register one-shot', () {
|
|
test('register returns session result and closes the session', () async {
|
|
final fakeClient = FakeEdgeRegistrationClient((config) {
|
|
return RegistrationResult.accepted('node-123', 'alias-123', {
|
|
'concurrency': 1,
|
|
});
|
|
});
|
|
|
|
final result = await fakeClient.register(validConfig);
|
|
|
|
expect(result.accepted, isTrue);
|
|
expect(result.nodeId, 'node-123');
|
|
expect(fakeClient.lastSession, isNotNull);
|
|
expect(fakeClient.lastSession!.closed, isTrue);
|
|
});
|
|
});
|
|
|
|
group('OtoServerRegistrationClient', () {
|
|
test('builds OTO Server registration request from config', () {
|
|
final client = OtoServerRegistrationClient(
|
|
commandTypes: ['Shell', 'Git'],
|
|
);
|
|
|
|
final request = client.buildRegisterRunnerRequest(validConfig);
|
|
|
|
expect(request.enrollmentToken, 'token-456');
|
|
expect(request.runnerId, 'agent-123');
|
|
expect(request.alias, 'my-agent');
|
|
expect(request.protocolVersion, otoRunnerProtocolVersion);
|
|
expect(request.capability.name, otoRunnerCapabilityName);
|
|
expect(request.commandCatalog.commandTypes, ['Shell', 'Git']);
|
|
});
|
|
|
|
test('posts registration request to OTO Server endpoint', () async {
|
|
final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0);
|
|
final captured = Completer<Map<String, dynamic>>();
|
|
final serverSub = server.listen((request) async {
|
|
expect(request.method, 'POST');
|
|
if (request.uri.path == '/api/v1/runners/register') {
|
|
final body = await utf8.decoder.bind(request).join();
|
|
captured.complete(jsonDecode(body) as Map<String, dynamic>);
|
|
request.response
|
|
..headers.contentType = ContentType.json
|
|
..write(
|
|
jsonEncode({
|
|
'accepted': true,
|
|
'runner_id': 'agent-123',
|
|
'alias': 'my-agent',
|
|
}),
|
|
);
|
|
} else {
|
|
request.response
|
|
..headers.contentType = ContentType.json
|
|
..write(jsonEncode({'success': true}));
|
|
}
|
|
await request.response.close();
|
|
});
|
|
|
|
try {
|
|
final config = AgentConfig(
|
|
agent: validConfig.agent,
|
|
server: ServerConnectionConfig(
|
|
url: 'http://${server.address.host}:${server.port}',
|
|
),
|
|
runtime: validConfig.runtime,
|
|
);
|
|
final client = OtoServerRegistrationClient(
|
|
commandTypes: ['Shell', 'Git'],
|
|
);
|
|
|
|
final result = await client.register(config);
|
|
|
|
expect(result.accepted, isTrue);
|
|
expect(result.runnerId, 'agent-123');
|
|
final payload = await captured.future;
|
|
expect(payload['enrollment_token'], 'token-456');
|
|
expect(payload['runner_id'], 'agent-123');
|
|
expect(payload['alias'], 'my-agent');
|
|
expect(payload['protocol_version'], otoRunnerProtocolVersion);
|
|
expect(payload['capability'], containsPair('name', 'oto-runner'));
|
|
expect(
|
|
payload['command_catalog'],
|
|
containsPair('command_types', ['Shell', 'Git']),
|
|
);
|
|
} finally {
|
|
await serverSub.cancel();
|
|
await server.close(force: true);
|
|
}
|
|
});
|
|
|
|
test('sends first heartbeat and disconnects on close', () async {
|
|
final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0);
|
|
final registerReceived = Completer<void>();
|
|
final heartbeatReceived = Completer<void>();
|
|
final disconnectReceived = Completer<void>();
|
|
|
|
final serverSub = server.listen((request) async {
|
|
if (request.method == 'POST' &&
|
|
request.uri.path == '/api/v1/runners/register') {
|
|
registerReceived.complete();
|
|
request.response
|
|
..headers.contentType = ContentType.json
|
|
..write(
|
|
jsonEncode({
|
|
'accepted': true,
|
|
'runner_id': 'agent-123',
|
|
'alias': 'my-agent',
|
|
}),
|
|
);
|
|
} else if (request.method == 'POST' &&
|
|
request.uri.path == '/api/v1/runners/agent-123/heartbeat') {
|
|
heartbeatReceived.complete();
|
|
request.response
|
|
..headers.contentType = ContentType.json
|
|
..write(jsonEncode({'success': true}));
|
|
} else if (request.method == 'POST' &&
|
|
request.uri.path == '/api/v1/runners/agent-123/disconnect') {
|
|
disconnectReceived.complete();
|
|
request.response
|
|
..headers.contentType = ContentType.json
|
|
..write(jsonEncode({'success': true}));
|
|
}
|
|
await request.response.close();
|
|
});
|
|
|
|
try {
|
|
final config = AgentConfig(
|
|
agent: validConfig.agent,
|
|
server: ServerConnectionConfig(
|
|
url: 'http://${server.address.host}:${server.port}',
|
|
),
|
|
runtime: validConfig.runtime,
|
|
);
|
|
final client = OtoServerRegistrationClient(
|
|
commandTypes: ['Shell', 'Git'],
|
|
heartbeatInterval: const Duration(milliseconds: 100),
|
|
);
|
|
|
|
final session = await client.openSession(config);
|
|
|
|
expect(session, isA<OtoServerJobSession>());
|
|
|
|
// First heartbeat should be sent automatically
|
|
await heartbeatReceived.future.timeout(const Duration(seconds: 2));
|
|
|
|
// Close session
|
|
await session.close();
|
|
|
|
// Disconnect should be sent
|
|
await disconnectReceived.future.timeout(const Duration(seconds: 2));
|
|
} finally {
|
|
await serverSub.cancel();
|
|
await server.close(force: true);
|
|
}
|
|
});
|
|
|
|
test(
|
|
'does not send heartbeat or disconnect on rejected registration',
|
|
() async {
|
|
final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0);
|
|
var heartbeatReceived = false;
|
|
var disconnectReceived = false;
|
|
|
|
final serverSub = server.listen((request) async {
|
|
if (request.method == 'POST' &&
|
|
request.uri.path == '/api/v1/runners/register') {
|
|
request.response
|
|
..headers.contentType = ContentType.json
|
|
..write(
|
|
jsonEncode({
|
|
'accepted': false,
|
|
'reject_reason': 'Invalid token',
|
|
}),
|
|
);
|
|
} else if (request.uri.path.contains('heartbeat')) {
|
|
heartbeatReceived = true;
|
|
request.response
|
|
..headers.contentType = ContentType.json
|
|
..write(jsonEncode({'success': true}));
|
|
} else if (request.uri.path.contains('disconnect')) {
|
|
disconnectReceived = true;
|
|
request.response
|
|
..headers.contentType = ContentType.json
|
|
..write(jsonEncode({'success': true}));
|
|
}
|
|
await request.response.close();
|
|
});
|
|
|
|
try {
|
|
final config = AgentConfig(
|
|
agent: validConfig.agent,
|
|
server: ServerConnectionConfig(
|
|
url: 'http://${server.address.host}:${server.port}',
|
|
),
|
|
runtime: validConfig.runtime,
|
|
);
|
|
final client = OtoServerRegistrationClient(
|
|
commandTypes: ['Shell', 'Git'],
|
|
heartbeatInterval: const Duration(milliseconds: 50),
|
|
);
|
|
|
|
final session = await client.openSession(config);
|
|
expect(session.result.accepted, isFalse);
|
|
|
|
// Wait a short duration to ensure no heartbeat is sent
|
|
await Future<void>.delayed(const Duration(milliseconds: 200));
|
|
expect(heartbeatReceived, isFalse);
|
|
|
|
await session.close();
|
|
// Wait a short duration to ensure no disconnect is sent
|
|
await Future<void>.delayed(const Duration(milliseconds: 100));
|
|
expect(disconnectReceived, isFalse);
|
|
} finally {
|
|
await serverSub.cancel();
|
|
await server.close(force: true);
|
|
}
|
|
},
|
|
);
|
|
|
|
test('session close cancels future heartbeats', () async {
|
|
final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0);
|
|
var heartbeatCount = 0;
|
|
|
|
final serverSub = server.listen((request) async {
|
|
if (request.method == 'POST' &&
|
|
request.uri.path == '/api/v1/runners/register') {
|
|
request.response
|
|
..headers.contentType = ContentType.json
|
|
..write(jsonEncode({'accepted': true, 'runner_id': 'agent-123'}));
|
|
} else if (request.uri.path.contains('heartbeat')) {
|
|
heartbeatCount++;
|
|
request.response
|
|
..headers.contentType = ContentType.json
|
|
..write(jsonEncode({'success': true}));
|
|
} else if (request.uri.path.contains('disconnect')) {
|
|
request.response
|
|
..headers.contentType = ContentType.json
|
|
..write(jsonEncode({'success': true}));
|
|
}
|
|
await request.response.close();
|
|
});
|
|
|
|
try {
|
|
final config = AgentConfig(
|
|
agent: validConfig.agent,
|
|
server: ServerConnectionConfig(
|
|
url: 'http://${server.address.host}:${server.port}',
|
|
),
|
|
runtime: validConfig.runtime,
|
|
);
|
|
final client = OtoServerRegistrationClient(
|
|
commandTypes: ['Shell', 'Git'],
|
|
heartbeatInterval: const Duration(milliseconds: 50),
|
|
);
|
|
|
|
final session = await client.openSession(config);
|
|
|
|
// Wait for first heartbeat(s)
|
|
await Future<void>.delayed(const Duration(milliseconds: 150));
|
|
final countBeforeClose = heartbeatCount;
|
|
expect(countBeforeClose, greaterThan(0));
|
|
|
|
await session.close();
|
|
|
|
// Wait longer to see if any more heartbeats occur
|
|
await Future<void>.delayed(const Duration(milliseconds: 150));
|
|
expect(heartbeatCount, equals(countBeforeClose));
|
|
} finally {
|
|
await serverSub.cancel();
|
|
await server.close(force: true);
|
|
}
|
|
});
|
|
|
|
test('job client claims jobs and reports build results', () async {
|
|
final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0);
|
|
final claimReceived = Completer<Map<String, dynamic>>();
|
|
final reportReceived = Completer<Map<String, dynamic>>();
|
|
final logReceived = Completer<Map<String, dynamic>>();
|
|
final artifactReceived = Completer<Map<String, dynamic>>();
|
|
|
|
final serverSub = server.listen((request) async {
|
|
final bodyText = await utf8.decoder.bind(request).join();
|
|
final body = bodyText.isEmpty
|
|
? <String, dynamic>{}
|
|
: jsonDecode(bodyText) as Map<String, dynamic>;
|
|
request.response.headers.contentType = ContentType.json;
|
|
if (request.method == 'POST' &&
|
|
request.uri.path == '/api/v1/runners/agent-123/jobs/claim') {
|
|
claimReceived.complete(body);
|
|
request.response.write(
|
|
jsonEncode({
|
|
'accepted': true,
|
|
'runner_id': 'agent-123',
|
|
'job_id': 'job-123',
|
|
'execution_id': 'exec-123',
|
|
'state': 'running',
|
|
}),
|
|
);
|
|
} else if (request.method == 'POST' &&
|
|
request.uri.path ==
|
|
'/api/v1/runners/agent-123/executions/exec-123/report') {
|
|
reportReceived.complete(body);
|
|
request.response.write(
|
|
jsonEncode({
|
|
'accepted': true,
|
|
'runner_id': 'agent-123',
|
|
'job_id': 'job-123',
|
|
'execution_id': 'exec-123',
|
|
'state': 'succeeded',
|
|
}),
|
|
);
|
|
} else if (request.method == 'POST' &&
|
|
request.uri.path ==
|
|
'/api/v1/runners/agent-123/executions/exec-123/logs') {
|
|
logReceived.complete(body);
|
|
request.response.statusCode = HttpStatus.created;
|
|
request.response.write(jsonEncode({'accepted': true}));
|
|
} else if (request.method == 'POST' &&
|
|
request.uri.path ==
|
|
'/api/v1/runners/agent-123/executions/exec-123/artifacts') {
|
|
artifactReceived.complete(body);
|
|
request.response.statusCode = HttpStatus.created;
|
|
request.response.write(jsonEncode({'accepted': true}));
|
|
} else {
|
|
request.response.statusCode = HttpStatus.notFound;
|
|
request.response.write(jsonEncode({'error': 'not found'}));
|
|
}
|
|
await request.response.close();
|
|
});
|
|
|
|
final client = OtoServerJobClient(
|
|
serverUrl: 'http://${server.address.host}:${server.port}',
|
|
runnerId: 'agent-123',
|
|
);
|
|
|
|
try {
|
|
final claim = await client.claimJob(
|
|
jobId: 'job-123',
|
|
executionId: 'exec-123',
|
|
);
|
|
expect(claim.accepted, isTrue);
|
|
expect(claim.state, 'running');
|
|
|
|
final result = BuildResult.success(
|
|
stepEvents: [
|
|
StepEvent(
|
|
stepId: 0,
|
|
workflowIndex: 0,
|
|
stepType: 'command',
|
|
event: 'completed',
|
|
timestamp: '2026-06-05T12:00:00Z',
|
|
commandId: 'build',
|
|
commandType: 'Shell',
|
|
),
|
|
],
|
|
);
|
|
final report = await client.reportExecution(
|
|
jobId: 'job-123',
|
|
executionId: 'exec-123',
|
|
result: result,
|
|
);
|
|
expect(report.accepted, isTrue);
|
|
expect(report.state, 'succeeded');
|
|
|
|
await client.appendLog(
|
|
executionId: 'exec-123',
|
|
line: 'build completed',
|
|
);
|
|
await client.reportArtifact(
|
|
executionId: 'exec-123',
|
|
name: 'binary',
|
|
path: '/dist/app',
|
|
);
|
|
|
|
final claimPayload = await claimReceived.future;
|
|
expect(claimPayload['runner_id'], 'agent-123');
|
|
expect(claimPayload['job_id'], 'job-123');
|
|
expect(claimPayload['execution_id'], 'exec-123');
|
|
|
|
final reportPayload = await reportReceived.future;
|
|
expect(reportPayload['success'], isTrue);
|
|
expect(reportPayload['exit_code'], 0);
|
|
expect(reportPayload['message'], 'Build completed successfully.');
|
|
expect(reportPayload['execution_result'], isA<Map<String, dynamic>>());
|
|
final stepEvents = reportPayload['step_events'] as List<dynamic>;
|
|
expect(stepEvents, hasLength(1));
|
|
expect(stepEvents.first, containsPair('commandId', 'build'));
|
|
|
|
final logPayload = await logReceived.future;
|
|
expect(logPayload['line'], 'build completed');
|
|
|
|
final artifactPayload = await artifactReceived.future;
|
|
expect(artifactPayload['name'], 'binary');
|
|
expect(artifactPayload['path'], '/dist/app');
|
|
} finally {
|
|
client.close();
|
|
await serverSub.cancel();
|
|
await server.close(force: true);
|
|
}
|
|
});
|
|
});
|
|
|
|
group('AgentRunner', () {
|
|
test('AgentRunner maps config token to register request', () async {
|
|
AgentConfig? capturedConfig;
|
|
final fakeClient = FakeEdgeRegistrationClient((config) {
|
|
capturedConfig = config;
|
|
return RegistrationResult.accepted('node-123', 'alias-123', {
|
|
'concurrency': 1,
|
|
});
|
|
});
|
|
|
|
final runner = DefaultAgentRunner(
|
|
client: fakeClient,
|
|
onLog: (msg) {},
|
|
waitForShutdown: () async {},
|
|
);
|
|
await runner.run(validConfig);
|
|
|
|
expect(capturedConfig, isNotNull);
|
|
expect(capturedConfig!.agent.enrollmentToken, 'token-456');
|
|
});
|
|
|
|
test('AgentRunner keeps session open until shutdown then closes', () async {
|
|
final fakeClient = FakeEdgeRegistrationClient((config) {
|
|
return RegistrationResult.accepted('node-123', 'alias-123', {
|
|
'concurrency': 1,
|
|
});
|
|
});
|
|
|
|
final shutdown = Completer<void>();
|
|
final runner = DefaultAgentRunner(
|
|
client: fakeClient,
|
|
onLog: (msg) {},
|
|
waitForShutdown: () => shutdown.future,
|
|
);
|
|
|
|
final runFuture = runner.run(validConfig);
|
|
// 종료 신호 전에는 session이 살아 있어야 한다.
|
|
await Future<void>.delayed(Duration.zero);
|
|
expect(fakeClient.lastSession, isNotNull);
|
|
expect(fakeClient.lastSession!.closed, isFalse);
|
|
|
|
shutdown.complete();
|
|
await runFuture;
|
|
expect(fakeClient.lastSession!.closed, isTrue);
|
|
});
|
|
|
|
test('AgentRunner runs injected job loop after registration', () async {
|
|
final fakeClient = FakeEdgeRegistrationClient((config) {
|
|
return RegistrationResult.accepted('node-123', 'alias-123', {
|
|
'concurrency': 1,
|
|
});
|
|
});
|
|
var jobLoopRan = false;
|
|
|
|
final runner = DefaultAgentRunner(
|
|
client: fakeClient,
|
|
onLog: (msg) {},
|
|
runJobs: (config, session) async {
|
|
jobLoopRan = true;
|
|
expect(config.agent.id, 'agent-123');
|
|
expect(session.result.accepted, isTrue);
|
|
},
|
|
waitForShutdown: () async {},
|
|
);
|
|
|
|
await runner.run(validConfig);
|
|
|
|
expect(jobLoopRan, isTrue);
|
|
});
|
|
|
|
test(
|
|
'AgentRunner reports rejected registration and closes session',
|
|
() async {
|
|
final fakeClient = FakeEdgeRegistrationClient((config) {
|
|
return RegistrationResult.rejected('Invalid token');
|
|
});
|
|
|
|
final runner = DefaultAgentRunner(
|
|
client: fakeClient,
|
|
onLog: (msg) {},
|
|
waitForShutdown: () async {},
|
|
);
|
|
await expectLater(
|
|
() => runner.run(validConfig),
|
|
throwsA(
|
|
isA<RegistrationException>().having(
|
|
(e) => e.message,
|
|
'message',
|
|
'Invalid token',
|
|
),
|
|
),
|
|
);
|
|
expect(fakeClient.lastSession, isNotNull);
|
|
expect(fakeClient.lastSession!.closed, isTrue);
|
|
},
|
|
);
|
|
});
|
|
}
|