source

package
v0.24.10 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 23, 2026 License: MIT Imports: 17 Imported by: 0

Documentation

Overview

Package source implements the unified sync-ETL Source adapter layer (#97).

Pipeline contract (epic #88):

Source.Fetch(ctx, cursor) → []Blob, Cursor, error

A Source yields new items since the opaque cursor and returns the cursor that resumes after the batch. The driver (Sync) layers durable checkpointing on top: a sha256 seen-set guarantees each item is emitted exactly once, and the checkpoint (cursor + seen-set) is persisted atomically to var/state/<Name>.json (path supplied by the caller from config). A mid-batch failure persists the seen-set so the next run resumes without re-processing already-consumed items.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func ImportPST

func ImportPST(ctx context.Context, o ImportOptions) (Stats, ConvStats, error)

ImportPST runs one PST import: source.Sync (extract via readpst, sha256 seen-set dedup, atomic checkpoint) and then the conversion of the extracted corpus root through the shared mailconv.FromEML path (folder = PST folder, e.g. "Posteingang"). Re-runs convert nothing twice.

func PlanPST

func PlanPST(o ImportOptions) []string

PlanPST returns the --dry-run lines for an import: every configured source with its corpus target, the converter resolution and the scratch/state paths. No filesystem writes happen.

Types

type Blob

type Blob struct {
	ID   string `json:"id"`
	Kind string `json:"kind,omitempty"`
	Path string `json:"path,omitempty"`
	Data []byte `json:"-"`
}

Blob is one unit yielded by a Source. ID is the stable, content-derived identity; the driver stores sha256(ID) in the seen-set, so two blobs with the same ID across runs are the same item and are emitted at most once. Kind selects the downstream ETL handler (mail/git/chat/...); Path locates the item for handlers that read from disk (Data stays nil — lazy loading, ETL #9).

type ConvStats

type ConvStats struct {
	OK   int
	Skip int
	Fail int
}

ConvStats reports the mailconv.FromEML conversion pass of one run.

type Cursor

type Cursor string

Cursor is opaque, source-specific resume state. The empty Cursor signals a first run. Sources are free to encode a page token, a delta link, or the last processed id.

type Disk

type Disk struct {
	Root string
}

Disk lists .eml files under Root and yields one Blob per file. It is the sync-ETL adapter for on-disk mail corpora (var/corpus/* directories).

The cursor never advances: Fetch re-lists the tree on every call and the driver's sha256 seen-set supplies idempotency (the seen key is derived from the file content, so a changed .eml is re-emitted once). This keeps the adapter stateless — the source scans, the driver dedups.

func (*Disk) Fetch

func (d *Disk) Fetch(_ context.Context, _ Cursor) ([]Blob, Cursor, error)

Fetch returns one Blob per .eml under Root, in deterministic order. Each blob's ID is the hex sha256 of the file content (content-addressed identity).

func (*Disk) Name

func (d *Disk) Name() string

type Git

type Git struct {
	Repo string
}

Git yields one Blob per commit in a repository's history (go-git, no git binary), newest first. It is the sync-ETL adapter for git-history sources.

The cursor is the newest commit already consumed (the HEAD of the previous run): Fetch returns every commit strictly newer than it, then advances the cursor to the newest commit of the returned batch. New commits pushed on top are picked up by the next run; a repeat Fetch with no new commits yields 0 blobs, so the adapter is idempotent on its own.

func (*Git) Fetch

func (g *Git) Fetch(_ context.Context, cursor Cursor) ([]Blob, Cursor, error)

Fetch returns commits strictly newer than cursor (or all commits on a first run). next is the newest commit of the batch — the next resume point.

func (*Git) Name

func (g *Git) Name() string

type ImportOptions

type ImportOptions struct {
	Sources   []PSTSource
	Staging   string
	Out       string
	ReadPST   string
	StatePath string
}

ImportOptions wires one PST import run (bin/mail/import-pst.go): the sources plus the staging/out/state dirs derived from config by the CLI.

type Options

type Options struct {
	// StatePath is the checkpoint file for the source (var/state/<Name>.json).
	// The caller supplies it from config; a relative path is allowed in tests
	// (temp dir). Required.
	StatePath string
}

Options configures a Sync run.

type PST

type PST struct {
	// Sources are the .pst files to convert, one label+path pair each (from
	// config pst.sources, see #79). The label names the corpus subdir.
	Sources []PSTSource
	// Staging is the scratch root readpst writes into; it is wiped per source
	// before every extraction (var/tmp/pst). Never point it at the corpus.
	Staging string
	// Out is the corpus root for extracted mail (var/corpus/mail/pst).
	Out string
	// ReadPST is the readpst binary path override; empty = PATH lookup, then
	// the repo-local var/dist toolchain dir, then an explicit config error.
	ReadPST string
}

PST is the sync-ETL adapter for Outlook .pst archives (#185): it runs readpst -e on each configured source into a wiped staging dir and yields one Blob per extracted .eml message, copied content-addressed into the corpus (<Out>/<label>/<folder>/<sha256:16>/<sha256:16>.eml — the same layout as the mbox splitter). It closes the PST import gap of #79.

The cursor never advances (like Disk): Fetch re-runs readpst and re-lists staging on every call; the driver's sha256 seen-set supplies idempotency, so a re-run converts nothing twice even though readpst overwrites the staging files. readpst extraction is deterministic (item order + folder names), so content IDs are stable across runs. Folders the #79 policy excludes (Drafts/Templates/Trash/Junk/Spam/Unsent, incl. German Outlook names) are skipped before the corpus copy.

func (*PST) Fetch

func (p *PST) Fetch(ctx context.Context, _ Cursor) ([]Blob, Cursor, error)

Fetch extracts every configured source and returns one Blob per extracted .eml message (policy-excluded folders skipped). Blob IDs are content hashes, so the driver seen-set skips unchanged messages on re-runs.

func (*PST) Name

func (p *PST) Name() string

Name is the checkpoint file stem: var/state/pst.json.

type PSTSource

type PSTSource struct {
	Label string
	Path  string
}

PSTSource is one .pst file to import: Label (corpus subdir / source tag) + Path (absolute .pst path from config, see #79).

type Source

type Source interface {
	// Name is the checkpoint file stem: var/state/<Name>.json. It must be
	// stable per source instance and safe as a file name.
	Name() string
	Fetch(ctx context.Context, cursor Cursor) ([]Blob, Cursor, error)
}

Source is the unified sync-ETL adapter contract (#97). Fetch returns the batch of blobs new since cursor plus the cursor that resumes after this batch; an empty batch with an unchanged cursor signals no new data. The driver drives Fetch sequentially, so implementations need not be internally concurrency-safe.

type Stats

type Stats struct {
	Fetched int // blobs returned by Fetch (incl. skipped duplicates)
	New     int // blobs handed to handle (not in the seen-set)
	Skipped int // already-seen blobs dropped
}

Stats reports what one Sync run did.

func Sync

func Sync(ctx context.Context, src Source, handle func(context.Context, Blob) error, o Options) (Stats, error)

Sync drives src: it loops Fetch → dedup via the sha256 seen-set → handle, persisting the checkpoint atomically after every successful batch and on any mid-batch failure. handle is called once per new blob.

Semantics:

  • a blob is emitted at most once (seen-set idempotency);
  • a batch whose cursor does not advance (or that returns no data) ends the run — "no new data" ⇒ 0 new blobs;
  • on a mid-batch failure the checkpoint keeps cursor at the batch start and the seen-set up to the failed item, so the next Sync resumes exactly where it stopped without re-processing already-consumed items.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL