Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 42 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -277,3 +277,45 @@ Debugging a live install (failing sends, re-pairing, the two-data-dir gotcha, si
## License

MIT

### Complete Google Contacts directory (optional, legacy read mode)

Google Messages supplies a limited contact suggestion list. Names in existing
threads often arrive independently through conversation updates; a successful
Google Messages connection does not prove the whole address book was imported.

To keep a complete directory, configure an authenticated Google Contacts MCP
server exposing the `contacts_list` tool with Google People API `connections`,
`nextPageToken`, and `totalItems` (or `totalPeople`) fields. For example,
[google-contacts-mcp](https://github.com/domdomegg/google-contacts-mcp) supports
this contract. Authorize that connector for the same Google account as the phone.
OpenMessage calls only its read-only listing tool; OAuth stays in the connector.

- `OPENMESSAGES_GOOGLE_CONTACTS_MCP_URL`: opt-in Streamable HTTP endpoint, such as
`http://127.0.0.1:3230/mcp`. HTTP requires a literal loopback IP; remote endpoints
require HTTPS. Redirects and credentials/query strings in the URL are rejected.
- `OPENMESSAGES_GOOGLE_CONTACTS_MCP_TOKEN_FILE`: optional private (0600) file
containing a connector bearer token, if required. Do not put a token in the URL.
- `OPENMESSAGES_CONTACTS_REFRESH_INTERVAL`: refresh interval, default `5m`,
minimum `1m`, maximum `24h`.

The daemon imports immediately and on the interval, even when the phone is
unavailable. All pages must arrive and agree with the reported total before a
single transaction replaces the previous directory. Failures preserve that last
complete snapshot. All returned phones, emails and organizations are retained;
full dialable phone numbers become compose suggestions. Email-only records remain
in the directory but do not create SMS routes or empty conversations.

For existing SMS threads, unambiguous numbers replace blank/numeric names, and
names supplied by the previous directory follow renames/deletions. Shared numbers
are left unresolved. Custom titles and group titles are preserved, as are message
history, phone-side contact IDs, read state and favorites. Later raw-number
conversation snapshots also consult the directory.

`GET /api/status` includes `contact_sync` with the source, running state, last
attempt/success times, people/phone counts and refresh errors. `complete` describes
the last successfully imported snapshot; inspect `last_error` and the success
age as well. `POST /api/contacts/sync` performs a manual refresh. With no connector
configured, the existing Google Messages suggestion fetch remains available and
is not reported as a complete directory. Demo, client-only and v2-primary modes
do not run the full-directory worker; v2 directory integration is future work.
5 changes: 5 additions & 0 deletions cmd/serve.go
Original file line number Diff line number Diff line change
Expand Up @@ -161,6 +161,7 @@ func RunServe(logger zerolog.Logger, args ...string) error {
if err != nil {
return fmt.Errorf("init app: %w", err)
}
a.ContactDirectoryDisabled = isDemo || v2Primary || !transports
defer a.Close()

interactiveTerminal := term.IsTerminal(int(os.Stdin.Fd()))
Expand Down Expand Up @@ -659,6 +660,9 @@ func RunServe(logger zerolog.Logger, args ...string) error {
// same store double-sends every due message.
if transports {
startLegacyScheduler(v2Primary, a.StartScheduler)
if !isDemo && !v2Primary {
a.StartContactDirectorySync()
}
}

v2Options := v2SendWebOptions(stack, v2Send)
Expand Down Expand Up @@ -717,6 +721,7 @@ func RunServe(logger zerolog.Logger, args ...string) error {
BackfillStatus: func() any { return a.GetBackfillProgress() },
BackfillPhone: a.BackfillConversationByPhone,
SyncGoogleContacts: a.SyncGoogleContacts,
ContactSyncStatus: func() any { return a.GetContactSyncStatus() },
})
} else {
httpHandler = web.ProtectLocalControl(controlAuth.Handler(mcpHTTPHandler))
Expand Down
5 changes: 5 additions & 0 deletions internal/app/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,10 @@ func (p *BackfillProgress) snapshot() BackfillSnapshot {
}

type App struct {
// Set before starting transports; directory writes currently target legacy reads.
ContactDirectoryDisabled bool

contactDirectory contactDirectoryState
clientMu sync.RWMutex
Client *client.Client
googleGeneration *GoogleGeneration
Expand Down Expand Up @@ -870,6 +874,7 @@ func (a *App) GetBackfillProgress() BackfillSnapshot {
}

func (a *App) Close() {
a.stopContactDirectorySync()
a.StopGoogleAvatarSync()
if cli := a.GetClient(); cli != nil {
cli.GM.Disconnect()
Expand Down
168 changes: 168 additions & 0 deletions internal/app/contact_directory.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,168 @@
package app

import (
"context"
"errors"
"os"
"strings"
"sync"
"time"

"github.com/maxghenis/openmessage/internal/contactsync"
)

type ContactSyncStatus struct {
Source string `json:"source"`
Running bool `json:"running"`
Complete bool `json:"complete"`
LastAttemptMS int64 `json:"last_attempt_ms"`
LastSuccessMS int64 `json:"last_success_ms"`
People int `json:"people"`
PhoneEntries int `json:"phone_entries"`
ThreadsUpdated int `json:"threads_updated"`
RefreshIntervalSeconds int `json:"refresh_interval_seconds"`
LastError string `json:"last_error,omitempty"`
}
type contactDirectoryState struct {
mu sync.Mutex
workMu sync.Mutex
wg sync.WaitGroup
cancel context.CancelFunc
ctx context.Context
closed bool
status ContactSyncStatus
}

func googleContactsMCPURL() string {
return strings.TrimSpace(os.Getenv("OPENMESSAGES_GOOGLE_CONTACTS_MCP_URL"))
}
func contactRefreshInterval() time.Duration {
if v, err := time.ParseDuration(strings.TrimSpace(os.Getenv("OPENMESSAGES_CONTACTS_REFRESH_INTERVAL"))); err == nil && v >= time.Minute && v <= 24*time.Hour {
return v
}
return 5 * time.Minute
}
func (a *App) GetContactSyncStatus() ContactSyncStatus {
d := &a.contactDirectory
d.mu.Lock()
defer d.mu.Unlock()
s := d.status
if a.ContactDirectoryDisabled {
s.Source = "disabled"
return s
}
if s.Source == "" {
s.Source = "google_messages_suggestions"
if googleContactsMCPURL() != "" {
s.Source = "google_contacts_mcp"
}
}
if s.Source == "google_contacts_mcp" {
s.RefreshIntervalSeconds = int(contactRefreshInterval() / time.Second)
}
return s
}

// StartContactDirectorySync is daemon-only and independent of phone connectivity.
// Reconnects cannot create additional workers. Close cancels and joins requests.
func (a *App) StartContactDirectorySync() {
if a == nil || a.ContactDirectoryDisabled || googleContactsMCPURL() == "" {
return
}
d := &a.contactDirectory
d.mu.Lock()
if d.closed || d.cancel != nil {
d.mu.Unlock()
return
}
ctx, cancel := context.WithCancel(context.Background())
d.ctx = ctx
d.cancel = cancel
d.wg.Add(1)
d.mu.Unlock()
go func() {
defer d.wg.Done()
runContactRefreshLoop(ctx, contactRefreshInterval(), func() {
if _, err := a.syncContactDirectory(ctx); err != nil && ctx.Err() == nil {
a.Logger.Warn().Msg("Google Contacts directory refresh failed; check contact_sync status")
}
})
}()
}
func runContactRefreshLoop(ctx context.Context, interval time.Duration, refresh func()) {
refresh()
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
refresh()
}
}
}
func (a *App) stopContactDirectorySync() {
d := &a.contactDirectory
d.mu.Lock()
d.closed = true
if d.cancel != nil {
d.cancel()
}
d.mu.Unlock()
d.wg.Wait()
// A manual refresh can be in flight even before the periodic worker starts.
d.workMu.Lock()
d.workMu.Unlock()
}
func (a *App) syncContactDirectory(parent context.Context) (int, error) {
if a.ContactDirectoryDisabled {
return 0, errors.New("contact directory sync requires a legacy daemon")
}
d := &a.contactDirectory
d.workMu.Lock()
defer d.workMu.Unlock()
d.mu.Lock()
if d.closed {
d.mu.Unlock()
return 0, errors.New("contact sync is stopped")
}
// Tie manual requests to daemon shutdown as well as their timeout.
if d.ctx != nil {
parent = d.ctx
}
d.status.Source = "google_contacts_mcp"
d.status.Running = true
d.status.LastAttemptMS = time.Now().UnixMilli()
d.mu.Unlock()
ctx, cancel := context.WithTimeout(parent, 90*time.Second)
defer cancel()
people, err := contactsync.Fetch(ctx, googleContactsMCPURL(), strings.TrimSpace(os.Getenv("OPENMESSAGES_GOOGLE_CONTACTS_MCP_TOKEN_FILE")))
phones, changed := 0, 0
if err == nil {
if ctx.Err() != nil {
err = ctx.Err()
} else {
phones, changed, err = a.Store.ReplaceContactDirectory(people)
}
}
d.mu.Lock()
d.status.Running = false
if err != nil {
d.status.LastError = "Contact directory refresh failed; verify connector availability and authorization, then retry"
} else {
d.status.LastSuccessMS = time.Now().UnixMilli()
d.status.Complete = true
d.status.People = len(people)
d.status.PhoneEntries = phones
d.status.ThreadsUpdated = changed
d.status.LastError = ""
}
d.mu.Unlock()
if err != nil {
return 0, err
}
a.Logger.Info().Int("people", len(people)).Int("phone_entries", phones).Int("threads_updated", changed).Msg("Google Contacts directory refreshed")
a.emitConversationsChange()
return phones, nil
}
73 changes: 73 additions & 0 deletions internal/app/contact_directory_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
package app

import (
"context"
"sync/atomic"
"testing"
"time"
)

func TestContactRefreshLoopRepeatsAndCancels(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
done := make(chan struct{})
var calls atomic.Int32
go func() {
defer close(done)
runContactRefreshLoop(ctx, time.Millisecond, func() {
if calls.Add(1) == 3 {
cancel()
}
})
}()
select {
case <-done:
case <-time.After(time.Second):
t.Fatal("worker did not stop")
}
if calls.Load() < 3 {
t.Fatal("refresh did not repeat")
}
}
func TestContactDirectoryWorkerClosesAndDoesNotRestart(t *testing.T) {
t.Setenv("OPENMESSAGES_GOOGLE_CONTACTS_MCP_URL", "http://127.0.0.1:1/mcp")
a := newTestApp(t, &mockGMClient{})
a.StartContactDirectorySync()
a.StartContactDirectorySync()
a.stopContactDirectorySync()
a.StartContactDirectorySync()
if _, err := a.SyncGoogleContacts(); err == nil {
t.Fatal("sync ran after shutdown")
}
if a.GetContactSyncStatus().Running {
t.Fatal("worker left running")
}
}
func TestContactRefreshInterval(t *testing.T) {
for _, s := range []string{"0s", "-1m", "30s", "invalid", "25h"} {
t.Setenv("OPENMESSAGES_CONTACTS_REFRESH_INTERVAL", s)
if contactRefreshInterval() != 5*time.Minute {
t.Fatal(s)
}
}
t.Setenv("OPENMESSAGES_CONTACTS_REFRESH_INTERVAL", "2m")
if contactRefreshInterval() != 2*time.Minute {
t.Fatal("valid interval ignored")
}
}

func TestContactDirectoryDisabledForNonLegacyDaemon(t *testing.T) {
t.Setenv("OPENMESSAGES_GOOGLE_CONTACTS_MCP_URL", "http://127.0.0.1:1/mcp")
a := newTestApp(t, &mockGMClient{})
a.ContactDirectoryDisabled = true
a.StartContactDirectorySync()
if a.contactDirectory.cancel != nil {
t.Fatal("started unsupported worker")
}
if _, err := a.SyncGoogleContacts(); err == nil {
t.Fatal("wrote unsupported read store")
}
if a.GetContactSyncStatus().Source != "disabled" {
t.Fatal("misleading status")
}
}
8 changes: 8 additions & 0 deletions internal/app/contacts.go
Original file line number Diff line number Diff line change
@@ -1,13 +1,18 @@
package app

import (
"context"
"fmt"
"strings"

"github.com/maxghenis/openmessage/internal/db"
)

func (a *App) StartGoogleContactSync() {
if a != nil && googleContactsMCPURL() != "" {
a.StartContactDirectorySync()
return
}
if a == nil || !googleAvatarSyncEnabled() {
return
}
Expand All @@ -30,6 +35,9 @@ func (a *App) StartGoogleContactSync() {
}

func (a *App) SyncGoogleContacts() (int, error) {
if googleContactsMCPURL() != "" {
return a.syncContactDirectory(context.Background())
}
gm := a.getGMClient()
if gm == nil {
return 0, fmt.Errorf("not connected to Google Messages")
Expand Down
Loading
Loading