Store secrets in separate db column

Signed-off-by: Erik Larsson <erik.larsson@schibsted.com>
This commit is contained in:
Erik Larsson
2021-04-09 23:33:15 +02:00
parent 008e8a8b1c
commit a41a912f2d
6 changed files with 73 additions and 11 deletions
@@ -0,0 +1,38 @@
/*
* Copyright 2020 Spotify AB
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
// @ts-check
/**
* @param {import('knex').Knex} knex
*/
exports.up = async function up(knex) {
await knex.schema.alterTable('tasks', table => {
table
.string('secrets')
.nullable()
.comment('JSON encoded secrets to authenticate tasks with');
});
};
/**
* @param {import('knex').Knex} knex
*/
exports.down = async function down(knex) {
await knex.schema.alterTable('tasks', table => {
table.dropColumn('secrets');
});
};
@@ -24,6 +24,7 @@ import {
DbTaskRow,
Status,
TaskEventType,
TaskSecrets,
TaskSpec,
TaskStore,
TaskStoreEmitOptions,
@@ -42,6 +43,7 @@ export type RawDbTaskRow = {
status: Status;
last_heartbeat_at?: string;
created_at: string;
secrets?: string;
};
export type RawDbTaskEventRow = {
@@ -71,23 +73,29 @@ export class DatabaseTaskStore implements TaskStore {
}
try {
const spec = JSON.parse(result.spec);
const secrets = result.secrets ? JSON.parse(result.secrets) : undefined;
return {
id: result.id,
spec,
status: result.status,
lastHeartbeatAt: result.last_heartbeat_at,
createdAt: result.created_at,
secrets,
};
} catch (error) {
throw new Error(`Failed to parse spec of task '${taskId}', ${error}`);
}
}
async createTask(spec: TaskSpec): Promise<{ taskId: string }> {
async createTask(
spec: TaskSpec,
secrets?: TaskSecrets,
): Promise<{ taskId: string }> {
const taskId = uuid();
await this.db<RawDbTaskRow>('tasks').insert({
id: taskId,
spec: JSON.stringify(spec),
secrets: secrets ? JSON.stringify(secrets) : undefined,
status: 'open',
});
return { taskId };
@@ -119,12 +127,14 @@ export class DatabaseTaskStore implements TaskStore {
try {
const spec = JSON.parse(task.spec);
const secrets = task.secrets ? JSON.parse(task.spec) : undefined;
return {
id: task.id,
spec,
status: 'processing',
lastHeartbeatAt: task.last_heartbeat_at,
createdAt: task.created_at,
secrets,
};
} catch (error) {
throw new Error(`Failed to parse spec of task '${task.id}', ${error}`);
@@ -209,6 +219,7 @@ export class DatabaseTaskStore implements TaskStore {
})
.update({
status,
secrets: undefined,
});
if (updateCount !== 1) {
throw new ConflictError(
@@ -18,6 +18,7 @@ import { Logger } from 'winston';
import {
CompletedTaskState,
Task,
TaskSecrets,
TaskSpec,
TaskStore,
TaskBroker,
@@ -136,8 +137,11 @@ export class StorageTaskBroker implements TaskBroker {
}
}
async dispatch(spec: TaskSpec): Promise<DispatchResult> {
const taskRow = await this.storage.createTask(spec);
async dispatch(
spec: TaskSpec,
secrets?: TaskSecrets,
): Promise<DispatchResult> {
const taskRow = await this.storage.createTask(spec, secrets);
this.signalDispatch();
return {
taskId: taskRow.taskId,
@@ -134,7 +134,7 @@ export class TaskWorker {
logger: taskLogger,
logStream: stream,
input,
token: task.spec.token,
token: task.secrets?.identityToken,
workspacePath,
async createTemporaryDirectory() {
const tmpDir = await fs.mkdtemp(
@@ -31,6 +31,7 @@ export type DbTaskRow = {
status: Status;
createdAt: string;
lastHeartbeatAt?: string;
secrets?: TaskSecrets;
};
export type TaskEventType = 'completion' | 'log';
@@ -44,7 +45,6 @@ export type DbTaskEventRow = {
export type TaskSpec = {
baseUrl?: string;
token?: string | undefined;
values: JsonObject;
steps: Array<{
id: string;
@@ -55,12 +55,17 @@ export type TaskSpec = {
output: { [name: string]: string };
};
export type TaskSecrets = {
identityToken: string | undefined;
};
export type DispatchResult = {
taskId: string;
};
export interface Task {
spec: TaskSpec;
secrets?: TaskSecrets;
done: boolean;
emitLog(message: string, metadata?: JsonValue): Promise<void>;
complete(result: CompletedTaskState, metadata?: JsonValue): Promise<void>;
@@ -69,7 +74,7 @@ export interface Task {
export interface TaskBroker {
claim(): Promise<Task>;
dispatch(spec: TaskSpec): Promise<DispatchResult>;
dispatch(spec: TaskSpec, secrets?: TaskSecrets): Promise<DispatchResult>;
vacuumTasks(timeoutS: { timeoutS: number }): Promise<void>;
observe(
options: {
@@ -94,7 +99,10 @@ export type TaskStoreGetEventsOptions = {
};
export interface TaskStore {
createTask(task: TaskSpec): Promise<{ taskId: string }>;
createTask(
task: TaskSpec,
secrets?: TaskSecrets,
): Promise<{ taskId: string }>;
getTask(taskId: string): Promise<DbTaskRow>;
claimTask(): Promise<DbTaskRow | undefined>;
completeTask(options: {
@@ -387,7 +387,6 @@ export async function createRouter(
taskSpec = {
baseUrl,
token,
values,
steps: template.spec.steps.map((step, index) => ({
...step,
@@ -404,7 +403,9 @@ export async function createRouter(
);
}
const result = await taskBroker.dispatch(taskSpec);
const result = await taskBroker.dispatch(taskSpec, {
identityToken: token,
});
res.status(201).json({ id: result.taskId });
})
@@ -414,8 +415,8 @@ export async function createRouter(
if (!task) {
throw new NotFoundError(`Task with id ${taskId} does not exist`);
}
// Do not disclose token
delete task.spec.token;
// Do not disclose secrets
delete task.secrets;
res.status(200).json(task);
})
.get('/v2/tasks/:taskId/eventstream', async (req, res) => {