diff --git a/cmd/reconcile.go b/cmd/reconcile.go
new file mode 100644
index 00000000..6ae11672
--- /dev/null
+++ b/cmd/reconcile.go
@@ -0,0 +1,419 @@
+package cmd
+
+import (
+ "context"
+ "encoding/json"
+ "errors"
+ "flag"
+ "fmt"
+ "io"
+ "net/http"
+ "os"
+ "path/filepath"
+ "strings"
+ "time"
+
+ "github.com/rs/zerolog"
+
+ "github.com/maxghenis/openmessage/internal/app"
+ "github.com/maxghenis/openmessage/internal/db"
+ "github.com/maxghenis/openmessage/internal/reconcile"
+ "github.com/maxghenis/openmessage/internal/storage/sqlite"
+)
+
+const (
+ reconcileSourceExitCode = 3
+ reconcileLockExitCode = 4
+ reconcileRunExitCode = 5
+)
+
+type reconcileSignalError struct {
+ code int
+ err error
+}
+
+func (e *reconcileSignalError) Error() string { return e.err.Error() }
+func (e *reconcileSignalError) Unwrap() error { return e.err }
+func (e *reconcileSignalError) ExitCode() int { return e.code }
+
+func newReconcileSignalError(code int, format string, args ...any) error {
+ return &reconcileSignalError{code: code, err: fmt.Errorf(format, args...)}
+}
+
+type reconcileSignalCommandOptions struct {
+ sourceDir string
+ sinceMS int64
+ dryRun bool
+}
+
+type reconcileSignalFunc func(
+ context.Context,
+ reconcile.Options,
+) (reconcile.Report, error)
+
+type reconcileSignalDependencies struct {
+ now func() time.Time
+ version string
+ commit string
+ probeURL string
+ httpClient *http.Client
+ openLegacy func(string) (*db.Store, error)
+ openV2 func(string) (*sqlite.Store, error)
+ openV2ReadOnly func(string) (*sqlite.Store, error)
+ reconcile reconcileSignalFunc
+}
+
+type reconcileSignalResult struct {
+ report reconcile.Report
+ hasReport bool
+}
+
+type reconcileSignalFailureOutput struct {
+ OK bool `json:"ok"`
+ ExitCode int `json:"exit_code"`
+ Error string `json:"error"`
+}
+
+// RunReconcileSignal handles
+// "openmessage v2 reconcile-signal [--from dir] [--since YYYY-MM-DD] [--dry-run] [--json]".
+//
+// stdout is always exactly one JSON value. Human summaries and safety guidance
+// go only to stderr.
+func RunReconcileSignal(logger zerolog.Logger, args ...string) error {
+ return runReconcileSignalCommand(
+ context.Background(),
+ logger,
+ args,
+ defaultReconcileSignalDependencies(),
+ os.Stdout,
+ os.Stderr,
+ )
+}
+
+func runReconcileSignalCommand(
+ ctx context.Context,
+ logger zerolog.Logger,
+ args []string,
+ deps reconcileSignalDependencies,
+ stdout io.Writer,
+ stderr io.Writer,
+) error {
+ fs := flag.NewFlagSet("v2 reconcile-signal", flag.ContinueOnError)
+ fs.SetOutput(stderr)
+ fs.Usage = func() {
+ fmt.Fprintln(stderr, "Usage: openmessage v2 reconcile-signal [--from
] [--since YYYY-MM-DD] [--dry-run] [--json]")
+ fmt.Fprintln(stderr, " --from defaults to the OpenMessage data directory")
+ fmt.Fprintln(stderr, " --since includes legacy messages at or after local midnight on that date")
+ fmt.Fprintln(stderr, " --dry-run performs all reads and derivations without mutating either message store")
+ fmt.Fprintln(stderr, " the reconciliation report is JSON on stdout; human guidance is on stderr")
+ }
+
+ options := reconcileSignalCommandOptions{sourceDir: defaultReconcileSignalDataDir()}
+ var since string
+ var explicitJSON bool
+ fs.StringVar(&options.sourceDir, "from", options.sourceDir, "OpenMessage data directory")
+ fs.StringVar(&since, "since", "", "earliest legacy message date (YYYY-MM-DD)")
+ fs.BoolVar(&options.dryRun, "dry-run", false, "report prospective changes without writing")
+ fs.BoolVar(&explicitJSON, "json", false, "emit the reconciliation report as JSON (always enabled)")
+ if err := fs.Parse(args); err != nil {
+ if errors.Is(err, flag.ErrHelp) {
+ return nil
+ }
+ commandErr := newReconcileSignalError(
+ reconcileSourceExitCode,
+ "parse reconcile-signal options: %w",
+ err,
+ )
+ return writeReconcileSignalFailure(stdout, stderr, commandErr)
+ }
+ _ = explicitJSON
+ if fs.NArg() != 0 {
+ commandErr := newReconcileSignalError(
+ reconcileSourceExitCode,
+ "unexpected reconcile-signal argument %q; use --from for the data directory",
+ fs.Arg(0),
+ )
+ return writeReconcileSignalFailure(stdout, stderr, commandErr)
+ }
+ sinceMS, err := parseDayBound(since, false)
+ if err != nil {
+ commandErr := newReconcileSignalError(
+ reconcileSourceExitCode,
+ "parse --since: %w",
+ err,
+ )
+ return writeReconcileSignalFailure(stdout, stderr, commandErr)
+ }
+ options.sinceMS = sinceMS
+
+ result, commandErr := executeReconcileSignal(ctx, logger, options, deps)
+ if result.hasReport {
+ if err := json.NewEncoder(stdout).Encode(result.report); err != nil {
+ return newReconcileSignalError(
+ reconcileRunExitCode,
+ "write Signal reconciliation report: %w",
+ err,
+ )
+ }
+ if commandErr == nil {
+ writeHumanReconcileSignalReport(stderr, result.report)
+ }
+ } else if commandErr != nil {
+ return writeReconcileSignalFailure(stdout, stderr, commandErr)
+ }
+ if commandErr != nil {
+ fmt.Fprintf(stderr, "Signal reconciliation failed: %v\n", commandErr)
+ }
+ return commandErr
+}
+
+func defaultReconcileSignalDependencies() reconcileSignalDependencies {
+ commit, _ := buildRevision()
+ return reconcileSignalDependencies{
+ now: time.Now,
+ version: Version(),
+ commit: commit,
+ probeURL: defaultBackendProbeURL(),
+ httpClient: &http.Client{
+ Timeout: 750 * time.Millisecond,
+ CheckRedirect: func(*http.Request, []*http.Request) error {
+ return http.ErrUseLastResponse
+ },
+ },
+ openLegacy: db.OpenReadOnly,
+ openV2: sqlite.Open,
+ openV2ReadOnly: sqlite.OpenReadOnly,
+ reconcile: reconcile.Signal,
+ }
+}
+
+func defaultReconcileSignalDataDir() string {
+ if dataDir := os.Getenv("OPENMESSAGES_DATA_DIR"); strings.TrimSpace(dataDir) != "" {
+ return dataDir
+ }
+ home, err := os.UserHomeDir()
+ if err != nil || strings.TrimSpace(home) == "" {
+ return app.DefaultDataDir()
+ }
+ return filepath.Join(home, "Library", "Application Support", "OpenMessage")
+}
+
+func executeReconcileSignal(
+ ctx context.Context,
+ logger zerolog.Logger,
+ options reconcileSignalCommandOptions,
+ deps reconcileSignalDependencies,
+) (result reconcileSignalResult, resultErr error) {
+ if ctx == nil {
+ return result, newReconcileSignalError(
+ reconcileRunExitCode,
+ "Signal reconciliation context is nil",
+ )
+ }
+ if deps.now == nil || deps.httpClient == nil || deps.openLegacy == nil ||
+ deps.openV2 == nil || deps.openV2ReadOnly == nil || deps.reconcile == nil {
+ return result, newReconcileSignalError(
+ reconcileRunExitCode,
+ "Signal reconciliation dependencies are incomplete",
+ )
+ }
+
+ sourceDir, err := canonicalExistingDirectory(options.sourceDir)
+ if err != nil {
+ return result, newReconcileSignalError(
+ reconcileSourceExitCode,
+ "resolve OpenMessage data directory %q: %w",
+ options.sourceDir,
+ err,
+ )
+ }
+ legacyPath := filepath.Join(sourceDir, legacyDatabaseName)
+ if err := requireRegularReconcileStore("legacy database", legacyPath); err != nil {
+ return result, err
+ }
+ v2Path := filepath.Join(sourceDir, "v2", v2StoreName)
+ if err := requireRegularReconcileStore("v2 database", v2Path); err != nil {
+ return result, err
+ }
+
+ createdAt := deps.now().UTC().Truncate(time.Second)
+ lockPath := filepath.Join(sourceDir, instanceLockName)
+ lock, err := acquireInstanceLock(lockPath, instanceLockRecord{
+ PID: os.Getpid(),
+ Process: "openmessage v2 reconcile-signal",
+ StartedAt: createdAt.Format(time.RFC3339),
+ BuildID: deps.version,
+ Commit: deps.commit,
+ CanonicalPath: sourceDir,
+ ProbeURL: deps.probeURL,
+ })
+ if err != nil {
+ if errors.Is(err, errInstanceLockHeld) {
+ return result, newReconcileSignalError(
+ reconcileLockExitCode,
+ "OpenMessage state is in use: %s is locked; stop the backend and retry",
+ lockPath,
+ )
+ }
+ return result, newReconcileSignalError(
+ reconcileLockExitCode,
+ "acquire OpenMessage instance lock %s: %w",
+ lockPath,
+ err,
+ )
+ }
+ defer func() {
+ if closeErr := lock.Close(); closeErr != nil && resultErr == nil {
+ resultErr = newReconcileSignalError(
+ reconcileRunExitCode,
+ "release OpenMessage instance lock: %w",
+ closeErr,
+ )
+ }
+ }()
+
+ running, probeErr := probeBackend(ctx, deps.httpClient, deps.probeURL)
+ if probeErr != nil {
+ return result, newReconcileSignalError(
+ reconcileLockExitCode,
+ "cannot safely rule out a running OpenMessage backend at %s: %w; stop the backend and retry",
+ deps.probeURL,
+ probeErr,
+ )
+ }
+ if running {
+ return result, newReconcileSignalError(
+ reconcileLockExitCode,
+ "OpenMessage backend is running at %s; stop the backend and retry",
+ deps.probeURL,
+ )
+ }
+
+ legacy, err := deps.openLegacy(legacyPath)
+ if err != nil {
+ return result, newReconcileSignalError(
+ reconcileSourceExitCode,
+ "open legacy database %s: %w",
+ legacyPath,
+ err,
+ )
+ }
+ defer func() {
+ if closeErr := legacy.Close(); closeErr != nil && resultErr == nil {
+ resultErr = newReconcileSignalError(
+ reconcileRunExitCode,
+ "close legacy database: %w",
+ closeErr,
+ )
+ }
+ }()
+ openV2 := deps.openV2
+ if options.dryRun {
+ openV2 = deps.openV2ReadOnly
+ }
+ v2, err := openV2(v2Path)
+ if err != nil {
+ return result, newReconcileSignalError(
+ reconcileSourceExitCode,
+ "open v2 database %s: %w",
+ v2Path,
+ err,
+ )
+ }
+ defer func() {
+ if closeErr := v2.Close(); closeErr != nil && resultErr == nil {
+ resultErr = newReconcileSignalError(
+ reconcileRunExitCode,
+ "close v2 database: %w",
+ closeErr,
+ )
+ }
+ }()
+
+ report, reconcileErr := deps.reconcile(ctx, reconcile.Options{
+ Legacy: legacy,
+ V2: v2,
+ SinceMS: options.sinceMS,
+ DryRun: options.dryRun,
+ Logger: logger,
+ })
+ result = reconcileSignalResult{report: report, hasReport: true}
+ if reconcileErr != nil {
+ return result, newReconcileSignalError(
+ reconcileRunExitCode,
+ "reconcile Signal history: %w",
+ reconcileErr,
+ )
+ }
+ return result, nil
+}
+
+func requireRegularReconcileStore(description string, path string) error {
+ info, err := os.Stat(path)
+ if err != nil {
+ if os.IsNotExist(err) {
+ return newReconcileSignalError(
+ reconcileSourceExitCode,
+ "%s not found at %s",
+ description,
+ path,
+ )
+ }
+ return newReconcileSignalError(
+ reconcileSourceExitCode,
+ "inspect %s %s: %w",
+ description,
+ path,
+ err,
+ )
+ }
+ if !info.Mode().IsRegular() {
+ return newReconcileSignalError(
+ reconcileSourceExitCode,
+ "%s path is not a regular file: %s",
+ description,
+ path,
+ )
+ }
+ return nil
+}
+
+func writeReconcileSignalFailure(
+ stdout io.Writer,
+ stderr io.Writer,
+ commandErr error,
+) error {
+ output := reconcileSignalFailureOutput{
+ OK: false,
+ ExitCode: ExitCode(commandErr),
+ Error: commandErr.Error(),
+ }
+ if err := json.NewEncoder(stdout).Encode(output); err != nil {
+ return newReconcileSignalError(
+ reconcileRunExitCode,
+ "write Signal reconciliation error report: %w",
+ err,
+ )
+ }
+ fmt.Fprintf(stderr, "Signal reconciliation failed: %v\n", commandErr)
+ return commandErr
+}
+
+func writeHumanReconcileSignalReport(writer io.Writer, report reconcile.Report) {
+ action := "Signal reconciliation complete"
+ if report.DryRun {
+ action = "Signal reconciliation dry run complete"
+ }
+ fmt.Fprintf(
+ writer,
+ "%s: conversations scanned=%d, created=%d; messages scanned=%d, imported=%d, already present=%d, media deferred=%d, skipped=%d\n",
+ action,
+ report.ConversationsScanned,
+ report.ConversationsCreated,
+ report.MessagesScanned,
+ report.MessagesImported,
+ report.MessagesAlreadyPresent,
+ report.MediaDeferred,
+ report.Skipped,
+ )
+}
diff --git a/cmd/reconcile_test.go b/cmd/reconcile_test.go
new file mode 100644
index 00000000..83338278
--- /dev/null
+++ b/cmd/reconcile_test.go
@@ -0,0 +1,398 @@
+package cmd
+
+import (
+ "bytes"
+ "context"
+ "encoding/json"
+ "errors"
+ "io"
+ "net/http"
+ "os"
+ "path/filepath"
+ "strings"
+ "syscall"
+ "testing"
+ "time"
+
+ "github.com/rs/zerolog"
+
+ "github.com/maxghenis/openmessage/internal/db"
+ "github.com/maxghenis/openmessage/internal/reconcile"
+ "github.com/maxghenis/openmessage/internal/storage/sqlite"
+)
+
+func TestReconcileSignalRefusesWhenBackendProbeResponds(t *testing.T) {
+ sourceDir := newReconcileSourceFixture(t, false)
+ legacyPath := filepath.Join(sourceDir, legacyDatabaseName)
+ v2Path := filepath.Join(sourceDir, "v2", v2StoreName)
+ legacyBefore, err := os.ReadFile(legacyPath)
+ if err != nil {
+ t.Fatal(err)
+ }
+ v2Before, err := os.ReadFile(v2Path)
+ if err != nil {
+ t.Fatal(err)
+ }
+
+ var probed bool
+ deps := reconcileSignalDependencies{
+ now: func() time.Time { return time.Date(2026, 7, 26, 12, 0, 0, 0, time.UTC) },
+ version: "test-version",
+ commit: "test-commit",
+ probeURL: "http://127.0.0.1:7007/api/status",
+ httpClient: &http.Client{Transport: reconcileRoundTripperFunc(func(request *http.Request) (*http.Response, error) {
+ probed = true
+ if request.Method != http.MethodGet || request.URL.Path != "/api/status" {
+ t.Fatalf("probe request = %s %s", request.Method, request.URL)
+ }
+ return &http.Response{
+ StatusCode: http.StatusServiceUnavailable,
+ Header: make(http.Header),
+ Body: io.NopCloser(strings.NewReader("backend owns the endpoint")),
+ Request: request,
+ }, nil
+ })},
+ openLegacy: func(string) (*db.Store, error) {
+ t.Fatal("legacy store opened while backend probe reported running")
+ return nil, nil
+ },
+ openV2: func(string) (*sqlite.Store, error) {
+ t.Fatal("v2 store opened while backend probe reported running")
+ return nil, nil
+ },
+ openV2ReadOnly: func(string) (*sqlite.Store, error) {
+ t.Fatal("read-only v2 store opened while backend probe reported running")
+ return nil, nil
+ },
+ reconcile: func(context.Context, reconcile.Options) (reconcile.Report, error) {
+ t.Fatal("reconciler called while backend probe reported running")
+ return reconcile.Report{}, nil
+ },
+ }
+
+ var stdout, stderr bytes.Buffer
+ err = runReconcileSignalCommand(
+ context.Background(),
+ zerolog.Nop(),
+ []string{"--from", sourceDir},
+ deps,
+ &stdout,
+ &stderr,
+ )
+ if err == nil {
+ t.Fatal("runReconcileSignalCommand() succeeded while backend was running")
+ }
+ if !probed {
+ t.Fatal("backend was not probed")
+ }
+ if got := ExitCode(err); got != reconcileLockExitCode {
+ t.Fatalf("ExitCode = %d, want %d (%v)", got, reconcileLockExitCode, err)
+ }
+ var output reconcileSignalFailureOutput
+ if err := json.Unmarshal(stdout.Bytes(), &output); err != nil {
+ t.Fatalf("decode failure output: %v\n%s", err, stdout.String())
+ }
+ if output.OK || output.ExitCode != reconcileLockExitCode ||
+ !strings.Contains(output.Error, "backend is running") {
+ t.Fatalf("failure output = %+v", output)
+ }
+ if !strings.Contains(stderr.String(), "stop the backend and retry") {
+ t.Fatalf("stderr is not actionable: %s", stderr.String())
+ }
+ legacyAfter, err := os.ReadFile(legacyPath)
+ if err != nil {
+ t.Fatal(err)
+ }
+ v2After, err := os.ReadFile(v2Path)
+ if err != nil {
+ t.Fatal(err)
+ }
+ if !bytes.Equal(legacyBefore, legacyAfter) || !bytes.Equal(v2Before, v2After) {
+ t.Fatal("backend refusal changed a message store")
+ }
+}
+
+func TestReconcileSignalParsesFlagsAndEmitsOneJSONReport(t *testing.T) {
+ sourceDir := newReconcileSourceFixture(t, true)
+ legacyPath := filepath.Join(sourceDir, legacyDatabaseName)
+ v2Path := filepath.Join(sourceDir, "v2", v2StoreName)
+ legacyBefore, err := os.ReadFile(legacyPath)
+ if err != nil {
+ t.Fatal(err)
+ }
+ v2Before, err := os.ReadFile(v2Path)
+ if err != nil {
+ t.Fatal(err)
+ }
+ deps := reconcileSignalDependencies{
+ now: func() time.Time { return time.Date(2026, 7, 26, 12, 0, 0, 0, time.UTC) },
+ version: "test-version",
+ commit: "test-commit",
+ probeURL: "http://127.0.0.1:7007/api/status",
+ httpClient: &http.Client{Transport: reconcileRoundTripperFunc(func(*http.Request) (*http.Response, error) {
+ return nil, syscall.ECONNREFUSED
+ })},
+ openLegacy: db.OpenReadOnly,
+ openV2: func(string) (*sqlite.Store, error) {
+ t.Fatal("writable v2 store opened for a dry run")
+ return nil, nil
+ },
+ openV2ReadOnly: sqlite.OpenReadOnly,
+ }
+ want := reconcile.Report{
+ DryRun: true,
+ ConversationsScanned: 2,
+ ConversationsCreated: 1,
+ MessagesScanned: 7,
+ MessagesImported: 3,
+ MessagesAlreadyPresent: 4,
+ MediaDeferred: 1,
+ Skipped: 0,
+ SkipReasons: map[string]int{},
+ }
+ var called bool
+ deps.reconcile = func(
+ _ context.Context,
+ options reconcile.Options,
+ ) (reconcile.Report, error) {
+ called = true
+ if options.Legacy == nil || options.V2 == nil {
+ t.Fatal("reconcile options omitted a store")
+ }
+ if !options.DryRun {
+ t.Fatal("DryRun = false, want true")
+ }
+ wantSince, err := parseDayBound("2026-07-15", false)
+ if err != nil {
+ t.Fatal(err)
+ }
+ if options.SinceMS != wantSince {
+ t.Fatalf("SinceMS = %d, want %d", options.SinceMS, wantSince)
+ }
+ want.SinceMS = wantSince
+ return want, nil
+ }
+
+ var stdout, stderr bytes.Buffer
+ err = runReconcileSignalCommand(
+ context.Background(),
+ zerolog.Nop(),
+ []string{
+ "--from", sourceDir,
+ "--since", "2026-07-15",
+ "--dry-run",
+ "--json",
+ },
+ deps,
+ &stdout,
+ &stderr,
+ )
+ if err != nil {
+ t.Fatalf("runReconcileSignalCommand(): %v\nstderr:\n%s", err, stderr.String())
+ }
+ if !called {
+ t.Fatal("reconciler was not called")
+ }
+ var got reconcile.Report
+ if err := json.Unmarshal(stdout.Bytes(), &got); err != nil {
+ t.Fatalf("stdout is not a reconciliation report: %v\n%s", err, stdout.String())
+ }
+ if got.MessagesImported != want.MessagesImported ||
+ got.MessagesAlreadyPresent != want.MessagesAlreadyPresent ||
+ got.SinceMS != want.SinceMS ||
+ !got.DryRun {
+ t.Fatalf("report = %+v, want %+v", got, want)
+ }
+ decoder := json.NewDecoder(bytes.NewReader(stdout.Bytes()))
+ if err := decoder.Decode(&got); err != nil {
+ t.Fatalf("decode first stdout value: %v", err)
+ }
+ var extra any
+ if err := decoder.Decode(&extra); err != io.EOF {
+ t.Fatalf("stdout has more than one JSON value: err=%v extra=%v", err, extra)
+ }
+ if !strings.Contains(stderr.String(), "Signal reconciliation dry run complete") {
+ t.Fatalf("human dry-run report missing from stderr: %s", stderr.String())
+ }
+ legacyAfter, err := os.ReadFile(legacyPath)
+ if err != nil {
+ t.Fatal(err)
+ }
+ v2After, err := os.ReadFile(v2Path)
+ if err != nil {
+ t.Fatal(err)
+ }
+ if !bytes.Equal(legacyBefore, legacyAfter) || !bytes.Equal(v2Before, v2After) {
+ t.Fatal("dry run changed a message store")
+ }
+}
+
+func TestDefaultReconcileSignalDataDir(t *testing.T) {
+ t.Run("macOS app directory", func(t *testing.T) {
+ t.Setenv("OPENMESSAGES_DATA_DIR", "")
+ home := filepath.Join(t.TempDir(), "home")
+ t.Setenv("HOME", home)
+
+ want := filepath.Join(home, "Library", "Application Support", "OpenMessage")
+ if got := defaultReconcileSignalDataDir(); got != want {
+ t.Fatalf("defaultReconcileSignalDataDir() = %q, want %q", got, want)
+ }
+ })
+
+ t.Run("environment override", func(t *testing.T) {
+ override := filepath.Join(t.TempDir(), "explicit data")
+ t.Setenv("OPENMESSAGES_DATA_DIR", override)
+
+ if got := defaultReconcileSignalDataDir(); got != override {
+ t.Fatalf("defaultReconcileSignalDataDir() = %q, want %q", got, override)
+ }
+ })
+}
+
+func TestReconcileSignalPartialFailureEmitsReportWithoutCompletion(t *testing.T) {
+ sourceDir := newReconcileSourceFixture(t, true)
+ sentinel := errors.New("injected reconciliation failure")
+ deps := reconcileSignalDependencies{
+ now: func() time.Time { return time.Date(2026, 7, 26, 12, 0, 0, 0, time.UTC) },
+ version: "test-version",
+ commit: "test-commit",
+ probeURL: "http://127.0.0.1:7007/api/status",
+ httpClient: &http.Client{Transport: reconcileRoundTripperFunc(func(*http.Request) (*http.Response, error) {
+ return nil, syscall.ECONNREFUSED
+ })},
+ openLegacy: db.OpenReadOnly,
+ openV2: func(string) (*sqlite.Store, error) {
+ t.Fatal("writable v2 store opened for a dry run")
+ return nil, nil
+ },
+ openV2ReadOnly: sqlite.OpenReadOnly,
+ reconcile: func(
+ context.Context,
+ reconcile.Options,
+ ) (reconcile.Report, error) {
+ return reconcile.Report{
+ DryRun: true,
+ MessagesScanned: 3,
+ MessagesImported: 2,
+ SkipReasons: map[string]int{},
+ }, sentinel
+ },
+ }
+
+ var stdout, stderr bytes.Buffer
+ err := runReconcileSignalCommand(
+ context.Background(),
+ zerolog.Nop(),
+ []string{"--from", sourceDir, "--dry-run"},
+ deps,
+ &stdout,
+ &stderr,
+ )
+ if !errors.Is(err, sentinel) {
+ t.Fatalf("runReconcileSignalCommand() error = %v, want wrapped sentinel", err)
+ }
+ if got := ExitCode(err); got != reconcileRunExitCode {
+ t.Fatalf("ExitCode = %d, want %d", got, reconcileRunExitCode)
+ }
+ var report reconcile.Report
+ if decodeErr := json.Unmarshal(stdout.Bytes(), &report); decodeErr != nil {
+ t.Fatalf("decode partial report: %v\n%s", decodeErr, stdout.String())
+ }
+ if report.MessagesScanned != 3 || report.MessagesImported != 2 {
+ t.Fatalf("partial report = %+v", report)
+ }
+ if !strings.Contains(stderr.String(), "Signal reconciliation failed") {
+ t.Fatalf("failure missing from stderr: %s", stderr.String())
+ }
+ if strings.Contains(stderr.String(), "complete") {
+ t.Fatalf("stderr contradicts partial failure: %s", stderr.String())
+ }
+}
+
+func TestReconcileSignalUsesWritableV2ForRepair(t *testing.T) {
+ sourceDir := newReconcileSourceFixture(t, true)
+ var writableOpened bool
+ deps := reconcileSignalDependencies{
+ now: func() time.Time { return time.Date(2026, 7, 26, 12, 0, 0, 0, time.UTC) },
+ version: "test-version",
+ commit: "test-commit",
+ probeURL: "http://127.0.0.1:7007/api/status",
+ httpClient: &http.Client{Transport: reconcileRoundTripperFunc(func(*http.Request) (*http.Response, error) {
+ return nil, syscall.ECONNREFUSED
+ })},
+ openLegacy: db.OpenReadOnly,
+ openV2: func(path string) (*sqlite.Store, error) {
+ writableOpened = true
+ return sqlite.Open(path)
+ },
+ openV2ReadOnly: func(string) (*sqlite.Store, error) {
+ t.Fatal("read-only v2 store opened for a repair")
+ return nil, nil
+ },
+ reconcile: func(
+ context.Context,
+ reconcile.Options,
+ ) (reconcile.Report, error) {
+ return reconcile.Report{SkipReasons: map[string]int{}}, nil
+ },
+ }
+
+ var stdout, stderr bytes.Buffer
+ if err := runReconcileSignalCommand(
+ context.Background(),
+ zerolog.Nop(),
+ []string{"--from", sourceDir},
+ deps,
+ &stdout,
+ &stderr,
+ ); err != nil {
+ t.Fatalf("runReconcileSignalCommand(): %v\nstderr:\n%s", err, stderr.String())
+ }
+ if !writableOpened {
+ t.Fatal("repair did not open the writable v2 store")
+ }
+}
+
+func newReconcileSourceFixture(t *testing.T, validStores bool) string {
+ t.Helper()
+ sourceDir := t.TempDir()
+ v2Dir := filepath.Join(sourceDir, "v2")
+ if err := os.Mkdir(v2Dir, 0o700); err != nil {
+ t.Fatal(err)
+ }
+ legacyPath := filepath.Join(sourceDir, legacyDatabaseName)
+ v2Path := filepath.Join(v2Dir, v2StoreName)
+ if !validStores {
+ if err := os.WriteFile(legacyPath, []byte("legacy sentinel"), 0o600); err != nil {
+ t.Fatal(err)
+ }
+ if err := os.WriteFile(v2Path, []byte("v2 sentinel"), 0o600); err != nil {
+ t.Fatal(err)
+ }
+ return sourceDir
+ }
+
+ legacy, err := db.New(legacyPath)
+ if err != nil {
+ t.Fatalf("db.New(): %v", err)
+ }
+ if err := legacy.Close(); err != nil {
+ t.Fatalf("legacy.Close(): %v", err)
+ }
+ v2, err := sqlite.Open(v2Path)
+ if err != nil {
+ t.Fatalf("sqlite.Open(): %v", err)
+ }
+ if err := v2.Close(); err != nil {
+ t.Fatalf("v2.Close(): %v", err)
+ }
+ return sourceDir
+}
+
+type reconcileRoundTripperFunc func(*http.Request) (*http.Response, error)
+
+func (function reconcileRoundTripperFunc) RoundTrip(
+ request *http.Request,
+) (*http.Response, error) {
+ return function(request)
+}
diff --git a/internal/db/db.go b/internal/db/db.go
index 20a48486..ba0ccd8c 100644
--- a/internal/db/db.go
+++ b/internal/db/db.go
@@ -3,6 +3,8 @@ package db
import (
"database/sql"
"fmt"
+ "net/url"
+ "path/filepath"
"strings"
_ "modernc.org/sqlite"
@@ -133,6 +135,44 @@ func New(dsn string) (*Store, error) {
return s, nil
}
+// OpenReadOnly opens an existing legacy database without running migrations,
+// repair statements, or WAL mode changes. Maintenance readers use this seam so
+// the legacy source remains immutable even when its schema predates the binary.
+func OpenReadOnly(path string) (*Store, error) {
+ if strings.TrimSpace(path) == "" {
+ return nil, fmt.Errorf("open read-only db: path is empty")
+ }
+ query := make(url.Values)
+ query.Set("mode", "ro")
+ dsn := (&url.URL{
+ Scheme: "file",
+ Path: filepath.ToSlash(path),
+ RawQuery: query.Encode(),
+ }).String()
+ database, err := sql.Open("sqlite", dsn)
+ if err != nil {
+ return nil, fmt.Errorf("open read-only db: %w", err)
+ }
+ database.SetMaxOpenConns(1)
+ closeWithError := func(openErr error) (*Store, error) {
+ _ = database.Close()
+ return nil, openErr
+ }
+ if err := database.Ping(); err != nil {
+ return closeWithError(fmt.Errorf("open read-only db connection: %w", err))
+ }
+ for _, pragma := range []string{
+ "PRAGMA query_only=ON",
+ "PRAGMA foreign_keys=ON",
+ "PRAGMA busy_timeout=5000",
+ } {
+ if _, err := database.Exec(pragma); err != nil {
+ return closeWithError(fmt.Errorf("configure read-only db: %w", err))
+ }
+ }
+ return &Store{db: database}, nil
+}
+
func (s *Store) Close() error {
return s.db.Close()
}
diff --git a/internal/db/readonly_test.go b/internal/db/readonly_test.go
new file mode 100644
index 00000000..23256c7b
--- /dev/null
+++ b/internal/db/readonly_test.go
@@ -0,0 +1,71 @@
+package db
+
+import (
+ "bytes"
+ "os"
+ "path/filepath"
+ "testing"
+)
+
+func TestOpenReadOnlyReadsWithoutMutatingLegacyStore(t *testing.T) {
+ path := filepath.Join(t.TempDir(), "messages.db")
+ writable, err := New(path)
+ if err != nil {
+ t.Fatalf("New(): %v", err)
+ }
+ if err := writable.UpsertConversation(&Conversation{
+ ConversationID: "signal:+16505550100",
+ Name: "Signal Peer",
+ LastMessageTS: 1_700_000_001_000,
+ SourcePlatform: "signal",
+ }); err != nil {
+ t.Fatalf("UpsertConversation(): %v", err)
+ }
+ if err := writable.Close(); err != nil {
+ t.Fatalf("Close(writable): %v", err)
+ }
+ before, err := os.ReadFile(path)
+ if err != nil {
+ t.Fatal(err)
+ }
+
+ readOnly, err := OpenReadOnly(path)
+ if err != nil {
+ t.Fatalf("OpenReadOnly(): %v", err)
+ }
+ conversations, err := readOnly.ListConversationsByPlatform("signal", 10)
+ if err != nil {
+ t.Fatalf("ListConversationsByPlatform(): %v", err)
+ }
+ if len(conversations) != 1 ||
+ conversations[0].ConversationID != "signal:+16505550100" {
+ t.Fatalf("read-only conversations = %+v", conversations)
+ }
+ if err := readOnly.UpsertConversation(&Conversation{
+ ConversationID: "must-not-write",
+ SourcePlatform: "signal",
+ }); err == nil {
+ t.Fatal("read-only store accepted a write")
+ }
+ if err := readOnly.Close(); err != nil {
+ t.Fatalf("Close(read-only): %v", err)
+ }
+
+ after, err := os.ReadFile(path)
+ if err != nil {
+ t.Fatal(err)
+ }
+ if !bytes.Equal(before, after) {
+ t.Fatal("OpenReadOnly changed legacy database bytes")
+ }
+}
+
+func TestOpenReadOnlyDoesNotCreateMissingLegacyStore(t *testing.T) {
+ path := filepath.Join(t.TempDir(), "missing.db")
+ if _, err := OpenReadOnly(path); err == nil {
+ t.Fatal("OpenReadOnly() succeeded for a missing database")
+ }
+ if _, err := os.Stat(path); !os.IsNotExist(err) {
+ t.Fatalf("missing path was created: %v", err)
+ }
+}
diff --git a/internal/ingest/reconcile_signal_test.go b/internal/ingest/reconcile_signal_test.go
new file mode 100644
index 00000000..ef4e920c
--- /dev/null
+++ b/internal/ingest/reconcile_signal_test.go
@@ -0,0 +1,162 @@
+package ingest
+
+import (
+ "context"
+ "errors"
+ "path/filepath"
+ "strconv"
+ "testing"
+ "time"
+
+ "github.com/rs/zerolog"
+
+ legacydb "github.com/maxghenis/openmessage/internal/db"
+ "github.com/maxghenis/openmessage/internal/reconcile"
+ "github.com/maxghenis/openmessage/internal/storage/sqlite"
+ "github.com/maxghenis/openmessage/internal/v2keys"
+)
+
+func TestSignalReconciledOutgoingAliasIsFoundByWorkerBareTimestampPath(t *testing.T) {
+ const (
+ remoteConversationID = "signal:+16505550100"
+ timestamp = int64(1_700_000_006_000)
+ )
+ ctx := context.Background()
+ harness := newSignalWorkerHarness(t, "reconcile-alias.sqlite3")
+ legacy, err := legacydb.New(filepath.Join(t.TempDir(), "messages.db"))
+ if err != nil {
+ t.Fatalf("db.New(): %v", err)
+ }
+ t.Cleanup(func() {
+ if err := legacy.Close(); err != nil {
+ t.Errorf("legacy.Close(): %v", err)
+ }
+ })
+ if err := legacy.UpsertConversation(&legacydb.Conversation{
+ ConversationID: remoteConversationID,
+ Name: "Signal Peer",
+ LastMessageTS: timestamp,
+ SourcePlatform: "signal",
+ }); err != nil {
+ t.Fatalf("UpsertConversation(): %v", err)
+ }
+ alias := v2keys.SignalLocalAlias(remoteConversationID, timestamp)
+ if err := legacy.UpsertMessage(&legacydb.Message{
+ MessageID: "signal:legacy-outgoing-alias",
+ ConversationID: remoteConversationID,
+ SenderName: "Me",
+ Body: "reconciled outgoing body",
+ TimestampMS: timestamp,
+ IsFromMe: true,
+ SourcePlatform: "signal",
+ SourceID: alias,
+ }); err != nil {
+ t.Fatalf("UpsertMessage(): %v", err)
+ }
+
+ report, err := reconcile.Signal(ctx, reconcile.Options{
+ Legacy: legacy,
+ V2: harness.store,
+ Logger: zerolog.Nop(),
+ })
+ if err != nil {
+ t.Fatalf("reconcile.Signal(): %v", err)
+ }
+ if report.MessagesImported != 1 || report.MessagesAlreadyPresent != 0 {
+ t.Fatalf("reconcile report = %+v", report)
+ }
+ conversation, err := harness.store.GetConversationByRemote(
+ signalDecoderAccountID,
+ remoteConversationID,
+ )
+ if err != nil {
+ t.Fatalf("GetConversationByRemote(): %v", err)
+ }
+ reconciled, err := harness.messages.GetMessageByRemote(
+ ctx,
+ signalDecoderAccountID,
+ conversation.ConversationID,
+ alias,
+ )
+ if err != nil {
+ t.Fatalf("GetMessageByRemote(alias): %v", err)
+ }
+ wantMessageID := v2keys.DeriveID(
+ "message",
+ signalDecoderAccountID,
+ remoteConversationID+"\x1f"+alias,
+ )
+ if reconciled.MessageID != wantMessageID {
+ t.Fatalf(
+ "reconciled message ID = %q, want %q",
+ reconciled.MessageID,
+ wantMessageID,
+ )
+ }
+
+ // Historical import uses the wall clock. Give the real replay a later clock
+ // so its update satisfies the store's monotonic timestamp constraints.
+ replayNow := time.UnixMilli(reconciled.UpdatedAtMS + 1)
+ liveMessages, err := sqlite.NewMessageRepository(
+ harness.store,
+ func() time.Time { return replayNow },
+ )
+ if err != nil {
+ t.Fatalf("NewMessageRepository(live replay): %v", err)
+ }
+ liveHarness := newSignalWorkerHarnessForStore(
+ t,
+ harness.path,
+ harness.store,
+ liveMessages,
+ )
+ liveHarness.start(t)
+ line := []byte(`{"account":"+15551230000","envelope":{"sourceNumber":"+15551230000","timestamp":1700000006000,"syncMessage":{"sentMessage":{"timestamp":1700000006000,"destinationServiceId":"7a81fd95-20f1-4437-86e2-d5c93ba18851","message":"live outgoing replay"}}}}`)
+ record := mustBuildSignalRecordWithResolutions(t, line, "", "+16505550100")
+ if err := liveHarness.sink.AppendIngress(ctx, record); err != nil {
+ t.Fatalf("AppendIngress(live replay): %v", err)
+ }
+ waitSignalCondition(t, "reconciled Signal outgoing live convergence", func() bool {
+ message, err := liveMessages.GetMessageByRemote(
+ ctx,
+ signalDecoderAccountID,
+ conversation.ConversationID,
+ alias,
+ )
+ return err == nil && message.Body == "live outgoing replay" &&
+ liveHarness.counters.Snapshot(signalDecoderAccountID).Projected == 1
+ })
+
+ converged, err := liveMessages.GetMessageByRemote(
+ ctx,
+ signalDecoderAccountID,
+ conversation.ConversationID,
+ alias,
+ )
+ if err != nil {
+ t.Fatalf("GetMessageByRemote(converged alias): %v", err)
+ }
+ if converged.MessageID != wantMessageID || converged.RemoteMessageID != alias {
+ t.Fatalf(
+ "converged message = %+v, want preserved PK %q and alias %q",
+ converged,
+ wantMessageID,
+ alias,
+ )
+ }
+ bareTimestamp := strconv.FormatInt(timestamp, 10)
+ if _, err := liveMessages.GetMessageByRemote(
+ ctx,
+ signalDecoderAccountID,
+ conversation.ConversationID,
+ bareTimestamp,
+ ); !errors.Is(err, sqlite.ErrNotFound) {
+ t.Fatalf("bare-timestamp lookup error = %v, want ErrNotFound", err)
+ }
+ if got := countSignalRows(t, harness.path, "inbox"); got != 1 {
+ t.Fatalf("inbox rows = %d, want one durable live replay", got)
+ }
+ if got := countSignalRows(t, harness.path, "messages"); got != 1 {
+ t.Fatalf("message rows = %d, want one converged reconciled row", got)
+ }
+}
diff --git a/internal/reconcile/signal.go b/internal/reconcile/signal.go
new file mode 100644
index 00000000..1b62f937
--- /dev/null
+++ b/internal/reconcile/signal.go
@@ -0,0 +1,661 @@
+// Package reconcile repairs gaps between legacy and v2 local stores.
+package reconcile
+
+import (
+ "context"
+ "errors"
+ "fmt"
+ "math"
+ "strconv"
+ "strings"
+ "time"
+
+ "github.com/rs/zerolog"
+
+ "github.com/maxghenis/openmessage/internal/db"
+ "github.com/maxghenis/openmessage/internal/storage/sqlite"
+ "github.com/maxghenis/openmessage/internal/v2keys"
+)
+
+const (
+ signalPlatform = "signal"
+ signalAccountID = "signal-primary"
+ signalBridgeKey = "signal_cli"
+
+ skipMissingConversationID = "missing_conversation_id"
+ skipMissingMessageKey = "missing_message_key"
+ skipNonPositiveTimestamp = "non_positive_timestamp"
+)
+
+// Options supplies the two local stores and bounds one Signal reconciliation.
+type Options struct {
+ Legacy *db.Store
+ V2 *sqlite.Store
+ SinceMS int64
+ DryRun bool
+ Logger zerolog.Logger
+}
+
+// Report is the machine-readable evidence emitted by one reconciliation.
+// During a dry run, MessagesImported and ConversationsCreated are planned
+// writes; no store mutation has occurred.
+type Report struct {
+ DryRun bool `json:"dry_run"`
+ SinceMS int64 `json:"since_ms,omitempty"`
+ ConversationsScanned int `json:"conversations_scanned"`
+ ConversationsCreated int `json:"conversations_created"`
+ MessagesScanned int `json:"messages_scanned"`
+ MessagesImported int `json:"messages_imported"`
+ MessagesAlreadyPresent int `json:"messages_already_present"`
+ MediaDeferred int `json:"media_deferred"`
+ Skipped int `json:"skipped"`
+ SkipReasons map[string]int `json:"skip_reasons"`
+}
+
+// Signal tops up v2 from legacy Signal history without touching transport,
+// credential, or network state. Each write uses the same natural keys as live
+// ingest, so interrupted and repeated runs converge.
+func Signal(ctx context.Context, opts Options) (Report, error) {
+ report := Report{
+ DryRun: opts.DryRun,
+ SinceMS: opts.SinceMS,
+ SkipReasons: make(map[string]int),
+ }
+ if ctx == nil {
+ return report, fmt.Errorf("reconcile Signal: context is nil")
+ }
+ if opts.Legacy == nil {
+ return report, fmt.Errorf("reconcile Signal: legacy store is nil")
+ }
+ if opts.V2 == nil {
+ return report, fmt.Errorf("reconcile Signal: v2 store is nil")
+ }
+ if opts.SinceMS < 0 {
+ return report, fmt.Errorf("reconcile Signal: since timestamp %d is negative", opts.SinceMS)
+ }
+ account, err := opts.V2.GetAccount(signalAccountID)
+ if err != nil {
+ return report, fmt.Errorf("reconcile Signal: resolve Signal account: %w", err)
+ }
+ if account.BridgeKey != signalBridgeKey {
+ return report, fmt.Errorf(
+ "reconcile Signal: account %q has bridge key %q, want %q",
+ signalAccountID,
+ account.BridgeKey,
+ signalBridgeKey,
+ )
+ }
+
+ messages, err := sqlite.NewMessageRepository(opts.V2, time.Now)
+ if err != nil {
+ return report, fmt.Errorf("reconcile Signal: create message repository: %w", err)
+ }
+ legacyConversations, err := opts.Legacy.ListConversationsByPlatform(
+ signalPlatform,
+ math.MaxInt,
+ )
+ if err != nil {
+ return report, fmt.Errorf("reconcile Signal: list legacy conversations: %w", err)
+ }
+
+ nowMS := time.Now().UnixMilli()
+ if nowMS <= 0 {
+ nowMS = 1
+ }
+ knownMessages := make(map[string]struct{})
+
+ for _, legacyConversation := range legacyConversations {
+ if err := ctx.Err(); err != nil {
+ return report, fmt.Errorf("reconcile Signal: %w", err)
+ }
+ report.ConversationsScanned++
+
+ legacyMessages, err := opts.Legacy.GetMessagesByConversationsRange(
+ []string{legacyConversation.ConversationID},
+ opts.SinceMS,
+ 0,
+ math.MaxInt,
+ )
+ if err != nil {
+ return report, fmt.Errorf(
+ "reconcile Signal conversation %q: list legacy messages: %w",
+ legacyConversation.ConversationID,
+ err,
+ )
+ }
+
+ remoteConversationID := v2keys.NormalizeRemoteConversationID(
+ signalPlatform,
+ legacyConversation.ConversationID,
+ )
+ if strings.TrimSpace(remoteConversationID) == "" {
+ for _, legacyMessage := range legacyMessages {
+ report.MessagesScanned++
+ if strings.TrimSpace(legacyMessage.MediaID) != "" {
+ report.MediaDeferred++
+ }
+ report.skip(skipMissingConversationID)
+ }
+ continue
+ }
+
+ conversation, created, err := resolveConversation(
+ opts,
+ legacyConversation,
+ legacyMessages,
+ remoteConversationID,
+ nowMS,
+ )
+ if err != nil {
+ return report, err
+ }
+ if created {
+ report.ConversationsCreated++
+ }
+
+ if signalConversationKind(remoteConversationID) == sqlite.ConversationKindDirect {
+ if err := ensureDirectPeer(
+ opts,
+ conversation,
+ remoteConversationID,
+ legacyConversation.Name,
+ conversation.CreatedAtMS,
+ ); err != nil {
+ return report, fmt.Errorf(
+ "reconcile Signal conversation %q: ensure direct peer: %w",
+ legacyConversation.ConversationID,
+ err,
+ )
+ }
+ }
+
+ remoteByLegacyID := make(map[string]string, len(legacyMessages))
+ for _, legacyMessage := range legacyMessages {
+ remoteByLegacyID[legacyMessage.MessageID] = deriveRemoteMessageID(legacyMessage)
+ }
+ for _, legacyMessage := range legacyMessages {
+ if err := ctx.Err(); err != nil {
+ return report, fmt.Errorf("reconcile Signal: %w", err)
+ }
+ report.MessagesScanned++
+ if strings.TrimSpace(legacyMessage.MediaID) != "" {
+ report.MediaDeferred++
+ }
+ if legacyMessage.TimestampMS <= 0 {
+ report.skip(skipNonPositiveTimestamp)
+ continue
+ }
+
+ remoteMessageID := deriveRemoteMessageID(legacyMessage)
+ if strings.TrimSpace(remoteMessageID) == "" {
+ report.skip(skipMissingMessageKey)
+ continue
+ }
+ naturalKey := conversation.ConversationID + "\x1f" + remoteMessageID
+ alreadyPresent := false
+ for _, candidate := range existingRemoteMessageIDs(
+ legacyMessage,
+ remoteConversationID,
+ remoteMessageID,
+ ) {
+ candidateKey := conversation.ConversationID + "\x1f" + candidate
+ if _, known := knownMessages[candidateKey]; known {
+ alreadyPresent = true
+ break
+ }
+ _, lookupErr := messages.GetMessageByRemote(
+ ctx,
+ signalAccountID,
+ conversation.ConversationID,
+ candidate,
+ )
+ switch {
+ case lookupErr == nil:
+ knownMessages[candidateKey] = struct{}{}
+ alreadyPresent = true
+ case !errors.Is(lookupErr, sqlite.ErrNotFound):
+ return report, fmt.Errorf(
+ "reconcile Signal message %q: check v2 natural key: %w",
+ legacyMessage.MessageID,
+ lookupErr,
+ )
+ }
+ if alreadyPresent {
+ break
+ }
+ }
+ if alreadyPresent {
+ knownMessages[naturalKey] = struct{}{}
+ report.MessagesAlreadyPresent++
+ if !opts.DryRun {
+ if err := opts.V2.BumpConversationRecency(
+ conversation.ConversationID,
+ legacyMessage.TimestampMS,
+ ); err != nil {
+ return report, fmt.Errorf(
+ "reconcile Signal message %q: bump conversation recency: %w",
+ legacyMessage.MessageID,
+ err,
+ )
+ }
+ }
+ continue
+ }
+
+ senderIdentityID, err := resolveInboundSender(
+ opts,
+ legacyMessage,
+ legacyMessage.TimestampMS,
+ )
+ if err != nil {
+ return report, fmt.Errorf(
+ "reconcile Signal message %q: resolve sender: %w",
+ legacyMessage.MessageID,
+ err,
+ )
+ }
+ replyToRemoteID, err := resolveReplyRemoteID(
+ opts.Legacy,
+ legacyMessage,
+ remoteByLegacyID,
+ )
+ if err != nil {
+ return report, fmt.Errorf(
+ "reconcile Signal message %q: resolve reply target: %w",
+ legacyMessage.MessageID,
+ err,
+ )
+ }
+
+ direction := sqlite.MessageDirectionIncoming
+ if legacyMessage.IsFromMe {
+ direction = sqlite.MessageDirectionOutgoing
+ }
+ projection := sqlite.MessageProjection{Message: sqlite.Message{
+ MessageID: v2keys.DeriveID(
+ "message",
+ signalAccountID,
+ remoteConversationID+"\x1f"+remoteMessageID,
+ ),
+ ConversationID: conversation.ConversationID,
+ AccountID: signalAccountID,
+ RemoteMessageID: remoteMessageID,
+ SenderIdentityID: senderIdentityID,
+ Direction: direction,
+ Body: legacyMessage.Body,
+ ReplyToRemoteID: replyToRemoteID,
+ State: sqlite.MessageStateActive,
+ OccurredAtMS: legacyMessage.TimestampMS,
+ }}
+ if !opts.DryRun {
+ if err := messages.ImportMessage(ctx, projection); err != nil {
+ return report, fmt.Errorf(
+ "reconcile Signal message %q: import: %w",
+ legacyMessage.MessageID,
+ err,
+ )
+ }
+ knownMessages[naturalKey] = struct{}{}
+ report.MessagesImported++
+ if err := opts.V2.BumpConversationRecency(
+ conversation.ConversationID,
+ legacyMessage.TimestampMS,
+ ); err != nil {
+ return report, fmt.Errorf(
+ "reconcile Signal message %q: bump conversation recency: %w",
+ legacyMessage.MessageID,
+ err,
+ )
+ }
+ } else {
+ knownMessages[naturalKey] = struct{}{}
+ report.MessagesImported++
+ }
+ }
+ }
+
+ opts.Logger.Debug().
+ Int("conversations_scanned", report.ConversationsScanned).
+ Int("messages_scanned", report.MessagesScanned).
+ Int("messages_imported", report.MessagesImported).
+ Int("messages_already_present", report.MessagesAlreadyPresent).
+ Bool("dry_run", opts.DryRun).
+ Msg("Signal reconciliation complete")
+ return report, nil
+}
+
+func (report *Report) skip(reason string) {
+ report.Skipped++
+ report.SkipReasons[reason]++
+}
+
+func resolveConversation(
+ opts Options,
+ legacyConversation *db.Conversation,
+ legacyMessages []*db.Message,
+ remoteConversationID string,
+ nowMS int64,
+) (sqlite.Conversation, bool, error) {
+ conversation, err := opts.V2.GetConversationByRemote(
+ signalAccountID,
+ remoteConversationID,
+ )
+ if err == nil {
+ kind := signalConversationKind(remoteConversationID)
+ if conversation.Kind == kind {
+ return conversation, false, nil
+ }
+ conversation.Kind = kind
+ conversation.UpdatedAtMS = maxInt64(conversation.UpdatedAtMS, nowMS)
+ if opts.DryRun {
+ return conversation, false, nil
+ }
+ if err := opts.V2.UpsertConversation(conversation); err != nil {
+ return sqlite.Conversation{}, false, fmt.Errorf(
+ "reconcile Signal conversation %q: correct v2 kind: %w",
+ legacyConversation.ConversationID,
+ err,
+ )
+ }
+ effective, err := opts.V2.GetConversationByRemote(
+ signalAccountID,
+ remoteConversationID,
+ )
+ if err != nil {
+ return sqlite.Conversation{}, false, fmt.Errorf(
+ "reconcile Signal conversation %q: reload corrected conversation: %w",
+ legacyConversation.ConversationID,
+ err,
+ )
+ }
+ return effective, false, nil
+ }
+ if !errors.Is(err, sqlite.ErrNotFound) {
+ return sqlite.Conversation{}, false, fmt.Errorf(
+ "reconcile Signal conversation %q: resolve v2 conversation: %w",
+ legacyConversation.ConversationID,
+ err,
+ )
+ }
+
+ createdAtMS, lastMessageAtMS := conversationTimes(
+ legacyConversation,
+ legacyMessages,
+ nowMS,
+ )
+ conversation = sqlite.Conversation{
+ ConversationID: v2keys.DeriveID(
+ "conversation",
+ signalAccountID,
+ remoteConversationID,
+ ),
+ AccountID: signalAccountID,
+ RemoteConversationID: remoteConversationID,
+ Kind: signalConversationKind(remoteConversationID),
+ Title: legacyConversation.Name,
+ NotificationMode: sqlite.NotificationModeAll,
+ LastMessageAtMS: lastMessageAtMS,
+ MetadataJSON: "{}",
+ CreatedAtMS: createdAtMS,
+ UpdatedAtMS: maxInt64(createdAtMS, lastMessageAtMS),
+ }
+ if opts.DryRun {
+ return conversation, true, nil
+ }
+ if err := opts.V2.UpsertConversation(conversation); err != nil {
+ return sqlite.Conversation{}, false, fmt.Errorf(
+ "reconcile Signal conversation %q: create v2 conversation: %w",
+ legacyConversation.ConversationID,
+ err,
+ )
+ }
+ // UpsertConversation preserves the first local primary key on a natural-key
+ // conflict. Resolve that effective row before importing child messages.
+ effective, err := opts.V2.GetConversationByRemote(
+ signalAccountID,
+ remoteConversationID,
+ )
+ if err != nil {
+ return sqlite.Conversation{}, false, fmt.Errorf(
+ "reconcile Signal conversation %q: reload created conversation: %w",
+ legacyConversation.ConversationID,
+ err,
+ )
+ }
+ return effective, true, nil
+}
+
+func signalConversationKind(remoteConversationID string) sqlite.ConversationKind {
+ if strings.HasPrefix(remoteConversationID, "signal-group:") {
+ return sqlite.ConversationKindGroup
+ }
+ return sqlite.ConversationKindDirect
+}
+
+func ensureDirectPeer(
+ opts Options,
+ conversation sqlite.Conversation,
+ remoteConversationID string,
+ displayName string,
+ atMS int64,
+) error {
+ raw := strings.TrimSpace(strings.TrimPrefix(remoteConversationID, "signal:"))
+ key, err := v2keys.IdentityKey(signalAccountID, signalPlatform, raw)
+ if err != nil {
+ return err
+ }
+ identityID, err := resolveIdentity(
+ opts,
+ key,
+ raw,
+ displayName,
+ atMS,
+ )
+ if err != nil {
+ return err
+ }
+ identity, err := opts.V2.GetIdentity(identityID)
+ if err == nil && identity.IsSelf {
+ return nil
+ }
+ if err != nil && !errors.Is(err, sqlite.ErrNotFound) {
+ return err
+ }
+ if opts.DryRun {
+ return nil
+ }
+ _, err = opts.V2.EnsureConversationParticipant(sqlite.ConversationParticipant{
+ AccountID: signalAccountID,
+ ConversationID: conversation.ConversationID,
+ IdentityID: identityID,
+ Role: sqlite.ParticipantRoleMember,
+ DisplayName: strings.TrimSpace(displayName),
+ IsActive: true,
+ })
+ return err
+}
+
+func resolveInboundSender(
+ opts Options,
+ message *db.Message,
+ atMS int64,
+) (*string, error) {
+ if message.IsFromMe {
+ return nil, nil
+ }
+ raw := strings.TrimSpace(message.SenderNumber)
+ // Live ingest preserves inbound messages whose transport supplies no usable
+ // sender, so historical repair degrades to a null sender in the same way.
+ if raw == "" {
+ return nil, nil
+ }
+ key, err := v2keys.IdentityKey(signalAccountID, signalPlatform, raw)
+ if err != nil {
+ return nil, nil
+ }
+ identityID, err := resolveIdentity(
+ opts,
+ key,
+ raw,
+ message.SenderName,
+ atMS,
+ )
+ if err != nil {
+ return nil, err
+ }
+ return &identityID, nil
+}
+
+func resolveIdentity(
+ opts Options,
+ key v2keys.Identity,
+ raw string,
+ displayName string,
+ atMS int64,
+) (string, error) {
+ identity, err := opts.V2.GetIdentityByCanonical(
+ key.AccountID,
+ sqlite.IdentityKind(key.Kind),
+ key.Canonical,
+ )
+ if err == nil {
+ return identity.IdentityID, nil
+ }
+ if !errors.Is(err, sqlite.ErrNotFound) {
+ return "", err
+ }
+
+ identityID := v2keys.DeriveID(
+ "identity",
+ key.AccountID,
+ key.Kind+"\x1f"+key.Canonical,
+ )
+ if opts.DryRun {
+ return identityID, nil
+ }
+ if err := opts.V2.UpsertIdentity(sqlite.Identity{
+ IdentityID: identityID,
+ AccountID: key.AccountID,
+ Kind: sqlite.IdentityKind(key.Kind),
+ CanonicalValue: key.Canonical,
+ RawValue: strings.TrimSpace(raw),
+ DisplayName: strings.TrimSpace(displayName),
+ MetadataJSON: "{}",
+ CreatedAtMS: atMS,
+ UpdatedAtMS: atMS,
+ }); err != nil {
+ return "", err
+ }
+ effective, err := opts.V2.GetIdentityByCanonical(
+ key.AccountID,
+ sqlite.IdentityKind(key.Kind),
+ key.Canonical,
+ )
+ if err != nil {
+ return "", err
+ }
+ return effective.IdentityID, nil
+}
+
+func resolveReplyRemoteID(
+ legacy *db.Store,
+ message *db.Message,
+ remoteByLegacyID map[string]string,
+) (*string, error) {
+ replyID := strings.TrimSpace(message.ReplyToID)
+ if replyID == "" {
+ return nil, nil
+ }
+ remoteID := remoteByLegacyID[replyID]
+ if strings.TrimSpace(remoteID) == "" {
+ target, err := legacy.GetMessageByID(replyID)
+ if err != nil {
+ return nil, err
+ }
+ if target != nil {
+ remoteID = deriveRemoteMessageID(target)
+ }
+ }
+ if strings.TrimSpace(remoteID) == "" {
+ remoteID = stripPlatformPrefix(replyID)
+ }
+ if strings.TrimSpace(remoteID) == "" {
+ return nil, nil
+ }
+ return &remoteID, nil
+}
+
+func deriveRemoteMessageID(message *db.Message) string {
+ if message.SourceID != "" {
+ return message.SourceID
+ }
+ return stripPlatformPrefix(message.MessageID)
+}
+
+// existingRemoteMessageIDs keeps the legacy source ID canonical while
+// recognizing an older live projection that may have retained Signal's bare
+// outgoing timestamp before a local alias existed to converge onto.
+func existingRemoteMessageIDs(
+ message *db.Message,
+ remoteConversationID string,
+ canonical string,
+) []string {
+ keys := []string{canonical}
+ if !message.IsFromMe || message.TimestampMS <= 0 {
+ return keys
+ }
+ if canonical != v2keys.SignalLocalAlias(remoteConversationID, message.TimestampMS) {
+ return keys
+ }
+ bareTimestamp := strconv.FormatInt(message.TimestampMS, 10)
+ if bareTimestamp != canonical {
+ keys = append(keys, bareTimestamp)
+ }
+ return keys
+}
+
+func stripPlatformPrefix(value string) string {
+ for _, prefix := range []string{"whatsapp:", "signal:", "gchat:", "imessage:"} {
+ if strings.HasPrefix(value, prefix) {
+ return strings.TrimPrefix(value, prefix)
+ }
+ }
+ return value
+}
+
+func conversationTimes(
+ conversation *db.Conversation,
+ messages []*db.Message,
+ fallbackMS int64,
+) (int64, int64) {
+ createdAtMS := int64(0)
+ lastMessageAtMS := conversation.LastMessageTS
+ for _, message := range messages {
+ if message.TimestampMS <= 0 {
+ continue
+ }
+ if createdAtMS == 0 || message.TimestampMS < createdAtMS {
+ createdAtMS = message.TimestampMS
+ }
+ if message.TimestampMS > lastMessageAtMS {
+ lastMessageAtMS = message.TimestampMS
+ }
+ }
+ if createdAtMS == 0 && conversation.LastMessageTS > 0 {
+ createdAtMS = conversation.LastMessageTS
+ }
+ if createdAtMS == 0 {
+ createdAtMS = fallbackMS
+ }
+ if lastMessageAtMS < 0 {
+ lastMessageAtMS = 0
+ }
+ return createdAtMS, lastMessageAtMS
+}
+
+func maxInt64(left, right int64) int64 {
+ if left > right {
+ return left
+ }
+ return right
+}
diff --git a/internal/reconcile/signal_test.go b/internal/reconcile/signal_test.go
new file mode 100644
index 00000000..aeacda33
--- /dev/null
+++ b/internal/reconcile/signal_test.go
@@ -0,0 +1,980 @@
+package reconcile
+
+import (
+ "context"
+ "crypto/sha1"
+ "encoding/hex"
+ "errors"
+ "path/filepath"
+ "strconv"
+ "strings"
+ "testing"
+ "time"
+
+ "github.com/rs/zerolog"
+
+ "github.com/maxghenis/openmessage/internal/db"
+ "github.com/maxghenis/openmessage/internal/storage/sqlite"
+ "github.com/maxghenis/openmessage/internal/v2keys"
+)
+
+const (
+ testSignalConversation = "signal:+16505550100"
+ testSignalPeer = "+16505550100"
+ testIncomingTimestamp = int64(1_700_000_001_000)
+ testOutgoingTimestamp = int64(1_700_000_006_000)
+)
+
+func TestSignalLegacySourceIDParity(t *testing.T) {
+ t.Parallel()
+
+ incoming := db.Message{
+ ConversationID: testSignalConversation,
+ SenderNumber: testSignalPeer,
+ TimestampMS: testIncomingTimestamp,
+ SourceID: legacySignalSourceID(
+ testSignalConversation,
+ testSignalPeer,
+ testIncomingTimestamp,
+ false,
+ ),
+ }
+ if got, want := incoming.SourceID, v2keys.SignalIncomingSourceID(
+ incoming.ConversationID,
+ incoming.SenderNumber,
+ incoming.TimestampMS,
+ ); got != want {
+ t.Fatalf("incoming legacy source_id = %q, want v2 key %q", got, want)
+ }
+
+ outgoing := db.Message{
+ ConversationID: testSignalConversation,
+ TimestampMS: testOutgoingTimestamp,
+ IsFromMe: true,
+ SourceID: legacySignalSourceID(
+ testSignalConversation,
+ "me",
+ testOutgoingTimestamp,
+ true,
+ ),
+ }
+ if got, want := outgoing.SourceID, v2keys.SignalLocalAlias(
+ outgoing.ConversationID,
+ outgoing.TimestampMS,
+ ); got != want {
+ t.Fatalf("outgoing legacy source_id = %q, want v2 alias %q", got, want)
+ }
+}
+
+func TestSignalImportsMissingMessagesAndIsIdempotent(t *testing.T) {
+ ctx := context.Background()
+ legacy, v2 := openSignalStores(t)
+ incomingSourceID := v2keys.SignalIncomingSourceID(
+ testSignalConversation,
+ testSignalPeer,
+ testIncomingTimestamp,
+ )
+ outgoingAlias := v2keys.SignalLocalAlias(
+ testSignalConversation,
+ testOutgoingTimestamp,
+ )
+ seedLegacySignalConversation(t, legacy, testOutgoingTimestamp)
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: "signal:legacy-incoming",
+ ConversationID: testSignalConversation,
+ SenderName: "Signal Peer",
+ SenderNumber: testSignalPeer,
+ Body: "missing incoming",
+ TimestampMS: testIncomingTimestamp,
+ MediaID: "legacy-photo-ref",
+ MimeType: "image/jpeg",
+ SourcePlatform: signalPlatform,
+ SourceID: incomingSourceID,
+ })
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: "signal:legacy-outgoing",
+ ConversationID: testSignalConversation,
+ SenderName: "Me",
+ Body: "missing outgoing",
+ TimestampMS: testOutgoingTimestamp,
+ IsFromMe: true,
+ SourcePlatform: signalPlatform,
+ SourceID: outgoingAlias,
+ })
+
+ first, err := Signal(ctx, Options{
+ Legacy: legacy,
+ V2: v2,
+ Logger: zerolog.Nop(),
+ })
+ if err != nil {
+ t.Fatalf("first Signal(): %v", err)
+ }
+ if first.ConversationsScanned != 1 || first.ConversationsCreated != 1 ||
+ first.MessagesScanned != 2 || first.MessagesImported != 2 ||
+ first.MessagesAlreadyPresent != 0 || first.MediaDeferred != 1 ||
+ first.Skipped != 0 {
+ t.Fatalf("first report = %+v", first)
+ }
+
+ conversation, err := v2.GetConversationByRemote(
+ signalAccountID,
+ testSignalConversation,
+ )
+ if err != nil {
+ t.Fatalf("GetConversationByRemote(): %v", err)
+ }
+ wantConversationID := v2keys.DeriveID(
+ "conversation",
+ signalAccountID,
+ testSignalConversation,
+ )
+ if conversation.ConversationID != wantConversationID ||
+ conversation.Kind != sqlite.ConversationKindDirect ||
+ conversation.LastMessageAtMS != testOutgoingTimestamp {
+ t.Fatalf("reconciled conversation = %+v", conversation)
+ }
+ participants, err := v2.ListParticipants(conversation.ConversationID)
+ if err != nil {
+ t.Fatalf("ListParticipants(): %v", err)
+ }
+ if len(participants) != 1 || participants[0].DisplayName != "Signal Peer" {
+ t.Fatalf("direct participants = %+v", participants)
+ }
+
+ repository := newSignalMessageRepository(t, v2)
+ incoming, err := repository.GetMessageByRemote(
+ ctx,
+ signalAccountID,
+ conversation.ConversationID,
+ incomingSourceID,
+ )
+ if err != nil {
+ t.Fatalf("GetMessageByRemote(incoming): %v", err)
+ }
+ if incoming.Body != "missing incoming" ||
+ incoming.Direction != sqlite.MessageDirectionIncoming ||
+ incoming.SenderIdentityID == nil {
+ t.Fatalf("incoming message = %+v", incoming)
+ }
+ // A fresh repair keeps the byte-identical legacy source ID. The worker
+ // resolves a later bare-timestamp delivery onto this local alias.
+ outgoing, err := repository.GetMessageByRemote(
+ ctx,
+ signalAccountID,
+ conversation.ConversationID,
+ outgoingAlias,
+ )
+ if err != nil {
+ t.Fatalf("GetMessageByRemote(outgoing): %v", err)
+ }
+ if outgoing.Body != "missing outgoing" ||
+ outgoing.Direction != sqlite.MessageDirectionOutgoing ||
+ outgoing.SenderIdentityID != nil {
+ t.Fatalf("outgoing message = %+v", outgoing)
+ }
+
+ second, err := Signal(ctx, Options{
+ Legacy: legacy,
+ V2: v2,
+ Logger: zerolog.Nop(),
+ })
+ if err != nil {
+ t.Fatalf("second Signal(): %v", err)
+ }
+ if second.ConversationsCreated != 0 || second.MessagesImported != 0 ||
+ second.MessagesAlreadyPresent != 2 || second.MessagesScanned != 2 {
+ t.Fatalf("second report = %+v", second)
+ }
+ stored, err := repository.ListMessagesByConversation(
+ ctx,
+ conversation.ConversationID,
+ 0,
+ "",
+ 10,
+ )
+ if err != nil {
+ t.Fatalf("ListMessagesByConversation(): %v", err)
+ }
+ if len(stored) != 2 {
+ t.Fatalf("stored messages = %d, want 2: %+v", len(stored), stored)
+ }
+ participants, err = v2.ListParticipants(conversation.ConversationID)
+ if err != nil {
+ t.Fatalf("ListParticipants(second run): %v", err)
+ }
+ if len(participants) != 1 {
+ t.Fatalf("participants after second run = %d, want 1", len(participants))
+ }
+}
+
+func TestSignalCountsExistingTwinWithoutReimporting(t *testing.T) {
+ ctx := context.Background()
+ legacy, v2 := openSignalStores(t)
+ sourceID := v2keys.SignalIncomingSourceID(
+ testSignalConversation,
+ testSignalPeer,
+ testIncomingTimestamp,
+ )
+ seedLegacySignalConversation(t, legacy, testIncomingTimestamp)
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: "signal:legacy-existing",
+ ConversationID: testSignalConversation,
+ SenderName: "Signal Peer",
+ SenderNumber: testSignalPeer,
+ Body: "legacy body must not overwrite",
+ TimestampMS: testIncomingTimestamp,
+ SourcePlatform: signalPlatform,
+ SourceID: sourceID,
+ })
+ seedSignalAccount(t, v2, testIncomingTimestamp)
+ conversationID := seedV2SignalConversation(
+ t,
+ v2,
+ "preexisting-conversation-pk",
+ testIncomingTimestamp-1,
+ )
+ repository := newSignalMessageRepository(t, v2)
+ if err := repository.ImportMessage(ctx, sqlite.MessageProjection{Message: sqlite.Message{
+ MessageID: "preexisting-message-pk",
+ ConversationID: conversationID,
+ AccountID: signalAccountID,
+ RemoteMessageID: sourceID,
+ Direction: sqlite.MessageDirectionIncoming,
+ Body: "v2 body remains",
+ State: sqlite.MessageStateActive,
+ OccurredAtMS: testIncomingTimestamp,
+ }}); err != nil {
+ t.Fatalf("seed v2 message: %v", err)
+ }
+
+ report, err := Signal(ctx, Options{
+ Legacy: legacy,
+ V2: v2,
+ Logger: zerolog.Nop(),
+ })
+ if err != nil {
+ t.Fatalf("Signal(): %v", err)
+ }
+ if report.MessagesImported != 0 || report.MessagesAlreadyPresent != 1 {
+ t.Fatalf("report = %+v", report)
+ }
+ got, err := repository.GetMessageByRemote(
+ ctx,
+ signalAccountID,
+ conversationID,
+ sourceID,
+ )
+ if err != nil {
+ t.Fatalf("GetMessageByRemote(): %v", err)
+ }
+ if got.MessageID != "preexisting-message-pk" || got.Body != "v2 body remains" {
+ t.Fatalf("existing twin was rewritten: %+v", got)
+ }
+ stored, err := repository.ListMessagesByConversation(ctx, conversationID, 0, "", 10)
+ if err != nil {
+ t.Fatalf("ListMessagesByConversation(): %v", err)
+ }
+ if len(stored) != 1 {
+ t.Fatalf("stored messages = %d, want 1", len(stored))
+ }
+}
+
+func TestSignalDryRunWritesNothing(t *testing.T) {
+ ctx := context.Background()
+ legacy, v2 := openSignalStores(t)
+ sourceID := v2keys.SignalIncomingSourceID(
+ testSignalConversation,
+ testSignalPeer,
+ testIncomingTimestamp,
+ )
+ seedLegacySignalConversation(t, legacy, testIncomingTimestamp)
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: "signal:legacy-dry-run",
+ ConversationID: testSignalConversation,
+ SenderName: "Signal Peer",
+ SenderNumber: testSignalPeer,
+ Body: "[Photo]",
+ TimestampMS: testIncomingTimestamp,
+ MediaID: "legacy-media",
+ SourcePlatform: signalPlatform,
+ SourceID: sourceID,
+ })
+ conversationID := v2keys.DeriveID(
+ "conversation",
+ signalAccountID,
+ testSignalConversation,
+ )
+ repository := newSignalMessageRepository(t, v2)
+ before, err := repository.ListMessagesByConversation(ctx, conversationID, 0, "", 10)
+ if err != nil {
+ t.Fatalf("list messages before dry run: %v", err)
+ }
+
+ report, err := Signal(ctx, Options{
+ Legacy: legacy,
+ V2: v2,
+ DryRun: true,
+ Logger: zerolog.Nop(),
+ })
+ if err != nil {
+ t.Fatalf("Signal(dry run): %v", err)
+ }
+ if !report.DryRun || report.ConversationsCreated != 1 ||
+ report.MessagesImported != 1 || report.MediaDeferred != 1 {
+ t.Fatalf("dry-run report = %+v", report)
+ }
+ after, err := repository.ListMessagesByConversation(ctx, conversationID, 0, "", 10)
+ if err != nil {
+ t.Fatalf("list messages after dry run: %v", err)
+ }
+ if len(after) != len(before) {
+ t.Fatalf("dry run changed message row count: before=%d after=%d", len(before), len(after))
+ }
+ if _, err := v2.GetConversationByRemote(
+ signalAccountID,
+ testSignalConversation,
+ ); !errors.Is(err, sqlite.ErrNotFound) {
+ t.Fatalf("dry run conversation lookup error = %v, want ErrNotFound", err)
+ }
+ accounts, err := v2.ListAccounts()
+ if err != nil {
+ t.Fatalf("ListAccounts(): %v", err)
+ }
+ if len(accounts) != 1 || accounts[0].AccountID != signalAccountID {
+ t.Fatalf("dry run changed accounts: %+v", accounts)
+ }
+ participants, err := v2.ListParticipants(conversationID)
+ if err != nil {
+ t.Fatalf("ListParticipants(): %v", err)
+ }
+ if len(participants) != 0 {
+ t.Fatalf("dry run created participants: %+v", participants)
+ }
+}
+
+func TestSignalRequiresExistingSignalAccount(t *testing.T) {
+ legacy, v2 := openSignalStoresWithoutAccount(t)
+
+ report, err := Signal(context.Background(), Options{
+ Legacy: legacy,
+ V2: v2,
+ Logger: zerolog.Nop(),
+ })
+ if !errors.Is(err, sqlite.ErrNotFound) {
+ t.Fatalf("Signal() error = %v, want wrapped ErrNotFound", err)
+ }
+ if report.ConversationsScanned != 0 || report.MessagesImported != 0 {
+ t.Fatalf("report after missing account = %+v", report)
+ }
+ accounts, listErr := v2.ListAccounts()
+ if listErr != nil {
+ t.Fatalf("ListAccounts(): %v", listErr)
+ }
+ if len(accounts) != 0 {
+ t.Fatalf("reconcile created transport account: %+v", accounts)
+ }
+}
+
+func TestSignalRejectsWrongAccountBridge(t *testing.T) {
+ legacy, v2 := openSignalStoresWithoutAccount(t)
+ if err := v2.UpsertAccount(sqlite.Account{
+ AccountID: signalAccountID,
+ BridgeKey: "not_signal",
+ Mode: sqlite.AccountModeLive,
+ Enabled: true,
+ ConfigJSON: "{}",
+ CreatedAtMS: testIncomingTimestamp,
+ UpdatedAtMS: testIncomingTimestamp,
+ }); err != nil {
+ t.Fatalf("UpsertAccount(): %v", err)
+ }
+
+ _, err := Signal(context.Background(), Options{
+ Legacy: legacy,
+ V2: v2,
+ Logger: zerolog.Nop(),
+ })
+ if err == nil || !strings.Contains(err.Error(), `want "signal_cli"`) {
+ t.Fatalf("Signal() error = %v, want bridge-key mismatch", err)
+ }
+}
+
+func TestSignalGroupPrefixCorrectsKindAndDegradesUnusableSenders(t *testing.T) {
+ const (
+ groupRemoteID = "signal-group:ZmFrZS1ncm91cA=="
+ nameOnlyMessage = "group-name-only"
+ invalidPeerMessage = "group-invalid-peer"
+ )
+ ctx := context.Background()
+ legacy, v2 := openSignalStores(t)
+ if err := legacy.UpsertConversation(&db.Conversation{
+ ConversationID: groupRemoteID,
+ Name: "Signal Group",
+ LastMessageTS: testIncomingTimestamp,
+ SourcePlatform: signalPlatform,
+ }); err != nil {
+ t.Fatalf("UpsertConversation(): %v", err)
+ }
+ sourceID := v2keys.SignalIncomingSourceID(
+ groupRemoteID,
+ "transport-source-without-number",
+ testIncomingTimestamp,
+ )
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: nameOnlyMessage,
+ ConversationID: groupRemoteID,
+ SenderName: "Display Name Only",
+ Body: "group history",
+ TimestampMS: testIncomingTimestamp,
+ SourcePlatform: signalPlatform,
+ SourceID: sourceID,
+ })
+ invalidPeerSourceID := v2keys.SignalIncomingSourceID(
+ groupRemoteID,
+ "+",
+ testOutgoingTimestamp,
+ )
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: invalidPeerMessage,
+ ConversationID: groupRemoteID,
+ SenderName: "Malformed Peer",
+ SenderNumber: "+",
+ Body: "valid history with unusable sender",
+ TimestampMS: testOutgoingTimestamp,
+ SourcePlatform: signalPlatform,
+ SourceID: invalidPeerSourceID,
+ })
+ const existingConversationID = "preexisting-misclassified-group"
+ if err := v2.UpsertConversation(sqlite.Conversation{
+ ConversationID: existingConversationID,
+ AccountID: signalAccountID,
+ RemoteConversationID: groupRemoteID,
+ Kind: sqlite.ConversationKindDirect,
+ Title: "Existing Signal Group",
+ NotificationMode: sqlite.NotificationModeAll,
+ LastMessageAtMS: testIncomingTimestamp,
+ MetadataJSON: "{}",
+ CreatedAtMS: testIncomingTimestamp,
+ UpdatedAtMS: testIncomingTimestamp,
+ }); err != nil {
+ t.Fatalf("UpsertConversation(): %v", err)
+ }
+
+ report, err := Signal(ctx, Options{
+ Legacy: legacy,
+ V2: v2,
+ Logger: zerolog.Nop(),
+ })
+ if err != nil {
+ t.Fatalf("Signal(): %v", err)
+ }
+ if report.MessagesImported != 2 || report.Skipped != 0 {
+ t.Fatalf("report = %+v", report)
+ }
+ corrected, err := v2.GetConversation(existingConversationID)
+ if err != nil {
+ t.Fatalf("GetConversation(): %v", err)
+ }
+ if corrected.Kind != sqlite.ConversationKindGroup {
+ t.Fatalf("corrected conversation kind = %q, want group", corrected.Kind)
+ }
+ repository := newSignalMessageRepository(t, v2)
+ for _, remoteMessageID := range []string{sourceID, invalidPeerSourceID} {
+ message, err := repository.GetMessageByRemote(
+ ctx,
+ signalAccountID,
+ existingConversationID,
+ remoteMessageID,
+ )
+ if err != nil {
+ t.Fatalf("GetMessageByRemote(%q): %v", remoteMessageID, err)
+ }
+ if message.SenderIdentityID != nil {
+ t.Fatalf("unusable sender created identity link: %+v", message)
+ }
+ }
+ participants, err := v2.ListParticipants(existingConversationID)
+ if err != nil {
+ t.Fatalf("ListParticipants(): %v", err)
+ }
+ if len(participants) != 0 {
+ t.Fatalf("group remote ID gained direct peer participant: %+v", participants)
+ }
+ identities, err := v2.ListIdentities(signalAccountID)
+ if err != nil {
+ t.Fatalf("ListIdentities(): %v", err)
+ }
+ if len(identities) != 0 {
+ t.Fatalf("name-only sender created identities: %+v", identities)
+ }
+}
+
+func TestSignalCountsDeferredMediaWhenConversationIDIsMissing(t *testing.T) {
+ legacy, v2 := openSignalStores(t)
+ if err := legacy.UpsertConversation(&db.Conversation{
+ ConversationID: "",
+ Name: "Malformed Signal Thread",
+ LastMessageTS: testIncomingTimestamp,
+ SourcePlatform: signalPlatform,
+ }); err != nil {
+ t.Fatalf("UpsertConversation(): %v", err)
+ }
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: "signal:missing-conversation",
+ ConversationID: "",
+ Body: "[Photo]",
+ TimestampMS: testIncomingTimestamp,
+ MediaID: "legacy-media-ref",
+ SourcePlatform: signalPlatform,
+ SourceID: "malformed-thread-message",
+ })
+
+ report, err := Signal(context.Background(), Options{
+ Legacy: legacy,
+ V2: v2,
+ Logger: zerolog.Nop(),
+ })
+ if err != nil {
+ t.Fatalf("Signal(): %v", err)
+ }
+ if report.MessagesScanned != 1 || report.MediaDeferred != 1 ||
+ report.MessagesImported != 0 || report.Skipped != 1 ||
+ report.SkipReasons[skipMissingConversationID] != 1 {
+ t.Fatalf("report = %+v", report)
+ }
+}
+
+func TestSignalRespectsSinceTimestamp(t *testing.T) {
+ ctx := context.Background()
+ legacy, v2 := openSignalStores(t)
+ seedLegacySignalConversation(t, legacy, testOutgoingTimestamp)
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: "signal:before-since",
+ ConversationID: testSignalConversation,
+ SenderNumber: testSignalPeer,
+ Body: "too old",
+ TimestampMS: testIncomingTimestamp,
+ SourcePlatform: signalPlatform,
+ SourceID: v2keys.SignalIncomingSourceID(
+ testSignalConversation,
+ testSignalPeer,
+ testIncomingTimestamp,
+ ),
+ })
+ afterSourceID := v2keys.SignalIncomingSourceID(
+ testSignalConversation,
+ testSignalPeer,
+ testOutgoingTimestamp,
+ )
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: "signal:at-since",
+ ConversationID: testSignalConversation,
+ SenderNumber: testSignalPeer,
+ Body: "included",
+ TimestampMS: testOutgoingTimestamp,
+ SourcePlatform: signalPlatform,
+ SourceID: afterSourceID,
+ })
+
+ report, err := Signal(ctx, Options{
+ Legacy: legacy,
+ V2: v2,
+ SinceMS: testOutgoingTimestamp,
+ Logger: zerolog.Nop(),
+ })
+ if err != nil {
+ t.Fatalf("Signal(): %v", err)
+ }
+ if report.MessagesScanned != 1 || report.MessagesImported != 1 {
+ t.Fatalf("report = %+v", report)
+ }
+ conversation, err := v2.GetConversationByRemote(
+ signalAccountID,
+ testSignalConversation,
+ )
+ if err != nil {
+ t.Fatalf("GetConversationByRemote(): %v", err)
+ }
+ repository := newSignalMessageRepository(t, v2)
+ stored, err := repository.ListMessagesByConversation(
+ ctx,
+ conversation.ConversationID,
+ 0,
+ "",
+ 10,
+ )
+ if err != nil {
+ t.Fatalf("ListMessagesByConversation(): %v", err)
+ }
+ if len(stored) != 1 || stored[0].RemoteMessageID != afterSourceID {
+ t.Fatalf("stored messages = %+v", stored)
+ }
+}
+
+func TestSignalFallbackKeyAndSkipReasons(t *testing.T) {
+ ctx := context.Background()
+ legacy, v2 := openSignalStores(t)
+ seedLegacySignalConversation(t, legacy, testOutgoingTimestamp)
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: "signal:zero-timestamp",
+ ConversationID: testSignalConversation,
+ Body: "invalid timestamp",
+ TimestampMS: 0,
+ IsFromMe: true,
+ SourcePlatform: signalPlatform,
+ SourceID: "key-with-invalid-timestamp",
+ })
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: "signal:",
+ ConversationID: testSignalConversation,
+ Body: "missing remote key",
+ TimestampMS: testIncomingTimestamp,
+ IsFromMe: true,
+ SourcePlatform: signalPlatform,
+ })
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: "signal:fallback-remote-key",
+ ConversationID: testSignalConversation,
+ Body: "fallback source",
+ TimestampMS: testOutgoingTimestamp,
+ IsFromMe: true,
+ SourcePlatform: signalPlatform,
+ })
+
+ report, err := Signal(ctx, Options{
+ Legacy: legacy,
+ V2: v2,
+ Logger: zerolog.Nop(),
+ })
+ if err != nil {
+ t.Fatalf("Signal(): %v", err)
+ }
+ if report.MessagesScanned != 3 || report.MessagesImported != 1 ||
+ report.Skipped != 2 ||
+ report.SkipReasons[skipNonPositiveTimestamp] != 1 ||
+ report.SkipReasons[skipMissingMessageKey] != 1 {
+ t.Fatalf("report = %+v", report)
+ }
+ conversation, err := v2.GetConversationByRemote(
+ signalAccountID,
+ testSignalConversation,
+ )
+ if err != nil {
+ t.Fatalf("GetConversationByRemote(): %v", err)
+ }
+ repository := newSignalMessageRepository(t, v2)
+ message, err := repository.GetMessageByRemote(
+ ctx,
+ signalAccountID,
+ conversation.ConversationID,
+ "fallback-remote-key",
+ )
+ if err != nil {
+ t.Fatalf("GetMessageByRemote(fallback): %v", err)
+ }
+ if message.Body != "fallback source" {
+ t.Fatalf("fallback message = %+v", message)
+ }
+}
+
+func TestSignalDoesNotLinkSelfIdentityAsDirectPeer(t *testing.T) {
+ legacy, v2 := openSignalStores(t)
+ seedLegacySignalConversation(t, legacy, testIncomingTimestamp)
+ key, err := v2keys.IdentityKey(signalAccountID, signalPlatform, testSignalPeer)
+ if err != nil {
+ t.Fatalf("IdentityKey(): %v", err)
+ }
+ identityID := v2keys.DeriveID(
+ "identity",
+ signalAccountID,
+ key.Kind+"\x1f"+key.Canonical,
+ )
+ if err := v2.UpsertIdentity(sqlite.Identity{
+ IdentityID: identityID,
+ AccountID: signalAccountID,
+ Kind: sqlite.IdentityKind(key.Kind),
+ CanonicalValue: key.Canonical,
+ RawValue: testSignalPeer,
+ DisplayName: "Self",
+ IsSelf: true,
+ MetadataJSON: "{}",
+ CreatedAtMS: testIncomingTimestamp,
+ UpdatedAtMS: testIncomingTimestamp,
+ }); err != nil {
+ t.Fatalf("UpsertIdentity(): %v", err)
+ }
+
+ if _, err := Signal(context.Background(), Options{
+ Legacy: legacy,
+ V2: v2,
+ Logger: zerolog.Nop(),
+ }); err != nil {
+ t.Fatalf("Signal(): %v", err)
+ }
+ conversation, err := v2.GetConversationByRemote(
+ signalAccountID,
+ testSignalConversation,
+ )
+ if err != nil {
+ t.Fatalf("GetConversationByRemote(): %v", err)
+ }
+ participants, err := v2.ListParticipants(conversation.ConversationID)
+ if err != nil {
+ t.Fatalf("ListParticipants(): %v", err)
+ }
+ if len(participants) != 0 {
+ t.Fatalf("self identity was linked as direct peer: %+v", participants)
+ }
+}
+
+func legacySignalSourceID(
+ conversationID string,
+ actor string,
+ timestampMS int64,
+ outgoing bool,
+) string {
+ sum := sha1.Sum([]byte(strings.Join([]string{
+ conversationID,
+ actor,
+ strconv.FormatInt(timestampMS, 10),
+ }, "\x1f")))
+ encoded := hex.EncodeToString(sum[:])
+ if outgoing {
+ return "local:" + encoded
+ }
+ return encoded
+}
+
+func openSignalStores(t *testing.T) (*db.Store, *sqlite.Store) {
+ t.Helper()
+ legacy, v2 := openSignalStoresWithoutAccount(t)
+ seedSignalAccount(t, v2, testIncomingTimestamp)
+ return legacy, v2
+}
+
+func openSignalStoresWithoutAccount(t *testing.T) (*db.Store, *sqlite.Store) {
+ t.Helper()
+ legacy, err := db.New(filepath.Join(t.TempDir(), "messages.db"))
+ if err != nil {
+ t.Fatalf("db.New(): %v", err)
+ }
+ t.Cleanup(func() {
+ if err := legacy.Close(); err != nil {
+ t.Errorf("legacy.Close(): %v", err)
+ }
+ })
+ v2, err := sqlite.Open(filepath.Join(t.TempDir(), "store.sqlite3"))
+ if err != nil {
+ t.Fatalf("sqlite.Open(): %v", err)
+ }
+ t.Cleanup(func() {
+ if err := v2.Close(); err != nil {
+ t.Errorf("v2.Close(): %v", err)
+ }
+ })
+ return legacy, v2
+}
+
+func seedLegacySignalConversation(t *testing.T, legacy *db.Store, lastMessageMS int64) {
+ t.Helper()
+ if err := legacy.UpsertConversation(&db.Conversation{
+ ConversationID: testSignalConversation,
+ Name: "Signal Peer",
+ LastMessageTS: lastMessageMS,
+ SourcePlatform: signalPlatform,
+ }); err != nil {
+ t.Fatalf("UpsertConversation(): %v", err)
+ }
+}
+
+func seedLegacySignalMessage(t *testing.T, legacy *db.Store, message db.Message) {
+ t.Helper()
+ if err := legacy.UpsertMessage(&message); err != nil {
+ t.Fatalf("UpsertMessage(%q): %v", message.MessageID, err)
+ }
+}
+
+func seedSignalAccount(t *testing.T, v2 *sqlite.Store, atMS int64) {
+ t.Helper()
+ if err := v2.UpsertAccount(sqlite.Account{
+ AccountID: signalAccountID,
+ BridgeKey: signalBridgeKey,
+ DisplayName: "Signal",
+ Mode: sqlite.AccountModeLive,
+ Enabled: true,
+ ConfigJSON: "{}",
+ CreatedAtMS: atMS,
+ UpdatedAtMS: atMS,
+ }); err != nil {
+ t.Fatalf("UpsertAccount(): %v", err)
+ }
+}
+
+func seedV2SignalConversation(
+ t *testing.T,
+ v2 *sqlite.Store,
+ conversationID string,
+ atMS int64,
+) string {
+ t.Helper()
+ if err := v2.UpsertConversation(sqlite.Conversation{
+ ConversationID: conversationID,
+ AccountID: signalAccountID,
+ RemoteConversationID: testSignalConversation,
+ Kind: sqlite.ConversationKindDirect,
+ Title: "Existing Signal Peer",
+ NotificationMode: sqlite.NotificationModeAll,
+ LastMessageAtMS: atMS,
+ MetadataJSON: "{}",
+ CreatedAtMS: atMS,
+ UpdatedAtMS: atMS,
+ }); err != nil {
+ t.Fatalf("UpsertConversation(): %v", err)
+ }
+ return conversationID
+}
+
+func newSignalMessageRepository(
+ t *testing.T,
+ v2 *sqlite.Store,
+) *sqlite.MessageRepository {
+ t.Helper()
+ repository, err := sqlite.NewMessageRepository(
+ v2,
+ func() time.Time { return time.UnixMilli(2_000_000_000_000) },
+ )
+ if err != nil {
+ t.Fatalf("NewMessageRepository(): %v", err)
+ }
+ return repository
+}
+
+// The duplication defect this reconciler must never reintroduce, reproduced
+// exactly as it appears on a real store: the live decoder keyed an outgoing
+// message by its bare send timestamp, while legacy keys the same message
+// "local:". Deduping on the legacy key alone misses the v2 row and
+// re-imports the user's own message as a duplicate — 50 of 50 outgoing rows on
+// the thread that prompted this work.
+func TestSignalDoesNotDuplicateOutgoingMessageAlreadyKeyedByTimestamp(t *testing.T) {
+ ctx := context.Background()
+ legacy, v2 := openSignalStores(t)
+ seedSignalAccount(t, v2, testIncomingTimestamp)
+ conversationID := seedV2SignalConversation(
+ t,
+ v2,
+ v2keys.DeriveID("conversation", signalAccountID, testSignalConversation),
+ testOutgoingTimestamp,
+ )
+
+ // v2 already holds the outgoing message under the live decoder's key.
+ bareKey := strconv.FormatInt(testOutgoingTimestamp, 10)
+ repository := newSignalMessageRepository(t, v2)
+ if err := repository.ImportMessage(ctx, sqlite.MessageProjection{
+ Message: sqlite.Message{
+ MessageID: v2keys.DeriveID(
+ "message",
+ signalAccountID,
+ testSignalConversation+"\x1f"+bareKey,
+ ),
+ ConversationID: conversationID,
+ AccountID: signalAccountID,
+ RemoteMessageID: bareKey,
+ Direction: sqlite.MessageDirectionOutgoing,
+ Body: "already projected by the live decoder",
+ State: sqlite.MessageStateActive,
+ OccurredAtMS: testOutgoingTimestamp,
+ },
+ }); err != nil {
+ t.Fatalf("seed live-decoder message: %v", err)
+ }
+
+ // Legacy holds the same message under the alias form.
+ seedLegacySignalConversation(t, legacy, testOutgoingTimestamp)
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: "signal:legacy-outgoing",
+ ConversationID: testSignalConversation,
+ TimestampMS: testOutgoingTimestamp,
+ IsFromMe: true,
+ Body: "already projected by the live decoder",
+ SourcePlatform: signalPlatform,
+ SourceID: v2keys.SignalLocalAlias(
+ testSignalConversation,
+ testOutgoingTimestamp,
+ ),
+ })
+
+ report, err := Signal(ctx, Options{Legacy: legacy, V2: v2, Logger: zerolog.Nop()})
+ if err != nil {
+ t.Fatalf("Signal(): %v", err)
+ }
+ if report.MessagesImported != 0 || report.MessagesAlreadyPresent != 1 {
+ t.Fatalf("report = %+v, want the message recognized as already present", report)
+ }
+ stored, err := repository.ListMessagesByConversation(ctx, conversationID, 0, "", 100)
+ if err != nil {
+ t.Fatalf("ListMessagesByConversation(): %v", err)
+ }
+ if len(stored) != 1 {
+ t.Fatalf("stored messages = %d (%+v), want exactly one — no duplicate", len(stored), stored)
+ }
+}
+
+// The mirror case: v2 holds the message under the migration-era "local:" alias
+// (what R2 wrote), and legacy has the same alias. Still one row.
+func TestSignalDoesNotDuplicateOutgoingMessageAlreadyKeyedByAlias(t *testing.T) {
+ ctx := context.Background()
+ legacy, v2 := openSignalStores(t)
+ seedSignalAccount(t, v2, testIncomingTimestamp)
+ conversationID := seedV2SignalConversation(
+ t,
+ v2,
+ v2keys.DeriveID("conversation", signalAccountID, testSignalConversation),
+ testOutgoingTimestamp,
+ )
+ alias := v2keys.SignalLocalAlias(testSignalConversation, testOutgoingTimestamp)
+ repository := newSignalMessageRepository(t, v2)
+ if err := repository.ImportMessage(ctx, sqlite.MessageProjection{
+ Message: sqlite.Message{
+ MessageID: v2keys.DeriveID(
+ "message",
+ signalAccountID,
+ testSignalConversation+"\x1f"+alias,
+ ),
+ ConversationID: conversationID,
+ AccountID: signalAccountID,
+ RemoteMessageID: alias,
+ Direction: sqlite.MessageDirectionOutgoing,
+ Body: "migrated by R2",
+ State: sqlite.MessageStateActive,
+ OccurredAtMS: testOutgoingTimestamp,
+ },
+ }); err != nil {
+ t.Fatalf("seed migrated message: %v", err)
+ }
+ seedLegacySignalConversation(t, legacy, testOutgoingTimestamp)
+ seedLegacySignalMessage(t, legacy, db.Message{
+ MessageID: "signal:legacy-outgoing-alias",
+ ConversationID: testSignalConversation,
+ TimestampMS: testOutgoingTimestamp,
+ IsFromMe: true,
+ Body: "migrated by R2",
+ SourcePlatform: signalPlatform,
+ SourceID: alias,
+ })
+
+ report, err := Signal(ctx, Options{Legacy: legacy, V2: v2, Logger: zerolog.Nop()})
+ if err != nil {
+ t.Fatalf("Signal(): %v", err)
+ }
+ if report.MessagesImported != 0 || report.MessagesAlreadyPresent != 1 {
+ t.Fatalf("report = %+v, want the alias row recognized", report)
+ }
+ stored, err := repository.ListMessagesByConversation(ctx, conversationID, 0, "", 100)
+ if err != nil {
+ t.Fatalf("ListMessagesByConversation(): %v", err)
+ }
+ if len(stored) != 1 {
+ t.Fatalf("stored messages = %d, want exactly one — no duplicate", len(stored))
+ }
+}
diff --git a/internal/storage/sqlite/readonly_test.go b/internal/storage/sqlite/readonly_test.go
new file mode 100644
index 00000000..1a3b994c
--- /dev/null
+++ b/internal/storage/sqlite/readonly_test.go
@@ -0,0 +1,79 @@
+package sqlite
+
+import (
+ "bytes"
+ "os"
+ "path/filepath"
+ "testing"
+)
+
+func TestOpenReadOnlyReadsWithoutMutatingV2Store(t *testing.T) {
+ path := filepath.Join(t.TempDir(), "store.sqlite3")
+ writable, err := Open(path)
+ if err != nil {
+ t.Fatalf("Open(): %v", err)
+ }
+ if err := writable.UpsertAccount(Account{
+ AccountID: "signal-primary",
+ BridgeKey: "signal_cli",
+ DisplayName: "Signal",
+ Mode: AccountModeLive,
+ Enabled: true,
+ ConfigJSON: "{}",
+ CreatedAtMS: 1_700_000_001_000,
+ UpdatedAtMS: 1_700_000_001_000,
+ }); err != nil {
+ t.Fatalf("UpsertAccount(): %v", err)
+ }
+ if err := writable.Close(); err != nil {
+ t.Fatalf("Close(writable): %v", err)
+ }
+ before, err := os.ReadFile(path)
+ if err != nil {
+ t.Fatal(err)
+ }
+
+ readOnly, err := OpenReadOnly(path)
+ if err != nil {
+ t.Fatalf("OpenReadOnly(): %v", err)
+ }
+ account, err := readOnly.GetAccount("signal-primary")
+ if err != nil {
+ t.Fatalf("GetAccount(): %v", err)
+ }
+ if account.BridgeKey != "signal_cli" {
+ t.Fatalf("read-only account = %+v", account)
+ }
+ if err := readOnly.UpsertAccount(Account{
+ AccountID: "must-not-write",
+ BridgeKey: "signal_cli",
+ Mode: AccountModeLive,
+ Enabled: true,
+ ConfigJSON: "{}",
+ CreatedAtMS: 1,
+ UpdatedAtMS: 1,
+ }); err == nil {
+ t.Fatal("read-only store accepted a write")
+ }
+ if err := readOnly.Close(); err != nil {
+ t.Fatalf("Close(read-only): %v", err)
+ }
+
+ after, err := os.ReadFile(path)
+ if err != nil {
+ t.Fatal(err)
+ }
+ if !bytes.Equal(before, after) {
+ t.Fatal("OpenReadOnly changed v2 database bytes")
+ }
+}
+
+func TestOpenReadOnlyDoesNotCreateMissingV2Store(t *testing.T) {
+ path := filepath.Join(t.TempDir(), "missing.sqlite3")
+ if _, err := OpenReadOnly(path); err == nil {
+ t.Fatal("OpenReadOnly() succeeded for a missing database")
+ }
+ if _, err := os.Stat(path); !os.IsNotExist(err) {
+ t.Fatalf("missing path was created: %v", err)
+ }
+}
diff --git a/internal/storage/sqlite/store.go b/internal/storage/sqlite/store.go
index 97ade349..2a9f81d0 100644
--- a/internal/storage/sqlite/store.go
+++ b/internal/storage/sqlite/store.go
@@ -60,6 +60,56 @@ func Open(path string) (*Store, error) {
return &Store{db: db}, nil
}
+// OpenReadOnly opens an existing store without changing journal mode or
+// running schema migrations. The read-only DSN and query_only pragma provide
+// independent guards against accidental writes in maintenance dry runs.
+func OpenReadOnly(path string) (*Store, error) {
+ if strings.TrimSpace(path) == "" {
+ return nil, fmt.Errorf("open read-only sqlite store: path is empty")
+ }
+ db, err := sql.Open("sqlite", readOnlyStoreDSN(path))
+ if err != nil {
+ return nil, fmt.Errorf("open read-only sqlite store: %w", err)
+ }
+ db.SetMaxOpenConns(1)
+ closeWithError := func(openErr error) (*Store, error) {
+ _ = db.Close()
+ return nil, openErr
+ }
+
+ ctx := context.Background()
+ if err := db.PingContext(ctx); err != nil {
+ return closeWithError(fmt.Errorf("open read-only sqlite store connection: %w", err))
+ }
+ checks := []struct {
+ pragma string
+ want int
+ }{
+ {pragma: "query_only", want: 1},
+ {pragma: "foreign_keys", want: 1},
+ {pragma: "busy_timeout", want: busyTimeoutMS},
+ }
+ for _, check := range checks {
+ var got int
+ if err := db.QueryRowContext(ctx, "PRAGMA "+check.pragma).Scan(&got); err != nil {
+ return closeWithError(fmt.Errorf(
+ "verify read-only sqlite %s pragma: %w",
+ check.pragma,
+ err,
+ ))
+ }
+ if got != check.want {
+ return closeWithError(fmt.Errorf(
+ "verify read-only sqlite %s pragma: got %d, want %d",
+ check.pragma,
+ got,
+ check.want,
+ ))
+ }
+ }
+ return &Store{db: db}, nil
+}
+
// Close closes the store's database connections.
func (s *Store) Close() error {
return s.db.Close()
@@ -100,6 +150,24 @@ func storeDSN(path string) string {
}).String()
}
+func readOnlyStoreDSN(path string) string {
+ normalizedPath := strings.ReplaceAll(path, `\`, "/")
+ if isWindowsAbsolutePath(normalizedPath) {
+ normalizedPath = "/" + normalizedPath
+ }
+
+ query := make(url.Values)
+ query.Set("mode", "ro")
+ query.Add("_pragma", fmt.Sprintf("busy_timeout(%d)", busyTimeoutMS))
+ query.Add("_pragma", "foreign_keys(ON)")
+ query.Add("_pragma", "query_only(ON)")
+ return (&url.URL{
+ Scheme: "file",
+ Path: normalizedPath,
+ RawQuery: query.Encode(),
+ }).String()
+}
+
func enableWAL(ctx context.Context, db *sql.DB) error {
var journalMode string
if err := db.QueryRowContext(ctx, `PRAGMA journal_mode = WAL`).Scan(&journalMode); err != nil {
diff --git a/main.go b/main.go
index efc4342d..793e9f57 100644
--- a/main.go
+++ b/main.go
@@ -23,13 +23,14 @@ func main() {
With().Timestamp().Logger().Level(level)
if len(os.Args) < 2 {
- fmt.Fprintln(os.Stderr, "Usage: openmessage ")
+ fmt.Fprintln(os.Stderr, "Usage: openmessage ")
fmt.Fprintln(os.Stderr, " pair [--google|--google-file path] - Pair with your phone via QR or Google account cookies")
fmt.Fprintln(os.Stderr, " serve [--demo] [--web|--no-web] [--mcp-sse|--no-mcp-sse] [--mcp-stdio] - Start explicit web/MCP transports")
fmt.Fprintln(os.Stderr, " demo - Start a seeded fake-data UI with live transports disabled")
fmt.Fprintln(os.Stderr, " backup [--to dir] [--json] - Create a verified legacy migration backup and manifest")
fmt.Fprintln(os.Stderr, " migrate [--check] [--from dir] [--to dir] [--json] - Transform the legacy store into a validated v2 store")
fmt.Fprintln(os.Stderr, " repair google-idspace --since [--apply] [--json] - Re-file v2 rows misrouted by a Google device id-space reset (dry run by default)")
+ fmt.Fprintln(os.Stderr, " v2 reconcile-signal [--from dir] [--since YYYY-MM-DD] [--dry-run] [--json] - Top up v2 from legacy Signal history")
fmt.Fprintln(os.Stderr, " read [--limit N] [--phone X] [--since YYYY-MM-DD] [--until YYYY-MM-DD] [--json] - Search the local store")
fmt.Fprintln(os.Stderr, " thread [--limit N] [--since YYYY-MM-DD] [--until YYYY-MM-DD] [--json] - Print a full conversation chronologically")
fmt.Fprintln(os.Stderr, " threads [--limit N] [--json] - List recent conversations (find an id/name for thread)")
@@ -58,6 +59,19 @@ func main() {
err = cmd.RunMigrate(logger, os.Args[2:]...)
case "repair":
err = cmd.RunRepair(logger, os.Args[2:]...)
+ case "v2":
+ if len(os.Args) < 3 {
+ fmt.Fprintln(os.Stderr, "Usage: openmessage v2 reconcile-signal [--from dir] [--since YYYY-MM-DD] [--dry-run] [--json]")
+ os.Exit(1)
+ }
+ switch os.Args[2] {
+ case "reconcile-signal":
+ err = cmd.RunReconcileSignal(logger, os.Args[3:]...)
+ default:
+ fmt.Fprintf(os.Stderr, "Unknown v2 command: %s\n", os.Args[2])
+ fmt.Fprintln(os.Stderr, "Usage: openmessage v2 reconcile-signal [--from dir] [--since YYYY-MM-DD] [--dry-run] [--json]")
+ os.Exit(1)
+ }
case "read", "search":
if len(os.Args) < 3 {
fmt.Fprintln(os.Stderr, "Usage: openmessage read [--limit N] [--phone NUMBER] [--json]")
@@ -130,7 +144,7 @@ func main() {
err = cmd.RunDebugMedia(logger, os.Args[2])
default:
fmt.Fprintf(os.Stderr, "Unknown command: %s\n", os.Args[1])
- fmt.Fprintln(os.Stderr, "Usage: openmessage ")
+ fmt.Fprintln(os.Stderr, "Usage: openmessage ")
os.Exit(1)
}