返回 DeepSeek-Reasonix
store.go
1 package draftstate
2
3 import (
4 "context"
5 "database/sql"
6 "errors"
7 "fmt"
8 "net/url"
9 "os"
10 "path/filepath"
11 "runtime"
12 "strings"
13 "sync"
14 "time"
15
16 sqlite "modernc.org/sqlite"
17 )
18
19 const SchemaVersion = 4
20 const SnapshotVersion = 5
21
22 var (
23 ErrConflict = errors.New("session draft revision conflict")
24 ErrNotFound = errors.New("session draft not found")
25 ErrConverted = errors.New("session draft was converted")
26 ErrOperationConflict = errors.New("session draft submission conflicts with the active operation")
27 ErrOperationNotFound = errors.New("session draft submission not found")
28 ErrUnsupportedVersion = errors.New("session draft schema version is unsupported")
29 )
30
31 type Draft struct {
32 ID string
33 WorkspaceID string
34 Scope string
35 WorkspaceRoot string
36 Revision uint64
37 ContentJSON string
38 SettingsJSON string
39 Status string
40 UpdatedAt time.Time
41 }
42
43 type Operation struct {
44 RequestID string
45 SourceDigest string
46 Revision uint64
47 ExecutionJSON string
48 ID string
49 DraftID string
50 WorkspaceID string
51 DraftRevision uint64
52 SessionID string
53 TopicID string
54 SubmissionID string
55 Fingerprint string
56 RequestJSON string
57 Phase string
58 Error string
59 CreatedAt time.Time
60 UpdatedAt time.Time
61 }
62
63 type Store struct {
64 path string
65 mu sync.Mutex
66 db *sql.DB
67 }
68
69 func New(path string) *Store { return &Store{path: filepath.Clean(path)} }
70
71 func (s *Store) Path() string {
72 if s == nil || s.path == "." {
73 return ""
74 }
75 return s.path
76 }
77
78 func (s *Store) Close() error {
79 if s == nil {
80 return nil
81 }
82 s.mu.Lock()
83 defer s.mu.Unlock()
84 if s.db == nil {
85 return nil
86 }
87 err := s.db.Close()
88 s.db = nil
89 return err
90 }
91
92 func (s *Store) openLocked() error {
93 if s.db != nil {
94 return nil
95 }
96 if s == nil || s.path == "" || s.path == "." {
97 return errors.New("session draft database path is unavailable")
98 }
99 if err := os.MkdirAll(filepath.Dir(s.path), 0o700); err != nil {
100 return err
101 }
102 db, err := sql.Open("sqlite", draftFileDSN(s.path))
103 if err != nil {
104 return err
105 }
106 db.SetMaxOpenConns(1)
107 var version int
108 if err = db.QueryRow(`PRAGMA user_version`).Scan(&version); err != nil {
109 _ = db.Close()
110 return err
111 }
112 if version < 0 || version > SchemaVersion {
113 _ = db.Close()
114 return fmt.Errorf("%w: %d", ErrUnsupportedVersion, version)
115 }
116 if version == 0 {
117 var tables int
118 if err = db.QueryRow(`SELECT count(*) FROM sqlite_master WHERE type='table' AND name NOT LIKE 'sqlite_%'`).Scan(&tables); err != nil {
119 _ = db.Close()
120 return err
121 }
122 if tables != 0 {
123 // A concurrent first opener may have committed the schema after our
124 // first read. Only an identifiable committed format may be opened.
125 if err = db.QueryRow(`PRAGMA user_version`).Scan(&version); err != nil {
126 _ = db.Close()
127 return err
128 }
129 if version < 1 || version > SchemaVersion {
130 _ = db.Close()
131 return fmt.Errorf("%w: %d", ErrUnsupportedVersion, version)
132 }
133 }
134 }
135 if version > 0 && version < SchemaVersion {
136 if err := backupBeforeUpgrade(db, s.path, version); err != nil {
137 _ = db.Close()
138 return err
139 }
140 }
141 // Even changing journal mode writes the database header. Check the version
142 // first so an older binary leaves an unknown future file byte-for-byte intact.
143 for _, statement := range []string{
144 `PRAGMA busy_timeout = 5000`,
145 `PRAGMA journal_mode = WAL`,
146 `PRAGMA foreign_keys = ON`,
147 } {
148 if _, err = db.Exec(statement); err != nil {
149 _ = db.Close()
150 return err
151 }
152 }
153 if version == 0 {
154 if err = initialize(db); err != nil {
155 _ = db.Close()
156 return err
157 }
158 } else {
159 if version == 1 {
160 if err = migrateV1ToV2(db); err != nil {
161 _ = db.Close()
162 return err
163 }
164 version = 2
165 }
166 if version == 2 {
167 if err = migrateV2ToV3(db); err != nil {
168 _ = db.Close()
169 return err
170 }
171 version = 3
172 }
173 if version == 3 {
174 if err = migrateV3ToV4(db); err != nil {
175 _ = db.Close()
176 return err
177 }
178 }
179 }
180 s.db = db
181 return nil
182 }
183
184 func initialize(db *sql.DB) error {
185 tx, err := db.Begin()
186 if err != nil {
187 return err
188 }
189 defer rollbackTransaction(tx)
190 statements := []string{
191 `CREATE TABLE IF NOT EXISTS drafts (
192 id TEXT PRIMARY KEY, workspace_id TEXT NOT NULL, scope TEXT NOT NULL,
193 workspace_root TEXT NOT NULL, revision INTEGER NOT NULL,
194 content_json TEXT NOT NULL, settings_json TEXT NOT NULL,
195 status TEXT NOT NULL, updated_at INTEGER NOT NULL
196 )`,
197 `CREATE UNIQUE INDEX IF NOT EXISTS one_active_draft_per_workspace ON drafts(workspace_id) WHERE status = 'active'`,
198 `CREATE TABLE IF NOT EXISTS operations (
199 id TEXT PRIMARY KEY, draft_id TEXT NOT NULL, workspace_id TEXT NOT NULL,
200 draft_revision INTEGER NOT NULL, session_id TEXT NOT NULL,
201 topic_id TEXT NOT NULL,
202 submission_id TEXT NOT NULL, fingerprint TEXT NOT NULL,
203 request_json TEXT NOT NULL, phase TEXT NOT NULL, error TEXT NOT NULL,
204 created_at INTEGER NOT NULL, updated_at INTEGER NOT NULL,
205 request_id TEXT NOT NULL DEFAULT '', source_digest TEXT NOT NULL DEFAULT '',
206 operation_revision INTEGER NOT NULL DEFAULT 1, execution_json TEXT NOT NULL DEFAULT ''
207 )`,
208 `CREATE INDEX IF NOT EXISTS operations_by_draft ON operations(draft_id, updated_at DESC)`,
209 `CREATE UNIQUE INDEX IF NOT EXISTS operations_by_request ON operations(draft_id, request_id) WHERE request_id <> ''`,
210 `CREATE TABLE IF NOT EXISTS conflicts (
211 id INTEGER PRIMARY KEY AUTOINCREMENT, draft_id TEXT NOT NULL,
212 expected_revision INTEGER NOT NULL, actual_revision INTEGER NOT NULL,
213 content_json TEXT NOT NULL, settings_json TEXT NOT NULL, created_at INTEGER NOT NULL
214 )`,
215 `CREATE TABLE IF NOT EXISTS restore_state (slot TEXT PRIMARY KEY, draft_id TEXT NOT NULL, updated_at INTEGER NOT NULL)`,
216 fmt.Sprintf(`PRAGMA user_version = %d`, SchemaVersion),
217 }
218 for _, statement := range statements {
219 if _, err := tx.Exec(statement); err != nil {
220 return err
221 }
222 }
223 return tx.Commit()
224 }
225
226 func migrateV3ToV4(db *sql.DB) error {
227 tx, err := db.Begin()
228 if err != nil {
229 return err
230 }
231 defer rollbackTransaction(tx)
232 for _, query := range []string{
233 `ALTER TABLE operations ADD COLUMN request_id TEXT NOT NULL DEFAULT ''`,
234 `ALTER TABLE operations ADD COLUMN source_digest TEXT NOT NULL DEFAULT ''`,
235 `ALTER TABLE operations ADD COLUMN operation_revision INTEGER NOT NULL DEFAULT 1`,
236 `ALTER TABLE operations ADD COLUMN execution_json TEXT NOT NULL DEFAULT ''`,
237 `CREATE UNIQUE INDEX operations_by_request ON operations(draft_id, request_id) WHERE request_id <> ''`,
238 `PRAGMA user_version = 4`,
239 } {
240 if _, err := tx.Exec(query); err != nil {
241 return err
242 }
243 }
244 return tx.Commit()
245 }
246
247 func migrateV1ToV2(db *sql.DB) error {
248 tx, err := db.Begin()
249 if err != nil {
250 return err
251 }
252 defer rollbackTransaction(tx)
253 if _, err := tx.Exec(`ALTER TABLE operations ADD COLUMN topic_id TEXT NOT NULL DEFAULT ''`); err != nil {
254 return err
255 }
256 if _, err := tx.Exec(`PRAGMA user_version = 2`); err != nil {
257 return err
258 }
259 return tx.Commit()
260 }
261
262 // v3 makes operation request_json a versioned, frozen execution snapshot.
263 // The shape is validated by the application because keeping the JSON opaque
264 // lets older completed records remain readable without rewriting user data.
265 func migrateV2ToV3(db *sql.DB) error {
266 tx, err := db.Begin()
267 if err != nil {
268 return err
269 }
270 defer rollbackTransaction(tx)
271 if _, err := tx.Exec(`PRAGMA user_version = 3`); err != nil {
272 return err
273 }
274 return tx.Commit()
275 }
276
277 func draftFileDSN(path string) string {
278 abs, err := filepath.Abs(path)
279 if err != nil {
280 abs = path
281 }
282 slash := filepath.ToSlash(abs)
283 if runtime.GOOS == "windows" && len(slash) >= 2 && slash[1] == ':' {
284 slash = "/" + slash
285 }
286 u := &url.URL{Scheme: "file", Path: slash}
287 return u.String() + "?_pragma=busy_timeout%285000%29&_pragma=foreign_keys%281%29"
288 }
289
290 func (s *Store) withDB(ctx context.Context, fn func(*sql.DB) error) error {
291 s.mu.Lock()
292 defer s.mu.Unlock()
293 deadline := time.Now().Add(5 * time.Second)
294 backoff := 5 * time.Millisecond
295 for {
296 if err := s.openLocked(); err != nil {
297 if !isSQLiteBusy(err) {
298 return err
299 }
300 } else if err := fn(s.db); err != nil {
301 if !isSQLiteBusy(err) {
302 return err
303 }
304 } else {
305 return nil
306 }
307
308 if time.Now().After(deadline) {
309 return fmt.Errorf("session draft database remained busy for 5s")
310 }
311 timer := time.NewTimer(backoff)
312 select {
313 case <-ctx.Done():
314 timer.Stop()
315 return ctx.Err()
316 case <-timer.C:
317 }
318 if backoff < 100*time.Millisecond {
319 backoff *= 2
320 }
321 }
322 }
323
324 func isSQLiteBusy(err error) bool {
325 var sqliteErr *sqlite.Error
326 if !errors.As(err, &sqliteErr) {
327 return false
328 }
329 code := sqliteErr.Code() & 0xff
330 return code == 5 || code == 6
331 }
332
333 func scanDraft(row interface{ Scan(...any) error }) (Draft, error) {
334 var draft Draft
335 var updated int64
336 err := row.Scan(&draft.ID, &draft.WorkspaceID, &draft.Scope, &draft.WorkspaceRoot,
337 &draft.Revision, &draft.ContentJSON, &draft.SettingsJSON, &draft.Status, &updated)
338 if errors.Is(err, sql.ErrNoRows) {
339 return Draft{}, ErrNotFound
340 }
341 draft.UpdatedAt = time.UnixMilli(updated).UTC()
342 return draft, err
343 }
344
345 const draftColumns = `id, workspace_id, scope, workspace_root, revision, content_json, settings_json, status, updated_at`
346
347 func (s *Store) Open(ctx context.Context, workspaceID, scope, root, draftID, settings string) (Draft, bool, error) {
348 var result Draft
349 created := false
350 err := s.withDB(ctx, func(db *sql.DB) error {
351 current, err := scanDraft(db.QueryRowContext(ctx, `SELECT `+draftColumns+` FROM drafts WHERE workspace_id = ? AND status = 'active'`, workspaceID))
352 if err == nil {
353 result = current
354 return nil
355 }
356 if !errors.Is(err, ErrNotFound) {
357 return err
358 }
359 now := time.Now().UTC().UnixMilli()
360 _, err = db.ExecContext(ctx, `INSERT INTO drafts(id, workspace_id, scope, workspace_root, revision, content_json, settings_json, status, updated_at) VALUES(?,?,?,?,1,'{}',?,'active',?)`, draftID, workspaceID, scope, root, settings, now)
361 if err != nil {
362 // A concurrent process may have won the unique active-draft insert.
363 current, lookupErr := scanDraft(db.QueryRowContext(ctx, `SELECT `+draftColumns+` FROM drafts WHERE workspace_id = ? AND status = 'active'`, workspaceID))
364 if lookupErr != nil {
365 return err
366 }
367 result = current
368 return nil
369 }
370 result = Draft{ID: draftID, WorkspaceID: workspaceID, Scope: scope, WorkspaceRoot: root, Revision: 1, ContentJSON: "{}", SettingsJSON: settings, Status: "active", UpdatedAt: time.UnixMilli(now).UTC()}
371 created = true
372 return nil
373 })
374 return result, created, err
375 }
376
377 func (s *Store) Get(ctx context.Context, draftID string) (Draft, error) {
378 var result Draft
379 err := s.withDB(ctx, func(db *sql.DB) error {
380 var err error
381 result, err = scanDraft(db.QueryRowContext(ctx, `SELECT `+draftColumns+` FROM drafts WHERE id = ?`, draftID))
382 return err
383 })
384 return result, err
385 }
386
387 func (s *Store) SetRestore(ctx context.Context, draftID string) error {
388 return s.withDB(ctx, func(db *sql.DB) error {
389 if strings.TrimSpace(draftID) == "" {
390 _, err := db.ExecContext(ctx, `DELETE FROM restore_state WHERE slot='active'`)
391 return err
392 }
393 var status string
394 if err := db.QueryRowContext(ctx, `SELECT status FROM drafts WHERE id=?`, draftID).Scan(&status); errors.Is(err, sql.ErrNoRows) {
395 return ErrNotFound
396 } else if err != nil {
397 return err
398 } else if status != "active" {
399 return ErrConverted
400 }
401 _, err := db.ExecContext(ctx, `INSERT INTO restore_state(slot,draft_id,updated_at) VALUES('active',?,?) ON CONFLICT(slot) DO UPDATE SET draft_id=excluded.draft_id, updated_at=excluded.updated_at`, draftID, time.Now().UTC().UnixMilli())
402 return err
403 })
404 }
405
406 func (s *Store) Restore(ctx context.Context) (Draft, error) {
407 var result Draft
408 err := s.withDB(ctx, func(db *sql.DB) error {
409 var err error
410 result, err = scanDraft(db.QueryRowContext(ctx, `SELECT `+draftColumns+` FROM drafts WHERE id=(SELECT draft_id FROM restore_state WHERE slot='active') AND status='active'`))
411 return err
412 })
413 return result, err
414 }
415
416 func (s *Store) ClearRestore(ctx context.Context, draftID string) error {
417 return s.withDB(ctx, func(db *sql.DB) error {
418 _, err := db.ExecContext(ctx, `DELETE FROM restore_state WHERE slot='active' AND draft_id=?`, draftID)
419 return err
420 })
421 }
422
423 func (s *Store) Save(ctx context.Context, draftID string, expected uint64, content, settings string, _ bool) (Draft, error) {
424 var result Draft
425 err := s.withDB(ctx, func(db *sql.DB) error {
426 tx, err := db.BeginTx(ctx, nil)
427 if err != nil {
428 return err
429 }
430 defer rollbackTransaction(tx)
431 current, err := scanDraft(tx.QueryRowContext(ctx, `SELECT `+draftColumns+` FROM drafts WHERE id = ?`, draftID))
432 if err != nil {
433 return err
434 }
435 if current.Status != "active" {
436 return ErrConverted
437 }
438 var activeOperations int
439 if err := tx.QueryRowContext(ctx, `SELECT COUNT(*) FROM operations WHERE draft_id=? AND phase NOT IN ('cancelled','terminal_failed')`, draftID).Scan(&activeOperations); err != nil {
440 return err
441 }
442 if activeOperations > 0 {
443 return ErrOperationConflict
444 }
445 if current.Revision != expected {
446 _, conflictErr := tx.ExecContext(ctx, `INSERT INTO conflicts(draft_id,expected_revision,actual_revision,content_json,settings_json,created_at) VALUES(?,?,?,?,?,?)`, draftID, expected, current.Revision, content, settings, time.Now().UTC().UnixMilli())
447 if conflictErr != nil {
448 return conflictErr
449 }
450 if err := tx.Commit(); err != nil {
451 return err
452 }
453 result = current
454 return ErrConflict
455 }
456 now := time.Now().UTC().UnixMilli()
457 next := current.Revision + 1
458 resultUpdate, err := tx.ExecContext(ctx, `UPDATE drafts SET revision=?,content_json=?,settings_json=?,updated_at=? WHERE id=? AND status='active' AND revision=?`, next, content, settings, now, draftID, current.Revision)
459 if err != nil {
460 return err
461 }
462 if changed, err := resultUpdate.RowsAffected(); err != nil || changed != 1 {
463 if err != nil {
464 return err
465 }
466 return ErrConflict
467 }
468 if err := tx.Commit(); err != nil {
469 return err
470 }
471 result = current
472 result.Revision = next
473 result.ContentJSON = content
474 result.SettingsJSON = settings
475 result.UpdatedAt = time.UnixMilli(now).UTC()
476 return nil
477 })
478 return result, err
479 }
480
481 func (s *Store) ListActive(ctx context.Context) ([]Draft, error) {
482 out := []Draft{}
483 err := s.withDB(ctx, func(db *sql.DB) error {
484 rows, err := db.QueryContext(ctx, `SELECT `+draftColumns+` FROM drafts WHERE status='active' ORDER BY updated_at DESC`)
485 if err != nil {
486 return err
487 }
488 defer rows.Close()
489 for rows.Next() {
490 draft, err := scanDraft(rows)
491 if err != nil {
492 return err
493 }
494 out = append(out, draft)
495 }
496 return rows.Err()
497 })
498 return out, err
499 }
500
501 func (s *Store) Discard(ctx context.Context, draftID string, expected uint64) error {
502 return s.withDB(ctx, func(db *sql.DB) error {
503 tx, err := db.BeginTx(ctx, nil)
504 if err != nil {
505 return err
506 }
507 defer rollbackTransaction(tx)
508 current, err := scanDraft(tx.QueryRowContext(ctx, `SELECT `+draftColumns+` FROM drafts WHERE id=?`, draftID))
509 if err != nil {
510 return err
511 }
512 if current.Status != "active" {
513 return ErrConverted
514 }
515 if current.Revision != expected {
516 return ErrConflict
517 }
518 var count int
519 if err := tx.QueryRowContext(ctx, `SELECT COUNT(*) FROM operations WHERE draft_id=? AND phase NOT IN ('cancelled','terminal_failed')`, draftID).Scan(&count); err != nil {
520 return err
521 }
522 if count > 0 {
523 return ErrOperationConflict
524 }
525 if _, err = tx.ExecContext(ctx, `UPDATE drafts SET status='discarded',revision=revision+1,updated_at=? WHERE id=?`, time.Now().UTC().UnixMilli(), draftID); err != nil {
526 return err
527 }
528 if _, err = tx.ExecContext(ctx, `DELETE FROM restore_state WHERE slot='active' AND draft_id=?`, draftID); err != nil {
529 return err
530 }
531 return tx.Commit()
532 })
533 }
534
535 func (s *Store) BeginOperation(ctx context.Context, op Operation) (Operation, bool, error) {
536 var result Operation
537 created := false
538 err := s.withDB(ctx, func(db *sql.DB) error {
539 tx, err := db.BeginTx(ctx, nil)
540 if err != nil {
541 return err
542 }
543 defer rollbackTransaction(tx)
544 // Resolve request identity before checking the editable slot: a lost
545 // response may be retried after receipt has converted the draft.
546 if op.RequestID != "" {
547 prior, lookupErr := scanOperation(tx.QueryRowContext(ctx, `SELECT `+operationColumns+` FROM operations WHERE draft_id=? AND request_id=?`, op.DraftID, op.RequestID))
548 if lookupErr == nil {
549 if prior.Fingerprint != op.Fingerprint || prior.SourceDigest != op.SourceDigest {
550 return ErrOperationConflict
551 }
552 result = prior
553 return nil
554 }
555 if !errors.Is(lookupErr, ErrOperationNotFound) {
556 return lookupErr
557 }
558 }
559 draft, err := scanDraft(tx.QueryRowContext(ctx, `SELECT `+draftColumns+` FROM drafts WHERE id=?`, op.DraftID))
560 if err != nil {
561 return err
562 }
563 if draft.Status != "active" {
564 return ErrConverted
565 }
566 if draft.Revision != op.DraftRevision {
567 return ErrConflict
568 }
569 if op.SourceDigest != "" {
570 digest, err := SnapshotDigest(draft.ContentJSON, draft.SettingsJSON)
571 if err != nil {
572 return err
573 }
574 if digest != op.SourceDigest {
575 return ErrConflict
576 }
577 }
578 prior, err := scanOperation(tx.QueryRowContext(ctx, `SELECT id,draft_id,workspace_id,draft_revision,session_id,topic_id,submission_id,fingerprint,request_json,phase,error,created_at,updated_at,request_id,source_digest,operation_revision,execution_json FROM operations WHERE draft_id=? ORDER BY rowid DESC LIMIT 1`, op.DraftID))
579 if err == nil && prior.Phase != "cancelled" && prior.Phase != "terminal_failed" {
580 if prior.Fingerprint != op.Fingerprint {
581 return ErrOperationConflict
582 }
583 result = prior
584 return nil
585 }
586 if err != nil && !errors.Is(err, ErrOperationNotFound) {
587 return err
588 }
589 if err == nil && prior.SessionID != "" {
590 op.SessionID = prior.SessionID
591 op.TopicID = prior.TopicID
592 }
593 now := time.Now().UTC()
594 op.CreatedAt, op.UpdatedAt, op.Phase, op.Error, op.Revision = now, now, "reserved", "", 1
595 _, err = tx.ExecContext(ctx, `INSERT INTO operations(id,draft_id,workspace_id,draft_revision,session_id,topic_id,submission_id,fingerprint,request_json,phase,error,created_at,updated_at,request_id,source_digest,operation_revision) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)`, op.ID, op.DraftID, op.WorkspaceID, op.DraftRevision, op.SessionID, op.TopicID, op.SubmissionID, op.Fingerprint, op.RequestJSON, op.Phase, op.Error, now.UnixMilli(), now.UnixMilli(), op.RequestID, op.SourceDigest, op.Revision)
596 if err != nil {
597 return err
598 }
599 if err := tx.Commit(); err != nil {
600 return err
601 }
602 result, created = op, true
603 return nil
604 })
605 return result, created, err
606 }
607
608 func scanOperation(row interface{ Scan(...any) error }) (Operation, error) {
609 var op Operation
610 var created, updated int64
611 err := row.Scan(&op.ID, &op.DraftID, &op.WorkspaceID, &op.DraftRevision, &op.SessionID, &op.TopicID, &op.SubmissionID, &op.Fingerprint, &op.RequestJSON, &op.Phase, &op.Error, &created, &updated, &op.RequestID, &op.SourceDigest, &op.Revision, &op.ExecutionJSON)
612 if errors.Is(err, sql.ErrNoRows) {
613 return Operation{}, ErrOperationNotFound
614 }
615 op.CreatedAt, op.UpdatedAt = time.UnixMilli(created).UTC(), time.UnixMilli(updated).UTC()
616 return op, err
617 }
618
619 func (s *Store) Operation(ctx context.Context, id string) (Operation, error) {
620 var result Operation
621 err := s.withDB(ctx, func(db *sql.DB) error {
622 var err error
623 result, err = scanOperation(db.QueryRowContext(ctx, `SELECT id,draft_id,workspace_id,draft_revision,session_id,topic_id,submission_id,fingerprint,request_json,phase,error,created_at,updated_at,request_id,source_digest,operation_revision,execution_json FROM operations WHERE id=?`, id))
624 return err
625 })
626 return result, err
627 }
628
629 func (s *Store) PendingOperations(ctx context.Context) ([]Operation, error) {
630 out := []Operation{}
631 err := s.withDB(ctx, func(db *sql.DB) error {
632 rows, err := db.QueryContext(ctx, `SELECT o.id,o.draft_id,o.workspace_id,o.draft_revision,o.session_id,o.topic_id,o.submission_id,o.fingerprint,o.request_json,o.phase,o.error,o.created_at,o.updated_at,o.request_id,o.source_digest,o.operation_revision,o.execution_json FROM operations o JOIN drafts d ON d.id=o.draft_id WHERE d.status='active' AND o.phase NOT IN ('cancelled','terminal_failed') ORDER BY o.created_at`)
633 if err != nil {
634 return err
635 }
636 defer rows.Close()
637 for rows.Next() {
638 op, err := scanOperation(rows)
639 if err != nil {
640 return err
641 }
642 out = append(out, op)
643 }
644 return rows.Err()
645 })
646 return out, err
647 }
648
649 func (s *Store) SetOperationPhase(ctx context.Context, id, phase, message string) (Operation, error) {
650 var result Operation
651 err := s.withDB(ctx, func(db *sql.DB) error {
652 now := time.Now().UTC().UnixMilli()
653 if _, err := db.ExecContext(ctx, `UPDATE operations SET phase=?,error=?,updated_at=?,operation_revision=operation_revision+1 WHERE id=?`, phase, message, now, id); err != nil {
654 return err
655 }
656 var err error
657 result, err = scanOperation(db.QueryRowContext(ctx, `SELECT id,draft_id,workspace_id,draft_revision,session_id,topic_id,submission_id,fingerprint,request_json,phase,error,created_at,updated_at,request_id,source_digest,operation_revision,execution_json FROM operations WHERE id=?`, id))
658 return err
659 })
660 return result, err
661 }
662
663 func (s *Store) ClaimOperationPhase(ctx context.Context, id string, from []string, phase string) (Operation, bool, error) {
664 return s.TransitionOperationPhase(ctx, id, from, phase, "")
665 }
666
667 func (s *Store) TransitionOperationPhase(ctx context.Context, id string, from []string, phase, message string) (Operation, bool, error) {
668 var result Operation
669 claimed := false
670 err := s.withDB(ctx, func(db *sql.DB) error {
671 if len(from) == 0 {
672 return ErrOperationConflict
673 }
674 placeholders := make([]string, len(from))
675 args := make([]any, 0, len(from)+4)
676 args = append(args, phase, message, time.Now().UTC().UnixMilli(), id)
677 for i, value := range from {
678 placeholders[i] = "?"
679 args = append(args, value)
680 }
681 updated, err := db.ExecContext(ctx, `UPDATE operations SET phase=?,error=?,updated_at=?,operation_revision=operation_revision+1 WHERE id=? AND phase IN (`+strings.Join(placeholders, ",")+`)`, args...)
682 if err != nil {
683 return err
684 }
685 count, err := updated.RowsAffected()
686 if err != nil {
687 return err
688 }
689 claimed = count == 1
690 result, err = scanOperation(db.QueryRowContext(ctx, `SELECT id,draft_id,workspace_id,draft_revision,session_id,topic_id,submission_id,fingerprint,request_json,phase,error,created_at,updated_at,request_id,source_digest,operation_revision,execution_json FROM operations WHERE id=?`, id))
691 return err
692 })
693 return result, claimed, err
694 }
695
696 func (s *Store) UpdateOperationRequest(ctx context.Context, id, phase, requestJSON string) (Operation, bool, error) {
697 var result Operation
698 updated := false
699 err := s.withDB(ctx, func(db *sql.DB) error {
700 now := time.Now().UTC().UnixMilli()
701 change, err := db.ExecContext(ctx, `UPDATE operations SET execution_json=?,updated_at=?,operation_revision=operation_revision+1 WHERE id=? AND phase=?`, requestJSON, now, id, phase)
702 if err != nil {
703 return err
704 }
705 count, err := change.RowsAffected()
706 if err != nil {
707 return err
708 }
709 updated = count == 1
710 result, err = scanOperation(db.QueryRowContext(ctx, `SELECT id,draft_id,workspace_id,draft_revision,session_id,topic_id,submission_id,fingerprint,request_json,phase,error,created_at,updated_at,request_id,source_digest,operation_revision,execution_json FROM operations WHERE id=?`, id))
711 return err
712 })
713 return result, updated, err
714 }
715
716 func (s *Store) EnsureOperationTopic(ctx context.Context, id, topicID string) (Operation, error) {
717 var result Operation
718 err := s.withDB(ctx, func(db *sql.DB) error {
719 if _, err := db.ExecContext(ctx, `UPDATE operations SET topic_id=?,updated_at=?,operation_revision=operation_revision+1 WHERE id=? AND topic_id=''`, topicID, time.Now().UTC().UnixMilli(), id); err != nil {
720 return err
721 }
722 var err error
723 result, err = scanOperation(db.QueryRowContext(ctx, `SELECT id,draft_id,workspace_id,draft_revision,session_id,topic_id,submission_id,fingerprint,request_json,phase,error,created_at,updated_at,request_id,source_digest,operation_revision,execution_json FROM operations WHERE id=?`, id))
724 return err
725 })
726 return result, err
727 }
728
729 // AcceptAndConvert publishes the durable submission receipt and releases the
730 // Workspace draft slot in one database transaction.
731 func (s *Store) AcceptAndConvert(ctx context.Context, draftID, operationID string) (Operation, error) {
732 var result Operation
733 err := s.withDB(ctx, func(db *sql.DB) error {
734 tx, err := db.BeginTx(ctx, nil)
735 if err != nil {
736 return err
737 }
738 defer rollbackTransaction(tx)
739 op, err := scanOperation(tx.QueryRowContext(ctx, `SELECT id,draft_id,workspace_id,draft_revision,session_id,topic_id,submission_id,fingerprint,request_json,phase,error,created_at,updated_at,request_id,source_digest,operation_revision,execution_json FROM operations WHERE id=? AND draft_id=?`, operationID, draftID))
740 if err != nil {
741 return err
742 }
743 if op.Phase != "accepted" && op.Phase != "dispatching" && op.Phase != "dispatch_unknown" && op.Phase != "dispatching_shell" {
744 return ErrOperationConflict
745 }
746 now := time.Now().UTC().UnixMilli()
747 if op.Phase != "accepted" {
748 changed, err := tx.ExecContext(ctx, `UPDATE operations SET phase='accepted',error='',updated_at=?,operation_revision=operation_revision+1 WHERE id=? AND phase=?`, now, operationID, op.Phase)
749 if err != nil {
750 return err
751 }
752 op.Revision++
753 if count, err := changed.RowsAffected(); err != nil || count != 1 {
754 if err != nil {
755 return err
756 }
757 return ErrOperationConflict
758 }
759 }
760 converted, err := tx.ExecContext(ctx, `UPDATE drafts SET status='converted',revision=revision+1,updated_at=? WHERE id=? AND status='active'`, now, draftID)
761 if err != nil {
762 return err
763 }
764 if count, err := converted.RowsAffected(); err != nil {
765 return err
766 } else if count == 0 {
767 var status string
768 if err := tx.QueryRowContext(ctx, `SELECT status FROM drafts WHERE id=?`, draftID).Scan(&status); err != nil {
769 return err
770 }
771 if status != "converted" {
772 return ErrOperationConflict
773 }
774 }
775 if _, err := tx.ExecContext(ctx, `DELETE FROM restore_state WHERE slot='active' AND draft_id=?`, draftID); err != nil {
776 return err
777 }
778 if err := tx.Commit(); err != nil {
779 return err
780 }
781 op.Phase, op.Error, op.UpdatedAt = "accepted", "", time.UnixMilli(now).UTC()
782 result = op
783 return nil
784 })
785 return result, err
786 }
787
788 func (s *Store) Convert(ctx context.Context, draftID, operationID string) error {
789 return s.withDB(ctx, func(db *sql.DB) error {
790 tx, err := db.BeginTx(ctx, nil)
791 if err != nil {
792 return err
793 }
794 defer rollbackTransaction(tx)
795 var phase string
796 if err := tx.QueryRowContext(ctx, `SELECT phase FROM operations WHERE id=? AND draft_id=?`, operationID, draftID).Scan(&phase); err != nil {
797 return err
798 }
799 if phase != "accepted" {
800 return ErrOperationConflict
801 }
802 if _, err := tx.ExecContext(ctx, `UPDATE drafts SET status='converted',revision=revision+1,updated_at=? WHERE id=? AND status='active'`, time.Now().UTC().UnixMilli(), draftID); err != nil {
803 return err
804 }
805 return tx.Commit()
806 })
807 }
808
808 lines GO