Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
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