---
title: "Tracing Data Lineage"
output: rmarkdown::html_vignette
vignette: >
  %\VignetteIndexEntry{Tracing Data Lineage}
  %\VignetteEngine{knitr::rmarkdown}
  %\VignetteEncoding{UTF-8}
---

```{r, include = FALSE}
knitr::opts_chunk$set(
  collapse = TRUE,
  comment  = "#>",
  eval     = FALSE
)
```

datom records two complementary lineage fields in every table's metadata:

- `parents` -- the **immediate** inputs used to derive this table
  (identified by `metadata_sha`, the datom version identifier).
- `source_lineage` -- the **transitive closure** of all raw tables that
  contributed data, stored as a flat list of `{project, table,
  version_sha}` entries where `version_sha` is the content SHA of the
  source parquet file.

Together they answer two different audit questions without walking a DAG:

| Question | Field | Tool |
|---|---|---|
| "What exact versions did my script read?" | `parents` | `datom_get_lineage(depth = "parents")` |
| "What raw data is ultimately in this table?" | `source_lineage` | `datom_get_lineage(depth = "source")` |

This article shows how lineage is recorded, how to query it, and how to
recompute it from existing reads to check consistency.

---

## Setup

The walkthrough below is illustrative (`eval = FALSE`) -- it shows the calls
without needing live credentials. As in the other articles, keep
machine-specific values up front:

```{r}
library(datom)

project_name <- "my-study"
bucket       <- "my-datom-bucket"
prefix       <- "lineage-demo"
region       <- "us-east-1"
access_key   <- keyring::key_get("AWS_ACCESS_KEY_ID")
secret_key   <- keyring::key_get("AWS_SECRET_ACCESS_KEY")
github_pat   <- keyring::key_get("GITHUB_PAT")
```

Writing derived tables (below) needs a **developer** connection. The lineage
*queries* -- `datom_get_parents()` and `datom_get_lineage()` -- are reads, and
the consistency recipe near the end is composed entirely of those reads plus
`datom_lineage_union()`, so it works just as well from a **reader** connection
(build the same store with `github_pat` omitted). This setup assumes the
project already exists:

```{r}
store <- datom_store(
  governance = NULL,
  data       = datom_store_s3(
    bucket     = bucket,
    prefix     = prefix,
    region     = region,
    access_key = access_key,
    secret_key = secret_key
  ),
  github_pat = github_pat            # omit for a read-only reader connection
)

conn <- datom_get_conn(path = "path/to/my-study-dev", store = store)
```

---

## A simple study: raw imports to derived analysis

Imagine a clinical study with two raw source tables:

```{r}
raw_dm <- datom_example_data("dm")      # Demographics
raw_lb <- datom_example_data("lb")      # Lab results
```

They are onboarded via `datom_sync()`, the file-based workflow shown in the
[Getting Started](getting-started.html) article: stage the raw tables as files
in the clone's
gitignored `input_files/` inbox, scan them into a manifest with
`datom_sync_manifest()`, then `datom_sync()`. datom automatically records
each imported table's own content SHA as its `source_lineage` -- a single
self-entry.

```{r}
# input_files/ lives inside the git clone and is the sync inbox.
input_dir <- file.path("path/to/my-study-dev", "input_files")

write.csv(raw_dm, file.path(input_dir, "dm.csv"), row.names = FALSE)
write.csv(raw_lb, file.path(input_dir, "lb.csv"), row.names = FALSE)

manifest <- datom_sync_manifest(conn)   # scans input_files/
datom_sync(conn, manifest)              # onboards; self-lineage recorded
```

After syncing, the imported tables carry their own lineage:

```{r}
datom_get_lineage(conn, "dm", depth = "source")
#> [[1]]
#> [[1]]$project
#> [1] "my-study"
#>
#> [[1]]$table
#> [1] "dm"
#>
#> [[1]]$version_sha
#> [1] "abc123..."
```

The self-entry means: "the raw content of `dm` is itself a source."

---

## Deriving a table

A downstream script reads both raw tables and produces a cleaned
demographics table:

```{r}
# Read the current versions
raw_dm_data <- datom_read(conn, "dm")
raw_lb_data <- datom_read(conn, "lb")

# ... cleaning logic ...
dm_clean <- raw_dm_data  # simplified

# Retrieve the metadata shas to identify exactly which versions were read
dm_version  <- datom_history(conn, "dm")[1, "version"]
lb_version  <- datom_history(conn, "lb")[1, "version"]
```

When writing the derived table, declare each parent with `datom_parent()`.
Each record reads the parent's authoritative `data_sha` and its
`source_lineage` from the parent's own versioned snapshot. `datom_write()`
then derives the derived table's `source_lineage` itself, as the
deduplicated union of the parents' lineages -- there is no public
`source_lineage` argument to supply or keep in sync:

```{r}
datom_write(
  conn,
  data    = dm_clean,
  name    = "dm_clean",
  parents = list(
    datom_parent(conn, "dm", dm_version),
    datom_parent(conn, "lb", lb_version)
  )
)
```

Now `dm_clean` knows which raw tables it came from:

```{r}
datom_get_lineage(conn, "dm_clean", depth = "source")
#> [[1]]
#> $project
#> [1] "my-study"
#> $table
#> [1] "dm"
#> $version_sha
#> [1] "abc123..."
#>
#> [[2]]
#> $project
#> [1] "my-study"
#> $table
#> [1] "lb"
#> $version_sha
#> [1] "def456..."
```

And which immediate versions it was derived from:

```{r}
datom_get_lineage(conn, "dm_clean", depth = "parents")
#> [[1]]
#> $source
#> [1] "my-study"
#> $table
#> [1] "dm"
#> $version
#> [1] "..."   # metadata_sha of dm at derivation time
```

---

## Propagating lineage further downstream

An analysis table derived from `dm_clean` propagates the lineage
automatically. Because `dm_clean` already encodes both `dm` and `lb`,
`datom_parent(conn, "dm_clean", clean_version)` captures that transitive
lineage, and `datom_write()` unions it into `analysis_pop` without
re-reading the original files:

```{r}
# Identify the dm_clean version this analysis was derived from.
clean_version <- datom_history(conn, "dm_clean")[1, "version"]

datom_write(
  conn,
  data    = analysis_pop,
  name    = "analysis_pop",
  parents = list(
    datom_parent(conn, "dm_clean", clean_version)
  )
)
```

---

## Recomputing lineage consistency

`datom_write()` derives a derived table's `source_lineage` from the union
of its parents' lineages at write time, and lineage is version-pinned. So a
recompute equals the recorded value in normal operation -- a difference
flags drift or corruption worth investigating. There is no dedicated
validator; instead you compose the existing reads with
`datom_lineage_union()`.

The recipe has four steps:

```{r}
# 1. Read the derived table's recorded parents. Each entry carries
#    source, table, version, and data_sha -- enough to pick the parent's
#    project connection and its pinned version.
parents <- datom_get_parents(conn, "dm_clean")

# 2. Read each parent's source_lineage through a connection scoped to that
#    parent's project. For same-project parents this is the same `conn`.
parent_lineages <- lapply(parents, function(p) {
  datom_get_lineage(conn, p$table, version = p$version, depth = "source")
})

# 3. Union the parents' lineages (dedup by {project, table, version_sha}).
recomputed <- datom_lineage_union(parent_lineages)

# 4. Compare against the derived table's recorded source_lineage.
recorded <- datom_get_lineage(conn, "dm_clean", depth = "source")
identical(recomputed, recorded)
#> [1] TRUE
```

Because the recipe reads each parent through its own connection, it extends
to **cross-project** parents without any change: open a connection scoped to
each parent's project (`p$source`) and read that parent through it. No single
connection is ever expected to reach across project stores.

```{r}
# Cross-project variant: resolve a connection per parent project.
parent_lineages <- lapply(parents, function(p) {
  parent_conn <- conn_for_project(p$source)   # your connection resolver
  datom_get_lineage(parent_conn, p$table, version = p$version,
                    depth = "source")
})
recomputed <- datom_lineage_union(parent_lineages)
```

---

## Key points

- `datom_sync()` auto-populates `source_lineage` for imported tables
  (a single self-entry using the file's content SHA).
- `datom_write()` derives a derived table's `source_lineage` from the union
  of its `datom_parent()` records -- there is no public `source_lineage`
  argument to keep in sync.
- `datom_get_lineage()` is a single-read operation -- no DAG traversal,
  no recursive network calls.
- Lineage consistency is a **composable recipe**: `datom_get_parents()` +
  per-parent `datom_get_lineage(depth = "source")` + `datom_lineage_union()`,
  compared against the recorded `source_lineage`. It is built from reads, so
  it runs from a reader connection and honors one connection per project.
- **Walker invariant**: `source_lineage` entries are terminal leaves.
  Lineage walkers must follow `parents`, never `source_lineage`.
