diff --git a/src/builtins/commands.zig b/src/builtins/commands.zig index 1c73e1eeb6..1615e243bb 100644 --- a/src/builtins/commands.zig +++ b/src/builtins/commands.zig @@ -384,6 +384,10 @@ pub const top_level_flags = [_]TopLevelFlag{ .usage = "--resume-", .description = "Resume a session by exact ID", }, + .{ + .usage = "--sessions-v2", + .description = "Use the experimental v2 session store, also set by FX_SESSIONS_V2=1", + }, .{ .usage = "-h, --help", .description = "Display this help and exit", diff --git a/src/core/app/app_callbacks.zig b/src/core/app/app_callbacks.zig index 9e55d4b653..a687ba3e32 100644 --- a/src/core/app/app_callbacks.zig +++ b/src/core/app/app_callbacks.zig @@ -343,6 +343,10 @@ pub fn Bindings(comptime App: type) type { null else null, + .append_turn_piece = if (comptime @hasField(App, "session_persistence")) + if (app.session_persistence.v2 != null) agentAppendTurnPiece else null + else + null, .propagate_grant = agentPropagateGrant, .push_event = agentPushEvent, .push_text = agentPushText, @@ -383,7 +387,7 @@ pub fn Bindings(comptime App: type) type { @hasField(@TypeOf(app.session_persistence), "writable")) { app.session.usage.configureCheckpointSink( - if (app.session_persistence.writable != null) + if (app.session_persistence.writable != null or app.session_persistence.v2 != null) .{ .context = @ptrCast(app), .allocator = app.alloc, @@ -1210,6 +1214,11 @@ pub fn Bindings(comptime App: type) type { try app_session_runtime.Runtime(App).clearRecoveryCheckpoint(app); } + fn agentAppendTurnPiece(ctx: *anyopaque, progress: agent_runtime.TurnProgress) anyerror!void { + const app: *App = @ptrCast(@alignCast(ctx)); + try app_session_runtime.Runtime(App).appendTurnPiece(app, progress); + } + fn agentPersistUsageCheckpoint( ctx: *anyopaque, snapshot: session_usage.Snapshot, diff --git a/src/core/app/app_entry_runtime.zig b/src/core/app/app_entry_runtime.zig index a01f217067..2f256d7e77 100644 --- a/src/core/app/app_entry_runtime.zig +++ b/src/core/app/app_entry_runtime.zig @@ -381,20 +381,30 @@ fn runInteractiveWithDeps(comptime App: type, comptime cooperative: bool, app: * if (handoff_value) |value| { var handoff = value; defer handoff.deinit(alloc); - var argv = [_][]const u8{ - request.executablePath(), - "resume", - handoff.session_id, - cli_surface.upgrade_relaunch_arg, - request.previousRevision() orelse "", - }; - const argv_slice = if (request.previousRevision() == null) argv[0..4] else argv[0..5]; + var argv: [6][]const u8 = undefined; + var argc: usize = 0; + argv[argc] = request.executablePath(); + argc += 1; + // A v2 session resumes only with the flag that saved it. + if (handoff.sessions_v2) { + argv[argc] = cli_surface.sessions_v2_arg; + argc += 1; + } + for ([_][]const u8{ "resume", handoff.session_id, cli_surface.upgrade_relaunch_arg }) |arg| { + argv[argc] = arg; + argc += 1; + } + if (request.previousRevision()) |revision| { + argv[argc] = revision; + argc += 1; + } + const argv_slice = argv[0..argc]; const replace_err = deps.replace_process( deps.replace_ctx, io_mod.getIo(), .{ .argv = argv_slice }, ); - writeUpgradeRelaunchFailure(deps, replace_err, handoff.session_id); + writeUpgradeRelaunchFailure(deps, replace_err, handoff.session_id, handoff.sessions_v2); } else { writeStderr( deps, @@ -412,9 +422,10 @@ fn runInteractiveWithDeps(comptime App: type, comptime cooperative: bool, app: * &message_buffer, handoff.session_id, resume_handoff_columns, + handoff.sessions_v2, ) else - formatResumeHandoff(&message_buffer, handoff.session_id)) catch return .returned; + formatResumeHandoff(&message_buffer, handoff.session_id, handoff.sessions_v2)) catch return .returned; deps.write_stdout(deps.stdout_ctx, message) catch {}; } return .returned; @@ -439,12 +450,13 @@ fn writeUpgradeRelaunchFailure( deps: RunDeps, err: std.process.ReplaceError, session_id: []const u8, + sessions_v2: bool, ) void { var buffer: [768]u8 = undefined; const message = std.fmt.bufPrint( &buffer, - "fx: upgrade installed, but relaunch failed: {s}\nContinue session with: fx --resume {s}\n", - .{ @errorName(err), session_id }, + "fx: upgrade installed, but relaunch failed: {s}\nContinue session with: fx {s}--resume {s}\n", + .{ @errorName(err), if (sessions_v2) "--sessions-v2 " else "", session_id }, ) catch "fx: upgrade installed, but relaunch failed; run `fx doctor`.\n"; writeStderr(deps, message); } @@ -515,11 +527,11 @@ fn writeRealStdout(_: ?*anyopaque, text: []const u8) !void { try std.Io.File.stdout().writeStreamingAll(io_mod.getIo(), text); } -fn formatResumeHandoff(buffer: []u8, session_id: []const u8) ![]const u8 { +fn formatResumeHandoff(buffer: []u8, session_id: []const u8, sessions_v2: bool) ![]const u8 { return std.fmt.bufPrint( buffer, - "Continue session with: fx --resume {s}\n", - .{session_id}, + "Continue session with: fx {s}--resume {s}\n", + .{ if (sessions_v2) "--sessions-v2 " else "", session_id }, ); } @@ -665,6 +677,7 @@ const TestCapture = struct { record_stderr_event: bool = false, record_stdout_event: bool = false, resume_handoff_id: ?[]const u8 = null, + resume_handoff_sessions_v2: bool = false, shutdown_failure: ?anyerror = null, raise_sigint_during_deinit: bool = false, upgrade_relaunch_path: ?[]const u8 = null, @@ -805,7 +818,7 @@ const TestApp = struct { self.deinit(); return .{ .failure = error.OutOfMemory }; }; - break :blk .{ .session_id = session_id }; + break :blk .{ .session_id = session_id, .sessions_v2 = active_capture.?.resume_handoff_sessions_v2 }; } else null; if (active_capture.?.raise_sigint_during_deinit) { _ = std.c.raise(std.posix.SIG.INT); @@ -1030,6 +1043,26 @@ test "app entry bounds graceful-exit SIGINT suppression to handoff lifetime" { try std.testing.expectEqual(@as(usize, 1), test_sigint_count.load(.seq_cst)); } +test "a v2 handoff relaunches and hints with --sessions-v2" { + const alloc = std.testing.allocator; + var capture = TestCapture.init(.{ .interactive = .{} }); + defer capture.deinit(); + capture.resume_handoff_id = "session-123"; + capture.resume_handoff_sessions_v2 = true; + capture.upgrade_relaunch_path = "/tmp/fx-upgraded"; + + const outcome = try runWithDeps(TestApp, alloc, &.{}, testConfig(), capture.deps()); + + try std.testing.expectEqual(@as(u8, 1), outcome.exit); + try std.testing.expectEqual(@as(usize, 5), capture.replace_arg_count); + try std.testing.expectEqualStrings("/tmp/fx-upgraded", capture.replaceArg(0)); + try std.testing.expectEqualStrings("--sessions-v2", capture.replaceArg(1)); + try std.testing.expectEqualStrings("resume", capture.replaceArg(2)); + try std.testing.expectEqualStrings("session-123", capture.replaceArg(3)); + try std.testing.expectEqualStrings("--upgrade-relaunch", capture.replaceArg(4)); + try std.testing.expect(std.mem.find(u8, capture.stderr.written(), "fx --sessions-v2 --resume session-123") != null); +} + test "app entry relaunches only after teardown with the validated handoff" { const alloc = std.testing.allocator; var capture = TestCapture.init(.{ .interactive = .{} }); diff --git a/src/core/app/app_session_runtime.zig b/src/core/app/app_session_runtime.zig index dbf1197125..f06efd9f30 100644 --- a/src/core/app/app_session_runtime.zig +++ b/src/core/app/app_session_runtime.zig @@ -1,5 +1,6 @@ const std = @import("std"); const worker_runtime = @import("../agent/worker_runtime.zig"); +const agent_runtime = @import("../agent/agent_runtime.zig"); const auto_classifier_context = @import("../permissions/auto_classifier_context.zig"); const builtin = @import("builtin"); const build_options = @import("build_options"); @@ -46,6 +47,7 @@ const tool_result_errors = @import("../tooling/tool_result_errors.zig"); const session_display_metadata = @import("../session/session_display_metadata.zig"); const session_title_generation = @import("../session/session_title_generation.zig"); const session_log = @import("../session/session_log.zig"); +const session_adapter = @import("../session/session_adapter.zig"); const session_store = @import("../session/session_store.zig"); const session_catalog_cache = @import("../session/session_catalog_cache.zig"); const session_summary_codec = @import("../session/session_summary_codec.zig"); @@ -404,6 +406,8 @@ test "resume handoff policy requires requested durable non-pristine state" { /// Owns `session_id`; callers must release it with `deinit`. pub const ResumeHandoff = struct { session_id: []u8, + /// The session is in `v2`: resuming it needs `--sessions-v2`. + sessions_v2: bool = false, pub fn deinit(self: *ResumeHandoff, alloc: Allocator) void { alloc.free(self.session_id); @@ -807,6 +811,19 @@ const SessionPickerLoad = struct { } }; + /// Where a listing comes from: v1's store, or the v2 store (D26). + const Source = union(enum) { + v1: *const session_store.Store, + v2: struct { store: *session_adapter.Store, workspace_root: []const u8 }, + + fn workspaceRoot(source: Source) []const u8 { + return switch (source) { + .v1 => |store| store.workspace_root, + .v2 => |v2| v2.workspace_root, + }; + } + }; + const Task = struct { thread: ?std.Thread = null, done: std.atomic.Value(bool) = std.atomic.Value(bool).init(false), @@ -818,6 +835,8 @@ const SessionPickerLoad = struct { request: PageRequest, catalog: ?session_catalog_cache.ActionableSessionCatalog = null, cache_writer: ?session_catalog_cache.Writer = null, + /// Shared with the app, which joins this thread before freeing it. + v2_store: ?*session_adapter.Store = null, failure: ?anyerror = null, fn requestStop(self: *Task) void { @@ -876,14 +895,14 @@ const SessionPickerLoad = struct { fn schedule( self: *SessionPickerLoad, - store: *const session_store.Store, + source: Source, request: PageRequest, ) !void { if (self.task != null) { self.replacePending(request); return; } - try self.start(store, request); + try self.start(source, request); } fn replacePending(self: *SessionPickerLoad, request: PageRequest) void { @@ -947,28 +966,31 @@ const SessionPickerLoad = struct { return null; } - fn startPending(self: *SessionPickerLoad, store: *const session_store.Store) !?u64 { + fn startPending(self: *SessionPickerLoad, source: Source) !?u64 { const request = self.pending orelse return null; self.pending = null; const generation = request.generation; - try self.start(store, request); + try self.start(source, request); return generation; } fn start( self: *SessionPickerLoad, - store: *const session_store.Store, + source: Source, request: PageRequest, ) !void { std.debug.assert(self.task == null); const alloc = std.heap.c_allocator; var owned_request = request; - const home_dir = alloc.dupe(u8, store.home_dir) catch |err| { + const home_dir = alloc.dupe(u8, switch (source) { + .v1 => |store| store.home_dir, + .v2 => |v2| v2.store.home, + }) catch |err| { owned_request.deinit(); return err; }; - const workspace_root = alloc.dupe(u8, store.workspace_root) catch |err| { + const workspace_root = alloc.dupe(u8, source.workspaceRoot()) catch |err| { alloc.free(home_dir); owned_request.deinit(); return err; @@ -984,10 +1006,13 @@ const SessionPickerLoad = struct { .workspace_root = workspace_root, .request = owned_request, }; - task.cache_writer = session_catalog_cache.Writer.init(store.*) catch |err| blk: { - debug_trace.logf("core", "session catalog cache writer unavailable err={s}", .{@errorName(err)}); - break :blk null; - }; + switch (source) { + .v1 => |store| task.cache_writer = session_catalog_cache.Writer.init(store.*) catch |err| blk: { + debug_trace.logf("core", "session catalog cache writer unavailable err={s}", .{@errorName(err)}); + break :blk null; + }, + .v2 => |v2| task.v2_store = v2.store, + } task.thread = std.Thread.spawn(.{}, threadMain, .{task}) catch |err| { if (task.cache_writer) |*writer| writer.deinit(); task.request.deinit(); @@ -1009,6 +1034,21 @@ const SessionPickerLoad = struct { fn threadMain(task: *Task) void { defer task.done.store(true, .release); const started = io_mod.nanoTimestamp(); + if (task.v2_store) |store| { + const summaries = session_adapter.listSummaries( + store, + std.heap.c_allocator, + task.request.active_id, + &task.cancel_requested, + ) catch |err| { + debug_trace.logf("core", "session picker catalog stopped backend=v2 err={s}", .{@errorName(err)}); + task.failure = err; + return; + }; + task.catalog = .{ .summaries = summaries }; + debug_trace.logf("core", "session picker catalog loaded backend=v2 sessions={d} elapsed_us={d}", .{ summaries.items.len, @divTrunc(io_mod.nanoTimestamp() - started, std.time.ns_per_us) }); + return; + } var read_only = session_store.Store.initReadOnlyFromHome( std.heap.c_allocator, task.home_dir, @@ -1262,12 +1302,17 @@ pub const Persistence = struct { pending_live_session_policy: ?BackgroundSessionPolicy = null, pending_live_session_wait: ?LiveSessionWait = null, shutdown_failure: ?anyerror = null, + /// `--sessions-v2` for this process; with it (or FX_SESSIONS_V2=1) the + /// session lives in `v2` and `store` and `writable` stay null (D26). + sessions_v2: bool = false, + v2_store: ?session_adapter.Store = null, + v2: ?*session_adapter.Session = null, /// Fieldwise initialization avoids retaining undefined optional payloads /// in a static release-binary template. pub fn initInto(storage: *Persistence) void { comptime { - if (std.meta.fields(Persistence).len != 25) { + if (std.meta.fields(Persistence).len != 28) { @compileError("update Persistence.initInto for the changed field set"); } } @@ -1297,6 +1342,9 @@ pub const Persistence = struct { storage.pending_live_session_policy = null; storage.pending_live_session_wait = null; storage.shutdown_failure = null; + storage.sessions_v2 = false; + storage.v2_store = null; + storage.v2 = null; } pub fn deinit(self: *Persistence, alloc: Allocator) void { @@ -1315,6 +1363,7 @@ pub const Persistence = struct { if (self.subagent_host) |host| host.deinit(); if (self.writable) |*loaded| loaded.deinit(alloc); if (self.store) |*store| store.deinit(alloc); + if (self.v2) |v2| v2.close(); if (self.workspace_preferences) |*preferences| preferences.deinit(alloc); if (self.session_preferences) |*preferences| preferences.deinit(alloc); if (self.js_host_session) |*owner| owner.deinit(alloc); @@ -1323,6 +1372,8 @@ pub const Persistence = struct { self.session_picker_load.deinit(); self.session_picker_cache.deinit(); self.title_generation.deinit(); + // After the picker thread, which lists through it, has joined. + if (self.v2_store) |*store| store.deinit(alloc); self.* = undefined; } }; @@ -1388,14 +1439,14 @@ pub fn Runtime(comptime App: type) type { } fn imageSnapshotStorageDir(app: *App) ![]u8 { - const sessions_dir = if (app.session_persistence.store) |*store| + // A v2 session keeps its images with its side files. + const sessions_dir = if (app.session_persistence.v2) |v2| + std.fs.path.dirname(v2.filesPath()) + else if (app.session_persistence.store) |*store| store.sessions_dir else null; - const session_id = if (app.session_persistence.writable) |*writable| - writable.active_id - else - null; + const session_id = activeSessionId(app); return session_store.imageSnapshotStorageDir( app.alloc, sessions_dir, @@ -1459,6 +1510,14 @@ pub fn Runtime(comptime App: type) type { required: bool, ) !void { if (comptime !runtime_profile.allows(App, .durable_sessions)) return; + // One backend per process (D26): with v2 on, v1's store stays closed. + if (session_adapter.enabled(app.session_persistence.sessions_v2)) { + app.session_persistence.v2_store = session_adapter.Store.openFromEnv(app.alloc) catch |err| { + if (required) return err; + return; + }; + return; + } var store = session_store.Store.init( app.alloc, app.workspace_root, @@ -1473,6 +1532,14 @@ pub fn Runtime(comptime App: type) type { pub fn enableSessionStores(app: *App) void { if (comptime !runtime_profile.allows(App, .durable_sessions)) return; + if (app.session_persistence.v2) |v2| { + if (comptime @hasDecl(@TypeOf(app.session), "configureWebFetchArtifacts")) { + app.session.configureWebFetchArtifacts(app.alloc, v2.filesPath()); + } + // Subagents move to the parent log with their own PR (D22). + debug_trace.logf("subagent", "interactive subagent host unavailable session={s} reason=sessions_v2", .{v2.id()}); + return; + } const loaded = if (app.session_persistence.writable) |*value| value else return; const capability = loaded.childCapability() catch |err| { debug_trace.logf( @@ -1520,6 +1587,7 @@ pub fn Runtime(comptime App: type) type { try beginFreshJsHostSession(app); } if (comptime !runtime_profile.allows(App, .durable_sessions)) return; + if (app.session_persistence.v2_store) |*v2_store| return beginFreshV2Session(app, v2_store); const store = app.session_persistence.store orelse return; const preferences = app.session_persistence.workspace_preferences orelse return error.SessionPreferencesUnavailable; @@ -1545,6 +1613,32 @@ pub fn Runtime(comptime App: type) type { app.total_web_search_requests = 0; } + /// A v2 session starts in memory; its folder appears with the first + /// turn (D2), so a launch that never prompts leaves nothing behind. + fn beginFreshV2Session(app: *App, v2_store: *session_adapter.Store) !void { + const preferences = app.session_persistence.workspace_preferences orelse + return error.SessionPreferencesUnavailable; + try replacePreferences( + app.alloc, + &app.session_persistence.session_preferences, + preferences, + ); + var permission_state = try app.session.snapshotPermissionState(app.alloc); + defer permission_state.deinit(app.alloc); + app.session_persistence.v2 = session_adapter.Session.create(app.alloc, v2_store, app.workspace_root, .app, .{ + .preferences = preferences, + .language = app.session.languageSnapshot(), + .permission_state = permission_state, + }) catch |err| { + try warnNonDurable(app, "session creation failed", err); + return; + }; + app.session_persistence.degraded_warning_emitted = false; + app.total_input_tokens = 0; + app.total_output_tokens = 0; + app.total_web_search_requests = 0; + } + fn beginFreshJsHostSession(app: *App) !void { const preferences = app.session_persistence.workspace_preferences orelse return error.SessionPreferencesUnavailable; @@ -1920,6 +2014,20 @@ pub fn Runtime(comptime App: type) type { app.requested_resume = null; defer target.deinit(app.alloc); + if (app.session_persistence.v2_store) |*v2_store| { + const v2_target: session_adapter.Target = switch (target) { + .pick => return openSessionPicker(app), + .last => .last, + // The manager records the session each host last opened. + .remembered => .last_opened, + .id => |session_id| .{ .id = session_id }, + }; + const v2 = try session_adapter.Session.resumeSession(app.alloc, v2_store, v2_target, app.workspace_root, .app); + try installResumedV2Session(app, v2, notice); + errdefer closeWritableSession(app); + try app.commitStartupResumeReplayAnchor(); + return; + } var remembered: ?[]u8 = null; defer if (remembered) |id| app.alloc.free(id); const resume_target: session_store.ResumeTarget = switch (target) { @@ -2105,6 +2213,19 @@ pub fn Runtime(comptime App: type) type { pub fn resumeSelectedSession(app: *App) !bool { const selected_id = app.session_persistence.session_picker.selectedId() orelse return false; + if (app.session_persistence.v2_store) |*v2_store| { + // A session open in another fx shows as busy at once, as + // v1's picker does (D38). + const v2 = try session_adapter.Session.resumeSessionWithoutWaiting(app.alloc, v2_store, .{ .id = selected_id }, app.workspace_root, .app); + var v2_owned = true; + errdefer if (v2_owned) v2.close(); + try app.prepareLiveSessionResume(); + v2_owned = false; + try installResumedV2Session(app, v2, .session); + startResumedSessionReconciliation(app); + try app.finishLiveSessionResume(); + return true; + } const log_options = session_log.Options{ .session_lock_deadline_ms = 0, }; @@ -2169,6 +2290,27 @@ pub fn Runtime(comptime App: type) type { enableSessionStores(app); } + /// Makes `v2` the open session, taking ownership, and restores the app + /// from it through v1's state, as `installResumedSession` does. + fn installResumedV2Session(app: *App, v2: *session_adapter.Session, notice: ResumeNotice) !void { + var v2_owned = true; + errdefer if (v2_owned) v2.close(); + var resumed = try v2.durableState(app.alloc, app.workspace_root); + defer resumed.deinit(app.alloc); + var display: session_display_metadata.DisplayMetadata = if (resumed.title) |title| + .{ .present = true, .title = try app.alloc.dupe(u8, title) } + else + try session_display_metadata.deriveFromHistory(app.alloc, resumed.state.history); + defer display.deinit(app.alloc); + + closeWritableSession(app); + app.session_persistence.v2 = v2; + v2_owned = false; + errdefer closeWritableSession(app); + try hydrateResumedSession(app, resumed.state, &display, notice); + enableSessionStores(app); + } + /// Record whether restored history references shell execution handles /// this process does not own. Registry membership, not the resume /// itself, decides staleness (see session_runtime.detectStaleShellHandles). @@ -2355,6 +2497,21 @@ pub fn Runtime(comptime App: type) type { return session_log.recoveryWasAsked(&loaded.log.dir); } + /// The picker lists from the backend this process runs on. + fn pickerSource(app: *App) ?SessionPickerLoad.Source { + if (app.session_persistence.v2_store) |*store| { + return .{ .v2 = .{ .store = store, .workspace_root = app.workspace_root } }; + } + if (app.session_persistence.store) |*store| return .{ .v1 = store }; + return null; + } + + /// The open session, which the picker leaves out. + fn pickerActiveId(app: *App) ?[]const u8 { + if (app.session_persistence.v2) |v2| return v2.id(); + return if (app.session_persistence.writable) |*loaded| loaded.active_id else null; + } + pub fn openSessionPicker(app: *App) !void { return openSessionPickerWithScope(app, .current_workspace); } @@ -2380,7 +2537,7 @@ pub fn Runtime(comptime App: type) type { loader.allocateGeneration(), active_id, ) catch return; - loader.schedule(store, request) catch |err| { + loader.schedule(.{ .v1 = store }, request) catch |err| { debug_trace.logf( "core", "session catalog preload unavailable err={s}", @@ -2410,18 +2567,13 @@ pub fn Runtime(comptime App: type) type { picker.load_state = .loading; picker.scope = scope; picker.setQuery(app.input_runtime.edit_state.input.items); - const store = if (app.session_persistence.store) |*value| - value - else { + const source = pickerSource(app) orelse { picker.load_state = .failed; try writeSessionPickerError(app, error.SessionStoreUnavailable); app.shell.render_requests.request(.footer); return; }; - const active_id = if (app.session_persistence.writable) |*loaded| - loaded.active_id - else - null; + const active_id = pickerActiveId(app); const limit = if (comptime @hasField(App, "shell")) resumePageLimitForRows(app.shell.layout.rows) else @@ -2430,8 +2582,9 @@ pub fn Runtime(comptime App: type) type { const matching = loader.matchingInitialGeneration(active_id); if (matching != previous_generation) loader.cancelGeneration(previous_generation); const cache = &app.session_persistence.session_picker_cache; - if (!cache.matches(active_id)) { - installStaleDiskCatalog(store, cache, active_id) catch |err| { + // v2 keeps no catalog file: its listing is the index itself. + if (source == .v1 and !cache.matches(active_id)) { + installStaleDiskCatalog(source.v1, cache, active_id) catch |err| { debug_trace.logf( "core", "session picker disk catalog unavailable err={s}", @@ -2445,7 +2598,7 @@ pub fn Runtime(comptime App: type) type { picker, app.alloc, cache, - store.workspace_root, + source.workspaceRoot(), null, limit, false, @@ -2470,7 +2623,7 @@ pub fn Runtime(comptime App: type) type { return; }; picker.generation = request.generation; - loader.schedule(store, request) catch |err| { + loader.schedule(source, request) catch |err| { if (!cache_visible) picker.load_state = .failed; try writeSessionPickerError(app, err); return; @@ -2494,8 +2647,8 @@ pub fn Runtime(comptime App: type) type { var task_owned = true; defer if (task_owned) task.deinit(); - const active_id = if (app.session_persistence.writable) |*loaded| loaded.active_id else null; - const valid = app.session_persistence.store != null and + const active_id = pickerActiveId(app); + const valid = pickerSource(app) != null and !task.abandoned and optionalStringEql(task.request.active_id, active_id); var cache_installed = false; if (valid and task.failure == null) { @@ -2532,7 +2685,7 @@ pub fn Runtime(comptime App: type) type { picker, app.alloc, cache, - app.session_persistence.store.?.workspace_root, + pickerSource(app).?.workspaceRoot(), null, limit, false, @@ -2545,12 +2698,9 @@ pub fn Runtime(comptime App: type) type { task.deinit(); task_owned = false; - const store = if (app.session_persistence.store) |*value| - value - else - return; + const source = pickerSource(app) orelse return; const pending_generation = loader.pendingGeneration(); - _ = loader.startPending(store) catch |err| { + _ = loader.startPending(source) catch |err| { if (pending_generation) |generation| { if (picker.active and picker.generation == generation) { if (picker.loading_more) { @@ -2594,10 +2744,8 @@ pub fn Runtime(comptime App: type) type { const picker = &app.session_persistence.session_picker; const continuation = picker.continuation orelse return error.SessionStoreUnavailable; - const active_id = if (app.session_persistence.writable) |*loaded| - loaded.active_id - else - null; + const source = pickerSource(app) orelse return error.SessionStoreUnavailable; + const active_id = pickerActiveId(app); const cache = &app.session_persistence.session_picker_cache; if (!cache.matches(active_id)) return error.SessionStoreUnavailable; const limit = if (comptime @hasField(App, "shell")) @@ -2610,7 +2758,7 @@ pub fn Runtime(comptime App: type) type { picker, app.alloc, cache, - app.session_persistence.store.?.workspace_root, + source.workspaceRoot(), continuation.view(), limit, true, @@ -2653,7 +2801,7 @@ pub fn Runtime(comptime App: type) type { pub fn cancelSessionPickerToComposer(app: *App) void { cancelSessionPicker(app); if (comptime !@hasField(App, "session")) return; - if (app.session_persistence.writable != null) return; + if (app.session_persistence.writable != null or app.session_persistence.v2 != null) return; beginFreshPersistedSession(app) catch |err| { debug_trace.logf( "session", @@ -2741,6 +2889,7 @@ pub fn Runtime(comptime App: type) type { } app.session_persistence.write_mutex.lockUncancelable(io_mod.getIo()); defer app.session_persistence.write_mutex.unlock(io_mod.getIo()); + if (app.session_persistence.v2) |v2| return v2.persistUsage(snapshot); const store = if (app.session_persistence.store) |value| value else @@ -3130,6 +3279,19 @@ pub fn Runtime(comptime App: type) type { var remember_failure: ?RememberFailure = null; defer if (remember_failure) |failure| reportRememberFailure(app, failure, false); app.session_persistence.write_mutex.lockUncancelable(io_mod.getIo()); + if (app.session_persistence.v2) |v2| { + defer app.session_persistence.write_mutex.unlock(io_mod.getIo()); + if (try commitV2HistoryTurn(app, v2, &prepared, mode)) |outcome| return outcome; + if (comptime @hasDecl(@TypeOf(app.session), "commitPreparedHistoryEntry")) { + app.session.commitPreparedHistoryEntry(app.alloc, prepared); + prepared_owned = false; + } else { + try app.session.appendHistoryEntry(app.alloc, prepared); + } + if (snapshot_file_ownership) |ownership| ownership.transfer(); + ensureCachedSessionTitle(app) catch {}; + return .committed; + } const loaded = if (app.session_persistence.writable) |*value| value else { @@ -3222,6 +3384,51 @@ pub fn Runtime(comptime App: type) type { return .committed; } + /// Saves a finished turn to v2 with v1's failure modes. Null means the + /// caller commits the turn in memory: it was saved, or its failure + /// was recorded and shown. + fn commitV2HistoryTurn( + app: *App, + v2: *session_adapter.Session, + prepared: *types.HistoryTurn, + mode: AppendHistoryMode, + ) !?HistoryAppendOutcome { + const failure: anyerror = saved: { + v2.prepareTurn(prepared) catch |err| break :saved err; + v2.commitTurn(prepared.*, app.session.languageSnapshot()) catch |err| break :saved err; + return null; + }; + switch (mode) { + .strict => return failure, + .visual_epoch => { + debug_trace.logf("session", "visual epoch history not committed err={s}", .{@errorName(failure)}); + return .uncommitted; + }, + .finished_prompt => { + recordShutdownFailure(app, failure); + if (comptime @hasDecl(App, "writeDomainNotice")) { + const body = try std.fmt.allocPrint( + app.alloc, + "Turn completed, but fx could not save it ({s}). The session keeps running; this turn may be missing after a resume.", + .{@errorName(failure)}, + ); + defer app.alloc.free(body); + app.writeDomainNotice(.{ .topic = "session", .tone = .@"error", .body = body }, true) catch {}; + } + return null; + }, + } + } + + /// `AgentRuntimeDeps.append_turn_piece` on v2: the pieces finished so + /// far reach the log before the next request (D12). + pub fn appendTurnPiece(app: *App, progress: agent_runtime.TurnProgress) !void { + app.session_persistence.write_mutex.lockUncancelable(io_mod.getIo()); + defer app.session_persistence.write_mutex.unlock(io_mod.getIo()); + const v2 = app.session_persistence.v2 orelse return; + try v2.appendProgress(progress.user, progress.execution); + } + pub fn commitRuntimePreferences( app: *App, patch: SessionPreferencePatch, @@ -3266,6 +3473,15 @@ pub fn Runtime(comptime App: type) type { app.session_persistence.write_mutex.lockUncancelable(io_mod.getIo()); defer app.session_persistence.write_mutex.unlock(io_mod.getIo()); + if (app.session_persistence.v2) |v2| { + if (result.session_error != null) return result; + const preferences = app.session_persistence.session_preferences orelse return result; + v2.setPreferences(preferences) catch |err| { + result.session_error = err; + warnDegraded(app, err) catch {}; + }; + return result; + } const loaded = if (app.session_persistence.writable) |*value| value else @@ -3291,6 +3507,7 @@ pub fn Runtime(comptime App: type) type { } pub fn activeSessionId(app: *App) ?[]const u8 { + if (app.session_persistence.v2) |v2| return v2.id(); if (app.session_persistence.writable) |*loaded| return loaded.active_id; if (app.session_persistence.js_host_session) |*owner| return owner.state.id; return null; @@ -3464,7 +3681,20 @@ pub fn Runtime(comptime App: type) type { { app.session_persistence.write_mutex.lockUncancelable(io_mod.getIo()); defer app.session_persistence.write_mutex.unlock(io_mod.getIo()); - if (app.session_persistence.writable) |*loaded| { + if (app.session_persistence.v2) |v2| { + installed = v2.installGeneratedTitle(app.session.agent.history.items, title) catch |err| { + debug_trace.logf( + "session", + "event=title_generation_apply result=failed session={s} err={s}", + .{ task.session_id, @errorName(err) }, + ); + app.session_persistence.title_generation.recordDropped(.install_failed, @errorName(err)); + return false; + }; + if (!installed) { + app.session_persistence.title_generation.recordDropped(.user_title_present, ""); + } + } else if (app.session_persistence.writable) |*loaded| { installed = session_title_generation.installGeneratedTitle( app.alloc, loaded, @@ -3522,12 +3752,17 @@ pub fn Runtime(comptime App: type) type { pub fn renameActiveSession(app: *App, raw: []const u8) !void { const title = try validateSessionTitle(raw); if (comptime !@hasField(App, "session_persistence")) return error.NoActiveSession; - if (app.session_persistence.writable == null) return error.NoActiveSession; + if (app.session_persistence.writable == null and app.session_persistence.v2 == null) return error.NoActiveSession; try setCachedSessionTitle(app, title); app.session_persistence.write_mutex.lockUncancelable(io_mod.getIo()); defer app.session_persistence.write_mutex.unlock(io_mod.getIo()); + if (app.session_persistence.v2) |v2| { + try v2.rename(title); + invalidateSessionPickerCaches(app); + return; + } const loaded = &app.session_persistence.writable.?; if (!try loaded.renameConversation(app.alloc, title)) { @@ -3540,6 +3775,7 @@ pub fn Runtime(comptime App: type) type { app: *App, alloc: Allocator, ) !?[]u8 { + if (app.session_persistence.v2) |v2| return try v2.folderPath(alloc); const loaded = if (app.session_persistence.writable) |*value| value else @@ -3553,11 +3789,15 @@ pub fn Runtime(comptime App: type) type { } pub fn childCapability(app: *App) ?*session_child_store.SessionChildCapability { - const loaded = if (app.session_persistence.writable) |*value| - value - else - return null; - return loaded.childCapability() catch null; + if (app.session_persistence.writable == null and app.session_persistence.v2 == null) return null; + return activeChildCapability(app) catch null; + } + + /// The open session's side-file capability, from either backend. + fn activeChildCapability(app: *App) !*session_child_store.SessionChildCapability { + if (app.session_persistence.v2) |v2| return v2.childCapability(); + const loaded = if (app.session_persistence.writable) |*value| value else return error.SessionPersistenceUnavailable; + return loaded.childCapability(); } pub fn subagentHost(app: *App) ?*subagent_tool_host.Runtime { @@ -3616,6 +3856,7 @@ pub fn Runtime(comptime App: type) type { pub fn prepareResumeHandoff(app: *App) !void { app.session_persistence.write_mutex.lockUncancelable(io_mod.getIo()); defer app.session_persistence.write_mutex.unlock(io_mod.getIo()); + if (app.session_persistence.v2) |v2| return settleV2Usage(app, v2); const loaded = if (app.session_persistence.writable) |*value| value else @@ -4831,11 +5072,8 @@ pub fn Runtime(comptime App: type) type { return false; }, }; - const loaded = if (app.session_persistence.writable) |*value| - value - else - return false; - const capability = loaded.childCapability() catch |err| { + if (app.session_persistence.writable == null and app.session_persistence.v2 == null) return false; + const capability = activeChildCapability(app) catch |err| { debug_trace.logf( "session", "resume command replay capability unavailable handle_bytes={d} err={s}", @@ -5055,11 +5293,8 @@ pub fn Runtime(comptime App: type) type { result: types.PersistedToolResult, ) !bool { const handle = result.output_handle orelse return false; - const loaded = if (app.session_persistence.writable) |*value| - value - else - return false; - const capability = loaded.childCapability() catch |err| { + if (app.session_persistence.writable == null and app.session_persistence.v2 == null) return false; + const capability = activeChildCapability(app) catch |err| { debug_trace.logf( "session", "resume command replay could not open result handle={s} err={s}", @@ -5265,6 +5500,7 @@ pub fn Runtime(comptime App: type) type { defer app.session_persistence.write_mutex.unlock(io_mod.getIo()); app.session_persistence.remember_fresh_session = false; discardAnyPendingCancelledCommand(app, "writable_session_close"); + if (app.session_persistence.v2) |v2| return closeV2Session(app, v2, handoff_intent); const loaded = if (app.session_persistence.writable) |*value| value else @@ -5337,6 +5573,45 @@ pub fn Runtime(comptime App: type) type { return handoff; } + /// Settles usage and closes; the manager ends an open turn as + /// `closed`, and a session with no turn leaves nothing on disk (D2). + fn closeV2Session(app: *App, v2: *session_adapter.Session, handoff_intent: ResumeHandoffIntent) ?ResumeHandoff { + var settled = true; + settleV2Usage(app, v2) catch |err| { + settled = false; + recordShutdownFailure(app, err); + debug_trace.logf("session", "final persistence settlement failed session={s} err={s}", .{ v2.id(), @errorName(err) }); + }; + const create_handoff = settled and shouldCreateResumeHandoff(.{ + .intent = handoff_intent, + .has_writable_session = true, + .is_pristine = !v2.saved(), + }); + const handoff: ?ResumeHandoff = if (create_handoff) blk: { + const session_id = app.alloc.dupe(u8, v2.id()) catch |err| { + debug_trace.logf("session", "resume handoff unavailable session={s} err={s}", .{ v2.id(), @errorName(err) }); + break :blk null; + }; + break :blk .{ .session_id = session_id, .sessions_v2 = true }; + } else null; + if (comptime @hasDecl(@TypeOf(app.session), "clearWebFetchArtifacts")) { + app.session.clearWebFetchArtifacts(); + } + disableSubagentHost(app); + v2.close(); + app.session_persistence.v2 = null; + return handoff; + } + + fn settleV2Usage(app: *App, v2: *session_adapter.Session) !void { + if (comptime !@hasField(@TypeOf(app.session), "usage")) return; + if (!app.session.usage.isDirty()) return; + var usage = try app.session.usage.snapshot(app.alloc); + defer usage.deinit(app.alloc); + try v2.persistUsage(usage); + app.session.usage.markClean(usage); + } + /// A fresh interactive session that never received durable work has /// nothing to resume. Discarding it on close keeps each launch from /// leaving an empty session directory behind, matching `fx ask` and @@ -5504,7 +5779,9 @@ pub fn Runtime(comptime App: type) type { if (comptime @hasField(App, "session_persistence")) { app.session_persistence.write_mutex.lockUncancelable(io_mod.getIo()); defer app.session_persistence.write_mutex.unlock(io_mod.getIo()); - if (app.session_persistence.writable) |*loaded| { + if (app.session_persistence.v2) |v2| { + try v2.commitCompaction(summary, active_prefix != null, retained_from); + } else if (app.session_persistence.writable) |*loaded| { _ = try loaded.commitContextCompaction( app.alloc, summary, @@ -5539,6 +5816,7 @@ pub fn Runtime(comptime App: type) type { ) !void { app.session_persistence.write_mutex.lockUncancelable(io_mod.getIo()); defer app.session_persistence.write_mutex.unlock(io_mod.getIo()); + if (app.session_persistence.v2) |v2| return v2.setPermissions(permission_state); const loaded = if (app.session_persistence.writable) |*value| value else diff --git a/src/core/cli/cli_surface.zig b/src/core/cli/cli_surface.zig index 1c4010ffbb..83f6d7f2af 100644 --- a/src/core/cli/cli_surface.zig +++ b/src/core/cli/cli_surface.zig @@ -94,6 +94,7 @@ const ResumeInvocation = struct { const resume_id_alias_prefix = "--resume-"; pub const upgrade_relaunch_arg = "--upgrade-relaunch"; +pub const sessions_v2_arg = "--sessions-v2"; pub const UpgradeRelaunch = struct { previous_revision: ?[]u8 = null, @@ -408,7 +409,7 @@ fn parseGlobalLaunchArgs( var index: usize = 0; while (index < args.len) { const arg = args[index]; - if (std.mem.eql(u8, arg, "--sessions-v2")) { + if (std.mem.eql(u8, arg, sessions_v2_arg)) { sessions_v2 = true; } else if (std.mem.eql(u8, arg, "--context-limit")) { index += 1; @@ -531,7 +532,7 @@ pub fn argsAfterGlobalLaunchArgs(args: []const [:0]const u8) []const [:0]const u !std.mem.eql(u8, arg, "--no-fast") and !std.mem.eql(u8, arg, "--provider-strict") and !std.mem.eql(u8, arg, "--no-provider-strict") and - !std.mem.eql(u8, arg, "--sessions-v2")) + !std.mem.eql(u8, arg, sessions_v2_arg)) { return args[index..]; } diff --git a/src/core/session/session_adapter.zig b/src/core/session/session_adapter.zig index ae582c23b3..ab0b1df4c4 100644 --- a/src/core/session/session_adapter.zig +++ b/src/core/session/session_adapter.zig @@ -124,6 +124,8 @@ pub const Target = union(enum) { id: []const u8, /// The newest updated root session in the workspace. last, + /// The session this host last opened in the workspace (`fx -c`). + last_opened, }; /// What resume gives back. Owns everything; free with `deinit`. @@ -133,10 +135,14 @@ pub const Restored = struct { preferences: ?session_codec.DurableSessionPreferences = null, permission_state: ?session_permission_state.State = null, usage: ?session_usage.Snapshot = null, + /// The stored title, generated or chosen by the user. + title: ?[]u8 = null, created_at_ms: i64, + updated_at_ms: i64 = 0, pub fn deinit(restored: *Restored, alloc: Allocator) void { types.freeHistoryTurnSlice(alloc, restored.history); + if (restored.title) |value| alloc.free(value); if (restored.preferences) |*value| value.deinit(alloc); if (restored.permission_state) |*value| value.deinit(alloc); if (restored.usage) |*value| value.deinit(alloc); @@ -144,6 +150,18 @@ pub const Restored = struct { } }; +/// A resumed session as v1's state. Owns everything; free with `deinit`. +pub const Resumed = struct { + state: session_codec.DurableSessionState, + title: ?[]u8, + + pub fn deinit(resumed: *Resumed, alloc: Allocator) void { + resumed.state.deinit(alloc); + if (resumed.title) |value| alloc.free(value); + resumed.* = undefined; + } +}; + pub const Session = struct { alloc: Allocator, store: *Store, @@ -194,13 +212,26 @@ pub const Session = struct { } pub fn resumeSession(alloc: Allocator, store: *Store, target: Target, workspace: []const u8, host: Host) !*Session { + return resumeWaiting(alloc, store, target, workspace, host, null); + } + + /// As `resumeSession`, but SessionBusy at once when another process has + /// the session open, for the session picker (D38). + pub fn resumeSessionWithoutWaiting(alloc: Allocator, store: *Store, target: Target, workspace: []const u8, host: Host) !*Session { + return resumeWaiting(alloc, store, target, workspace, host, 0); + } + + /// `lock_wait_ms` null waits the manager's 2 s for the writer lock. + fn resumeWaiting(alloc: Allocator, store: *Store, target: Target, workspace: []const u8, host: Host, lock_wait_ms: ?u64) !*Session { const handle = store.manager.openResume(.{ .target = switch (target) { .id => |session_id| .{ .id = session_id }, .last => .last, + .last_opened => .{ .last_opened = host }, }, .workspace = workspace, .host = host, + .lock_wait_ms = lock_wait_ms, }) catch |err| return resumeError(err, target); errdefer handle.release(); const self = try init(alloc, store, handle, false); @@ -228,6 +259,17 @@ pub const Session = struct { return self.last_turn > 0; } + /// The session's folder under `~/.fx/sessions/v2`, for display. Caller owns it. + pub fn folderPath(self: *const Session, alloc: Allocator) ![]u8 { + return std.fs.path.join(alloc, &.{ + self.store.home, + profile_paths.root_dir_name, + profile_paths.sessions_dir_name, + session_layout.sessions_v2_dir, + self.id(), + }); + } + /// `~/.fx/session-files/{id}`: borrowed until `close`. pub fn filesPath(self: *const Session) []const u8 { return self.files_path; @@ -553,6 +595,17 @@ pub const Session = struct { _ = try self.handle.append(&.{.{ .set = .{ .key = .permissions, .value = value } }}); } + /// A title the user chose. It differs from the derived one, so a + /// generated title never replaces it. + pub fn rename(self: *Session, title: []const u8) !void { + self.mutex.lockUncancelable(io_mod.getIo()); + defer self.mutex.unlock(io_mod.getIo()); + const value = try jsonString(self.alloc, title); + defer self.alloc.free(value); + _ = try self.handle.append(&.{.{ .set = .{ .key = .title, .value = value } }}); + self.titled = true; + } + /// v1's rule: a generated title never replaces one the user chose, only /// none or the one derived from the first message. pub fn installGeneratedTitle(self: *Session, history: []const types.HistoryTurn, title: []const u8) !bool { @@ -653,8 +706,10 @@ pub const Session = struct { .history = &.{}, .language = language, .created_at_ms = std.math.cast(i64, state.created_ms) orelse 0, + .updated_at_ms = std.math.cast(i64, state.updated_ms) orelse 0, }; errdefer restored.deinit(alloc); + if (state.title) |raw| restored.title = try alloc.dupe(u8, try std.json.parseFromSliceLeaky([]const u8, sa, raw, .{})); if (state.prefs) |raw| restored.preferences = try decodePreferences(alloc, raw); if (state.permissions) |raw| restored.permission_state = try session_codec.decodePermissionState(alloc, raw); if (state.usage) |raw| { @@ -703,6 +758,45 @@ pub const Session = struct { return restored; } + /// The session as v1's `DurableSessionState`, for hosts that restore + /// through it, and its stored title. Caller owns both. + pub fn durableState(self: *Session, alloc: Allocator, workspace: []const u8) !Resumed { + var restored = try self.restore(alloc); + defer restored.deinit(alloc); + const id_copy = try alloc.dupe(u8, self.id()); + errdefer alloc.free(id_copy); + const origin = try alloc.dupe(u8, workspace); + errdefer alloc.free(origin); + const workspace_copy = try alloc.dupe(u8, workspace); + errdefer alloc.free(workspace_copy); + // Every session starts with its preferences (`create`). + const preferences = restored.preferences orelse return error.InvalidSessionFormat; + restored.preferences = null; + const resumed: Resumed = .{ + .state = .{ + .id = id_copy, + .origin_workspace_root = origin, + .workspace_root = workspace_copy, + .created_at_ms = restored.created_at_ms, + .updated_at_ms = restored.updated_at_ms, + .conversation_language = restored.language, + .preferences = preferences, + .history = restored.history, + // v1 reloads its totals as zero as well. + .total_input_tokens = 0, + .total_output_tokens = 0, + .permission_state = restored.permission_state orelse .{}, + .usage = restored.usage, + }, + .title = restored.title, + }; + restored.history = &.{}; + restored.permission_state = null; + restored.usage = null; + restored.title = null; + return resumed; + } + /// The cursor of `turn`'s `turn_started`, reading back from `before`. fn findTurnStart(self: *Session, sa: Allocator, before: sm.Cursor, turn: u64) !sm.Cursor { var from: sm.From = .{ .at = before }; @@ -846,6 +940,68 @@ fn withoutCreatedAt(a: Allocator, bytes: []const u8) ![]u8 { return jsonValue(a, value); } +// --------------------------------------------------------------------------- +// Listing + +const list_page_size: usize = 256; + +/// Every saved root session, newest first, as v1's picker summaries, leaving +/// out `active_id`. A saved session always has a turn (D2), so each one can +/// be resumed. Stops with `error.Cancelled` once `cancel` is set. Caller +/// owns the list and every summary; safe from any thread. +pub fn listSummaries( + store: *Store, + alloc: Allocator, + active_id: ?[]const u8, + cancel: *const std.atomic.Value(bool), +) !std.ArrayList(session_store.SessionSummary) { + var list: std.ArrayList(session_store.SessionSummary) = .empty; + errdefer { + for (list.items) |*summary| summary.deinit(alloc); + list.deinit(alloc); + } + var scratch = std.heap.ArenaAllocator.init(alloc); + defer scratch.deinit(); + var cursor: ?sm.ListCursor = null; + while (true) { + if (cancel.load(.acquire)) return error.Cancelled; + var page = try store.manager.list(alloc, .all, cursor, list_page_size); + defer page.deinit(); + for (page.items) |item| { + if (item.role != .root) continue; + if (active_id) |active| if (std.mem.eql(u8, active, item.id)) continue; + try list.append(alloc, try summaryOf(alloc, scratch.allocator(), item)); + } + cursor = page.next orelse break; + } + return list; +} + +fn summaryOf(alloc: Allocator, scratch: Allocator, item: sm.Summary) !session_store.SessionSummary { + const id = try alloc.dupe(u8, item.id); + errdefer alloc.free(id); + const workspace = try alloc.dupe(u8, item.workspace); + errdefer alloc.free(workspace); + const origin = try alloc.dupe(u8, item.workspace); + errdefer alloc.free(origin); + const title: ?[]u8 = if (item.title) |raw| + try alloc.dupe(u8, try std.json.parseFromSliceLeaky([]const u8, scratch, raw, .{})) + else + null; + errdefer if (title) |value| alloc.free(value); + return .{ + .id = id, + .workspace_root = workspace, + .origin_workspace_root = origin, + .title = title, + .display_metadata_present = title != null, + .created_at_ms = std.math.cast(i64, item.created_ms) orelse 0, + .updated_at_ms = std.math.cast(i64, item.updated_ms) orelse 0, + .conversation_language = if (item.language) |raw| try decodeLanguage(scratch, raw) else types.ConversationLanguage.default(), + .history_len = @max(item.turns, 1), + }; +} + // --------------------------------------------------------------------------- // Usage recovery: what the profile's readers need from v2 sessions @@ -937,13 +1093,14 @@ fn newestUsage(store: *Store, alloc: Allocator, id: []const u8) !?UsageCheckpoin /// v1's error names for a failed resume, so every host reports the same /// error whichever backend is on. -const ResumeError = error{ SessionNotFound, NoSavedSessions, SessionBusy, InvalidSessionFormat, UnsupportedSessionFormat } || sm.OpenError; +const ResumeError = error{ SessionNotFound, NoSavedSessions, NoRememberedSession, SessionBusy, InvalidSessionFormat, UnsupportedSessionFormat } || sm.OpenError; fn resumeError(err: sm.OpenError, target: Target) ResumeError { return switch (err) { error.NotFound => switch (target) { .id => error.SessionNotFound, .last => error.NoSavedSessions, + .last_opened => error.NoRememberedSession, }, // A child is resumed only through its parent, as in v1. error.ChildSession => error.SessionNotFound, @@ -1481,6 +1638,91 @@ test "a failed resume reports v1's error names" { try testing.expectEqual(error.Io, resumeError(error.Io, .last)); } +test "a resumed session gives v1's state and the title the user chose" { + var t: TestHome = undefined; + try t.init(); + defer t.deinit(); + var model = "m".*; + const s = try Session.create(testing.allocator, &t.store, "/w", .app, testSeed(&model)); + try s.commitTurn(assistantTurn("q", "a"), types.ConversationLanguage.default()); + try s.rename("Chosen title"); + const id = try testing.allocator.dupe(u8, s.id()); + defer testing.allocator.free(id); + s.close(); + + const r = try Session.resumeSession(testing.allocator, &t.store, .{ .id = id }, "/w", .app); + defer r.close(); + var resumed = try r.durableState(testing.allocator, "/w"); + defer resumed.deinit(testing.allocator); + try testing.expectEqualStrings(id, resumed.state.id); + try testing.expectEqualStrings("/w", resumed.state.workspace_root); + try testing.expectEqual(@as(usize, 1), resumed.state.history.len); + try testing.expectEqualStrings("m", resumed.state.preferences.model); + try testing.expectEqualStrings("Chosen title", resumed.title.?); + try testing.expect(resumed.state.created_at_ms > 0); + try testing.expect(resumed.state.updated_at_ms >= resumed.state.created_at_ms); + // A generated title does not replace the one the user chose. + try testing.expect(!try r.installGeneratedTitle(resumed.state.history, "Generated")); +} + +test "the picker lists saved root sessions newest first, without the open one" { + var t: TestHome = undefined; + try t.init(); + defer t.deinit(); + var model = "m".*; + const first = try Session.create(testing.allocator, &t.store, "/w", .app, testSeed(&model)); + try first.commitTurn(assistantTurn("first", "a"), types.ConversationLanguage.default()); + const first_id = try testing.allocator.dupe(u8, first.id()); + defer testing.allocator.free(first_id); + first.close(); + // A session with no turn is not saved, so it is never listed (D2). + const empty = try Session.create(testing.allocator, &t.store, "/w", .app, testSeed(&model)); + empty.close(); + io_mod.sleep(2 * std.time.ns_per_ms); + const second = try Session.create(testing.allocator, &t.store, "/w", .app, testSeed(&model)); + defer second.close(); + try second.commitTurn(assistantTurn("second", "b"), types.ConversationLanguage.default()); + + var cancel = std.atomic.Value(bool).init(false); + var all = try listSummaries(&t.store, testing.allocator, null, &cancel); + defer { + for (all.items) |*summary| summary.deinit(testing.allocator); + all.deinit(testing.allocator); + } + try testing.expectEqual(@as(usize, 2), all.items.len); + try testing.expectEqualStrings(second.id(), all.items[0].id); + try testing.expectEqualStrings("/w", all.items[0].workspace_root.?); + try testing.expect(all.items[0].hasResumableContent()); + + var others = try listSummaries(&t.store, testing.allocator, second.id(), &cancel); + defer { + for (others.items) |*summary| summary.deinit(testing.allocator); + others.deinit(testing.allocator); + } + try testing.expectEqual(@as(usize, 1), others.items.len); + try testing.expectEqualStrings(first_id, others.items[0].id); + + cancel.store(true, .release); + try testing.expectError(error.Cancelled, listSummaries(&t.store, testing.allocator, null, &cancel)); +} + +test "-c resumes the session this host last opened" { + var t: TestHome = undefined; + try t.init(); + defer t.deinit(); + var model = "m".*; + const s = try Session.create(testing.allocator, &t.store, "/w", .app, testSeed(&model)); + try s.commitTurn(assistantTurn("q", "a"), types.ConversationLanguage.default()); + const id = try testing.allocator.dupe(u8, s.id()); + defer testing.allocator.free(id); + s.close(); + + const again = try Session.resumeSession(testing.allocator, &t.store, .last_opened, "/w", .app); + defer again.close(); + try testing.expectEqualStrings(id, again.id()); + try testing.expectError(error.NoRememberedSession, Session.resumeSession(testing.allocator, &t.store, .last_opened, "/w", .acp)); +} + test "side files live in a private session-files folder" { var t: TestHome = undefined; try t.init(); diff --git a/src/core/session_manager/api.zig b/src/core/session_manager/api.zig index 6b9bda7ddf..4413618006 100644 --- a/src/core/session_manager/api.zig +++ b/src/core/session_manager/api.zig @@ -108,6 +108,10 @@ pub const ResumeOptions = struct { workspace: []const u8, host: Host, parent: ?[]const u8 = null, + /// How long this open waits for the session's writer lock before Busy; + /// null uses the manager's `lock_wait_ms`. A picker passes 0 so a + /// session open elsewhere shows as busy at once (D38). + lock_wait_ms: ?u64 = null, }; pub const ForkOptions = struct { @@ -252,6 +256,7 @@ pub const Manager = struct { .workspace = options.workspace, .host = options.host, .parent = options.parent, + .lock_wait_ms = options.lock_wait_ms, }) catch |err| return switch (err) { error.NotFound, error.Busy, error.ChildSession, error.Corrupt, error.UnsupportedVersion, error.Io, error.OutOfMemory => |e| e, }; @@ -1342,6 +1347,36 @@ const api_tests = struct { try testing.expectError(error.SessionClosed, s.read(gpa, .end, .backward, 10)); s.release(); } + + test "a resume sets its own wait for a held session, and none is Busy at once (D38)" { + var f: Fixture = undefined; + try f.init(); + defer f.deinit(); + const m = f.manager; + const held = try m.openNew(.{ .workspace = "/w", .host = .app }); + _ = try held.append(&.{ .turn_started, piece, .turn_committed }); + const id = try gpa.dupe(u8, held.id()); + defer gpa.free(id); + var released = false; + defer if (!released) held.release(); + + const Timed = struct { + fn busyAfterMs(manager: *api.Manager, session_id: []const u8, wait: ?u64) !i64 { + const started = std.Io.Timestamp.now(io, .awake).toMilliseconds(); + try testing.expectError(error.Busy, manager.openResume(.{ .target = .{ .id = session_id }, .workspace = "/w", .host = .app, .lock_wait_ms = wait })); + return std.Io.Timestamp.now(io, .awake).toMilliseconds() - started; + } + }; + // The fixture's manager waits 50 ms; each open's own wait wins. + try testing.expect(try Timed.busyAfterMs(m, id, 400) >= 400); + try testing.expect(try Timed.busyAfterMs(m, id, 0) < 400); + try testing.expect(try Timed.busyAfterMs(m, id, null) < 400); + + held.release(); + released = true; + const reopened = try m.openResume(.{ .target = .{ .id = id }, .workspace = "/w", .host = .app, .lock_wait_ms = 0 }); + reopened.release(); + } }; test { diff --git a/src/core/session_manager/session.zig b/src/core/session_manager/session.zig index 72b702068f..794d7c3750 100644 --- a/src/core/session_manager/session.zig +++ b/src/core/session_manager/session.zig @@ -923,6 +923,9 @@ pub const ResumeOptions = struct { host: schema.Host, /// Required to resume a child session. parent: ?[]const u8 = null, + /// How long to wait for the flock before Busy; null uses the + /// environment's `lock_wait_ms` (D38). + lock_wait_ms: ?u64 = null, }; /// Opens a published session for writing: the flock (else Busy after the @@ -987,7 +990,7 @@ fn openParts(env: *const Env, options: ResumeOptions) OpenError!Parts { else => error.Io, }; errdefer s.closeFile(lock); - if (!env.isPlanted(.skip_flock)) try acquireLock(env, lock); + if (!env.isPlanted(.skip_flock)) try acquireLock(env, lock, options.lock_wait_ms orelse env.options.lock_wait_ms); var opened = Log.open(gpa, s, dir, "log.jsonl", .read_write, .{}) catch |err| return switch (err) { error.NotFound => error.Corrupt, @@ -1036,13 +1039,13 @@ fn openParts(env: *const Env, options: ResumeOptions) OpenError!Parts { return .{ .dir = dir, .lock = lock, .log = opened.log, .loaded = loaded }; } -fn acquireLock(env: *const Env, lock: storage.File) OpenError!void { +fn acquireLock(env: *const Env, lock: storage.File, wait_ms: u64) OpenError!void { const io = env.s.io; const step_ms: u64 = 10; var waited: u64 = 0; while (true) { if (env.s.tryLock(lock) catch return error.Io) return; - if (waited >= env.options.lock_wait_ms) return error.Busy; + if (waited >= wait_ms) return error.Busy; io.sleep(.fromMilliseconds(@intCast(step_ms)), .awake) catch return error.Busy; waited += step_ms; } diff --git a/src/main.zig b/src/main.zig index bffd0fd689..d35c7f2731 100644 --- a/src/main.zig +++ b/src/main.zig @@ -669,6 +669,7 @@ const App = struct { launch.requested_resume = null; } errdefer if (app.requested_resume) |*target| target.deinit(alloc); + app.session_persistence.sessions_v2 = launch.modifiers.sessions_v2; try BootstrapAppRuntime.bootstrap( &app, footer_rows, @@ -914,8 +915,9 @@ const App = struct { buffer: []u8, session_id: []const u8, terminal_cols: u16, + sessions_v2: bool, ) ![]const u8 { - return ui_render.formatResumeHandoff(buffer, session_id, terminal_cols); + return ui_render.formatResumeHandoff(buffer, session_id, terminal_cols, sessions_v2); } /// Full teardown for hosts that keep running after the shell ends, such as diff --git a/src/ui/render.zig b/src/ui/render.zig index b2267f4582..4090fb15ae 100644 --- a/src/ui/render.zig +++ b/src/ui/render.zig @@ -790,9 +790,11 @@ pub fn formatResumeHandoff( buffer: []u8, session_id: []const u8, terminal_cols: u16, + sessions_v2: bool, ) ![]const u8 { const label = "Continue session with:"; - const command = "fx --resume "; + // A v2 session resumes only with the flag that saved it. + const command = if (sessions_v2) "fx --sessions-v2 --resume " else "fx --resume "; const single_row_width = label.len + 1 + command.len + session_id.len; const separator = if (single_row_width <= terminal_cols) " " else "\n "; return std.fmt.bufPrint( @@ -861,18 +863,25 @@ test "resume handoff uses one row only when the full instruction fits" { const single_row = "Continue session with: fx --resume session-123"; var exact_buffer: [128]u8 = undefined; - const exact = try formatResumeHandoff(&exact_buffer, "session-123", single_row.len); + const exact = try formatResumeHandoff(&exact_buffer, "session-123", single_row.len, false); try std.testing.expectEqualStrings( "\x1b[38;5;245mContinue session with: fx --resume session-123\x1b[0m\n", exact, ); var narrow_buffer: [128]u8 = undefined; - const narrow = try formatResumeHandoff(&narrow_buffer, "session-123", single_row.len - 1); + const narrow = try formatResumeHandoff(&narrow_buffer, "session-123", single_row.len - 1, false); try std.testing.expectEqualStrings( "\x1b[38;5;245mContinue session with:\n fx --resume session-123\x1b[0m\n", narrow, ); + + var v2_buffer: [128]u8 = undefined; + const v2 = try formatResumeHandoff(&v2_buffer, "session-123", 80, true); + try std.testing.expectEqualStrings( + "\x1b[38;5;245mContinue session with: fx --sessions-v2 --resume session-123\x1b[0m\n", + v2, + ); } test "resume handoff follows the active muted theme shade" { @@ -880,7 +889,7 @@ test "resume handoff follows the active muted theme shade" { defer initTheme(false, null); var buffer: [128]u8 = undefined; - const message = try formatResumeHandoff(&buffer, "session-123", 80); + const message = try formatResumeHandoff(&buffer, "session-123", 80, false); try std.testing.expectEqualStrings( "\x1b[38;5;247mContinue session with: fx --resume session-123\x1b[0m\n", message, diff --git a/tests/e2e/sessions-v2.test.ts b/tests/e2e/sessions-v2.test.ts index b1244d959c..25e688c74f 100644 --- a/tests/e2e/sessions-v2.test.ts +++ b/tests/e2e/sessions-v2.test.ts @@ -22,6 +22,8 @@ import { fakeShellRun, startDynamicFakeGateway, startFakeGateway, + TmuxSession, + tmuxAvailable, } from "./tmux-helpers"; // Sessions v2 behind FX_SESSIONS_V2 and --sessions-v2: every `fx ask` entry @@ -595,3 +597,233 @@ test("a full disk fails the turn cleanly and the session resumes after", async ( rmSync(fixture.root, { recursive: true, force: true }); } }, TIMEOUT * 3); + +// --------------------------------------------------------------------------- +// The interactive app + +/// Answers the newest prompt in the request, in the order given: every +/// request carries the earlier prompts, and a fresh session also asks for a +/// title. A null answer holds that request open and calls `held`. +function replyToLatest(pairs: [string, string | null][], held?: () => void) { + return startDynamicFakeGateway(async (body) => { + for (let index = pairs.length - 1; index >= 0; index -= 1) { + const [prompt, answer] = pairs[index]!; + if (!body.includes(prompt)) continue; + if (answer !== null) return fakeGatewayFinalText(answer); + held?.(); + return new Promise(() => {}); + } + return fakeGatewayFinalText("UNEXPECTED_REQUEST"); + }); +} + +async function startApp(fixture: Fixture, gateway: any, args: string[], waitForComposer = true) { + const stderrPath = join(fixture.root, "stderr.log"); + writeFileSync(stderrPath, ""); + const session = await TmuxSession.create({ + cmd: `${FX_BIN} --sessions-v2 ${args.join(" ")}`.trim(), + cwd: fixture.workspace, + env: { ...env(fixture, gateway, false), NO_COLOR: "1" }, + stderrPath, + }); + if (waitForComposer) await session.waitForComposer(TIMEOUT); + return { session, stderrPath }; +} + +async function quitApp(app: { session: TmuxSession; stderrPath: string }) { + await app.session.sendText("/quit"); + expect(await app.session.waitForSessionEnd()).toBe(true); + await app.session.kill(); + expect(readFileSync(app.stderrPath, "utf8")).toBe(""); +} + +function savedSessions(fixture: Fixture): string[] { + return readdirSync(v2Root(fixture)).filter((name) => !name.startsWith(".") && !name.startsWith("index")); +} + +function onlySession(fixture: Fixture): string { + const ids = savedSessions(fixture); + expect(ids.length).toBe(1); + return ids[0]!; +} + +async function scrollbackContains(session: TmuxSession, marker: string) { + const deadline = Date.now() + TIMEOUT; + let latest = ""; + while (Date.now() < deadline) { + latest = await session.captureFullScrollback(); + if (latest.includes(marker)) return latest; + await Bun.sleep(100); + } + throw new Error(`scrollback never showed ${marker}`); +} + +test.skipIf(!tmuxAvailable())("the interactive app saves to v2 and resumes with -c, --resume last and the id", async () => { + const fixture = createFixture("fx-v2-app-"); + const gateway = replyToLatest([ + ["First interactive question.", "APP_V2_FIRST"], + ["Continue with -c.", "APP_V2_CONTINUE"], + ["Continue with --resume last.", "APP_V2_LAST"], + ["Continue with the id.", "APP_V2_BY_ID"], + ]); + try { + const first = await startApp(fixture, gateway, []); + await first.session.sendText("First interactive question."); + await first.session.waitForText("APP_V2_FIRST", TIMEOUT); + await quitApp(first); + const id = onlySession(fixture); + expectNoV1Sessions(fixture); + + const runs: [string[], string, string, string][] = [ + [["-c"], "APP_V2_FIRST", "Continue with -c.", "APP_V2_CONTINUE"], + [["--resume", "last"], "APP_V2_CONTINUE", "Continue with --resume last.", "APP_V2_LAST"], + [["--resume", id], "APP_V2_LAST", "Continue with the id.", "APP_V2_BY_ID"], + ]; + for (const [args, restored, prompt, answer] of runs) { + const app = await startApp(fixture, gateway, args); + expect(await scrollbackContains(app.session, restored)).toContain(restored); + await app.session.sendText(prompt); + await app.session.waitForText(answer, TIMEOUT); + await quitApp(app); + expect(onlySession(fixture)).toBe(id); + } + expect(gateway.requests.at(-1)!.body).toContain("First interactive question."); + const kinds = logLines(fixture, id).map((line) => line.kind); + expect(kinds.filter((kind) => kind === "turn_committed").length).toBe(4); + expectWholeLog(fixture, id); + expectNoV1Sessions(fixture); + } finally { + gateway.stop(); + rmSync(fixture.root, { recursive: true, force: true }); + } +}, TIMEOUT * 6); + +test.skipIf(!tmuxAvailable())("a killed interactive turn comes back interrupted and -c continues it", async () => { + const fixture = createFixture("fx-v2-app-kill-"); + let held: () => void = () => {}; + const reachedHold = new Promise((resolve) => (held = resolve)); + const gateway = replyToLatest( + [ + ["Before the interactive kill.", "APP_BEFORE_KILL"], + ["This interactive turn is killed.", null], + ["After the interactive kill.", "APP_AFTER_KILL"], + ], + () => held(), + ); + try { + const app = await startApp(fixture, gateway, []); + await app.session.sendText("Before the interactive kill."); + await app.session.waitForText("APP_BEFORE_KILL", TIMEOUT); + await app.session.sendText("This interactive turn is killed."); + await reachedHold; + Bun.spawnSync(["kill", "-9", String(app.session.processPid())]); + await app.session.kill(); + const id = onlySession(fixture); + + const resumed = await startApp(fixture, gateway, ["-c"]); + expect(await scrollbackContains(resumed.session, "APP_BEFORE_KILL")).toContain("APP_BEFORE_KILL"); + await resumed.session.sendText("After the interactive kill."); + await resumed.session.waitForText("APP_AFTER_KILL", TIMEOUT); + await quitApp(resumed); + expect(onlySession(fixture)).toBe(id); + // The killed prompt was saved before its request, and the model sees it. + expect(gateway.requests.at(-1)!.body).toContain("This interactive turn is killed."); + const lines = logLines(fixture, id); + expect(lines.filter((line) => line.kind === "turn_interrupted").map((line) => line.reason)).toEqual(["crash"]); + expectWholeLog(fixture, id); + expectNoV1Sessions(fixture); + } finally { + gateway.stop(); + rmSync(fixture.root, { recursive: true, force: true }); + } +}, TIMEOUT * 4); + +test.skipIf(!tmuxAvailable())("the picker lists v2 sessions, /rename sticks, and /new starts another", async () => { + const fixture = createFixture("fx-v2-app-picker-"); + const gateway = replyToLatest([ + ["Picker session one.", "PICK_ONE"], + ["Picker session two.", "PICK_TWO"], + ["Back in session one.", "PICK_BACK"], + ]); + try { + const app = await startApp(fixture, gateway, []); + await app.session.sendText("Picker session one."); + await app.session.waitForText("PICK_ONE", TIMEOUT); + await app.session.sendText("/rename Renamed picker session"); + await app.session.waitForText('renamed to "Renamed picker session"', TIMEOUT); + await app.session.waitForStableComposer(); + await app.session.sendText("/new"); + await app.session.waitForStableComposer(); + await app.session.sendText("Picker session two."); + await app.session.waitForText("PICK_TWO", TIMEOUT); + await quitApp(app); + const ids = savedSessions(fixture); + expect(ids.length).toBe(2); + + const picker = await startApp(fixture, gateway, ["-r"], false); + await picker.session.waitForPane((pane) => pane.includes("Renamed picker session") && pane.includes("enter resume"), TIMEOUT); + await picker.session.sendLiteralText("Renamed"); + await picker.session.waitForPane((pane) => pane.includes("Renamed picker session"), TIMEOUT); + await picker.session.sendKeys("Enter"); + await picker.session.waitForComposer(TIMEOUT); + expect(await scrollbackContains(picker.session, "PICK_ONE")).toContain("PICK_ONE"); + await picker.session.sendText("Back in session one."); + await picker.session.waitForText("PICK_BACK", TIMEOUT); + await quitApp(picker); + + // The resumed session is the renamed one: its log holds both of its turns. + const renamed = ids.find((id) => readFileSync(join(v2Root(fixture), id, "log.jsonl"), "utf8").includes("Renamed picker session"))!; + const log = readFileSync(join(v2Root(fixture), renamed, "log.jsonl"), "utf8"); + expect(log).toContain("Back in session one."); + expect(log).not.toContain("Picker session two."); + for (const id of ids) expectWholeLog(fixture, id); + expectNoV1Sessions(fixture); + } finally { + gateway.stop(); + rmSync(fixture.root, { recursive: true, force: true }); + } +}, TIMEOUT * 6); + +test.skipIf(!tmuxAvailable())("the picker shows a session open in another fx as busy at once, and opens it once the owner quits", async () => { + const fixture = createFixture("fx-v2-app-picker-busy-"); + const gateway = replyToLatest([["Hold this session open.", "PICKER_BUSY_SAVED"]]); + const busy = "This session is open in another fx. Close it there, then press enter to retry."; + let contender: TmuxSession | null = null; + try { + const owner = await startApp(fixture, gateway, []); + await owner.session.sendText("Hold this session open."); + await owner.session.waitForText("PICKER_BUSY_SAVED", TIMEOUT); + const id = onlySession(fixture); + + const contenderStderr = join(fixture.root, "contender-stderr.log"); + writeFileSync(contenderStderr, ""); + contender = await TmuxSession.create({ + cmd: `${FX_BIN} --sessions-v2 -r`, + cwd: fixture.workspace, + env: { ...env(fixture, gateway, false), NO_COLOR: "1" }, + stderrPath: contenderStderr, + }); + await contender.waitForPane((pane) => pane.includes("Hold this session open.") && pane.includes("enter resume"), TIMEOUT); + // v1's picker takes no lock wait, and neither does v2's (D38). + const pressed = Date.now(); + await contender.sendKeys("Enter"); + await contender.waitForPane((pane) => pane.includes(busy), 1_000); + expect(Date.now() - pressed).toBeLessThan(1_000); + expect(owner.session.isPaneAlive()).toBe(true); + + await quitApp(owner); + await contender.sendKeys("Enter"); + await contender.waitForComposer(TIMEOUT); + expect(await scrollbackContains(contender, "PICKER_BUSY_SAVED")).toContain("PICKER_BUSY_SAVED"); + await contender.sendText("/quit"); + expect(await contender.waitForSessionEnd()).toBe(true); + expect(readFileSync(contenderStderr, "utf8")).toBe(""); + expect(onlySession(fixture)).toBe(id); + expectWholeLog(fixture, id); + expectNoV1Sessions(fixture); + } finally { + if (contender) await contender.kill(); + gateway.stop(); + rmSync(fixture.root, { recursive: true, force: true }); + } +}, TIMEOUT * 4);