diff --git a/snprc_ehr/resources/etls/ExportBiocontainmentObservations.xml b/snprc_ehr/resources/etls/ExportBiocontainmentObservations.xml new file mode 100644 index 000000000..468571492 --- /dev/null +++ b/snprc_ehr/resources/etls/ExportBiocontainmentObservations.xml @@ -0,0 +1,29 @@ + + + ExportBiocontainmentObservations + Sync Biocontainment Observations edits to CAMP + + + Copy to target + + + + + + + + + + + + + Fan merged edits out to CAMP's live BiocontainmentObservation table + + + + + + + + + \ No newline at end of file diff --git a/snprc_ehr/resources/queries/snprc_ehr/ExportBiocontainmentObservations.sql b/snprc_ehr/resources/queries/snprc_ehr/ExportBiocontainmentObservations.sql new file mode 100644 index 000000000..d58e3b32a --- /dev/null +++ b/snprc_ehr/resources/queries/snprc_ehr/ExportBiocontainmentObservations.sql @@ -0,0 +1,54 @@ + + +/******************************************************** + Copyright (c) 2026 SNPRC + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. + + This Query Selects BiocontainmentObservations rows edited in TAC so they can be pushed back to + CAMP's dbo.BiocontainmentObservation. + ********************************************************/ + +-- ========================================================================================== +-- Author: Ram +-- Create date: 09/10/26 +-- Description: Selects BiocontainmentObservations rows edited in TAC so they can be pushed back to +-- CAMP's dbo.BiocontainmentObservation via the TAC to CAMP export ETL +-- (resources/etls/ExportBiocontainmentObservations.xml) +-- ========================================================================================== +SELECT + o.Id, + o.date AS ObservationDate, + o.Location, + o.WeightLoss, + o.WeightLossCO AS WeightLoss_CO, + o.TemperatureChange, + o.TemperatureChangeCO AS TemperatureChange_CO, + o.Responsiveness, + o.HairCoat, + o.Respiration, + o.Petechia, + o.PetechiaCO AS Petechia_CO, + o.Bleeding, + o.NasalDischarge, + o.FeedEaten, + o.FeedEatenCO AS FeedEaten_CO, + o.FoodEnrichment, + o.FoodEnrichmentCO AS FoodEnrichment_CO, + o.Stool, + o.FluidIntake, + o.Dehydration, + o.DehydrationCO AS Dehydration_CO, + o.Comment, + o.objectid, + o.modified, + o.modifiedBy, + o.modifiedBy.email AS modifiedByEmail, + o.created, + o.createdBy, + o.Container +FROM study.BiocontainmentObservations o +WHERE o.modifiedBy.email IS NULL OR o.modifiedBy.email NOT LIKE '%@noreply-txbiomed.org' \ No newline at end of file diff --git a/snprc_ehr/resources/source_queries/create_CAMP_BiocontainmentObservations.sql b/snprc_ehr/resources/source_queries/create_CAMP_BiocontainmentObservations.sql new file mode 100644 index 000000000..cbc415fe3 --- /dev/null +++ b/snprc_ehr/resources/source_queries/create_CAMP_BiocontainmentObservations.sql @@ -0,0 +1,227 @@ +/******************************************************************* +Creates a TAC_src landing table + a stored procedure to receive Biocontainment +Observations edits made in TAC, for fan-out into CAMP's live +dbo.BiocontainmentObservations table + + *******************************************************************/ + +IF NOT EXISTS (SELECT name FROM sys.schemas WHERE name = N'TAC_src') +EXEC('CREATE SCHEMA [TAC_src] AUTHORIZATION [DBO]'); + +DROP TABLE IF EXISTS TAC_src.BiocontainmentObservations; + +CREATE TABLE TAC_src.BiocontainmentObservations +( + tid INT IDENTITY, + Id VARCHAR(6) NOT NULL, + ObservationDate DATETIME NOT NULL, + Location DECIMAL(6,2) NOT NULL, + WeightLoss INT NULL, + WeightLoss_CO INT NULL, + TemperatureChange INT NULL, + TemperatureChange_CO INT NULL, + Responsiveness INT NULL, + HairCoat INT NULL, + Respiration INT NULL, + Petechia INT NULL, + Petechia_CO INT NULL, + Bleeding INT NULL, + NasalDischarge INT NULL, + FeedEaten INT NULL, + FeedEaten_CO INT NULL, + FoodEnrichment INT NULL, + FoodEnrichment_CO INT NULL, + Stool VARCHAR(2) NULL, + FluidIntake INT NULL, + Dehydration INT NULL, + Dehydration_CO INT NULL, + Comment VARCHAR(80) NULL, + ETLDownDate DATETIME DEFAULT GETDATE(), + Created DATETIME NULL, + CreatedBy INT NULL, + CreatedByEmail VARCHAR(128) NULL, + Modified DATETIME NULL, + ModifiedBy INT NULL, + ModifiedByEmail VARCHAR(128) NULL, + Container UNIQUEIDENTIFIER NOT NULL, + objectid UNIQUEIDENTIFIER NOT NULL, + FannedOutDate DATETIME NULL, + CONSTRAINT PK_TAC_BIOCONTAINMENTOBSERVATIONS + PRIMARY KEY CLUSTERED (tid ASC) +) +GO + + +DROP PROCEDURE IF EXISTS [TAC_src].[usp_FanOutBiocontainmentObservations]; +GO + +CREATE PROCEDURE [TAC_src].[usp_FanOutBiocontainmentObservations] AS +-- ========================================================================================== +-- Author: Ram +-- Create date: 09/10/26 +-- Description: Fans TAC edits out of the TAC_src.BiocontainmentObservations landing table into CAMP's +-- live dbo.BiocontainmentObservation table. Called as step2 of the TAC to CAMP +-- export ETL (resources/etls/ExportBiocontainmentObservations.xml), after step1's merge +-- into TAC_src.BiocontainmentObservations has already committed. +-- ========================================================================================== +BEGIN + SET NOCOUNT ON; + + DECLARE @errno INT, + @errmsg VARCHAR(255); + + DECLARE @matched TABLE + ( + targetTid INT NOT NULL PRIMARY KEY, + landingTid INT NOT NULL, + landingObjectId UNIQUEIDENTIFIER NOT NULL, + capturedTimestamp BINARY(8) NOT NULL, + Location DECIMAL(6,2) NULL, + WeightLoss INT NULL, + WeightLoss_CO INT NULL, + TemperatureChange INT NULL, + TemperatureChange_CO INT NULL, + Responsiveness INT NULL, + HairCoat INT NULL, + Respiration INT NULL, + Petechia INT NULL, + Petechia_CO INT NULL, + Bleeding INT NULL, + NasalDischarge INT NULL, + FeedEaten INT NULL, + FeedEaten_CO INT NULL, + FoodEnrichment INT NULL, + FoodEnrichment_CO INT NULL, + Stool VARCHAR(2) NULL, + FluidIntake INT NULL, + Dehydration INT NULL, + Dehydration_CO INT NULL, + Comment VARCHAR(80) NULL + ); + + -- capturedTimestamp is CAMP's rowversion for this row at match time. The + -- final UPDATE below only applies if the row's rowversion is still this + -- value, so a concurrent edit to the same CAMP row (e.g. from the Android + -- app) between this SELECT and the UPDATE loses the race safely instead + -- of being silently overwritten - the loser is left unmatched, its + -- FannedOutDate stays NULL, and it's picked up again on the next run. + INSERT INTO @matched + SELECT b.tid, l.tid, l.objectid, b.timestamp, l.Location, l.WeightLoss, l.WeightLoss_CO, + l.TemperatureChange, l.TemperatureChange_CO, l.Responsiveness, + l.HairCoat, l.Respiration, l.Petechia, l.Petechia_CO, l.Bleeding, + l.NasalDischarge, l.FeedEaten, l.FeedEaten_CO, l.FoodEnrichment, + l.FoodEnrichment_CO, l.Stool, l.FluidIntake, l.Dehydration, + l.Dehydration_CO, l.Comment + FROM TAC_src.BiocontainmentObservations l + INNER JOIN dbo.BiocontainmentObservation b ON b.ObjectId = l.objectid + WHERE l.FannedOutDate IS NULL OR l.FannedOutDate < l.Modified; + + -- Log (non-fatal) any not-yet-fanned-out landing rows that didn't match + -- a CAMP row. + IF EXISTS (SELECT 1 FROM TAC_src.BiocontainmentObservations l + WHERE (l.FannedOutDate IS NULL OR l.FannedOutDate < l.Modified) + AND NOT EXISTS (SELECT 1 FROM @matched m WHERE m.landingObjectId = l.objectid)) + BEGIN + DECLARE @unmatchedCount INT = + ( + SELECT COUNT(*) FROM TAC_src.BiocontainmentObservations l + WHERE (l.FannedOutDate IS NULL OR l.FannedOutDate < l.Modified) + AND NOT EXISTS (SELECT 1 FROM @matched m WHERE m.landingObjectId = l.objectid) + ); + RAISERROR('TAC_src.usp_FanOutBiocontainmentObservations: %d row(s) had no matching ObjectId in dbo.BiocontainmentObservation and were skipped', 10, 1, @unmatchedCount) WITH NOWAIT; + END; + + IF NOT EXISTS (SELECT 1 FROM @matched) + RETURN; + + + DECLARE @updated TABLE (targetTid INT NOT NULL PRIMARY KEY, landingTid INT NOT NULL); + + BEGIN TRANSACTION; + + -- Optimistic concurrency: only apply to rows whose rowversion still + -- matches what we captured above. A row that was concurrently modified + -- (e.g. by the Android app) between the match and here is silently + -- excluded from the OUTPUT, so it gets neither an editCommentBuffer + -- entry nor its FannedOutDate set - it stays eligible and is retried + -- on the next run against CAMP's now-current data. + UPDATE b + SET Location = COALESCE(m.Location, b.Location), + WeightLoss = COALESCE(m.WeightLoss, b.WeightLoss), + WeightLoss_CO = COALESCE(m.WeightLoss_CO, b.WeightLoss_CO), + TemperatureChange = COALESCE(m.TemperatureChange, b.TemperatureChange), + TemperatureChange_CO = COALESCE(m.TemperatureChange_CO, b.TemperatureChange_CO), + Responsiveness = COALESCE(m.Responsiveness, b.Responsiveness), + HairCoat = COALESCE(m.HairCoat, b.HairCoat), + Respiration = COALESCE(m.Respiration, b.Respiration), + Petechia = COALESCE(m.Petechia, b.Petechia), + Petechia_CO = COALESCE(m.Petechia_CO, b.Petechia_CO), + Bleeding = COALESCE(m.Bleeding, b.Bleeding), + NasalDischarge = COALESCE(m.NasalDischarge, b.NasalDischarge), + FeedEaten = COALESCE(m.FeedEaten, b.FeedEaten), + FeedEaten_CO = COALESCE(m.FeedEaten_CO, b.FeedEaten_CO), + FoodEnrichment = COALESCE(m.FoodEnrichment, b.FoodEnrichment), + FoodEnrichment_CO = COALESCE(m.FoodEnrichment_CO, b.FoodEnrichment_CO), + Stool = COALESCE(m.Stool, b.Stool), + FluidIntake = COALESCE(m.FluidIntake, b.FluidIntake), + Dehydration = COALESCE(m.Dehydration, b.Dehydration), + Dehydration_CO = COALESCE(m.Dehydration_CO, b.Dehydration_CO), + Comment = COALESCE(m.Comment, b.Comment) + OUTPUT inserted.tid, m.landingTid INTO @updated (targetTid, landingTid) + FROM dbo.BiocontainmentObservation b + INNER JOIN @matched m ON m.targetTid = b.tid AND b.timestamp = m.capturedTimestamp; + + IF @@ERROR <> 0 + BEGIN + SELECT @errno = 30021, @errmsg = 'Error occurred fanning TAC_src.BiocontainmentObservations edits out into dbo.BiocontainmentObservation.'; + GOTO error; + END; + + IF EXISTS (SELECT 1 FROM @matched m WHERE NOT EXISTS (SELECT 1 FROM @updated u WHERE u.targetTid = m.targetTid)) + BEGIN + DECLARE @conflictCount INT = (SELECT COUNT(*) FROM @matched m WHERE NOT EXISTS (SELECT 1 FROM @updated u WHERE u.targetTid = m.targetTid)); + RAISERROR('TAC_src.usp_FanOutBiocontainmentObservations: %d row(s) skipped due to a concurrent CAMP-side modification (rowversion mismatch); will retry next run', 10, 1, @conflictCount) WITH NOWAIT; + END; + + IF NOT EXISTS (SELECT 1 FROM @updated) + BEGIN + COMMIT TRANSACTION; + RETURN; + END; + + INSERT INTO dbo.editCommentBuffer (tableName, tid, editComment) + SELECT 'BiocontainmentObservation', targetTid, 'TAC ETL sync' + FROM @updated; + + IF @@ERROR <> 0 + BEGIN + SELECT @errno = 30022, @errmsg = 'Error occurred inserting into dbo.editCommentBuffer.'; + GOTO error; + END; + + UPDATE l + SET FannedOutDate = GETDATE() + FROM TAC_src.BiocontainmentObservations l + INNER JOIN @updated u ON u.landingTid = l.tid; + + IF @@ERROR <> 0 + BEGIN + SELECT @errno = 30023, @errmsg = 'Error occurred marking TAC_src.BiocontainmentObservations rows as fanned out.'; + GOTO error; + END; + + COMMIT TRANSACTION; + RETURN; + + error: + IF @@TRANCOUNT > 0 + ROLLBACK TRANSACTION; + RAISERROR('%d: %s', 16, 0, @errno, @errmsg); +END +GO + + +GRANT DELETE, INSERT, REFERENCES, SELECT, UPDATE ON TAC_src.BiocontainmentObservations TO z_labkey; +GRANT VIEW DEFINITION ON TAC_src.BiocontainmentObservations TO z_labkey; +GRANT EXECUTE ON [TAC_src].[usp_FanOutBiocontainmentObservations] TO z_labkey; +GO \ No newline at end of file diff --git a/snprc_ehr/src/org/labkey/snprc_ehr/steps/BiocontainmentObservationsFanOutTask.java b/snprc_ehr/src/org/labkey/snprc_ehr/steps/BiocontainmentObservationsFanOutTask.java new file mode 100644 index 000000000..a55c562f2 --- /dev/null +++ b/snprc_ehr/src/org/labkey/snprc_ehr/steps/BiocontainmentObservationsFanOutTask.java @@ -0,0 +1,78 @@ +package org.labkey.snprc_ehr.steps; + +import org.jetbrains.annotations.NotNull; +import org.labkey.api.data.Container; +import org.labkey.api.data.DbSchema; +import org.labkey.api.data.DbScope; +import org.labkey.api.data.SqlExecutor; +import org.labkey.api.data.TableInfo; +import org.labkey.api.di.TaskRefTaskImpl; +import org.labkey.api.pipeline.PipelineJob; +import org.labkey.api.pipeline.PipelineJobException; +import org.labkey.api.pipeline.RecordedActionSet; +import org.labkey.api.query.QueryService; +import org.labkey.api.query.UserSchema; + +/** + * Calls TAC_src.usp_FanOutBiocontainmentObservations (in the animal database, via the + * tacSrc external schema) after step1 of ExportBiocontainmentObservations.xml has + * already merged edited rows into TAC_src.BiocontainmentObservations and committed. + * + * This exists as a separate post-merge step, rather than a trigger on + * TAC_src.BiocontainmentObservations, because LabKey's bulkLoad="true" + * targetOption="merge" write uses an OUTPUT clause (without INTO) to + * report rows copied - and SQL Server disallows OUTPUT without INTO on any + * statement whose target table has an enabled trigger for that action. + * A trigger there failed every run with "Optimistic concurrency exception: + * Table deleted" (confirmed by toggling the trigger disabled/enabled - + * disabled succeeds every time, enabled fails every time, including on a + * 0-row merge). Running the fan-out as a stored procedure after step1 has + * already committed avoids that conflict entirely + */ +public class BiocontainmentObservationsFanOutTask extends TaskRefTaskImpl +{ + private static final String TAC_SRC_SCHEMA_NAME = "tacSrc"; + private static final String LANDING_TABLE_NAME = "BiocontainmentObservations"; + + private void runFanOut(PipelineJob job) throws PipelineJobException + { + Container c = job.getContainer(); + + UserSchema schema = QueryService.get().getUserSchema(job.getUser(), c, TAC_SRC_SCHEMA_NAME); + if (schema == null) + throw new PipelineJobException("Could not find schema " + TAC_SRC_SCHEMA_NAME + " in " + c.getPath()); + + TableInfo ti = schema.getTable(LANDING_TABLE_NAME); + if (ti == null) + throw new PipelineJobException("Could not find table " + LANDING_TABLE_NAME + " in schema " + TAC_SRC_SCHEMA_NAME + " in " + c.getPath()); + + DbSchema dbSchema = ti.getSchema(); + DbScope scope = dbSchema.getScope(); + + job.getLogger().info("Calling TAC_src.usp_FanOutBiocontainmentObservations"); + new SqlExecutor(scope).execute("EXEC TAC_src.usp_FanOutBiocontainmentObservations"); + job.getLogger().info("Fan-out to CAMP complete"); + } + + @Override + public RecordedActionSet run(@NotNull PipelineJob job) + { + try + { + runFanOut(job); + } + catch (Exception e) + { + // Unlike BiocontainmentObservationsAuditLogTask (where a logging failure + // is acceptable to swallow), a failure here means TAC edits did + // not actually reach CAMP - this must fail the pipeline step + // rather than report success, so a real problem (e.g. a CHECK + // constraint rejecting an out-of-range value, a missing + // dependency like editCommentBuffer) is visible to whoever + // monitors these ETL jobs instead of being silently swallowed. + job.getLogger().error(e.getMessage(), e); + throw new RuntimeException("BiocontainmentObservationsFanOutTask failed: " + e.getMessage(), e); + } + return new RecordedActionSet(makeRecordedAction()); + } +}