Semogram Docs
Data Pipelines

First pipeline

Copy two known database rows through a complete source-to-destination pipeline

This example copies two orders from one PostgreSQL table into another. It gives you a result you can verify without a transform, ontology package or model. It uses the Postgres read and write capabilities and two workspace endpoints in one project.

Arrows show data movement during a real run. The source and destination are workspace resources; the pipeline processing belongs to a project.Test database2 known ordersIngressdemo_ordersEgressdemo_orders_copyOutput table2 matching orders
Saving the graph configures this flow; a run performs the work

What you need

You need a Semogram account with workspace membership, a project in that workspace and permission to create and run pipelines. Reading or writing also requires access to the selected endpoints and external systems. Creating a pipeline does not grant those permissions.

Use a reachable test PostgreSQL database, its SQL client and a database identity authorized to create/read the source fixture and write the dedicated output. Programmatic creation uses pipelines:write; execution uses pipelines:execute; run inspection uses runs:read. Endpoint creation and plugin installation need their own permissions.

Prepare the source

Run this in a test database. Use a fresh table name if the table already exists, and update the endpoint target accordingly.

Create the two-row source
CREATE TABLE public.pipeline_demo_orders (
  order_id integer PRIMARY KEY,
  region text NOT NULL,
  total numeric(12,2) NOT NULL
);
INSERT INTO public.pipeline_demo_orders VALUES
  (1, 'North', 120.00),
  (2, 'South', 75.50);

Open workspace Plugins → Explore, find Postgres (@craven/postgres), inspect its installed contract and install the read capability. Fill these connection fields and run its supported connection check:

Installation fieldValue
HostYour reachable database host
Port5432, or your configured port
DatabaseThe test database containing the fixture
UserA database identity with the required access
PasswordEnter in the protected installation field
SslEnable when the database requires TLS

Open workspace Data Endpoints → New data endpoint. Use name demo_orders, namespace pipelines, direction Ingress, the installed Postgres read capability and Stream (source_stream) as the source contract. In the target fields set Schema to public and Table to pipeline_demo_orders. Save and inspect its schema/sample. demo_orders is a name you choose in Semogram; public.pipeline_demo_orders is the actual database table. Keep its returned endpoint UUID for programmatic examples.

The Postgres installation guide and endpoint guide provide additional detail; the required setup for this fixture is included here.

Prepare the destination

In Plugins → Explore, install Postgres's write capability with the same connection fields above and an identity authorized for the test output. Read and write can use different database users. The writer can perform table/schema setup and enable row-level security; give it the needed setup privileges on this dedicated test target.

Create the output table with your SQL client:

Create the empty output
CREATE TABLE public.pipeline_demo_orders_copy (
  order_id integer PRIMARY KEY,
  region text NOT NULL,
  total numeric(12,2) NOT NULL
);

Create a second workspace endpoint: name demo_orders_copy, namespace pipelines, direction Egress, installed Postgres write capability, contract write, target Schema public, Table pipeline_demo_orders_copy, Write mode append. Save it. It is a destination, not another source endpoint. Keep its UUID for programmatic creation.

Create the graph

Assistant prompt
Create a pipeline in this project named Copy two demo orders. Read the workspace endpoint demo_orders in full mode using parser extraction. Connect that output to an egress node using demo_orders_copy in append mode. Do not add other data sources or a schedule. Show me the two endpoint targets and proposed graph before saving.

Open the project’s Pipeline Studio → New pipeline. Use the authoring assistant to propose the graph, then inspect/edit the two nodes in Studio:

NodeSelected fields
Read demo ordersIngress; Data endpoint demo_orders; Load mode Full load; extraction strategy parser; output orders
Write demo ordersEgress; Data endpoint demo_orders_copy; Write mode Append; source orders

Connect the ingress output to the egress input. The output/input port keys in the serialized graph are out and in. Review the selected endpoints' actual database targets; their friendly names alone are insufficient.

Use a workspace API key with pipelines:write. The complete request body contains the same two-node graph; save it as pipeline-create.json. Replace the four example UUIDs: metadata workspaceId, metadata projectId, ingress dataEndpointId and egress dataEndpointId with the workspace, project and endpoint IDs created above. Both metadata IDs must match the authenticated workspace and URL project.

Create a pipeline
curl --request POST "https://platform.semogram.com/api/v1/projects/<PROJECT_ID_UUID>/pipelines" \
  --header "Authorization: Bearer ${SEMOGRAM_API_KEY}" \
  --header "Idempotency-Key: demo-order-copy-create-001" \
  --header "Content-Type: application/json" \
  --data-binary @pipeline-create.json

The response includes pipeline.id, versionId, document and draft. Keep the pipeline ID. Creating this resource records an initial version and draft; it does not start a run.

The UUIDs below are illustrative. Replace workspace/project metadata and the two endpoint IDs with the resources created above. Use the matching projectId argument through your initialized workspace MCP connection.

Send through an authenticated, initialized MCP client. projectId selects the project within the connected workspace.

MCP request
{
  "jsonrpc": "2.0",
  "id": 1,
  "method": "tools/call",
  "params": {
    "name": "pipeline_create",
    "arguments": {
      "projectId": "22222222-2222-4222-8222-222222222222",
      "document": {
        "version": 1,
        "kind": "studio-plan",
        "metadata": {
          "name": "Copy two demo orders",
          "description": "Copy the two-row test fixture into a dedicated output table",
          "projectId": "22222222-2222-4222-8222-222222222222",
          "workspaceId": "11111111-1111-4111-8111-111111111111"
        },
        "graph": {
          "id": "demo-order-copy",
          "nodes": [
            {
              "id": "read-orders",
              "kind": "ingress",
              "label": "Read demo orders",
              "inputs": [],
              "outputs": [
                {
                  "key": "out",
                  "label": "Orders",
                  "direction": "output",
                  "interface": {
                    "shape": "collection"
                  }
                }
              ],
              "policy": {},
              "ingress": {
                "dataEndpointId": "33333333-3333-4333-8333-333333333333",
                "mode": "full",
                "output": "orders",
                "strategy": "parser"
              }
            },
            {
              "id": "write-orders",
              "kind": "egress",
              "label": "Write demo orders",
              "inputs": [
                {
                  "key": "in",
                  "label": "Orders",
                  "direction": "input",
                  "interface": {
                    "shape": "collection"
                  }
                }
              ],
              "outputs": [],
              "policy": {},
              "egress": {
                "dataEndpointId": "44444444-4444-4444-8444-444444444444",
                "source": "orders",
                "writeMode": "append"
              }
            }
          ],
          "edges": [
            {
              "id": "orders-to-output",
              "fromNodeId": "read-orders",
              "fromOutput": "out",
              "toNodeId": "write-orders",
              "toInput": "in"
            }
          ]
        }
      },
      "commitMessage": "Initial two-row fixture",
      "idempotencyKey": "demo-order-copy-create-001"
    }
  }
}

Validate and save

In Studio, Validate the draft. Resolve errors, inspect warnings and Save version after edits. Confirm the selected version contains these two endpoints and append mode. Preflight does not copy either row or test the external write.

API clients can validate the actor's draft with POST /api/v1/projects/<PROJECT_ID>/pipelines/<PIPELINE_ID>/draft/validate (pipelines:read) and commit the full document with POST .../versions (pipelines:write). MCP uses pipeline_validate with pipelineId and projectId; saving a draft through pipeline_draft_save does not commit a version. The initial version returned by creation can be inspected before launch.

Run once

Select Run in Studio after saving and validation. Open the returned run and inspect both steps.

Assistant prompt
Validate Copy two demo orders and show the active saved version and both targets. After I approve the write, run that saved pipeline once. Report the run ID and inspect its final step results; do not schedule it.

Use a workspace API key with pipelines:execute, access to this project and an accountable current member where the operation requires one. Set SEMOGRAM_API_KEY; replace UUID placeholders with actual IDs.

HTTP API request
curl --request POST "https://platform.semogram.com/api/v1/projects/<PROJECT_ID_UUID>/pipelines/<PIPELINE_ID_UUID>/execute" \
  --header "Authorization: Bearer ${SEMOGRAM_API_KEY}" \
  --header "Idempotency-Key: pipeline-demo-operation-001"

Send through an authenticated, initialized MCP client. projectId selects the project within the connected workspace.

MCP request
{
  "jsonrpc": "2.0",
  "id": 1,
  "method": "tools/call",
  "params": {
    "name": "pipeline_execute",
    "arguments": {
      "pipelineId": "<PIPELINE_ID_UUID>",
      "projectId": "<PROJECT_ID_UUID>",
      "idempotencyKey": "demo-order-copy-run-001"
    }
  }
}

Keep runId from the accepted response. Poll GET /api/v1/projects/<PROJECT_ID>/runs/<RUN_ID> with runs:read, or call MCP job_get with jobId set to that run ID and projectId. Disconnecting your client does not make a queued run disappear.

Check the result

Use your database SQL client:

Verify the external target
SELECT order_id, region, total
FROM public.pipeline_demo_orders_copy
ORDER BY order_id;
order_idregiontotal
1North120.00
2South75.50

Inspect the source too: it should still contain the two original rows. Check that the run's ingress and egress steps succeeded and inspect their errors/counts instead of accepting only an assistant summary.

Append does not turn a repeat into an upsert. With these primary keys, rerunning can fail on duplicate keys; a failed write may have committed some rows. Inspect the output before retrying. To repeat this fixture, empty only the dedicated test output using your database tools and run once again. Do not clear a shared business table.

If something fails

SymptomCheck
Endpoint missingBoth endpoints are in this workspace and the graph contains their actual IDs
Validation errorPort connection, read/write capability and saved graph configuration
Read failsSource table, grants, TLS and runtime network reachability
Write failsDestination privileges, row-level security, target contract and duplicate primary keys
Accepted but no resultInspect the durable run ID, dispatch state and step errors

Once this fixture works, try Transform records. Do not enable a recurring schedule until the expected behavior and replay policy are established.