-
Notifications
You must be signed in to change notification settings - Fork 718
feat(sync): retain admitted prompt identity across replicas #1466
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. Weβll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
1d6fc86
8f547a5
36b630d
7ab5aa3
4350477
c2b056b
6e870d1
77ecbd3
da77c73
ed3e050
d167859
0850928
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
|
|
@@ -292,12 +292,13 @@ type UpdateObservationParams struct { | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| type Prompt struct { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ID int64 `json:"id"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SyncID string `json:"sync_id"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SessionID string `json:"session_id"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Content string `json:"content"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Project string `json:"project,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| CreatedAt string `json:"created_at"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SourceInboxID string `json:"source_inbox_id,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ID int64 `json:"id"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SyncID string `json:"sync_id"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SessionID string `json:"session_id"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Content string `json:"content"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Project string `json:"project,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| CreatedAt string `json:"created_at"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| type AddPromptParams struct { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
@@ -603,14 +604,15 @@ type syncObservationPayload struct { | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| type syncPromptPayload struct { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SyncID string `json:"sync_id"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SessionID string `json:"session_id"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Content string `json:"content"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Project *string `json:"project,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| CreatedAt string `json:"created_at,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Deleted bool `json:"deleted,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| DeletedAt *string `json:"deleted_at,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| HardDelete bool `json:"hard_delete,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SourceInboxID string `json:"source_inbox_id,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SyncID string `json:"sync_id"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SessionID string `json:"session_id"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Content string `json:"content"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Project *string `json:"project,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| CreatedAt string `json:"created_at,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Deleted bool `json:"deleted,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| DeletedAt *string `json:"deleted_at,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| HardDelete bool `json:"hard_delete,omitempty"` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // syncRelationPayload is the wire format for a memory_relations row sent over | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
@@ -2385,9 +2387,9 @@ func (s *Store) evaluateCloudUpgradeLegacyMutationTx(tx *sql.Tx, mutation SyncMu | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if op == SyncOpUpsert { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| var local syncPromptPayload | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| err := tx.QueryRow( | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| `SELECT sync_id, session_id, content, project, created_at FROM user_prompts WHERE sync_id = ? ORDER BY id DESC LIMIT 1`, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| `SELECT sync_id, session_id, content, project, created_at, ifnull(source_inbox_id, '') FROM user_prompts WHERE sync_id = ? ORDER BY id DESC LIMIT 1`, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| body.SyncID, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ).Scan(&local.SyncID, &local.SessionID, &local.Content, &local.Project, &local.CreatedAt) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ).Scan(&local.SyncID, &local.SessionID, &local.Content, &local.Project, &local.CreatedAt, &local.SourceInboxID) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err != nil && !errors.Is(err, sql.ErrNoRows) { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return cloudUpgradeLegacyMutationEvaluation{}, err | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
@@ -2399,6 +2401,10 @@ func (s *Store) evaluateCloudUpgradeLegacyMutationTx(tx *sql.Tx, mutation SyncMu | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| body.Content = strings.TrimSpace(local.Content) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| changed = true | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if body.SourceInboxID == "" && err == nil && local.SourceInboxID != "" { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| body.SourceInboxID = local.SourceInboxID | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| changed = true | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| missing := []string{} | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if strings.TrimSpace(body.SessionID) == "" { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| missing = append(missing, "session_id") | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
@@ -3831,11 +3837,12 @@ func (s *Store) AddPromptWithResult(p AddPromptParams) (int64, bool, error) { | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return err | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return s.enqueueSyncMutationTx(tx, SyncEntityPrompt, syncID, SyncOpUpsert, syncPromptPayload{ | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SyncID: syncID, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SessionID: p.SessionID, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Content: content, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Project: nullableString(p.Project), | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| CreatedAt: createdAt, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SyncID: syncID, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SessionID: p.SessionID, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Content: content, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Project: nullableString(p.Project), | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| CreatedAt: createdAt, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SourceInboxID: p.SourceInboxID, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| }) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| }) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err != nil { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
@@ -5665,7 +5672,7 @@ func (s *Store) exportWithProjectScope(project string) (_ *ExportData, err error | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // Prompts | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| promptQuery := "SELECT id, ifnull(sync_id, '') as sync_id, session_id, content, ifnull(project, '') as project, created_at FROM user_prompts" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| promptQuery := "SELECT id, ifnull(sync_id, '') as sync_id, session_id, content, ifnull(project, '') as project, created_at, ifnull(source_inbox_id, '') FROM user_prompts" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| promptArgs := []any{} | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if project != "" { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| promptQuery += ` WHERE id IN (SELECT id FROM user_prompts WHERE project = ? | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
@@ -5682,7 +5689,7 @@ func (s *Store) exportWithProjectScope(project string) (_ *ExportData, err error | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| defer promptRows.Close() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| for promptRows.Next() { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| var p Prompt | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err := promptRows.Scan(&p.ID, &p.SyncID, &p.SessionID, &p.Content, &p.Project, &p.CreatedAt); err != nil { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err := promptRows.Scan(&p.ID, &p.SyncID, &p.SessionID, &p.Content, &p.Project, &p.CreatedAt, &p.SourceInboxID); err != nil { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return nil, err | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| data.Prompts = append(data.Prompts, p) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
@@ -5863,10 +5870,43 @@ func (s *Store) Import(data *ExportData) (*ImportResult, error) { | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } else if err != sql.ErrNoRows { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return nil, fmt.Errorf("import prompt %d: %w", p.ID, err) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if p.SourceInboxID != "" { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| var existingID int64 | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| var existingSession, existingIdentity, existingProject, sessionProject string | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| err := tx.QueryRow(`SELECT p.id, p.session_id, ifnull(p.source_inbox_id, ''), ifnull(p.project, ''), ifnull(s.project, '') FROM user_prompts p LEFT JOIN sessions s ON s.id = p.session_id WHERE p.sync_id = ? ORDER BY p.id DESC LIMIT 1`, syncID).Scan(&existingID, &existingSession, &existingIdentity, &existingProject, &sessionProject) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err != nil && err != sql.ErrNoRows { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return nil, fmt.Errorf("import prompt %d: lookup identity: %w", p.ID, err) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err == nil { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| localProject := strings.TrimSpace(existingProject) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| incomingProject := strings.TrimSpace(p.Project) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if localProject == "" { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| localProject = strings.TrimSpace(sessionProject) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if incomingProject == "" { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| incomingProject = strings.TrimSpace(sessionProject) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| localProject, _ = NormalizeProject(localProject) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| incomingProject, _ = NormalizeProject(incomingProject) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if existingSession != p.SessionID || localProject != incomingProject || (existingIdentity != "" && existingIdentity != p.SourceInboxID) { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return nil, fmt.Errorf("import prompt %d: conflicting inbox identity", p.ID) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if existingIdentity == "" { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| res, err := s.execHook(tx, `UPDATE user_prompts SET source_inbox_id = ? WHERE id = ? AND (source_inbox_id IS NULL OR source_inbox_id = '') AND NOT EXISTS (SELECT 1 FROM user_prompts WHERE session_id = ? AND source_inbox_id = ?)`, p.SourceInboxID, existingID, p.SessionID, p.SourceInboxID) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err != nil { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return nil, fmt.Errorf("import prompt %d: adopt identity: %w", p.ID, err) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| updated, err := res.RowsAffected() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err != nil || updated != 1 { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
Comment on lines
+5873
to
+5900
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. ποΈ Data Integrity & Integration | π‘ Minor | β‘ Quick win π Supported by static analysisπ Script executed: sed -n '5720,5910p' internal/store/store.go
rg -n 'CREATE TABLE.*user_prompts|idx_prompts_project|source_inbox_id|ImportProject' internal/store/migrations* internal/store/store.go | head -110Repository: Gentleman-Programming/engram Length of output: 12450 π Script executed: #!/bin/bash
set -e
printf '%s\n' '--- Import/Export declarations and call sites ---'
rg -n -C 4 'func \(s \*Store\) (Import|Export)|\.Import\(|\.Export\(|ImportProject|ExportData|Prompts' --glob '*.go' .
printf '%s\n' '--- prompt schema and indexes ---'
sed -n '1215,1290p' internal/store/store.go
sed -n '1370,1420p' internal/store/store.go
printf '%s\n' '--- prompt export and nearby import context ---'
sed -n '5630,5715p' internal/store/store.go
sed -n '5840,5920p' internal/store/store.go
printf '%s\n' '--- relevant tests ---'
rg -n -C 5 'source_inbox_id|conflicting inbox identity|adopt identity|Import\(' --glob '*_test.go' .
printf '%s\n' '--- PR diff for this file ---'
git diff --unified=30 b44a88f771c6020bacd2b425e9bc3443ca7f095b 77ecbd39e32ceba4af60e8e865802a4f386217df -- internal/store/store.go | sed -n '1,360p'Repository: Gentleman-Programming/engram Length of output: 45679 π€ get_repo_knowledge executed:
Length of output: 44892 π Script executed: #!/bin/bash
set -e
printf '%s\n' '--- export/import declarations ---'
rg -n '^func \(s \*Store\) (Export|ExportProject|exportWithProjectScope|Import)\b' internal/store/store.go
sed -n '5580,5685p' internal/store/store.go
sed -n '5860,5910p' internal/store/store.go
printf '%s\n' '--- schema and indexes ---'
sed -n '1228,1285p' internal/store/store.go
sed -n '1380,1415p' internal/store/store.go
printf '%s\n' '--- CLI import/export reachability ---'
sed -n '1935,2010p' cmd/engram/main.go
printf '%s\n' '--- identity/project tests only ---'
rg -n -C 12 'SourceInboxID|source_inbox_id|inbox identity|conflicting inbox|adopt identity|ExportProject|project.*prompt|prompt.*project' internal/store/store_test.go cmd/engram/main_test.go cmd/engram/main_extra_test.goRepository: Gentleman-Programming/engram Length of output: 42338 Scope identity adoption to the imported project. If an imported prompt has the same Suggested fix- res, err := s.execHook(tx, `UPDATE user_prompts SET source_inbox_id = ? WHERE id = ? AND (source_inbox_id IS NULL OR source_inbox_id = '') AND NOT EXISTS (SELECT 1 FROM user_prompts WHERE session_id = ? AND source_inbox_id = ?)`, p.SourceInboxID, existingID, p.SessionID, p.SourceInboxID)
+ res, err := s.execHook(tx, `UPDATE user_prompts SET source_inbox_id = ? WHERE id = ? AND (source_inbox_id IS NULL OR source_inbox_id = '') AND ifnull(project, '') = ifnull(?, '') AND NOT EXISTS (SELECT 1 FROM user_prompts WHERE session_id = ? AND source_inbox_id = ?)`, p.SourceInboxID, existingID, p.Project, p.SessionID, p.SourceInboxID)π Committable suggestion
Suggested change
π€ Prompt for AI Agents |
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return nil, fmt.Errorf("import prompt %d: inbox identity already owned: %v", p.ID, err) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| res, err := s.execHook(tx, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| `INSERT INTO user_prompts (sync_id, session_id, content, project, created_at) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SELECT ?, ?, ?, ?, ? WHERE NOT EXISTS (SELECT 1 FROM user_prompts WHERE sync_id = ?)`, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| syncID, p.SessionID, p.Content, p.Project, p.CreatedAt, syncID, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| `INSERT INTO user_prompts (sync_id, session_id, content, project, created_at, source_inbox_id) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SELECT ?, ?, ?, ?, ?, ? WHERE NOT EXISTS (SELECT 1 FROM user_prompts WHERE sync_id = ? OR (session_id = ? AND source_inbox_id = ?))`, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| syncID, p.SessionID, p.Content, p.Project, p.CreatedAt, nullableString(p.SourceInboxID), syncID, p.SessionID, nullableString(p.SourceInboxID), | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
coderabbitai[bot] marked this conversation as resolved.
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err != nil { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return nil, fmt.Errorf("import prompt %d: %w", p.ID, err) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
@@ -9126,8 +9166,8 @@ func (s *Store) enqueueRescuedProjectMutationsTx(tx *sql.Tx, target string, sess | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| for _, id := range p.PromptIDs { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| var payload syncPromptPayload | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| err := tx.QueryRow(`SELECT sync_id, session_id, content, project, created_at FROM user_prompts WHERE id = ? AND project = ?`, id, target). | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Scan(&payload.SyncID, &payload.SessionID, &payload.Content, &payload.Project, &payload.CreatedAt) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| err := tx.QueryRow(`SELECT sync_id, session_id, content, project, created_at, ifnull(source_inbox_id, '') FROM user_prompts WHERE id = ? AND project = ?`, id, target). | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Scan(&payload.SyncID, &payload.SessionID, &payload.Content, &payload.Project, &payload.CreatedAt, &payload.SourceInboxID) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if errors.Is(err, sql.ErrNoRows) { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| continue | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
@@ -9682,7 +9722,7 @@ func (s *Store) backfillPromptSyncMutationsTx(tx *sql.Tx, project string, source | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| mutationSource := backfillMutationSource(source) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // ββ Live prompts ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| rows, err := s.queryItHook(tx, ` | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SELECT p.sync_id, p.session_id, p.content, p.project, p.created_at | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SELECT p.sync_id, p.session_id, p.content, p.project, p.created_at, ifnull(p.source_inbox_id, '') | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| FROM user_prompts p | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| LEFT JOIN sessions s ON s.id = p.session_id | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| WHERE ( | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
@@ -9709,7 +9749,7 @@ func (s *Store) backfillPromptSyncMutationsTx(tx *sql.Tx, project string, source | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| var pending []syncPromptPayload | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| for rows.Next() { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| var payload syncPromptPayload | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err := rows.Scan(&payload.SyncID, &payload.SessionID, &payload.Content, &payload.Project, &payload.CreatedAt); err != nil { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err := rows.Scan(&payload.SyncID, &payload.SessionID, &payload.Content, &payload.Project, &payload.CreatedAt, &payload.SourceInboxID); err != nil { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return closeRowsWithError(rows, err) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| pending = append(pending, payload) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
@@ -11003,17 +11043,31 @@ func (s *Store) applyPromptUpsertTx(tx *sql.Tx, payload syncPromptPayload) error | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| var existingID int64 | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| err = tx.QueryRow(`SELECT id FROM user_prompts WHERE sync_id = ? ORDER BY id DESC LIMIT 1`, payload.SyncID).Scan(&existingID) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if payload.SourceInboxID != "" { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| var owner string | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| err := tx.QueryRow(`SELECT sync_id FROM user_prompts WHERE session_id = ? AND source_inbox_id = ?`, payload.SessionID, payload.SourceInboxID).Scan(&owner) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err == nil && owner != payload.SyncID { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return fmt.Errorf("prompt inbox identity conflict for session %q and inbox ID %q", payload.SessionID, payload.SourceInboxID) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err != nil && !errors.Is(err, sql.ErrNoRows) { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return err | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| var existingSessionID, existingSourceInboxID string | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| err = tx.QueryRow(`SELECT id, session_id, ifnull(source_inbox_id, '') FROM user_prompts WHERE sync_id = ? ORDER BY id DESC LIMIT 1`, payload.SyncID).Scan(&existingID, &existingSessionID, &existingSourceInboxID) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err == nil && existingSourceInboxID != "" && (existingSessionID != payload.SessionID || (payload.SourceInboxID != "" && existingSourceInboxID != payload.SourceInboxID)) { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return fmt.Errorf("prompt inbox identity conflict for sync ID %q: existing session %q and inbox ID %q, received session %q and inbox ID %q", payload.SyncID, existingSessionID, existingSourceInboxID, payload.SessionID, payload.SourceInboxID) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err == sql.ErrNoRows { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if strings.TrimSpace(payload.CreatedAt) == "" { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| _, err = s.execHook(tx, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| `INSERT INTO user_prompts (sync_id, session_id, content, project) VALUES (?, ?, ?, ?)`, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| payload.SyncID, payload.SessionID, payload.Content, payload.Project, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| `INSERT INTO user_prompts (sync_id, session_id, content, project, source_inbox_id) VALUES (?, ?, ?, ?, ?)`, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| payload.SyncID, payload.SessionID, payload.Content, payload.Project, nullableString(payload.SourceInboxID), | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } else { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| _, err = s.execHook(tx, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| `INSERT INTO user_prompts (sync_id, session_id, content, project, created_at) VALUES (?, ?, ?, ?, ?)`, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| payload.SyncID, payload.SessionID, payload.Content, payload.Project, payload.CreatedAt, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| `INSERT INTO user_prompts (sync_id, session_id, content, project, created_at, source_inbox_id) VALUES (?, ?, ?, ?, ?, ?)`, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| payload.SyncID, payload.SessionID, payload.Content, payload.Project, payload.CreatedAt, nullableString(payload.SourceInboxID), | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return err | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
@@ -11026,9 +11080,10 @@ func (s *Store) applyPromptUpsertTx(tx *sql.Tx, payload syncPromptPayload) error | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| SET session_id = ?, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| content = ?, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| project = ?, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| source_inbox_id = COALESCE(?, source_inbox_id), | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| created_at = CASE WHEN ? = '' THEN created_at ELSE ? END | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| WHERE id = ?`, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| payload.SessionID, payload.Content, payload.Project, strings.TrimSpace(payload.CreatedAt), payload.CreatedAt, existingID, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| payload.SessionID, payload.Content, payload.Project, nullableString(payload.SourceInboxID), strings.TrimSpace(payload.CreatedAt), payload.CreatedAt, existingID, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return err | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
Uh oh!
There was an error while loading. Please reload this page.