From af46a22d1d1febda5f6c293651685e7174e4c819 Mon Sep 17 00:00:00 2001 From: Max Ghenis Date: Fri, 25 Sep 2026 11:14:55 -0400 Subject: [PATCH] Add `openmessage v2 reconcile-signal` to top up v2 from legacy Signal history A local-only maintenance command for the gap the 7/23-7/25 Signal v2 projection stall left: the legacy store kept ingesting while v2 did not. It imports legacy Signal messages that v2 lacks through the historical-import seam with live ingest's natural keys (idempotent; repeat and interrupted runs converge), requires the existing signal-primary/signal_cli account, corrects conversation kind from the Signal remote prefix, keeps malformed-sender messages under a null sender identity, and defers media (counted in the report). Safety: takes the data directory's instance lock and refuses (exit 4) if it is held or if a backend answers on OPENMESSAGES_HOST:OPENMESSAGES_PORT; opens legacy through a new repair-free read-only path (db.OpenReadOnly) and v2 read-only under --dry-run (sqlite.OpenReadOnly). stdout is exactly one JSON report; guidance goes to stderr. Ported from wip/signal-reconcile (7eac32f, 2026-07-26), which was built on the #155 branch before #155 landed as the squash b02aabb (#156). This commit is the net diff b95775c..7eac32f (b95775c = the #155 tip) without the agent's PROGRESS.md ledger, applied to 2b09b61. main.go conflicted only because main added `repair` in the same spots; both commands are kept. Co-Authored-By: Claude Fable 5 Co-Authored-By: Claude Opus 5.5 --- cmd/reconcile.go | 419 ++++++++++ cmd/reconcile_test.go | 398 +++++++++ internal/db/db.go | 40 + internal/db/readonly_test.go | 71 ++ internal/ingest/reconcile_signal_test.go | 162 ++++ internal/reconcile/signal.go | 661 +++++++++++++++ internal/reconcile/signal_test.go | 980 +++++++++++++++++++++++ internal/storage/sqlite/readonly_test.go | 79 ++ internal/storage/sqlite/store.go | 68 ++ main.go | 18 +- 10 files changed, 2894 insertions(+), 2 deletions(-) create mode 100644 cmd/reconcile.go create mode 100644 cmd/reconcile_test.go create mode 100644 internal/db/readonly_test.go create mode 100644 internal/ingest/reconcile_signal_test.go create mode 100644 internal/reconcile/signal.go create mode 100644 internal/reconcile/signal_test.go create mode 100644 internal/storage/sqlite/readonly_test.go 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) }