From 223f4986f9b076ef73621ccb68c44e2b19446cc5 Mon Sep 17 00:00:00 2001 From: Ben Beckford Date: Wed, 29 Jul 2026 13:50:58 -0700 Subject: [PATCH 1/4] wip(server): asset added to album workflow trigger --- i18n/en.json | 2 + open-api/immich-openapi-specs.json | 7 +- packages/plugin-sdk/src/types.ts | 8 ++ packages/sdk/src/fetch-client.ts | 6 +- server/src/enum.ts | 2 + server/src/queries/workflow.repository.sql | 2 + server/src/repositories/event.repository.ts | 1 + .../src/repositories/workflow.repository.ts | 1 + server/src/services/album.service.ts | 17 ++- .../services/workflow-execution.service.ts | 125 +++++++++++++----- server/src/types.ts | 1 + server/src/utils/workflow.ts | 6 +- web/src/lib/utils/workflow.ts | 6 + 13 files changed, 145 insertions(+), 39 deletions(-) diff --git a/i18n/en.json b/i18n/en.json index 1dbd50e294..86a3bba166 100644 --- a/i18n/en.json +++ b/i18n/en.json @@ -2136,6 +2136,8 @@ "trash_page_info": "Trashed items will be permanently deleted after {days} days", "trashed_items_will_be_permanently_deleted_after": "Trashed items will be permanently deleted after {days, plural, one {# day} other {# days}}.", "trigger": "Trigger", + "trigger_album_asset_added": "Asset Added to Album", + "trigger_album_asset_added_description": "Triggered when an asset is added to an album", "trigger_asset_metadata_extraction": "Asset Metadata Extraction", "trigger_asset_metadata_extraction_description": "Triggered when the EXIF metadata of an asset is extracted", "trigger_asset_uploaded": "Asset Upload", diff --git a/open-api/immich-openapi-specs.json b/open-api/immich-openapi-specs.json index dd5d4dce45..ff7b6eb698 100644 --- a/open-api/immich-openapi-specs.json +++ b/open-api/immich-openapi-specs.json @@ -19316,6 +19316,7 @@ "OcrQueueAll", "Ocr", "WorkflowAssetTrigger", + "WorkflowAlbumAssetTrigger", "IntegrityUntrackedFilesQueueAll", "IntegrityUntrackedFiles", "IntegrityUntrackedRefresh", @@ -28031,7 +28032,8 @@ "description": "Plugin trigger type", "enum": [ "AssetCreate", - "AssetMetadataExtraction" + "AssetMetadataExtraction", + "AlbumAssetAdded" ], "type": "string" }, @@ -28058,7 +28060,8 @@ "WorkflowType": { "description": "Workflow type", "enum": [ - "AssetV1" + "AssetV1", + "AlbumAssetV1" ], "type": "string" }, diff --git a/packages/plugin-sdk/src/types.ts b/packages/plugin-sdk/src/types.ts index 9756d652c3..ce4c81f84b 100644 --- a/packages/plugin-sdk/src/types.ts +++ b/packages/plugin-sdk/src/types.ts @@ -10,6 +10,7 @@ type DeepPartial = T extends Date export type WorkflowEventMap = { [WorkflowType.AssetV1]: AssetV1; + [WorkflowType.AlbumAssetV1]: AlbumAssetV1; // [WorkflowType.AssetPersonV1]: AssetPersonV1; } & { [K in WorkflowType]: unknown }; @@ -18,6 +19,7 @@ export type WorkflowEventData = WorkflowEventMap[T]; export enum WorkflowTrigger { AssetCreate = 'AssetCreate', AssetMetadataExtraction = 'AssetMetadataExtraction', + AlbumAssetAdded = 'AlbumAssetAdded', // PersonRecognized = 'PersonRecognized', } @@ -123,6 +125,12 @@ export type AssetV1 = { }; }; +export type AlbumAssetV1 = AssetV1 & { + album: { + id: string; + }; +}; + // export type AssetPersonV1 = AssetV1 & { // person: { // id: string; diff --git a/packages/sdk/src/fetch-client.ts b/packages/sdk/src/fetch-client.ts index a7e0d9511e..aab9199c9d 100644 --- a/packages/sdk/src/fetch-client.ts +++ b/packages/sdk/src/fetch-client.ts @@ -7419,11 +7419,13 @@ export enum PartnerDirection { SharedWith = "shared-with" } export enum WorkflowType { - AssetV1 = "AssetV1" + AssetV1 = "AssetV1", + AlbumAssetV1 = "AlbumAssetV1" } export enum WorkflowTrigger { AssetCreate = "AssetCreate", - AssetMetadataExtraction = "AssetMetadataExtraction" + AssetMetadataExtraction = "AssetMetadataExtraction", + AlbumAssetAdded = "AlbumAssetAdded" } export enum QueueJobStatus { Active = "active", diff --git a/server/src/enum.ts b/server/src/enum.ts index 0d29244e09..66e36b7121 100644 --- a/server/src/enum.ts +++ b/server/src/enum.ts @@ -903,6 +903,7 @@ export enum JobName { // Workflow WorkflowAssetTrigger = 'WorkflowAssetTrigger', + WorkflowAlbumAssetTrigger = 'WorkflowAlbumAssetTrigger', // Integrity IntegrityUntrackedFilesQueueAll = 'IntegrityUntrackedFilesQueueAll', @@ -1225,6 +1226,7 @@ export const WorkflowTriggerSchema = z export enum WorkflowType { AssetV1 = 'AssetV1', + AlbumAssetV1 = 'AlbumAssetV1', // AssetPersonV1 = 'AssetPersonV1', } diff --git a/server/src/queries/workflow.repository.sql b/server/src/queries/workflow.repository.sql index f44ba54066..78ab21385f 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"."ownerId", ( select coalesce(json_agg(agg), '[]') @@ -43,6 +44,7 @@ select "workflow"."enabled", "workflow"."createdAt", "workflow"."updatedAt", + "workflow"."ownerId", ( select coalesce(json_agg(agg), '[]') diff --git a/server/src/repositories/event.repository.ts b/server/src/repositories/event.repository.ts index 416f823952..7a92446c77 100644 --- a/server/src/repositories/event.repository.ts +++ b/server/src/repositories/event.repository.ts @@ -40,6 +40,7 @@ type EventMap = { // album events AlbumUpdate: [{ id: string; userIds: string[]; recipientIds: string[] }]; AlbumInvite: [{ id: string; userId: string; senderName: string }]; + AlbumAssetsAdded: [{ albumId: string; userIds: string[]; recipientIds: string[]; assetIds: string[] }]; // asset events AssetCreate: [{ asset: Pick; file?: UploadFile }]; diff --git a/server/src/repositories/workflow.repository.ts b/server/src/repositories/workflow.repository.ts index 9c946fe583..9c4e435d00 100644 --- a/server/src/repositories/workflow.repository.ts +++ b/server/src/repositories/workflow.repository.ts @@ -26,6 +26,7 @@ export class WorkflowRepository { 'workflow.enabled', 'workflow.createdAt', 'workflow.updatedAt', + 'workflow.ownerId', ]) .select((eb) => [ jsonArrayFrom( diff --git a/server/src/services/album.service.ts b/server/src/services/album.service.ts index 071c5d5aec..313c8affb1 100644 --- a/server/src/services/album.service.ts +++ b/server/src/services/album.service.ts @@ -193,6 +193,12 @@ export class AlbumService extends BaseService { const userIds = album.albumUsers.map(({ user }) => user.id); const recipientIds = userIds.filter((userId) => userId !== auth.user.id); await this.eventRepository.emit('AlbumUpdate', { id, userIds, recipientIds }); + await this.eventRepository.emit('AlbumAssetsAdded', { + albumId: id, + userIds, + recipientIds, + assetIds: results.filter((a) => a.success).map((a) => a.id), + }); } return results; @@ -221,7 +227,8 @@ export class AlbumService extends BaseService { } const albumAssetValues: { albumId: string; assetId: string }[] = []; - const events: { id: string; userIds: string[]; recipientIds: string[] }[] = []; + const updateEvents: { id: string; userIds: string[]; recipientIds: string[] }[] = []; + const addedEvents: { albumId: string; userIds: string[]; recipientIds: string[]; assetIds: string[] }[] = []; for (const albumId of allowedAlbumIds) { const existingAssetIds = await this.albumRepository.getAssetIds(albumId, [...allowedAssetIds]); const notPresentAssetIds = [...allowedAssetIds.difference(existingAssetIds)]; @@ -246,13 +253,17 @@ export class AlbumService extends BaseService { ); const userIds = album.albumUsers.map(({ user }) => user.id); const recipientIds = userIds.filter((userId) => userId !== auth.user.id); - events.push({ id: albumId, userIds, recipientIds }); + updateEvents.push({ id: albumId, userIds, recipientIds }); + addedEvents.push({ albumId, userIds, recipientIds, assetIds: notPresentAssetIds }); } await this.albumRepository.addAssetIdsToAlbums(albumAssetValues); - for (const event of events) { + for (const event of updateEvents) { await this.eventRepository.emit('AlbumUpdate', event); } + for (const event of addedEvents) { + await this.eventRepository.emit('AlbumAssetsAdded', event); + } return results; } diff --git a/server/src/services/workflow-execution.service.ts b/server/src/services/workflow-execution.service.ts index 318f36a4b6..ebe12523a2 100644 --- a/server/src/services/workflow-execution.service.ts +++ b/server/src/services/workflow-execution.service.ts @@ -27,7 +27,7 @@ import { ArgOf } from 'src/repositories/event.repository'; import { AlbumService } from 'src/services/album.service'; import { AssetService } from 'src/services/asset.service'; import { BaseService } from 'src/services/base.service'; -import { JobOf } from 'src/types'; +import { JobItem, JobOf } from 'src/types'; const dummy = () => { throw new Error( @@ -42,6 +42,8 @@ type ExecuteOptions = { type AssetTrigger = { userId: string; assetId: string; trigger: WorkflowTrigger }; +type AlbumAssetTrigger = { userIds: string[]; albumId: string; assetIds: string[]; trigger: WorkflowTrigger }; + type HostContext = { allowedHosts: string[]; }; @@ -309,6 +311,16 @@ export class WorkflowExecutionService extends BaseService { return this.onAssetTrigger({ userId, assetId, trigger: WorkflowTrigger.AssetMetadataExtraction }); } + @OnEvent({ name: 'AlbumAssetsAdded' }) + onAlbumAssetsAdded({ albumId, userIds, recipientIds, assetIds }: ArgOf<'AlbumAssetsAdded'>) { + return this.onAlbumAssetTrigger({ + albumId, + userIds: [...userIds, ...recipientIds], + assetIds, + trigger: WorkflowTrigger.AlbumAssetAdded, + }); + } + private async onAssetTrigger({ userId, assetId, trigger }: AssetTrigger) { const items = await this.workflowRepository.search({ userId, trigger }); await this.jobRepository.queueAll( @@ -319,11 +331,64 @@ export class WorkflowExecutionService extends BaseService { ); } + private async onAlbumAssetTrigger({ albumId, userIds, assetIds, trigger }: AlbumAssetTrigger) { + let jobs: JobItem[] = []; + for (const userId of userIds) { + const items = await this.workflowRepository.search({ userId, trigger }); + if (!items.length) { + continue; + } + + let batch: JobItem[] = items.flatMap(({ id: workflowId }) => + assetIds.map((assetId) => ({ + name: JobName.WorkflowAlbumAssetTrigger, + data: { workflowId, assetId, albumId, trigger }, + })), + ); + jobs.push(...batch); + } + + await this.jobRepository.queueAll(jobs); + } + + private writeAssetV1(assetId: string, type: WorkflowType) { + const assetService = BaseService.create(AssetService, this); + + return async (auth: AuthDto, changes: WorkflowChanges) => { + const asset = changes.asset; + if (!asset) { + return; + } + + await assetService.update(auth, assetId, { + isFavorite: asset.isFavorite, + visibility: asset.visibility, + dateTimeOriginal: asset.exifInfo?.dateTimeOriginal ?? undefined, + // TODO allow setting to null + longitude: asset.exifInfo?.longitude ?? undefined, + // TODO allow setting to null + latitude: asset.exifInfo?.latitude ?? undefined, + // TODO allow setting to null + description: asset.exifInfo?.description ?? undefined, + rating: asset.exifInfo?.rating, + + // TODO add to update dto + // make: asset.exifInfo?.make, + // model: asset.exifInfo?.model, + // city: asset.exifInfo?.city, + // state: asset.exifInfo?.state, + // country: asset.exifInfo?.country, + // lensModel: asset.exifInfo?.lensModel, + // fNumber: asset.exifInfo?.fNumber, + // fps: asset.exifInfo?.fps, + // iso: asset.exifInfo?.iso, + }); + }; + } + @OnJob({ name: JobName.WorkflowAssetTrigger, queue: QueueName.Workflow }) handleAssetTrigger({ workflowId, assetId }: JobOf) { return this.execute(workflowId, (type) => { - const assetService = BaseService.create(AssetService, this); - switch (type) { case WorkflowType.AssetV1: { return { @@ -334,35 +399,35 @@ export class WorkflowExecutionService extends BaseService { authUserId: asset.ownerId, }; }, + write: this.writeAssetV1(assetId, type), + } satisfies ExecuteOptions; + } + } + }); + } + + @OnJob({ name: JobName.WorkflowAlbumAssetTrigger, queue: QueueName.Workflow }) + handleAlbumAssetTrigger({ workflowId, assetId, albumId }: JobOf) { + return this.execute(workflowId, (type) => { + switch (type) { + case WorkflowType.AssetV1: { + return { + read: async () => { + const asset = await this.workflowRepository.getForAssetV1(assetId); + const workflow = await this.workflowRepository.get(workflowId); + + return { + data: { asset, album: { id: albumId } } as any, + authUserId: workflow!.ownerId, + }; + }, write: async (auth, changes) => { - const asset = changes.asset; - if (!asset) { - return; + const asset = await this.workflowRepository.getForAssetV1(assetId); + const workflow = await this.workflowRepository.get(workflowId); + + if (asset.ownerId === workflow!.ownerId) { + await this.writeAssetV1(assetId, type)(auth, changes); } - - await assetService.update(auth, assetId, { - isFavorite: asset.isFavorite, - visibility: asset.visibility, - dateTimeOriginal: asset.exifInfo?.dateTimeOriginal ?? undefined, - // TODO allow setting to null - longitude: asset.exifInfo?.longitude ?? undefined, - // TODO allow setting to null - latitude: asset.exifInfo?.latitude ?? undefined, - // TODO allow setting to null - description: asset.exifInfo?.description ?? undefined, - rating: asset.exifInfo?.rating, - - // TODO add to update dto - // make: asset.exifInfo?.make, - // model: asset.exifInfo?.model, - // city: asset.exifInfo?.city, - // state: asset.exifInfo?.state, - // country: asset.exifInfo?.country, - // lensModel: asset.exifInfo?.lensModel, - // fNumber: asset.exifInfo?.fNumber, - // fps: asset.exifInfo?.fps, - // iso: asset.exifInfo?.iso, - }); }, } satisfies ExecuteOptions; } diff --git a/server/src/types.ts b/server/src/types.ts index dca3119c76..c738fd2559 100644 --- a/server/src/types.ts +++ b/server/src/types.ts @@ -458,6 +458,7 @@ export type JobItem = // Workflow | { name: JobName.WorkflowAssetTrigger; data: { workflowId: string; assetId: string } } + | { name: JobName.WorkflowAlbumAssetTrigger; data: { workflowId: string; assetId: string; albumId: string } } // Integrity | { name: JobName.IntegrityUntrackedFilesQueueAll; data?: IIntegrityJob } diff --git a/server/src/utils/workflow.ts b/server/src/utils/workflow.ts index c892239c80..49fd1b170d 100644 --- a/server/src/utils/workflow.ts +++ b/server/src/utils/workflow.ts @@ -6,6 +6,7 @@ export const triggerMap: Record = { [WorkflowTrigger.AssetCreate]: [WorkflowType.AssetV1], // [WorkflowTrigger.PersonRecognized]: [WorkflowType.AssetPersonV1], [WorkflowTrigger.AssetMetadataExtraction]: [WorkflowType.AssetV1], + [WorkflowTrigger.AlbumAssetAdded]: [WorkflowType.AlbumAssetV1], }; export const getWorkflowTriggers = () => @@ -14,6 +15,7 @@ export const getWorkflowTriggers = () => /** some types extend other types and have implied compatibility */ const inferredMap: Record = { [WorkflowType.AssetV1]: [], + [WorkflowType.AlbumAssetV1]: [WorkflowType.AssetV1], // [WorkflowType.AssetPersonV1]: [WorkflowType.AssetV1], }; @@ -29,8 +31,8 @@ const withImpliedItems = (type: WorkflowType): WorkflowType[] => { export const isMethodCompatible = (pluginMethod: { types: WorkflowType[] }, trigger: WorkflowTrigger) => { const validTypes = triggerMap[trigger]; - const pluginCompatibility = pluginMethod.types.map((type) => withImpliedItems(type)); - for (const requested of validTypes) { + const pluginCompatibility = validTypes.map((type) => withImpliedItems(type)); + for (const requested of pluginMethod.types) { for (const pluginCompatibilityGroup of pluginCompatibility) { if (pluginCompatibilityGroup.includes(requested)) { return true; diff --git a/web/src/lib/utils/workflow.ts b/web/src/lib/utils/workflow.ts index b8dc44cc67..5b4e477f22 100644 --- a/web/src/lib/utils/workflow.ts +++ b/web/src/lib/utils/workflow.ts @@ -13,6 +13,9 @@ export const getTriggerName = ($t: MessageFormatter, type: WorkflowTrigger) => { case WorkflowTrigger.AssetMetadataExtraction: { return $t('trigger_asset_metadata_extraction'); } + case WorkflowTrigger.AlbumAssetAdded: { + return $t('trigger_album_asset_added'); + } default: { return type; } @@ -30,6 +33,9 @@ export const getTriggerDescription = ($t: MessageFormatter, type: WorkflowTrigge case WorkflowTrigger.AssetMetadataExtraction: { return $t('trigger_asset_metadata_extraction_description'); } + case WorkflowTrigger.AlbumAssetAdded: { + return $t('trigger_album_asset_added_description'); + } } }; From 630b7f1dee85fb5fc34c5d8d0b122aeecf4635d6 Mon Sep 17 00:00:00 2001 From: Ben Beckford Date: Mon, 3 Aug 2026 04:01:15 -0700 Subject: [PATCH 2/4] feat(server): asset added to album workflow trigger --- packages/plugin-core/manifest.json | 22 ++++ packages/plugin-core/src/index.ts | 4 + packages/sdk/src/fetch-client.ts | 1 + server/src/enum.ts | 1 + server/src/repositories/event.repository.ts | 2 +- .../src/repositories/workflow.repository.ts | 26 ++++- server/src/services/album.service.ts | 13 +-- .../services/workflow-execution.service.ts | 106 ++++++++++-------- server/src/types.ts | 5 +- server/src/utils/workflow.ts | 2 +- 10 files changed, 118 insertions(+), 64 deletions(-) diff --git a/packages/plugin-core/manifest.json b/packages/plugin-core/manifest.json index a72d60c04d..6a2559362b 100644 --- a/packages/plugin-core/manifest.json +++ b/packages/plugin-core/manifest.json @@ -469,6 +469,28 @@ "required": ["albumIds"] } }, + { + "name": "assetAlbumFilter", + "title": "Filter by album(s)", + "description": "Filter by which album(s) triggered the workflow", + "types": ["AlbumAssetV1"], + "schema": { + "type": "object", + "properties": { + "albumIds": { + "type": "string", + "title": "Album IDs", + "array": true, + "description": "Allowed album IDs", + "uiHint": { + "type": "AlbumId" + } + } + }, + "required": ["albumIds"] + }, + "uiHints": ["Filter"] + }, { "name": "webhook", "title": "Trigger Webhook", diff --git a/packages/plugin-core/src/index.ts b/packages/plugin-core/src/index.ts index 7b91ed7111..9c4d81670a 100644 --- a/packages/plugin-core/src/index.ts +++ b/packages/plugin-core/src/index.ts @@ -66,6 +66,8 @@ const methods = wrapper({ return {}; }, + assetAlbumFilter: ({ config, data }) => ({ workflow: { continue: config.albumIds.includes(data.album.id) } }), + assetArchive: ({ config, data }) => { if (!config.inverse && data.asset.visibility !== AssetVisibility.Archive) { return { changes: { asset: { visibility: AssetVisibility.Archive } } }; @@ -209,6 +211,7 @@ const methods = wrapper({ const { assetAddToAlbums, + assetAlbumFilter, assetArchive, assetFavorite, assetFileFilter, @@ -227,6 +230,7 @@ const { export { assetAddToAlbums, + assetAlbumFilter, assetArchive, assetFavorite, assetFileFilter, diff --git a/packages/sdk/src/fetch-client.ts b/packages/sdk/src/fetch-client.ts index 81c2090f43..9f8764f9e5 100644 --- a/packages/sdk/src/fetch-client.ts +++ b/packages/sdk/src/fetch-client.ts @@ -7492,6 +7492,7 @@ export enum JobName { OcrQueueAll = "OcrQueueAll", Ocr = "Ocr", WorkflowAssetTrigger = "WorkflowAssetTrigger", + WorkflowAlbumAssetTrigger = "WorkflowAlbumAssetTrigger", IntegrityUntrackedFilesQueueAll = "IntegrityUntrackedFilesQueueAll", IntegrityUntrackedFiles = "IntegrityUntrackedFiles", IntegrityUntrackedRefresh = "IntegrityUntrackedRefresh", diff --git a/server/src/enum.ts b/server/src/enum.ts index 66e36b7121..803cad4410 100644 --- a/server/src/enum.ts +++ b/server/src/enum.ts @@ -344,6 +344,7 @@ export enum SystemMetadataKey { VersionCheckState = 'version-check-state', License = 'license', IntegrityChecksumCheckpoint = 'integrity-checksum-checkpoint', + AlbumAssetWorkflowCheckpoint = 'album-asset-workflow-checkpoint', } export enum UserMetadataKey { diff --git a/server/src/repositories/event.repository.ts b/server/src/repositories/event.repository.ts index 7a92446c77..cbd7d9ed62 100644 --- a/server/src/repositories/event.repository.ts +++ b/server/src/repositories/event.repository.ts @@ -40,7 +40,7 @@ type EventMap = { // album events AlbumUpdate: [{ id: string; userIds: string[]; recipientIds: string[] }]; AlbumInvite: [{ id: string; userId: string; senderName: string }]; - AlbumAssetsAdded: [{ albumId: string; userIds: string[]; recipientIds: string[]; assetIds: string[] }]; + AlbumAssetsAdded: []; // asset events AssetCreate: [{ asset: Pick; file?: UploadFile }]; diff --git a/server/src/repositories/workflow.repository.ts b/server/src/repositories/workflow.repository.ts index 9c4e435d00..62cb1f4aa4 100644 --- a/server/src/repositories/workflow.repository.ts +++ b/server/src/repositories/workflow.repository.ts @@ -1,5 +1,5 @@ import { Injectable } from '@nestjs/common'; -import { Insertable, Kysely, Updateable } from 'kysely'; +import { Insertable, Kysely, SelectQueryBuilder, Updateable } from 'kysely'; import { jsonArrayFrom, jsonObjectFrom } from 'kysely/helpers/postgres'; import { InjectKysely } from 'nestjs-kysely'; import { columns } from 'src/database'; @@ -139,8 +139,26 @@ export class WorkflowRepository { } getForAssetV1(assetId: string) { + return this.assetV1Query(this.db.selectFrom('asset').where('id', '=', assetId)).executeTakeFirstOrThrow(); + } + + getForAlbumAssetV1(after: string) { return this.db - .selectFrom('asset') + .selectFrom('album_asset') + .select(['albumId', 'assetId', 'album_asset.updateId']) + .orderBy('album_asset.updateId', 'asc') + .where('album_asset.updateId', '>', after) + .limit(2000) + .select((eb) => + jsonObjectFrom(this.assetV1Query(eb.selectFrom('asset').whereRef('asset.id', '=', 'album_asset.assetId'))).as( + 'asset', + ), + ) + .execute(); + } + + private assetV1Query(qb: SelectQueryBuilder) { + return qb .leftJoin('asset_exif', 'asset_exif.assetId', 'asset.id') .select((eb) => [ ...columns.workflowAssetV1, @@ -181,8 +199,6 @@ export class WorkflowRepository { ]) .whereRef('asset_exif.assetId', '=', 'asset.id'), ).as('exifInfo'), - ]) - .where('id', '=', assetId) - .executeTakeFirstOrThrow(); + ]); } } diff --git a/server/src/services/album.service.ts b/server/src/services/album.service.ts index 313c8affb1..55e3682e38 100644 --- a/server/src/services/album.service.ts +++ b/server/src/services/album.service.ts @@ -193,12 +193,7 @@ export class AlbumService extends BaseService { const userIds = album.albumUsers.map(({ user }) => user.id); const recipientIds = userIds.filter((userId) => userId !== auth.user.id); await this.eventRepository.emit('AlbumUpdate', { id, userIds, recipientIds }); - await this.eventRepository.emit('AlbumAssetsAdded', { - albumId: id, - userIds, - recipientIds, - assetIds: results.filter((a) => a.success).map((a) => a.id), - }); + await this.eventRepository.emit('AlbumAssetsAdded'); } return results; @@ -228,7 +223,6 @@ export class AlbumService extends BaseService { const albumAssetValues: { albumId: string; assetId: string }[] = []; const updateEvents: { id: string; userIds: string[]; recipientIds: string[] }[] = []; - const addedEvents: { albumId: string; userIds: string[]; recipientIds: string[]; assetIds: string[] }[] = []; for (const albumId of allowedAlbumIds) { const existingAssetIds = await this.albumRepository.getAssetIds(albumId, [...allowedAssetIds]); const notPresentAssetIds = [...allowedAssetIds.difference(existingAssetIds)]; @@ -254,16 +248,13 @@ export class AlbumService extends BaseService { const userIds = album.albumUsers.map(({ user }) => user.id); const recipientIds = userIds.filter((userId) => userId !== auth.user.id); updateEvents.push({ id: albumId, userIds, recipientIds }); - addedEvents.push({ albumId, userIds, recipientIds, assetIds: notPresentAssetIds }); } await this.albumRepository.addAssetIdsToAlbums(albumAssetValues); for (const event of updateEvents) { await this.eventRepository.emit('AlbumUpdate', event); } - for (const event of addedEvents) { - await this.eventRepository.emit('AlbumAssetsAdded', event); - } + await this.eventRepository.emit('AlbumAssetsAdded'); return results; } diff --git a/server/src/services/workflow-execution.service.ts b/server/src/services/workflow-execution.service.ts index ebe12523a2..1ba7dd32de 100644 --- a/server/src/services/workflow-execution.service.ts +++ b/server/src/services/workflow-execution.service.ts @@ -21,6 +21,7 @@ import { JobName, JobStatus, QueueName, + SystemMetadataKey, WorkflowType, } from 'src/enum'; import { ArgOf } from 'src/repositories/event.repository'; @@ -28,6 +29,7 @@ import { AlbumService } from 'src/services/album.service'; import { AssetService } from 'src/services/asset.service'; import { BaseService } from 'src/services/base.service'; import { JobItem, JobOf } from 'src/types'; +import { withImpliedItems } from 'src/utils/workflow'; const dummy = () => { throw new Error( @@ -42,8 +44,6 @@ type ExecuteOptions = { type AssetTrigger = { userId: string; assetId: string; trigger: WorkflowTrigger }; -type AlbumAssetTrigger = { userIds: string[]; albumId: string; assetIds: string[]; trigger: WorkflowTrigger }; - type HostContext = { allowedHosts: string[]; }; @@ -312,13 +312,8 @@ export class WorkflowExecutionService extends BaseService { } @OnEvent({ name: 'AlbumAssetsAdded' }) - onAlbumAssetsAdded({ albumId, userIds, recipientIds, assetIds }: ArgOf<'AlbumAssetsAdded'>) { - return this.onAlbumAssetTrigger({ - albumId, - userIds: [...userIds, ...recipientIds], - assetIds, - trigger: WorkflowTrigger.AlbumAssetAdded, - }); + onAlbumAssetsAdded() { + return this.onAlbumAssetTrigger(WorkflowTrigger.AlbumAssetAdded); } private async onAssetTrigger({ userId, assetId, trigger }: AssetTrigger) { @@ -331,30 +326,53 @@ export class WorkflowExecutionService extends BaseService { ); } - private async onAlbumAssetTrigger({ albumId, userIds, assetIds, trigger }: AlbumAssetTrigger) { - let jobs: JobItem[] = []; - for (const userId of userIds) { - const items = await this.workflowRepository.search({ userId, trigger }); - if (!items.length) { - continue; - } + private async onAlbumAssetTrigger(trigger: WorkflowTrigger) { + let checkpoint = await this.systemMetadataRepository.get(SystemMetadataKey.AlbumAssetWorkflowCheckpoint); + const now = await this.syncCheckpointRepository.getNow(); - let batch: JobItem[] = items.flatMap(({ id: workflowId }) => - assetIds.map((assetId) => ({ - name: JobName.WorkflowAlbumAssetTrigger, - data: { workflowId, assetId, albumId, trigger }, - })), - ); - jobs.push(...batch); + if (!checkpoint) { + checkpoint = { lastUuid: now.nowId }; + await this.systemMetadataRepository.set(SystemMetadataKey.AlbumAssetWorkflowCheckpoint, checkpoint); } - await this.jobRepository.queueAll(jobs); + const workflows = new Map(); + + while (checkpoint.lastUuid < now.nowId) { + const albumAssets = await this.workflowRepository.getForAlbumAssetV1(checkpoint.lastUuid); + if (albumAssets.length === 0) { + break; + } + + const jobs: JobItem[] = []; + for (const albumAsset of albumAssets) { + const userId = albumAsset.asset?.ownerId; + + if (!workflows.has(userId)) { + workflows.set(userId, await this.workflowRepository.search({ userId, trigger })); + } + + for (const workflow of workflows.get(userId)) { + jobs.push({ + name: JobName.WorkflowAlbumAssetTrigger, + data: { + workflowId: workflow.id, + albumAsset: { asset: albumAsset.asset as any, album: { id: albumAsset.albumId } }, + userId: workflow.ownerId, + }, + }); + } + } + + await this.jobRepository.queueAll(jobs); + checkpoint!.lastUuid = albumAssets[0].updateId; + await this.systemMetadataRepository.set(SystemMetadataKey.AlbumAssetWorkflowCheckpoint, checkpoint); + } } - private writeAssetV1(assetId: string, type: WorkflowType) { + private writeAssetV1(assetId: string) { const assetService = BaseService.create(AssetService, this); - return async (auth: AuthDto, changes: WorkflowChanges) => { + return async (auth: AuthDto, changes: WorkflowChanges) => { const asset = changes.asset; if (!asset) { return; @@ -399,38 +417,37 @@ export class WorkflowExecutionService extends BaseService { authUserId: asset.ownerId, }; }, - write: this.writeAssetV1(assetId, type), + write: this.writeAssetV1(assetId), } satisfies ExecuteOptions; } + default: { + return; + } } }); } @OnJob({ name: JobName.WorkflowAlbumAssetTrigger, queue: QueueName.Workflow }) - handleAlbumAssetTrigger({ workflowId, assetId, albumId }: JobOf) { + handleAlbumAssetTrigger({ workflowId, userId, albumAsset }: JobOf) { return this.execute(workflowId, (type) => { switch (type) { - case WorkflowType.AssetV1: { + case WorkflowType.AlbumAssetV1: { return { - read: async () => { - const asset = await this.workflowRepository.getForAssetV1(assetId); - const workflow = await this.workflowRepository.get(workflowId); - - return { - data: { asset, album: { id: albumId } } as any, - authUserId: workflow!.ownerId, - }; - }, + read: () => + Promise.resolve({ + data: albumAsset, + authUserId: userId, + }), write: async (auth, changes) => { - const asset = await this.workflowRepository.getForAssetV1(assetId); - const workflow = await this.workflowRepository.get(workflowId); - - if (asset.ownerId === workflow!.ownerId) { - await this.writeAssetV1(assetId, type)(auth, changes); + if (albumAsset.asset.ownerId === userId) { + await this.writeAssetV1(albumAsset.asset.id)(auth, changes); } }, } satisfies ExecuteOptions; } + default: { + return; + } } }); } @@ -447,7 +464,8 @@ export class WorkflowExecutionService extends BaseService { // TODO infer from steps let type: T | undefined; for (const targetType of Object.values(WorkflowType)) { - const isMissing = workflow.steps.some((step) => !step.types.includes(targetType)); + const implied = withImpliedItems(targetType); + const isMissing = workflow.steps.some((step) => step.types.every((type) => !implied.includes(type))); if (!isMissing) { type = targetType as unknown as T; break; diff --git a/server/src/types.ts b/server/src/types.ts index c738fd2559..521aba6a05 100644 --- a/server/src/types.ts +++ b/server/src/types.ts @@ -1,4 +1,4 @@ -import { WorkflowTrigger } from '@immich/plugin-sdk'; +import { AlbumAssetV1, WorkflowTrigger } from '@immich/plugin-sdk'; import { ShallowDehydrateObject } from 'kysely'; import { SystemConfig } from 'src/config'; import { VECTOR_EXTENSIONS } from 'src/constants'; @@ -458,7 +458,7 @@ export type JobItem = // Workflow | { name: JobName.WorkflowAssetTrigger; data: { workflowId: string; assetId: string } } - | { name: JobName.WorkflowAlbumAssetTrigger; data: { workflowId: string; assetId: string; albumId: string } } + | { name: JobName.WorkflowAlbumAssetTrigger; data: { workflowId: string; userId: string; albumAsset: AlbumAssetV1 } } // Integrity | { name: JobName.IntegrityUntrackedFilesQueueAll; data?: IIntegrityJob } @@ -572,6 +572,7 @@ export interface SystemMetadata extends Record = { // [WorkflowType.AssetPersonV1]: [WorkflowType.AssetV1], }; -const withImpliedItems = (type: WorkflowType): WorkflowType[] => { +export const withImpliedItems = (type: WorkflowType): WorkflowType[] => { const childTypes = inferredMap[type]; const results = [type]; for (const child of childTypes) { From 71b461c72dedfb9b0202a00b9eafc02e89c99851 Mon Sep 17 00:00:00 2001 From: Ben Beckford Date: Mon, 10 Aug 2026 00:41:06 -0700 Subject: [PATCH 3/4] feat(server): queue batched workflow runs --- open-api/immich-openapi-specs.json | 3 +- packages/sdk/src/fetch-client.ts | 3 +- server/src/enum.ts | 14 ++- server/src/queries/workflow.repository.sql | 1 + .../src/repositories/workflow.repository.ts | 15 ++- server/src/schema/index.ts | 2 + .../1786346220385-AddWorkflowQueue.ts | 16 +++ server/src/schema/tables/workflow-queue.ts | 14 +++ .../services/workflow-execution.service.ts | 115 +++++++++++------- server/src/types.ts | 8 +- 10 files changed, 140 insertions(+), 51 deletions(-) create mode 100644 server/src/schema/migrations/1786346220385-AddWorkflowQueue.ts create mode 100644 server/src/schema/tables/workflow-queue.ts diff --git a/open-api/immich-openapi-specs.json b/open-api/immich-openapi-specs.json index 0628fc30b8..c03966ba34 100644 --- a/open-api/immich-openapi-specs.json +++ b/open-api/immich-openapi-specs.json @@ -19338,8 +19338,9 @@ "VersionCheck", "OcrQueueAll", "Ocr", + "WorkflowScan", + "WorkflowRun", "WorkflowAssetTrigger", - "WorkflowAlbumAssetTrigger", "IntegrityUntrackedFilesQueueAll", "IntegrityUntrackedFiles", "IntegrityUntrackedRefresh", diff --git a/packages/sdk/src/fetch-client.ts b/packages/sdk/src/fetch-client.ts index 9f8764f9e5..df7a2e237a 100644 --- a/packages/sdk/src/fetch-client.ts +++ b/packages/sdk/src/fetch-client.ts @@ -7491,8 +7491,9 @@ export enum JobName { VersionCheck = "VersionCheck", OcrQueueAll = "OcrQueueAll", Ocr = "Ocr", + WorkflowScan = "WorkflowScan", + WorkflowRun = "WorkflowRun", WorkflowAssetTrigger = "WorkflowAssetTrigger", - WorkflowAlbumAssetTrigger = "WorkflowAlbumAssetTrigger", IntegrityUntrackedFilesQueueAll = "IntegrityUntrackedFilesQueueAll", IntegrityUntrackedFiles = "IntegrityUntrackedFiles", IntegrityUntrackedRefresh = "IntegrityUntrackedRefresh", diff --git a/server/src/enum.ts b/server/src/enum.ts index 803cad4410..480602e665 100644 --- a/server/src/enum.ts +++ b/server/src/enum.ts @@ -344,7 +344,7 @@ export enum SystemMetadataKey { VersionCheckState = 'version-check-state', License = 'license', IntegrityChecksumCheckpoint = 'integrity-checksum-checkpoint', - AlbumAssetWorkflowCheckpoint = 'album-asset-workflow-checkpoint', + WorkflowCheckpoint = 'workflow-checkpoint', } export enum UserMetadataKey { @@ -903,8 +903,9 @@ export enum JobName { Ocr = 'Ocr', // Workflow + WorkflowScan = 'WorkflowScan', + WorkflowRun = 'WorkflowRun', WorkflowAssetTrigger = 'WorkflowAssetTrigger', - WorkflowAlbumAssetTrigger = 'WorkflowAlbumAssetTrigger', // Integrity IntegrityUntrackedFilesQueueAll = 'IntegrityUntrackedFilesQueueAll', @@ -1233,6 +1234,15 @@ export enum WorkflowType { export const WorkflowTypeSchema = z.enum(WorkflowType).describe('Workflow type').meta({ id: 'WorkflowType' }); +export enum WorkflowScanType { + AlbumAsset = 'AlbumAsset', +} + +export const WorkflowScanTypeSchema = z + .enum(WorkflowScanType) + .describe('Workflow scan type') + .meta({ id: 'WorkflowScanType' }); + export enum CalendarHeatmapType { Upload = 'Upload', Taken = 'Taken', diff --git a/server/src/queries/workflow.repository.sql b/server/src/queries/workflow.repository.sql index 78ab21385f..661c6a3063 100644 --- a/server/src/queries/workflow.repository.sql +++ b/server/src/queries/workflow.repository.sql @@ -75,6 +75,7 @@ select "workflow"."id", "workflow"."name", "workflow"."trigger", + "workflow"."ownerId", ( select coalesce(json_agg(agg), '[]') diff --git a/server/src/repositories/workflow.repository.ts b/server/src/repositories/workflow.repository.ts index 62cb1f4aa4..152a2e689d 100644 --- a/server/src/repositories/workflow.repository.ts +++ b/server/src/repositories/workflow.repository.ts @@ -6,6 +6,7 @@ import { columns } from 'src/database'; import { DummyValue, GenerateSql } from 'src/decorators'; import { WorkflowSearchDto } from 'src/dtos/workflow.dto'; import { DB } from 'src/schema'; +import { WorkflowQueueTable } from 'src/schema/tables/workflow-queue'; import { WorkflowStepTable } from 'src/schema/tables/workflow-step.table'; import { WorkflowTable } from 'src/schema/tables/workflow.table'; @@ -66,7 +67,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.ownerId']) .select((eb) => [ jsonArrayFrom( eb @@ -157,6 +158,18 @@ export class WorkflowRepository { .execute(); } + addToQueue(dto: Insertable) { + return this.db.insertInto('workflow_queue').values(dto).returning(['id']).executeTakeFirstOrThrow(); + } + + getQueue(id: string) { + return this.db.selectFrom('workflow_queue').selectAll().where('id', '=', id).executeTakeFirstOrThrow(); + } + + removeFromQueue(workflowId: string) { + return this.db.deleteFrom('workflow_queue').where('id', '=', workflowId).execute(); + } + private assetV1Query(qb: SelectQueryBuilder) { return qb .leftJoin('asset_exif', 'asset_exif.assetId', 'asset.id') diff --git a/server/src/schema/index.ts b/server/src/schema/index.ts index 4e55ff4cf6..cfd1c24084 100644 --- a/server/src/schema/index.ts +++ b/server/src/schema/index.ts @@ -86,6 +86,7 @@ import { VideoStreamSessionTable, VideoStreamVariantTable, } from 'src/schema/tables/video-stream.table'; +import { WorkflowQueueTable } from 'src/schema/tables/workflow-queue'; import { WorkflowStepTable } from 'src/schema/tables/workflow-step.table'; import { WorkflowTable } from 'src/schema/tables/workflow.table'; @@ -277,4 +278,5 @@ export interface DB { workflow: WorkflowTable; workflow_step: WorkflowStepTable; + workflow_queue: WorkflowQueueTable; } diff --git a/server/src/schema/migrations/1786346220385-AddWorkflowQueue.ts b/server/src/schema/migrations/1786346220385-AddWorkflowQueue.ts new file mode 100644 index 0000000000..3bae1dcd83 --- /dev/null +++ b/server/src/schema/migrations/1786346220385-AddWorkflowQueue.ts @@ -0,0 +1,16 @@ +import { Kysely, sql } from 'kysely'; + +export async function up(db: Kysely): Promise { + await sql`CREATE TABLE "workflow_queue" ( + "id" uuid NOT NULL DEFAULT uuid_generate_v4(), + "workflowId" uuid NOT NULL, + "data" jsonb NOT NULL, + CONSTRAINT "workflow_queue_workflowId_fkey" FOREIGN KEY ("workflowId") REFERENCES "workflow" ("id") ON UPDATE CASCADE ON DELETE CASCADE, + CONSTRAINT "workflow_queue_pkey" PRIMARY KEY ("id") +);`.execute(db); + await sql`CREATE INDEX "workflow_queue_workflowId_idx" ON "workflow_queue" ("workflowId");`.execute(db); +} + +export async function down(db: Kysely): Promise { + await sql`DROP TABLE "workflow_queue";`.execute(db); +} diff --git a/server/src/schema/tables/workflow-queue.ts b/server/src/schema/tables/workflow-queue.ts new file mode 100644 index 0000000000..14e29061ab --- /dev/null +++ b/server/src/schema/tables/workflow-queue.ts @@ -0,0 +1,14 @@ +import { Column, ForeignKeyColumn, Generated, PrimaryGeneratedColumn, Table } from '@immich/sql-tools'; +import { WorkflowTable } from 'src/schema/tables/workflow.table'; + +@Table('workflow_queue') +export class WorkflowQueueTable { + @PrimaryGeneratedColumn('uuid') + id!: Generated; + + @ForeignKeyColumn(() => WorkflowTable, { onDelete: 'CASCADE', onUpdate: 'CASCADE' }) + workflowId!: string; + + @Column({ type: 'jsonb' }) + data!: unknown[]; +} diff --git a/server/src/services/workflow-execution.service.ts b/server/src/services/workflow-execution.service.ts index 1ba7dd32de..e3537157b7 100644 --- a/server/src/services/workflow-execution.service.ts +++ b/server/src/services/workflow-execution.service.ts @@ -1,5 +1,6 @@ import { CurrentPlugin } from '@extism/extism'; import { + AlbumAssetV1, WorkflowChanges, WorkflowEventData, WorkflowEventPayload, @@ -22,13 +23,14 @@ import { JobStatus, QueueName, SystemMetadataKey, + WorkflowScanType, WorkflowType, } from 'src/enum'; import { ArgOf } from 'src/repositories/event.repository'; import { AlbumService } from 'src/services/album.service'; import { AssetService } from 'src/services/asset.service'; import { BaseService } from 'src/services/base.service'; -import { JobItem, JobOf } from 'src/types'; +import { JobOf } from 'src/types'; import { withImpliedItems } from 'src/utils/workflow'; const dummy = () => { @@ -50,6 +52,7 @@ type HostContext = { export class WorkflowExecutionService extends BaseService { private jwtSecret!: string; + private scanning = false; @OnEvent({ name: 'AppBootstrap', priority: BootstrapEventPriority.PluginSync, workers: [ImmichWorker.Microservices] }) async onPluginSync() { @@ -313,7 +316,7 @@ export class WorkflowExecutionService extends BaseService { @OnEvent({ name: 'AlbumAssetsAdded' }) onAlbumAssetsAdded() { - return this.onAlbumAssetTrigger(WorkflowTrigger.AlbumAssetAdded); + return this.jobRepository.queue({ name: JobName.WorkflowScan, data: { type: WorkflowScanType.AlbumAsset } }); } private async onAssetTrigger({ userId, assetId, trigger }: AssetTrigger) { @@ -326,47 +329,65 @@ export class WorkflowExecutionService extends BaseService { ); } - private async onAlbumAssetTrigger(trigger: WorkflowTrigger) { - let checkpoint = await this.systemMetadataRepository.get(SystemMetadataKey.AlbumAssetWorkflowCheckpoint); + @OnJob({ name: JobName.WorkflowScan, queue: QueueName.Workflow }) + private async scan({ type }: JobOf) { + if (this.scanning) { + return JobStatus.Skipped; + } + + this.scanning = true; + + if (type !== WorkflowScanType.AlbumAsset) { + return; + } + + let checkpoint = await this.systemMetadataRepository.get(SystemMetadataKey.WorkflowCheckpoint); const now = await this.syncCheckpointRepository.getNow(); if (!checkpoint) { - checkpoint = { lastUuid: now.nowId }; - await this.systemMetadataRepository.set(SystemMetadataKey.AlbumAssetWorkflowCheckpoint, checkpoint); + checkpoint = { albumAssetUuid: now.nowId }; + await this.systemMetadataRepository.set(SystemMetadataKey.WorkflowCheckpoint, checkpoint); } const workflows = new Map(); - while (checkpoint.lastUuid < now.nowId) { - const albumAssets = await this.workflowRepository.getForAlbumAssetV1(checkpoint.lastUuid); + while (checkpoint.albumAssetUuid < now.nowId) { + const albumAssets = await this.workflowRepository.getForAlbumAssetV1(checkpoint.albumAssetUuid); if (albumAssets.length === 0) { break; } - const jobs: JobItem[] = []; + const jobs = new Map(); for (const albumAsset of albumAssets) { const userId = albumAsset.asset?.ownerId; if (!workflows.has(userId)) { - workflows.set(userId, await this.workflowRepository.search({ userId, trigger })); + workflows.set( + userId, + await this.workflowRepository.search({ userId, trigger: WorkflowTrigger.AlbumAssetAdded }), + ); } for (const workflow of workflows.get(userId)) { - jobs.push({ - name: JobName.WorkflowAlbumAssetTrigger, - data: { - workflowId: workflow.id, - albumAsset: { asset: albumAsset.asset as any, album: { id: albumAsset.albumId } }, - userId: workflow.ownerId, - }, - }); + if (!jobs.has(workflow.id)) { + jobs.set(workflow.id, []); + } + + jobs.get(workflow.id)!.push({ asset: albumAsset.asset as any, album: { id: albumAsset.albumId } }); } } - await this.jobRepository.queueAll(jobs); - checkpoint!.lastUuid = albumAssets[0].updateId; - await this.systemMetadataRepository.set(SystemMetadataKey.AlbumAssetWorkflowCheckpoint, checkpoint); + for (const workflowId of jobs.keys()) { + const { id } = await this.workflowRepository.addToQueue({ workflowId, data: jobs.get(workflowId)! }); + await this.jobRepository.queue({ name: JobName.WorkflowRun, data: { queueId: id } }); + } + + checkpoint!.albumAssetUuid = albumAssets[0].updateId; + await this.systemMetadataRepository.set(SystemMetadataKey.WorkflowCheckpoint, checkpoint); } + + this.scanning = false; + return JobStatus.Success; } private writeAssetV1(assetId: string) { @@ -427,29 +448,37 @@ export class WorkflowExecutionService extends BaseService { }); } - @OnJob({ name: JobName.WorkflowAlbumAssetTrigger, queue: QueueName.Workflow }) - handleAlbumAssetTrigger({ workflowId, userId, albumAsset }: JobOf) { - return this.execute(workflowId, (type) => { - switch (type) { - case WorkflowType.AlbumAssetV1: { - return { - read: () => - Promise.resolve({ - data: albumAsset, - authUserId: userId, - }), - write: async (auth, changes) => { - if (albumAsset.asset.ownerId === userId) { - await this.writeAssetV1(albumAsset.asset.id)(auth, changes); - } - }, - } satisfies ExecuteOptions; + @OnJob({ name: JobName.WorkflowRun, queue: QueueName.Workflow }) + async runQueue({ queueId }: JobOf) { + const queue = await this.workflowRepository.getQueue(queueId); + + for (const item of queue.data) { + await this.execute(queue.workflowId, (type) => { + switch (type) { + case WorkflowType.AssetV1: + case WorkflowType.AlbumAssetV1: { + return { + read: async () => { + const workflow = await this.workflowRepository.getForWorkflowRun(queue.workflowId); + return { + data: item as any, + authUserId: workflow!.ownerId, + }; + }, + write: async (auth, changes) => { + const workflow = await this.workflowRepository.getForWorkflowRun(queue.workflowId); + if ((item as AlbumAssetV1).asset.ownerId === workflow?.ownerId) { + await this.writeAssetV1((item as AlbumAssetV1).asset.id)(auth, changes); + } + }, + } satisfies ExecuteOptions; + } + default: { + return; + } } - default: { - return; - } - } - }); + }); + } } private async execute( diff --git a/server/src/types.ts b/server/src/types.ts index 7b51639404..b62f9b90b7 100644 --- a/server/src/types.ts +++ b/server/src/types.ts @@ -1,4 +1,4 @@ -import { AlbumAssetV1, WorkflowTrigger } from '@immich/plugin-sdk'; +import { WorkflowTrigger } from '@immich/plugin-sdk'; import { ShallowDehydrateObject } from 'kysely'; import { SystemConfig } from 'src/config'; import { VECTOR_EXTENSIONS } from 'src/constants'; @@ -30,6 +30,7 @@ import { SystemMetadataKey, TranscodeTarget, UserMetadataKey, + WorkflowScanType, WorkflowType, } from 'src/enum'; import { Mocked } from 'vitest'; @@ -459,7 +460,8 @@ export type JobItem = // Workflow | { name: JobName.WorkflowAssetTrigger; data: { workflowId: string; assetId: string } } - | { name: JobName.WorkflowAlbumAssetTrigger; data: { workflowId: string; userId: string; albumAsset: AlbumAssetV1 } } + | { name: JobName.WorkflowRun; data: { queueId: string } } + | { name: JobName.WorkflowScan; data: { type: WorkflowScanType } } // Integrity | { name: JobName.IntegrityUntrackedFilesQueueAll; data?: IIntegrityJob } @@ -573,7 +575,7 @@ export interface SystemMetadata extends Record Date: Wed, 12 Aug 2026 14:32:21 -0700 Subject: [PATCH 4/4] chore(server): batch add workflow queues to db --- server/src/repositories/workflow.repository.ts | 4 ++-- server/src/services/workflow-execution.service.ts | 11 +++++++---- 2 files changed, 9 insertions(+), 6 deletions(-) diff --git a/server/src/repositories/workflow.repository.ts b/server/src/repositories/workflow.repository.ts index 152a2e689d..3dbd67c6a1 100644 --- a/server/src/repositories/workflow.repository.ts +++ b/server/src/repositories/workflow.repository.ts @@ -158,8 +158,8 @@ export class WorkflowRepository { .execute(); } - addToQueue(dto: Insertable) { - return this.db.insertInto('workflow_queue').values(dto).returning(['id']).executeTakeFirstOrThrow(); + addToQueue(dto: Insertable[]) { + return this.db.insertInto('workflow_queue').values(dto).returning(['id']).execute(); } getQueue(id: string) { diff --git a/server/src/services/workflow-execution.service.ts b/server/src/services/workflow-execution.service.ts index e3537157b7..c7b027368f 100644 --- a/server/src/services/workflow-execution.service.ts +++ b/server/src/services/workflow-execution.service.ts @@ -377,10 +377,13 @@ export class WorkflowExecutionService extends BaseService { } } - for (const workflowId of jobs.keys()) { - const { id } = await this.workflowRepository.addToQueue({ workflowId, data: jobs.get(workflowId)! }); - await this.jobRepository.queue({ name: JobName.WorkflowRun, data: { queueId: id } }); - } + const queues = await this.workflowRepository.addToQueue( + jobs + .entries() + .map(([workflowId, data]) => ({ workflowId, data })) + .toArray(), + ); + await this.jobRepository.queueAll(queues.map(({ id }) => ({ name: JobName.WorkflowRun, data: { queueId: id } }))); checkpoint!.albumAssetUuid = albumAssets[0].updateId; await this.systemMetadataRepository.set(SystemMetadataKey.WorkflowCheckpoint, checkpoint);