refactor(catalog-backend): make stitchTicket required in performStitching
Now that immediate mode is removed, the stitch ticket is always provided by the deferred worker. Remove the optional marker and all the conditional guards that checked for its presence. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com> Signed-off-by: Fredrik Adelöw <freben@gmail.com>
This commit is contained in:
@@ -349,12 +349,20 @@ describe.each(databases.eachSupportedId())(
|
||||
},
|
||||
]);
|
||||
|
||||
await markForStitching({ knex, entityRefs: ['k:ns/n'] });
|
||||
|
||||
const stitchLogger = mockServices.logger.mock();
|
||||
await expect(
|
||||
performStitching({
|
||||
knex,
|
||||
logger: stitchLogger,
|
||||
entityRef: 'k:ns/n',
|
||||
stitchTicket: (
|
||||
await knex('stitch_queue')
|
||||
.where('entity_ref', 'k:ns/n')
|
||||
.select('stitch_ticket')
|
||||
.first()
|
||||
)?.stitch_ticket,
|
||||
}),
|
||||
).resolves.toBe('changed');
|
||||
|
||||
|
||||
@@ -52,10 +52,9 @@ export async function performStitching(options: {
|
||||
knex: Knex | Knex.Transaction;
|
||||
logger: LoggerService;
|
||||
entityRef: string;
|
||||
stitchTicket?: string;
|
||||
stitchTicket: string;
|
||||
}): Promise<'changed' | 'unchanged' | 'abandoned'> {
|
||||
const { knex, logger, entityRef } = options;
|
||||
const stitchTicket = options.stitchTicket;
|
||||
const { knex, logger, entityRef, stitchTicket } = options;
|
||||
|
||||
// The entity is removed from the stitch queue on ANY completion, except when
|
||||
// an exception is thrown. In the latter case, the entity will be retried at a
|
||||
@@ -205,7 +204,7 @@ export async function performStitching(options: {
|
||||
// stitch cycle).
|
||||
const isMySQL = String(knex.client.config.client).includes('mysql');
|
||||
|
||||
if (stitchTicket && isMySQL) {
|
||||
if (isMySQL) {
|
||||
const ticketValid = await knex<DbStitchQueueRow>('stitch_queue')
|
||||
.where('entity_ref', entityRef)
|
||||
.where('stitch_ticket', stitchTicket)
|
||||
@@ -229,7 +228,7 @@ export async function performStitching(options: {
|
||||
.onConflict('entity_id')
|
||||
.merge(['final_entity', 'hash', 'last_updated_at']);
|
||||
|
||||
if (stitchTicket && !isMySQL) {
|
||||
if (!isMySQL) {
|
||||
upsert = upsert.where(
|
||||
knex.raw(
|
||||
'exists (select 1 from stitch_queue where entity_ref = ? and stitch_ticket = ?)',
|
||||
@@ -244,7 +243,7 @@ export async function performStitching(options: {
|
||||
// database engines (row IDs vs row counts vs empty arrays), so we
|
||||
// check the hash directly — we already know hash !== previousHash
|
||||
// from the check above, so a mismatch means the write was blocked.
|
||||
if (stitchTicket && !isMySQL) {
|
||||
if (!isMySQL) {
|
||||
const written = await knex<DbFinalEntitiesRow>('final_entities')
|
||||
.where('entity_id', entityId)
|
||||
.where('hash', hash)
|
||||
@@ -265,7 +264,7 @@ export async function performStitching(options: {
|
||||
removeFromStitchQueueOnCompletion = false;
|
||||
throw error;
|
||||
} finally {
|
||||
if (removeFromStitchQueueOnCompletion && stitchTicket) {
|
||||
if (removeFromStitchQueueOnCompletion) {
|
||||
await markDeferredStitchCompleted({
|
||||
knex: knex,
|
||||
entityRef,
|
||||
|
||||
@@ -128,7 +128,7 @@ export class DefaultStitcher {
|
||||
|
||||
async #stitchOne(options: {
|
||||
entityRef: string;
|
||||
stitchTicket?: string;
|
||||
stitchTicket: string;
|
||||
stitchRequestedAt?: DateTime;
|
||||
}) {
|
||||
const track = this.tracker.stitchStart({
|
||||
|
||||
Reference in New Issue
Block a user