Own transaction when raising an event
This commit is contained in:
@@ -26,26 +26,26 @@ module.exports = async()=>{ // called by server.js
|
||||
admin.on (['CREATE','UPDATE'], 'Books', async ({data:{supplier_ID: supplierId}}, next) => {
|
||||
// Using Promise.all(...) to parallelize local write, i.e. next(), and replication
|
||||
|
||||
/*
|
||||
// ERROR: Reference integrity is violated for association "supplier"
|
||||
if (supplierId) return Promise.all ([ next(), async()=>{
|
||||
let replicated = await db.exists (Suppliers, supplierId)
|
||||
if (!replicated) await replicate (supplierId, 'initial')
|
||||
}])
|
||||
else return next() //> don't forget to pass down the interceptor stack
|
||||
*/
|
||||
const replicateIfNotExists = async()=>{
|
||||
let replicated = await db.exists (Suppliers, supplierId);
|
||||
if (!replicated) await replicate (supplierId, 'initial');
|
||||
};
|
||||
|
||||
let replicated = await db.exists (Suppliers, supplierId);
|
||||
if (!replicated) await replicate (supplierId, 'initial');
|
||||
return next();
|
||||
if (supplierId) return Promise.all ([ next(), replicateIfNotExists() ])
|
||||
else return next() //> don't forget to pass down the interceptor stack
|
||||
})
|
||||
|
||||
})
|
||||
|
||||
// Subscribe to changes in the S4 origin of Suppliers data
|
||||
// REVISIT: cds context is still from the UPDAT method when running in same programm, but should
|
||||
// be a separate
|
||||
S4bupa.on ('BusinessPartners/Changed', async msg => { //> would be great if we had batch events from S/4
|
||||
let replicas = await SELECT('ID').from (Suppliers) .where ('ID in', msg.businessPartners)
|
||||
return replicate (replicas.map(each => each.ID))
|
||||
await new Promise( resolve => setTimeout( resolve, 1000 ));
|
||||
const tx = cds.db.tx(msg);
|
||||
let replicas = await tx.run(SELECT('ID').from (Suppliers) .where ('ID in', msg.data.businessPartners));
|
||||
await replicateTx(replicas.map(each => each.ID), undefined, tx, msg);
|
||||
await tx.commit();
|
||||
})
|
||||
|
||||
/**
|
||||
@@ -57,6 +57,7 @@ module.exports = async()=>{ // called by server.js
|
||||
if (!Array.isArray(IDs)) IDs = [ IDs ]
|
||||
// TODO: Doesn't work when running in same process with mocked API_BUSINESS_PARTNER
|
||||
|
||||
// TODO: Issue
|
||||
let suppliers = await S4bupa.read (Suppliers).where(...([[]].concat(IDs).reduce( (where, id, index ) => { where.push(`${index>1 ? "OR ":""}ID = `, id); return where })));
|
||||
//let suppliers = await S4bupa.read (Suppliers).where('ID in',IDs)
|
||||
if (_initial) return db.insert (suppliers) .into (Suppliers) //> using bulk insert
|
||||
@@ -65,4 +66,17 @@ module.exports = async()=>{ // called by server.js
|
||||
))
|
||||
}
|
||||
|
||||
async function replicateTx (IDs,_initial, tx, msg) {
|
||||
if (!Array.isArray(IDs)) IDs = [ IDs ]
|
||||
// TODO: Doesn't work when running in same process with mocked API_BUSINESS_PARTNER
|
||||
|
||||
// TODO: Issue
|
||||
let suppliers = await S4bupa.tx(msg).read (Suppliers).where(...([[]].concat(IDs).reduce( (where, id, index ) => { where.push(`${index>1 ? "OR ":""}ID = `, id); return where })));
|
||||
//let suppliers = await S4bupa.read (Suppliers).where('ID in',IDs)
|
||||
if (_initial) return tx.insert (suppliers) .into (Suppliers) //> using bulk insert
|
||||
else return Promise.all(suppliers.map ( //> parallelizing updates
|
||||
each => tx.update (Suppliers,each.ID) .with (each)
|
||||
))
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user