Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion bps/LSSTCam/bps_Daytime.yaml
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
pipelineYaml: $AP_PIPE_DIR/pipelines/LSSTCam/ApPipe.yaml
pipelineYaml: $AP_PIPE_DIR/pipelines/LSSTCam/ApPipeDaytime.yaml

project: ApPipe
campaign: AP-daytime
Expand Down
5 changes: 1 addition & 4 deletions bps/clustering/clustering_Daytime.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -21,9 +21,6 @@ cluster:
preload:
pipetasks: mpSkyEphemerisQuery,getRegionTimeFromVisit
dimensions: group,detector
preloadApdb:
pipetasks: loadDiaCatalogs,analyzeLoadDiaCatalogsMetrics
dimensions: group,detector
singleFrame:
pipetasks: isr,calibrateImage,analyzePreliminarySummaryStats
dimensions: visit,detector
Expand Down Expand Up @@ -52,6 +49,6 @@ ordering:
associationOrder:
ordering_type: sort
findDependencyMethod: sink
labels: preloadApdb,association
labels: association
dimensions: visit
blocking: False
5 changes: 5 additions & 0 deletions doc/lsst.ap.pipe/pipeline-overview.rst
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,11 @@ to verify the output.
:doc:`ap_pipe <index>` is entirely written in Python. Key contents include:

- :file:`ApPipe.yaml`: a `~lsst.pipe.base.Pipeline` configuration for running the entire AP Pipeline.
- :file:`LSSTCam/ApPipeDaytime.yaml`: the variant used for daytime (non-real-time) LSSTCam processing.
It drops ``loadDiaCatalogs`` and has `~lsst.ap.association.DiaPipelineTask` read the DIAObject and
DIASource history from the APDB during association, so that the duplicate DIASource check sees rows
written by any earlier pass over the same image.
Prompt Processing must not use it; the preload is what keeps the APDB out of its latency-critical path.

By default the pipeline is limited to running on data taken in filter bands whose names match those used by the Rubin Observatory LSST Camera (that is `ugrizy`).
In order to run on bands outside of these filters, one must add the associated columns to the `~lsst.dax.apdb.Apdb` schema and add the band names to the config of `~lsst.ap.association.DiaPipelineTask`.
50 changes: 50 additions & 0 deletions pipelines/LSSTCam/ApPipeDaytime.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
description: >-
AP pipeline for daytime (non-real-time) processing of LSSTCam data.

Association reads the DiaObject, DiaSource, and DiaForcedSource history
directly from the APDB rather than consuming preloaded catalogs, so the
duplicate DiaSource check in associateApdb can see rows written by an
earlier pass over the same image. The preloaded catalogs are built before
the image is processed and so cannot contain those rows; reusing them let
a reprocessing run write duplicate diaSources to the APDB (DM-55633).
loadDiaCatalogs is therefore dropped from this pipeline entirely, which
also stops --skip-existing-in from reviving a stale preload.

Prompt Processing must not use this pipeline. Its preload exists to keep
the APDB out of the latency-critical path, and it guards its own retries
with Apdb.containsVisitDetector.
instrument: lsst.obs.lsst.LsstCam
imports:
- location: $AP_PIPE_DIR/pipelines/LSSTCam/ApPipe.yaml

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
- location: $AP_PIPE_DIR/pipelines/LSSTCam/ApPipe.yaml
- location: $AP_PIPE_DIR/pipelines/LSSTCam/ApPipe.yaml
labeledSubsetModifyMode: EDIT

Claude suggests that you can avoid redefining the subsets below (which is fragile to future changes):

The comment says excluding a task deletes every subset that named it — true for the default DROP, but the import supports labeledSubsetModifyMode: EDIT, which removes just the missing labels. drp_pipe/pipelines/LSSTCam/DRP.yaml:5 already uses it. This would eliminate the duplicated list, which will otherwise silently drift when tasks are added to the base apPipe subset (nothing in test_pipelines.py checks subset completeness — self.synonyms isn't actually asserted against anything in the visible test code). Worth checking EDIT doesn't leave any other base subset (prompt, etc.) in a misleading state.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Wow! That is great, I did not know about EDIT mode before. I will test it, and definitely use it if it works as advertised.

exclude:
- loadDiaCatalogs
- analyzeLoadDiaCatalogsMetrics
# EDIT drops the excluded labels from the inherited subsets. The
# default, DROP, would delete preload and promptQaMetrics outright
# and force apPipe to be restated here.
labeledSubsetModifyMode: EDIT

tasks:
associateApdb:
class: lsst.ap.association.DiaPipelineTask
config:
# Load all of the APDB catalogs in association, instead of relying on
# preloaded catalogs from loadDiaCatalogs
doReloadAllApdbCatalogs: True
analyzeDiaSourceAssociationMetrics:
class: lsst.analysis.tools.tasks.TaskMetadataAnalysisTask
config:
# Publish the APDB read timings that analyzeLoadDiaCatalogsMetrics used
# to provide. The metric names match the ones that task emitted, so
# they reach the same Sasquatch topics as in Prompt Processing; only
# the dataset type carrying them changes, since the timings now come
# from a visit-dimensioned quantum.
# Setting `metrics` replaces the dict, so the inherited entries are
# repeated here.
atools.associationMetadataMetrics.metrics:
numTotalSolarSystemObjects: ct
numAssociatedSsObjects: ct
writeToApdbDuration: s
loadDiaObjectsDuration: s
loadDiaSourcesDuration: s
loadDiaForcedSourcesDuration: s
10 changes: 8 additions & 2 deletions pipelines/_ingredients/ApPipe.yaml

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Claude suggests you don' t need to change these contracts:

pipe_base's subset_from_labels prunes any contract whose text mentions an excluded label
(pipelineIR.py:970). That's why the untouched loadDiaCatalogs.apdb_config_url == associateApdb.apdb_config_url contract (line 316) doesn't break Daytime. So the associateApdb.doReloadAllApdbCatalogs or (...) guards only matter if someone runs both loadDiaCatalogs and doReloadAllApdbCatalogs: True — a configuration that wastes the preload and should probably fail loudly rather than be blessed.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It looks like you are right, but this should also become moot if switching to EDIT mode works.

Original file line number Diff line number Diff line change
Expand Up @@ -312,8 +312,14 @@ subsets:
Requires prompt subset to be run first.
contracts:
- detectAndMeasureDiaSource.doSkySources == filterDiaSource.doRemoveSkySources
# Both loadDiaCatalogs and associateApdb connect to the APDB, so make sure they use the same configuration
- loadDiaCatalogs.apdb_config_url == associateApdb.apdb_config_url
# Both loadDiaCatalogs and associateApdb connect to the APDB, so make sure
# they use the same configuration. associateApdb must not also reload the
# full history, which would leave the preload unused; ApPipeDaytime.yaml
# drops loadDiaCatalogs instead.
- contract: loadDiaCatalogs.apdb_config_url == associateApdb.apdb_config_url
and not associateApdb.doReloadAllApdbCatalogs
msg: "loadDiaCatalogs and associateApdb must share an APDB config, and
associateApdb.doReloadAllApdbCatalogs must be False when loadDiaCatalogs runs"
# to reduce latency, we need two calls to the sattle service when active
- calibrateImage.run_sattle == detectAndMeasureDiaSource.run_sattle
# Inputs and outputs must match. For consistency, contracts are written in execution order:
Expand Down
2 changes: 1 addition & 1 deletion scripts/LSSTCam/submit_ap_daytime.sh
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ BLOCKS_SQL="($(printf "'%s'," $BLOCKS | sed 's/,$//'))"

# Pipeline and butler config must mirror bps_Daytime.yaml — we replicate them
# here because we build the quantum graph ourselves before calling BPS.
PIPELINE_YAML="${AP_PIPE_DIR}/pipelines/LSSTCam/ApPipe.yaml"
PIPELINE_YAML="${AP_PIPE_DIR}/pipelines/LSSTCam/ApPipeDaytime.yaml"
APDB_CONFIG="s3://embargo@rubin-summit-users/apdb_config/cassandra/pp_apdb_lsstcam.yaml"
BUTLER_CONFIG="embargo"
INPUT_COLLECTIONS="LSSTCam/defaults,LSSTCam/templates,LSSTCam/runs/prompt-${DAY_OBS}"
Expand Down
1 change: 1 addition & 0 deletions tests/test_pipelines.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ def setUp(self):
# Each pipeline file should have a subset that represents it in
# higher-level pipelines.
self.synonyms = {"ApPipe.yaml": "apPipe",
"ApPipeDaytime.yaml": "apPipe",
"ApPipeWithIsrTaskLSST.yaml": "apPipe",
"ApPipeWithPreconvolution.yaml": "apPipe",
"ApPipeWithFakes.yaml": "apPipe",
Expand Down
Loading