Runs one execution of a configured integration pipeline: a chain of PipelineStep instances built from an account's pipeline XML by PipelineManager and driven by start(), each step passing data to the next until the chain completes or throws. A Pipeline is created for one EndPointMapping and website, carries the state that steps read and write as they run (attributes, current position, warnings/failures/info messages, output stream), and is available to JavaScript steps such as JsRowStep as the pipeline variable so scripts can call back into it. On completion it records a PipelineExecution row (when the endpoint mapping requests it) and fires a PipelineProcessEvent so other parts of the platform can react to the run.
Properties
| Property | Returns | Description |
|---|---|---|
| attributes | Map<String,Object> | Mutable map of name/value pairs shared between all steps in this run, used to pass data such as fromDate and toDate along the pipeline and to build expression contexts for MVEL templates. |
| cache | Map<String,Object> | A map for caching data generated during pipeline processing, such as lookups a step wants to avoid repeating for every row. Lazily created on first access, and may be cleared at any time by clearSession. |
| currentPosition | String | Description of where processing is currently up to within the input data, such as the file, worksheet, column and row, set by steps as they work. Attached to warning, failure and info messages as they are recorded. |
| currentProfile | Profile | The profile this pipeline is running as, used for permissions and for the PipelineProcessEvent fired when the run finishes. Defaults to the current request's principal profile if not explicitly set before start. |
| dataSession | DataSession | The branch containing this pipeline's XML definition and any source code its steps load, such as JsRowStep scripts. |
| destinationAddress | String | The address data is being sent to for this run, such as an output file path or endpoint URL, as set by steps or by setDestinationAddress. Recorded on the PipelineExecution row. |
| endPointMapping | EndPointMapping | The endpoint configuration this pipeline is executing, which supplies its credentials, file naming, duplicate-prevention and execution-recording settings. |
| executionRecord | PipelineExecution | The PipelineExecution row recorded for this run, set once finishRecordExecution has completed. |
| expressionAttributes | Map<String,Object> | |
| failures | List<IntegrationMessage> | The failure messages recorded so far by addFailure, each a definite problem encountered while running this pipeline. |
| fileHash | String | Hash of the uploaded or input file this pipeline is processing, used to detect duplicate uploads and recorded on the PipelineExecution row. |
| fileName | String | Name of the file this pipeline is producing or consuming, evaluated from the endpoint mapping's file name template on first access if it has not already been set explicitly. |
| head | PipelineStep | The first step in this pipeline's chain, which start() prepares and executes. |
| infos | List<IntegrationMessage> | The informational messages recorded so far by addInfo. |
| orgRoot | OrganisationRootFolder | Looks up the OrganisationRootFolder for this pipeline's organisation. |
| out | OutputStream | The output stream that pipeline steps write their result to, as set by setOut or by start. |
| pipelineExecutionId | Long | Database id of the PipelineExecution record for this run, once one has been created. Null before prepareRecordExecution has run, and always null if the endpoint mapping does not record executions. |
| pipelineManager | PipelineManager | The manager that built this pipeline and tracks it while it runs. |
| pipelinePath | String | Path, within the website's branch, of the pipeline XML definition this pipeline was built from. |
| processTaskName | String | Name of the trackable processable task this pipeline is running under, when it was started from one. |
| resultContentType | String | The output response content type set by a step, used for export jobs to tell the caller what kind of file is being produced. |
| rollbackOnly | boolean | Whether an exception has occurred during this run, signalling to transactions that they must roll back rather than commit. |
| running | Boolean | Whether this pipeline is currently executing. |
| sourceAddress | String | The address data is being read from for this run, such as an input file path or sender address, as set by steps or by setSourceAddress. Recorded on the PipelineExecution row. |
| thisOrg | OrgData | Lightweight summary of the organisation this pipeline is running for. Null if the pipeline was created without a website. |
| warnings | List<IntegrationMessage> | The warning messages recorded so far by addWarning, each an issue that should be reviewed but did not stop the run. |
| website | WebsiteRootFolder | The website this integration job is defined within. If none was supplied when the pipeline was created, falls back to the current request's root folder and caches it for subsequent calls. |
Methods
getPipelineExecutionId() · setCurrentPosition(String s) · getThisOrg() · exec(PipelineStep s, Object args) · getOut() · getHead() · getFileHash() · getAttributes() · getTextSourceFile(Path p) · getDataSession() · getPipelineManager() · getEndPointMapping() · start(OutputStream out, Object args) · finished(PipelineStep next) · prepare(PipelineStep next) · org() · getWebsite() · getOrgRoot() · addWarning(String code, String message) · addFailure(String code, String message) · addInfo(String code, String message) · getFailures() · getWarnings() · getInfos() · getDestinationAddress() · getSourceAddress() · getPipelinePath() · isRollbackOnly() · getProcessTaskName() · getCurrentPosition() · getCurrentProfile() · isRunning() · stop() · getResultContentType() · setResultContentType(String resultContentType) · getFileName() · getExecutionRecord() · getCache()
getPipelineExecutionId()
Returns: Long
Database id of the PipelineExecution record for this run, once one has been created. Null before prepareRecordExecution has run, and always null if the endpoint mapping does not record executions.
setCurrentPosition(String s)
Returns: void
Sets the description of where processing is currently up to within the input data, such as the file, worksheet, column and row. Steps call this as they work so that subsequent warning, failure and info messages are attached to a useful position.
| Parameter | Description |
|---|---|
s | the current position description |
getThisOrg()
Returns: OrgData
Lightweight summary of the organisation this pipeline is running for. Null if the pipeline was created without a website.
exec(PipelineStep s, Object args)
Returns: void
Called by a step to hand its output on to the next step in the chain, doing nothing if there is no next step.
| Parameter | Description |
|---|---|
s | the next step to run, or null if this is the last step |
args | the data being passed to the next step, whose meaning is defined by the calling step |
getOut()
Returns: OutputStream
The output stream that pipeline steps write their result to, as set by setOut or by start.
getHead()
Returns: PipelineStep
The first step in this pipeline's chain, which start() prepares and executes.
getFileHash()
Returns: String
Hash of the uploaded or input file this pipeline is processing, used to detect duplicate uploads and recorded on the PipelineExecution row.
getAttributes()
Returns: Map<String,Object>
Mutable map of name/value pairs shared between all steps in this run, used to pass data such as fromDate and toDate along the pipeline and to build expression contexts for MVEL templates.
getTextSourceFile(Path p)
Returns: String
Reads a text file from the website's branch, the same source used to load pipeline and step configuration.
| Parameter | Description |
|---|---|
p | the path of the file within the website's data session |
getDataSession()
Returns: DataSession
The branch containing this pipeline's XML definition and any source code its steps load, such as JsRowStep scripts.
getPipelineManager()
Returns: PipelineManager
The manager that built this pipeline and tracks it while it runs.
getEndPointMapping()
Returns: EndPointMapping
The endpoint configuration this pipeline is executing, which supplies its credentials, file naming, duplicate-prevention and execution-recording settings.
start(OutputStream out, Object args)
Returns: void
Runs the pipeline from its head step to completion, firing STARTED/COMPLETED/FAILED PipelineProcessEvents, recording a PipelineExecution row if the endpoint mapping requests it, and always calling finished() and unregistering the pipeline from the PipelineManager afterwards. Any exception thrown by a step is logged, recorded as a failure, and rethrown wrapped in a RuntimeException.
| Parameter | Description |
|---|---|
out | the stream that steps write their output to |
args | initial arguments passed to the head step's exec method |
finished(PipelineStep next)
Returns: void
Called by a step's own finished method to propagate the finished notification to the next step in the chain, doing nothing if there is no next step.
| Parameter | Description |
|---|---|
next | the next step to notify, or null if this is the last step |
prepare(PipelineStep next)
Returns: void
Called by a step's own prepare method to propagate the prepare call to the next step in the chain, doing nothing if there is no next step.
| Parameter | Description |
|---|---|
next | the next step to prepare, or null if this is the last step |
org()
Returns: Organisation
The organisation this pipeline is running for.
getWebsite()
Returns: WebsiteRootFolder
The website this integration job is defined within. If none was supplied when the pipeline was created, falls back to the current request's root folder and caches it for subsequent calls.
getOrgRoot()
Returns: OrganisationRootFolder
Looks up the OrganisationRootFolder for this pipeline's organisation.
addWarning(String code, String message)
Returns: void
Records a warning against the current position in the pipeline, an event that should be reviewed by an administrator because it might indicate a problem, but does not stop the run.
| Parameter | Description |
|---|---|
code | short code identifying the kind of warning |
message | human-readable description of the warning |
addFailure(String code, String message)
Returns: void
Records a failure against the current position in the pipeline, a definite problem that needs to be resolved. Failures are surfaced on the PipelineExecution record and logged, but do not by themselves stop the run.
| Parameter | Description |
|---|---|
code | short code identifying the kind of failure |
message | human-readable description of the failure |
addInfo(String code, String message)
Returns: void
Records an informational message against the current position in the pipeline, for context that is neither a warning nor a failure.
| Parameter | Description |
|---|---|
code | short code identifying the kind of information |
message | human-readable description of the information |
getFailures()
Returns: List<IntegrationMessage>
The failure messages recorded so far by addFailure, each a definite problem encountered while running this pipeline.
getWarnings()
Returns: List<IntegrationMessage>
The warning messages recorded so far by addWarning, each an issue that should be reviewed but did not stop the run.
getInfos()
Returns: List<IntegrationMessage>
The informational messages recorded so far by addInfo.
getDestinationAddress()
Returns: String
The address data is being sent to for this run, such as an output file path or endpoint URL, as set by steps or by setDestinationAddress. Recorded on the PipelineExecution row.
getSourceAddress()
Returns: String
The address data is being read from for this run, such as an input file path or sender address, as set by steps or by setSourceAddress. Recorded on the PipelineExecution row.
getPipelinePath()
Returns: String
Path, within the website's branch, of the pipeline XML definition this pipeline was built from.
isRollbackOnly()
Returns: boolean
Whether an exception has occurred during this run, signalling to transactions that they must roll back rather than commit.
getProcessTaskName()
Returns: String
Name of the trackable processable task this pipeline is running under, when it was started from one.
getCurrentPosition()
Returns: String
Description of where processing is currently up to within the input data, such as the file, worksheet, column and row, set by steps as they work. Attached to warning, failure and info messages as they are recorded.
getCurrentProfile()
Returns: Profile
The profile this pipeline is running as, used for permissions and for the PipelineProcessEvent fired when the run finishes. Defaults to the current request's principal profile if not explicitly set before start.
isRunning()
Returns: Boolean
Whether this pipeline is currently executing.
stop()
Returns: void
Aborts the pipeline: marks it as needing a rollback and no longer running. Called by steps that decide the run should not continue, such as an export step that has finished producing its output.
getResultContentType()
Returns: String
The output response content type set by a step, used for export jobs to tell the caller what kind of file is being produced.
setResultContentType(String resultContentType)
Returns: void
Sets the output response content type for this run, used for export jobs to tell the caller what kind of file is being produced.
| Parameter | Description |
|---|---|
resultContentType | the content type of the pipeline's output |
getFileName()
Returns: String
Name of the file this pipeline is producing or consuming, evaluated from the endpoint mapping's file name template on first access if it has not already been set explicitly.
getExecutionRecord()
Returns: PipelineExecution
The PipelineExecution row recorded for this run, set once finishRecordExecution has completed.
getCache()
Returns: Map<String,Object>
A map for caching data generated during pipeline processing, such as lookups a step wants to avoid repeating for every row. Lazily created on first access, and may be cleared at any time by clearSession.