feat: asset face v3 (#31591)

This commit is contained in:
Jason Rasmussen
2026-09-16 16:27:01 -04:00
committed by GitHub
parent 810454a343
commit 54e7ffa4e6
12 changed files with 336 additions and 39 deletions
@@ -325,6 +325,8 @@ class SyncStreamService {
return _syncStreamRepository.updateAssetFacesV1(data.cast());
case SyncEntityType.assetFaceV2:
return _syncStreamRepository.updateAssetFacesV2(data.cast());
case SyncEntityType.assetFaceV3:
throw UnimplementedError('SyncEntityType.assetFaceV3 is not implemented yet');
case SyncEntityType.assetFaceDeleteV1:
return _syncStreamRepository.deleteAssetFacesV1(data.cast());
case SyncEntityType.assetOcrV1:
@@ -194,7 +194,7 @@ const _kResponseMap = <SyncEntityType, Function(Object)>{
SyncEntityType.personV1: SyncPersonV1.fromJson,
SyncEntityType.personDeleteV1: SyncPersonDeleteV1.fromJson,
SyncEntityType.assetFaceV1: SyncAssetFaceV1.fromJson,
SyncEntityType.assetFaceV2: SyncAssetFaceV2.fromJson,
SyncEntityType.assetFaceV2: SyncAssetFaceV3.fromJson,
SyncEntityType.assetFaceDeleteV1: SyncAssetFaceDeleteV1.fromJson,
SyncEntityType.assetOcrV1: SyncAssetOcrV1.fromJson,
SyncEntityType.assetOcrDeleteV1: SyncAssetOcrDeleteV1.fromJson,
@@ -840,7 +840,7 @@ class SyncStreamRepository extends DatabaseAccessor<Drift> with $SyncStreamRepos
}
}
Future<void> updateAssetFacesV2(Iterable<SyncAssetFaceV2> data) async {
Future<void> updateAssetFacesV2(Iterable<SyncAssetFaceV3> data) async {
try {
await _db.batch((batch) {
for (final assetFace in data) {
+3 -1
View File
@@ -30089,7 +30089,7 @@
],
"type": "object"
},
"SyncAssetFaceV2": {
"SyncAssetFaceV3": {
"properties": {
"assetId": {
"description": "Asset ID",
@@ -30863,6 +30863,7 @@
"PersonDeleteV1",
"AssetFaceV1",
"AssetFaceV2",
"AssetFaceV3",
"AssetFaceDeleteV1",
"UserMetadataV1",
"UserMetadataDeleteV1",
@@ -31188,6 +31189,7 @@
"PeopleV1",
"AssetFacesV1",
"AssetFacesV2",
"AssetFacesV3",
"UserMetadataV1"
],
"type": "string"
+3 -1
View File
@@ -3407,7 +3407,7 @@ export type SyncAssetFaceV1 = {
/** Source type */
sourceType: string;
};
export type SyncAssetFaceV2 = {
export type SyncAssetFaceV3 = {
/** Asset ID */
assetId: string;
/** Bounding box X1 */
@@ -8386,6 +8386,7 @@ export enum SyncEntityType {
PersonDeleteV1 = "PersonDeleteV1",
AssetFaceV1 = "AssetFaceV1",
AssetFaceV2 = "AssetFaceV2",
AssetFaceV3 = "AssetFaceV3",
AssetFaceDeleteV1 = "AssetFaceDeleteV1",
UserMetadataV1 = "UserMetadataV1",
UserMetadataDeleteV1 = "UserMetadataDeleteV1",
@@ -8421,6 +8422,7 @@ export enum SyncRequestType {
PeopleV1 = "PeopleV1",
AssetFacesV1 = "AssetFacesV1",
AssetFacesV2 = "AssetFacesV2",
AssetFacesV3 = "AssetFacesV3",
UserMetadataV1 = "UserMetadataV1"
}
export enum AssetOrderBy {
+14
View File
@@ -471,6 +471,20 @@ export const columns = {
syncStack: ['stack.id', 'stack.createdAt', 'stack.updatedAt', 'stack.primaryAssetId', 'stack.ownerId'],
syncUser: ['id', 'name', 'email', 'avatarColor', 'deletedAt', 'updateId', 'profileImagePath', 'profileChangedAt'],
stack: ['stack.id', 'stack.primaryAssetId', 'ownerId'],
syncAssetFace: [
'asset_face.id',
'asset_face.assetId',
'asset_face.personGroupId as personId',
'asset_face.imageWidth',
'asset_face.imageHeight',
'asset_face.boundingBoxX1',
'asset_face.boundingBoxY1',
'asset_face.boundingBoxX2',
'asset_face.boundingBoxY2',
'asset_face.sourceType',
'asset_face.isVisible',
'asset_face.deletedAt',
],
syncAssetExif: [
'asset_exif.assetId',
'asset_exif.description',
+6 -4
View File
@@ -371,10 +371,11 @@ const SyncAssetFaceV1Schema = z
})
.meta({ id: 'SyncAssetFaceV1' });
const SyncAssetFaceV2Schema = SyncAssetFaceV1Schema.extend({
// same shape as V2, but scoped to the whole cluster group instead of the user's own assets
const SyncAssetFaceV3Schema = SyncAssetFaceV1Schema.extend({
deletedAt: isoDatetimeToDate.nullable().describe('Face deleted at'),
isVisible: z.boolean().describe('Is the face visible in the asset'),
}).meta({ id: 'SyncAssetFaceV2' });
}).meta({ id: 'SyncAssetFaceV3' });
const SyncAssetFaceDeleteV1Schema = z
.object({ assetFaceId: z.uuidv4().describe('Asset face ID') })
@@ -453,7 +454,7 @@ class SyncPersonDeleteV1 extends createZodDto(SyncPersonDeleteV1Schema) {}
@ExtraModel()
class SyncAssetFaceV1 extends createZodDto(SyncAssetFaceV1Schema) {}
@ExtraModel()
class SyncAssetFaceV2 extends createZodDto(SyncAssetFaceV2Schema) {}
class SyncAssetFaceV3 extends createZodDto(SyncAssetFaceV3Schema) {}
@ExtraModel()
class SyncAssetFaceDeleteV1 extends createZodDto(SyncAssetFaceDeleteV1Schema) {}
@ExtraModel()
@@ -515,7 +516,8 @@ export type SyncItem = {
[SyncEntityType.PersonV1]: SyncPersonV1;
[SyncEntityType.PersonDeleteV1]: SyncPersonDeleteV1;
[SyncEntityType.AssetFaceV1]: SyncAssetFaceV1;
[SyncEntityType.AssetFaceV2]: SyncAssetFaceV2;
[SyncEntityType.AssetFaceV2]: SyncAssetFaceV3;
[SyncEntityType.AssetFaceV3]: SyncAssetFaceV3;
[SyncEntityType.AssetFaceDeleteV1]: SyncAssetFaceDeleteV1;
[SyncEntityType.UserMetadataV1]: SyncUserMetadataV1;
[SyncEntityType.UserMetadataDeleteV1]: SyncUserMetadataDeleteV1;
+5
View File
@@ -1038,7 +1038,9 @@ export enum SyncRequestType {
PeopleV1 = 'PeopleV1',
/** @deprecated */
AssetFacesV1 = 'AssetFacesV1',
/** @deprecated */
AssetFacesV2 = 'AssetFacesV2',
AssetFacesV3 = 'AssetFacesV3',
UserMetadataV1 = 'UserMetadataV1',
}
@@ -1119,8 +1121,11 @@ export enum SyncEntityType {
PersonV1 = 'PersonV1',
PersonDeleteV1 = 'PersonDeleteV1',
/** @deprecated */
AssetFaceV1 = 'AssetFaceV1',
/** @deprecated */
AssetFaceV2 = 'AssetFaceV2',
AssetFaceV3 = 'AssetFaceV3',
AssetFaceDeleteV1 = 'AssetFaceDeleteV1',
UserMetadataV1 = 'UserMetadataV1',
+67 -12
View File
@@ -518,7 +518,7 @@ where
order by
"asset_edit"."updateId" asc
-- SyncRepository.assetFace.getDeletes
-- SyncRepository.assetFace.getDeletesV2
select
"asset_face_audit"."id",
"assetFaceId"
@@ -532,19 +532,41 @@ where
order by
"asset_face_audit"."id" asc
-- SyncRepository.assetFace.getUpserts
-- SyncRepository.assetFace.getDeletesV3
select
"asset_face_audit"."id",
"assetFaceId"
from
"asset_face_audit" as "asset_face_audit"
inner join "asset" on "asset"."id" = "asset_face_audit"."assetId"
inner join "user" as "owner" on "owner"."id" = "asset"."ownerId"
where
"asset_face_audit"."id" < $1
and "asset_face_audit"."id" > $2
and "owner"."clusterGroupId" = (
select
"user"."clusterGroupId"
from
"user"
where
"user"."id" = $3
)
order by
"asset_face_audit"."id" asc
-- SyncRepository.assetFace.getUpsertsV2
select
"asset_face"."id",
"assetId",
"personGroupId" as "personId",
"imageWidth",
"imageHeight",
"boundingBoxX1",
"boundingBoxY1",
"boundingBoxX2",
"boundingBoxY2",
"sourceType",
"isVisible",
"asset_face"."assetId",
"asset_face"."personGroupId" as "personId",
"asset_face"."imageWidth",
"asset_face"."imageHeight",
"asset_face"."boundingBoxX1",
"asset_face"."boundingBoxY1",
"asset_face"."boundingBoxX2",
"asset_face"."boundingBoxY2",
"asset_face"."sourceType",
"asset_face"."isVisible",
"asset_face"."deletedAt",
"asset_face"."updateId"
from
@@ -557,6 +579,39 @@ where
order by
"asset_face"."updateId" asc
-- SyncRepository.assetFace.getUpsertsV3
select
"asset_face"."id",
"asset_face"."assetId",
"asset_face"."personGroupId" as "personId",
"asset_face"."imageWidth",
"asset_face"."imageHeight",
"asset_face"."boundingBoxX1",
"asset_face"."boundingBoxY1",
"asset_face"."boundingBoxX2",
"asset_face"."boundingBoxY2",
"asset_face"."sourceType",
"asset_face"."isVisible",
"asset_face"."deletedAt",
"asset_face"."updateId"
from
"asset_face" as "asset_face"
inner join "asset" on "asset"."id" = "asset_face"."assetId"
inner join "user" as "owner" on "owner"."id" = "asset"."ownerId"
where
"asset_face"."updateId" < $1
and "asset_face"."updateId" > $2
and "owner"."clusterGroupId" = (
select
"user"."clusterGroupId"
from
"user"
where
"user"."id" = $3
)
order by
"asset_face"."updateId" asc
-- SyncRepository.assetMetadata.getDeletes
select
"asset_metadata_audit"."id",
+31 -17
View File
@@ -461,8 +461,9 @@ class PersonGroupSync extends BaseSync {
}
class AssetFaceSync extends BaseSync {
// TODO(v5) drop when AssetFacesV2 is removed
@GenerateSql({ params: [dummyQueryOptions], stream: true })
getDeletes(options: SyncQueryOptions) {
getDeletesV2(options: SyncQueryOptions) {
return this.auditQuery('asset_face_audit', options)
.select(['asset_face_audit.id', 'assetFaceId'])
.leftJoin('asset', 'asset.id', 'asset_face_audit.assetId')
@@ -470,32 +471,45 @@ class AssetFaceSync extends BaseSync {
.stream();
}
@GenerateSql({ params: [dummyQueryOptions], stream: true })
getDeletesV3(options: SyncQueryOptions) {
return this.auditQuery('asset_face_audit', options)
.select(['asset_face_audit.id', 'assetFaceId'])
.innerJoin('asset', 'asset.id', 'asset_face_audit.assetId')
.innerJoin('user as owner', 'owner.id', 'asset.ownerId')
.where('owner.clusterGroupId', '=', ({ selectFrom }) =>
selectFrom('user').select('user.clusterGroupId').where('user.id', '=', options.userId),
)
.stream();
}
cleanupAuditTable(daysAgo: number) {
return this.auditCleanup('asset_face_audit', daysAgo);
}
// TODO(v5) drop when AssetFacesV2 is removed
@GenerateSql({ params: [dummyQueryOptions], stream: true })
getUpserts(options: SyncQueryOptions) {
getUpsertsV2(options: SyncQueryOptions) {
return this.upsertQuery('asset_face', options)
.select([
'asset_face.id',
'assetId',
'personGroupId as personId',
'imageWidth',
'imageHeight',
'boundingBoxX1',
'boundingBoxY1',
'boundingBoxX2',
'boundingBoxY2',
'sourceType',
'isVisible',
'asset_face.deletedAt',
'asset_face.updateId',
])
.select(columns.syncAssetFace)
.select('asset_face.updateId')
.leftJoin('asset', 'asset.id', 'asset_face.assetId')
.where('asset.ownerId', '=', options.userId)
.stream();
}
@GenerateSql({ params: [dummyQueryOptions], stream: true })
getUpsertsV3(options: SyncQueryOptions) {
return this.upsertQuery('asset_face', options)
.select(columns.syncAssetFace)
.select('asset_face.updateId')
.innerJoin('asset', 'asset.id', 'asset_face.assetId')
.innerJoin('user as owner', 'owner.id', 'asset.ownerId')
.where('owner.clusterGroupId', '=', ({ selectFrom }) =>
selectFrom('user').select('user.clusterGroupId').where('user.id', '=', options.userId),
)
.stream();
}
}
class AssetExifSync extends BaseSync {
+19 -2
View File
@@ -87,6 +87,7 @@ export const SYNC_TYPES_ORDER = [
SyncRequestType.PeopleV1,
SyncRequestType.AssetFacesV1,
SyncRequestType.AssetFacesV2,
SyncRequestType.AssetFacesV3,
SyncRequestType.UserMetadataV1,
SyncRequestType.AssetMetadataV1,
SyncRequestType.AssetEditsV1,
@@ -214,6 +215,7 @@ export class SyncService extends BaseService {
[SyncRequestType.PartnerStacksV1]: () => this.syncPartnerStackV1(options, response, checkpointMap, session.id),
[SyncRequestType.PeopleV1]: () => this.syncPeopleV1(options, response, checkpointMap),
[SyncRequestType.AssetFacesV2]: () => this.syncAssetFacesV2(options, response, checkpointMap),
[SyncRequestType.AssetFacesV3]: () => this.syncAssetFacesV3(options, response, checkpointMap),
[SyncRequestType.UserMetadataV1]: () => this.syncUserMetadataV1(options, response, checkpointMap),
[SyncRequestType.AssetOcrV1]: () => this.syncAssetOcrV1(options, response, checkpointMap, auth),
} as const;
@@ -896,15 +898,30 @@ export class SyncService extends BaseService {
);
}
// TODO(v5) drop when AssetFacesV2 is removed
private async syncAssetFacesV2(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) {
const deleteType = SyncEntityType.AssetFaceDeleteV1;
const deletes = this.syncRepository.assetFace.getDeletes({ ...options, ack: checkpointMap[deleteType] });
const deletes = this.syncRepository.assetFace.getDeletesV2({ ...options, ack: checkpointMap[deleteType] });
for await (const { id, ...data } of deletes) {
await send(response, { type: deleteType, ids: [id], data });
}
const upsertType = SyncEntityType.AssetFaceV2;
const upserts = this.syncRepository.assetFace.getUpserts({ ...options, ack: checkpointMap[upsertType] });
const upserts = this.syncRepository.assetFace.getUpsertsV2({ ...options, ack: checkpointMap[upsertType] });
for await (const { updateId, ...data } of upserts) {
await send(response, { type: upsertType, ids: [updateId], data });
}
}
private async syncAssetFacesV3(options: SyncQueryOptions, response: Writable, checkpointMap: CheckpointMap) {
const deleteType = SyncEntityType.AssetFaceDeleteV1;
const deletes = this.syncRepository.assetFace.getDeletesV3({ ...options, ack: checkpointMap[deleteType] });
for await (const { id, ...data } of deletes) {
await send(response, { type: deleteType, ids: [id], data });
}
const upsertType = SyncEntityType.AssetFaceV3;
const upserts = this.syncRepository.assetFace.getUpsertsV3({ ...options, ack: checkpointMap[upsertType] });
for await (const { updateId, ...data } of upserts) {
await send(response, { type: upsertType, ids: [updateId], data });
}
@@ -228,3 +228,187 @@ describe(SyncEntityType.AssetFaceV2, () => {
await ctx.assertSyncIsComplete(auth, [SyncRequestType.AssetFacesV2]);
});
});
describe(SyncEntityType.AssetFaceV3, () => {
it('should detect and sync the first asset face', async () => {
const { auth, ctx } = await setup();
const { asset } = await ctx.newAsset({ ownerId: auth.user.id });
const { person } = await ctx.newPerson({ ownerId: auth.user.id });
const { assetFace } = await ctx.newAssetFace({ assetId: asset.id, personGroupId: person.personGroupId });
const response = await ctx.syncStream(auth, [SyncRequestType.AssetFacesV3]);
expect(response).toEqual([
{
ack: expect.any(String),
data: {
id: assetFace.id,
assetId: asset.id,
personId: person.personGroupId,
imageWidth: assetFace.imageWidth,
imageHeight: assetFace.imageHeight,
boundingBoxX1: assetFace.boundingBoxX1,
boundingBoxY1: assetFace.boundingBoxY1,
boundingBoxX2: assetFace.boundingBoxX2,
boundingBoxY2: assetFace.boundingBoxY2,
sourceType: assetFace.sourceType,
deletedAt: null,
isVisible: true,
},
type: SyncEntityType.AssetFaceV3,
},
expect.objectContaining({ type: SyncEntityType.SyncCompleteV1 }),
]);
await ctx.syncAckAll(auth, response);
await ctx.assertSyncIsComplete(auth, [SyncRequestType.AssetFacesV3]);
});
it('should sync an asset face without a person', async () => {
const { auth, ctx } = await setup();
const { asset } = await ctx.newAsset({ ownerId: auth.user.id });
const { assetFace } = await ctx.newAssetFace({ assetId: asset.id });
const response = await ctx.syncStream(auth, [SyncRequestType.AssetFacesV3]);
expect(response).toEqual([
{
ack: expect.any(String),
data: expect.objectContaining({ id: assetFace.id, assetId: asset.id, personId: null }),
type: SyncEntityType.AssetFaceV3,
},
expect.objectContaining({ type: SyncEntityType.SyncCompleteV1 }),
]);
});
it('should sync an asset face belonging to another user in the same cluster group', async () => {
const { auth, user, ctx } = await setup();
const { user: user2 } = await ctx.newUser({ clusterGroupId: user.clusterGroupId });
const { asset } = await ctx.newAsset({ ownerId: user2.id });
const { person } = await ctx.newPerson({ ownerId: user2.id });
const { assetFace } = await ctx.newAssetFace({ assetId: asset.id, personGroupId: person.personGroupId });
const response = await ctx.syncStream(auth, [SyncRequestType.AssetFacesV3]);
expect(response).toEqual([
{
ack: expect.any(String),
data: expect.objectContaining({ id: assetFace.id, assetId: asset.id, personId: person.personGroupId }),
type: SyncEntityType.AssetFaceV3,
},
expect.objectContaining({ type: SyncEntityType.SyncCompleteV1 }),
]);
await ctx.syncAckAll(auth, response);
await ctx.assertSyncIsComplete(auth, [SyncRequestType.AssetFacesV3]);
});
it('should detect and sync an updated asset face for another user in the same cluster group', async () => {
const { auth, user, ctx } = await setup();
const personRepo = ctx.get(PersonRepository);
const { user: user2 } = await ctx.newUser({ clusterGroupId: user.clusterGroupId });
const { asset } = await ctx.newAsset({ ownerId: user2.id });
const { person } = await ctx.newPerson({ ownerId: user2.id });
const { assetFace } = await ctx.newAssetFace({ assetId: asset.id, personGroupId: person.personGroupId });
const response = await ctx.syncStream(auth, [SyncRequestType.AssetFacesV3]);
expect(response).toEqual([
expect.objectContaining({ type: SyncEntityType.AssetFaceV3 }),
expect.objectContaining({ type: SyncEntityType.SyncCompleteV1 }),
]);
await ctx.syncAckAll(auth, response);
await ctx.assertSyncIsComplete(auth, [SyncRequestType.AssetFacesV3]);
await personRepo.softDeleteAssetFaces(assetFace.id);
expect(await ctx.syncStream(auth, [SyncRequestType.AssetFacesV3])).toEqual([
{
ack: expect.any(String),
data: expect.objectContaining({ id: assetFace.id, deletedAt: expect.any(String) }),
type: SyncEntityType.AssetFaceV3,
},
expect.objectContaining({ type: SyncEntityType.SyncCompleteV1 }),
]);
});
it('should detect and sync a deleted asset face', async () => {
const { auth, ctx } = await setup();
const personRepo = ctx.get(PersonRepository);
const { asset } = await ctx.newAsset({ ownerId: auth.user.id });
const { person } = await ctx.newPerson({ ownerId: auth.user.id });
const { assetFace } = await ctx.newAssetFace({ assetId: asset.id, personGroupId: person.personGroupId });
await personRepo.deleteAssetFace(assetFace.id);
const response = await ctx.syncStream(auth, [SyncRequestType.AssetFacesV3]);
expect(response).toEqual([
{
ack: expect.any(String),
data: {
assetFaceId: assetFace.id,
},
type: SyncEntityType.AssetFaceDeleteV1,
},
expect.objectContaining({ type: SyncEntityType.SyncCompleteV1 }),
]);
await ctx.syncAckAll(auth, response);
await ctx.assertSyncIsComplete(auth, [SyncRequestType.AssetFacesV3]);
});
it('should detect and sync a deleted asset face belonging to another user in the same cluster group', async () => {
const { auth, user, ctx } = await setup();
const personRepo = ctx.get(PersonRepository);
const { user: user2 } = await ctx.newUser({ clusterGroupId: user.clusterGroupId });
const { asset } = await ctx.newAsset({ ownerId: user2.id });
const { person } = await ctx.newPerson({ ownerId: user2.id });
const { assetFace } = await ctx.newAssetFace({ assetId: asset.id, personGroupId: person.personGroupId });
await personRepo.deleteAssetFace(assetFace.id);
const response = await ctx.syncStream(auth, [SyncRequestType.AssetFacesV3]);
expect(response).toEqual([
{
ack: expect.any(String),
data: {
assetFaceId: assetFace.id,
},
type: SyncEntityType.AssetFaceDeleteV1,
},
expect.objectContaining({ type: SyncEntityType.SyncCompleteV1 }),
]);
});
it('should not sync an asset face or asset face delete for a user in a different cluster group', async () => {
const { auth, ctx } = await setup();
const personRepo = ctx.get(PersonRepository);
const { user: user2 } = await ctx.newUser();
const { session } = await ctx.newSession({ userId: user2.id });
const { asset } = await ctx.newAsset({ ownerId: user2.id });
const { person } = await ctx.newPerson({ ownerId: user2.id });
const { assetFace } = await ctx.newAssetFace({ assetId: asset.id, personGroupId: person.personGroupId });
const auth2 = factory.auth({ session, user: user2 });
expect(await ctx.syncStream(auth2, [SyncRequestType.AssetFacesV3])).toEqual([
expect.objectContaining({ type: SyncEntityType.AssetFaceV3 }),
expect.objectContaining({ type: SyncEntityType.SyncCompleteV1 }),
]);
await ctx.assertSyncIsComplete(auth, [SyncRequestType.AssetFacesV3]);
await personRepo.deleteAssetFace(assetFace.id);
expect(await ctx.syncStream(auth2, [SyncRequestType.AssetFacesV3])).toEqual([
expect.objectContaining({ type: SyncEntityType.AssetFaceDeleteV1 }),
expect.objectContaining({ type: SyncEntityType.SyncCompleteV1 }),
]);
await ctx.assertSyncIsComplete(auth, [SyncRequestType.AssetFacesV3]);
});
it('should detect and sync a deleted asset face without a person', async () => {
const { auth, ctx } = await setup();
const personRepo = ctx.get(PersonRepository);
const { asset } = await ctx.newAsset({ ownerId: auth.user.id });
const { assetFace } = await ctx.newAssetFace({ assetId: asset.id });
await personRepo.deleteAssetFace(assetFace.id);
expect(await ctx.syncStream(auth, [SyncRequestType.AssetFacesV3])).toEqual([
{ ack: expect.any(String), data: { assetFaceId: assetFace.id }, type: SyncEntityType.AssetFaceDeleteV1 },
expect.objectContaining({ type: SyncEntityType.SyncCompleteV1 }),
]);
});
});