diff --git a/agent-task/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/CODE_REVIEW-cloud-G07.md b/agent-task/archive/2026/06/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/code_review_cloud_G07_0.log similarity index 64% rename from agent-task/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/CODE_REVIEW-cloud-G07.md rename to agent-task/archive/2026/06/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/code_review_cloud_G07_0.log index 6fa74d4..a7919f8 100644 --- a/agent-task/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/CODE_REVIEW-cloud-G07.md +++ b/agent-task/archive/2026/06/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/code_review_cloud_G07_0.log @@ -53,42 +53,46 @@ task=m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke, plan=0, | 항목 | 완료 여부 | |------|---------| -| [API-1] smoke 문서와 Makefile/test grouping 갱신 | [ ] | -| [API-2] runner socket smoke 대표성 보강 | [ ] | -| [API-3] full local smoke 기록 | [ ] | +| [API-1] smoke 문서와 Makefile/test grouping 갱신 | [x] | +| [API-2] runner socket smoke 대표성 보강 | [x] | +| [API-3] full local smoke 기록 | [x] | ## 구현 체크리스트 -- [ ] 선행 `02+01_dispatch_state`, `03+01,02_runner_actions`, `04+01_compat_boundary` PASS complete evidence를 확인한다. -- [ ] [API-1] `agent-test/local/agent-smoke.md`와 Makefile/test grouping을 socket-first 기준으로 갱신한다. 검증: `cd apps/runner && dart analyze` -- [ ] [API-2] runner socket smoke가 lifecycle/dispatch/action/compat boundary evidence를 대표하도록 보강한다. 검증: `cd apps/runner && dart test test/oto_server_connection_smoke_test.dart` -- [ ] [API-3] full local smoke 명령을 실행하고 결과를 기록한다. 검증: Go full test와 runner required test list -- [ ] CODE_REVIEW-*-G??.md의 구현 에이전트 소유 섹션을 실제 구현 내용과 검증 출력으로 채운다. 이 항목이 완료되기 전에는 구현이 완료된 것이 아니다. +- [x] 선행 `02+01_dispatch_state`, `03+01,02_runner_actions`, `04+01_compat_boundary` PASS complete evidence를 확인한다. +- [x] [API-1] `agent-test/local/agent-smoke.md`와 Makefile/test grouping을 socket-first 기준으로 갱신한다. 검증: `cd apps/runner && dart analyze` +- [x] [API-2] runner socket smoke가 lifecycle/dispatch/action/compat boundary evidence를 대표하도록 보강한다. 검증: `cd apps/runner && dart test test/oto_server_connection_smoke_test.dart` +- [x] [API-3] full local smoke 명령을 실행하고 결과를 기록한다. 검증: Go full test와 runner required test list +- [x] CODE_REVIEW-*-G??.md의 구현 에이전트 소유 섹션을 실제 구현 내용과 검증 출력으로 채운다. 이 항목이 완료되기 전에는 구현이 완료된 것이 아니다. ## 코드리뷰 전용 체크리스트 > **[REVIEW AGENT ONLY]** 이 체크리스트는 코드리뷰 에이전트만 사용한다. > 구현 에이전트는 이 섹션을 수정하거나 체크하지 않는다. -- [ ] `코드리뷰 결과`에 `PASS`, `WARN`, `FAIL` 중 하나의 판정을 append한다. -- [ ] 판정과 `차원별 평가`, Required/Suggested/Nit 분류가 서로 일치한다. -- [ ] active `CODE_REVIEW-*-G??.md`를 `code_review_cloud_G07_N.log`로 아카이브한다. -- [ ] active `PLAN-*-G??.md`를 `plan_cloud_G07_M.log`로 아카이브한다. -- [ ] `.gitignore`의 Agent-Ops 관리 block이 `agent-task/**/*.md`와 `agent-task/**/*.log`를 unignore하고 `agent-roadmap/current.md`를 ignore하는지 확인한다. -- [ ] PASS이면 `agent-ops/skills/common/code-review/templates/complete-log-template.md` 기준으로 `complete.log`를 작성하고 active `.md` 파일을 남기지 않는다. -- [ ] PASS이면 active task 디렉터리 `agent-task/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/`를 `agent-task/archive/YYYY/MM/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/`로 이동하고 최종 archive 경로에서 이 체크리스트를 갱신한다. -- [ ] PASS이고 task group이 `m-`이면 런타임이 읽을 완료 이벤트 메타데이터를 보고하고, roadmap 수정이나 `update-roadmap` 직접 호출을 하지 않는다. -- [ ] PASS split 작업이면 이동 후 빈 active parent `agent-task/m-runner-proto-socket-transport-hardening/`를 제거하거나, 남은 sibling/file이 있어 유지했다고 확인한다. +- [x] `코드리뷰 결과`에 `PASS`, `WARN`, `FAIL` 중 하나의 판정을 append한다. +- [x] 판정과 `차원별 평가`, Required/Suggested/Nit 분류가 서로 일치한다. +- [x] active `CODE_REVIEW-*-G??.md`를 `code_review_cloud_G07_N.log`로 아카이브한다. +- [x] active `PLAN-*-G??.md`를 `plan_cloud_G07_M.log`로 아카이브한다. +- [x] `.gitignore`의 Agent-Ops 관리 block이 `agent-task/**/*.md`와 `agent-task/**/*.log`를 unignore하고 `agent-roadmap/current.md`를 ignore하는지 확인한다. +- [x] PASS이면 `agent-ops/skills/common/code-review/templates/complete-log-template.md` 기준으로 `complete.log`를 작성하고 active `.md` 파일을 남기지 않는다. +- [x] PASS이면 active task 디렉터리 `agent-task/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/`를 `agent-task/archive/YYYY/MM/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/`로 이동하고 최종 archive 경로에서 이 체크리스트를 갱신한다. +- [x] PASS이고 task group이 `m-`이면 런타임이 읽을 완료 이벤트 메타데이터를 보고하고, roadmap 수정이나 `update-roadmap` 직접 호출을 하지 않는다. +- [x] PASS split 작업이면 이동 후 빈 active parent `agent-task/m-runner-proto-socket-transport-hardening/`를 제거하거나, 남은 sibling/file이 있어 유지했다고 확인한다. - [ ] WARN/FAIL이고 user-review gate가 트리거되지 않았으면 다음 active plan/review 파일 또는 follow-up plan을 작성하고 `complete.log`를 작성하지 않는다. - [ ] USER_REVIEW이면 `agent-ops/skills/common/code-review/templates/user-review-template.md` 기준으로 `USER_REVIEW.md`를 작성하고 active `PLAN-*.md`, `CODE_REVIEW-*.md`, `complete.log`를 남기지 않는다. ## 계획 대비 변경 사항 -_구현 에이전트가 계획과 다르게 구현한 부분을 이유와 함께 기록한다._ +- `oto_agent_registration_test.dart` 수정 없음: 이미 `group('OtoServerRegistrationClient (Compatibility Fallback HTTP)', ...)` 구조로 compat-boundary 단계에서 정리되어 있었다. +- `SDD.md` 수정 없음: SDD Evidence Map은 이미 S06 evidence 경로를 올바르게 가리키고 있어 변경 불필요. +- Makefile 수정 없음: `RUNNER_REQUIRED_TESTS`에 `oto_server_connection_smoke_test.dart`가 이미 포함되어 있고 smoke 문서의 agent unit 명령과 일치함을 확인. ## 주요 설계 결정 -_구현 에이전트가 주요 설계 결정 사항을 기록한다._ +- `oto_server_connection_smoke_test.dart` `main()` 내부를 `group('socket-first transport smoke', ...)`과 `group('compatibility fallback HTTP', ...)`로 분리하고 socket 그룹을 앞에 배치했다. 기존 테스트 코드는 내용 변경 없이 indentation(2→4 spaces) 조정과 그룹 래핑만 수행했다. +- 헬퍼 함수 영역의 섹션 주석을 `// ─── Socket session lifecycle tests`에서 `// ─── Test helpers`로 변경해 실제 내용(helper 함수들)과 일치하도록 했다. +- `agent-test/local/agent-smoke.md`에 `services/core/internal/runnersocket/**` 읽기 조건 추가, socket-first transport smoke group 통과를 필수 검증/판정 기준에 명시했다. ## 사용자 리뷰 요청 @@ -111,38 +115,52 @@ _기본값은 `없음`이다. 구현 중 새 결정이 필요해 보여도 직 ## 검증 결과 -_구현 에이전트가 각 중간 검증 및 최종 검증 명령 실행 후 출력을 여기에 붙여 넣는다._ - ### API-1 중간 검증 ```bash $ cd apps/runner && dart analyze -(output) +Analyzing runner... +No issues found! ``` ### API-2 중간 검증 ```bash $ cd apps/runner && dart test test/oto_server_connection_smoke_test.dart -(output) +00:03 +3: socket-first transport smoke OtoServerSocketRegistrationClient socket session close is idempotent after server disconnect +00:04 +4: compatibility fallback HTTP OTO Dart runner registers with Go OTO Server via compatibility fallback HTTP, goes online, and disconnects on close +... +00:10 +12: All tests passed! ``` -### API-3 중간 검증 +### API-3 중간 검증 (최종 검증 겸용) ```bash $ cd services/core && go test -count=1 ./... -(output) +ok github.com/toki/oto/services/core/internal/cicdstate 0.005s +ok github.com/toki/oto/services/core/internal/httpserver 0.078s +ok github.com/toki/oto/services/core/internal/runnerregistry 0.004s +ok github.com/toki/oto/services/core/internal/runnersocket 0.049s + $ cd apps/runner && dart analyze -(output) +Analyzing runner... +No issues found! + $ cd apps/runner && dart test test/oto_agent_migration_plan_test.dart test/oto_agent_bootstrap_script_test.dart test/oto_agent_cli_test.dart test/oto_agent_config_test.dart test/oto_agent_registration_test.dart test/oto_server_connection_smoke_test.dart -(output) +00:14 +96: All tests passed! ``` ### 최종 검증 ```bash $ cd services/core && go test -count=1 ./... -(output) +ok github.com/toki/oto/services/core/internal/cicdstate 0.005s +ok github.com/toki/oto/services/core/internal/httpserver 0.078s +ok github.com/toki/oto/services/core/internal/runnerregistry 0.004s +ok github.com/toki/oto/services/core/internal/runnersocket 0.049s + $ cd apps/runner && dart analyze -(output) +Analyzing runner... +No issues found! + $ cd apps/runner && dart test test/oto_agent_migration_plan_test.dart test/oto_agent_bootstrap_script_test.dart test/oto_agent_cli_test.dart test/oto_agent_config_test.dart test/oto_agent_registration_test.dart test/oto_server_connection_smoke_test.dart -(output) +00:14 +96: All tests passed! ``` --- @@ -150,3 +168,24 @@ $ cd apps/runner && dart test test/oto_agent_migration_plan_test.dart test/oto_a > **[IMPLEMENTING AGENT — BEFORE SAVING] Have you filled in every implementation-owned section: completion table, implementation checklist, changes from plan, design decisions, and verification output?** > If anything is blank, go back and fill it in before saving this file. > Leave review-agent-only sections unchanged. + +## 코드리뷰 결과 + +- 종합 판정: PASS +- 차원별 평가: + - correctness: Pass + - completeness: Pass + - test coverage: Pass + - API contract: Pass + - code quality: Pass + - plan deviation: Pass + - verification trust: Pass + - spec conformance: Pass +- 발견된 문제: 없음 +- 검증 확인: + - `git diff --check` - PASS; whitespace error 없음. + - `cd apps/runner && dart analyze` - PASS; `No issues found!`. + - `cd services/core && go test -count=1 ./...` - PASS; core packages 통과. + - `cd apps/runner && dart test test/oto_agent_migration_plan_test.dart test/oto_agent_bootstrap_script_test.dart test/oto_agent_cli_test.dart test/oto_agent_config_test.dart test/oto_agent_registration_test.dart test/oto_server_connection_smoke_test.dart` - PASS; `+96: All tests passed!`. + - `cd apps/runner && dart test -r expanded test/oto_server_connection_smoke_test.dart` - PASS; socket-first transport smoke 3개와 compatibility fallback HTTP group이 분리되어 `+12: All tests passed!`. +- 다음 단계: PASS 절차로 active plan/review를 log로 아카이브하고 `complete.log` 작성 후 task directory를 archive로 이동한다. diff --git a/agent-task/archive/2026/06/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/complete.log b/agent-task/archive/2026/06/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/complete.log new file mode 100644 index 0000000..4d4f7af --- /dev/null +++ b/agent-task/archive/2026/06/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/complete.log @@ -0,0 +1,51 @@ +# Complete - m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke + +## 완료 일시 + +2026-06-21 + +## 요약 + +Runner proto-socket transport hardening의 socket smoke 정리 작업을 1회 리뷰 루프로 완료했다. 최종 판정은 PASS다. + +## 루프 이력 + +| Plan | Review | Verdict | 메모 | +|------|--------|---------|------| +| `plan_cloud_G07_0.log` | `code_review_cloud_G07_0.log` | PASS | socket-first smoke group과 HTTP compatibility fallback group 분리, agent smoke 기준, full local smoke 검증이 충족됨 | + +## 구현/정리 내용 + +- `agent-test/local/agent-smoke.md`에 runnersocket 변경 시 socket-first transport smoke group 통과 기준을 추가했다. +- `apps/runner/test/oto_server_connection_smoke_test.dart`에서 socket-first transport smoke 3개를 앞쪽 그룹으로 분리하고 HTTP compatibility fallback tests를 별도 그룹으로 유지했다. +- 선행 `dispatch-state`, `runner-actions`, `compat-boundary` split PASS complete evidence를 확인하고 S06 smoke evidence로 종합했다. + +## 최종 검증 + +- `git diff --check` - PASS; whitespace error 없음. +- `cd apps/runner && dart analyze` - PASS; `No issues found!`. +- `cd services/core && go test -count=1 ./...` - PASS; `internal/cicdstate`, `internal/httpserver`, `internal/runnerregistry`, `internal/runnersocket` 통과. +- `cd apps/runner && dart test test/oto_agent_migration_plan_test.dart test/oto_agent_bootstrap_script_test.dart test/oto_agent_cli_test.dart test/oto_agent_config_test.dart test/oto_agent_registration_test.dart test/oto_server_connection_smoke_test.dart` - PASS; `+96: All tests passed!`. +- `cd apps/runner && dart test -r expanded test/oto_server_connection_smoke_test.dart` - PASS; socket-first transport smoke 3개와 compatibility fallback HTTP tests가 분리되어 `+12: All tests passed!`. + +## Roadmap Completion + +- Milestone: `agent-roadmap/phase/control-plane-product-surface/milestones/runner-proto-socket-transport-hardening.md` +- Completed task ids: + - `socket-smoke`: PASS; evidence=`agent-task/archive/2026/06/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/plan_cloud_G07_0.log`, `agent-task/archive/2026/06/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/code_review_cloud_G07_0.log`; verification=`git diff --check`, `cd apps/runner && dart analyze`, `cd services/core && go test -count=1 ./...`, `cd apps/runner && dart test test/oto_agent_migration_plan_test.dart test/oto_agent_bootstrap_script_test.dart test/oto_agent_cli_test.dart test/oto_agent_config_test.dart test/oto_agent_registration_test.dart test/oto_server_connection_smoke_test.dart`, `cd apps/runner && dart test -r expanded test/oto_server_connection_smoke_test.dart` +- Not completed task ids: 없음 + +## Spec Completion + +- SDD: `agent-roadmap/sdd/control-plane-product-surface/runner-proto-socket-transport-hardening/SDD.md` +- Completed scenario ids: + - `S06`: PASS; task=`socket-smoke`; evidence=`agent-task/archive/2026/06/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/plan_cloud_G07_0.log`, `agent-task/archive/2026/06/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/code_review_cloud_G07_0.log`; verification=`git diff --check`, `cd apps/runner && dart analyze`, `cd services/core && go test -count=1 ./...`, `cd apps/runner && dart test test/oto_agent_migration_plan_test.dart test/oto_agent_bootstrap_script_test.dart test/oto_agent_cli_test.dart test/oto_agent_config_test.dart test/oto_agent_registration_test.dart test/oto_server_connection_smoke_test.dart`, `cd apps/runner && dart test -r expanded test/oto_server_connection_smoke_test.dart` +- Not completed scenario ids: 없음 + +## 잔여 Nit + +- 없음 + +## 후속 작업 + +- 없음 diff --git a/agent-task/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/PLAN-cloud-G07.md b/agent-task/archive/2026/06/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/plan_cloud_G07_0.log similarity index 100% rename from agent-task/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/PLAN-cloud-G07.md rename to agent-task/archive/2026/06/m-runner-proto-socket-transport-hardening/05+02,03,04_socket_smoke/plan_cloud_G07_0.log diff --git a/agent-test/local/agent-smoke.md b/agent-test/local/agent-smoke.md index 8cfcd86..e63cd8f 100644 --- a/agent-test/local/agent-smoke.md +++ b/agent-test/local/agent-smoke.md @@ -10,7 +10,7 @@ last_rule_updated_at: 2026-06-12 ## 읽기 조건 -- `apps/runner/lib/oto/agent/**`, `apps/runner/lib/cli/commands/command_agent.dart`, `apps/runner/assets/script/shell/oto_agent_bootstrap.sh`, `services/core/cmd/oto-core/main.go`, `services/core/internal/httpserver/**`, 또는 `services/core/internal/runnerregistry/**` 변경 검증이 필요한 경우 +- `apps/runner/lib/oto/agent/**`, `apps/runner/lib/cli/commands/command_agent.dart`, `apps/runner/assets/script/shell/oto_agent_bootstrap.sh`, `services/core/cmd/oto-core/main.go`, `services/core/internal/httpserver/**`, `services/core/internal/runnersocket/**`, 또는 `services/core/internal/runnerregistry/**` 변경 검증이 필요한 경우 ## 적용 범위 @@ -70,6 +70,7 @@ last_rule_updated_at: 2026-06-12 - runner agent 도메인 변경 후 agent 관련 `cd apps/runner && dart test ...` 명령을 실행한다. - services/core registration, bootstrap command, release URL validation 변경 후 `cd services/core && go test ./...`를 실행한다. - OS별 bootstrap matrix 또는 release asset 이름을 바꾸면 `README.md`, `apps/runner/test/oto_agent_bootstrap_script_test.dart`, 이 문서의 asset 이름이 함께 갱신되었는지 확인한다. +- `apps/runner/lib/oto/agent/registration_client.dart` 또는 `services/core/internal/runnersocket/**` 변경 후 `oto_server_connection_smoke_test.dart` socket-first transport smoke group(heartbeat/lifecycle, cancel dispatch, idempotent close)이 통과해야 한다. ## 보조 검증 @@ -82,6 +83,7 @@ last_rule_updated_at: 2026-06-12 - `cd apps/runner && dart analyze`가 issue 없이 종료한다. - 지정한 `cd apps/runner && dart test`가 모두 통과한다. - 지정한 `cd services/core && go test ./...`가 모두 통과한다. +- `oto_server_connection_smoke_test.dart` socket-first transport smoke group의 heartbeat/lifecycle, cancel dispatch, idempotent close tests가 모두 통과한다. HTTP compatibility fallback HTTP group은 보조 evidence이며 기본 smoke 판정 기준이 아니다. ## 기준 출력 예시 diff --git a/apps/runner/test/oto_server_connection_smoke_test.dart b/apps/runner/test/oto_server_connection_smoke_test.dart index 8dda7ff..90601cb 100644 --- a/apps/runner/test/oto_server_connection_smoke_test.dart +++ b/apps/runner/test/oto_server_connection_smoke_test.dart @@ -18,391 +18,907 @@ const _runnerId = 'oto-smoke-runner'; const _runnerAlias = 'oto-smoke-alias'; void main() { - test( - 'OTO Dart runner registers with Go OTO Server via compatibility fallback HTTP, goes online, and disconnects on close', - () async { - final port = await _freePort(); - final serverAddr = '$_host:$port'; + group('socket-first transport smoke', () { + test( + 'OtoServerSocketRegistrationClient socket session sends heartbeat and disconnects on close', + () async { + final port = await _freePort(); + final serverAddr = '$_host:$port'; - // Start Go OTO Server in background - final serverProcess = await _startCore(serverAddr); + final (:process, :socketAddr) = await _startCoreWithSocket(serverAddr); - final output = StringBuffer(); - final stdoutSub = serverProcess.stdout - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - final stderrSub = serverProcess.stderr - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); + final output = StringBuffer(); + final stdoutSub = process.stdout + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + final stderrSub = process.stderr + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); - try { - // Wait for server port to listen - await _waitForPort(_host, port, serverProcess, output); - - final agentConfig = AgentConfig( - agent: const AgentIdentityConfig( - id: _runnerId, - alias: _runnerAlias, - enrollmentToken: _token, - ), - server: ServerConnectionConfig(url: 'http://$serverAddr'), - runtime: const AgentRuntimeConfig( - installDir: '/tmp/install', - workspaceRoot: '/tmp/workspace', - logDir: '/tmp/log', - ), - ); - - // Open OTO Server session (with short heartbeat interval for testing) - final client = OtoServerRegistrationClient( - commandTypes: ['Shell', 'Git'], - heartbeatInterval: const Duration(milliseconds: 200), - ); - - final session = await client.openSession(agentConfig); - final result = session.result; - - expect(result.accepted, isTrue); - expect(result.runnerId, _runnerId); - expect(result.alias, _runnerAlias); - - // Poll the server's GET endpoint until status becomes 'online' - final httpClient = http.Client(); try { - final statusUrl = Uri.parse( - 'http://$serverAddr/api/v1/runners/$_runnerId', - ); - var isOnline = false; - final deadline = DateTime.now().add(const Duration(seconds: 10)); + await _waitForPort(_host, port, process, output); - while (DateTime.now().isBefore(deadline)) { - final response = await httpClient.get(statusUrl); - if (response.statusCode == 200) { - final data = jsonDecode(response.body); - if (data['status'] == 'online') { - isOnline = true; - break; - } - } - await Future.delayed(const Duration(milliseconds: 200)); - } - - expect( - isOnline, - isTrue, - reason: 'Runner did not transition to online state in registry', + final agentConfig = AgentConfig( + agent: const AgentIdentityConfig( + id: _runnerId, + alias: _runnerAlias, + enrollmentToken: _token, + ), + server: ServerConnectionConfig( + url: 'http://$serverAddr', + socketUrl: 'tcp://$socketAddr', + ), + runtime: const AgentRuntimeConfig( + installDir: '/tmp/install', + workspaceRoot: '/tmp/workspace', + logDir: '/tmp/log', + ), ); - // Close the session to trigger disconnect - await session.close(); - - // Verify status transitions to disconnected - var isDisconnected = false; - final disconnectDeadline = DateTime.now().add( - const Duration(seconds: 5), + final client = OtoServerSocketRegistrationClient( + commandTypes: ['Shell', 'Git'], + heartbeatInterval: const Duration(milliseconds: 200), ); - while (DateTime.now().isBefore(disconnectDeadline)) { - final response = await httpClient.get(statusUrl); - if (response.statusCode == 200) { - final data = jsonDecode(response.body); - if (data['status'] == 'disconnected') { - isDisconnected = true; - break; - } - } - await Future.delayed(const Duration(milliseconds: 200)); - } - - expect( - isDisconnected, - isTrue, - reason: - 'Runner did not transition to disconnected state in registry', - ); - } finally { - httpClient.close(); - } - } finally { - serverProcess.kill(ProcessSignal.sigterm); - await serverProcess.exitCode.timeout( - const Duration(seconds: 5), - onTimeout: () { - serverProcess.kill(ProcessSignal.sigkill); - return serverProcess.exitCode; - }, - ); - await stdoutSub.cancel(); - await stderrSub.cancel(); - } - }, - timeout: const Timeout(Duration(seconds: 30)), - ); - - test( - 'OTO Server issues runner bootstrap command and serves bootstrap script', - () async { - final port = await _freePort(); - final serverAddr = '$_host:$port'; - - // Start Go OTO Server in background - final serverProcess = await _startCore( - serverAddr, - releaseBaseUrl: 'https://example.com/releases', - ); - - final output = StringBuffer(); - final stdoutSub = serverProcess.stdout - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - final stderrSub = serverProcess.stderr - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - - try { - // Wait for server port to listen - await _waitForPort(_host, port, serverProcess, output); - - final httpClient = http.Client(); - try { - // 1. Verify bootstrap command endpoint - final bootstrapCmdUrl = Uri.parse( - 'http://$serverAddr/api/v1/runners/bootstrap-command', - ); - final response = await httpClient.post( - bootstrapCmdUrl, - headers: {'content-type': 'application/json'}, - body: jsonEncode({ - 'runner_id': 'test-runner-id', - 'enrollment_token': 'test-token-123', - }), - ); - - expect(response.statusCode, equals(200)); - final data = jsonDecode(response.body); - expect(data['bootstrap_command'], isNotNull); - final bootstrapCommand = data['bootstrap_command'] as String; - - // Verify command contains server URL, runner ID, enrollment token - expect(bootstrapCommand, contains('http://$serverAddr')); - expect(bootstrapCommand, contains('test-runner-id')); - expect(bootstrapCommand, contains('test-token-123')); - expect(bootstrapCommand, contains('--server-url')); - expect(bootstrapCommand, contains('--agent-id')); - expect(bootstrapCommand, contains('--enrollment-token')); - - // 2. Full script execution test using fake curl/tar - final tempDir = await Directory.systemTemp.createTemp( - 'oto_smoke_script_run_', - ); - final tempHome = Directory('${tempDir.path}/home'); - await tempHome.create(recursive: true); - - final scriptFile = File('assets/script/shell/oto_agent_bootstrap.sh'); - - // Create fake OTO executable source - final fakeOtoSource = Directory('${tempDir.path}/fake_oto_src'); - await fakeOtoSource.create(recursive: true); - final fakeOtoFile = File('${fakeOtoSource.path}/oto'); - await fakeOtoFile.writeAsString( - '#!/bin/sh\necho "fake-oto-started"\n', - ); - await Process.run('chmod', ['+x', fakeOtoFile.path]); - - // Create fake tar.gz archive - final fakeTarGz = File('${tempDir.path}/oto-linux-x64.tar.gz'); - await Process.run('tar', [ - '-czf', - fakeTarGz.path, - '-C', - fakeOtoSource.path, - 'oto', - ]); - - // Create fake bin dir and fake curl - final fakeBinDir = Directory('${tempDir.path}/bin'); - await fakeBinDir.create(recursive: true); - final fakeCurl = File('${fakeBinDir.path}/curl'); - await fakeCurl.writeAsString('''#!/bin/sh -out_file="" -while [ \$# -gt 0 ]; do - if [ "\$1" = "-o" ]; then - out_file="\$2" - shift 2 - elif [ "\$1" = "-fsSL" ]; then - shift 1 - else - shift 1 - fi -done - -if [ -n "\$out_file" ]; then - cp "${fakeTarGz.path}" "\$out_file" -else - cat "${scriptFile.absolute.path}" -fi -'''); - await Process.run('chmod', ['+x', fakeCurl.path]); - - // Append custom paths and options to the command string - final customConfigPath = '${tempDir.path}/smoke-config.yaml'; - final customInstallDir = '${tempDir.path}/install_dir'; - final customWorkspaceRoot = '${tempDir.path}/workspace'; - final customLogDir = '${tempDir.path}/log'; - - final fullCommand = [ - bootstrapCommand, - "--config-path '$customConfigPath'", - "--install-dir '$customInstallDir'", - "--workspace-root '$customWorkspaceRoot'", - "--log-dir '$customLogDir'", - '--no-background', - ].join(' '); - - // Set up environment with fake curl path and custom HOME - final env = Map.from(Platform.environment); - env['PATH'] = '${fakeBinDir.path}:${env['PATH']}'; - env['HOME'] = tempHome.path; - - // Run the full pipeline via bash - final scriptResult = await Process.run('bash', [ - '-c', - fullCommand, - ], environment: env); - - expect( - scriptResult.exitCode, - equals(0), - reason: - 'Script failed: ${scriptResult.stderr}\nStdout: ${scriptResult.stdout}', - ); - - // Assertions: - // - Generated config exists and has correct values - final configFile = File(customConfigPath); - expect( - await configFile.exists(), - isTrue, - reason: 'Config file not generated', - ); - final configContent = await configFile.readAsString(); - expect(configContent, contains('server:')); - expect(configContent, contains('url: "http://$serverAddr"')); - expect(configContent, contains('id: "test-runner-id"')); - expect(configContent, contains('enrollment_token: "test-token-123"')); - - // - Installed binary exists and has execution permissions - final installedOto = File('$customInstallDir/oto'); - expect( - await installedOto.exists(), - isTrue, - reason: 'OTO binary not installed', - ); - final stat = await installedOto.stat(); - expect( - stat.mode & 0x49, - isNot(0), - reason: 'Installed binary not executable', - ); - - await tempDir.delete(recursive: true); - - // 3. Verify bootstrap script hosting endpoint - final scriptUrl = Uri.parse( - 'http://$serverAddr/bootstrap/oto-agent.sh', - ); - final scriptResponse = await httpClient.get(scriptUrl); - expect(scriptResponse.statusCode, equals(200)); - expect( - scriptResponse.headers['content-type'], - contains('application/x-sh'), - ); - expect(scriptResponse.body, isNotEmpty); - expect(scriptResponse.body, contains('--server-url')); - } finally { - httpClient.close(); - } - } finally { - serverProcess.kill(ProcessSignal.sigterm); - await serverProcess.exitCode.timeout( - const Duration(seconds: 5), - onTimeout: () { - serverProcess.kill(ProcessSignal.sigkill); - return serverProcess.exitCode; - }, - ); - await stdoutSub.cancel(); - await stderrSub.cancel(); - } - }, - timeout: const Timeout(Duration(seconds: 30)), - ); - - test( - 'OTO Server owns job execution logs and artifacts reported by runner via compatibility fallback HTTP', - () async { - final port = await _freePort(); - final serverAddr = '$_host:$port'; - const jobId = 'oto-smoke-job'; - const executionId = 'oto-smoke-execution'; - - final serverProcess = await _startCore(serverAddr); - - final output = StringBuffer(); - final stdoutSub = serverProcess.stdout - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - final stderrSub = serverProcess.stderr - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - - try { - await _waitForPort(_host, port, serverProcess, output); - - final agentConfig = AgentConfig( - agent: const AgentIdentityConfig( - id: _runnerId, - alias: _runnerAlias, - enrollmentToken: _token, - ), - server: ServerConnectionConfig(url: 'http://$serverAddr'), - runtime: const AgentRuntimeConfig( - installDir: '/tmp/install', - workspaceRoot: '/tmp/workspace', - logDir: '/tmp/log', - ), - ); - - final registrationClient = OtoServerRegistrationClient( - commandTypes: ['Shell', 'Git'], - heartbeatInterval: const Duration(milliseconds: 200), - ); - final session = await registrationClient.openSession(agentConfig); - try { + final session = await client.openSession(agentConfig); + expect(session, isA()); expect(session.result.accepted, isTrue); - expect(session, isA()); - final jobs = (session as OtoServerJobSession).jobs; + expect(session.result.runnerId, _runnerId); + + // Capture stream completion after close(). + final pushSession = session as OtoServerPushJobSession; + final runRequestsDone = Completer(); + final cancelRequestsDone = Completer(); + final runSub = pushSession.runRequests.listen( + (_) {}, + onDone: runRequestsDone.complete, + ); + final cancelSub = pushSession.cancelRequests.listen( + (_) {}, + onDone: cancelRequestsDone.complete, + ); final httpClient = http.Client(); try { - const inlineYaml = 'commands:\n - type: Shell'; - final createJobResponse = await httpClient.post( + final isOnline = await _pollRunnerStatus( + httpClient, + serverAddr, + _runnerId, + 'online', + ); + expect( + isOnline, + isTrue, + reason: 'socket runner did not go online in registry', + ); + + // close() must cancel the heartbeat timer and close socket streams. + await session.close(); + await runRequestsDone.future.timeout( + const Duration(seconds: 2), + onTimeout: () => + fail('runRequests stream did not close after session.close()'), + ); + await cancelRequestsDone.future.timeout( + const Duration(seconds: 2), + onTimeout: () => fail( + 'cancelRequests stream did not close after session.close()', + ), + ); + await runSub.cancel(); + await cancelSub.cancel(); + + final isDisconnected = await _pollRunnerStatus( + httpClient, + serverAddr, + _runnerId, + 'disconnected', + timeout: const Duration(seconds: 5), + ); + expect( + isDisconnected, + isTrue, + reason: 'socket runner did not transition to disconnected', + ); + } finally { + httpClient.close(); + } + } finally { + process.kill(ProcessSignal.sigterm); + await process.exitCode.timeout( + const Duration(seconds: 5), + onTimeout: () { + process.kill(ProcessSignal.sigkill); + return process.exitCode; + }, + ); + await stdoutSub.cancel(); + await stderrSub.cancel(); + } + }, + timeout: const Timeout(Duration(seconds: 30)), + ); + + test( + 'OtoServerSocketRegistrationClient receives socket cancel requests from HTTP cancel action', + () async { + final port = await _freePort(); + final serverAddr = '$_host:$port'; + const jobId = 'smoke-job-socket-cancel'; + + final (:process, :socketAddr) = await _startCoreWithSocket(serverAddr); + + final output = StringBuffer(); + final stdoutSub = process.stdout + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + final stderrSub = process.stderr + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + + try { + await _waitForPort(_host, port, process, output); + + final agentConfig = AgentConfig( + agent: const AgentIdentityConfig( + id: _runnerId, + alias: _runnerAlias, + enrollmentToken: _token, + ), + server: ServerConnectionConfig( + url: 'http://$serverAddr', + socketUrl: 'tcp://$socketAddr', + ), + runtime: const AgentRuntimeConfig( + installDir: '/tmp/install', + workspaceRoot: '/tmp/workspace', + logDir: '/tmp/log', + ), + ); + + final client = OtoServerSocketRegistrationClient( + commandTypes: ['Shell'], + heartbeatInterval: const Duration(milliseconds: 200), + ); + final session = await client.openSession(agentConfig); + expect(session.result.accepted, isTrue); + final pushSession = session as OtoServerPushJobSession; + + final runReceived = Completer(); + final runSub = pushSession.runRequests.listen((request) { + if (!runReceived.isCompleted) { + expect(request.runnerId, _runnerId); + expect(request.jobId, jobId); + expect(request.executionId, isNotEmpty); + runReceived.complete(request.executionId); + } + }); + var expectedExecutionId = ''; + final cancelReceived = Completer(); + final cancelSub = pushSession.cancelRequests.listen((request) { + if (!cancelReceived.isCompleted) { + expect(request.runnerId, _runnerId); + expect(request.executionId, expectedExecutionId); + expect(request.reason, 'socket cancel smoke'); + cancelReceived.complete(request.executionId); + } + }); + + final httpClient = http.Client(); + final jobs = OtoServerJobClient( + serverUrl: 'http://$serverAddr', + runnerId: _runnerId, + client: httpClient, + ); + try { + final createResp = await httpClient.post( Uri.parse('http://$serverAddr/api/v1/jobs'), headers: {'content-type': 'application/json'}, body: jsonEncode({ 'id': jobId, - 'name': 'smoke build', + 'name': 'socket cancel smoke', 'run_request': { - 'pipeline_yaml': inlineYaml, - 'variables': {'FLAVOR': 'release'}, - 'command_types': ['Shell', 'Git'], + 'pipeline_yaml': 'commands:\n - type: Shell', + 'command_types': ['Shell'], + }, + }), + ); + expect(createResp.statusCode, equals(201)); + + expectedExecutionId = await runReceived.future.timeout( + const Duration(seconds: 3), + onTimeout: () => fail('socket RunRequest was not received'), + ); + + final cancelRes = await jobs.cancelRun( + executionId: expectedExecutionId, + reason: 'socket cancel smoke', + ); + expect(cancelRes.success, isTrue); + + await cancelReceived.future.timeout( + const Duration(seconds: 3), + onTimeout: () => fail('socket CancelRunRequest was not received'), + ); + + final execResp = await httpClient.get( + Uri.parse( + 'http://$serverAddr/api/v1/executions/$expectedExecutionId', + ), + ); + final execJson = jsonDecode(execResp.body) as Map; + expect(execJson['state'], equals('canceled')); + } finally { + await runSub.cancel(); + await cancelSub.cancel(); + httpClient.close(); + await session.close(); + } + } finally { + process.kill(ProcessSignal.sigterm); + await process.exitCode.timeout( + const Duration(seconds: 5), + onTimeout: () { + process.kill(ProcessSignal.sigkill); + return process.exitCode; + }, + ); + await stdoutSub.cancel(); + await stderrSub.cancel(); + } + }, + timeout: const Timeout(Duration(seconds: 30)), + ); + + test( + 'OtoServerSocketRegistrationClient socket session close is idempotent after server disconnect', + () async { + final port = await _freePort(); + final serverAddr = '$_host:$port'; + + final (:process, :socketAddr) = await _startCoreWithSocket(serverAddr); + + final output = StringBuffer(); + final stdoutSub = process.stdout + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + final stderrSub = process.stderr + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + + try { + await _waitForPort(_host, port, process, output); + + final agentConfig = AgentConfig( + agent: const AgentIdentityConfig( + id: _runnerId, + alias: _runnerAlias, + enrollmentToken: _token, + ), + server: ServerConnectionConfig( + url: 'http://$serverAddr', + socketUrl: 'tcp://$socketAddr', + ), + runtime: const AgentRuntimeConfig( + installDir: '/tmp/install', + workspaceRoot: '/tmp/workspace', + logDir: '/tmp/log', + ), + ); + + final client = OtoServerSocketRegistrationClient( + commandTypes: ['Shell'], + // Short interval so in-flight heartbeats exercise the catch block. + heartbeatInterval: const Duration(milliseconds: 50), + ); + + final session = await client.openSession(agentConfig); + expect(session.result.accepted, isTrue); + + // Graceful server shutdown: closes all proto-socket clients cleanly, + // sending EOF to the Dart side. In-flight heartbeat sendRequest calls + // receive StateError('connection closed'), which is swallowed by the + // _sendHeartbeat catch block. + process.kill(ProcessSignal.sigterm); + await process.exitCode.timeout(const Duration(seconds: 5)); + + // Allow the EOF to propagate and the transport to auto-close. + await Future.delayed(const Duration(milliseconds: 300)); + + // session.close() must complete without throwing even though the + // proto-socket transport is already closed by the EOF handler. + await expectLater(session.close(), completes); + } finally { + process.kill(ProcessSignal.sigkill); + await process.exitCode.timeout( + const Duration(seconds: 2), + onTimeout: () => -1, + ); + await stdoutSub.cancel(); + await stderrSub.cancel(); + } + }, + timeout: const Timeout(Duration(seconds: 30)), + ); + + }); + + group('compatibility fallback HTTP', () { + test( + 'OTO Dart runner registers with Go OTO Server via compatibility fallback HTTP, goes online, and disconnects on close', + () async { + final port = await _freePort(); + final serverAddr = '$_host:$port'; + + // Start Go OTO Server in background + final serverProcess = await _startCore(serverAddr); + + final output = StringBuffer(); + final stdoutSub = serverProcess.stdout + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + final stderrSub = serverProcess.stderr + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + + try { + // Wait for server port to listen + await _waitForPort(_host, port, serverProcess, output); + + final agentConfig = AgentConfig( + agent: const AgentIdentityConfig( + id: _runnerId, + alias: _runnerAlias, + enrollmentToken: _token, + ), + server: ServerConnectionConfig(url: 'http://$serverAddr'), + runtime: const AgentRuntimeConfig( + installDir: '/tmp/install', + workspaceRoot: '/tmp/workspace', + logDir: '/tmp/log', + ), + ); + + // Open OTO Server session (with short heartbeat interval for testing) + final client = OtoServerRegistrationClient( + commandTypes: ['Shell', 'Git'], + heartbeatInterval: const Duration(milliseconds: 200), + ); + + final session = await client.openSession(agentConfig); + final result = session.result; + + expect(result.accepted, isTrue); + expect(result.runnerId, _runnerId); + expect(result.alias, _runnerAlias); + + // Poll the server's GET endpoint until status becomes 'online' + final httpClient = http.Client(); + try { + final statusUrl = Uri.parse( + 'http://$serverAddr/api/v1/runners/$_runnerId', + ); + var isOnline = false; + final deadline = DateTime.now().add(const Duration(seconds: 10)); + + while (DateTime.now().isBefore(deadline)) { + final response = await httpClient.get(statusUrl); + if (response.statusCode == 200) { + final data = jsonDecode(response.body); + if (data['status'] == 'online') { + isOnline = true; + break; + } + } + await Future.delayed(const Duration(milliseconds: 200)); + } + + expect( + isOnline, + isTrue, + reason: 'Runner did not transition to online state in registry', + ); + + // Close the session to trigger disconnect + await session.close(); + + // Verify status transitions to disconnected + var isDisconnected = false; + final disconnectDeadline = DateTime.now().add( + const Duration(seconds: 5), + ); + + while (DateTime.now().isBefore(disconnectDeadline)) { + final response = await httpClient.get(statusUrl); + if (response.statusCode == 200) { + final data = jsonDecode(response.body); + if (data['status'] == 'disconnected') { + isDisconnected = true; + break; + } + } + await Future.delayed(const Duration(milliseconds: 200)); + } + + expect( + isDisconnected, + isTrue, + reason: + 'Runner did not transition to disconnected state in registry', + ); + } finally { + httpClient.close(); + } + } finally { + serverProcess.kill(ProcessSignal.sigterm); + await serverProcess.exitCode.timeout( + const Duration(seconds: 5), + onTimeout: () { + serverProcess.kill(ProcessSignal.sigkill); + return serverProcess.exitCode; + }, + ); + await stdoutSub.cancel(); + await stderrSub.cancel(); + } + }, + timeout: const Timeout(Duration(seconds: 30)), + ); + + test( + 'OTO Server issues runner bootstrap command and serves bootstrap script', + () async { + final port = await _freePort(); + final serverAddr = '$_host:$port'; + + // Start Go OTO Server in background + final serverProcess = await _startCore( + serverAddr, + releaseBaseUrl: 'https://example.com/releases', + ); + + final output = StringBuffer(); + final stdoutSub = serverProcess.stdout + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + final stderrSub = serverProcess.stderr + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + + try { + // Wait for server port to listen + await _waitForPort(_host, port, serverProcess, output); + + final httpClient = http.Client(); + try { + // 1. Verify bootstrap command endpoint + final bootstrapCmdUrl = Uri.parse( + 'http://$serverAddr/api/v1/runners/bootstrap-command', + ); + final response = await httpClient.post( + bootstrapCmdUrl, + headers: {'content-type': 'application/json'}, + body: jsonEncode({ + 'runner_id': 'test-runner-id', + 'enrollment_token': 'test-token-123', + }), + ); + + expect(response.statusCode, equals(200)); + final data = jsonDecode(response.body); + expect(data['bootstrap_command'], isNotNull); + final bootstrapCommand = data['bootstrap_command'] as String; + + // Verify command contains server URL, runner ID, enrollment token + expect(bootstrapCommand, contains('http://$serverAddr')); + expect(bootstrapCommand, contains('test-runner-id')); + expect(bootstrapCommand, contains('test-token-123')); + expect(bootstrapCommand, contains('--server-url')); + expect(bootstrapCommand, contains('--agent-id')); + expect(bootstrapCommand, contains('--enrollment-token')); + + // 2. Full script execution test using fake curl/tar + final tempDir = await Directory.systemTemp.createTemp( + 'oto_smoke_script_run_', + ); + final tempHome = Directory('${tempDir.path}/home'); + await tempHome.create(recursive: true); + + final scriptFile = File('assets/script/shell/oto_agent_bootstrap.sh'); + + // Create fake OTO executable source + final fakeOtoSource = Directory('${tempDir.path}/fake_oto_src'); + await fakeOtoSource.create(recursive: true); + final fakeOtoFile = File('${fakeOtoSource.path}/oto'); + await fakeOtoFile.writeAsString( + '#!/bin/sh\necho "fake-oto-started"\n', + ); + await Process.run('chmod', ['+x', fakeOtoFile.path]); + + // Create fake tar.gz archive + final fakeTarGz = File('${tempDir.path}/oto-linux-x64.tar.gz'); + await Process.run('tar', [ + '-czf', + fakeTarGz.path, + '-C', + fakeOtoSource.path, + 'oto', + ]); + + // Create fake bin dir and fake curl + final fakeBinDir = Directory('${tempDir.path}/bin'); + await fakeBinDir.create(recursive: true); + final fakeCurl = File('${fakeBinDir.path}/curl'); + await fakeCurl.writeAsString('''#!/bin/sh + out_file="" + while [ \$# -gt 0 ]; do + if [ "\$1" = "-o" ]; then + out_file="\$2" + shift 2 + elif [ "\$1" = "-fsSL" ]; then + shift 1 + else + shift 1 + fi + done + + if [ -n "\$out_file" ]; then + cp "${fakeTarGz.path}" "\$out_file" + else + cat "${scriptFile.absolute.path}" + fi + '''); + await Process.run('chmod', ['+x', fakeCurl.path]); + + // Append custom paths and options to the command string + final customConfigPath = '${tempDir.path}/smoke-config.yaml'; + final customInstallDir = '${tempDir.path}/install_dir'; + final customWorkspaceRoot = '${tempDir.path}/workspace'; + final customLogDir = '${tempDir.path}/log'; + + final fullCommand = [ + bootstrapCommand, + "--config-path '$customConfigPath'", + "--install-dir '$customInstallDir'", + "--workspace-root '$customWorkspaceRoot'", + "--log-dir '$customLogDir'", + '--no-background', + ].join(' '); + + // Set up environment with fake curl path and custom HOME + final env = Map.from(Platform.environment); + env['PATH'] = '${fakeBinDir.path}:${env['PATH']}'; + env['HOME'] = tempHome.path; + + // Run the full pipeline via bash + final scriptResult = await Process.run('bash', [ + '-c', + fullCommand, + ], environment: env); + + expect( + scriptResult.exitCode, + equals(0), + reason: + 'Script failed: ${scriptResult.stderr}\nStdout: ${scriptResult.stdout}', + ); + + // Assertions: + // - Generated config exists and has correct values + final configFile = File(customConfigPath); + expect( + await configFile.exists(), + isTrue, + reason: 'Config file not generated', + ); + final configContent = await configFile.readAsString(); + expect(configContent, contains('server:')); + expect(configContent, contains('url: "http://$serverAddr"')); + expect(configContent, contains('id: "test-runner-id"')); + expect(configContent, contains('enrollment_token: "test-token-123"')); + + // - Installed binary exists and has execution permissions + final installedOto = File('$customInstallDir/oto'); + expect( + await installedOto.exists(), + isTrue, + reason: 'OTO binary not installed', + ); + final stat = await installedOto.stat(); + expect( + stat.mode & 0x49, + isNot(0), + reason: 'Installed binary not executable', + ); + + await tempDir.delete(recursive: true); + + // 3. Verify bootstrap script hosting endpoint + final scriptUrl = Uri.parse( + 'http://$serverAddr/bootstrap/oto-agent.sh', + ); + final scriptResponse = await httpClient.get(scriptUrl); + expect(scriptResponse.statusCode, equals(200)); + expect( + scriptResponse.headers['content-type'], + contains('application/x-sh'), + ); + expect(scriptResponse.body, isNotEmpty); + expect(scriptResponse.body, contains('--server-url')); + } finally { + httpClient.close(); + } + } finally { + serverProcess.kill(ProcessSignal.sigterm); + await serverProcess.exitCode.timeout( + const Duration(seconds: 5), + onTimeout: () { + serverProcess.kill(ProcessSignal.sigkill); + return serverProcess.exitCode; + }, + ); + await stdoutSub.cancel(); + await stderrSub.cancel(); + } + }, + timeout: const Timeout(Duration(seconds: 30)), + ); + + test( + 'OTO Server owns job execution logs and artifacts reported by runner via compatibility fallback HTTP', + () async { + final port = await _freePort(); + final serverAddr = '$_host:$port'; + const jobId = 'oto-smoke-job'; + const executionId = 'oto-smoke-execution'; + + final serverProcess = await _startCore(serverAddr); + + final output = StringBuffer(); + final stdoutSub = serverProcess.stdout + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + final stderrSub = serverProcess.stderr + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + + try { + await _waitForPort(_host, port, serverProcess, output); + + final agentConfig = AgentConfig( + agent: const AgentIdentityConfig( + id: _runnerId, + alias: _runnerAlias, + enrollmentToken: _token, + ), + server: ServerConnectionConfig(url: 'http://$serverAddr'), + runtime: const AgentRuntimeConfig( + installDir: '/tmp/install', + workspaceRoot: '/tmp/workspace', + logDir: '/tmp/log', + ), + ); + + final registrationClient = OtoServerRegistrationClient( + commandTypes: ['Shell', 'Git'], + heartbeatInterval: const Duration(milliseconds: 200), + ); + final session = await registrationClient.openSession(agentConfig); + try { + expect(session.result.accepted, isTrue); + expect(session, isA()); + final jobs = (session as OtoServerJobSession).jobs; + + final httpClient = http.Client(); + try { + const inlineYaml = 'commands:\n - type: Shell'; + final createJobResponse = await httpClient.post( + Uri.parse('http://$serverAddr/api/v1/jobs'), + headers: {'content-type': 'application/json'}, + body: jsonEncode({ + 'id': jobId, + 'name': 'smoke build', + 'run_request': { + 'pipeline_yaml': inlineYaml, + 'variables': {'FLAVOR': 'release'}, + 'command_types': ['Shell', 'Git'], + }, + }), + ); + expect( + createJobResponse.statusCode, + equals(201), + reason: 'Create job failed: ${createJobResponse.body}', + ); + + final claim = await jobs.claimJob( + jobId: jobId, + executionId: executionId, + ); + expect(claim.accepted, isTrue); + expect(claim.state, 'running'); + expect(claim.runRequest, isNotNull); + expect(claim.runRequest!.pipelineYaml, inlineYaml); + expect(claim.runRequest!.jobId, jobId); + expect(claim.runRequest!.executionId, executionId); + + await jobs.appendLog( + executionId: executionId, + line: 'runner started smoke build', + ); + await jobs.reportArtifact( + executionId: executionId, + name: 'smoke-report', + path: '/tmp/oto-smoke/report.txt', + ); + final report = await jobs.reportExecution( + jobId: jobId, + executionId: executionId, + result: BuildResult.success( + stepEvents: [ + StepEvent( + stepId: 0, + workflowIndex: 0, + stepType: 'command', + event: 'started', + timestamp: DateTime.now().toUtc().toIso8601String(), + commandId: 'smoke-build', + commandType: 'Shell', + ), + StepEvent( + stepId: 0, + workflowIndex: 0, + stepType: 'command', + event: 'completed', + timestamp: DateTime.now().toUtc().toIso8601String(), + commandId: 'smoke-build', + commandType: 'Shell', + ), + ], + ), + ); + expect(report.accepted, isTrue); + expect(report.state, 'succeeded'); + + final job = await _getJson( + httpClient, + 'http://$serverAddr/api/v1/jobs/$jobId', + ); + expect(job['state'], 'succeeded'); + expect(job['execution_id'], executionId); + + final execution = await _getJson( + httpClient, + 'http://$serverAddr/api/v1/executions/$executionId', + ); + expect(execution['state'], 'succeeded'); + expect(execution['job_id'], jobId); + + final logs = await _getJson( + httpClient, + 'http://$serverAddr/api/v1/executions/$executionId/logs', + ); + final logLines = (logs['logs'] as List) + .map( + (entry) => + (entry as Map)['Line'] ?? entry['line'], + ) + .map((line) => line.toString()) + .toList(); + expect(logLines, contains('runner started smoke build')); + expect(logLines, contains('Build completed successfully.')); + + final artifacts = await _getJson( + httpClient, + 'http://$serverAddr/api/v1/executions/$executionId/artifacts', + ); + final artifactRows = artifacts['artifacts'] as List; + expect(artifactRows, hasLength(1)); + final artifact = artifactRows.single as Map; + expect(artifact['Name'] ?? artifact['name'], 'smoke-report'); + expect( + artifact['Path'] ?? artifact['path'], + '/tmp/oto-smoke/report.txt', + ); + } finally { + httpClient.close(); + } + } finally { + await session.close(); + } + } finally { + serverProcess.kill(ProcessSignal.sigterm); + await serverProcess.exitCode.timeout( + const Duration(seconds: 5), + onTimeout: () { + serverProcess.kill(ProcessSignal.sigkill); + return serverProcess.exitCode; + }, + ); + await stdoutSub.cancel(); + await stderrSub.cancel(); + } + }, + timeout: const Timeout(Duration(seconds: 30)), + ); + + test( + 'remoteRunExecutor runOnce executes job via Go OTO Server and reports step events (compatibility fallback HTTP)', + () async { + final port = await _freePort(); + final serverAddr = '$_host:$port'; + const jobId = 'remote-smoke-job'; + const executionId = 'remote-smoke-exec'; + + final serverProcess = await _startCore(serverAddr); + + final output = StringBuffer(); + final stdoutSub = serverProcess.stdout + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + final stderrSub = serverProcess.stderr + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + + try { + await _waitForPort(_host, port, serverProcess, output); + + final agentConfig = AgentConfig( + agent: const AgentIdentityConfig( + id: _runnerId, + alias: _runnerAlias, + enrollmentToken: _token, + ), + server: ServerConnectionConfig(url: 'http://$serverAddr'), + runtime: const AgentRuntimeConfig( + installDir: '/tmp/install', + workspaceRoot: '/tmp/workspace', + logDir: '/tmp/log', + ), + ); + + final regClient = OtoServerRegistrationClient( + commandTypes: ['Shell', 'Git'], + ); + final session = await regClient.openSession(agentConfig); + expect( + session.result.accepted, + isTrue, + reason: 'Registration was rejected', + ); + final jobs = (session as OtoServerJobSession).jobs; + + final httpClient = http.Client(); + try { + // Create a job with valid pipeline YAML + final createResp = await httpClient.post( + Uri.parse('http://$serverAddr/api/v1/jobs'), + headers: {'content-type': 'application/json'}, + body: jsonEncode({ + 'id': jobId, + 'name': 'remote run smoke', + 'run_request': { + 'pipeline_yaml': ''' + property: + workspace: . + commands: + - command: Print + id: say-hello + param: + message: hello from smoke + pipeline: + id: smoke-pipe + workflow: + - exe: say-hello + ''', + 'command_types': ['Print'], }, }), ); expect( - createJobResponse.statusCode, + createResp.statusCode, equals(201), - reason: 'Create job failed: ${createJobResponse.body}', + reason: 'Create job failed: ${createResp.body}', ); final claim = await jobs.claimJob( @@ -410,78 +926,608 @@ fi executionId: executionId, ); expect(claim.accepted, isTrue); - expect(claim.state, 'running'); expect(claim.runRequest, isNotNull); - expect(claim.runRequest!.pipelineYaml, inlineYaml); - expect(claim.runRequest!.jobId, jobId); - expect(claim.runRequest!.executionId, executionId); - await jobs.appendLog( - executionId: executionId, - line: 'runner started smoke build', + // Execute via RemoteRunExecutor.runOnce + final parser = RunRequestParser(workspaceRoot: '/tmp/workspace'); + final executor = RemoteRunExecutor( + parser: parser, + jobClient: jobs, + onLog: (msg) {}, ); - await jobs.reportArtifact( - executionId: executionId, - name: 'smoke-report', - path: '/tmp/oto-smoke/report.txt', + + final buildResult = await executor.runOnce( + jobId: claim.jobId, + executionId: claim.executionId, + runRequest: claim.runRequest!, ); - final report = await jobs.reportExecution( + + expect( + buildResult.success, + isTrue, + reason: 'Execution failed: ${buildResult.message}', + ); + expect(buildResult.exitCode, 0); + expect(buildResult.stepEvents, isNotEmpty); + + final events = buildResult.stepEvents.map((e) => e.event).toList(); + expect(events, contains('started')); + expect(events, contains('completed')); + + // Verify logs on server + final logsResp = await httpClient.get( + Uri.parse('http://$serverAddr/api/v1/executions/$executionId/logs'), + ); + final logsJson = jsonDecode(logsResp.body) as Map; + final logsList = logsJson['logs'] as List; + final lines = logsList + .map((e) => (e as Map)['Line'] ?? e['line']) + .map((l) => l.toString()) + .toList(); + expect(lines, contains('run started for job $jobId')); + expect(lines, contains('Build completed successfully.')); + + // Verify execution state + final execResp = await httpClient.get( + Uri.parse('http://$serverAddr/api/v1/executions/$executionId'), + ); + expect(execResp.statusCode, 200); + final execJson = jsonDecode(execResp.body) as Map; + expect(execJson['state'], equals('succeeded')); + expect(execJson['job_id'], equals(jobId)); + + httpClient.close(); + } finally { + await session.close(); + } + } finally { + serverProcess.kill(ProcessSignal.sigterm); + await serverProcess.exitCode.timeout( + const Duration(seconds: 5), + onTimeout: () { + serverProcess.kill(ProcessSignal.sigkill); + return serverProcess.exitCode; + }, + ); + await stdoutSub.cancel(); + await stderrSub.cancel(); + } + }, + timeout: const Timeout(Duration(seconds: 30)), + ); + + test( + 'remoteRunExecutor runOnce merges variables into pipeline property (compatibility fallback HTTP) (REVIEW_API-3)', + () async { + final port = await _freePort(); + final serverAddr = '$_host:$port'; + const jobId = 'var-merge-job'; + const executionId = 'var-merge-exec'; + + final serverProcess = await _startCore(serverAddr); + + final output = StringBuffer(); + final stdoutSub = serverProcess.stdout + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + final stderrSub = serverProcess.stderr + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + + try { + await _waitForPort(_host, port, serverProcess, output); + + final agentConfig = AgentConfig( + agent: const AgentIdentityConfig( + id: _runnerId, + alias: _runnerAlias, + enrollmentToken: _token, + ), + server: ServerConnectionConfig(url: 'http://$serverAddr'), + runtime: const AgentRuntimeConfig( + installDir: '/tmp/install', + workspaceRoot: '/tmp/workspace', + logDir: '/tmp/log', + ), + ); + + final regClient = OtoServerRegistrationClient(commandTypes: ['Print']); + final session = await regClient.openSession(agentConfig); + expect( + session.result.accepted, + isTrue, + reason: 'Registration rejected', + ); + final jobs = (session as OtoServerJobSession).jobs; + + final httpClient = http.Client(); + try { + // YAML uses tag; variables must override it. + final createResp = await httpClient.post( + Uri.parse('http://$serverAddr/api/v1/jobs'), + headers: {'content-type': 'application/json'}, + body: jsonEncode({ + 'id': jobId, + 'name': 'variable merge smoke', + 'run_request': { + 'pipeline_yaml': ''' + property: + workspace: . + FLAVOR: original + commands: + - command: Print + id: print-flavor + param: + message: "" + pipeline: + id: var-pipe + workflow: + - exe: print-flavor + ''', + 'variables': {'FLAVOR': 'release'}, + 'command_types': ['Print'], + }, + }), + ); + expect( + createResp.statusCode, + equals(201), + reason: 'Create job failed: ${createResp.body}', + ); + + final claim = await jobs.claimJob( jobId: jobId, executionId: executionId, - result: BuildResult.success( - stepEvents: [ - StepEvent( - stepId: 0, - workflowIndex: 0, - stepType: 'command', - event: 'started', - timestamp: DateTime.now().toUtc().toIso8601String(), - commandId: 'smoke-build', - commandType: 'Shell', - ), - StepEvent( - stepId: 0, - workflowIndex: 0, - stepType: 'command', - event: 'completed', - timestamp: DateTime.now().toUtc().toIso8601String(), - commandId: 'smoke-build', - commandType: 'Shell', - ), - ], - ), ); - expect(report.accepted, isTrue); - expect(report.state, 'succeeded'); + expect(claim.accepted, isTrue); + expect(claim.runRequest, isNotNull); - final job = await _getJson( - httpClient, - 'http://$serverAddr/api/v1/jobs/$jobId', + final parser = RunRequestParser(workspaceRoot: '/tmp/workspace'); + final capturedLogs = []; + final executor = RemoteRunExecutor( + parser: parser, + jobClient: jobs, + onLog: capturedLogs.add, ); - expect(job['state'], 'succeeded'); - expect(job['execution_id'], executionId); - final execution = await _getJson( - httpClient, - 'http://$serverAddr/api/v1/executions/$executionId', + final buildResult = await executor.runOnce( + jobId: claim.jobId, + executionId: claim.executionId, + runRequest: claim.runRequest!, ); - expect(execution['state'], 'succeeded'); - expect(execution['job_id'], jobId); - final logs = await _getJson( - httpClient, - 'http://$serverAddr/api/v1/executions/$executionId/logs', + expect( + buildResult.success, + isTrue, + reason: 'Execution failed: ${buildResult.message}', + ); + expect( + buildResult.stepEvents, + isNotEmpty, + reason: 'Expected step events from the build', ); - final logLines = (logs['logs'] as List) - .map( - (entry) => - (entry as Map)['Line'] ?? entry['line'], - ) - .map((line) => line.toString()) - .toList(); - expect(logLines, contains('runner started smoke build')); - expect(logLines, contains('Build completed successfully.')); + // REVIEW_API2-2: assert the remote variable actually overrode the + // YAML value. After a successful build, Application.instance.property + // holds the merged, tag-replaced property map that the pipeline used + // to resolve ``. It must be the remote 'release', + // not the YAML default 'original'. + final mergedFlavor = Application.instance.property['FLAVOR']; + expect( + mergedFlavor, + equals('release'), + reason: + 'remote variable FLAVOR was not substituted; got: $mergedFlavor', + ); + expect( + mergedFlavor, + isNot(equals('original')), + reason: + 'YAML default "original" leaked through instead of the remote override', + ); + + expect( + capturedLogs.any((l) => l.contains('success=true')), + isTrue, + reason: 'Expected successful execution report; logs: $capturedLogs', + ); + } finally { + httpClient.close(); + await session.close(); + } + } finally { + serverProcess.kill(ProcessSignal.sigterm); + await serverProcess.exitCode.timeout( + const Duration(seconds: 5), + onTimeout: () { + serverProcess.kill(ProcessSignal.sigkill); + return serverProcess.exitCode; + }, + ); + await stdoutSub.cancel(); + await stderrSub.cancel(); + } + }, + timeout: const Timeout(Duration(seconds: 30)), + ); + + test( + 'remoteRunExecutor runOnce reports failure via compatibility fallback HTTP when property is invalid despite variables (REVIEW_API3-1)', + () async { + final port = await _freePort(); + final serverAddr = '$_host:$port'; + const jobId = 'bad-prop-job'; + const executionId = 'bad-prop-exec'; + + final serverProcess = await _startCore(serverAddr); + + final output = StringBuffer(); + final stdoutSub = serverProcess.stdout + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + final stderrSub = serverProcess.stderr + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + + try { + await _waitForPort(_host, port, serverProcess, output); + + final agentConfig = AgentConfig( + agent: const AgentIdentityConfig( + id: _runnerId, + alias: _runnerAlias, + enrollmentToken: _token, + ), + server: ServerConnectionConfig(url: 'http://$serverAddr'), + runtime: const AgentRuntimeConfig( + installDir: '/tmp/install', + workspaceRoot: '/tmp/workspace', + logDir: '/tmp/log', + ), + ); + + final regClient = OtoServerRegistrationClient(commandTypes: ['Print']); + final session = await regClient.openSession(agentConfig); + expect( + session.result.accepted, + isTrue, + reason: 'Registration rejected', + ); + final jobs = (session as OtoServerJobSession).jobs; + + final httpClient = http.Client(); + try { + // `property` is a list (invalid). Remote variables must NOT coerce it + // into a valid map; Application.build validation must still reject it. + final createResp = await httpClient.post( + Uri.parse('http://$serverAddr/api/v1/jobs'), + headers: {'content-type': 'application/json'}, + body: jsonEncode({ + 'id': jobId, + 'name': 'invalid property smoke', + 'run_request': { + 'pipeline_yaml': ''' + property: + - bad + commands: + - command: Print + id: print-flavor + param: + message: "" + pipeline: + id: bad-pipe + workflow: + - exe: print-flavor + ''', + 'variables': {'FLAVOR': 'release'}, + 'command_types': ['Print'], + }, + }), + ); + expect( + createResp.statusCode, + equals(201), + reason: 'Create job failed: ${createResp.body}', + ); + + final claim = await jobs.claimJob( + jobId: jobId, + executionId: executionId, + ); + expect(claim.accepted, isTrue); + expect(claim.runRequest, isNotNull); + + final parser = RunRequestParser(workspaceRoot: '/tmp/workspace'); + final executor = RemoteRunExecutor( + parser: parser, + jobClient: jobs, + onLog: (msg) {}, + ); + + final buildResult = await executor.runOnce( + jobId: claim.jobId, + executionId: claim.executionId, + runRequest: claim.runRequest!, + ); + + // Build must fail on the invalid property; the merge must not have + // rewritten it into a valid map. + expect( + buildResult.success, + isFalse, + reason: 'Expected failure from invalid property type', + ); + + // Server job should be in failed state because reportExecution ran. + final jobResp = await httpClient.get( + Uri.parse('http://$serverAddr/api/v1/jobs/$jobId'), + ); + expect(jobResp.statusCode, 200); + final jobJson = jsonDecode(jobResp.body) as Map; + expect( + jobJson['state'], + equals('failed'), + reason: + 'Server did not receive failure report; got: ${jobJson['state']}', + ); + } finally { + httpClient.close(); + await session.close(); + } + } finally { + serverProcess.kill(ProcessSignal.sigterm); + await serverProcess.exitCode.timeout( + const Duration(seconds: 5), + onTimeout: () { + serverProcess.kill(ProcessSignal.sigkill); + return serverProcess.exitCode; + }, + ); + await stdoutSub.cancel(); + await stderrSub.cancel(); + } + }, + timeout: const Timeout(Duration(seconds: 30)), + ); + + test( + 'remoteRunExecutor runOnce sends failure report via compatibility fallback HTTP when parse fails (REVIEW_API-5)', + () async { + final port = await _freePort(); + final serverAddr = '$_host:$port'; + const jobId = 'parse-fail-job'; + const executionId = 'parse-fail-exec'; + + final serverProcess = await _startCore(serverAddr); + + final output = StringBuffer(); + final stdoutSub = serverProcess.stdout + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + final stderrSub = serverProcess.stderr + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + + try { + await _waitForPort(_host, port, serverProcess, output); + + final agentConfig = AgentConfig( + agent: const AgentIdentityConfig( + id: _runnerId, + alias: _runnerAlias, + enrollmentToken: _token, + ), + server: ServerConnectionConfig(url: 'http://$serverAddr'), + runtime: const AgentRuntimeConfig( + installDir: '/tmp/install', + workspaceRoot: '/tmp/workspace', + logDir: '/tmp/log', + ), + ); + + final regClient = OtoServerRegistrationClient(commandTypes: ['Shell']); + final session = await regClient.openSession(agentConfig); + expect( + session.result.accepted, + isTrue, + reason: 'Registration rejected', + ); + final jobs = (session as OtoServerJobSession).jobs; + + final httpClient = http.Client(); + try { + // Create a job with a path that escapes the workspace root + // (the parser will throw ArgumentError, triggering REVIEW_API-5 path). + final createResp = await httpClient.post( + Uri.parse('http://$serverAddr/api/v1/jobs'), + headers: {'content-type': 'application/json'}, + body: jsonEncode({ + 'id': jobId, + 'name': 'parse fail smoke', + 'run_request': { + 'pipeline_yaml_path': '/etc/passwd', + 'command_types': ['Shell'], + }, + }), + ); + expect( + createResp.statusCode, + equals(201), + reason: 'Create job failed: ${createResp.body}', + ); + + final claim = await jobs.claimJob( + jobId: jobId, + executionId: executionId, + ); + expect(claim.accepted, isTrue); + expect(claim.runRequest, isNotNull); + + // workspaceRoot=/tmp/workspace; /etc/passwd is outside → parse error. + final parser = RunRequestParser(workspaceRoot: '/tmp/workspace'); + final executor = RemoteRunExecutor( + parser: parser, + jobClient: jobs, + onLog: (msg) {}, + ); + + final buildResult = await executor.runOnce( + jobId: claim.jobId, + executionId: claim.executionId, + runRequest: claim.runRequest!, + ); + + // Build should have failed. + expect( + buildResult.success, + isFalse, + reason: 'Expected failure from parse error', + ); + + // Server job should be in failed state because reportExecution was called. + final jobResp = await httpClient.get( + Uri.parse('http://$serverAddr/api/v1/jobs/$jobId'), + ); + expect(jobResp.statusCode, 200); + final jobJson = jsonDecode(jobResp.body) as Map; + expect( + jobJson['state'], + equals('failed'), + reason: + 'Server did not receive failure report; got: ${jobJson['state']}', + ); + } finally { + httpClient.close(); + await session.close(); + } + } finally { + serverProcess.kill(ProcessSignal.sigterm); + await serverProcess.exitCode.timeout( + const Duration(seconds: 5), + onTimeout: () { + serverProcess.kill(ProcessSignal.sigkill); + return serverProcess.exitCode; + }, + ); + await stdoutSub.cancel(); + await stderrSub.cancel(); + } + }, + timeout: const Timeout(Duration(seconds: 30)), + ); + + test( + 'remoteRunExecutor runOnce reports declared artifacts via compatibility fallback HTTP (API-3)', + () async { + final port = await _freePort(); + final serverAddr = '$_host:$port'; + const jobId = 'artifact-smoke-job'; + const executionId = 'artifact-smoke-exec'; + + final serverProcess = await _startCore(serverAddr); + + final output = StringBuffer(); + final stdoutSub = serverProcess.stdout + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + final stderrSub = serverProcess.stderr + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + + try { + await _waitForPort(_host, port, serverProcess, output); + + final agentConfig = AgentConfig( + agent: const AgentIdentityConfig( + id: _runnerId, + alias: _runnerAlias, + enrollmentToken: _token, + ), + server: ServerConnectionConfig(url: 'http://$serverAddr'), + runtime: const AgentRuntimeConfig( + installDir: '/tmp/install', + workspaceRoot: '/tmp/workspace', + logDir: '/tmp/log', + ), + ); + + final regClient = OtoServerRegistrationClient(commandTypes: ['Print']); + final session = await regClient.openSession(agentConfig); + expect( + session.result.accepted, + isTrue, + reason: 'Registration rejected', + ); + final jobs = (session as OtoServerJobSession).jobs; + + final httpClient = http.Client(); + try { + // The pipeline declares an artifact via property.artifacts; the runner + // must report that metadata to the server after the execution report. + final createResp = await httpClient.post( + Uri.parse('http://$serverAddr/api/v1/jobs'), + headers: {'content-type': 'application/json'}, + body: jsonEncode({ + 'id': jobId, + 'name': 'artifact event smoke', + 'run_request': { + 'pipeline_yaml': ''' + property: + workspace: . + artifacts: + - name: smoke-report + path: dist/report.json + commands: + - command: Print + id: say-hello + param: + message: hello from artifact smoke + pipeline: + id: artifact-pipe + workflow: + - exe: say-hello + ''', + 'command_types': ['Print'], + }, + }), + ); + expect( + createResp.statusCode, + equals(201), + reason: 'Create job failed: ${createResp.body}', + ); + + final claim = await jobs.claimJob( + jobId: jobId, + executionId: executionId, + ); + expect(claim.accepted, isTrue); + expect(claim.runRequest, isNotNull); + + // workspaceRoot confines the declared relative artifact path. + final parser = RunRequestParser(workspaceRoot: '/tmp/workspace'); + final executor = RemoteRunExecutor( + parser: parser, + jobClient: jobs, + onLog: (msg) {}, + ); + + final buildResult = await executor.runOnce( + jobId: claim.jobId, + executionId: claim.executionId, + runRequest: claim.runRequest!, + ); + + expect( + buildResult.success, + isTrue, + reason: 'Execution failed: ${buildResult.message}', + ); + expect(buildResult.artifacts, hasLength(1)); + expect(buildResult.artifacts.first.name, 'smoke-report'); + + // The declared artifact must be persisted on the server's artifacts + // endpoint by the remote loop, not only carried in the build result. final artifacts = await _getJson( httpClient, 'http://$serverAddr/api/v1/executions/$executionId/artifacts', @@ -490,1193 +1536,153 @@ fi expect(artifactRows, hasLength(1)); final artifact = artifactRows.single as Map; expect(artifact['Name'] ?? artifact['name'], 'smoke-report'); - expect( - artifact['Path'] ?? artifact['path'], - '/tmp/oto-smoke/report.txt', - ); + expect(artifact['Path'] ?? artifact['path'], 'dist/report.json'); } finally { httpClient.close(); + await session.close(); } } finally { - await session.close(); + serverProcess.kill(ProcessSignal.sigterm); + await serverProcess.exitCode.timeout( + const Duration(seconds: 5), + onTimeout: () { + serverProcess.kill(ProcessSignal.sigkill); + return serverProcess.exitCode; + }, + ); + await stdoutSub.cancel(); + await stderrSub.cancel(); } - } finally { - serverProcess.kill(ProcessSignal.sigterm); - await serverProcess.exitCode.timeout( - const Duration(seconds: 5), - onTimeout: () { - serverProcess.kill(ProcessSignal.sigkill); - return serverProcess.exitCode; - }, - ); - await stdoutSub.cancel(); - await stderrSub.cancel(); - } - }, - timeout: const Timeout(Duration(seconds: 30)), - ); + }, + timeout: const Timeout(Duration(seconds: 30)), + ); - test( - 'remoteRunExecutor runOnce executes job via Go OTO Server and reports step events (compatibility fallback HTTP)', - () async { - final port = await _freePort(); - final serverAddr = '$_host:$port'; - const jobId = 'remote-smoke-job'; - const executionId = 'remote-smoke-exec'; + test( + 'OTO Server cancel, status, and self-update via compatibility fallback HTTP', + () async { + final port = await _freePort(); + final serverAddr = '$_host:$port'; + const jobId = 'smoke-job-cancel'; + const executionId = 'smoke-exec-cancel'; - final serverProcess = await _startCore(serverAddr); + final serverProcess = await _startCore(serverAddr); - final output = StringBuffer(); - final stdoutSub = serverProcess.stdout - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - final stderrSub = serverProcess.stderr - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); + final output = StringBuffer(); + final stdoutSub = serverProcess.stdout + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); + final stderrSub = serverProcess.stderr + .transform(systemEncoding.decoder) + .listen(output.write, onError: output.write); - try { - await _waitForPort(_host, port, serverProcess, output); - - final agentConfig = AgentConfig( - agent: const AgentIdentityConfig( - id: _runnerId, - alias: _runnerAlias, - enrollmentToken: _token, - ), - server: ServerConnectionConfig(url: 'http://$serverAddr'), - runtime: const AgentRuntimeConfig( - installDir: '/tmp/install', - workspaceRoot: '/tmp/workspace', - logDir: '/tmp/log', - ), - ); - - final regClient = OtoServerRegistrationClient( - commandTypes: ['Shell', 'Git'], - ); - final session = await regClient.openSession(agentConfig); - expect( - session.result.accepted, - isTrue, - reason: 'Registration was rejected', - ); - final jobs = (session as OtoServerJobSession).jobs; - - final httpClient = http.Client(); try { - // Create a job with valid pipeline YAML - final createResp = await httpClient.post( - Uri.parse('http://$serverAddr/api/v1/jobs'), - headers: {'content-type': 'application/json'}, - body: jsonEncode({ - 'id': jobId, - 'name': 'remote run smoke', - 'run_request': { - 'pipeline_yaml': ''' -property: - workspace: . -commands: - - command: Print - id: say-hello - param: - message: hello from smoke -pipeline: - id: smoke-pipe - workflow: - - exe: say-hello -''', - 'command_types': ['Print'], - }, - }), - ); - expect( - createResp.statusCode, - equals(201), - reason: 'Create job failed: ${createResp.body}', - ); - - final claim = await jobs.claimJob( - jobId: jobId, - executionId: executionId, - ); - expect(claim.accepted, isTrue); - expect(claim.runRequest, isNotNull); - - // Execute via RemoteRunExecutor.runOnce - final parser = RunRequestParser(workspaceRoot: '/tmp/workspace'); - final executor = RemoteRunExecutor( - parser: parser, - jobClient: jobs, - onLog: (msg) {}, - ); - - final buildResult = await executor.runOnce( - jobId: claim.jobId, - executionId: claim.executionId, - runRequest: claim.runRequest!, - ); - - expect( - buildResult.success, - isTrue, - reason: 'Execution failed: ${buildResult.message}', - ); - expect(buildResult.exitCode, 0); - expect(buildResult.stepEvents, isNotEmpty); - - final events = buildResult.stepEvents.map((e) => e.event).toList(); - expect(events, contains('started')); - expect(events, contains('completed')); - - // Verify logs on server - final logsResp = await httpClient.get( - Uri.parse('http://$serverAddr/api/v1/executions/$executionId/logs'), - ); - final logsJson = jsonDecode(logsResp.body) as Map; - final logsList = logsJson['logs'] as List; - final lines = logsList - .map((e) => (e as Map)['Line'] ?? e['line']) - .map((l) => l.toString()) - .toList(); - expect(lines, contains('run started for job $jobId')); - expect(lines, contains('Build completed successfully.')); - - // Verify execution state - final execResp = await httpClient.get( - Uri.parse('http://$serverAddr/api/v1/executions/$executionId'), - ); - expect(execResp.statusCode, 200); - final execJson = jsonDecode(execResp.body) as Map; - expect(execJson['state'], equals('succeeded')); - expect(execJson['job_id'], equals(jobId)); - - httpClient.close(); - } finally { - await session.close(); - } - } finally { - serverProcess.kill(ProcessSignal.sigterm); - await serverProcess.exitCode.timeout( - const Duration(seconds: 5), - onTimeout: () { - serverProcess.kill(ProcessSignal.sigkill); - return serverProcess.exitCode; - }, - ); - await stdoutSub.cancel(); - await stderrSub.cancel(); - } - }, - timeout: const Timeout(Duration(seconds: 30)), - ); - - test( - 'remoteRunExecutor runOnce merges variables into pipeline property (compatibility fallback HTTP) (REVIEW_API-3)', - () async { - final port = await _freePort(); - final serverAddr = '$_host:$port'; - const jobId = 'var-merge-job'; - const executionId = 'var-merge-exec'; - - final serverProcess = await _startCore(serverAddr); - - final output = StringBuffer(); - final stdoutSub = serverProcess.stdout - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - final stderrSub = serverProcess.stderr - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - - try { - await _waitForPort(_host, port, serverProcess, output); - - final agentConfig = AgentConfig( - agent: const AgentIdentityConfig( - id: _runnerId, - alias: _runnerAlias, - enrollmentToken: _token, - ), - server: ServerConnectionConfig(url: 'http://$serverAddr'), - runtime: const AgentRuntimeConfig( - installDir: '/tmp/install', - workspaceRoot: '/tmp/workspace', - logDir: '/tmp/log', - ), - ); - - final regClient = OtoServerRegistrationClient(commandTypes: ['Print']); - final session = await regClient.openSession(agentConfig); - expect( - session.result.accepted, - isTrue, - reason: 'Registration rejected', - ); - final jobs = (session as OtoServerJobSession).jobs; - - final httpClient = http.Client(); - try { - // YAML uses tag; variables must override it. - final createResp = await httpClient.post( - Uri.parse('http://$serverAddr/api/v1/jobs'), - headers: {'content-type': 'application/json'}, - body: jsonEncode({ - 'id': jobId, - 'name': 'variable merge smoke', - 'run_request': { - 'pipeline_yaml': ''' -property: - workspace: . - FLAVOR: original -commands: - - command: Print - id: print-flavor - param: - message: "" -pipeline: - id: var-pipe - workflow: - - exe: print-flavor -''', - 'variables': {'FLAVOR': 'release'}, - 'command_types': ['Print'], - }, - }), - ); - expect( - createResp.statusCode, - equals(201), - reason: 'Create job failed: ${createResp.body}', - ); - - final claim = await jobs.claimJob( - jobId: jobId, - executionId: executionId, - ); - expect(claim.accepted, isTrue); - expect(claim.runRequest, isNotNull); - - final parser = RunRequestParser(workspaceRoot: '/tmp/workspace'); - final capturedLogs = []; - final executor = RemoteRunExecutor( - parser: parser, - jobClient: jobs, - onLog: capturedLogs.add, - ); - - final buildResult = await executor.runOnce( - jobId: claim.jobId, - executionId: claim.executionId, - runRequest: claim.runRequest!, - ); - - expect( - buildResult.success, - isTrue, - reason: 'Execution failed: ${buildResult.message}', - ); - expect( - buildResult.stepEvents, - isNotEmpty, - reason: 'Expected step events from the build', - ); - - // REVIEW_API2-2: assert the remote variable actually overrode the - // YAML value. After a successful build, Application.instance.property - // holds the merged, tag-replaced property map that the pipeline used - // to resolve ``. It must be the remote 'release', - // not the YAML default 'original'. - final mergedFlavor = Application.instance.property['FLAVOR']; - expect( - mergedFlavor, - equals('release'), - reason: - 'remote variable FLAVOR was not substituted; got: $mergedFlavor', - ); - expect( - mergedFlavor, - isNot(equals('original')), - reason: - 'YAML default "original" leaked through instead of the remote override', - ); - - expect( - capturedLogs.any((l) => l.contains('success=true')), - isTrue, - reason: 'Expected successful execution report; logs: $capturedLogs', - ); - } finally { - httpClient.close(); - await session.close(); - } - } finally { - serverProcess.kill(ProcessSignal.sigterm); - await serverProcess.exitCode.timeout( - const Duration(seconds: 5), - onTimeout: () { - serverProcess.kill(ProcessSignal.sigkill); - return serverProcess.exitCode; - }, - ); - await stdoutSub.cancel(); - await stderrSub.cancel(); - } - }, - timeout: const Timeout(Duration(seconds: 30)), - ); - - test( - 'remoteRunExecutor runOnce reports failure via compatibility fallback HTTP when property is invalid despite variables (REVIEW_API3-1)', - () async { - final port = await _freePort(); - final serverAddr = '$_host:$port'; - const jobId = 'bad-prop-job'; - const executionId = 'bad-prop-exec'; - - final serverProcess = await _startCore(serverAddr); - - final output = StringBuffer(); - final stdoutSub = serverProcess.stdout - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - final stderrSub = serverProcess.stderr - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - - try { - await _waitForPort(_host, port, serverProcess, output); - - final agentConfig = AgentConfig( - agent: const AgentIdentityConfig( - id: _runnerId, - alias: _runnerAlias, - enrollmentToken: _token, - ), - server: ServerConnectionConfig(url: 'http://$serverAddr'), - runtime: const AgentRuntimeConfig( - installDir: '/tmp/install', - workspaceRoot: '/tmp/workspace', - logDir: '/tmp/log', - ), - ); - - final regClient = OtoServerRegistrationClient(commandTypes: ['Print']); - final session = await regClient.openSession(agentConfig); - expect( - session.result.accepted, - isTrue, - reason: 'Registration rejected', - ); - final jobs = (session as OtoServerJobSession).jobs; - - final httpClient = http.Client(); - try { - // `property` is a list (invalid). Remote variables must NOT coerce it - // into a valid map; Application.build validation must still reject it. - final createResp = await httpClient.post( - Uri.parse('http://$serverAddr/api/v1/jobs'), - headers: {'content-type': 'application/json'}, - body: jsonEncode({ - 'id': jobId, - 'name': 'invalid property smoke', - 'run_request': { - 'pipeline_yaml': ''' -property: - - bad -commands: - - command: Print - id: print-flavor - param: - message: "" -pipeline: - id: bad-pipe - workflow: - - exe: print-flavor -''', - 'variables': {'FLAVOR': 'release'}, - 'command_types': ['Print'], - }, - }), - ); - expect( - createResp.statusCode, - equals(201), - reason: 'Create job failed: ${createResp.body}', - ); - - final claim = await jobs.claimJob( - jobId: jobId, - executionId: executionId, - ); - expect(claim.accepted, isTrue); - expect(claim.runRequest, isNotNull); - - final parser = RunRequestParser(workspaceRoot: '/tmp/workspace'); - final executor = RemoteRunExecutor( - parser: parser, - jobClient: jobs, - onLog: (msg) {}, - ); - - final buildResult = await executor.runOnce( - jobId: claim.jobId, - executionId: claim.executionId, - runRequest: claim.runRequest!, - ); - - // Build must fail on the invalid property; the merge must not have - // rewritten it into a valid map. - expect( - buildResult.success, - isFalse, - reason: 'Expected failure from invalid property type', - ); - - // Server job should be in failed state because reportExecution ran. - final jobResp = await httpClient.get( - Uri.parse('http://$serverAddr/api/v1/jobs/$jobId'), - ); - expect(jobResp.statusCode, 200); - final jobJson = jsonDecode(jobResp.body) as Map; - expect( - jobJson['state'], - equals('failed'), - reason: - 'Server did not receive failure report; got: ${jobJson['state']}', - ); - } finally { - httpClient.close(); - await session.close(); - } - } finally { - serverProcess.kill(ProcessSignal.sigterm); - await serverProcess.exitCode.timeout( - const Duration(seconds: 5), - onTimeout: () { - serverProcess.kill(ProcessSignal.sigkill); - return serverProcess.exitCode; - }, - ); - await stdoutSub.cancel(); - await stderrSub.cancel(); - } - }, - timeout: const Timeout(Duration(seconds: 30)), - ); - - test( - 'remoteRunExecutor runOnce sends failure report via compatibility fallback HTTP when parse fails (REVIEW_API-5)', - () async { - final port = await _freePort(); - final serverAddr = '$_host:$port'; - const jobId = 'parse-fail-job'; - const executionId = 'parse-fail-exec'; - - final serverProcess = await _startCore(serverAddr); - - final output = StringBuffer(); - final stdoutSub = serverProcess.stdout - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - final stderrSub = serverProcess.stderr - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - - try { - await _waitForPort(_host, port, serverProcess, output); - - final agentConfig = AgentConfig( - agent: const AgentIdentityConfig( - id: _runnerId, - alias: _runnerAlias, - enrollmentToken: _token, - ), - server: ServerConnectionConfig(url: 'http://$serverAddr'), - runtime: const AgentRuntimeConfig( - installDir: '/tmp/install', - workspaceRoot: '/tmp/workspace', - logDir: '/tmp/log', - ), - ); - - final regClient = OtoServerRegistrationClient(commandTypes: ['Shell']); - final session = await regClient.openSession(agentConfig); - expect( - session.result.accepted, - isTrue, - reason: 'Registration rejected', - ); - final jobs = (session as OtoServerJobSession).jobs; - - final httpClient = http.Client(); - try { - // Create a job with a path that escapes the workspace root - // (the parser will throw ArgumentError, triggering REVIEW_API-5 path). - final createResp = await httpClient.post( - Uri.parse('http://$serverAddr/api/v1/jobs'), - headers: {'content-type': 'application/json'}, - body: jsonEncode({ - 'id': jobId, - 'name': 'parse fail smoke', - 'run_request': { - 'pipeline_yaml_path': '/etc/passwd', - 'command_types': ['Shell'], - }, - }), - ); - expect( - createResp.statusCode, - equals(201), - reason: 'Create job failed: ${createResp.body}', - ); - - final claim = await jobs.claimJob( - jobId: jobId, - executionId: executionId, - ); - expect(claim.accepted, isTrue); - expect(claim.runRequest, isNotNull); - - // workspaceRoot=/tmp/workspace; /etc/passwd is outside → parse error. - final parser = RunRequestParser(workspaceRoot: '/tmp/workspace'); - final executor = RemoteRunExecutor( - parser: parser, - jobClient: jobs, - onLog: (msg) {}, - ); - - final buildResult = await executor.runOnce( - jobId: claim.jobId, - executionId: claim.executionId, - runRequest: claim.runRequest!, - ); - - // Build should have failed. - expect( - buildResult.success, - isFalse, - reason: 'Expected failure from parse error', - ); - - // Server job should be in failed state because reportExecution was called. - final jobResp = await httpClient.get( - Uri.parse('http://$serverAddr/api/v1/jobs/$jobId'), - ); - expect(jobResp.statusCode, 200); - final jobJson = jsonDecode(jobResp.body) as Map; - expect( - jobJson['state'], - equals('failed'), - reason: - 'Server did not receive failure report; got: ${jobJson['state']}', - ); - } finally { - httpClient.close(); - await session.close(); - } - } finally { - serverProcess.kill(ProcessSignal.sigterm); - await serverProcess.exitCode.timeout( - const Duration(seconds: 5), - onTimeout: () { - serverProcess.kill(ProcessSignal.sigkill); - return serverProcess.exitCode; - }, - ); - await stdoutSub.cancel(); - await stderrSub.cancel(); - } - }, - timeout: const Timeout(Duration(seconds: 30)), - ); - - test( - 'remoteRunExecutor runOnce reports declared artifacts via compatibility fallback HTTP (API-3)', - () async { - final port = await _freePort(); - final serverAddr = '$_host:$port'; - const jobId = 'artifact-smoke-job'; - const executionId = 'artifact-smoke-exec'; - - final serverProcess = await _startCore(serverAddr); - - final output = StringBuffer(); - final stdoutSub = serverProcess.stdout - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - final stderrSub = serverProcess.stderr - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - - try { - await _waitForPort(_host, port, serverProcess, output); - - final agentConfig = AgentConfig( - agent: const AgentIdentityConfig( - id: _runnerId, - alias: _runnerAlias, - enrollmentToken: _token, - ), - server: ServerConnectionConfig(url: 'http://$serverAddr'), - runtime: const AgentRuntimeConfig( - installDir: '/tmp/install', - workspaceRoot: '/tmp/workspace', - logDir: '/tmp/log', - ), - ); - - final regClient = OtoServerRegistrationClient(commandTypes: ['Print']); - final session = await regClient.openSession(agentConfig); - expect( - session.result.accepted, - isTrue, - reason: 'Registration rejected', - ); - final jobs = (session as OtoServerJobSession).jobs; - - final httpClient = http.Client(); - try { - // The pipeline declares an artifact via property.artifacts; the runner - // must report that metadata to the server after the execution report. - final createResp = await httpClient.post( - Uri.parse('http://$serverAddr/api/v1/jobs'), - headers: {'content-type': 'application/json'}, - body: jsonEncode({ - 'id': jobId, - 'name': 'artifact event smoke', - 'run_request': { - 'pipeline_yaml': ''' -property: - workspace: . - artifacts: - - name: smoke-report - path: dist/report.json -commands: - - command: Print - id: say-hello - param: - message: hello from artifact smoke -pipeline: - id: artifact-pipe - workflow: - - exe: say-hello -''', - 'command_types': ['Print'], - }, - }), - ); - expect( - createResp.statusCode, - equals(201), - reason: 'Create job failed: ${createResp.body}', - ); - - final claim = await jobs.claimJob( - jobId: jobId, - executionId: executionId, - ); - expect(claim.accepted, isTrue); - expect(claim.runRequest, isNotNull); - - // workspaceRoot confines the declared relative artifact path. - final parser = RunRequestParser(workspaceRoot: '/tmp/workspace'); - final executor = RemoteRunExecutor( - parser: parser, - jobClient: jobs, - onLog: (msg) {}, - ); - - final buildResult = await executor.runOnce( - jobId: claim.jobId, - executionId: claim.executionId, - runRequest: claim.runRequest!, - ); - - expect( - buildResult.success, - isTrue, - reason: 'Execution failed: ${buildResult.message}', - ); - expect(buildResult.artifacts, hasLength(1)); - expect(buildResult.artifacts.first.name, 'smoke-report'); - - // The declared artifact must be persisted on the server's artifacts - // endpoint by the remote loop, not only carried in the build result. - final artifacts = await _getJson( - httpClient, - 'http://$serverAddr/api/v1/executions/$executionId/artifacts', - ); - final artifactRows = artifacts['artifacts'] as List; - expect(artifactRows, hasLength(1)); - final artifact = artifactRows.single as Map; - expect(artifact['Name'] ?? artifact['name'], 'smoke-report'); - expect(artifact['Path'] ?? artifact['path'], 'dist/report.json'); - } finally { - httpClient.close(); - await session.close(); - } - } finally { - serverProcess.kill(ProcessSignal.sigterm); - await serverProcess.exitCode.timeout( - const Duration(seconds: 5), - onTimeout: () { - serverProcess.kill(ProcessSignal.sigkill); - return serverProcess.exitCode; - }, - ); - await stdoutSub.cancel(); - await stderrSub.cancel(); - } - }, - timeout: const Timeout(Duration(seconds: 30)), - ); - - test( - 'OTO Server cancel, status, and self-update via compatibility fallback HTTP', - () async { - final port = await _freePort(); - final serverAddr = '$_host:$port'; - const jobId = 'smoke-job-cancel'; - const executionId = 'smoke-exec-cancel'; - - final serverProcess = await _startCore(serverAddr); - - final output = StringBuffer(); - final stdoutSub = serverProcess.stdout - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - final stderrSub = serverProcess.stderr - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - - try { - await _waitForPort(_host, port, serverProcess, output); - - final agentConfig = AgentConfig( - agent: const AgentIdentityConfig( - id: _runnerId, - alias: _runnerAlias, - enrollmentToken: _token, - ), - server: ServerConnectionConfig(url: 'http://$serverAddr'), - runtime: const AgentRuntimeConfig( - installDir: '/tmp/install', - workspaceRoot: '/tmp/workspace', - logDir: '/tmp/log', - ), - ); - - final regClient = OtoServerRegistrationClient( - commandTypes: ['Shell', 'Git'], - ); - final session = await regClient.openSession(agentConfig); - expect(session.result.accepted, isTrue); - final jobs = (session as OtoServerJobSession).jobs; - - final httpClient = http.Client(); - try { - // 1. Create a job - final createResp = await httpClient.post( - Uri.parse('http://$serverAddr/api/v1/jobs'), - headers: {'content-type': 'application/json'}, - body: jsonEncode({ - 'id': jobId, - 'name': 'smoke build', - 'run_request': { - 'pipeline_yaml': 'commands:\n - type: Shell', - 'command_types': ['Shell'], - }, - }), - ); - expect(createResp.statusCode, equals(201)); - - // 2. Claim job - final claim = await jobs.claimJob( - jobId: jobId, - executionId: executionId, - ); - expect(claim.accepted, isTrue); - - // 3. Check status (should be running, and current execution should be set) - final statusRes = await jobs.fetchRunnerStatus(); - expect(statusRes.accepted, isTrue); - expect(statusRes.runnerId, _runnerId); - expect(statusRes.currentExecutionId, executionId); - expect(statusRes.currentJobId, jobId); - - // 4. Request self-update (should be deferred because execution is running) - final updateRes = await jobs.requestSelfUpdate( - version: 'v2.0.0', - downloadUrl: 'https://example.com/binary', - ); - expect(updateRes.accepted, isFalse); - expect(updateRes.deferred, isTrue); - - // 5. Cancel running job - final cancelRes = await jobs.cancelRun( - executionId: executionId, - reason: 'smoke testing cancel', - ); - expect(cancelRes.success, isTrue); - - // Verify execution is canceled in store - final execResp = await httpClient.get( - Uri.parse('http://$serverAddr/api/v1/executions/$executionId'), - ); - final execJson = jsonDecode(execResp.body) as Map; - expect(execJson['state'], equals('canceled')); - - // 6. Request self-update again (should be accepted now that execution is canceled/idle) - final updateRes2 = await jobs.requestSelfUpdate( - version: 'v2.0.0', - downloadUrl: 'https://example.com/binary', - ); - expect(updateRes2.accepted, isTrue); - expect(updateRes2.deferred, isFalse); - } finally { - httpClient.close(); - await session.close(); - } - } finally { - serverProcess.kill(ProcessSignal.sigterm); - await serverProcess.exitCode.timeout( - const Duration(seconds: 5), - onTimeout: () { - serverProcess.kill(ProcessSignal.sigkill); - return serverProcess.exitCode; - }, - ); - await stdoutSub.cancel(); - await stderrSub.cancel(); - } - }, - timeout: const Timeout(Duration(seconds: 30)), - ); - - test( - 'OtoServerSocketRegistrationClient socket session sends heartbeat and disconnects on close', - () async { - final port = await _freePort(); - final serverAddr = '$_host:$port'; - - final (:process, :socketAddr) = await _startCoreWithSocket(serverAddr); - - final output = StringBuffer(); - final stdoutSub = process.stdout - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - final stderrSub = process.stderr - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - - try { - await _waitForPort(_host, port, process, output); - - final agentConfig = AgentConfig( - agent: const AgentIdentityConfig( - id: _runnerId, - alias: _runnerAlias, - enrollmentToken: _token, - ), - server: ServerConnectionConfig( - url: 'http://$serverAddr', - socketUrl: 'tcp://$socketAddr', - ), - runtime: const AgentRuntimeConfig( - installDir: '/tmp/install', - workspaceRoot: '/tmp/workspace', - logDir: '/tmp/log', - ), - ); - - final client = OtoServerSocketRegistrationClient( - commandTypes: ['Shell', 'Git'], - heartbeatInterval: const Duration(milliseconds: 200), - ); - - final session = await client.openSession(agentConfig); - expect(session, isA()); - expect(session.result.accepted, isTrue); - expect(session.result.runnerId, _runnerId); - - // Capture stream completion after close(). - final pushSession = session as OtoServerPushJobSession; - final runRequestsDone = Completer(); - final cancelRequestsDone = Completer(); - final runSub = pushSession.runRequests.listen( - (_) {}, - onDone: runRequestsDone.complete, - ); - final cancelSub = pushSession.cancelRequests.listen( - (_) {}, - onDone: cancelRequestsDone.complete, - ); - - final httpClient = http.Client(); - try { - final isOnline = await _pollRunnerStatus( - httpClient, - serverAddr, - _runnerId, - 'online', - ); - expect( - isOnline, - isTrue, - reason: 'socket runner did not go online in registry', - ); - - // close() must cancel the heartbeat timer and close socket streams. - await session.close(); - await runRequestsDone.future.timeout( - const Duration(seconds: 2), - onTimeout: () => - fail('runRequests stream did not close after session.close()'), - ); - await cancelRequestsDone.future.timeout( - const Duration(seconds: 2), - onTimeout: () => fail( - 'cancelRequests stream did not close after session.close()', + await _waitForPort(_host, port, serverProcess, output); + + final agentConfig = AgentConfig( + agent: const AgentIdentityConfig( + id: _runnerId, + alias: _runnerAlias, + enrollmentToken: _token, + ), + server: ServerConnectionConfig(url: 'http://$serverAddr'), + runtime: const AgentRuntimeConfig( + installDir: '/tmp/install', + workspaceRoot: '/tmp/workspace', + logDir: '/tmp/log', ), ); - await runSub.cancel(); - await cancelSub.cancel(); - final isDisconnected = await _pollRunnerStatus( - httpClient, - serverAddr, - _runnerId, - 'disconnected', - timeout: const Duration(seconds: 5), + final regClient = OtoServerRegistrationClient( + commandTypes: ['Shell', 'Git'], ); - expect( - isDisconnected, - isTrue, - reason: 'socket runner did not transition to disconnected', - ); - } finally { - httpClient.close(); - } - } finally { - process.kill(ProcessSignal.sigterm); - await process.exitCode.timeout( - const Duration(seconds: 5), - onTimeout: () { - process.kill(ProcessSignal.sigkill); - return process.exitCode; - }, - ); - await stdoutSub.cancel(); - await stderrSub.cancel(); - } - }, - timeout: const Timeout(Duration(seconds: 30)), - ); + final session = await regClient.openSession(agentConfig); + expect(session.result.accepted, isTrue); + final jobs = (session as OtoServerJobSession).jobs; - test( - 'OtoServerSocketRegistrationClient receives socket cancel requests from HTTP cancel action', - () async { - final port = await _freePort(); - final serverAddr = '$_host:$port'; - const jobId = 'smoke-job-socket-cancel'; + final httpClient = http.Client(); + try { + // 1. Create a job + final createResp = await httpClient.post( + Uri.parse('http://$serverAddr/api/v1/jobs'), + headers: {'content-type': 'application/json'}, + body: jsonEncode({ + 'id': jobId, + 'name': 'smoke build', + 'run_request': { + 'pipeline_yaml': 'commands:\n - type: Shell', + 'command_types': ['Shell'], + }, + }), + ); + expect(createResp.statusCode, equals(201)); - final (:process, :socketAddr) = await _startCoreWithSocket(serverAddr); + // 2. Claim job + final claim = await jobs.claimJob( + jobId: jobId, + executionId: executionId, + ); + expect(claim.accepted, isTrue); - final output = StringBuffer(); - final stdoutSub = process.stdout - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - final stderrSub = process.stderr - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); + // 3. Check status (should be running, and current execution should be set) + final statusRes = await jobs.fetchRunnerStatus(); + expect(statusRes.accepted, isTrue); + expect(statusRes.runnerId, _runnerId); + expect(statusRes.currentExecutionId, executionId); + expect(statusRes.currentJobId, jobId); - try { - await _waitForPort(_host, port, process, output); + // 4. Request self-update (should be deferred because execution is running) + final updateRes = await jobs.requestSelfUpdate( + version: 'v2.0.0', + downloadUrl: 'https://example.com/binary', + ); + expect(updateRes.accepted, isFalse); + expect(updateRes.deferred, isTrue); - final agentConfig = AgentConfig( - agent: const AgentIdentityConfig( - id: _runnerId, - alias: _runnerAlias, - enrollmentToken: _token, - ), - server: ServerConnectionConfig( - url: 'http://$serverAddr', - socketUrl: 'tcp://$socketAddr', - ), - runtime: const AgentRuntimeConfig( - installDir: '/tmp/install', - workspaceRoot: '/tmp/workspace', - logDir: '/tmp/log', - ), - ); + // 5. Cancel running job + final cancelRes = await jobs.cancelRun( + executionId: executionId, + reason: 'smoke testing cancel', + ); + expect(cancelRes.success, isTrue); - final client = OtoServerSocketRegistrationClient( - commandTypes: ['Shell'], - heartbeatInterval: const Duration(milliseconds: 200), - ); - final session = await client.openSession(agentConfig); - expect(session.result.accepted, isTrue); - final pushSession = session as OtoServerPushJobSession; + // Verify execution is canceled in store + final execResp = await httpClient.get( + Uri.parse('http://$serverAddr/api/v1/executions/$executionId'), + ); + final execJson = jsonDecode(execResp.body) as Map; + expect(execJson['state'], equals('canceled')); - final runReceived = Completer(); - final runSub = pushSession.runRequests.listen((request) { - if (!runReceived.isCompleted) { - expect(request.runnerId, _runnerId); - expect(request.jobId, jobId); - expect(request.executionId, isNotEmpty); - runReceived.complete(request.executionId); + // 6. Request self-update again (should be accepted now that execution is canceled/idle) + final updateRes2 = await jobs.requestSelfUpdate( + version: 'v2.0.0', + downloadUrl: 'https://example.com/binary', + ); + expect(updateRes2.accepted, isTrue); + expect(updateRes2.deferred, isFalse); + } finally { + httpClient.close(); + await session.close(); } - }); - var expectedExecutionId = ''; - final cancelReceived = Completer(); - final cancelSub = pushSession.cancelRequests.listen((request) { - if (!cancelReceived.isCompleted) { - expect(request.runnerId, _runnerId); - expect(request.executionId, expectedExecutionId); - expect(request.reason, 'socket cancel smoke'); - cancelReceived.complete(request.executionId); - } - }); - - final httpClient = http.Client(); - final jobs = OtoServerJobClient( - serverUrl: 'http://$serverAddr', - runnerId: _runnerId, - client: httpClient, - ); - try { - final createResp = await httpClient.post( - Uri.parse('http://$serverAddr/api/v1/jobs'), - headers: {'content-type': 'application/json'}, - body: jsonEncode({ - 'id': jobId, - 'name': 'socket cancel smoke', - 'run_request': { - 'pipeline_yaml': 'commands:\n - type: Shell', - 'command_types': ['Shell'], - }, - }), - ); - expect(createResp.statusCode, equals(201)); - - expectedExecutionId = await runReceived.future.timeout( - const Duration(seconds: 3), - onTimeout: () => fail('socket RunRequest was not received'), - ); - - final cancelRes = await jobs.cancelRun( - executionId: expectedExecutionId, - reason: 'socket cancel smoke', - ); - expect(cancelRes.success, isTrue); - - await cancelReceived.future.timeout( - const Duration(seconds: 3), - onTimeout: () => fail('socket CancelRunRequest was not received'), - ); - - final execResp = await httpClient.get( - Uri.parse( - 'http://$serverAddr/api/v1/executions/$expectedExecutionId', - ), - ); - final execJson = jsonDecode(execResp.body) as Map; - expect(execJson['state'], equals('canceled')); } finally { - await runSub.cancel(); - await cancelSub.cancel(); - httpClient.close(); - await session.close(); + serverProcess.kill(ProcessSignal.sigterm); + await serverProcess.exitCode.timeout( + const Duration(seconds: 5), + onTimeout: () { + serverProcess.kill(ProcessSignal.sigkill); + return serverProcess.exitCode; + }, + ); + await stdoutSub.cancel(); + await stderrSub.cancel(); } - } finally { - process.kill(ProcessSignal.sigterm); - await process.exitCode.timeout( - const Duration(seconds: 5), - onTimeout: () { - process.kill(ProcessSignal.sigkill); - return process.exitCode; - }, - ); - await stdoutSub.cancel(); - await stderrSub.cancel(); - } - }, - timeout: const Timeout(Duration(seconds: 30)), - ); + }, + timeout: const Timeout(Duration(seconds: 30)), + ); - test( - 'OtoServerSocketRegistrationClient socket session close is idempotent after server disconnect', - () async { - final port = await _freePort(); - final serverAddr = '$_host:$port'; - - final (:process, :socketAddr) = await _startCoreWithSocket(serverAddr); - - final output = StringBuffer(); - final stdoutSub = process.stdout - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - final stderrSub = process.stderr - .transform(systemEncoding.decoder) - .listen(output.write, onError: output.write); - - try { - await _waitForPort(_host, port, process, output); - - final agentConfig = AgentConfig( - agent: const AgentIdentityConfig( - id: _runnerId, - alias: _runnerAlias, - enrollmentToken: _token, - ), - server: ServerConnectionConfig( - url: 'http://$serverAddr', - socketUrl: 'tcp://$socketAddr', - ), - runtime: const AgentRuntimeConfig( - installDir: '/tmp/install', - workspaceRoot: '/tmp/workspace', - logDir: '/tmp/log', - ), - ); - - final client = OtoServerSocketRegistrationClient( - commandTypes: ['Shell'], - // Short interval so in-flight heartbeats exercise the catch block. - heartbeatInterval: const Duration(milliseconds: 50), - ); - - final session = await client.openSession(agentConfig); - expect(session.result.accepted, isTrue); - - // Graceful server shutdown: closes all proto-socket clients cleanly, - // sending EOF to the Dart side. In-flight heartbeat sendRequest calls - // receive StateError('connection closed'), which is swallowed by the - // _sendHeartbeat catch block. - process.kill(ProcessSignal.sigterm); - await process.exitCode.timeout(const Duration(seconds: 5)); - - // Allow the EOF to propagate and the transport to auto-close. - await Future.delayed(const Duration(milliseconds: 300)); - - // session.close() must complete without throwing even though the - // proto-socket transport is already closed by the EOF handler. - await expectLater(session.close(), completes); - } finally { - process.kill(ProcessSignal.sigkill); - await process.exitCode.timeout( - const Duration(seconds: 2), - onTimeout: () => -1, - ); - await stdoutSub.cancel(); - await stderrSub.cancel(); - } - }, - timeout: const Timeout(Duration(seconds: 30)), - ); + }); } -// ─── Socket session lifecycle tests ────────────────────────────────────────── +// ─── Test helpers ────────────────────────────────────────────────────────────── typedef _CoreWithSocket = ({Process process, String socketAddr});