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
75 changes: 72 additions & 3 deletions apps/mobile/src/connection/environment-cache-store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ import * as Option from "effect/Option";

import * as MobileDatabase from "../persistence/mobile-database";
import { make } from "./environment-cache-store";
import { encodeStoredShellSnapshot } from "./shell-cache-encoding";
import { makeStoredShellSnapshotEncoder } from "./shell-cache-encoding";

const ENVIRONMENT_ID = EnvironmentId.make("environment-1");
const PROJECT_ID = ProjectId.make("project-1");
Expand Down Expand Up @@ -409,6 +409,9 @@ describe("mobile SQLite environment cache store", () => {
});

describe("cooperative shell cache encoding", () => {
// A fresh encoder per call keeps these cases cold; the cache has its own cases below.
const encodeStoredShellSnapshot = (stored: typeof StoredOrchestrationShellSnapshot.Type) =>
makeStoredShellSnapshotEncoder()(stored);
const encodeOriginal = Schema.encodeEffect(
Schema.fromJsonString(StoredOrchestrationShellSnapshot),
);
Expand Down Expand Up @@ -511,8 +514,10 @@ describe("cooperative shell cache encoding", () => {
...stored,
snapshot: {
...stored.snapshot,
threads: Array.from({ length: 65 }, () => stored.snapshot.threads[0]!),
archivedThreads: Array.from({ length: 33 }, () => stored.snapshot.archivedThreads[0]!),
threads: Array.from({ length: 65 }, () => ({ ...stored.snapshot.threads[0]! })),
archivedThreads: Array.from({ length: 33 }, () => ({
...stored.snapshot.archivedThreads[0]!,
})),
},
});
// Five row chunks (3 active + 2 archived): four yields between them, one before the envelope.
Expand Down Expand Up @@ -592,4 +597,68 @@ describe("cooperative shell cache encoding", () => {
}
}),
);

it.effect("reuses encoded rows across saves without changing the payload", () =>
Effect.gen(function* () {
const encode = makeStoredShellSnapshotEncoder();
expect(yield* encode(stored)).toBe(yield* encodeOriginal(stored));

const setTimer = vi.spyOn(globalThis, "setTimeout");
try {
// Same row references, new envelope: no row work, so no host yields.
const resaved = {
...stored,
snapshot: { ...stored.snapshot, snapshotSequence: 4 },
};
expect(yield* encode(resaved)).toBe(yield* encodeOriginal(resaved));
expect(setTimer).not.toHaveBeenCalled();
} finally {
vi.restoreAllMocks();
}
}),
);

it.effect("encodes replaced, added, and moved-to-archive rows on a warm encoder", () =>
Effect.gen(function* () {
const encode = makeStoredShellSnapshotEncoder();
yield* encode(stored);
const [first, second, ...rest] = stored.snapshot.threads;
const next = {
...stored,
snapshot: {
...stored.snapshot,
threads: [
{ ...first!, title: "Renamed", itemCount: 9 },
...rest,
{ ...first!, id: ThreadId.make("thread-new"), branch: " feature " },
],
archivedThreads: [
...stored.snapshot.archivedThreads,
{ ...second!, archivedAt: DateTime.add(NOW, { hours: 1 }) },
],
},
};
const actual = yield* encode(next);
expect(actual).toBe(yield* encodeOriginal(next));
const parsed = JSON.parse(actual);
expect(parsed.snapshot.threads[0].title).toBe("Renamed");
expect(parsed.snapshot.threads.at(-1).branch).toBe("feature");
expect(parsed.snapshot.archivedThreads).toHaveLength(2);
}),
);

it.effect("does not cache rows from a failed encode", () =>
Effect.gen(function* () {
const encode = makeStoredShellSnapshotEncoder();
const invalidRow = { ...stored.snapshot.threads[0]!, itemCount: -1 };
const invalid = {
...stored,
snapshot: { ...stored.snapshot, threads: [...stored.snapshot.threads, invalidRow] },
};
expect(yield* Effect.isFailure(encode(invalid))).toBe(true);
expect(yield* Effect.isFailure(encode(invalid))).toBe(true);
// Valid rows from the failed save still encode correctly afterwards.
expect(yield* encode(stored)).toBe(yield* encodeOriginal(stored));
}),
);
});
3 changes: 2 additions & 1 deletion apps/mobile/src/connection/environment-cache-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ import * as Option from "effect/Option";
import * as Schema from "effect/Schema";

import * as MobileDatabase from "../persistence/mobile-database";
import { encodeStoredShellSnapshot } from "./shell-cache-encoding";
import { makeStoredShellSnapshotEncoder } from "./shell-cache-encoding";
import {
attachProjectFaviconDatabase,
projectFaviconDatabaseCache,
Expand Down Expand Up @@ -99,6 +99,7 @@ function loadDecodedCache<A, B>(input: {
export const make = Effect.fn("MobileEnvironmentCacheStore.make")(function* () {
const database = yield* MobileDatabase.MobileDatabase;
attachProjectFaviconDatabase(database);
const encodeStoredShellSnapshot = makeStoredShellSnapshotEncoder();
return Persistence.EnvironmentCacheStore.of({
loadShell: Effect.fn("MobileEnvironmentCache.loadShell")((environmentId) =>
loadDecodedCache({
Expand Down
74 changes: 47 additions & 27 deletions apps/mobile/src/connection/shell-cache-encoding.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,33 +30,53 @@ const encodeEnvelope = Schema.encodeEffect(
);

/**
* Encode the shell cache payload exactly like `fromJsonString(StoredOrchestrationShellSnapshot)`,
* yielding to the host between bounded chunks of thread rows. The final envelope encode and
* JSON stringify still run synchronously over the whole payload.
* Make an encoder for the shell cache payload that produces exactly what
* `fromJsonString(StoredOrchestrationShellSnapshot)` would, yielding to the host
* between bounded chunks of thread rows. The final envelope encode and JSON
* stringify still run synchronously over the whole payload.
*
* Shell rows are immutable and the shell reducer keeps unchanged rows by
* reference, so each encoder remembers the canonical encoding of rows it has
* already encoded successfully and only runs the row codec for new rows.
*/
export const encodeStoredShellSnapshot = Effect.fnUntraced(function* (
stored: typeof StoredOrchestrationShellSnapshot.Type,
) {
let hasWorked = false;
const yieldBetweenChunks = Effect.suspend(() => {
if (hasWorked) return yieldToHost;
hasWorked = true;
return Effect.void;
});
const encodeRows = (rows: ReadonlyArray<OrchestrationV2ThreadShellJson>) =>
Effect.gen(function* () {
const encoded: Array<unknown> = [];
for (let start = 0; start < rows.length; start += ROWS_PER_CHUNK) {
yield* yieldBetweenChunks;
const chunk = yield* encodeThreadChunk(rows.slice(start, start + ROWS_PER_CHUNK));
for (const row of chunk) encoded.push(row);
}
return encoded;
export function makeStoredShellSnapshotEncoder() {
const encodedRows = new WeakMap<OrchestrationV2ThreadShellJson, unknown>();

return Effect.fnUntraced(function* (stored: typeof StoredOrchestrationShellSnapshot.Type) {
let hasWorked = false;
const yieldBetweenChunks = Effect.suspend(() => {
if (hasWorked) return yieldToHost;
hasWorked = true;
return Effect.void;
});
const encodeMisses = (misses: ReadonlyArray<OrchestrationV2ThreadShellJson>) =>
Effect.gen(function* () {
yield* yieldBetweenChunks;
const chunk = yield* encodeThreadChunk(misses);
chunk.forEach((row, index) => encodedRows.set(misses[index]!, row));
});
const encodeRows = (rows: ReadonlyArray<OrchestrationV2ThreadShellJson>) =>
Effect.gen(function* () {
let misses: Array<OrchestrationV2ThreadShellJson> = [];
for (const row of rows) {
if (encodedRows.has(row)) continue;
misses.push(row);
if (misses.length === ROWS_PER_CHUNK) {
yield* encodeMisses(misses);
misses = [];
}
}
if (misses.length > 0) yield* encodeMisses(misses);
return rows.map((row) => encodedRows.get(row));
});

const { snapshot } = stored;
const threads = yield* encodeRows(snapshot.threads);
const archivedThreads = yield* encodeRows(snapshot.archivedThreads);
yield* yieldBetweenChunks;
return yield* encodeEnvelope({ ...stored, snapshot: { ...snapshot, threads, archivedThreads } });
});
const { snapshot } = stored;
const threads = yield* encodeRows(snapshot.threads);
const archivedThreads = yield* encodeRows(snapshot.archivedThreads);
if (hasWorked) yield* yieldToHost;
return yield* encodeEnvelope({
...stored,
snapshot: { ...snapshot, threads, archivedThreads },
});
});
}
Loading