> ## Documentation Index
> Fetch the complete documentation index at: https://docs.synq.io/llms.txt
> Use this file to discover all available pages before exploring further.

# Replication lineage

### replication\_lineage

Find full example [here](https://github.com/getsynq/api/tree/main/examples/golang/replication_lineage)

#### main.go

```go theme={null}
// Lineage for a replication tool: a table copied from one database into a
// table in a different warehouse, with column lineage across the hop.
//
//	Postgres table                   (already in Coalesce Quality)
//	  |  SQL: INSERT INTO <target> (...) SELECT ... FROM <source>
//	  v
//	replication stream               custom entity, SqlDefinition + references
//	  |  the same SQL's write target
//	  v
//	Snowflake table                  (already in Coalesce Quality)
//
// Two tables cannot be joined by a relationship directly: a relationship needs a
// custom entity on at least one side. That is not a limitation to work around,
// it is the model — the rows do not move by themselves, something copies them,
// and that something is what fails, lags or drops a column. So it gets an
// entity, one per replicated table (a "stream"), and the two tables hang off it.
//
// One entity per stream, never one for the whole connection: with a single
// entity in the middle, every source table would read as upstream of every
// target table the connection writes.
//
// The hop is described as SQL rather than as a list of edges, because SQL says
// which column lands in which, including renames, and the parser classifies each
// one (passthrough, renamed, transformed). The two tables live in different
// warehouses, so neither name in the SQL can be looked up by warehouse address
// from the other side. Both are BOUND instead: a reference pins the name as the
// SQL writes it to the entity it stands for, whatever platform that is on.
//
// Everything written here is idempotent: re-running the program produces the
// same estate, and a stream that disappears from the list is deleted.
package main

import (
	"context"
	"crypto/tls"
	"fmt"
	"os"
	"strings"
	"time"

	entitiescustomv1grpc "buf.build/gen/go/getsynq/api/grpc/go/synq/entities/custom/v1/customv1grpc"
	lineagev1grpc "buf.build/gen/go/getsynq/api/grpc/go/synq/entities/lineage/v1/lineagev1grpc"
	customfeaturesv1 "buf.build/gen/go/getsynq/api/protocolbuffers/go/synq/entities/custom/features/v1"
	customv1 "buf.build/gen/go/getsynq/api/protocolbuffers/go/synq/entities/custom/v1"
	lineagev1 "buf.build/gen/go/getsynq/api/protocolbuffers/go/synq/entities/lineage/v1"
	entitiesv1 "buf.build/gen/go/getsynq/api/protocolbuffers/go/synq/entities/v1"
	"golang.org/x/oauth2/clientcredentials"
	"google.golang.org/grpc"
	"google.golang.org/grpc/credentials"
	"google.golang.org/grpc/credentials/oauth"
	"google.golang.org/protobuf/types/known/timestamppb"
)

// Custom type ids are workspace-wide and yours to allocate: 1..1000. Pick a free
// one and keep it — an entity's type is its id, so reusing a number silently
// reclassifies everything that had it.
const typeStream = 60

// groupID makes a re-sync self-cleaning: a stream removed from the replication
// tool simply stops being sent, and the server deletes its entity.
const groupID = "replication"

// stream is one replicated table, as the replication tool describes it.
type stream struct {
	name string

	source *entitiesv1.Identifier
	target *entitiesv1.Identifier

	// Source column -> target column. Most tools copy names through unchanged;
	// list the columns anyway, because the list is what says which of the
	// target's columns came from the source and which the tool added itself.
	columns [][2]string
}

// A real integration reads these from the replication tool's own catalog.
var streams = []stream{
	{
		name:   "app_db.public.customers",
		source: postgresTable("app-db.internal.example.com", "app_db", "public", "customers"),
		target: snowflakeTable("xy12345.eu-west-1", "RAW", "APP_DB", "CUSTOMERS"),
		columns: [][2]string{
			{"id", "id"},
			{"email", "email"},
			{"created_at", "created_at"},
		},
	},
	{
		name:   "app_db.public.orders",
		source: postgresTable("app-db.internal.example.com", "app_db", "public", "orders"),
		target: snowflakeTable("xy12345.eu-west-1", "RAW", "APP_DB", "ORDERS"),
		columns: [][2]string{
			{"id", "id"},
			{"customer_id", "customer_id"},
			{"total", "total_amount"},
			{"created_at", "created_at"},
		},
	},
}

func main() {
	ctx := context.Background()

	// developer.synq.io (EU) is the default. The other deployments are
	// api.us.synq.io (US) and api.au.synq.io (AU).
	endpoint := env("QUALITY_API_ENDPOINT", "developer.synq.io")
	clientID := os.Getenv("QUALITY_CLIENT_ID")
	clientSecret := os.Getenv("QUALITY_CLIENT_SECRET")
	if clientID == "" || clientSecret == "" {
		fmt.Println("set QUALITY_CLIENT_ID and QUALITY_CLIENT_SECRET (see README.md)")
		os.Exit(1)
	}

	oauthConfig := &clientcredentials.Config{
		ClientID:     clientID,
		ClientSecret: clientSecret,
		TokenURL:     fmt.Sprintf("https://%s/oauth2/token", endpoint),
	}
	conn, err := grpc.NewClient(
		fmt.Sprintf("%s:443", endpoint),
		grpc.WithTransportCredentials(credentials.NewTLS(&tls.Config{InsecureSkipVerify: false})),
		grpc.WithPerRPCCredentials(oauth.TokenSource{TokenSource: oauthConfig.TokenSource(ctx)}),
		grpc.WithAuthority(endpoint),
	)
	if err != nil {
		panic(err)
	}
	defer conn.Close()

	types := entitiescustomv1grpc.NewTypesServiceClient(conn)
	entities := entitiescustomv1grpc.NewEntitiesServiceClient(conn)
	features := entitiescustomv1grpc.NewFeaturesServiceClient(conn)
	groups := entitiescustomv1grpc.NewGroupsServiceClient(conn)
	lineage := lineagev1grpc.NewLineageServiceClient(conn)

	fmt.Println("== declare the stream type")
	_, err = types.UpsertType(ctx, &customv1.UpsertTypeRequest{Type: &entitiesv1.Type{
		TypeId:  typeStream,
		Name:    "Replication stream",
		SvgIcon: iconArrows,
	}})
	must(err)

	fmt.Println("== declare one entity per stream")
	for _, s := range streams {
		_, err := entities.UpsertEntity(ctx, &customv1.UpsertEntityRequest{Entity: &entitiesv1.Entity{
			Id:          streamID(s),
			TypeId:      typeStream,
			Name:        s.name,
			Description: "Replicated by the replication tool.",
		}})
		must(err)
		fmt.Printf("   %s\n", s.name)
	}

	fmt.Println("== describe each hop as SQL")
	for _, s := range streams {
		_, err := features.UpsertEntityFeature(ctx, &customv1.UpsertEntityFeatureRequest{
			Feature: &customv1.Feature{
				EntityId:  streamID(s),
				FeatureId: "sql",
				Feature:   &customv1.Feature_SqlDefinition{SqlDefinition: streamSql(s)},
			},
		})
		must(err)
		fmt.Printf("   %s\n", streamSql(s).GetSql())
	}

	fmt.Println("== reconcile the group")
	ids := make([]*entitiesv1.Identifier, 0, len(streams))
	for _, s := range streams {
		ids = append(ids, streamID(s))
	}
	resp, err := groups.UpsertEntitiesGroup(ctx, &customv1.UpsertEntitiesGroupRequest{
		Group: &customv1.Group{GroupId: groupID, EntityIds: ids},
	})
	must(err)
	for _, deleted := range resp.GetDeletedIds() {
		fmt.Printf("   deleted (no longer replicated): %s\n", deleted.GetCustom().GetId())
	}

	fmt.Println("== verify")
	verify(ctx, lineage, streams[1], 2*time.Minute)
}

// streamSql is the hop, written as the SQL a replication would run if source and
// target were in one database.
//
// The dialect is the SOURCE side's, and nothing in the SQL is ever executed: it
// is parsed for lineage only. The two names it uses are placeholders — they only
// have to match the references' object_name — which is what lets the target sit
// in a different warehouse from the source.
func streamSql(s stream) *customfeaturesv1.SqlDefinition {
	targetColumns := make([]string, 0, len(s.columns))
	selectList := make([]string, 0, len(s.columns))
	for _, c := range s.columns {
		selectList = append(selectList, c[0])
		targetColumns = append(targetColumns, c[1])
	}

	return &customfeaturesv1.SqlDefinition{
		StateAt: timestamppb.Now(),
		Dialect: entitiesv1.SqlDialect_SQL_DIALECT_POSTGRESQL,
		// An explicit column list on the INSERT is what pairs each selected
		// column with its target column. Without one the pairing is positional
		// and needs the target's shape to resolve.
		Sql: fmt.Sprintf("INSERT INTO target_table (%s) SELECT %s FROM source_table",
			strings.Join(targetColumns, ", "), strings.Join(selectList, ", ")),
		// Bound, not looked up: the source and the target are on different
		// platforms, and a reference resolves to its entity on whichever one it
		// is. A binding to a table Coalesce Quality has not ingested yet is kept,
		// and the edge appears once the table does.
		References: []*customfeaturesv1.SqlTableReference{
			{ObjectName: "source_table", Entity: s.source},
			{ObjectName: "target_table", Entity: s.target},
		},
	}
}

// verify reads the column lineage back from the target side.
//
// The way this goes wrong is silent: a reference to a table that is not in
// Coalesce Quality under that identifier leaves a valid write and an empty
// graph. Lineage is computed asynchronously from what was written, so this polls.
func verify(ctx context.Context, api lineagev1grpc.LineageServiceClient, s stream, budget time.Duration) {
	columns := make([]string, 0, len(s.columns))
	for _, c := range s.columns {
		columns = append(columns, c[1])
	}

	deadline := time.Now().Add(budget)
	for {
		resp, err := api.GetLineage(ctx, &lineagev1.GetLineageRequest{
			LineageDirection: lineagev1.LineageDirection_LINEAGE_DIRECTION_UPSTREAM,
			StartPoint: &lineagev1.GetLineageStartPoint{
				From: &lineagev1.GetLineageStartPoint_EntityColumns{
					EntityColumns: &lineagev1.EntityColumnsStartPoint{Id: s.target, ColumnNames: columns},
				},
			},
		})
		if err != nil {
			fmt.Printf("   (lineage read failed: %v)\n", err)
			return
		}
		lin := resp.GetLineage()
		// Two hops per column: source -> stream -> target.
		if len(lin.GetColumnDependencies()) >= 2*len(columns) {
			for _, d := range lin.GetColumnDependencies() {
				fmt.Printf("   %s.%s -> %s.%s\n",
					lin.GetNodes()[d.GetSourceNodeIdx()].GetIds()[0].GetEntityId(), d.GetSourceNodeColumnId(),
					lin.GetNodes()[d.GetTargetNodeIdx()].GetIds()[0].GetEntityId(), d.GetTargetNodeColumnId(),
				)
			}
			return
		}
		if time.Now().After(deadline) {
			fmt.Printf("   (gave up after %s: check both tables exist in Coalesce Quality under these identifiers)\n", budget)
			return
		}
		time.Sleep(10 * time.Second)
	}
}

func streamID(s stream) *entitiesv1.Identifier {
	return &entitiesv1.Identifier{Id: &entitiesv1.Identifier_Custom{
		Custom: &entitiesv1.CustomIdentifier{Id: "replication::stream::" + s.name},
	}}
}

func postgresTable(host, database, schema, table string) *entitiesv1.Identifier {
	return &entitiesv1.Identifier{Id: &entitiesv1.Identifier_PostgresTable{
		PostgresTable: &entitiesv1.PostgresTableIdentifier{Host: host, Database: database, Schema: schema, Table: table},
	}}
}

func snowflakeTable(account, database, schema, table string) *entitiesv1.Identifier {
	return &entitiesv1.Identifier{Id: &entitiesv1.Identifier_SnowflakeTable{
		SnowflakeTable: &entitiesv1.SnowflakeTableIdentifier{Account: account, Database: database, Schema: schema, Table: table},
	}}
}

func env(name, fallback string) string {
	if v := os.Getenv(name); v != "" {
		return v
	}
	return fallback
}

func must(err error) {
	if err != nil {
		panic(err)
	}
}

var iconArrows = []byte(`<svg xmlns="http://www.w3.org/2000/svg" viewBox="0 0 512 512"><path d="M32 160h352l-64-64 32-32 128 128-128 128-32-32 64-64H32z"/><path d="M480 352H128l64 64-32 32L32 320l128-128 32 32-64 64h352z"/></svg>`)
```


This documentation is built and hosted on [Mintlify](https://mintlify.com), a developer documentation platform.