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

PropertyReturnsDescription
attributesMap<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.
cacheMap<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.
currentPositionStringDescription 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.
currentProfileProfileThe 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.
dataSessionDataSessionThe branch containing this pipeline's XML definition and any source code its steps load, such as JsRowStep scripts.
destinationAddressStringThe 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.
endPointMappingEndPointMappingThe endpoint configuration this pipeline is executing, which supplies its credentials, file naming, duplicate-prevention and execution-recording settings.
executionRecordPipelineExecutionThe PipelineExecution row recorded for this run, set once finishRecordExecution has completed.
expressionAttributesMap<String,Object>
failuresList<IntegrationMessage>The failure messages recorded so far by addFailure, each a definite problem encountered while running this pipeline.
fileHashStringHash of the uploaded or input file this pipeline is processing, used to detect duplicate uploads and recorded on the PipelineExecution row.
fileNameStringName 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.
headPipelineStepThe first step in this pipeline's chain, which start() prepares and executes.
infosList<IntegrationMessage>The informational messages recorded so far by addInfo.
orgRootOrganisationRootFolderLooks up the OrganisationRootFolder for this pipeline's organisation.
outOutputStreamThe output stream that pipeline steps write their result to, as set by setOut or by start.
pipelineExecutionIdLongDatabase 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.
pipelineManagerPipelineManagerThe manager that built this pipeline and tracks it while it runs.
pipelinePathStringPath, within the website's branch, of the pipeline XML definition this pipeline was built from.
processTaskNameStringName of the trackable processable task this pipeline is running under, when it was started from one.
resultContentTypeStringThe output response content type set by a step, used for export jobs to tell the caller what kind of file is being produced.
rollbackOnlybooleanWhether an exception has occurred during this run, signalling to transactions that they must roll back rather than commit.
runningBooleanWhether this pipeline is currently executing.
sourceAddressStringThe 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.
thisOrgOrgDataLightweight summary of the organisation this pipeline is running for. Null if the pipeline was created without a website.
warningsList<IntegrationMessage>The warning messages recorded so far by addWarning, each an issue that should be reviewed but did not stop the run.
websiteWebsiteRootFolderThe 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.

ParameterDescription
sthe 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.

ParameterDescription
sthe next step to run, or null if this is the last step
argsthe 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.

ParameterDescription
pthe 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.

ParameterDescription
outthe stream that steps write their output to
argsinitial 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.

ParameterDescription
nextthe 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.

ParameterDescription
nextthe 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.

ParameterDescription
codeshort code identifying the kind of warning
messagehuman-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.

ParameterDescription
codeshort code identifying the kind of failure
messagehuman-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.

ParameterDescription
codeshort code identifying the kind of information
messagehuman-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.

ParameterDescription
resultContentTypethe 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.

To get full access to the Kademi Hub existing customers can login here, or new customers can register here.