Skip to content
5 changes: 4 additions & 1 deletion app_dart/bin/gae_server.dart
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,10 @@ Future<void> main() async {
}

final cache = CacheService.redis();
final firestore = await FirestoreService.from(const GoogleAuthProvider());
final firestore = await FirestoreService.from(
const GoogleAuthProvider(),
cache: cache,
);
final bigQuery = await BigQueryService.from(const GoogleAuthProvider());

// Start with a fresh copy of the DynamicConfig. If this throws, the server
Expand Down
71 changes: 49 additions & 22 deletions app_dart/lib/src/model/firestore/task.dart
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,9 @@
/// @docImport 'commit.dart';
library;

import 'dart:convert';
import 'dart:typed_data';

import 'package:buildbucket/buildbucket_pb.dart' as bbv2;
import 'package:cocoon_common/task_status.dart';
import 'package:googleapis/firestore/v1.dart' hide Status;
Expand Down Expand Up @@ -113,6 +116,7 @@ final class Task extends AppDocument<Task> {
static const fieldStatus = 'status';
static const fieldTestFlaky = 'testFlaky';
static const fieldAttempt = 'attempt';
static const fieldRevisionId = 'revisionId';

/// Returns a document ID for a task from the given parameters.
static AppDocumentId<Task> documentIdFor({
Expand All @@ -136,17 +140,27 @@ final class Task extends AppDocument<Task> {
fromDocument: Task.fromDocument,
);

/// Serializes a [Task] document into bytes for caching.
static Uint8List serialize(Task task) {
return utf8.encode(
json.encode(Document(name: task.name, fields: task.fields).toJson()),
);
}

/// Deserializes bytes into a [Task] document.
static Task deserialize(Uint8List data) {
final jsonMap = json.decode(utf8.decode(data)) as Map<String, Object?>;
return Task.fromDocument(Document.fromJson(jsonMap));
}

/// Lookup [Task] from Firestore.
///
/// `documentName` follows `/projects/{project}/databases/{database}/documents/{document_path}`
static Future<Task> fromFirestore(
FirestoreService firestoreService,
AppDocumentId<Task> id,
) async {
final document = await firestoreService.getDocument(
p.posix.join(kDatabase, 'documents', kTaskCollectionId, id.documentId),
);
return Task.fromDocument(document);
return firestoreService.getTask(id);
}

factory Task({
Expand Down Expand Up @@ -178,6 +192,7 @@ final class Task extends AppDocument<Task> {
fieldStatus: status.value.toValue(),
fieldTestFlaky: testFlaky.toValue(),
fieldAttempt: currentAttempt.toValue(),
fieldRevisionId: 1.toValue(),
},
name: p.posix.join(
kDatabase,
Expand Down Expand Up @@ -214,27 +229,17 @@ final class Task extends AppDocument<Task> {
..name = name;
}

/// Returns a Firestore [Write] that patches the [status] field for [id].
@useResult
static Write patchStatus(AppDocumentId<Task> id, TaskStatus status) {
return Write(
currentDocument: Precondition(exists: true),
update: Document(
name: p.posix.join(
kDatabase,
'documents',
kTaskCollectionId,
id.documentId,
),
fields: {fieldStatus: Value(stringValue: status.value)},
),
updateMask: DocumentMask(fieldPaths: [fieldStatus]),
);
}

/// The task was run successfully.
static const statusSucceeded = TaskStatus.succeeded;

int get revisionId => fields.containsKey(fieldRevisionId)
? int.parse(fields[fieldRevisionId]!.integerValue!)
: 1;

void incrementRevisionId() {
fields[fieldRevisionId] = (revisionId + 1).toValue();
}

/// The timestamp (in milliseconds since the Epoch) that this task was
/// created.
///
Expand Down Expand Up @@ -308,14 +313,17 @@ final class Task extends AppDocument<Task> {

void setStatus(TaskStatus status) {
fields[fieldStatus] = status.value.toValue();
incrementRevisionId();
}

void setEndTimestamp(int endTimestamp) {
fields[fieldEndTimestamp] = endTimestamp.toValue();
incrementRevisionId();
}

void setTestFlaky(bool testFlaky) {
fields[fieldTestFlaky] = testFlaky.toValue();
incrementRevisionId();
}

void updateFromBuild(bbv2.Build build) {
Expand All @@ -335,6 +343,23 @@ final class Task extends AppDocument<Task> {
.toValue();

_setStatusFromLuciStatus(build);
incrementRevisionId();
}

Task createRetry({DateTime? now}) {
now ??= DateTime.now();
return Task(
builderName: taskName,
currentAttempt: currentAttempt + 1,
commitSha: commitSha,
bringup: bringup,
createTimestamp: now.millisecondsSinceEpoch,
startTimestamp: 0,
endTimestamp: 0,
status: TaskStatus.waitingForBackfill,
testFlaky: false,
buildNumber: null,
);
}

void resetAsRetry({int? attempt, DateTime? now}) {
Expand All @@ -360,11 +385,13 @@ final class Task extends AppDocument<Task> {
fieldTestFlaky: false.toValue(),
fieldCommitSha: commitSha.toValue(),
fieldAttempt: attempt.toValue(),
fieldRevisionId: 1.toValue(),
};
}

void setBuildNumber(int buildNumber) {
fields[fieldBuildNumber] = buildNumber.toValue();
incrementRevisionId();
}

void _setStatusFromLuciStatus(bbv2.Build build) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ import 'dart:convert';
import 'package:buildbucket/buildbucket_pb.dart' as bbv2;
import 'package:cocoon_common/is_dart_internal.dart';
import 'package:cocoon_server/logging.dart';
import 'package:googleapis/firestore/v1.dart';

import '../../cocoon_service.dart';
import '../model/bbv2_extension.dart';
Expand Down Expand Up @@ -89,10 +88,7 @@ final class DartInternalSubscription extends SubscriptionHandler {
}

log.info('Inserting Task into Firestore: ${fsTask.toString()}');
await _firestore.batchWriteDocuments(
BatchWriteRequest(writes: documentsToWrites([fsTask])),
kDatabase,
);
await _firestore.updateTasks([fsTask]);

return Response.json(fsTask.toString());
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,6 @@ import 'package:archive/archive.dart';
import 'package:buildbucket/buildbucket_pb.dart' as bbv2;
import 'package:cocoon_common/is_release_branch.dart';
import 'package:cocoon_server/logging.dart';
import 'package:googleapis/firestore/v1.dart' hide Status;

import '../../ci_yaml.dart';
import '../model/firestore/commit.dart' as fs;
Expand Down Expand Up @@ -150,10 +149,7 @@ final class PostsubmitLuciSubscription extends SubscriptionHandler {

Future<void> _updateFirestore(fs.Task fsTask, bbv2.Build build) async {
fsTask.updateFromBuild(build);
await _firestore.batchWriteDocuments(
BatchWriteRequest(writes: documentsToWrites([fsTask], exists: true)),
kDatabase,
);
await _firestore.updateTasks([fsTask]);
}

// No need to update task in Firestore if
Expand Down
37 changes: 16 additions & 21 deletions app_dart/lib/src/request_handlers/rerun_prod_task.dart
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@
import 'package:cocoon_common/is_dart_internal.dart';
import 'package:cocoon_common/task_status.dart';
import 'package:cocoon_server/logging.dart';
import 'package:googleapis/firestore/v1.dart' as g;
import 'package:meta/meta.dart';

import '../model/ci_yaml/ci_yaml.dart';
Expand Down Expand Up @@ -184,9 +183,9 @@ final class RerunProdTask extends ApiRequestHandler {
).withMostRecentTaskOnly().tasks;

final wasMarkedNew = <String>[];
final documentWrites = <g.Write>[];
final tasksToUpdate = <fs.Task>[];
final createdTasks = <fs.Task>[];

// Wait for cancellations?
final Future<void> cancelRunningTasks;
if (statusesToRerun.contains(TaskStatus.inProgress)) {
cancelRunningTasks = _luciBuildService.cancelBuildsBySha(
Expand All @@ -208,32 +207,28 @@ final class RerunProdTask extends ApiRequestHandler {
}

// If it appears the task was in progress, cancel any running builders
// and crease a _new_ task (to represent a new run).
// and create a _new_ task (to represent a new run).
if (task.status == TaskStatus.inProgress) {
// Mark cancelled.
documentWrites.add(
fs.Task.patchStatus(
fs.TaskId(
commitSha: task.commitSha,
currentAttempt: task.currentAttempt,
taskName: task.taskName,
),
TaskStatus.cancelled,
),
);
// Mark current attempt cancelled with incremented revisionId
task.setStatus(TaskStatus.cancelled);
tasksToUpdate.add(task);
}

// Start a new task.
task.resetAsRetry(now: _now());
documentWrites.add(
g.Write(currentDocument: g.Precondition(exists: false), update: task),
);
// Start a new task attempt
final retryTask = task.createRetry(now: _now());
tasksToUpdate.add(retryTask);
createdTasks.add(retryTask);
}

final writes = documentsToWrites(tasksToUpdate);
await Future.wait([
cancelRunningTasks,
_firestore.commit(transaction, documentWrites),
_firestore.commit(transaction, writes),
]);
await _firestore.cacheTaskPayloads(tasksToUpdate);
if (createdTasks.isNotEmpty) {
await _firestore.updateCacheForCreatedTasks(createdTasks);
}

return wasMarkedNew;
}
Expand Down
42 changes: 17 additions & 25 deletions app_dart/lib/src/request_handlers/scheduler/batch_backfiller.dart
Original file line number Diff line number Diff line change
Expand Up @@ -174,31 +174,23 @@ final class BatchBackfiller extends ApiRequestHandler {
Iterable<BackfillTask> schedule,
Iterable<SkippableTask> skip,
) async {
log.debug('Querying ${schedule.length} tasks in Firestore...');
await _firestore.writeViaTransaction([
...schedule.map((toUpdate) {
final BackfillTask(:task) = toUpdate;
return fs.Task.patchStatus(
fs.TaskId(
commitSha: task.commitSha,
taskName: task.name,
currentAttempt: task.currentAttempt,
),
TaskStatus.inProgress,
);
}),
...skip.map((toSkip) {
final SkippableTask(:task) = toSkip;
return fs.Task.patchStatus(
fs.TaskId(
commitSha: task.commitSha,
taskName: task.name,
currentAttempt: task.currentAttempt,
),
TaskStatus.skipped,
);
}),
]);
final taskStatusMap = <fs.TaskId, TaskStatus>{
for (final item in schedule)
fs.TaskId(
commitSha: item.task.commitSha,
taskName: item.task.name,
currentAttempt: item.task.currentAttempt,
): TaskStatus.inProgress,
for (final item in skip)
fs.TaskId(
commitSha: item.task.commitSha,
taskName: item.task.name,
currentAttempt: item.task.currentAttempt,
): TaskStatus.skipped,
};
if (taskStatusMap.isEmpty) return;

await _firestore.updateTaskStatuses(taskStatusMap);
log.debug('Wrote to Firestore for backfill');
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ import 'package:cocoon_common/is_release_branch.dart';
import 'package:cocoon_common/task_status.dart';
import 'package:cocoon_server/logging.dart';
import 'package:github/github.dart' as gh;
import 'package:googleapis/firestore/v1.dart';

import 'package:meta/meta.dart';

import '../../../cocoon_service.dart';
Expand Down Expand Up @@ -139,10 +139,7 @@ final class VacuumStaleTasks extends ApiRequestHandler {
}
tasks.add(task);
}
await _firestore.batchWriteDocuments(
BatchWriteRequest(writes: documentsToWrites(tasks)),
kDatabase,
);
await _firestore.updateTasks(tasks);
}
}

Expand Down
Loading
Loading