diff --git a/app_dart/lib/cocoon_service.dart b/app_dart/lib/cocoon_service.dart index 6c6e325f03..6852a94737 100644 --- a/app_dart/lib/cocoon_service.dart +++ b/app_dart/lib/cocoon_service.dart @@ -28,6 +28,7 @@ export 'src/request_handlers/github/webhook_subscription.dart'; export 'src/request_handlers/github_rate_limit_status.dart'; export 'src/request_handlers/github_webhook.dart'; export 'src/request_handlers/postsubmit_luci_subscription.dart'; +export 'src/request_handlers/presubmit_guard_update_subscription.dart'; export 'src/request_handlers/presubmit_luci_subscription.dart'; export 'src/request_handlers/push_build_status_to_github.dart'; export 'src/request_handlers/push_gold_status_to_github.dart'; diff --git a/app_dart/lib/server.dart b/app_dart/lib/server.dart index 55e21b53a1..3cf395bb31 100644 --- a/app_dart/lib/server.dart +++ b/app_dart/lib/server.dart @@ -120,6 +120,12 @@ Server createServer({ ciYamlFetcher: ciYamlFetcher, firestore: firestore, ), + '/api/v2/presubmit-guard-update-subscription': + PresubmitGuardUpdateSubscription( + cache: cache, + config: config, + scheduler: scheduler, + ), '/api/v2/postsubmit-luci-subscription': PostsubmitLuciSubscription( cache: cache, config: config, diff --git a/app_dart/lib/src/request_handlers/get_engine_artifacts_ready.dart b/app_dart/lib/src/request_handlers/get_engine_artifacts_ready.dart index 6a83e6ed80..fbbf91e6c7 100644 --- a/app_dart/lib/src/request_handlers/get_engine_artifacts_ready.dart +++ b/app_dart/lib/src/request_handlers/get_engine_artifacts_ready.dart @@ -76,8 +76,10 @@ final class GetEngineArtifactsReady extends PublicApiRequestHandler { final engineGuard = guards .where((g) => g.stage == CiStage.fusionEngineBuild) .firstOrNull; - remaining = engineGuard?.remainingJobs ?? 0; - failed = engineGuard?.failedJobs ?? 0; + if (engineGuard != null) { + failed = engineGuard.failedJobs; + remaining = engineGuard.remainingJobs; + } } } diff --git a/app_dart/lib/src/request_handlers/get_presubmit_guard.dart b/app_dart/lib/src/request_handlers/get_presubmit_guard.dart index eb906fb661..d266ea68ce 100644 --- a/app_dart/lib/src/request_handlers/get_presubmit_guard.dart +++ b/app_dart/lib/src/request_handlers/get_presubmit_guard.dart @@ -104,12 +104,24 @@ final class GetPresubmitGuard extends PublicApiRequestHandler { // Consolidate metadata from the first record. final first = guards.first; - final totalFailed = guards.fold(0, (sum, g) => sum + g.failedJobs); - final totalRemaining = guards.fold( - 0, - (sum, g) => sum + g.remainingJobs, - ); - final totalBuilds = guards.fold(0, (sum, g) => sum + g.jobs.length); + var totalFailed = 0; + var totalRemaining = 0; + var totalBuilds = 0; + final stages = []; + + for (final g in guards) { + totalFailed += g.failedJobs; + totalRemaining += g.remainingJobs; + totalBuilds += g.jobs.length; + + stages.add( + rpc_model.PresubmitGuardStage( + name: g.stage.name, + createdAt: g.creationTime, + jobs: g.jobs, + ), + ); + } final guardStatus = GuardStatus.calculate( failedBuilds: totalFailed, @@ -123,14 +135,7 @@ final class GetPresubmitGuard extends PublicApiRequestHandler { author: first.author, guardStatus: guardStatus, enableGeminiLogAnalysis: config.flags.enableGeminiLogAnalysis, - stages: [ - for (final g in guards) - rpc_model.PresubmitGuardStage( - name: g.stage.name, - createdAt: g.creationTime, - jobs: g.jobs, - ), - ], + stages: stages, ); return Response.json(response); diff --git a/app_dart/lib/src/request_handlers/get_presubmit_guard_summaries.dart b/app_dart/lib/src/request_handlers/get_presubmit_guard_summaries.dart index 4de9181525..929868c7b3 100644 --- a/app_dart/lib/src/request_handlers/get_presubmit_guard_summaries.dart +++ b/app_dart/lib/src/request_handlers/get_presubmit_guard_summaries.dart @@ -94,18 +94,16 @@ final class GetPresubmitGuardSummaries extends PublicApiRequestHandler { final sha = entry.key; final shaGuards = entry.value; - final totalFailed = shaGuards.fold( - 0, - (int sum, PresubmitGuard g) => sum + g.failedJobs, - ); - final totalRemaining = shaGuards.fold( - 0, - (int sum, PresubmitGuard g) => sum + g.remainingJobs, - ); - final totalBuilds = shaGuards.fold( - 0, - (int sum, PresubmitGuard g) => sum + g.jobs.length, - ); + var totalFailed = 0; + var totalRemaining = 0; + var totalBuilds = 0; + + for (final g in shaGuards) { + totalFailed += g.failedJobs; + totalRemaining += g.remainingJobs; + totalBuilds += g.jobs.length; + } + final earliestCreationTime = shaGuards.fold( // assuming creation time is always in the past :) DateTime.now().millisecondsSinceEpoch, diff --git a/app_dart/lib/src/request_handlers/presubmit_guard_update_subscription.dart b/app_dart/lib/src/request_handlers/presubmit_guard_update_subscription.dart new file mode 100644 index 0000000000..526f931bf8 --- /dev/null +++ b/app_dart/lib/src/request_handlers/presubmit_guard_update_subscription.dart @@ -0,0 +1,68 @@ +// Copyright 2026 The Flutter Authors. All rights reserved. +// Use of this source code is governed by a BSD-style license that can be +// found in the LICENSE file. + +import 'dart:async'; +import 'dart:convert'; + +import 'package:cocoon_server/logging.dart'; +import 'package:meta/meta.dart'; + +import '../../cocoon_service.dart'; +import '../request_handling/subscription_handler.dart'; + +/// An endpoint for listening to PubSub updates to debounce and synchronize +/// [PresubmitGuard] state in Firestore and process stage completions. +@immutable +final class PresubmitGuardUpdateSubscription extends SubscriptionHandler { + /// Creates an endpoint for processing debounced [PresubmitGuard] updates. + const PresubmitGuardUpdateSubscription({ + required super.cache, + required super.config, + required Scheduler scheduler, + super.authProvider, + }) : _scheduler = scheduler, + super(subscriptionName: 'presubmit-guard-update-sub'); + + final Scheduler _scheduler; + + @override + Future post(Request request) async { + if (message.data == null || message.data!.isEmpty) { + log.info('presubmit-guard-update: No data in message.'); + return Response.emptyOk; + } + + late final Map messageJson; + try { + messageJson = jsonDecode(message.data!) as Map; + } catch (e) { + log.warn( + 'presubmit-guard-update: Failed to decode json data: ${message.data}', + e, + ); + return Response.emptyOk; + } + + final guardDocumentName = messageJson['guard_document_name'] as String?; + if (guardDocumentName == null || guardDocumentName.isEmpty) { + log.warn( + 'presubmit-guard-update: Missing guard_document_name in message data.', + ); + return Response.emptyOk; + } + + try { + await _scheduler.processPresubmitGuardUpdate(guardDocumentName); + } catch (e, s) { + log.error( + 'presubmit-guard-update: Error updating PresubmitGuard $guardDocumentName', + e, + s, + ); + rethrow; + } + + return Response.emptyOk; + } +} diff --git a/app_dart/lib/src/service/firestore/unified_check_run.dart b/app_dart/lib/src/service/firestore/unified_check_run.dart index 7d621829a3..93454c50d1 100644 --- a/app_dart/lib/src/service/firestore/unified_check_run.dart +++ b/app_dart/lib/src/service/firestore/unified_check_run.dart @@ -5,6 +5,8 @@ /// @docImport 'unified_check_run.dart'; library; +import 'dart:typed_data'; + import 'package:cocoon_common/task_status.dart'; import 'package:cocoon_server/logging.dart'; import 'package:collection/collection.dart'; @@ -19,6 +21,8 @@ import '../../model/firestore/base.dart'; import '../../model/firestore/ci_staging.dart'; import '../../model/firestore/presubmit_guard.dart'; import '../../model/firestore/presubmit_job.dart'; +import '../../request_handling/pubsub.dart'; +import '../cache_service.dart'; import '../config.dart'; import '../firestore.dart'; @@ -274,6 +278,55 @@ final class UnifiedCheckRun { ); } + /// Calculates the current live job statuses, remaining count, and failed count + /// for a given [PresubmitGuard]. + static Future getLatestJobStatusesForGuard({ + required FirestoreService firestoreService, + required PresubmitGuard guard, + PresubmitJob? overrideJob, + }) async { + final allJobs = await queryAllPresubmitJobsForGuard( + firestoreService: firestoreService, + checkRunId: guard.checkRunId, + ); + + final latestJobs = {}; + for (final job in allJobs) { + final current = latestJobs[job.jobName]; + if (current == null || job.attemptNumber > current.attemptNumber) { + latestJobs[job.jobName] = job; + } + } + + if (overrideJob != null) { + latestJobs[overrideJob.jobName] = overrideJob; + } + + var remaining = 0; + var failed = 0; + final jobStatuses = {}; + + for (final jobName in guard.jobs.keys) { + final job = latestJobs[jobName]; + final status = + job?.status ?? guard.jobs[jobName] ?? TaskStatus.waitingForBackfill; + jobStatuses[jobName] = status; + if (!status.isComplete) { + remaining++; + } + if (status.isFailure) { + failed++; + } + } + + return PresubmitGuardJobStatus( + remaining: remaining, + failed: failed, + jobStatuses: jobStatuses, + latestJobs: latestJobs, + ); + } + /// Returns check for the specified github [checkRunId] and /// [jobName] and [attemptNumber]. static Future queryPresubmitJob({ @@ -513,199 +566,197 @@ final class UnifiedCheckRun { /// both valid and recorded successfully, the record's `remaining` value /// signals how many more tests are running. Returns the record (valid: false) /// otherwise. - static Future markConclusion({ + /// Updates the corresponding [PresubmitJob] document and asynchronously schedules + /// a debounced [PresubmitGuard] synchronization via PubSub. + static Future markConclusion({ required FirestoreService firestoreService, required PresubmitGuardId guardId, required PresubmitJobState state, + required CacheService cacheService, + required PubSub pubsub, }) async { final changeCrumb = '${guardId.slug.owner}_${guardId.slug.name}_${guardId.prNum}_${guardId.checkRunId}'; final logCrumb = 'markConclusion(${changeCrumb}_${guardId.stage}, ${state.jobName}, ${state.status}, ${state.attemptNumber})'; - // Marking needs to happen while in a transaction to ensure `remaining` is - // updated correctly. For that to happen correctly; we need to perform a - // read of the document in the transaction as well. So start the transaction - // first thing. - final transaction = await firestoreService.beginTransaction(); + final checkDocName = PresubmitJob.documentNameFor( + slug: guardId.slug, + checkRunId: guardId.checkRunId, + jobName: state.jobName, + attemptNumber: state.attemptNumber, + ); - var remaining = -1; - var failed = -1; + // 1. Update the PresubmitJob document in a transaction. + final transaction = await firestoreService.beginTransaction(); var valid = false; - - late final PresubmitGuard presubmitGuard; - late final PresubmitJob presubmitJob; - // transaction block try { - // First: read the fields we want to change. - final presubmitGuardDocumentName = PresubmitGuard.documentNameFor( - slug: guardId.slug, - prNum: guardId.prNum, - checkRunId: guardId.checkRunId, - stage: guardId.stage, - ); - final presubmitGuardDocument = await firestoreService.getDocument( - presubmitGuardDocumentName, - transaction: transaction, - ); - presubmitGuard = PresubmitGuard.fromDocument(presubmitGuardDocument); - - // Check if the build is present in the guard before trying to load it. - if (presubmitGuard.jobs[state.jobName] == null) { - log.info( - '$logCrumb: ${state.jobName} with attemptNumber ${state.attemptNumber} not present for $transaction / ${presubmitGuardDocument.fields}', - ); - await firestoreService.rollback(transaction); - return PresubmitGuardConclusion( - result: PresubmitGuardConclusionResult.missing, - remaining: presubmitGuard.remainingJobs, - checkRunGuard: presubmitGuard.checkRunJson, - failed: presubmitGuard.failedJobs, - summary: - 'Check run "${state.jobName}" not present in ${guardId.stage} CI stage', - details: 'Change $changeCrumb', - ); - } - - final checkDocName = PresubmitJob.documentNameFor( - slug: guardId.slug, - checkRunId: guardId.checkRunId, - jobName: state.jobName, - attemptNumber: state.attemptNumber, - ); final presubmitJobDocument = await firestoreService.getDocument( checkDocName, transaction: transaction, ); - presubmitJob = PresubmitJob.fromDocument(presubmitJobDocument); - - remaining = presubmitGuard.remainingJobs; - failed = presubmitGuard.failedJobs; - final jobs = presubmitGuard.jobs; - var status = jobs[state.jobName]!; + final presubmitJob = PresubmitJob.fromDocument(presubmitJobDocument); - // If job is waiting for backfill, that means its initiated by github - // or re-run. So no processing needed, we should only update appropriate - // checks with that [TaskStatus] if (state.status == TaskStatus.waitingForBackfill) { - status = state.status; + presubmitJob.status = state.status; valid = true; - // If job is in progress, we should update apropriate checks with start - // time and their status to that [TaskStatus] only if the job is not - // completed. } else if (state.status == TaskStatus.inProgress) { presubmitJob.startTime = state.startTime!; presubmitJob.buildNumber = state.buildNumber; presubmitJob.buildId = state.buildId; - // If the job is not completed, update the status. - if (!status.isComplete) { - status = state.status; + if (!presubmitJob.status.isComplete) { + presubmitJob.status = state.status; } valid = true; } else { - // If job already compleated remaining and failed should not updated. - if (!status.isComplete) { - // "remaining" should go down if job is succeeded or failed. - // "failed_count" can go up or down depending on: - // attemptNumber > 1 && jobSuccessed: down (-1) - // attemptNumber = 1 && jobFailed: up (+1) - // So if the test existed and either remaining or failed_count is changed; - // the response is valid. - if (state.status.isComplete) { - // Guard against going negative and log enough info so we can debug. - if (remaining == 0) { - throw '$logCrumb: field "${PresubmitGuard.fieldRemainingJobs}" is already zero for $transaction / ${presubmitGuardDocument.fields}'; - } - remaining -= 1; - valid = true; - } - - if (state.status.isFailure) { - log.info('$logCrumb: test failed'); - failed += 1; - valid = true; - } - status = state.status; - // All checks pass. "valid" is only set to true if there was a change in either the remaining or failed count. - log.info( - '$logCrumb: setting remaining to $remaining, failed to $failed', - ); - presubmitGuard.remainingJobs = remaining; - presubmitGuard.failedJobs = failed; + if (!presubmitJob.status.isComplete) { + presubmitJob.status = state.status; presubmitJob.endTime = state.endTime!; presubmitJob.summary = state.summary; presubmitJob.buildNumber = state.buildNumber; presubmitJob.buildId = state.buildId; + valid = true; } else { - status = state.status; + presubmitJob.status = state.status; valid = true; } } - jobs[state.jobName] = status; - presubmitGuard.jobs = jobs; - presubmitJob.status = status; + + if (valid) { + await firestoreService.commit( + transaction, + documentsToWrites([presubmitJob], exists: true), + ); + } else { + await firestoreService.rollback(transaction); + } } on DetailedApiRequestError catch (e, stack) { + await firestoreService.rollback(transaction); if (e.status == 404) { - // An attempt to read a document not in firestore should not be retried. log.info( - '$logCrumb: ${PresubmitJob.collectionId} document not found for $transaction', - ); - await firestoreService.rollback(transaction); - return PresubmitGuardConclusion( - result: PresubmitGuardConclusionResult.internalError, - remaining: -1, - checkRunGuard: null, - failed: failed, - summary: 'Internal server error', - details: - ''' -${PresubmitJob.collectionId} document not found for stage "${guardId.stage}" for $changeCrumb. Got 404 from Firestore. -Error: ${e.toString()} -$stack -''', + '$logCrumb: ${PresubmitJob.collectionId} document not found for $transaction\n$stack', ); + return; } - // All other errors should bubble up and be retried. - await firestoreService.rollback(transaction); rethrow; } catch (e) { - // All other errors should bubble up and be retried. await firestoreService.rollback(transaction); rethrow; } - // Commit this write firebase and if no one else was writing at the same time, return success. - // If this commit fails, that means someone else modified firestore and the caller should try again. - // We do not need to rollback the transaction; firebase documentation says a failed commit takes care of that. + + if (!valid) { + log.info('$logCrumb: Not a valid state transition for ${state.jobName}'); + return; + } + + final presubmitGuardDocumentName = PresubmitGuard.documentNameFor( + slug: guardId.slug, + prNum: guardId.prNum, + checkRunId: guardId.checkRunId, + stage: guardId.stage, + ); + + // 2. Asynchronously trigger the debounced guard update via pubsub, using + // setIfNotExists so only the first job to mark the guard dirty publishes. + final wasSet = await cacheService.setIfNotExists( + 'presubmit_guard_dirty', + presubmitGuardDocumentName, + Uint8List.fromList([1]), + ttl: const Duration(minutes: 15), + ); + if (wasSet) { + try { + await pubsub.publish('presubmit-guard-update', { + 'guard_document_name': presubmitGuardDocumentName, + }); + } catch (e) { + log.warn( + '$logCrumb: Failed to publish presubmit-guard-update via pubsub', + e, + ); + await cacheService.purge( + 'presubmit_guard_dirty', + presubmitGuardDocumentName, + ); + } + } + } + + /// Asynchronously updates a [PresubmitGuard] document based on the live set + /// of [PresubmitJob] records. + /// + /// Used by the debounced PubSub subscription (`presubmit-guard-update`). + static Future<(PresubmitGuardConclusion, PresubmitGuard)?> + updatePresubmitGuard({ + required FirestoreService firestoreService, + required CacheService cacheService, + required String guardDocumentName, + }) async { + // Clear dirty flag right before querying/updating to allow new arrivals to re-debounce. + await cacheService.purge('presubmit_guard_dirty', guardDocumentName); + + Document latestGuardDoc; try { - final response = await firestoreService.commit( - transaction, - documentsToWrites([presubmitGuard, presubmitJob], exists: true), - ); - log.info( - '$logCrumb: results = ${response.writeResults?.map((e) => e.toJson())}', - ); - return PresubmitGuardConclusion( - result: valid - ? PresubmitGuardConclusionResult.ok - : PresubmitGuardConclusionResult.internalError, - remaining: remaining, - checkRunGuard: presubmitGuard.checkRunJson, - failed: failed, - summary: valid - ? 'Successfully updated presubmit guard status' - : 'Not a valid state transition for ${state.jobName}', - details: valid - ? ''' -For CI stage ${guardId.stage}: - Pending: $remaining - Failed: $failed -''' - : 'Attempted to set the state of job ${state.jobName} ' - 'to "${state.status.name}".', - ); - } catch (e) { - log.info('$logCrumb: failed to update presubmit job', e); + latestGuardDoc = await firestoreService.getDocument(guardDocumentName); + } on DetailedApiRequestError catch (e) { + if (e.status == 404) { + log.info( + 'PresubmitGuard $guardDocumentName not found in firestore (404), skipping.', + ); + return null; + } rethrow; } + + final latestGuard = PresubmitGuard.fromDocument(latestGuardDoc); + final guardStatusInfo = await getLatestJobStatusesForGuard( + firestoreService: firestoreService, + guard: latestGuard, + ); + + latestGuard.remainingJobs = guardStatusInfo.remaining; + latestGuard.failedJobs = guardStatusInfo.failed; + + final jobsMap = latestGuard.jobs; + for (final jobName in jobsMap.keys) { + if (guardStatusInfo.latestJobs[jobName] case final job?) { + jobsMap[jobName] = job.status; + } + } + latestGuard.jobs = jobsMap; + + final response = await firestoreService.writeViaTransaction( + documentsToWrites([latestGuard], exists: true), + ); + log.info( + 'updatePresubmitGuard($guardDocumentName): results = ${response.writeResults?.map((e) => e.toJson())}', + ); + + final conclusion = PresubmitGuardConclusion( + result: PresubmitGuardConclusionResult.ok, + remaining: guardStatusInfo.remaining, + checkRunGuard: latestGuard.checkRunJson, + failed: guardStatusInfo.failed, + summary: 'Successfully updated presubmit guard status', + details: + 'For CI stage ${latestGuard.stage}:\n Pending: ${guardStatusInfo.remaining}\n Failed: ${guardStatusInfo.failed}\n', + ); + return (conclusion, latestGuard); } } + +/// Holds aggregated job status information for a [PresubmitGuard]. +@immutable +final class PresubmitGuardJobStatus { + const PresubmitGuardJobStatus({ + required this.remaining, + required this.failed, + required this.jobStatuses, + required this.latestJobs, + }); + + final int remaining; + final int failed; + final Map jobStatuses; + final Map latestJobs; +} diff --git a/app_dart/lib/src/service/scheduler.dart b/app_dart/lib/src/service/scheduler.dart index ced9831a60..f8d5307b5c 100644 --- a/app_dart/lib/src/service/scheduler.dart +++ b/app_dart/lib/src/service/scheduler.dart @@ -34,6 +34,7 @@ import '../model/github/checks.dart' show MergeGroup; import '../model/github/workflow_job.dart'; import '../model/proto/internal/scheduler.pb.dart' as pb; import '../request_handling/http_utils.dart'; +import '../request_handling/pubsub.dart'; import 'big_query.dart'; import 'cache_service.dart'; import 'config.dart'; @@ -69,7 +70,10 @@ class Scheduler { required ContentAwareHashService contentAwareHash, required FirestoreService firestore, required BigQueryService bigQuery, - }) : _luciBuildService = luciBuildService, + PubSub pubsub = const PubSub(), + }) : _cache = cache, + _pubsub = pubsub, + _luciBuildService = luciBuildService, _githubChecksService = githubChecksService, _config = config, _getFilesChanged = getFilesChanged, @@ -83,6 +87,8 @@ class Scheduler { config: config, ); + final CacheService _cache; + final PubSub _pubsub; final GetFilesChanged _getFilesChanged; final Config _config; final GithubChecksService _githubChecksService; @@ -1139,11 +1145,11 @@ detailsUrl: $detailsUrl late PresubmitGuardConclusion stagingConclusion; if (check.isUnifiedCheckRun) { - stage = check.stage!; - stagingConclusion = await _markUnifiedCheckRunConclusion( + await _markUnifiedCheckRunConclusion( guardId: check.guardId, state: check.state, ); + return true; } else { // for github flow check runs are processed only if the build succeeded or // some kind of failure occurred. @@ -1583,7 +1589,7 @@ $stacktrace }); } - Future _markUnifiedCheckRunConclusion({ + Future _markUnifiedCheckRunConclusion({ required PresubmitGuardId guardId, required PresubmitJobState state, }) async { @@ -1591,17 +1597,16 @@ $stacktrace 'checkCompleted(${state.jobName}, ${state.buildNumber}, ${guardId.stage}, ${guardId.slug}, ${state.status})'; log.info('$logCrumb: ${guardId.documentId}'); - // We're doing a transactional update, which could fail if multiple tasks - // are running at the same time so retry a sane amount of times before - // giving up. const r = RetryOptions(maxAttempts: 10, maxDelay: Duration(minutes: 2)); try { - return await r.retry(() { + await r.retry(() { return UnifiedCheckRun.markConclusion( firestoreService: _firestore, guardId: guardId, state: state, + cacheService: _cache, + pubsub: _pubsub, ); }); } on Exception catch (e, s) { @@ -1610,6 +1615,120 @@ $stacktrace } } + /// Processes an asynchronous debounced update for a [PresubmitGuard] document. + /// + /// Invoked via PubSub (`/api/v2/presubmit-guard-update-subscription`). + Future processPresubmitGuardUpdate(String guardDocumentName) async { + final updateResult = await UnifiedCheckRun.updatePresubmitGuard( + firestoreService: _firestore, + cacheService: _cache, + guardDocumentName: guardDocumentName, + ); + if (updateResult == null) { + return; + } + + final (stagingConclusion, presubmitGuard) = updateResult; + if (!stagingConclusion.isOk) { + return; + } + + final guardCheckRun = checkRunFromString(presubmitGuard.checkRunJson); + final isMergeGroup = guardCheckRun.name == Config.kMergeQueueLockName; + final slug = presubmitGuard.slug; + final sha = presubmitGuard.commitSha; + final stage = presubmitGuard.stage; + final logCrumb = + 'processPresubmitGuardUpdate(${slug.fullName}, $sha, $stage, isMergeGroup=$isMergeGroup)'; + + if (stagingConclusion.result == + PresubmitGuardConclusionResult.internalError) { + if (isMergeGroup) { + await _completeArtifacts(sha, false); + await failGuardForMergeGroup( + slug: slug, + lock: guardCheckRun, + headSha: sha, + summary: stagingConclusion.summary, + details: stagingConclusion.details, + ); + } + return; + } + + if (stagingConclusion.isPending) { + log.info( + '$logCrumb: not progressing, remaining work count: ${stagingConclusion.remaining}', + ); + return; + } + + if (stagingConclusion.isFailed) { + if (isMergeGroup) { + await _completeArtifacts(sha, false); + await failGuardForMergeGroup( + slug: slug, + lock: guardCheckRun, + headSha: sha, + summary: stagingConclusion.summary, + details: stagingConclusion.details, + ); + } else { + final detailsUrl = + 'https://flutter-dashboard.appspot.com/#/presubmit?repo=${slug.name}&sha=$sha'; + await _requireActionForGuard( + slug: slug, + lock: guardCheckRun, + headSha: sha, + summary: _githubChecksService.getGithubSummaryWithHeader(''' +**[Failed Checks Details]($detailsUrl)** + +''', kDashboardChecksDescription), + details: + 'For $stage CI stage ${stagingConclusion.failed} checks failed', + detailsUrl: detailsUrl, + ); + } + return; + } + + switch (stage) { + case CiStage.fusionEngineBuild: + if (isMergeGroup) { + await _completeArtifacts(sha, true); + await _closeMergeQueue( + mergeQueueGuard: stagingConclusion.checkRunGuard!, + slug: slug, + sha: sha, + stage: CiStage.fusionEngineBuild, + logCrumb: logCrumb, + ); + } else { + final checkRunMap = + json.decode(presubmitGuard.checkRunJson) as Map; + if (checkRunMap['check_suite'] is Map) { + final suiteMap = checkRunMap['check_suite'] as Map; + suiteMap['pull_requests'] ??= []; + } + await _closeSuccessfulEngineBuildStage( + checkRun: cocoon_checks.CheckRun.fromJson(checkRunMap), + mergeQueueGuard: stagingConclusion.checkRunGuard!, + slug: slug, + sha: sha, + logCrumb: logCrumb, + ); + } + case CiStage.fusionTests: + case CiStage.genericTests: + await _closeSuccessfulTestStage( + mergeQueueGuard: stagingConclusion.checkRunGuard!, + slug: slug, + sha: sha, + logCrumb: logCrumb, + ); + } + } + /// Reschedules a failed build using a [CheckRunEvent]. The CheckRunEvent is /// generated when someone clicks the re-run button from a failed build from /// the Github UI. diff --git a/app_dart/test/request_handlers/get_engine_artifacts_ready_test.dart b/app_dart/test/request_handlers/get_engine_artifacts_ready_test.dart index 2f1ab173c8..a76e1a2dd6 100644 --- a/app_dart/test/request_handlers/get_engine_artifacts_ready_test.dart +++ b/app_dart/test/request_handlers/get_engine_artifacts_ready_test.dart @@ -4,6 +4,7 @@ import 'dart:convert'; +import 'package:cocoon_common/task_status.dart'; import 'package:cocoon_integration_test/testing.dart'; import 'package:cocoon_server_test/test_logging.dart'; import 'package:cocoon_service/src/model/common/firestore_extensions.dart'; @@ -144,6 +145,7 @@ void main() { slug: Config.flutterSlug, headSha: 'abc123', stage: CiStage.fusionEngineBuild, + jobs: {'some_job': TaskStatus.succeeded}, remainingJobs: 0, failedJobs: 0, ); @@ -169,6 +171,7 @@ void main() { slug: Config.flutterSlug, headSha: 'abc123', stage: CiStage.fusionEngineBuild, + jobs: {'some_job': TaskStatus.inProgress}, remainingJobs: 1, failedJobs: 0, ); @@ -194,6 +197,7 @@ void main() { slug: Config.flutterSlug, headSha: 'abc123', stage: CiStage.fusionEngineBuild, + jobs: {'some_job': TaskStatus.failed}, remainingJobs: 0, failedJobs: 1, ); diff --git a/app_dart/test/request_handlers/presubmit_guard_update_subscription_test.dart b/app_dart/test/request_handlers/presubmit_guard_update_subscription_test.dart new file mode 100644 index 0000000000..595f97977a --- /dev/null +++ b/app_dart/test/request_handlers/presubmit_guard_update_subscription_test.dart @@ -0,0 +1,148 @@ +// Copyright 2026 The Flutter Authors. All rights reserved. +// Use of this source code is governed by a BSD-style license that can be +// found in the LICENSE file. + +import 'dart:convert'; +import 'dart:typed_data'; + +import 'package:cocoon_common/task_status.dart'; +import 'package:cocoon_integration_test/testing.dart'; +import 'package:cocoon_server_test/test_logging.dart'; +import 'package:cocoon_service/cocoon_service.dart'; +import 'package:cocoon_service/src/model/luci/pubsub_message.dart'; +import 'package:github/github.dart'; +import 'package:test/test.dart'; + +import '../src/request_handling/subscription_tester.dart'; + +void main() { + useTestLoggerPerTest(); + + late PresubmitGuardUpdateSubscription handler; + late FakeConfig config; + late FakeHttpRequest request; + late SubscriptionTester tester; + late FakeFirestoreService firestore; + late CacheService cache; + + setUp(() { + firestore = FakeFirestoreService(); + config = FakeConfig(); + cache = CacheService.inMemory(); + + handler = PresubmitGuardUpdateSubscription( + cache: cache, + config: config, + scheduler: FakeScheduler( + config: config, + firestore: firestore, + bigQuery: MockBigQueryService(), + cache: cache, + ), + ); + + request = FakeHttpRequest(); + tester = SubscriptionTester(request: request); + }); + + test('returns emptyOk when message data is null', () async { + tester.message = const PushMessage(data: null); + await tester.post(handler); + }); + + test('skips update when guard is not marked dirty in cache', () async { + const guardName = + 'projects/flutter-dashboard/databases/cocoon/documents/presubmit_guards/flutter_flutter_123_456_fusionTests'; + tester.message = PushMessage( + data: jsonEncode({'guard_document_name': guardName}), + ); + + await tester.post(handler); + }); + + test('updates PresubmitGuard with latest jobs when dirty in cache', () async { + final slug = RepositorySlug('flutter', 'flutter'); + const prNum = 123; + const checkRunId = 456; + const stage = CiStage.fusionTests; + + final guard = + generatePresubmitGuard( + slug: slug, + prNum: prNum, + stage: stage, + checkRun: generateCheckRun(checkRunId), + remainingJobs: 2, + failedJobs: 0, + jobs: { + 'linux_test': TaskStatus.waitingForBackfill, + 'mac_test': TaskStatus.waitingForBackfill, + }, + ) + ..name = PresubmitGuard.documentNameFor( + slug: slug, + prNum: prNum, + checkRunId: checkRunId, + stage: stage, + ); + await firestore.writeViaTransaction( + documentsToWrites([guard], exists: false), + ); + + final job1 = + PresubmitJob( + slug: slug, + checkRunId: checkRunId, + jobName: 'linux_test', + status: TaskStatus.succeeded, + attemptNumber: 1, + creationTime: 1000, + ) + ..name = PresubmitJob.documentNameFor( + slug: slug, + checkRunId: checkRunId, + jobName: 'linux_test', + attemptNumber: 1, + ); + final job2 = + PresubmitJob( + slug: slug, + checkRunId: checkRunId, + jobName: 'mac_test', + status: TaskStatus.inProgress, + attemptNumber: 1, + creationTime: 1000, + ) + ..name = PresubmitJob.documentNameFor( + slug: slug, + checkRunId: checkRunId, + jobName: 'mac_test', + attemptNumber: 1, + ); + await firestore.writeViaTransaction( + documentsToWrites([job1, job2], exists: false), + ); + + final guardName = guard.name!; + await cache.set( + 'presubmit_guard_dirty', + guardName, + Uint8List.fromList([1]), + ); + + tester.message = PushMessage( + data: jsonEncode({'guard_document_name': guardName}), + ); + + await tester.post(handler); + + final updatedDoc = await firestore.getDocument(guardName); + final updatedGuard = PresubmitGuard.fromDocument(updatedDoc); + + expect(updatedGuard.remainingJobs, 1); + expect(updatedGuard.failedJobs, 0); + expect(updatedGuard.jobs['linux_test'], TaskStatus.succeeded); + expect(updatedGuard.jobs['mac_test'], TaskStatus.inProgress); + expect(await cache.get('presubmit_guard_dirty', guardName), isNull); + }); +} diff --git a/app_dart/test/service/firestore/unified_check_run_test.dart b/app_dart/test/service/firestore/unified_check_run_test.dart index 1f8e0f2fc1..88b3a0530f 100644 --- a/app_dart/test/service/firestore/unified_check_run_test.dart +++ b/app_dart/test/service/firestore/unified_check_run_test.dart @@ -10,6 +10,7 @@ import 'package:cocoon_service/src/model/common/presubmit_job_state.dart'; import 'package:cocoon_service/src/model/firestore/base.dart'; import 'package:cocoon_service/src/model/firestore/presubmit_guard.dart'; import 'package:cocoon_service/src/model/firestore/presubmit_job.dart'; +import 'package:cocoon_service/src/service/cache_service.dart'; import 'package:cocoon_service/src/service/firestore.dart'; import 'package:cocoon_service/src/service/firestore/unified_check_run.dart'; import 'package:cocoon_service/src/service/flags/dynamic_config.dart'; @@ -23,10 +24,14 @@ void main() { late FakeConfig config; late FakeFirestoreService firestoreService; + late CacheService cacheService; + late FakePubSub pubsub; setUp(() { config = FakeConfig(); firestoreService = FakeFirestoreService(); + cacheService = CacheService.inMemory(); + pubsub = FakePubSub(); }); group('UnifiedCheckRun', () { @@ -180,16 +185,14 @@ void main() { buildId: Int64.MAX_VALUE, ); - final result = await UnifiedCheckRun.markConclusion( + await UnifiedCheckRun.markConclusion( firestoreService: firestoreService, guardId: guardId, state: state, + cacheService: cacheService, + pubsub: pubsub, ); - expect(result.result, PresubmitGuardConclusionResult.ok); - expect(result.remaining, 1); - expect(result.failed, 0); - final checkDoc = await PresubmitJob.fromFirestore( firestoreService, PresubmitJobId( @@ -203,12 +206,40 @@ void main() { expect(checkDoc.endTime, 3000); expect(checkDoc.buildNumber, 456); expect(checkDoc.buildId, Int64.MAX_VALUE); + + final guardDocName = PresubmitGuard.documentNameFor( + slug: guardId.slug, + prNum: guardId.prNum, + checkRunId: guardId.checkRunId, + stage: guardId.stage, + ); + expect(pubsub.messages.length, 1); + expect( + await cacheService.get('presubmit_guard_dirty', guardDocName), + isNotNull, + ); + + final result = (await UnifiedCheckRun.updatePresubmitGuard( + firestoreService: firestoreService, + cacheService: cacheService, + guardDocumentName: guardDocName, + ))?.$1; + expect(result!.result, PresubmitGuardConclusionResult.ok); + expect(result.remaining, 1); + expect(result.failed, 0); }); test( 'update all check status to succeeded lead to complete guard', () async { - final result1 = await UnifiedCheckRun.markConclusion( + final guardDocName = PresubmitGuard.documentNameFor( + slug: guardId.slug, + prNum: guardId.prNum, + checkRunId: guardId.checkRunId, + stage: guardId.stage, + ); + + await UnifiedCheckRun.markConclusion( firestoreService: firestoreService, guardId: guardId, state: const PresubmitJobState( @@ -218,15 +249,23 @@ void main() { startTime: 2000, endTime: 3000, ), + cacheService: cacheService, + pubsub: pubsub, ); - expect(result1.remaining, 1); + final result1 = (await UnifiedCheckRun.updatePresubmitGuard( + firestoreService: firestoreService, + cacheService: cacheService, + guardDocumentName: guardDocName, + ))?.$1; + + expect(result1!.remaining, 1); expect(result1.failed, 0); expect(result1.isOk, true); expect(result1.isComplete, false); expect(result1.isPending, true); - final result2 = await UnifiedCheckRun.markConclusion( + await UnifiedCheckRun.markConclusion( firestoreService: firestoreService, guardId: guardId, state: const PresubmitJobState( @@ -236,9 +275,17 @@ void main() { startTime: 2000, endTime: 3000, ), + cacheService: cacheService, + pubsub: pubsub, ); - expect(result2.remaining, 0); + final result2 = (await UnifiedCheckRun.updatePresubmitGuard( + firestoreService: firestoreService, + cacheService: cacheService, + guardDocumentName: guardDocName, + ))?.$1; + + expect(result2!.remaining, 0); expect(result2.failed, 0); expect(result2.isOk, true); expect(result2.isComplete, true); @@ -267,13 +314,28 @@ void main() { endTime: 3000, ); - final result = await UnifiedCheckRun.markConclusion( + final guardDocName = PresubmitGuard.documentNameFor( + slug: guardId.slug, + prNum: guardId.prNum, + checkRunId: guardId.checkRunId, + stage: guardId.stage, + ); + + await UnifiedCheckRun.markConclusion( firestoreService: firestoreService, guardId: guardId, state: state, + cacheService: cacheService, + pubsub: pubsub, ); - expect(result.result, PresubmitGuardConclusionResult.ok); + final result = (await UnifiedCheckRun.updatePresubmitGuard( + firestoreService: firestoreService, + cacheService: cacheService, + guardDocumentName: guardDocName, + ))?.$1; + + expect(result!.result, PresubmitGuardConclusionResult.ok); expect(result.remaining, 1); expect(result.failed, 1); }); @@ -285,14 +347,17 @@ void main() { attemptNumber: 1, ); - final result = await UnifiedCheckRun.markConclusion( + await UnifiedCheckRun.markConclusion( firestoreService: firestoreService, guardId: guardId, state: state, + cacheService: cacheService, + pubsub: pubsub, ); - expect(result.result, PresubmitGuardConclusionResult.missing); + expect(pubsub.messages.isEmpty, true); }); + test('updates check status and build number on inProgress', () async { final state = const PresubmitJobState( jobName: 'linux', @@ -303,14 +368,14 @@ void main() { buildId: Int64.MAX_VALUE, ); - final result = await UnifiedCheckRun.markConclusion( + await UnifiedCheckRun.markConclusion( firestoreService: firestoreService, guardId: guardId, state: state, + cacheService: cacheService, + pubsub: pubsub, ); - expect(result.result, PresubmitGuardConclusionResult.ok); - final checkDoc = await PresubmitJob.fromFirestore( firestoreService, PresubmitJobId( diff --git a/app_dart/test/service/scheduler_test.dart b/app_dart/test/service/scheduler_test.dart index 7f849316c5..0938bff2fd 100644 --- a/app_dart/test/service/scheduler_test.dart +++ b/app_dart/test/service/scheduler_test.dart @@ -4246,6 +4246,14 @@ targets: expect(await scheduler.processCheckRunCompleted(check), isTrue); + final guardName = PresubmitGuard.documentNameFor( + slug: pullRequest.base!.repo!.slug(), + prNum: pullRequest.number!, + checkRunId: checkRunGuard.id!, + stage: CiStage.fusionEngineBuild, + ); + await scheduler.processPresubmitGuardUpdate(guardName); + // Should schedule tests for the next stage (fusionTests) expect(fakeLuciBuildService.scheduledTryBuilds, isNotEmpty); expect(fakeLuciBuildService.stage, CiStage.fusionTests); @@ -4326,6 +4334,14 @@ targets: expect(await scheduler.processCheckRunCompleted(check), isTrue); + final guardName = PresubmitGuard.documentNameFor( + slug: pullRequest.base!.repo!.slug(), + prNum: pullRequest.number!, + checkRunId: checkRunGuard.id!, + stage: CiStage.fusionTests, + ); + await scheduler.processPresubmitGuardUpdate(guardName); + verify( mockGithubChecksUtil.updateCheckRun( any, @@ -4416,6 +4432,14 @@ targets: expect(await scheduler.processCheckRunCompleted(check), isTrue); + final guardName = PresubmitGuard.documentNameFor( + slug: pullRequest.base!.repo!.slug(), + prNum: pullRequest.number!, + checkRunId: checkRunGuard.id!, + stage: CiStage.fusionTests, + ); + await scheduler.processPresubmitGuardUpdate(guardName); + verify( mockGithubChecksUtil.updateCheckRun( any, @@ -4503,6 +4527,14 @@ targets: expect(await scheduler.processCheckRunCompleted(check), isTrue); + final guardName = PresubmitGuard.documentNameFor( + slug: pullRequest.base!.repo!.slug(), + prNum: pullRequest.number!, + checkRunId: checkRunGuard.id!, + stage: CiStage.genericTests, + ); + await scheduler.processPresubmitGuardUpdate(guardName); + verify( mockGithubChecksUtil.updateCheckRun( any, diff --git a/packages/cocoon_integration_test/lib/src/fakes/fake_scheduler.dart b/packages/cocoon_integration_test/lib/src/fakes/fake_scheduler.dart index 5233a84f1b..7ea6b43dca 100644 --- a/packages/cocoon_integration_test/lib/src/fakes/fake_scheduler.dart +++ b/packages/cocoon_integration_test/lib/src/fakes/fake_scheduler.dart @@ -5,6 +5,7 @@ import 'package:cocoon_service/src/foundation/github_checks_util.dart'; import 'package:cocoon_service/src/model/ci_yaml/ci_yaml.dart'; import 'package:cocoon_service/src/model/proto/protos.dart' as pb; +import 'package:cocoon_service/src/request_handling/pubsub.dart'; import 'package:cocoon_service/src/service/build_bucket_client.dart'; import 'package:cocoon_service/src/service/cache_service.dart'; import 'package:cocoon_service/src/service/config.dart'; @@ -22,6 +23,7 @@ import 'fake_content_aware_hash_service.dart'; import 'fake_get_files_changed.dart'; import 'fake_github_service.dart'; import 'fake_luci_build_service.dart'; +import 'fake_pubsub.dart'; /// Fake for [Scheduler] to use for tests that rely on it. class FakeScheduler extends Scheduler { @@ -36,8 +38,11 @@ class FakeScheduler extends Scheduler { ContentAwareHashService? contentAwareHash, required super.firestore, required super.bigQuery, + CacheService? cache, + PubSub? pubsub, }) : super( - cache: CacheService.inMemory(), + cache: cache ?? CacheService.inMemory(), + pubsub: pubsub ?? FakePubSub(), githubChecksService: GithubChecksService( config, githubChecksUtil: githubChecksUtil, diff --git a/packages/cocoon_integration_test/lib/src/utilities/entity_generators.dart b/packages/cocoon_integration_test/lib/src/utilities/entity_generators.dart index 5856df24ae..d7d25b48d0 100644 --- a/packages/cocoon_integration_test/lib/src/utilities/entity_generators.dart +++ b/packages/cocoon_integration_test/lib/src/utilities/entity_generators.dart @@ -188,7 +188,10 @@ github.CheckRun generateCheckRun( 'id': i, 'name': name, 'started_at': startedAt.toIso8601String(), - 'check_suite': {'id': checkSuite}, + 'check_suite': { + 'id': checkSuite, + 'pull_requests': [], + }, }); }