diff --git a/i18n/en.json b/i18n/en.json index 1022a6f851..0acfe68fa3 100644 --- a/i18n/en.json +++ b/i18n/en.json @@ -1788,6 +1788,7 @@ "restore_all": "Restore all", "restore_user": "Restore user", "restored_asset": "Restored asset", + "result": "Result", "resume": "Resume", "resume_paused_jobs": "Resume {count, plural, one {# paused job} other {# paused jobs}}", "review_duplicates": "Review duplicates", @@ -2278,6 +2279,13 @@ "workflow_info": "Workflow info", "workflow_json": "Workflow JSON", "workflow_json_help": "Edit the workflow configuration in JSON format. Changes will sync to the visual builder.", + "workflow_logging_completed": "Completed", + "workflow_logging_disable": "Disable logging", + "workflow_logging_disabled_description": "Logging is currently disabled for this workflow.", + "workflow_logging_enable": "Enable logging", + "workflow_logging_error_step": "Error on step #{step}", + "workflow_logging_halted": "Halted", + "workflow_logging_halted_step": "Halted on step #{step}", "workflow_name": "Workflow name", "workflow_navigation_prompt": "Are you sure you want to leave without saving your changes?", "workflow_summary": "Workflow summary", diff --git a/mobile/lib/utils/openapi_patching.dart b/mobile/lib/utils/openapi_patching.dart index e92b2afd12..9df2eea462 100644 --- a/mobile/lib/utils/openapi_patching.dart +++ b/mobile/lib/utils/openapi_patching.dart @@ -41,6 +41,7 @@ final Map> openApiPatches = { 'SyncAssetV1': {'isEdited': false}, 'ServerFeaturesDto': {'ocr': false, 'realtimeTranscoding': false}, 'MemoriesResponse': {'duration': 5}, + 'WorkflowResponseDto': {'logging': false}, }; void upgradeDto(dynamic value, String targetType) { diff --git a/open-api/immich-openapi-specs.json b/open-api/immich-openapi-specs.json index 3596eaf4b1..d00faab9a0 100644 --- a/open-api/immich-openapi-specs.json +++ b/open-api/immich-openapi-specs.json @@ -15873,6 +15873,15 @@ "type": "string" } }, + { + "name": "logging", + "required": false, + "in": "query", + "description": "Workflow logs run results", + "schema": { + "type": "boolean" + } + }, { "name": "name", "required": false, @@ -16186,6 +16195,95 @@ "x-immich-state": "Deprecated" } }, + "/workflows/{id}/logs": { + "get": { + "description": "Retrieve logs of a workflows runs by ID", + "operationId": "getWorkflowLogs", + "parameters": [ + { + "name": "before", + "required": false, + "in": "query", + "description": "Filter by runs before a date/time", + "schema": { + "format": "date-time", + "pattern": "^(?:(?:\\d\\d[2468][048]|\\d\\d[13579][26]|\\d\\d0[48]|[02468][048]00|[13579][26]00)-02-29|\\d{4}-(?:(?:0[13578]|1[02])-(?:0[1-9]|[12]\\d|3[01])|(?:0[469]|11)-(?:0[1-9]|[12]\\d|30)|(?:02)-(?:0[1-9]|1\\d|2[0-8])))T(?:(?:[01]\\d|2[0-3]):[0-5]\\d(?::[0-5]\\d(?:\\.\\d+)?)?(?:Z|([+-](?:[01]\\d|2[0-3]):[0-5]\\d)))$", + "example": "2024-01-01T00:00:00.000Z", + "type": "string" + } + }, + { + "name": "id", + "required": true, + "in": "path", + "schema": { + "format": "uuid", + "pattern": "^([0-9a-fA-F]{8}-[0-9a-fA-F]{4}-4[0-9a-fA-F]{3}-[89abAB][0-9a-fA-F]{3}-[0-9a-fA-F]{12})$", + "type": "string" + } + }, + { + "name": "limit", + "required": false, + "in": "query", + "description": "Maximum number of logs", + "schema": { + "maximum": 9007199254740991, + "exclusiveMinimum": true, + "default": 50, + "type": "integer", + "minimum": 0 + } + }, + { + "name": "result", + "required": false, + "in": "query", + "description": "Filter by run result", + "schema": { + "$ref": "#/components/schemas/WorkflowResult" + } + } + ], + "responses": { + "200": { + "content": { + "application/json": { + "schema": { + "items": { + "$ref": "#/components/schemas/WorkflowLogEntryDto" + }, + "type": "array" + } + } + }, + "description": "" + } + }, + "security": [ + { + "bearer": [] + }, + { + "cookie": [] + }, + { + "api_key": [] + } + ], + "summary": "Retrieve workflow logs", + "tags": [ + "Workflows" + ], + "x-immich-history": [ + { + "version": "v3.0.0", + "state": "Added" + } + ], + "x-immich-permission": "workflow.logs" + } + }, "/workflows/{id}/share": { "get": { "description": "Retrieve a workflow details without ids, default values, etc.", @@ -21057,6 +21155,7 @@ "workflow.read", "workflow.update", "workflow.delete", + "workflow.logs", "adminUser.create", "adminUser.read", "adminUser.update", @@ -27909,6 +28008,10 @@ "description": "Workflow enabled", "type": "boolean" }, + "logging": { + "description": "Workflow logs run results", + "type": "boolean" + }, "name": { "description": "Workflow name", "nullable": true, @@ -27930,6 +28033,59 @@ ], "type": "object" }, + "WorkflowLogEntryDto": { + "properties": { + "at": { + "description": "Workflow run date/time", + "example": "2024-01-01T00:00:00.000Z", + "format": "date-time", + "pattern": "^(?:(?:\\d\\d[2468][048]|\\d\\d[13579][26]|\\d\\d0[48]|[02468][048]00|[13579][26]00)-02-29|\\d{4}-(?:(?:0[13578]|1[02])-(?:0[1-9]|[12]\\d|3[01])|(?:0[469]|11)-(?:0[1-9]|[12]\\d|30)|(?:02)-(?:0[1-9]|1\\d|2[0-8])))T(?:(?:[01]\\d|2[0-3]):[0-5]\\d(?::[0-5]\\d(?:\\.\\d+)?)?(?:Z|([+-](?:[01]\\d|2[0-3]):[0-5]\\d)))$", + "type": "string" + }, + "id": { + "description": "Workflow log entry ID", + "format": "uuid", + "pattern": "^([0-9a-fA-F]{8}-[0-9a-fA-F]{4}-4[0-9a-fA-F]{3}-[89abAB][0-9a-fA-F]{3}-[0-9a-fA-F]{12})$", + "type": "string" + }, + "lastStep": { + "description": "Last step ran, if the workflow ended early", + "properties": { + "index": { + "description": "Index of the step in the workflow", + "exclusiveMinimum": true, + "maximum": 9007199254740991, + "minimum": 0, + "type": "integer" + }, + "method": { + "description": "Method of the step", + "type": "string" + } + }, + "required": [ + "method", + "index" + ], + "type": "object" + }, + "result": { + "$ref": "#/components/schemas/WorkflowResult" + }, + "triggerDataId": { + "description": "Workflow trigger data ID", + "format": "uuid", + "pattern": "^([0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[1-8][0-9a-fA-F]{3}-[89abAB][0-9a-fA-F]{3}-[0-9a-fA-F]{12}|00000000-0000-0000-0000-000000000000|ffffffff-ffff-ffff-ffff-ffffffffffff)$", + "type": "string" + } + }, + "required": [ + "at", + "id", + "result" + ], + "type": "object" + }, "WorkflowResponseDto": { "properties": { "createdAt": { @@ -27951,6 +28107,10 @@ "pattern": "^([0-9a-fA-F]{8}-[0-9a-fA-F]{4}-4[0-9a-fA-F]{3}-[89abAB][0-9a-fA-F]{3}-[0-9a-fA-F]{12})$", "type": "string" }, + "logging": { + "description": "Workflow logs run results", + "type": "boolean" + }, "name": { "description": "Workflow name", "nullable": true, @@ -27977,6 +28137,7 @@ "description", "enabled", "id", + "logging", "name", "steps", "trigger", @@ -27984,6 +28145,15 @@ ], "type": "object" }, + "WorkflowResult": { + "description": "Workflow run result", + "enum": [ + "completed", + "halted", + "error" + ], + "type": "string" + }, "WorkflowShareResponseDto": { "properties": { "description": { @@ -28109,6 +28279,10 @@ "description": "Workflow enabled", "type": "boolean" }, + "logging": { + "description": "Workflow logs run results", + "type": "boolean" + }, "name": { "description": "Workflow name", "nullable": true, diff --git a/packages/sdk/src/fetch-client.ts b/packages/sdk/src/fetch-client.ts index af24b2d21b..b2f98b58ef 100644 --- a/packages/sdk/src/fetch-client.ts +++ b/packages/sdk/src/fetch-client.ts @@ -2791,6 +2791,8 @@ export type WorkflowResponseDto = { enabled: boolean; /** Workflow ID */ id: string; + /** Workflow logs run results */ + logging: boolean; /** Workflow name */ name: string | null; /** Workflow steps */ @@ -2805,6 +2807,8 @@ export type WorkflowCreateDto = { description?: string | null; /** Workflow enabled */ enabled?: boolean; + /** Workflow logs run results */ + logging?: boolean; /** Workflow name */ name?: string | null; steps?: WorkflowStepDto[]; @@ -2822,12 +2826,30 @@ export type WorkflowUpdateDto = { description?: string | null; /** Workflow enabled */ enabled?: boolean; + /** Workflow logs run results */ + logging?: boolean; /** Workflow name */ name?: string | null; steps?: WorkflowStepDto[]; /** Workflow trigger type */ trigger?: WorkflowTrigger; }; +export type WorkflowLogEntryDto = { + /** Workflow run date/time */ + at: string; + /** Workflow log entry ID */ + id: string; + /** Last step ran, if the workflow ended early */ + lastStep?: { + /** Index of the step in the workflow */ + index: number; + /** Method of the step */ + method: string; + }; + result: WorkflowResult; + /** Workflow trigger data ID */ + triggerDataId?: string; +}; export type WorkflowShareStepDto = { /** Step configuration */ config: { @@ -6977,10 +6999,11 @@ export function getUniqueOriginalPaths(opts?: Oazapfts.RequestOpts) { /** * List all workflows */ -export function searchWorkflows({ description, enabled, id, name, trigger }: { +export function searchWorkflows({ description, enabled, id, logging, name, trigger }: { description?: string; enabled?: boolean; id?: string; + logging?: boolean; name?: string; trigger?: WorkflowTrigger; }, opts?: Oazapfts.RequestOpts) { @@ -6991,6 +7014,7 @@ export function searchWorkflows({ description, enabled, id, name, trigger }: { description, enabled, id, + logging, name, trigger }))}`, { @@ -7063,6 +7087,26 @@ export function updateWorkflow({ id, workflowUpdateDto }: { body: workflowUpdateDto }))); } +/** + * Retrieve workflow logs + */ +export function getWorkflowLogs({ before, id, limit, result }: { + before?: string; + id: string; + limit?: number; + result?: WorkflowResult; +}, opts?: Oazapfts.RequestOpts) { + return oazapfts.ok(oazapfts.fetchJson<{ + status: 200; + data: WorkflowLogEntryDto[]; + }>(`/workflows/${encodeURIComponent(id)}/logs${QS.query(QS.explode({ + before, + limit, + result + }))}`, { + ...opts + })); +} /** * Retrieve a workflow */ @@ -7310,6 +7354,7 @@ export enum Permission { WorkflowRead = "workflow.read", WorkflowUpdate = "workflow.update", WorkflowDelete = "workflow.delete", + WorkflowLogs = "workflow.logs", AdminUserCreate = "adminUser.create", AdminUserRead = "adminUser.read", AdminUserUpdate = "adminUser.update", @@ -7687,6 +7732,11 @@ export enum AssetOrderBy { TakenAt = "takenAt", CreatedAt = "createdAt" } +export enum WorkflowResult { + Completed = "completed", + Halted = "halted", + Error = "error" +} export enum ReleaseType { Major = "major", Premajor = "premajor", diff --git a/server/src/controllers/workflow.controller.ts b/server/src/controllers/workflow.controller.ts index b8d9fd2b93..b0d025763e 100644 --- a/server/src/controllers/workflow.controller.ts +++ b/server/src/controllers/workflow.controller.ts @@ -4,6 +4,8 @@ import { Endpoint, HistoryBuilder } from 'src/decorators'; import { AuthDto } from 'src/dtos/auth.dto'; import { WorkflowCreateDto, + WorkflowGetLogsDto, + WorkflowLogEntryDto, WorkflowResponseDto, WorkflowSearchDto, WorkflowShareResponseDto, @@ -113,4 +115,19 @@ export class WorkflowController { deleteWorkflow(@Auth() auth: AuthDto, @Param() { id }: UUIDParamDto): Promise { return this.service.delete(auth, id); } + + @Get(':id/logs') + @Authenticated({ permission: Permission.WorkflowLogs }) + @Endpoint({ + summary: 'Retrieve workflow logs', + description: 'Retrieve logs of a workflows runs by ID', + history: HistoryBuilder.v3(), + }) + getWorkflowLogs( + @Auth() auth: AuthDto, + @Param() { id }: UUIDParamDto, + @Query() dto: WorkflowGetLogsDto, + ): Promise { + return this.service.getLogs(auth, id, dto); + } } diff --git a/server/src/dtos/workflow.dto.ts b/server/src/dtos/workflow.dto.ts index b4269ea0c9..990cae79a1 100644 --- a/server/src/dtos/workflow.dto.ts +++ b/server/src/dtos/workflow.dto.ts @@ -1,6 +1,7 @@ import type { WorkflowStepConfig, WorkflowTrigger } from '@immich/plugin-sdk'; import { createZodDto } from 'nestjs-zod'; -import { WorkflowTriggerSchema, WorkflowTypeSchema } from 'src/enum'; +import { WorkflowResultSchema, WorkflowTriggerSchema, WorkflowTypeSchema } from 'src/enum'; +import { isoDatetimeToDate } from 'src/validation'; import z from 'zod'; const WorkflowTriggerResponseSchema = z @@ -17,6 +18,7 @@ const WorkflowSearchSchema = z name: z.string().optional().describe('Workflow name'), description: z.string().optional().describe('Workflow description'), enabled: z.boolean().optional().describe('Workflow enabled'), + logging: z.boolean().optional().describe('Workflow logs run results'), }) .meta({ id: 'WorkflowSearchDto' }); @@ -42,6 +44,7 @@ const WorkflowCreateSchema = z name: z.string().nullable().optional().describe('Workflow name'), description: z.string().nullable().optional().describe('Workflow description'), enabled: z.boolean().optional().describe('Workflow enabled'), + logging: z.boolean().optional().describe('Workflow logs run results'), steps: z.array(WorkflowStepSchema).optional(), }) .meta({ id: 'WorkflowCreateDto' }); @@ -52,6 +55,7 @@ const WorkflowUpdateSchema = z name: z.string().nullable().optional().describe('Workflow name'), description: z.string().nullable().optional().describe('Workflow description'), enabled: z.boolean().optional().describe('Workflow enabled'), + logging: z.boolean().optional().describe('Workflow logs run results'), steps: z.array(WorkflowStepSchema).optional(), }) .meta({ id: 'WorkflowUpdateDto' }); @@ -65,6 +69,7 @@ const WorkflowResponseSchema = z createdAt: z.string().describe('Creation date'), updatedAt: z.string().describe('Update date'), enabled: z.boolean().describe('Workflow enabled'), + logging: z.boolean().describe('Workflow logs run results'), steps: z.array(WorkflowStepSchema).describe('Workflow steps'), }) .meta({ id: 'WorkflowResponseDto' }); @@ -78,12 +83,36 @@ const WorkflowShareResponseSchema = z }) .meta({ id: 'WorkflowShareResponseDto' }); +const WorkflowLogEntrySchema = z + .object({ + id: z.uuidv4().describe('Workflow log entry ID'), + at: isoDatetimeToDate.describe('Workflow run date/time'), + result: WorkflowResultSchema.describe('Workflow run result'), + triggerDataId: z.uuid().optional().describe('Workflow trigger data ID'), + lastStep: z + .object({ + method: z.string().describe('Method of the step'), + index: z.int().positive().describe('Index of the step in the workflow'), + }) + .optional() + .describe('Last step ran, if the workflow ended early'), + }) + .meta({ id: 'WorkflowLogEntryDto' }); + +const WorkflowGetLogsSchema = z.object({ + result: WorkflowResultSchema.optional().describe('Filter by run result'), + before: isoDatetimeToDate.optional().describe('Filter by runs before a date/time'), + limit: z.coerce.number().int().positive().default(50).describe('Maximum number of logs'), +}); + export class WorkflowTriggerResponseDto extends createZodDto(WorkflowTriggerResponseSchema) {} export class WorkflowSearchDto extends createZodDto(WorkflowSearchSchema) {} export class WorkflowCreateDto extends createZodDto(WorkflowCreateSchema) {} export class WorkflowUpdateDto extends createZodDto(WorkflowUpdateSchema) {} export class WorkflowResponseDto extends createZodDto(WorkflowResponseSchema) {} export class WorkflowShareResponseDto extends createZodDto(WorkflowShareResponseSchema) {} +export class WorkflowLogEntryDto extends createZodDto(WorkflowLogEntrySchema) {} +export class WorkflowGetLogsDto extends createZodDto(WorkflowGetLogsSchema) {} type Workflow = { id: string; @@ -93,6 +122,7 @@ type Workflow = { name: string | null; description: string | null; enabled: boolean; + logging: boolean; }; type WorkflowStep = { @@ -107,6 +137,7 @@ export const mapWorkflow = (workflow: Workflow & { steps: WorkflowStep[] }): Wor id: workflow.id, enabled: workflow.enabled, trigger: workflow.trigger, + logging: workflow.logging, name: workflow.name, description: workflow.description, createdAt: workflow.createdAt.toISOString(), diff --git a/server/src/enum.ts b/server/src/enum.ts index 0d29244e09..8882b554f3 100644 --- a/server/src/enum.ts +++ b/server/src/enum.ts @@ -298,6 +298,7 @@ export enum Permission { WorkflowRead = 'workflow.read', WorkflowUpdate = 'workflow.update', WorkflowDelete = 'workflow.delete', + WorkflowLogs = 'workflow.logs', AdminUserCreate = 'adminUser.create', AdminUserRead = 'adminUser.read', @@ -1235,6 +1236,17 @@ export enum CalendarHeatmapType { Taken = 'Taken', } +export enum WorkflowResult { + Completed = 'completed', + Halted = 'halted', + Error = 'error', +} + +export const WorkflowResultSchema = z + .enum(WorkflowResult) + .describe('Workflow run result') + .meta({ id: 'WorkflowResult' }); + export enum SearchOrderField { FileCreatedAt = 'fileCreatedAt', LocalDateTime = 'localDateTime', diff --git a/server/src/queries/workflow.repository.sql b/server/src/queries/workflow.repository.sql index f44ba54066..293d053d81 100644 --- a/server/src/queries/workflow.repository.sql +++ b/server/src/queries/workflow.repository.sql @@ -9,6 +9,7 @@ select "workflow"."enabled", "workflow"."createdAt", "workflow"."updatedAt", + "workflow"."logging", ( select coalesce(json_agg(agg), '[]') @@ -43,6 +44,7 @@ select "workflow"."enabled", "workflow"."createdAt", "workflow"."updatedAt", + "workflow"."logging", ( select coalesce(json_agg(agg), '[]') @@ -73,6 +75,7 @@ select "workflow"."id", "workflow"."name", "workflow"."trigger", + "workflow"."logging", ( select coalesce(json_agg(agg), '[]') @@ -100,6 +103,39 @@ where "id" = $2 and "enabled" = $3 +-- WorkflowRepository.getLogs +select + "workflow_log"."id", + "workflow_log"."createdAt", + "workflow_log"."result", + "workflow_log"."workflowId", + "workflow_log"."workflowStepId", + "workflow_log"."triggerDataId", + ( + select + to_json(obj) + from + ( + select + "plugin_method"."pluginId", + "plugin_method"."name" as "methodName", + "workflow_step"."order" + from + "workflow_step" + inner join "plugin_method" on "plugin_method"."id" = "workflow_step"."pluginMethodId" + where + "workflow_step"."id" = "workflow_log"."workflowStepId" + ) as obj + ) as "step" +from + "workflow_log" +where + "workflow_log"."workflowId" = $1 +order by + "workflow_log"."createdAt" desc +limit + $2 + -- WorkflowRepository.delete delete from "workflow" where diff --git a/server/src/repositories/workflow.repository.ts b/server/src/repositories/workflow.repository.ts index 888ca94440..8344896199 100644 --- a/server/src/repositories/workflow.repository.ts +++ b/server/src/repositories/workflow.repository.ts @@ -4,8 +4,9 @@ import { jsonArrayFrom, jsonObjectFrom } from 'kysely/helpers/postgres'; import { InjectKysely } from 'nestjs-kysely'; import { columns } from 'src/database'; import { DummyValue, GenerateSql } from 'src/decorators'; -import { WorkflowSearchDto } from 'src/dtos/workflow.dto'; +import { WorkflowGetLogsDto, WorkflowSearchDto } from 'src/dtos/workflow.dto'; import { DB } from 'src/schema'; +import { WorkflowLogTable } from 'src/schema/tables/workflow-log.table'; import { WorkflowStepTable } from 'src/schema/tables/workflow-step.table'; import { WorkflowTable } from 'src/schema/tables/workflow.table'; import { withTags } from 'src/utils/database'; @@ -27,6 +28,7 @@ export class WorkflowRepository { 'workflow.enabled', 'workflow.createdAt', 'workflow.updatedAt', + 'workflow.logging', ]) .select((eb) => [ jsonArrayFrom( @@ -66,7 +68,7 @@ export class WorkflowRepository { getForWorkflowRun(id: string) { return this.db .selectFrom('workflow') - .select(['workflow.id', 'workflow.name', 'workflow.trigger']) + .select(['workflow.id', 'workflow.name', 'workflow.trigger', 'workflow.logging']) .select((eb) => [ jsonArrayFrom( eb @@ -99,6 +101,9 @@ export class WorkflowRepository { update(id: string, dto: Updateable, steps?: WorkflowStepUpsert[]) { return this.db.transaction().execute(async (tx) => { + if (dto.logging === false) { + await tx.deleteFrom('workflow_log').where('workflowId', '=', id).execute(); + } if (Object.values(dto).some((prop) => prop !== undefined)) { await tx.updateTable('workflow').set(dto).where('id', '=', id).executeTakeFirstOrThrow(); } @@ -106,6 +111,39 @@ export class WorkflowRepository { }); } + @GenerateSql({ params: [DummyValue.UUID, { result: undefined }] }) + getLogs(id: string, dto: WorkflowGetLogsDto) { + return this.db + .selectFrom('workflow_log') + .select([ + 'workflow_log.id', + 'workflow_log.createdAt', + 'workflow_log.result', + 'workflow_log.workflowId', + 'workflow_log.workflowStepId', + 'workflow_log.triggerDataId', + ]) + .where('workflow_log.workflowId', '=', id) + .select((eb) => [ + jsonObjectFrom( + eb + .selectFrom('workflow_step') + .whereRef('workflow_step.id', '=', 'workflow_log.workflowStepId') + .innerJoin('plugin_method', 'plugin_method.id', 'workflow_step.pluginMethodId') + .select(['plugin_method.pluginId', 'plugin_method.name as methodName', 'workflow_step.order']), + ).as('step'), + ]) + .$if(dto.result !== undefined, (qb) => qb.where('workflow_log.result', '=', dto.result!)) + .$if(dto.before !== undefined, (qb) => qb.where('workflow_log.createdAt', '<', dto.before!)) + .orderBy('workflow_log.createdAt', 'desc') + .limit(dto.limit) + .execute(); + } + + log(dto: Insertable) { + return this.db.insertInto('workflow_log').values(dto).execute(); + } + async updateStep(id: string, dto: Updateable) { await this.db.updateTable('workflow_step').where('workflow_step.id', '=', id).set(dto).execute(); } diff --git a/server/src/schema/index.ts b/server/src/schema/index.ts index 171fefb50f..45423f7c40 100644 --- a/server/src/schema/index.ts +++ b/server/src/schema/index.ts @@ -87,6 +87,7 @@ import { VideoStreamSessionTable, VideoStreamVariantTable, } from 'src/schema/tables/video-stream.table'; +import { WorkflowLogTable } from 'src/schema/tables/workflow-log.table'; import { WorkflowStepTable } from 'src/schema/tables/workflow-step.table'; import { WorkflowTable } from 'src/schema/tables/workflow.table'; @@ -279,4 +280,5 @@ export interface DB { workflow: WorkflowTable; workflow_step: WorkflowStepTable; + workflow_log: WorkflowLogTable; } diff --git a/server/src/schema/migrations/1783930557118-AddWorkflowLogsTable.ts b/server/src/schema/migrations/1783930557118-AddWorkflowLogsTable.ts new file mode 100644 index 0000000000..13fa5f5f97 --- /dev/null +++ b/server/src/schema/migrations/1783930557118-AddWorkflowLogsTable.ts @@ -0,0 +1,24 @@ +import { Kysely, sql } from 'kysely'; + +export async function up(db: Kysely): Promise { + await sql`ALTER TABLE "workflow" ADD "logging" boolean NOT NULL DEFAULT false;`.execute(db); + await sql`CREATE TABLE "workflow_log" ( + "id" uuid NOT NULL DEFAULT uuid_generate_v4(), + "createdAt" timestamp with time zone NOT NULL DEFAULT now(), + "workflowId" uuid NOT NULL, + "result" character varying NOT NULL, + "workflowStepId" uuid, + "triggerDataId" uuid, + "runId" uuid NOT NULL, + CONSTRAINT "workflow_log_workflowId_fkey" FOREIGN KEY ("workflowId") REFERENCES "workflow" ("id") ON UPDATE CASCADE ON DELETE CASCADE, + CONSTRAINT "workflow_log_workflowStepId_fkey" FOREIGN KEY ("workflowStepId") REFERENCES "workflow_step" ("id") ON UPDATE CASCADE ON DELETE SET NULL, + CONSTRAINT "workflow_log_pkey" PRIMARY KEY ("id") +);`.execute(db); + await sql`CREATE INDEX "workflow_log_workflowId_idx" ON "workflow_log" ("workflowId");`.execute(db); + await sql`CREATE INDEX "workflow_log_workflowStepId_idx" ON "workflow_log" ("workflowStepId");`.execute(db); +} + +export async function down(db: Kysely): Promise { + await sql`ALTER TABLE "workflow" DROP COLUMN "logging";`.execute(db); + await sql`DROP TABLE "workflow_log";`.execute(db); +} diff --git a/server/src/schema/tables/workflow-log.table.ts b/server/src/schema/tables/workflow-log.table.ts new file mode 100644 index 0000000000..e720cb2a37 --- /dev/null +++ b/server/src/schema/tables/workflow-log.table.ts @@ -0,0 +1,36 @@ +import { + Column, + CreateDateColumn, + ForeignKeyColumn, + Generated, + PrimaryGeneratedColumn, + Table, + Timestamp, +} from '@immich/sql-tools'; +import { WorkflowResult } from 'src/enum'; +import { WorkflowStepTable } from 'src/schema/tables/workflow-step.table'; +import { WorkflowTable } from 'src/schema/tables/workflow.table'; + +@Table('workflow_log') +export class WorkflowLogTable { + @PrimaryGeneratedColumn() + id!: Generated; + + @CreateDateColumn() + createdAt!: Generated; + + @ForeignKeyColumn(() => WorkflowTable, { onUpdate: 'CASCADE', onDelete: 'CASCADE', index: true }) + workflowId!: string; + + @Column() + result!: WorkflowResult; + + @ForeignKeyColumn(() => WorkflowStepTable, { onDelete: 'SET NULL', onUpdate: 'CASCADE', nullable: true }) + workflowStepId!: string | null; + + @Column({ type: 'uuid', nullable: true }) + triggerDataId!: string | null; + + @Column({ type: 'uuid' }) + runId!: string; +} diff --git a/server/src/schema/tables/workflow.table.ts b/server/src/schema/tables/workflow.table.ts index 944fffd9d5..c0c9dc94e4 100644 --- a/server/src/schema/tables/workflow.table.ts +++ b/server/src/schema/tables/workflow.table.ts @@ -41,4 +41,7 @@ export class WorkflowTable { @Column({ type: 'boolean', default: true }) enabled!: Generated; + + @Column({ type: 'boolean', default: false }) + logging!: Generated; } diff --git a/server/src/services/workflow-execution.service.ts b/server/src/services/workflow-execution.service.ts index 4995eada9c..193fab16de 100644 --- a/server/src/services/workflow-execution.service.ts +++ b/server/src/services/workflow-execution.service.ts @@ -22,6 +22,7 @@ import { JobName, JobStatus, QueueName, + WorkflowResult, WorkflowType, } from 'src/enum'; import { ArgOf } from 'src/repositories/event.repository'; @@ -38,7 +39,7 @@ const dummy = () => { }; type ExecuteOptions = { - read: (type: T) => Promise<{ authUserId: string; data: WorkflowEventData }>; + read: (type: T) => Promise<{ authUserId: string; data: WorkflowEventData; entityId?: string }>; write: (auth: AuthDto, changes: WorkflowChanges) => Promise; }; @@ -345,6 +346,7 @@ export class WorkflowExecutionService extends BaseService { return { data: { asset } as any, authUserId: asset.ownerId, + entityId: asset.id, }; }, write: async (auth, changes) => { @@ -412,11 +414,13 @@ export class WorkflowExecutionService extends BaseService { return; } - try { - const { read, write } = handler; - const readResult = await read(type); - let data = readResult.data; - for (const step of workflow.steps) { + const { read, write } = handler; + const readResult = await read(type); + let data = readResult.data; + const runId = crypto.randomUUID(); + + for (const step of workflow.steps) { + try { const payload: WorkflowEventPayload = { trigger: workflow.trigger, type, @@ -467,14 +471,45 @@ export class WorkflowExecutionService extends BaseService { const shouldContinue = result?.workflow?.continue ?? true; if (!shouldContinue) { - break; - } - } + if (workflow.logging) { + await this.workflowRepository.log({ + workflowId, + result: WorkflowResult.Halted, + workflowStepId: step.id, + triggerDataId: readResult.entityId, + runId, + }); + } - this.logger.debug(`Workflow ${workflowId} executed successfully`); - } catch (error) { - this.logger.error(`Error executing workflow ${workflowId}:`, error); - return JobStatus.Failed; + this.logger.debug(`Workflow ${workflowId} run ${runId} stopped on step ${step.id}`); + return; + } + } catch (error) { + this.logger.error(`Error executing workflow ${workflowId} run ${runId}:`, error); + + if (workflow.logging) { + await this.workflowRepository.log({ + workflowId, + result: WorkflowResult.Error, + workflowStepId: step.id, + triggerDataId: readResult.entityId, + runId, + }); + } + + return JobStatus.Failed; + } } + + if (workflow.logging) { + await this.workflowRepository.log({ + workflowId, + result: WorkflowResult.Completed, + triggerDataId: readResult.entityId, + runId, + }); + } + + this.logger.debug(`Workflow ${workflowId} run ${runId} executed successfully`); } } diff --git a/server/src/services/workflow.service.ts b/server/src/services/workflow.service.ts index 1d4ecb89a0..e68e6aa94e 100644 --- a/server/src/services/workflow.service.ts +++ b/server/src/services/workflow.service.ts @@ -5,6 +5,8 @@ import { mapWorkflow, mapWorkflowShare, WorkflowCreateDto, + WorkflowGetLogsDto, + WorkflowLogEntryDto, WorkflowResponseDto, WorkflowSearchDto, WorkflowShareResponseDto, @@ -83,6 +85,23 @@ export class WorkflowService extends BaseService { await this.workflowRepository.delete(id); } + async getLogs(auth: AuthDto, id: string, dto: WorkflowGetLogsDto): Promise { + await this.requireAccess({ auth, permission: Permission.WorkflowLogs, ids: [id] }); + const logs = await this.workflowRepository.getLogs(id, dto); + return logs.map((entry) => ({ + id: entry.id, + at: entry.createdAt, + result: entry.result, + triggerDataId: entry.triggerDataId ?? undefined, + lastStep: entry.step + ? { + index: entry.step.order, + method: `${entry.step.pluginId}#${entry.step.methodName}`, + } + : undefined, + })); + } + private async resolveAndValidateSteps(steps: T[], trigger: WorkflowTrigger) { const methods = await this.pluginRepository.getForValidation(); const results: Array = []; diff --git a/server/src/utils/access.ts b/server/src/utils/access.ts index ce8e614565..0d6f4b4eeb 100644 --- a/server/src/utils/access.ts +++ b/server/src/utils/access.ts @@ -324,7 +324,8 @@ const checkOtherAccess = async (access: AccessRepository, request: OtherAccessRe case Permission.WorkflowRead: case Permission.WorkflowUpdate: - case Permission.WorkflowDelete: { + case Permission.WorkflowDelete: + case Permission.WorkflowLogs: { return access.workflow.checkOwnerAccess(auth.user.id, ids); } diff --git a/web/src/lib/modals/WorkflowLogsModal.svelte b/web/src/lib/modals/WorkflowLogsModal.svelte new file mode 100644 index 0000000000..b897d4b437 --- /dev/null +++ b/web/src/lib/modals/WorkflowLogsModal.svelte @@ -0,0 +1,171 @@ + + + + + {#if workflow.logging} + + + {$t('date')} + {$t('result')} + + + {#each entries as entry (entry.id)} + + + +

+ {DateTime.fromISO(entry.at).toLocaleString(DateTime.DATETIME_MED, { + locale: $locale, + })} +

+ {#if entry.triggerDataId} + + + + {/if} +
+
+ + + {#if entry.result === WorkflowResult.Completed} +

+ {$t('workflow_logging_completed')} +

+ {:else if entry.result === WorkflowResult.Halted} +

+ {#if entry.lastStep} + {$t('workflow_logging_halted_step', { values: { step: entry.lastStep.index + 1 } })} + {:else} + {$t('workflow_logging_halted')} + {/if} +

+ {:else} +

+ {#if entry.lastStep} + {$t('workflow_logging_error_step', { values: { step: entry.lastStep.index + 1 } })} + {:else} + {$t('error')} + {/if} +

+ {/if} +
+
+
+ {/each} + {#if hasNext} +
...
+ {/if} +
+
+
+