Skip to content

Commit 1965ef9

Browse files
committed
feat(tray): restore browser parity and stream memory events
1 parent 3cb9086 commit 1965ef9

8 files changed

Lines changed: 1832 additions & 148 deletions

bridge_sync_worker.py

Lines changed: 43 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,7 @@
4848
export_index_json,
4949
mark_tombstones_pushed,
5050
load_remote_tasks_for_merge,
51+
task_source_event_ids,
5152
content_length,
5253
has_meaningful_content,
5354
is_archived_duplicate_redirect_task,
@@ -67,6 +68,7 @@
6768
export_provenance_links,
6869
export_knowledge_links,
6970
export_memory_events,
71+
write_memory_events_file_streaming,
7072
export_memory_audit_issues,
7173
export_memory_artifacts,
7274
export_memory_conflicts,
@@ -581,9 +583,13 @@ def _export_knowledge_ratings(conn: sqlite3.Connection) -> list:
581583
return []
582584

583585

584-
def _export_extended_memory(conn: sqlite3.Connection) -> dict[str, list]:
586+
def _export_extended_memory(
587+
conn: sqlite3.Connection,
588+
*,
589+
include_memory_events: bool = True,
590+
) -> dict[str, list]:
585591
"""Export append-only/event/provenance memory artifacts for cross-device sync."""
586-
return {
592+
exported = {
587593
"context_chunks": export_context_chunks(conn),
588594
"context_annotations": export_context_annotations(conn),
589595
"context_questions": export_context_questions(conn),
@@ -592,12 +598,14 @@ def _export_extended_memory(conn: sqlite3.Connection) -> dict[str, list]:
592598
"canonical_facts": export_canonical_facts(conn),
593599
"provenance_links": export_provenance_links(conn),
594600
"knowledge_links": export_knowledge_links(conn),
595-
"memory_events": export_memory_events(conn),
596601
"memory_audit_issues": export_memory_audit_issues(conn),
597602
"memory_artifacts": export_memory_artifacts(conn),
598603
"memory_conflicts": export_memory_conflicts(conn),
599604
"memory_audit_state": export_memory_audit_state(conn),
600605
}
606+
if include_memory_events:
607+
exported["memory_events"] = export_memory_events(conn)
608+
return exported
601609

602610

603611
def _merge_remote_tasks(tasks_out: list[dict], existing_data: dict) -> list[dict]:
@@ -834,9 +842,26 @@ def _main_locked(
834842
except (json.JSONDecodeError, OSError, TypeError) as exc:
835843
log.warning("shared.json read failed for merge: %s", exc)
836844

845+
# Resolve task artifacts before the extended-memory import so the streaming
846+
# ledger reader retains only event IDs actually referenced by task LWW
847+
# metadata (plus one causal head per relevant aggregate).
848+
remote_tasks, _loaded_from_index = load_remote_tasks_for_merge(
849+
bridge_dir,
850+
remote_payload,
851+
log,
852+
)
853+
remote_event_subset: list[dict] = []
854+
837855
with get_conn(_db_path) as conn:
838856
_progress(progress_callback, 10, "Importing remote entities...")
839-
br = import_remote_bridge_data(conn, bridge_dir, remote_payload, log)
857+
br = import_remote_bridge_data(
858+
conn,
859+
bridge_dir,
860+
remote_payload,
861+
log,
862+
remote_task_event_ids=task_source_event_ids(remote_tasks),
863+
event_subset_out=remote_event_subset,
864+
)
840865
if br["entities"] or br["relations"]:
841866
log.info(
842867
"Imported %d remote entities and %d relations",
@@ -847,19 +872,14 @@ def _main_locked(
847872
log.info("Imported %d remote knowledge ratings", br["ratings"])
848873
with get_conn(_db_path) as conn:
849874
_progress(progress_callback, 15, "Importing remote tasks...")
850-
remote_tasks, _loaded_from_index = load_remote_tasks_for_merge(
851-
bridge_dir,
852-
remote_payload,
853-
log,
854-
)
855875
merge_failed = False
856876
if remote_tasks:
857877
try:
858878
new_t, upd_t = merge_import_tasks(
859879
conn,
860880
remote_tasks,
861881
import_content=True,
862-
remote_events=remote_payload.get("memory_events", []),
882+
remote_events=remote_event_subset,
863883
)
864884
sync_task_attachments_from_remote(conn, remote_tasks, bridge_dir)
865885
log.info("LWW merged %d new tasks, %d field updates", new_t, upd_t)
@@ -1051,7 +1071,16 @@ def _main_locked(
10511071
_progress(progress_callback, 55, "Exporting knowledge ratings...")
10521072
kr_out = _export_knowledge_ratings(conn)
10531073
_progress(progress_callback, 58, "Exporting memory ledger...")
1054-
extended_memory = _export_extended_memory(conn)
1074+
# memory_events is hundreds of MB on the live ledger. Stream it to the
1075+
# same atomic JSON transport file instead of retaining 453k dicts plus
1076+
# one giant serialized string in the long-lived Qt process.
1077+
_memory_events_path, memory_event_count = write_memory_events_file_streaming(
1078+
conn, bridge_dir
1079+
)
1080+
log.info("streamed %d memory events", memory_event_count)
1081+
extended_memory = _export_extended_memory(
1082+
conn, include_memory_events=False
1083+
)
10551084
# Export transaction closed
10561085

10571086
# Phase 4: Build payload + write files + git ops (no transaction)
@@ -1065,7 +1094,9 @@ def _main_locked(
10651094
"tasks": tasks_out,
10661095
}
10671096
# v5: write extended memory to separate files (keeps shared.json under CF Pages 25 MB limit)
1068-
write_extended_memory_files(bridge_dir, extended_memory)
1097+
write_extended_memory_files(
1098+
bridge_dir, extended_memory, skip_keys={"memory_events"}
1099+
)
10691100
for key in EXTENDED_MEMORY_KEYS:
10701101
payload[key] = [] # empty placeholders for backward compat
10711102
if pub_entities or pub_tasks:

0 commit comments

Comments
 (0)