openapi: 3.0.3 info: title: Feldera Input Connectors Output Connectors API description: "\nWith Feldera, users create data pipelines out of SQL programs.\nA SQL program comprises tables and views, and includes as well the definition of\ninput and output connectors for each respectively. A connector defines a data\nsource or data sink to feed input data into tables or receive output data\ncomputed by the views respectively.\n\n## Pipeline\n\nThe API is centered around the **pipeline**, which most importantly consists\nout of the SQL program, but also has accompanying metadata and configuration parameters\n(e.g., compilation profile, number of workers, etc.).\n\n* A pipeline is identified and referred to by its user-provided unique name.\n* The pipeline program is asynchronously compiled when the pipeline is first created or\n when its program is subsequently updated.\n* Pipeline deployment start is only able to proceed to provisioning once the program is successfully\n compiled.\n* A pipeline cannot be updated while it is deployed.\n\n## Concurrency\n\nEach pipeline has a version, which is incremented each time its core fields are updated.\nThe version is monotonically increasing. There is additionally a program version which covers\nonly the program-related core fields, and is used by the compiler to discern when to recompile.\n\n## Client request handling\n\n### Request outcome expectations\n\nThe outcome of a request is that it either fails (e.g., DNS lookup failed) without any response\n(no status code nor body), or it succeeds and gets back a response status code and body.\n\nIn case of a response, usually it is the Feldera endpoint that generated it:\n- If it is success (2xx), it will return whichever body belongs to the success response.\n- Otherwise, if it is an error (4xx, 5xx), it will return a Feldera error response JSON body\n which will have an application-level `error_code`.\n\nHowever, there are two notable exceptions when the response is not generated by the Feldera\nendpoint:\n- If the HTTP server, to which the endpoint belongs, encountered an issue, it might return\n 4xx (e.g., for an unknown endpoint) or 5xx error codes by itself (e.g., when it is initializing).\n- If the Feldera API server is behind a (reverse) proxy, the proxy can return error codes by itself,\n for example BAD GATEWAY (502) or GATEWAY TIMEOUT (504).\n\nAs such, it is not guaranteed that the (4xx, 5xx) will have a Feldera error response JSON body\nin these latter cases.\n\n### Error handling and retrying\n\nThe error type returned by the client should distinguish between the error responses generated\nby Feldera endpoints themselves (which have a Feldera error response body) and those that are\ngenerated by other sources.\n\nIn order for a client operation (e.g., `pipeline.resume()`) to be robust (i.e., not fail due to\na single HTTP request not succeeding) the client should use a retry mechanism if the operation\nis idempotent. The retry mechanism must however have a time limit, after which it times out.\nThis guarantees that the client operation is eventually responsive, which enables the script\nit is a part of to not hang indefinitely on Feldera operations and instead be able to decide\nby itself whether and how to proceed. If no response is returned, the mechanism should generally\nretry. When a response is returned, the decision whether to retry can generally depend on the status\ncode: especially the status codes 408, 502, 503 and 504 should be considered as transient errors.\nFiner grained retry decisions should be made by taking into account the application-level\n`error_code` if the response body was indeed a Feldera error response body.\n\n## Feldera client errors (4xx)\n\n_Client behavior:_ clients should generally return with an error when they get back a 4xx status\ncode, as it usually means the request will likely not succeed even if it is sent again. Certain\nrequests might make use of a timed retry mechanism when the client error is transient without\nrequiring any user intervention to overcome, for instance a transaction already being in progress\nleading to a temporary CONFLICT (409) error.\n\n- **BAD REQUEST (400)**: invalid user request (general).\n - _Example:_ the new pipeline name `example1@~` contains invalid characters.\n\n- **UNAUTHORIZED (401)**: the user is not authorized to issue the request.\n - _Example:_ an invalid API key is provided.\n\n- **NOT FOUND (404)**: a resource required to exist in order to process the request was not found.\n - _Example:_ a pipeline named `example` does not exist when trying to update it.\n\n- **CONFLICT (409)**: there is a conflict between the request and a relevant resource.\n - _Example:_ a pipeline named `example` already exists.\n - _Example:_ another transaction is already in process.\n\n## Feldera server errors (5xx)\n\n- **INTERNAL SERVER ERROR (500)**: the server is unexpectedly unable to process the request\n (general).\n - _Example:_ unable to reach the database.\n - _Client behavior:_ immediately return with an error.\n\n- **NOT IMPLEMENTED (501)**: the server does not implement functionality required to process the\n request.\n - _Example:_ making a request to an enterprise-only endpoint in the OSS edition.\n - _Client behavior:_ immediately return with an error.\n\n- **SERVICE UNAVAILABLE (503)**: the server is not (yet) able to process the request.\n - _Example:_ pausing a pipeline which is still provisioning.\n - _Client behavior:_ depending on the type of request, client may use a timed retry mechanism.\n\n## Feldera error response body\n\nWhen the Feldera API returns an HTTP error status code (4xx, 5xx), the body will contain the\nfollowing JSON object:\n\n```json\n{\n \"message\": \"Human-readable explanation.\",\n \"error_code\": \"CodeSpecifyingError\",\n \"details\": {\n\n }\n}\n```\n\nIt contains the following fields:\n- **message (string)**: human-readable explanation of the error that occurred and potentially\n hinting what can be done about it.\n- **error_code (string)**: application-level code about the error that occurred, written in CamelCase.\n For example: `UnknownPipelineName`, `DuplicateName`, `PauseWhileNotProvisioned`, ... .\n- **details (object)**: JSON object corresponding to the `error_code` with fields that provide\n details relevant to it. For example: if a name is unknown, a field with the unknown name in\n question.\n" contact: name: Feldera Team email: dev@feldera.com license: name: MIT OR Apache-2.0 version: 0.323.0 tags: - name: Output Connectors paths: /v0/pipelines/{pipeline_name}/egress/{table_name}: post: tags: - Output Connectors summary: Subscribe to View description: 'Subscribe to a stream of updates from a SQL view or table. The pipeline responds with a continuous stream of changes to the specified table or view. The stream is configurable two ways: - Simple configuration of the format may be provided using query parameters. Specify `backpressure` to specify behavior when the HTTP client cannot keep up. Use `format` to specify `csv` or `json` output. For `json` output format, `update_format` and `json_flavor` may be provided (with the same possible values as in JSON format configuration for connectors). - Comprehensive configuration may be provided by providing a connector configuration as a JSON body. In this case, no query parameters are allowed. Updates are split into `Chunk`s. The pipeline continues sending updates until the client closes the connection or the pipeline is stopped.' operationId: http_output parameters: - name: pipeline_name in: path description: Unique pipeline name required: true schema: type: string - name: table_name in: path description: SQL table name. Unquoted SQL names have to be capitalized. Quoted SQL names have to exactly match the case from the SQL program. required: true schema: type: string - name: format in: query description: Output data format, either 'csv' or 'json'. required: true schema: type: string - name: send_snapshot in: query description: 'Set to `true` to send a full snapshot of a materialized view before streaming incremental updates. The default is `false`. Works on a paused pipeline: the snapshot is delivered from the latest cached view state without requiring the pipeline to be running.' required: false schema: type: boolean nullable: true - name: array in: query description: Set to `true` to group updates in this stream into JSON arrays (used in conjunction with `format=json`). The default value is `false` required: false schema: type: boolean nullable: true - name: backpressure in: query description: "Apply backpressure on the pipeline when the HTTP client cannot receive data fast enough.\n When this flag is set to false (the default), the HTTP connector drops data chunks if the client is not keeping up with its output. This prevents a slow HTTP client from slowing down the entire pipeline.\n When the flag is set to true, the connector waits for the client to receive each chunk and blocks the pipeline if the client cannot keep up." required: false schema: type: boolean nullable: true responses: '200': description: Connection to the endpoint successfully established. The body of the response contains a stream of data chunks. content: application/json: schema: $ref: '#/components/schemas/Chunk' '400': description: '' content: application/json: schema: $ref: '#/components/schemas/ErrorResponse' '404': description: Pipeline and/or table/view with that name does not exist content: application/json: schema: $ref: '#/components/schemas/ErrorResponse' examples: Pipeline with that name does not exist: value: message: Unknown pipeline name 'non-existent-pipeline' error_code: UnknownPipelineName details: pipeline_name: non-existent-pipeline '500': description: '' content: application/json: schema: $ref: '#/components/schemas/ErrorResponse' '503': description: '' content: application/json: schema: $ref: '#/components/schemas/ErrorResponse' examples: Disconnected during response: value: message: 'Error sending HTTP request to pipeline: the pipeline disconnected while it was processing this HTTP request. This could be because the pipeline either (a) encountered a fatal error or panic, (b) was stopped, or (c) experienced network issues -- retrying might help in the last case. Alternatively, check the pipeline logs. Failed request: /pause pipeline-id=N/A pipeline-name="my_pipeline"' error_code: PipelineInteractionUnreachable details: pipeline_name: my_pipeline request: /pause error: the pipeline disconnected while it was processing this HTTP request. This could be because the pipeline either (a) encountered a fatal error or panic, (b) was stopped, or (c) experienced network issues -- retrying might help in the last case. Alternatively, check the pipeline logs. Pipeline is currently unavailable: value: message: 'Error sending HTTP request to pipeline: deployment status is currently ''unavailable'' -- wait for it to become ''running'' or ''paused'' again Failed request: /pause pipeline-id=N/A pipeline-name="my_pipeline"' error_code: PipelineInteractionUnreachable details: pipeline_name: my_pipeline request: /pause error: deployment status is currently 'unavailable' -- wait for it to become 'running' or 'paused' again Pipeline is not deployed: value: message: Unable to interact with pipeline because the deployment status (stopped) indicates it is not (yet) fully provisioned pipeline-id=N/A pipeline-name="my_pipeline" error_code: PipelineInteractionNotDeployed details: pipeline_name: my_pipeline status: Stopped desired_status: Provisioned Response timeout: value: message: 'Error sending HTTP request to pipeline: timeout (10s) was reached: this means the pipeline took too long to respond -- this can simply be because the request was too difficult to process in time, or other reasons (e.g., deadlock): the pipeline logs might contain additional information (original send request error: Timeout while waiting for response) Failed request: /pause pipeline-id=N/A pipeline-name="my_pipeline"' error_code: PipelineInteractionUnreachable details: pipeline_name: my_pipeline request: /pause error: 'timeout (10s) was reached: this means the pipeline took too long to respond -- this can simply be because the request was too difficult to process in time, or other reasons (e.g., deadlock): the pipeline logs might contain additional information (original send request error: Timeout while waiting for response)' security: - JSON web token (JWT) or API key: [] /v0/pipelines/{pipeline_name}/views/{view_name}/connectors/{connector_name}/stats: get: tags: - Output Connectors summary: Get Output Status description: Retrieve the status of an output connector. operationId: get_pipeline_output_connector_status parameters: - name: pipeline_name in: path description: Unique pipeline name required: true schema: type: string - name: view_name in: path description: SQL view name required: true schema: type: string - name: connector_name in: path description: Output connector name required: true schema: type: string responses: '200': description: Output connector status retrieved successfully content: application/json: schema: $ref: '#/components/schemas/OutputEndpointStatus' '404': description: Pipeline, view and/or output connector with that name does not exist content: application/json: schema: $ref: '#/components/schemas/ErrorResponse' examples: Pipeline with that name does not exist: value: message: Unknown pipeline name 'non-existent-pipeline' error_code: UnknownPipelineName details: pipeline_name: non-existent-pipeline '500': description: '' content: application/json: schema: $ref: '#/components/schemas/ErrorResponse' '503': description: '' content: application/json: schema: $ref: '#/components/schemas/ErrorResponse' examples: Disconnected during response: value: message: 'Error sending HTTP request to pipeline: the pipeline disconnected while it was processing this HTTP request. This could be because the pipeline either (a) encountered a fatal error or panic, (b) was stopped, or (c) experienced network issues -- retrying might help in the last case. Alternatively, check the pipeline logs. Failed request: /pause pipeline-id=N/A pipeline-name="my_pipeline"' error_code: PipelineInteractionUnreachable details: pipeline_name: my_pipeline request: /pause error: the pipeline disconnected while it was processing this HTTP request. This could be because the pipeline either (a) encountered a fatal error or panic, (b) was stopped, or (c) experienced network issues -- retrying might help in the last case. Alternatively, check the pipeline logs. Pipeline is currently unavailable: value: message: 'Error sending HTTP request to pipeline: deployment status is currently ''unavailable'' -- wait for it to become ''running'' or ''paused'' again Failed request: /pause pipeline-id=N/A pipeline-name="my_pipeline"' error_code: PipelineInteractionUnreachable details: pipeline_name: my_pipeline request: /pause error: deployment status is currently 'unavailable' -- wait for it to become 'running' or 'paused' again Pipeline is not deployed: value: message: Unable to interact with pipeline because the deployment status (stopped) indicates it is not (yet) fully provisioned pipeline-id=N/A pipeline-name="my_pipeline" error_code: PipelineInteractionNotDeployed details: pipeline_name: my_pipeline status: Stopped desired_status: Provisioned Response timeout: value: message: 'Error sending HTTP request to pipeline: timeout (10s) was reached: this means the pipeline took too long to respond -- this can simply be because the request was too difficult to process in time, or other reasons (e.g., deadlock): the pipeline logs might contain additional information (original send request error: Timeout while waiting for response) Failed request: /pause pipeline-id=N/A pipeline-name="my_pipeline"' error_code: PipelineInteractionUnreachable details: pipeline_name: my_pipeline request: /pause error: 'timeout (10s) was reached: this means the pipeline took too long to respond -- this can simply be because the request was too difficult to process in time, or other reasons (e.g., deadlock): the pipeline logs might contain additional information (original send request error: Timeout while waiting for response)' security: - JSON web token (JWT) or API key: [] components: schemas: ShortEndpointConfig: type: object description: Schema definition for endpoint config that only includes the stream field. required: - stream properties: stream: type: string description: The name of the stream. OutputEndpointMetrics: type: object description: Performance metrics for an output endpoint. required: - transmitted_records - transmitted_bytes - queued_records - queued_batches - buffered_records - buffered_batches - num_encode_errors - num_transport_errors - total_processed_input_records - total_processed_steps - memory properties: batch_records_written: type: integer format: int64 description: 'Number of records written so far while the connector is processing a batch of updates. Resets to 0 after the batch is committed. `None` when the connector does not support batch-progress reporting.' nullable: true minimum: 0 buffered_batches: type: integer format: int64 description: Number of batches in the buffer. minimum: 0 buffered_records: type: integer format: int64 description: Number of records pushed to the output buffer. minimum: 0 memory: type: integer format: int64 description: Extra memory in use beyond that used for queuing records. minimum: 0 num_encode_errors: type: integer format: int64 description: Number of encoding errors. minimum: 0 num_transport_errors: type: integer format: int64 description: Number of transport errors. minimum: 0 queued_batches: type: integer format: int64 description: Number of queued batches. minimum: 0 queued_records: type: integer format: int64 description: Number of queued records. minimum: 0 total_processed_input_records: type: integer format: int64 description: 'The number of input records processed by the circuit. This metric tracks the end-to-end progress of the pipeline: the output of this endpoint is equal to the output of the circuit after processing `total_processed_input_records` records. In a multihost pipeline, this count reflects only the input records processed on the same host as the output endpoint, which is not usually meaningful.' minimum: 0 total_processed_steps: type: integer format: int64 description: 'The number of steps whose input records have been processed by the endpoint. This is meaningful in a multihost pipeline because steps are synchronized across all of the hosts. # Interpretation This is a count, not a step number. If `total_processed_steps` is 0, no steps have been processed to completion. If `total_processed_steps > 0`, then the last step whose input records have been processed to completion is `total_processed_steps - 1`. A record that was ingested in step `n` is fully processed when `total_processed_steps > n`.' minimum: 0 transmitted_bytes: type: integer format: int64 description: Bytes sent on the underlying transport. minimum: 0 transmitted_records: type: integer format: int64 description: Records sent on the underlying transport. minimum: 0 ConnectorHealthStatus: type: string enum: - Healthy - Unhealthy ErrorResponse: type: object description: Information returned by REST API endpoints on error. required: - message - error_code - details properties: details: description: 'Detailed error metadata. The contents of this field is determined by `error_code`.' error_code: type: string description: Error code is a string that specifies this error type. example: CodeSpecifyingErrorType message: type: string description: Human-readable error message. example: Explanation of the error that occurred. ConnectorHealth: type: object required: - status properties: description: type: string nullable: true status: $ref: '#/components/schemas/ConnectorHealthStatus' OutputEndpointStatus: type: object description: Output endpoint status information. required: - endpoint_name - config - metrics properties: config: $ref: '#/components/schemas/ShortEndpointConfig' encode_errors: type: array items: $ref: '#/components/schemas/ConnectorError' description: Recent encoding errors on this endpoint. nullable: true endpoint_name: type: string description: Endpoint name. fatal_error: type: string description: The first fatal error that occurred at the endpoint. nullable: true health: allOf: - $ref: '#/components/schemas/ConnectorHealth' nullable: true metrics: $ref: '#/components/schemas/OutputEndpointMetrics' transport_errors: type: array items: $ref: '#/components/schemas/ConnectorError' description: Recent transport errors on this endpoint. nullable: true ConnectorError: type: object required: - timestamp - index - message properties: index: type: integer format: int64 description: 'Sequence number of the error. The client can use this field to detect gaps in the error list reported by the pipeline. When the connector reports a large number of errors, the pipeline will only preserve and report the most recent errors of each kind.' minimum: 0 message: type: string description: Error message. tag: type: string description: 'Optional tag for the error. The tag is used to group errors by their type.' nullable: true timestamp: type: string format: date-time description: Timestamp when the error occurred, serialized as RFC3339 with microseconds. Chunk: type: object description: 'A set of updates to a SQL table or view. The `sequence_number` field stores the offset of the chunk relative to the start of the stream and can be used to implement reliable delivery. The payload is stored in the `bin_data`, `text_data`, or `json_data` field depending on the data format used.' required: - sequence_number - snapshot properties: bin_data: type: string format: binary description: Base64 encoded binary payload, e.g., bincode. nullable: true json_data: type: object description: JSON payload. nullable: true sequence_number: type: integer format: int64 minimum: 0 snapshot: type: boolean description: '`true` when this chunk is part of the initial snapshot delivered in `send_snapshot` mode; `false` for incremental delta updates.' text_data: type: string description: Text payload, e.g., CSV. nullable: true securitySchemes: JSON web token (JWT) or API key: type: http scheme: bearer bearerFormat: JWT description: "Use a JWT token obtained via an OAuth2/OIDC\n login workflow or an API key obtained via\n the `/v0/api-keys` endpoint."