| 1 | package agent |
| 2 | |
| 3 | import ( |
| 4 | "bufio" |
| 5 | "context" |
| 6 | "database/sql" |
| 7 | "encoding/json" |
| 8 | "errors" |
| 9 | "fmt" |
| 10 | "io" |
| 11 | "os" |
| 12 | "strings" |
| 13 | |
| 14 | "reasonix/internal/fileops" |
| 15 | "reasonix/internal/historywork" |
| 16 | ) |
| 17 | |
| 18 | type displayImportCheckpoint struct { |
| 19 | Version int `json:"version"` |
| 20 | SourceKey string `json:"sourceKey"` |
| 21 | SourceOffset int64 `json:"sourceOffset"` |
| 22 | Count int `json:"count"` |
| 23 | Users int `json:"users"` |
| 24 | Offset int64 `json:"offset"` |
| 25 | Header map[string]json.RawMessage `json:"header"` |
| 26 | } |
| 27 | |
| 28 | func (progress *displayIndexImportProgress) save(ctx context.Context, tx *sql.Tx) error { |
| 29 | body, err := json.Marshal(displayImportCheckpoint{1, progress.sourceKey, progress.sourceOffset, progress.count, progress.users, progress.offset, progress.header}) |
| 30 | if err != nil { |
| 31 | return err |
| 32 | } |
| 33 | _, err = tx.ExecContext(ctx, `INSERT OR REPLACE INTO metadata VALUES('display_import_progress',?)`, string(body)) |
| 34 | return err |
| 35 | } |
| 36 | |
| 37 | func restoreDisplayImportProgress(ctx context.Context, db *sql.DB, key string, size int64) (*displayIndexImportProgress, error) { |
| 38 | p := &displayIndexImportProgress{sourceKey: key, header: map[string]json.RawMessage{}} |
| 39 | var raw string |
| 40 | err := db.QueryRowContext(ctx, `SELECT value FROM metadata WHERE key='display_import_progress'`).Scan(&raw) |
| 41 | if err != nil && !errors.Is(err, sql.ErrNoRows) { |
| 42 | return nil, err |
| 43 | } |
| 44 | var saved displayImportCheckpoint |
| 45 | valid := err == nil && json.Unmarshal([]byte(raw), &saved) == nil && saved.Version == 1 && saved.SourceKey == key && |
| 46 | saved.SourceOffset > 0 && saved.SourceOffset < size && saved.Count > 0 && saved.Offset > 0 && saved.Users >= 0 && saved.Users <= saved.Count && saved.Header != nil |
| 47 | var count int |
| 48 | var end int64 |
| 49 | if err := db.QueryRowContext(ctx, `SELECT COUNT(*),COALESCE(MAX(offset+length),0) FROM entries`).Scan(&count, &end); err != nil { |
| 50 | return nil, err |
| 51 | } |
| 52 | if valid && count == saved.Count && end == saved.Offset { |
| 53 | p.count, p.users, p.offset = saved.Count, saved.Users, saved.Offset |
| 54 | p.sourceOffset, p.header = saved.SourceOffset, saved.Header |
| 55 | return p, nil |
| 56 | } |
| 57 | // The unfinished generation is disposable. Missing or damaged progress |
| 58 | // cannot certify the rows that happen to remain in the staging database. |
| 59 | _, err = db.ExecContext(ctx, `DELETE FROM entries; DELETE FROM metadata`) |
| 60 | return p, err |
| 61 | } |
| 62 | |
| 63 | func displayImportFileKey(path string, f *os.File) (string, error) { |
| 64 | info, err := f.Stat() |
| 65 | if err != nil { |
| 66 | return "", err |
| 67 | } |
| 68 | current, err := os.Stat(path) |
| 69 | if err != nil { |
| 70 | return "", err |
| 71 | } |
| 72 | target, version := fileops.DiskHandleSnapshot(path, f, info) |
| 73 | pathTarget, pathVersion := fileops.DiskSnapshot(path, current) |
| 74 | if !os.SameFile(info, current) || target.Key != pathTarget.Key || version != pathVersion { |
| 75 | return "", ErrDisplaySourceChanged |
| 76 | } |
| 77 | return fmt.Sprintf("%s:%s", target.Key, version), nil |
| 78 | } |
| 79 | |
| 80 | // A checkpoint is immediately after a complete array element. Recreate only |
| 81 | // the JSON framing, retaining all preceding header fields from the same batch. |
| 82 | // Decoder offsets are translated back to actual file offsets for the next |
| 83 | // durable checkpoint; buffered read-ahead is never mistaken for consumption. |
| 84 | func displayImportDecoder(ctx context.Context, f *os.File, offset int64) (*json.Decoder, int64, error) { |
| 85 | if _, err := f.Seek(offset, io.SeekStart); err != nil { |
| 86 | return nil, 0, err |
| 87 | } |
| 88 | reader := bufio.NewReaderSize(&historywork.Reader{Context: ctx, Source: f}, historywork.ReadChunk) |
| 89 | if offset == 0 { |
| 90 | return json.NewDecoder(reader), 0, nil |
| 91 | } |
| 92 | comma := false |
| 93 | for { |
| 94 | prefix, err := reader.Peek(1) |
| 95 | if err != nil { |
| 96 | return nil, 0, err |
| 97 | } |
| 98 | switch prefix[0] { |
| 99 | case ' ', '\t', '\r', '\n': |
| 100 | case ',': |
| 101 | if comma { |
| 102 | return nil, 0, errors.New("duplicate display entry separator") |
| 103 | } |
| 104 | comma = true |
| 105 | default: |
| 106 | if !comma && prefix[0] != ']' || comma && prefix[0] != '{' { |
| 107 | return nil, 0, errors.New("invalid resumed display entry separator") |
| 108 | } |
| 109 | const framing = `{"entries":[` |
| 110 | return json.NewDecoder(io.MultiReader(strings.NewReader(framing), reader)), offset - int64(len(framing)), nil |
| 111 | } |
| 112 | _, _ = reader.ReadByte() |
| 113 | offset++ |
| 114 | } |
| 115 | } |
| 116 |