| 1 | package agent |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "database/sql" |
| 6 | "encoding/json" |
| 7 | "errors" |
| 8 | "fmt" |
| 9 | "os" |
| 10 | "reasonix/internal/fileops" |
| 11 | "reasonix/internal/projectiondb" |
| 12 | "reasonix/internal/store" |
| 13 | ) |
| 14 | |
| 15 | type displayPagerEventSource struct { |
| 16 | fingerprint string |
| 17 | eventInfo os.FileInfo |
| 18 | eventVersion fileops.Version |
| 19 | dag, schemaOne, plain bool |
| 20 | } |
| 21 | |
| 22 | func observeDisplayPagerEvents(ctx context.Context, source, head string, forceSource, plain bool, indexInfo os.FileInfo, fingerprint, stored string, storedErr error) (*displayPagerEventSource, error) { |
| 23 | var err error |
| 24 | var eventInfo os.FileInfo |
| 25 | var eventVersion fileops.Version |
| 26 | dag := false |
| 27 | schemaOne := false |
| 28 | hasDAG := false |
| 29 | if eventInfo, err = os.Stat(store.SessionEventLog(source)); err == nil && eventInfo.Size() > 0 { |
| 30 | eventTarget, version := fileops.DiskSnapshot(store.SessionEventLog(source), eventInfo) |
| 31 | eventVersion = version |
| 32 | fingerprint += fmt.Sprintf(":event:%s:%s", eventTarget.Key, eventVersion) |
| 33 | // Reuse the proven kind for unchanged schema-1 sources. A legacy JSON |
| 34 | // writer may place its identifying fields after a huge messages array; |
| 35 | // rediscovering that header on every cached open would reread the array. |
| 36 | if plain && head == "" && storedErr == nil && stored == fingerprint+":schema1" { |
| 37 | schemaOne = true |
| 38 | } else { |
| 39 | schema, kind, known, probeErr := validatedDisplayEventKind(ctx, source) |
| 40 | if probeErr != nil { |
| 41 | return nil, probeErr |
| 42 | } |
| 43 | hasDAG = known && schema == sessionDAGSchemaVersion |
| 44 | dag = hasDAG && (plain || head != "" || forceSource) |
| 45 | schemaOne = known && schema == sessionEventSchemaVersion && (kind == sessionEventTypeReplace || kind == sessionEventTypeAppend) && plain |
| 46 | } |
| 47 | if dag { |
| 48 | fingerprint += fmt.Sprintf(":dag:%d:%d:%s", eventInfo.Size(), eventInfo.ModTime().UnixNano(), head) |
| 49 | } else if schemaOne { |
| 50 | fingerprint += ":schema1" |
| 51 | } |
| 52 | } else if err != nil && !os.IsNotExist(err) { |
| 53 | return nil, err |
| 54 | } else { |
| 55 | // Missing and empty logs carry no authority over a checkpoint. Record |
| 56 | // that absence explicitly so Validate can detect a newly written log. |
| 57 | eventInfo = nil |
| 58 | } |
| 59 | if head != "" && !dag { |
| 60 | return nil, errors.New("requested branch has no DAG source") |
| 61 | } |
| 62 | if dag || schemaOne { |
| 63 | plain = false |
| 64 | } |
| 65 | if plain { |
| 66 | // A checkpoint is only authoritative in the absence of an event log. |
| 67 | // Event/DAG readers must establish the selected view before indexing. |
| 68 | if err := validateCheckpointDisplaySource(source); err != nil { |
| 69 | return nil, err |
| 70 | } |
| 71 | fingerprint += ":checkpoint" |
| 72 | } else if !dag && !schemaOne { |
| 73 | fingerprint += fmt.Sprintf(":%d:%d", indexInfo.Size(), indexInfo.ModTime().UnixNano()) |
| 74 | } |
| 75 | return &displayPagerEventSource{fingerprint, eventInfo, eventVersion, dag, schemaOne, plain}, nil |
| 76 | } |
| 77 | |
| 78 | func validatedDisplayEventKind(ctx context.Context, source string) (int, string, bool, error) { |
| 79 | f, openErr := fileops.OpenReplaceableRead(store.SessionEventLog(source)) |
| 80 | if openErr != nil { |
| 81 | return 0, "", false, openErr |
| 82 | } |
| 83 | schema, kind, known, probeErr := probeDisplayEventHeader(ctx, f) |
| 84 | _ = f.Close() |
| 85 | if probeErr != nil { |
| 86 | return 0, "", false, probeErr |
| 87 | } |
| 88 | if known && (schema > sessionDAGSchemaVersion || schema == sessionEventSchemaVersion && kind != sessionEventTypeReplace && kind != sessionEventTypeAppend) { |
| 89 | return 0, "", false, fmt.Errorf("%w: event schema %d type %q", ErrDisplayFormatUnsupported, schema, kind) |
| 90 | } |
| 91 | return schema, kind, known, nil |
| 92 | } |
| 93 | |
| 94 | func (p *DisplayPager) loadMatchingHeader(ctx context.Context, info os.FileInfo, identity PersistedState, known bool, head string) (bool, error) { |
| 95 | var header string |
| 96 | if err := p.DB.QueryRowContext(ctx, `SELECT value FROM metadata WHERE key='header'`).Scan(&header); err != nil { |
| 97 | return false, err |
| 98 | } |
| 99 | if err := json.Unmarshal([]byte(header), &p.Header); err != nil { |
| 100 | return false, err |
| 101 | } |
| 102 | return !(p.Header.TranscriptSize != info.Size() || known && head == "" && (p.Header.ContentDigest != identity.DigestHex || p.Header.RevisionKnown != identity.RevisionKnown || p.Header.RevisionKnown && p.Header.Revision != identity.Revision)), nil |
| 103 | } |
| 104 | |
| 105 | func validateCheckpointDisplaySource(source string) error { |
| 106 | if logInfo, logErr := os.Stat(store.SessionEventLog(source)); logErr == nil && logInfo.Size() > 0 { |
| 107 | return ErrDisplayFormatUnsupported |
| 108 | } else if logErr != nil && !os.IsNotExist(logErr) { |
| 109 | return logErr |
| 110 | } |
| 111 | return nil |
| 112 | } |
| 113 | |
| 114 | func (s *displayPagerEventSource) build(ctx context.Context, db *sql.DB, source, indexPath, head string, checkpointSize int64) error { |
| 115 | if s.dag { |
| 116 | return buildDAGDisplayPager(ctx, db, source, s.fingerprint, head, checkpointSize) |
| 117 | } |
| 118 | if s.schemaOne { |
| 119 | return buildEventDisplayPager(ctx, db, source, s.fingerprint, checkpointSize) |
| 120 | } |
| 121 | if s.plain { |
| 122 | return buildCheckpointDisplayPager(ctx, db, source, s.fingerprint) |
| 123 | } |
| 124 | return importDisplayPager(ctx, db, indexPath, s.fingerprint) |
| 125 | } |
| 126 | |
| 127 | func (s *displayPagerEventSource) rebuild(ctx context.Context, opts projectiondb.OpenOptions, source, indexPath, head string, size int64) error { |
| 128 | if s.plain { |
| 129 | opts.ResumeKey = "checkpoint-v1:" + s.fingerprint |
| 130 | } else if s.schemaOne { |
| 131 | opts.ResumeKey = "event-v1:" + s.fingerprint |
| 132 | } else if s.dag { |
| 133 | opts.ResumeKey = "dag-v1:" + s.fingerprint |
| 134 | } else if !s.dag && !s.schemaOne { |
| 135 | opts.ResumeKey = "display-import-v1:" + s.fingerprint |
| 136 | } |
| 137 | return projectiondb.Rebuild(ctx, opts, func(ctx context.Context, db *sql.DB) error { |
| 138 | return s.build(ctx, db, source, indexPath, head, size) |
| 139 | }) |
| 140 | } |
| 141 |