processJob function

Future<bool> processJob(
  1. ApiClient apiClient,
  2. String workerId, {
  3. void onJobFound()?,
})

Implementation

Future<bool> processJob(
  ApiClient apiClient,
  String workerId, {
  void Function()? onJobFound,
}) async {
  final buildJob = await apiClient.claimNextJob(null);
  if (buildJob == null) return false;

  onJobFound?.call();

  final buildJobId = buildJob.id;
  final runId = _uuid.v4();

  // Initialize Logger
  setLoggerApiClient(apiClient);

  await apiClient.createRun(buildJobId, runId);
  await apiClient.updateCheckRun(buildJob, 'in_progress');

  // Clean up orphaned VMs and zombie processes from previous runs to free memory/locks
  try {
    await cleanupOrphanedVms(workerId);
  } catch (e) {
    await logWarning(
      buildJobId,
      runId,
      'Failed to run VM cleanup before job: $e',
    );
  }

  await logInfo(
    buildJobId,
    runId,
    'Processing job: $buildJobId for ${buildJob.owner}/${buildJob.repo} [v$version]',
  );

  // Resolve Installation Token
  String token;
  try {
    final tokenResp = await apiClient.resolveInstallationToken(buildJobId);
    token = tokenResp['token'] as String;
  } catch (e) {
    await logError(
      buildJobId,
      runId,
      'Failed to resolve GitHub App Installation Token: $e',
    );
    await apiClient.updateRunStatus(
      buildJobId: buildJobId,
      runId: runId,
      status: 'completed',
      conclusion: 'failure',
    );
    await apiClient.completeJob(buildJobId, 'FAILURE');
    await apiClient.updateCheckRun(
      buildJob,
      'completed',
      conclusion: 'failure',
    );
    await apiClient.handleBuildJobStatusChange(buildJob, 'FAILURE');
    return true;
  }

  final owner = buildJob.owner;
  final repo = buildJob.repo;
  final commitSha = buildJob.commitSha ?? '';

  final vmName = currentVmName(workerId: workerId, buildJobId: buildJobId);
  lume.LumeVM? vm;

  Future<void> execCommand(String command) => execVmCommand(
    vmName: vmName,
    command: command,
    buildJobId: buildJobId,
    runId: runId,
    token: token,
    ipAddress: vm?.ipAddress,
  );

  Future<bool> isCancelled() async {
    try {
      return await apiClient.isJobCancelled(buildJobId);
    } catch (_) {
      return false;
    }
  }

  try {
    final workflowFileName = buildJob.workflowFileName;
    if (workflowFileName == null || workflowFileName.isEmpty) {
      throw Exception('workflowFileName is missing');
    }

    await logInfo(buildJobId, runId, 'Workflow: $workflowFileName');

    await cloneVm(
      baseVmName: baseVmName,
      vmName: vmName,
      buildJobId: buildJobId,
      runId: runId,
      workerId: workerId,
    );

    await logInfo(buildJobId, runId, 'Booting macOS VM via Lume...');
    vm = await runVm(vmName);
    await logInfo(buildJobId, runId, 'VM booted successfully!');

    await setupDirectSsh(vm);
    final vmIp = vm.ipAddress;
    await logInfo(buildJobId, runId, 'VM IP: $vmIp. VM is ready!');

    await logInfo(buildJobId, runId, 'Cloning repository $owner/$repo...');
    final githubHost = buildJob.githubBaseUrl != null
        ? Uri.parse(buildJob.githubBaseUrl!).host
        : 'github.com';
    final cloneUrl =
        'https://x-access-token:$token@$githubHost/$owner/$repo.git';

    var cloneAttempt = 0;
    await retry(
      () => execCommand('git clone --depth 1 --no-checkout $cloneUrl'),
      delayFactor: const Duration(seconds: 5),
      randomizationFactor: 0,
      maxAttempts: 3,
      onRetry: (e) {
        cloneAttempt++;
        logInfo(
          buildJobId,
          runId,
          'git clone failed (attempt $cloneAttempt/3). Retrying...',
        );
      },
    );

    final pullRequestNumber = buildJob.pullRequestNumber;

    await logInfo(buildJobId, runId, 'Fetching commit $commitSha...');
    var fetchAttempt = 0;
    await retry(
      () async {
        try {
          await execCommand('git -C $repo fetch --depth 1 origin $commitSha');
        } catch (_) {
          if (pullRequestNumber != null) {
            await logInfo(
              buildJobId,
              runId,
              'Direct fetch failed, trying PR ref pull/$pullRequestNumber/head...',
            );
            await execCommand(
              'git -C $repo fetch --depth 1 origin pull/$pullRequestNumber/head',
            );
          } else {
            rethrow;
          }
        }
      },
      delayFactor: const Duration(seconds: 5),
      randomizationFactor: 0,
      maxAttempts: 3,
      onRetry: (e) {
        fetchAttempt++;
        logInfo(
          buildJobId,
          runId,
          'git fetch failed (attempt $fetchAttempt/3). Retrying...',
        );
      },
    );

    await logInfo(buildJobId, runId, 'Checking out commit $commitSha...');
    await execCommand('git -C $repo checkout $commitSha');
    await logInfo(buildJobId, runId, 'Repository cloned successfully');

    // Build Environment variables
    final envVars = await buildEnvVars(
      apiClient: apiClient,
      buildJob: buildJob,
      projectId: apiClient.projectId,
      buildJobId: buildJobId,
      runId: runId,
    );

    // Build Secrets (filtered by workflow references)
    final secretVars = await buildSecretVars(
      apiClient: apiClient,
      token: token,
      buildJobId: buildJobId,
      runId: runId,
      buildJob: buildJob,
    );

    final envFileLines = <String>[];
    final secretFileLines = <String>[];

    for (final entry in envVars.entries) {
      final escaped = entry.value.replaceAll('\n', '\\n');
      envFileLines.add('${entry.key}=$escaped');
    }
    for (final entry in secretVars.entries) {
      final escaped = entry.value.replaceAll('\n', '\\n');
      secretFileLines.add('${entry.key}=$escaped');
    }

    final envFileContent = envFileLines.join('\n');
    final secretFileContent = secretFileLines.join('\n');

    await writeFileToVm(vmIp, '/tmp/openci-env', envFileContent);
    await writeFileToVm(vmIp, '/tmp/openci-secrets', secretFileContent);
    await writeFileToVm(
      vmIp,
      '/tmp/openci-event.json',
      buildEventPayload(buildJob),
    );
    await logInfo(buildJobId, runId, 'Environment variables written');

    await logInfo(buildJobId, runId, 'Running workflow with act...');

    final eventType = pullRequestNumber != null ? 'pull_request' : 'push';
    final jobKey = buildJob.workflowJobKey ?? buildJob.jobKey;
    final jobFlag = jobKey != null ? '-j $jobKey ' : '';

    final matrixArgs = <String>[];
    final buildJobMatrix = buildJob.matrix;
    if (buildJobMatrix != null && buildJobMatrix.isNotEmpty) {
      for (final entry in buildJobMatrix.entries) {
        matrixArgs.add('--matrix "${entry.key}:${entry.value}"');
      }
    }
    final matrixFlag = matrixArgs.isNotEmpty ? '${matrixArgs.join(' ')} ' : '';

    // Use the VM's real home. Each build runs in its own fresh VM, so a unique
    // per-run HOME is unnecessary for isolation. Critically, macOS 26 (Tahoe)
    // will NOT treat a code-signing keychain stored outside the user's real
    // home (e.g. under /tmp) as a valid signing identity, so build keychains
    // must live under /Users/admin/Library/Keychains for `find-identity -v`.
    final actScript = [
      'set -e',
      'export HOME=/Users/admin',
      'export PATH="/Users/admin/flutter/bin:/opt/homebrew/bin:\$PATH"',
      'cd $repo',
      'act $eventType -W .openci/$workflowFileName '
          '$jobFlag'
          '$matrixFlag'
          '-P macos-latest=-self-hosted '
          '-P macos-14=-self-hosted '
          '-P macos-15=-self-hosted '
          '-P ubuntu-latest=-self-hosted '
          '-e /tmp/openci-event.json '
          '--env-file /tmp/openci-env '
          '--secret-file /tmp/openci-secrets',
    ].join('\n');

    await writeFileToVm(vmIp, '/tmp/openci-act.sh', actScript);
    await execCommand('chmod +x /tmp/openci-act.sh');

    try {
      await execCommandStreaming(
        ['/bin/zsh', '-l', '/tmp/openci-act.sh'],
        vmIp,
        buildJobId,
        runId,
        token,
        isCancelled: isCancelled,
      );

      await Future.delayed(const Duration(seconds: 5));

      await logInfo(buildJobId, runId, 'Build completed successfully');
      await updateJobFinalStatus(
        apiClient: apiClient,
        buildJob: buildJob,
        runId: runId,
        status: BuildJobStatus.SUCCESS,
        conclusion: 'success',
      );
    } on TimeoutException catch (timeoutError) {
      await logError(
        buildJobId,
        runId,
        'Job execution timed out: $timeoutError',
      );
      await updateJobFinalStatus(
        apiClient: apiClient,
        buildJob: buildJob,
        runId: runId,
        status: BuildJobStatus.TIMED_OUT,
        conclusion: 'timed_out',
      );
      return true;
    } catch (actError) {
      if (await isCancelled()) {
        await logInfo(buildJobId, runId, 'Build was cancelled by user');
        await updateJobFinalStatus(
          apiClient: apiClient,
          buildJob: buildJob,
          runId: runId,
          status: BuildJobStatus.CANCELLED,
          conclusion: 'cancelled',
        );
        return true;
      }
      await logWarning(buildJobId, runId, 'Act build failed: $actError');
      await updateJobFinalStatus(
        apiClient: apiClient,
        buildJob: buildJob,
        runId: runId,
        status: BuildJobStatus.FAILURE,
        conclusion: 'failure',
      );
      return true;
    }
  } on TimeoutException catch (e, s) {
    await logError(
      buildJobId,
      runId,
      'Job timed out: $e',
      stackTrace: s.toString(),
    );
    await updateJobFinalStatus(
      apiClient: apiClient,
      buildJob: buildJob,
      runId: runId,
      status: BuildJobStatus.TIMED_OUT,
      conclusion: 'timed_out',
    );
    return true;
  } catch (e, s) {
    await logError(
      buildJobId,
      runId,
      'Job failed: $e',
      stackTrace: s.toString(),
    );
    await updateJobFinalStatus(
      apiClient: apiClient,
      buildJob: buildJob,
      runId: runId,
      status: BuildJobStatus.FAILURE,
      conclusion: 'failure',
    );
    rethrow;
  } finally {
    await flushRemainingLogs(runId: runId);
    await stopVm(vm);
    await deleteVm(vmName);
    await pruneStaleVms(buildJobId, runId, workerId: workerId);
  }

  return true;
}