Skip to content
Merged
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
Original file line number Diff line number Diff line change
@@ -0,0 +1,107 @@
# Server Cursor Int64

## Problem

Authoritative server transaction progress is represented by unrelated integers and
versioned strings. The string comparison rejects normally ordered pull batches
across decimal digit boundaries, and repeated conversions do not enforce numeric
range or non-negativity consistently. The server JSON number representation also
requires an explicit JavaScript safe-integer bound.

## Decision

Own one abstract, non-negative int64 cursor in logseq_db_types.Server_cursor,
the existing common dependency of storage and overlay. Expose the identical
abstract type through logseq_overlay_db.Server_cursor and Types.Server_cursor.
The overlay spec signature keeps the zero/of_int64/to_int64/equal/compare API,
with type identity tied to the common abstract type, never to public int64.

Migrate authoritative progress in protocol, reducer, worker, snapshot/checkpoint,
overlay and persistence to this type. Convert to/from int64 only in I/O codecs
and checked interval arithmetic. JSON numbers must be in [0, 9007199254740991].
Preserve existing numeric versioned cursor tokens only inside legacy persistence
codecs. Do not migrate local changefeed cursors, Datascript transaction entity IDs,
submission counts/ordinals, graph or lifecycle generations.

The only currently identified necessary Dune change is in
logseq_overlay_db/spec/dune:

(modules server_cursor types database)
(virtual_modules server_cursor types database)

Both shared types and the overlay implementation discover the new implementation
modules through their existing standard module selection. No new package or
storage-to-overlay dependency is needed. Further necessary Dune changes require
separate explicit permission.


Adopt the existing common db_types dependency as the owner of the abstract
non-negative int64 cursor and expose the identical type through overlay.
The user authorized the migration and exactly the two spec Dune module-list
changes above. Other Dune files remain unchanged.

## Alternatives considered

### Define the type only in overlay

Storage already sits below the overlay implementation, so using an overlay-owned
cursor in the storage checkpoint API introduces a package implementation cycle.

### Keep storage and protocol as raw int64

This permits unconstrained values outside I/O and fails the requested shared
progress identity. Only codec boundaries should expose raw numeric values.

### Retain VERSIONED_TOKEN and change its comparator

This keeps unrelated token factories and malformed numeric states constructible,
and does not unify protocol/checkpoint progress or protect numeric ranges.

## Acceptance criteria

- Normal public reducer pull batches crossing 9 to 10 and 99 to 100 are admitted;
duplicates and reverse order are rejected.
- Cursor construction rejects negatives; I/O rejects overflow and values beyond
the server's JSON safe-integer range without silent narrowing.
- SQLite checkpoints and existing outbox/receipt tokens recover the same progress.
- Direct production consumers share the same abstract type, with no public raw
integer/string constructors bypassing of_int64 validation.
- Relevant tests, all direct consumer builds, independent diff review and exact
PR head CI complete; no merge or deployment is performed.

## Risks

- Legacy persistence syntax remains a compatibility boundary rather than a public
domain constructor; malformed legacy numeric values will be rejected explicitly.
- Non-negative int64 is a broader domain than the JavaScript wire can represent;
wire codecs must impose the additional upper bound.
- Existing arithmetic that adds ordinals must reject int64 overflow rather than
wrap; full domain values must not pass through OCaml int conversions.
- This change does not address the independent WebSocket reconnect finding.

## Consequences

Authoritative progress now has one abstract identity across storage, overlay,
protocol, reducer and worker; negative values cannot enter domain state.
Legacy numeric token bytes remain unchanged at the persistence boundary, and
wire encoders/decoders impose the server's JavaScript safe-integer ceiling.
Checked interval arithmetic rejects an overflowing submission before freezing or
persisting it, including an independent suffix after a definitive rejection.
The native UI consumers build without changing their cursor display semantics.

The HTTP snapshot baseline retains its existing 64 MiB response budget, shared
between its transport request and decoder. WebSocket response decoding retains
its independent 262,144-byte budget. A public Core.step bootstrap regression uses
ordinary cursor values and a valid larger HTTP pull response; decimal-boundary
bug regressions remain exclusively at that public reducer boundary.

The migration does not run live account or Simulator acceptance. Repository-wide
decision validation still reports three pre-existing unrelated document errors;
this decision is independently validated.

## Questions

- May logseq_overlay_db/spec/dune add server_cursor to its modules and
virtual_modules lists as shown above? AGENTS.md prohibits Dune edits unless
explicitly requested. The user explicitly approved exactly these two lines on
2026-10-09; no other Dune modification is authorized.
28 changes: 22 additions & 6 deletions logseq_db_storage/lib/sync_checkpoint_store.ml
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,11 @@ let initialize_database db checkpoint =
(Graph_types.Uuid.to_string checkpoint.graph_id));
Sqlite3.Rc.check (Sqlite3.bind_int statement 3 checkpoint.schema.major);
Sqlite3.Rc.check (Sqlite3.bind_int statement 4 checkpoint.schema.minor);
Sqlite3.Rc.check (Sqlite3.bind_int statement 5 checkpoint.applied_server_t);
Sqlite3.Rc.check
(Sqlite3.bind_int64
statement
5
(Server_cursor.to_int64 checkpoint.applied_server_t));
Sqlite3.Rc.check (Sqlite3.bind_text statement 6 checkpoint.checksum);
Sqlite3.Rc.check
(Sqlite3.bind_text
Expand Down Expand Up @@ -104,7 +108,14 @@ let read_database db =
let graph_id = Sqlite3.column_text statement 1 in
let major = Sqlite3.column_int statement 2 in
let minor = Sqlite3.column_int statement 3 in
let applied_server_t = Sqlite3.column_int statement 4 in
let applied_server_t =
match Sqlite3.column statement 4 with
| Sqlite3.Data.INT value ->
Server_cursor.of_int64 value
|> Result.map_error (fun `Negative_cursor ->
"sync_meta cursor is negative")
| _ -> Error "sync_meta cursor is not an integer"
in
let checksum = Sqlite3.column_text statement 5 in
let status =
match Sqlite3.column statement 6 with
Expand All @@ -124,9 +135,10 @@ let read_database db =
match Graph_types.Uuid.of_string graph_id with
| Error _ -> Error "sync_meta graph id is invalid"
| Ok graph_id ->
(match status, last_error with
| Error message, _ | _, Error message -> Error message
| Ok status, Ok last_error ->
(match applied_server_t, status, last_error with
| Error message, _, _ | _, Error message, _ | _, _, Error message ->
Error message
| Ok applied_server_t, Ok status, Ok last_error ->
(match
Sync_checkpoint.create_full
~graph_id
Expand Down Expand Up @@ -169,7 +181,11 @@ let update_database db checkpoint =
(Graph_types.Uuid.to_string checkpoint.graph_id));
Sqlite3.Rc.check (Sqlite3.bind_int statement 3 checkpoint.schema.major);
Sqlite3.Rc.check (Sqlite3.bind_int statement 4 checkpoint.schema.minor);
Sqlite3.Rc.check (Sqlite3.bind_int statement 5 checkpoint.applied_server_t);
Sqlite3.Rc.check
(Sqlite3.bind_int64
statement
5
(Server_cursor.to_int64 checkpoint.applied_server_t));
Sqlite3.Rc.check (Sqlite3.bind_text statement 6 checkpoint.checksum);
Sqlite3.Rc.check
(Sqlite3.bind_text
Expand Down
1 change: 1 addition & 0 deletions logseq_db_types/lib/limits.ml
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
let maximum_request_bytes = 1_048_576
let maximum_response_bytes = 262_144
let maximum_snapshot_baseline_response_bytes = 64 * 1024 * 1024
let maximum_push_bytes = 65_536
let maximum_changed_uuids = 1_024
1 change: 1 addition & 0 deletions logseq_db_types/lib/limits.mli
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
val maximum_request_bytes : int
val maximum_response_bytes : int
val maximum_snapshot_baseline_response_bytes : int
val maximum_push_bytes : int
val maximum_changed_uuids : int
7 changes: 7 additions & 0 deletions logseq_db_types/lib/server_cursor.ml
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
type t = int64

let zero = 0L
let of_int64 value = if value < 0L then Error `Negative_cursor else Ok value
let to_int64 value = value
let equal = Int64.equal
let compare = Int64.compare
8 changes: 8 additions & 0 deletions logseq_db_types/lib/server_cursor.mli
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
(** A non-negative authoritative server transaction cursor. *)
type t

val zero : t
val of_int64 : int64 -> (t, [ `Negative_cursor ]) result
val to_int64 : t -> int64
val equal : t -> t -> bool
val compare : t -> t -> int
6 changes: 2 additions & 4 deletions logseq_db_types/lib/sync_checkpoint.ml
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ type t =
{ format_version : int
; graph_id : Graph_types.Uuid.t
; schema : Graph_types.schema_version
; applied_server_t : int
; applied_server_t : Server_cursor.t
; checksum : string
; status : status
; last_error : string option
Expand All @@ -26,8 +26,6 @@ let valid_checksum value =
let create_full ~graph_id ~schema ~applied_server_t ~checksum ~status ~last_error =
if schema.Graph_types.major < 0 || schema.minor < 0
then Error "sync schema version must be non-negative"
else if applied_server_t < 0
then Error "applied server t must be non-negative"
else if not (valid_checksum checksum)
then Error "sync checksum must be 16 lowercase hexadecimal characters"
else if status = Active && Option.is_some last_error
Expand Down Expand Up @@ -61,7 +59,7 @@ let equal left right =
left.format_version = right.format_version
&& Graph_types.Uuid.equal left.graph_id right.graph_id
&& left.schema = right.schema
&& left.applied_server_t = right.applied_server_t
&& Server_cursor.equal left.applied_server_t right.applied_server_t
&& String.equal left.checksum right.checksum
&& left.status = right.status
&& left.last_error = right.last_error
Expand Down
6 changes: 3 additions & 3 deletions logseq_db_types/lib/sync_checkpoint.mli
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ type t =
{ format_version : int
; graph_id : Graph_types.Uuid.t
; schema : Graph_types.schema_version
; applied_server_t : int
; applied_server_t : Server_cursor.t
; checksum : string
; status : status
; last_error : string option
Expand All @@ -17,14 +17,14 @@ val format_version : int
val create
: graph_id:Graph_types.Uuid.t
-> schema:Graph_types.schema_version
-> applied_server_t:int
-> applied_server_t:Server_cursor.t
-> checksum:string
-> (t, string) result

val create_full
: graph_id:Graph_types.Uuid.t
-> schema:Graph_types.schema_version
-> applied_server_t:int
-> applied_server_t:Server_cursor.t
-> checksum:string
-> status:status
-> last_error:string option
Expand Down
6 changes: 1 addition & 5 deletions logseq_db_worker/lib/effect_runner/effect_runner.ml
Original file line number Diff line number Diff line change
Expand Up @@ -1458,11 +1458,7 @@ let handle_sync_worker_effect t = function
(match Hashtbl.find_opt t.inspections key with
| None -> Error (effect_error "The snapshot mirror inspection is stale.")
| Some inspection ->
let cursor =
Overlay.Server_cursor.of_string
(Printf.sprintf "server-cursor:v1:%d" request.applied_server_t)
|> Result.get_ok
in
let cursor = request.applied_server_t in
(match
Database.prepare_snapshot_activation
t.dependencies.overlay
Expand Down
2 changes: 1 addition & 1 deletion logseq_db_worker/lui/logseq_db_worker_lui_service.ml
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,7 @@ type snapshot =
{ sync_phase : sync_phase
; catalog : graph list
; selected_graph : graph_id option
; applied_server_t : int option
; applied_server_t : Logseq_db_types.Server_cursor.t option
; timeline_presentation_pending : bool
; startup : startup_facts
; last_error : string option
Expand Down
2 changes: 1 addition & 1 deletion logseq_db_worker/lui/logseq_db_worker_lui_service.mli
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ type snapshot =
{ sync_phase : sync_phase
; catalog : graph list
; selected_graph : graph_id option
; applied_server_t : int option
; applied_server_t : Logseq_db_types.Server_cursor.t option
; timeline_presentation_pending : bool
; startup : startup_facts
; last_error : string option
Expand Down
Loading
Loading