diff --git a/apps/daemon/drizzle/0015_fast_martin_li.sql b/apps/daemon/drizzle/0015_fast_martin_li.sql new file mode 100644 index 000000000..8775fc152 --- /dev/null +++ b/apps/daemon/drizzle/0015_fast_martin_li.sql @@ -0,0 +1,14 @@ +CREATE TABLE `worktree_sessions` ( + `worktree_path` text NOT NULL, + `session_id` text NOT NULL, + `created_at` integer NOT NULL, + PRIMARY KEY(`worktree_path`, `session_id`), + FOREIGN KEY (`worktree_path`) REFERENCES `worktrees`(`worktree_path`) ON UPDATE no action ON DELETE cascade +); +--> statement-breakpoint +CREATE UNIQUE INDEX `worktree_sessions_session_unique` ON `worktree_sessions` (`session_id`);--> statement-breakpoint +INSERT INTO `worktree_sessions` (`worktree_path`, `session_id`, `created_at`) +SELECT `worktree_path`, `session_id`, `created_at` FROM `worktrees` +WHERE `session_id` NOT LIKE 'orphan-worktree-%';--> statement-breakpoint +DROP INDEX `worktrees_session_id_unique`;--> statement-breakpoint +ALTER TABLE `worktrees` DROP COLUMN `session_id`; \ No newline at end of file diff --git a/apps/daemon/drizzle/meta/0015_snapshot.json b/apps/daemon/drizzle/meta/0015_snapshot.json new file mode 100644 index 000000000..1a5b10cab --- /dev/null +++ b/apps/daemon/drizzle/meta/0015_snapshot.json @@ -0,0 +1,1611 @@ +{ + "version": "6", + "dialect": "sqlite", + "id": "8bffb46d-1e09-4f6a-8ff3-fb733b195d41", + "prevId": "a9ae9955-d7e4-45a9-a33c-59f8b342efd7", + "tables": { + "attachment_blobs": { + "name": "attachment_blobs", + "columns": { + "attachment_id": { + "name": "attachment_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "variant": { + "name": "variant", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "blob_id": { + "name": "blob_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "attachment_blobs_blob_idx": { + "name": "attachment_blobs_blob_idx", + "columns": ["blob_id"], + "isUnique": false + } + }, + "foreignKeys": { + "attachment_blobs_attachment_id_attachments_attachment_id_fk": { + "name": "attachment_blobs_attachment_id_attachments_attachment_id_fk", + "tableFrom": "attachment_blobs", + "tableTo": "attachments", + "columnsFrom": ["attachment_id"], + "columnsTo": ["attachment_id"], + "onDelete": "cascade", + "onUpdate": "no action" + }, + "attachment_blobs_blob_id_blobs_blob_id_fk": { + "name": "attachment_blobs_blob_id_blobs_blob_id_fk", + "tableFrom": "attachment_blobs", + "tableTo": "blobs", + "columnsFrom": ["blob_id"], + "columnsTo": ["blob_id"], + "onDelete": "no action", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": { + "attachment_blobs_attachment_id_variant_pk": { + "columns": ["attachment_id", "variant"], + "name": "attachment_blobs_attachment_id_variant_pk" + } + }, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "attachments": { + "name": "attachments", + "columns": { + "attachment_id": { + "name": "attachment_id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "kind": { + "name": "kind", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "name": { + "name": "name", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "mime_type": { + "name": "mime_type", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "size_bytes": { + "name": "size_bytes", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "metadata_json": { + "name": "metadata_json", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": {}, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "blobs": { + "name": "blobs", + "columns": { + "blob_id": { + "name": "blob_id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "size_bytes": { + "name": "size_bytes", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": {}, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "conversation_operations": { + "name": "conversation_operations", + "columns": { + "operation_id": { + "name": "operation_id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "session_id": { + "name": "session_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "kind": { + "name": "kind", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "state": { + "name": "state", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "turn_id": { + "name": "turn_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "error_code": { + "name": "error_code", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "error_message": { + "name": "error_message", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "resolved_at": { + "name": "resolved_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + } + }, + "indexes": { + "conversation_operations_session_idx": { + "name": "conversation_operations_session_idx", + "columns": ["session_id"], + "isUnique": false + }, + "conversation_operations_open_session_unique": { + "name": "conversation_operations_open_session_unique", + "columns": ["session_id"], + "isUnique": true, + "where": "state = 'open'" + } + }, + "foreignKeys": { + "conversation_operations_session_id_sessions_session_id_fk": { + "name": "conversation_operations_session_id_sessions_session_id_fk", + "tableFrom": "conversation_operations", + "tableTo": "sessions", + "columnsFrom": ["session_id"], + "columnsTo": ["session_id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "conversation_turns": { + "name": "conversation_turns", + "columns": { + "turn_id": { + "name": "turn_id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "session_id": { + "name": "session_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "parent_turn_id": { + "name": "parent_turn_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "sibling_ordinal": { + "name": "sibling_ordinal", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "input_type": { + "name": "input_type", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "prompt_id": { + "name": "prompt_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "command_name": { + "name": "command_name", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "command_arguments": { + "name": "command_arguments", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "shell_command": { + "name": "shell_command", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "run_id": { + "name": "run_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "state": { + "name": "state", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "conversation_turns_session_idx": { + "name": "conversation_turns_session_idx", + "columns": ["session_id"], + "isUnique": false + }, + "conversation_turns_sibling_unique": { + "name": "conversation_turns_sibling_unique", + "columns": ["session_id", "parent_turn_id", "sibling_ordinal"], + "isUnique": true, + "where": "parent_turn_id IS NOT NULL" + }, + "conversation_turns_root_sibling_unique": { + "name": "conversation_turns_root_sibling_unique", + "columns": ["session_id", "sibling_ordinal"], + "isUnique": true, + "where": "parent_turn_id IS NULL" + } + }, + "foreignKeys": { + "conversation_turns_session_id_sessions_session_id_fk": { + "name": "conversation_turns_session_id_sessions_session_id_fk", + "tableFrom": "conversation_turns", + "tableTo": "sessions", + "columnsFrom": ["session_id"], + "columnsTo": ["session_id"], + "onDelete": "cascade", + "onUpdate": "no action" + }, + "conversation_turns_parent_turn_id_conversation_turns_turn_id_fk": { + "name": "conversation_turns_parent_turn_id_conversation_turns_turn_id_fk", + "tableFrom": "conversation_turns", + "tableTo": "conversation_turns", + "columnsFrom": ["parent_turn_id"], + "columnsTo": ["turn_id"], + "onDelete": "cascade", + "onUpdate": "no action" + }, + "conversation_turns_prompt_id_prompts_prompt_id_fk": { + "name": "conversation_turns_prompt_id_prompts_prompt_id_fk", + "tableFrom": "conversation_turns", + "tableTo": "prompts", + "columnsFrom": ["prompt_id"], + "columnsTo": ["prompt_id"], + "onDelete": "no action", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "loop_iterations": { + "name": "loop_iterations", + "columns": { + "loop_id": { + "name": "loop_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "index": { + "name": "index", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "status": { + "name": "status", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "worker_session_id": { + "name": "worker_session_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "verifier_session_id": { + "name": "verifier_session_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "checks_json": { + "name": "checks_json", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "verdict_json": { + "name": "verdict_json", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "error": { + "name": "error", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "started_at": { + "name": "started_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "ended_at": { + "name": "ended_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + } + }, + "indexes": {}, + "foreignKeys": { + "loop_iterations_loop_id_loops_loop_id_fk": { + "name": "loop_iterations_loop_id_loops_loop_id_fk", + "tableFrom": "loop_iterations", + "tableTo": "loops", + "columnsFrom": ["loop_id"], + "columnsTo": ["loop_id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": { + "loop_iterations_loop_id_index_pk": { + "columns": ["loop_id", "index"], + "name": "loop_iterations_loop_id_index_pk" + } + }, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "loops": { + "name": "loops", + "columns": { + "loop_id": { + "name": "loop_id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "spec_json": { + "name": "spec_json", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "status": { + "name": "status", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "iteration_count": { + "name": "iteration_count", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": 0 + }, + "error": { + "name": "error", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "summary": { + "name": "summary", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "started_at": { + "name": "started_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "updated_at": { + "name": "updated_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "ended_at": { + "name": "ended_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + } + }, + "indexes": {}, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "prompt_attachment_refs": { + "name": "prompt_attachment_refs", + "columns": { + "prompt_id": { + "name": "prompt_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "attachment_id": { + "name": "attachment_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": {}, + "foreignKeys": { + "prompt_attachment_refs_prompt_id_prompts_prompt_id_fk": { + "name": "prompt_attachment_refs_prompt_id_prompts_prompt_id_fk", + "tableFrom": "prompt_attachment_refs", + "tableTo": "prompts", + "columnsFrom": ["prompt_id"], + "columnsTo": ["prompt_id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": { + "prompt_attachment_refs_prompt_id_attachment_id_pk": { + "columns": ["prompt_id", "attachment_id"], + "name": "prompt_attachment_refs_prompt_id_attachment_id_pk" + } + }, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "prompts": { + "name": "prompts", + "columns": { + "prompt_id": { + "name": "prompt_id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "blocks_json": { + "name": "blocks_json", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "context_attachment_ids_json": { + "name": "context_attachment_ids_json", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": {}, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "provider_turn_bindings": { + "name": "provider_turn_bindings", + "columns": { + "turn_id": { + "name": "turn_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "run_id": { + "name": "run_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "history_id": { + "name": "history_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "checkpoint": { + "name": "checkpoint", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "captured_from": { + "name": "captured_from", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": {}, + "foreignKeys": { + "provider_turn_bindings_turn_id_conversation_turns_turn_id_fk": { + "name": "provider_turn_bindings_turn_id_conversation_turns_turn_id_fk", + "tableFrom": "provider_turn_bindings", + "tableTo": "conversation_turns", + "columnsFrom": ["turn_id"], + "columnsTo": ["turn_id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": { + "provider_turn_bindings_turn_id_history_id_pk": { + "columns": ["turn_id", "history_id"], + "name": "provider_turn_bindings_turn_id_history_id_pk" + } + }, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "schedule_runs": { + "name": "schedule_runs", + "columns": { + "run_id": { + "name": "run_id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "schedule_id": { + "name": "schedule_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "status": { + "name": "status", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "trigger": { + "name": "trigger", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "session_id": { + "name": "session_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "error": { + "name": "error", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "summary": { + "name": "summary", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "started_at": { + "name": "started_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "ended_at": { + "name": "ended_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + } + }, + "indexes": { + "schedule_runs_schedule_started_idx": { + "name": "schedule_runs_schedule_started_idx", + "columns": ["schedule_id", "started_at"], + "isUnique": false + } + }, + "foreignKeys": { + "schedule_runs_schedule_id_schedules_schedule_id_fk": { + "name": "schedule_runs_schedule_id_schedules_schedule_id_fk", + "tableFrom": "schedule_runs", + "tableTo": "schedules", + "columnsFrom": ["schedule_id"], + "columnsTo": ["schedule_id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "schedules": { + "name": "schedules", + "columns": { + "schedule_id": { + "name": "schedule_id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "name": { + "name": "name", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "prompt": { + "name": "prompt", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "cadence_type": { + "name": "cadence_type", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "cron_expression": { + "name": "cron_expression", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "cron_timezone": { + "name": "cron_timezone", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "interval_ms": { + "name": "interval_ms", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "target_type": { + "name": "target_type", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "target_session_id": { + "name": "target_session_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "target_config_json": { + "name": "target_config_json", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "status": { + "name": "status", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "completed_reason": { + "name": "completed_reason", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "misfire_policy": { + "name": "misfire_policy", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "next_run_at": { + "name": "next_run_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "last_run_at": { + "name": "last_run_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "run_count": { + "name": "run_count", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": 0 + }, + "max_runs": { + "name": "max_runs", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "expires_at": { + "name": "expires_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "updated_at": { + "name": "updated_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "schedules_next_run_at_idx": { + "name": "schedules_next_run_at_idx", + "columns": ["next_run_at"], + "isUnique": false + } + }, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "session_resources": { + "name": "session_resources", + "columns": { + "resource_id": { + "name": "resource_id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "session_id": { + "name": "session_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "direction": { + "name": "direction", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "name": { + "name": "name", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "kind": { + "name": "kind", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "status": { + "name": "status", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "locator_type": { + "name": "locator_type", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "locator": { + "name": "locator", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "normalized_locator_key": { + "name": "normalized_locator_key", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "attachment_id": { + "name": "attachment_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "mime_type": { + "name": "mime_type", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "size_bytes": { + "name": "size_bytes", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "error": { + "name": "error", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "updated_at": { + "name": "updated_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "session_resources_session_idx": { + "name": "session_resources_session_idx", + "columns": ["session_id"], + "isUnique": false + }, + "session_resources_locator_idx": { + "name": "session_resources_locator_idx", + "columns": ["session_id", "normalized_locator_key"], + "isUnique": true + } + }, + "foreignKeys": { + "session_resources_session_id_sessions_session_id_fk": { + "name": "session_resources_session_id_sessions_session_id_fk", + "tableFrom": "session_resources", + "tableTo": "sessions", + "columnsFrom": ["session_id"], + "columnsTo": ["session_id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "session_runs": { + "name": "session_runs", + "columns": { + "id": { + "name": "id", + "type": "integer", + "primaryKey": true, + "notNull": true, + "autoincrement": true + }, + "session_id": { + "name": "session_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "seq": { + "name": "seq", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "run_id": { + "name": "run_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "base_turn_id": { + "name": "base_turn_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "history_id": { + "name": "history_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "account_id": { + "name": "account_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "model": { + "name": "model", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "effort": { + "name": "effort", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "approval_policy_id": { + "name": "approval_policy_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "started_at": { + "name": "started_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "ended_at": { + "name": "ended_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + } + }, + "indexes": { + "session_runs_session_id_idx": { + "name": "session_runs_session_id_idx", + "columns": ["session_id"], + "isUnique": false + }, + "session_runs_run_id_unique": { + "name": "session_runs_run_id_unique", + "columns": ["run_id"], + "isUnique": true + } + }, + "foreignKeys": { + "session_runs_session_id_sessions_session_id_fk": { + "name": "session_runs_session_id_sessions_session_id_fk", + "tableFrom": "session_runs", + "tableTo": "sessions", + "columnsFrom": ["session_id"], + "columnsTo": ["session_id"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "sessions": { + "name": "sessions", + "columns": { + "session_id": { + "name": "session_id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "kind": { + "name": "kind", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "cwd": { + "name": "cwd", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "title": { + "name": "title", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "origin_type": { + "name": "origin_type", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "origin_history_id": { + "name": "origin_history_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "origin_imported_at": { + "name": "origin_imported_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "origin_source_session_id": { + "name": "origin_source_session_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "origin_source_turn_id": { + "name": "origin_source_turn_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "origin_forked_at": { + "name": "origin_forked_at", + "type": "integer", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "created_via": { + "name": "created_via", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "automation_kind": { + "name": "automation_kind", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "automation_id": { + "name": "automation_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "active_leaf_turn_id": { + "name": "active_leaf_turn_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "graph_revision": { + "name": "graph_revision", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": 0 + }, + "event_epoch": { + "name": "event_epoch", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": 0 + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "updated_at": { + "name": "updated_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "sessions_updated_at_idx": { + "name": "sessions_updated_at_idx", + "columns": ["updated_at"], + "isUnique": false + } + }, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "upload_leases": { + "name": "upload_leases", + "columns": { + "upload_id": { + "name": "upload_id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "declared_sha256": { + "name": "declared_sha256", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "declared_size": { + "name": "declared_size", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "name": { + "name": "name", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "mime_type": { + "name": "mime_type", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "kind": { + "name": "kind", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "blob_id": { + "name": "blob_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "attachment_id": { + "name": "attachment_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "expires_at": { + "name": "expires_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "upload_leases_expires_at_idx": { + "name": "upload_leases_expires_at_idx", + "columns": ["expires_at"], + "isUnique": false + } + }, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "workspaces": { + "name": "workspaces", + "columns": { + "workspace_id": { + "name": "workspace_id", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "cwd": { + "name": "cwd", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "name": { + "name": "name", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "kind": { + "name": "kind", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false, + "default": "'project'" + }, + "parent_workspace_id": { + "name": "parent_workspace_id", + "type": "text", + "primaryKey": false, + "notNull": false, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "last_used_at": { + "name": "last_used_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "workspaces_cwd_unique": { + "name": "workspaces_cwd_unique", + "columns": ["cwd"], + "isUnique": true + }, + "workspaces_last_used_at_idx": { + "name": "workspaces_last_used_at_idx", + "columns": ["last_used_at"], + "isUnique": false + } + }, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "worktree_sessions": { + "name": "worktree_sessions", + "columns": { + "worktree_path": { + "name": "worktree_path", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "session_id": { + "name": "session_id", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "worktree_sessions_session_unique": { + "name": "worktree_sessions_session_unique", + "columns": ["session_id"], + "isUnique": true + } + }, + "foreignKeys": { + "worktree_sessions_worktree_path_worktrees_worktree_path_fk": { + "name": "worktree_sessions_worktree_path_worktrees_worktree_path_fk", + "tableFrom": "worktree_sessions", + "tableTo": "worktrees", + "columnsFrom": ["worktree_path"], + "columnsTo": ["worktree_path"], + "onDelete": "cascade", + "onUpdate": "no action" + } + }, + "compositePrimaryKeys": { + "worktree_sessions_worktree_path_session_id_pk": { + "columns": ["worktree_path", "session_id"], + "name": "worktree_sessions_worktree_path_session_id_pk" + } + }, + "uniqueConstraints": {}, + "checkConstraints": {} + }, + "worktrees": { + "name": "worktrees", + "columns": { + "worktree_path": { + "name": "worktree_path", + "type": "text", + "primaryKey": true, + "notNull": true, + "autoincrement": false + }, + "repo_root": { + "name": "repo_root", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "branch": { + "name": "branch", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "created_at": { + "name": "created_at", + "type": "integer", + "primaryKey": false, + "notNull": true, + "autoincrement": false + }, + "state": { + "name": "state", + "type": "text", + "primaryKey": false, + "notNull": true, + "autoincrement": false + } + }, + "indexes": { + "worktrees_repo_root_branch_unique": { + "name": "worktrees_repo_root_branch_unique", + "columns": ["repo_root", "branch"], + "isUnique": true + } + }, + "foreignKeys": {}, + "compositePrimaryKeys": {}, + "uniqueConstraints": {}, + "checkConstraints": {} + } + }, + "views": {}, + "enums": {}, + "_meta": { + "schemas": {}, + "tables": {}, + "columns": {} + }, + "internal": { + "indexes": {} + } +} diff --git a/apps/daemon/drizzle/meta/_journal.json b/apps/daemon/drizzle/meta/_journal.json index 14611f1c9..5e9d6ef06 100644 --- a/apps/daemon/drizzle/meta/_journal.json +++ b/apps/daemon/drizzle/meta/_journal.json @@ -106,6 +106,13 @@ "when": 1788416689309, "tag": "0014_elite_mister_fear", "breakpoints": true + }, + { + "idx": 15, + "version": "6", + "when": 1788854811240, + "tag": "0015_fast_martin_li", + "breakpoints": true } ] } diff --git a/apps/daemon/src/__tests__/worktree-store.test.ts b/apps/daemon/src/__tests__/worktree-store.test.ts new file mode 100644 index 000000000..f5f4f9167 --- /dev/null +++ b/apps/daemon/src/__tests__/worktree-store.test.ts @@ -0,0 +1,167 @@ +import { mkdtemp, rm } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { WorktreeUnavailableError } from '@linkcode/engine'; +import type { SessionId, WorktreeRecord } from '@linkcode/schema'; +import Sqlite from 'better-sqlite3'; +import { readMigrationFiles } from 'drizzle-orm/migrator'; +import { afterEach, describe, expect, it } from 'vitest'; +import { daemonMigrationsFolder } from '../database-migrations'; +import type { DaemonDatabase } from '../db/database'; +import { openDaemonDatabase } from '../db/database'; +import { createWorktreeStore } from '../worktree-store'; + +const temporaryDirectories: string[] = []; +const openDatabases = new Set(); + +const s1 = 's-1' as SessionId; +const s2 = 's-2' as SessionId; +const s3 = 's-3' as SessionId; + +afterEach(async () => { + for (const database of openDatabases) database.close(); + openDatabases.clear(); + await Promise.all( + temporaryDirectories.splice(0).map((path) => rm(path, { recursive: true, force: true })), + ); +}); + +async function databasePath(): Promise { + const directory = await mkdtemp(join(tmpdir(), 'linkcode-worktree-store-')); + temporaryDirectories.push(directory); + return join(directory, 'daemon.db'); +} + +async function openStore() { + const database = openDaemonDatabase(await databasePath()); + openDatabases.add(database); + return { database, store: createWorktreeStore(database.client) }; +} + +function record(worktreePath: string, branch = 'feature'): WorktreeRecord { + return { worktreePath, repoRoot: '/repo', branch, createdAt: 1, state: 'active' }; +} + +/** Migrate a fresh file up to (excluding) the lease migration, the way drizzle's migrator would + * have left a daemon that shut down before it shipped. */ +function openPreLeaseDatabase(path: string): Sqlite.Database { + const sqlite = new Sqlite(path); + sqlite.exec( + 'CREATE TABLE IF NOT EXISTS "__drizzle_migrations" (id SERIAL PRIMARY KEY, hash text NOT NULL, created_at numeric)', + ); + const applied = sqlite.prepare( + 'INSERT INTO "__drizzle_migrations" ("hash", "created_at") VALUES (?, ?)', + ); + const migrations = readMigrationFiles({ migrationsFolder: daemonMigrationsFolder }); + for (let i = 0, len = migrations.length; i < len; i++) { + const migration = migrations[i]; + if (migration.sql.some((statement) => statement.includes('worktree_sessions'))) break; + for (let j = 0, statements = migration.sql.length; j < statements; j++) { + sqlite.exec(migration.sql[j]); + } + applied.run(migration.hash, migration.folderMillis); + } + return sqlite; +} + +describe('SQLite worktree store', () => { + it('round-trips records and updates them by worktree path', async () => { + const { store } = await openStore(); + const active = record('/wt/a'); + const orphan: WorktreeRecord = { ...record('/wt/old', 'old'), state: 'orphaned' }; + await store.save(active); + await store.save(orphan); + expect((await store.load()).worktrees).toEqual([active, orphan]); + await store.save({ ...active, state: 'orphaned' }); + expect((await store.load()).worktrees).toEqual([{ ...active, state: 'orphaned' }, orphan]); + }); + + it('marks the worktree deleting with its last release and refuses new leases from then on', async () => { + const { store } = await openStore(); + await store.save(record('/wt/a')); + await store.acquireLease('/wt/a', s1); + await store.acquireLease('/wt/a', s1); + await store.acquireLease('/wt/a', s2); + expect((await store.load()).leases.map((lease) => lease.sessionId)).toEqual([s1, s2]); + + expect(await store.releaseLease(s1)).toEqual({ worktreePath: '/wt/a', last: false }); + expect((await store.load()).worktrees).toEqual([record('/wt/a')]); + expect(await store.releaseLease(s1)).toBeUndefined(); + + expect(await store.releaseLease(s2)).toEqual({ worktreePath: '/wt/a', last: true }); + expect((await store.load()).worktrees).toMatchObject([{ state: 'deleting' }]); + await expect(store.acquireLease('/wt/a', s3)).rejects.toBeInstanceOf(WorktreeUnavailableError); + }); + + it('refuses a lease on an unknown worktree or a second worktree for one session', async () => { + const { store } = await openStore(); + await expect(store.acquireLease('/wt/missing', s1)).rejects.toBeInstanceOf( + WorktreeUnavailableError, + ); + await store.save(record('/wt/a')); + await store.save(record('/wt/b', 'other')); + await store.acquireLease('/wt/a', s1); + await expect(store.acquireLease('/wt/b', s1)).rejects.toThrow('Failed to lease worktree'); + expect((await store.load()).leases).toMatchObject([{ worktreePath: '/wt/a', sessionId: s1 }]); + }); + + it('keeps one worktree per repository branch and drops the leases with the worktree', async () => { + const { store } = await openStore(); + await store.save(record('/wt/a')); + await expect(store.save(record('/wt/b'))).rejects.toThrow('Failed to save worktree'); + await store.acquireLease('/wt/a', s1); + await store.delete('/wt/a'); + expect(await store.load()).toEqual({ worktrees: [], leases: [] }); + }); + + it('backfills leases from the owning-session column when migrating a pre-lease database', async () => { + const path = await databasePath(); + const legacy = openPreLeaseDatabase(path); + legacy + .prepare( + 'INSERT INTO worktrees (worktree_path, repo_root, branch, session_id, created_at, state) VALUES (?, ?, ?, ?, ?, ?)', + ) + .run('/wt/held', '/repo', 'feature', 's-legacy', 5, 'active'); + legacy + .prepare( + 'INSERT INTO worktrees (worktree_path, repo_root, branch, session_id, created_at, state) VALUES (?, ?, ?, ?, ?, ?)', + ) + .run('/wt/orphan', '/repo', 'stale', 'orphan-worktree-0123456789ab', 6, 'orphaned'); + legacy.close(); + + const database = openDaemonDatabase(path); + openDatabases.add(database); + expect(await createWorktreeStore(database.client).load()).toEqual({ + worktrees: [ + { + worktreePath: '/wt/held', + repoRoot: '/repo', + branch: 'feature', + createdAt: 5, + state: 'active', + }, + { + worktreePath: '/wt/orphan', + repoRoot: '/repo', + branch: 'stale', + createdAt: 6, + state: 'orphaned', + }, + ], + leases: [{ worktreePath: '/wt/held', sessionId: 's-legacy', createdAt: 5 }], + }); + + // The physical schema, not just the projection: the column and its index are gone, so a + // migration that forgot either DROP would be caught here rather than passing through load(). + const raw = new Sqlite(path); + const columns = (raw.prepare('PRAGMA table_info(worktrees)').all() as Array<{ name: string }>) + .map((column) => column.name) + .sort(); + expect(columns).toEqual(['branch', 'created_at', 'repo_root', 'state', 'worktree_path']); + const indexes = ( + raw.prepare('PRAGMA index_list(worktrees)').all() as Array<{ name: string }> + ).map((index) => index.name); + expect(indexes).not.toContain('worktrees_session_id_unique'); + raw.close(); + }); +}); diff --git a/apps/daemon/src/db/schema.ts b/apps/daemon/src/db/schema.ts index 78c2c8f52..76458146e 100644 --- a/apps/daemon/src/db/schema.ts +++ b/apps/daemon/src/db/schema.ts @@ -380,20 +380,36 @@ export const workspaces = sqliteTable( (table) => [index('workspaces_last_used_at_idx').on(table.lastUsedAt)], ); -/** Managed git worktrees. Session ids intentionally have no FK: rows survive session deletion until - * the dedicated cleanup lifecycle owns removal. */ +/** Managed git worktrees (`WorktreeRecord`); the sessions holding one are `worktree_sessions` + * rows. Written by ../worktree-store.ts on the shared connection: releasing the last lease and + * marking the row `deleting` must be one transaction. */ export const worktrees = sqliteTable( 'worktrees', { worktreePath: text('worktree_path').primaryKey(), repoRoot: text('repo_root').notNull(), branch: text('branch').notNull(), + createdAt: integer('created_at').notNull(), + state: text('state', { enum: ['active', 'orphaned', 'deleting'] }).notNull(), + }, + (table) => [uniqueIndex('worktrees_repo_root_branch_unique').on(table.repoRoot, table.branch)], +); + +/** Session leases on managed worktrees (`WorktreeLease`). Session ids intentionally have no FK: + * the cleanup lifecycle owns a lease's removal (the last one marks the worktree `deleting` before + * any filesystem work), and boot reconcile sweeps leases whose session is gone. */ +export const worktreeSessions = sqliteTable( + 'worktree_sessions', + { + worktreePath: text('worktree_path') + .notNull() + .references(() => worktrees.worktreePath, { onDelete: 'cascade' }), sessionId: text('session_id').notNull(), createdAt: integer('created_at').notNull(), - state: text('state', { enum: ['active', 'orphaned'] }).notNull(), }, (table) => [ - uniqueIndex('worktrees_repo_root_branch_unique').on(table.repoRoot, table.branch), - uniqueIndex('worktrees_session_id_unique').on(table.sessionId), + primaryKey({ columns: [table.worktreePath, table.sessionId] }), + // A session works in one directory; the lease is where it works. + uniqueIndex('worktree_sessions_session_unique').on(table.sessionId), ], ); diff --git a/apps/daemon/src/index.ts b/apps/daemon/src/index.ts index 5f69a83a3..2fbfb53a3 100644 --- a/apps/daemon/src/index.ts +++ b/apps/daemon/src/index.ts @@ -287,7 +287,7 @@ async function main(): Promise { scheduleStore: createScheduleStore(databasePath()), loopStore: createLoopStore(databasePath()), workspaceStore: createWorkspaceStore(databasePath()), - worktreeStore: createWorktreeStore(databasePath()), + worktreeStore: createWorktreeStore(database.client), worktreeRoot: worktreeRoot(), previewRoutes, browserToolsEnabled: process.env.LINKCODE_BROWSER_TOOLS === '1', diff --git a/apps/daemon/src/worktree-store.ts b/apps/daemon/src/worktree-store.ts index 5361e4264..5febb73a7 100644 --- a/apps/daemon/src/worktree-store.ts +++ b/apps/daemon/src/worktree-store.ts @@ -1,32 +1,37 @@ -import { mkdirSync } from 'node:fs'; -import { dirname } from 'node:path'; -import { fileURLToPath } from 'node:url'; -import type { WorktreeStore } from '@linkcode/engine'; -import type { WorktreeRecord } from '@linkcode/schema'; -import { WorktreeRecordSchema } from '@linkcode/schema'; -import Sqlite from 'better-sqlite3'; -import { eq } from 'drizzle-orm'; -import { drizzle } from 'drizzle-orm/better-sqlite3'; -import { migrate } from 'drizzle-orm/better-sqlite3/migrator'; -import { worktrees } from './db/schema'; - -export function createWorktreeStore(dbPath: string): WorktreeStore { - if (dbPath !== ':memory:') mkdirSync(dirname(dbPath), { recursive: true }); - const sqlite = new Sqlite(dbPath); - sqlite.pragma('journal_mode = WAL'); - const db = drizzle(sqlite); - migrate(db, { migrationsFolder: fileURLToPath(new URL('../drizzle', import.meta.url)) }); +import type { WorktreeLeaseRelease, WorktreeStore, WorktreeStoreSnapshot } from '@linkcode/engine'; +import { WorktreeUnavailableError } from '@linkcode/engine'; +import type { SessionId, WorktreeRecord } from '@linkcode/schema'; +import { WorktreeLeaseSchema, WorktreeRecordSchema } from '@linkcode/schema'; +import { count, eq } from 'drizzle-orm'; +import type { DaemonDatabaseClient } from './db/database'; +import { worktreeSessions, worktrees } from './db/schema'; +/** + * SQLite-backed `WorktreeStore` on the daemon's shared graph/session connection, so a lease + * release and the `deleting` mark it may cause are one transaction. Rows are validated back + * through the zod schemas on load. + */ +export function createWorktreeStore(db: DaemonDatabaseClient): WorktreeStore { return { - load(): Promise { - return Promise.resolve( - db - .select() - .from(worktrees) - .all() - .map((row) => WorktreeRecordSchema.parse(row)), - ); + load(): Promise { + try { + return Promise.resolve({ + worktrees: db + .select() + .from(worktrees) + .all() + .map((row) => WorktreeRecordSchema.parse(row)), + leases: db + .select() + .from(worktreeSessions) + .all() + .map((row) => WorktreeLeaseSchema.parse(row)), + }); + } catch (error) { + return Promise.reject(new Error('Failed to load worktrees', { cause: error })); + } }, + save(record: WorktreeRecord): Promise { try { db.insert(worktrees) @@ -38,9 +43,76 @@ export function createWorktreeStore(dbPath: string): WorktreeStore { return Promise.reject(new Error('Failed to save worktree', { cause: error })); } }, + delete(worktreePath): Promise { - db.delete(worktrees).where(eq(worktrees.worktreePath, worktreePath)).run(); - return Promise.resolve(); + try { + // Leases cascade with the row. + db.delete(worktrees).where(eq(worktrees.worktreePath, worktreePath)).run(); + return Promise.resolve(); + } catch (error) { + return Promise.reject(new Error('Failed to delete worktree', { cause: error })); + } + }, + + acquireLease(worktreePath: string, sessionId: SessionId): Promise { + try { + db.transaction((tx) => { + const row = tx + .select({ state: worktrees.state }) + .from(worktrees) + .where(eq(worktrees.worktreePath, worktreePath)) + .get(); + if (row === undefined || row.state === 'deleting') { + throw new WorktreeUnavailableError(worktreePath); + } + // Idempotent for the same worktree; the session index refuses a second worktree. + tx.insert(worktreeSessions) + .values({ worktreePath, sessionId, createdAt: Date.now() }) + .onConflictDoNothing({ + target: [worktreeSessions.worktreePath, worktreeSessions.sessionId], + }) + .run(); + }); + return Promise.resolve(); + } catch (error) { + return Promise.reject( + error instanceof WorktreeUnavailableError + ? error + : new Error('Failed to lease worktree', { cause: error }), + ); + } + }, + + releaseLease(sessionId: SessionId): Promise { + try { + const released = db.transaction((tx) => { + const lease = tx + .select({ worktreePath: worktreeSessions.worktreePath }) + .from(worktreeSessions) + .where(eq(worktreeSessions.sessionId, sessionId)) + .get(); + if (lease === undefined) return; + tx.delete(worktreeSessions).where(eq(worktreeSessions.sessionId, sessionId)).run(); + const remaining = tx + .select({ value: count() }) + .from(worktreeSessions) + .where(eq(worktreeSessions.worktreePath, lease.worktreePath)) + .get(); + const last = (remaining?.value ?? 0) === 0; + // The `deleting` mark lands with the release: nothing can lease the directory that + // cleanup is about to remove. + if (last) { + tx.update(worktrees) + .set({ state: 'deleting' }) + .where(eq(worktrees.worktreePath, lease.worktreePath)) + .run(); + } + return { worktreePath: lease.worktreePath, last }; + }); + return Promise.resolve(released); + } catch (error) { + return Promise.reject(new Error('Failed to release worktree lease', { cause: error })); + } }, }; } diff --git a/apps/daemon/tests/integration/worktree-store.test.ts b/apps/daemon/tests/integration/worktree-store.test.ts deleted file mode 100644 index 680e713f1..000000000 --- a/apps/daemon/tests/integration/worktree-store.test.ts +++ /dev/null @@ -1,57 +0,0 @@ -import type { WorktreeRecord } from '@linkcode/schema'; -import { SessionIdSchema } from '@linkcode/schema'; -import { describe, expect, it } from 'vitest'; -import { createWorktreeStore } from '../../src/worktree-store'; - -function record(values: Partial = {}): WorktreeRecord { - return { - worktreePath: '/managed/repo-feature', - repoRoot: '/repo', - branch: 'feature', - sessionId: SessionIdSchema.parse('sess-1'), - createdAt: 123, - state: 'active', - ...values, - }; -} - -describe('daemon sqlite worktree store', () => { - it('round-trips active and orphaned rows without requiring a session row', async () => { - const store = createWorktreeStore(':memory:'); - const rows = [ - record(), - record({ - worktreePath: '/managed/repo-old', - branch: 'old', - sessionId: SessionIdSchema.parse('sess-missing'), - state: 'orphaned', - }), - ]; - await Promise.all(rows.map((row) => store.save(row))); - expect(await store.load()).toEqual(rows); - }); - - it('enforces one row per normalized repository and branch', async () => { - const store = createWorktreeStore(':memory:'); - await store.save(record()); - await expect( - store.save( - record({ - worktreePath: '/managed/duplicate', - sessionId: SessionIdSchema.parse('sess-2'), - }), - ), - ).rejects.toThrow(); - }); - - it('updates and deletes a record by its worktree path', async () => { - const store = createWorktreeStore(':memory:'); - const active = record(); - await store.save(active); - await store.save({ ...active, state: 'orphaned' }); - expect(await store.load()).toEqual([{ ...active, state: 'orphaned' }]); - - await store.delete(active.worktreePath); - expect(await store.load()).toEqual([]); - }); -}); diff --git a/packages/foundation/schema/src/model/worktree.ts b/packages/foundation/schema/src/model/worktree.ts index 83ab603c3..397969b2b 100644 --- a/packages/foundation/schema/src/model/worktree.ts +++ b/packages/foundation/schema/src/model/worktree.ts @@ -1,16 +1,26 @@ import { z } from 'zod'; import { SessionIdSchema, TimestampSchema } from './primitives'; -export const WorktreeStateSchema = z.enum(['active', 'orphaned']); +/** `deleting`: the last lease was released and cleanup owns the directory — no new lease lands. */ +export const WorktreeStateSchema = z.enum(['active', 'orphaned', 'deleting']); export type WorktreeState = z.infer; -/** Durable ownership record for a LinkCode-managed git worktree. */ +/** Durable record of a LinkCode-managed git worktree; the sessions holding it are its leases. */ export const WorktreeRecordSchema = z.object({ worktreePath: z.string().min(1), repoRoot: z.string().min(1), branch: z.string().min(1), - sessionId: SessionIdSchema, createdAt: TimestampSchema, state: WorktreeStateSchema, }); export type WorktreeRecord = z.infer; + +/** A session's hold on a managed worktree. A worktree lives while any lease does; a fork shares + * its source's, and releasing the last one marks the worktree `deleting` before any filesystem + * work. A session holds at most one. */ +export const WorktreeLeaseSchema = z.object({ + worktreePath: z.string().min(1), + sessionId: SessionIdSchema, + createdAt: TimestampSchema, +}); +export type WorktreeLease = z.infer; diff --git a/packages/host/engine/src/__tests__/worktree-store.test.ts b/packages/host/engine/src/__tests__/worktree-store.test.ts new file mode 100644 index 000000000..a98bd2d55 --- /dev/null +++ b/packages/host/engine/src/__tests__/worktree-store.test.ts @@ -0,0 +1,58 @@ +import type { SessionId, WorktreeRecord } from '@linkcode/schema'; +import { describe, expect, it } from 'vitest'; +import { InMemoryWorktreeStore, WorktreeUnavailableError } from '../worktree/worktree-store'; + +const s1 = 's-1' as SessionId; +const s2 = 's-2' as SessionId; +const s3 = 's-3' as SessionId; + +function record(worktreePath: string, branch = 'feature'): WorktreeRecord { + return { worktreePath, repoRoot: '/repo', branch, createdAt: 1, state: 'active' }; +} + +describe('InMemoryWorktreeStore leases', () => { + it('keeps one worktree per repository branch', async () => { + const store = new InMemoryWorktreeStore(); + await store.save(record('/wt/a')); + await expect(store.save(record('/wt/b'))).rejects.toThrow('worktree already exists'); + await store.save({ ...record('/wt/a'), state: 'orphaned' }); + expect((await store.load()).worktrees).toMatchObject([{ state: 'orphaned' }]); + }); + + it('marks the worktree deleting with its last release and refuses new leases from then on', async () => { + const store = new InMemoryWorktreeStore(); + await store.save(record('/wt/a')); + await store.acquireLease('/wt/a', s1); + await store.acquireLease('/wt/a', s1); + await store.acquireLease('/wt/a', s2); + expect((await store.load()).leases.map((lease) => lease.sessionId)).toEqual([s1, s2]); + + expect(await store.releaseLease(s1)).toEqual({ worktreePath: '/wt/a', last: false }); + expect((await store.load()).worktrees).toMatchObject([{ state: 'active' }]); + expect(await store.releaseLease(s1)).toBeUndefined(); + + expect(await store.releaseLease(s2)).toEqual({ worktreePath: '/wt/a', last: true }); + expect((await store.load()).worktrees).toMatchObject([{ state: 'deleting' }]); + await expect(store.acquireLease('/wt/a', s3)).rejects.toBeInstanceOf(WorktreeUnavailableError); + }); + + it('refuses a lease on an unknown worktree or a second worktree for one session', async () => { + const store = new InMemoryWorktreeStore(); + await expect(store.acquireLease('/wt/missing', s1)).rejects.toBeInstanceOf( + WorktreeUnavailableError, + ); + await store.save(record('/wt/a')); + await store.save(record('/wt/b', 'other')); + await store.acquireLease('/wt/a', s1); + await expect(store.acquireLease('/wt/b', s1)).rejects.toThrow('already holds a worktree'); + }); + + it('drops the leases with the worktree', async () => { + const store = new InMemoryWorktreeStore(); + await store.save(record('/wt/a')); + await store.acquireLease('/wt/a', s1); + await store.delete('/wt/a'); + expect(await store.load()).toEqual({ worktrees: [], leases: [] }); + expect(await store.releaseLease(s1)).toBeUndefined(); + }); +}); diff --git a/packages/host/engine/src/engine.ts b/packages/host/engine/src/engine.ts index a60faa260..306241900 100644 --- a/packages/host/engine/src/engine.ts +++ b/packages/host/engine/src/engine.ts @@ -215,6 +215,12 @@ export const createEngineRuntime = Effect.fn('Engine.create')(function* ( attachmentStore, ); const conversationJournals = new ConversationLiveJournals(); + const git = deps.git ?? (yield* GitService.make()); + const worktrees = new WorktreeService( + deps.worktreeStore ?? new InMemoryWorktreeStore(), + deps.worktreeRoot, + git, + ); const sessions = new SessionOrchestrator( transport, factory, @@ -231,6 +237,7 @@ export const createEngineRuntime = Effect.fn('Engine.create')(function* ( conversationTurns, conversationJournals, ingest, + worktrees, deps.browserToolsEnabled ? () => new BrowserReplHost((op, args) => browserBroker.dispatch(op, args)) : undefined, @@ -254,12 +261,6 @@ export const createEngineRuntime = Effect.fn('Engine.create')(function* ( ); const workspaces = new WorkspaceRegistry(deps.workspaceStore ?? new InMemoryWorkspaceStore()); const workspaceRequests = new WorkspaceRequestHandler(transport, workspaces, responder); - const git = deps.git ?? (yield* GitService.make()); - const worktrees = new WorktreeService( - deps.worktreeStore ?? new InMemoryWorktreeStore(), - deps.worktreeRoot, - git, - ); const gitRequests = new GitRequestHandler(transport, git, responder); const fileSuggest = deps.fileSuggest ?? (yield* FileSuggestService.make()); const fileRequests = new FileRequestHandler( diff --git a/packages/host/engine/src/index.ts b/packages/host/engine/src/index.ts index 863cf0423..9dc939f18 100644 --- a/packages/host/engine/src/index.ts +++ b/packages/host/engine/src/index.ts @@ -47,4 +47,9 @@ export type { SimulatorMcpProvider } from './simulator/mcp'; export { SimulatorService } from './simulator/service'; export type { PtyBackend, PtyOpenOptions, PtyProcess } from './terminal/pty-backend'; export type { WorkspaceStore } from './workspace/workspace-store'; -export type { WorktreeStore } from './worktree/worktree-store'; +export { + type WorktreeLeaseRelease, + type WorktreeStore, + type WorktreeStoreSnapshot, + WorktreeUnavailableError, +} from './worktree/worktree-store'; diff --git a/packages/host/engine/src/session/fork-service.ts b/packages/host/engine/src/session/fork-service.ts index 1e226271a..e37a460c1 100644 --- a/packages/host/engine/src/session/fork-service.ts +++ b/packages/host/engine/src/session/fork-service.ts @@ -57,6 +57,10 @@ interface AdmittedFork { readonly path: ConversationTurn[]; readonly cut: ForkCut; readonly operation: OpenForkOperation; + /** The source's managed worktree, captured while the source was provably present under the + * semaphore; the child leases exactly this path, so a source deleted before the launch is a + * typed `conflict` at the store, not an unleased start in a vanishing directory. */ + readonly sourceWorktreePath: string | undefined; } /** @@ -165,8 +169,14 @@ export class SessionForkService { new RequestError({ code: 'busy', message: `Session is busy: ${source.sessionId}` }), ); } + // Captured synchronously while the source is provably present: a delete interleaving the + // async admit below cannot make the launch skip the child's lease. + const sourceWorktreePath = this.worktrees.get(source.sessionId)?.worktreePath; const { checkpoints, sessions, turns, worktrees } = this; return Effect.gen(function* () { + // The child starts in that directory: gone from disk (`orphaned`, or not yet marked so) + // it fails typed here, exactly as resuming the source does, not as a spawn error. + yield* worktrees.verifyResume(source.sessionId); if (yield* turns.hasOpenOperation(source.sessionId)) { return yield* Effect.fail( new RequestError({ @@ -195,16 +205,6 @@ export class SessionForkService { new RequestError({ code: 'conflict', message: 'The conversation graph has moved' }), ); } - // A managed worktree has exactly one owning session until worktree leases land; the - // child could not hold the working tree it would share. - if (worktrees.get(source.sessionId) !== undefined) { - return yield* Effect.fail( - new RequestError({ - code: 'unsupported', - message: 'Forking a session on a managed worktree is not supported yet', - }), - ); - } if (sessions.historyCapabilitiesOf(source.kind).forkAfterTurn !== true) { return yield* Effect.fail( new RequestError({ @@ -235,7 +235,7 @@ export class SessionForkService { createdAt: Date.now(), }; yield* turns.persistOperation(operation); - return { source, through, path, cut, operation }; + return { source, through, path, cut, operation, sourceWorktreePath }; }); }), ); @@ -254,10 +254,10 @@ export class SessionForkService { { readonly sessionId: SessionId; readonly mcpWarnings: readonly McpWarning[] }, EngineFailure > { - const { history, lifecycle, records, sessions, turns } = this; + const { history, lifecycle, records, sessions, turns, worktrees } = this; const abandon = this.abandon.bind(this); return Effect.gen(function* () { - const { source, through, path, cut, operation } = admitted; + const { source, through, path, cut, operation, sourceWorktreePath } = admitted; const childId = lifecycle.nextSessionId(); // The source's pins, resolved for the child: per-session resources such as the simulator // MCP endpoint token must belong to the child, or its tools act as the source's. @@ -309,6 +309,14 @@ export class SessionForkService { eventEpoch: 0, }; records.registerProvisional(child); + // The child leases its source's managed worktree before the adapter starts there. The path + // was captured at admit, so a source deleted meanwhile makes this acquire fail typed + // `conflict` (the store refuses a `deleting`/removed worktree) instead of starting the child + // in a directory being torn down. + const lease = + sourceWorktreePath === undefined + ? Effect.void + : worktrees.acquire(childId, sourceWorktreePath); const start = sessions .startLive( undefined, @@ -359,7 +367,8 @@ export class SessionForkService { ), Effect.uninterruptible, ); - yield* start.pipe( + yield* lease.pipe( + Effect.andThen(start), Effect.andThen(commit), Effect.onExit((exit) => Exit.isFailure(exit) && records.isProvisional(childId) ? abandon(child) : Effect.void, @@ -371,7 +380,7 @@ export class SessionForkService { /** A fork that never committed: stop the child adapter if it started, forget the record. */ private abandon(child: SessionRecord): Effect.Effect { - const { records, sessions } = this; + const { records, sessions, worktrees } = this; return Effect.suspend(() => { const stop = sessions.liveRunId(child.sessionId) === undefined @@ -388,6 +397,20 @@ export class SessionForkService { ), ); return stop.pipe( + // The child's lease goes with it; when the source vanished meanwhile this was the last one. + Effect.andThen( + worktrees + .cleanupDeletedSession(child.sessionId) + .pipe( + Effect.catch((error) => + Effect.logError( + 'Failed to release the abandoned fork child worktree lease', + { sessionId: child.sessionId }, + error.cause, + ), + ), + ), + ), Effect.andThen( Effect.sync(() => { records.discardProvisional(child.sessionId); diff --git a/packages/host/engine/src/session/lifecycle-service.ts b/packages/host/engine/src/session/lifecycle-service.ts index ae34aed70..5cf4e906b 100644 --- a/packages/host/engine/src/session/lifecycle-service.ts +++ b/packages/host/engine/src/session/lifecycle-service.ts @@ -412,13 +412,19 @@ export class SessionLifecycleService { startAdapter = (adapter) => history.branch(adapter, cut, resolved.options); } const runId = mintRunId(); - const intent = yield* turns.persistIntent({ - sessionId: sourceSessionId, - operationId: mintOperationId(), - runId, - parentTurnId, - input: { type: 'prompt', blocks: yield* ingest.promptBlocks(content) }, - }); + const blocks = yield* ingest.promptBlocks(content); + // The worktree gate: editing a prompt relaunches this session's adapter, so it is a + // turn start and must wait for a co-leaseholder's running turn like any other. + const intent = yield* sessions.admitTurn( + sourceSessionId, + turns.persistIntent({ + sessionId: sourceSessionId, + operationId: mintOperationId(), + runId, + parentTurnId, + input: { type: 'prompt', blocks }, + }), + ); yield* Effect.gen(function* () { yield* sessions.stopForReplacement(sourceSessionId); yield* launchRun(replyTo, source, resolved, startAdapter, { @@ -631,153 +637,163 @@ export class SessionLifecycleService { new RequestError({ code: 'busy', message: `Session is busy: ${request.sessionId}` }), ); } - // Seam: the worktree co-leaseholder busy gate joins this critical section later. const { checkpoints, records, sessions, turns } = this; const admitAttachments = this.admitPromptBlocks.bind(this); - return Effect.gen(function* () { - if (yield* turns.hasOpenOperation(request.sessionId)) { - return yield* Effect.fail( - new RequestError({ - code: 'busy', - message: 'Another operation is open on this session', - }), - ); - } - let parentTurnId: TurnId | null; - let launch: TurnLaunch; - if (request.parentTurnId === undefined) { - // Plain send: no guards — targets the current active leaf under the busy rules alone. - parentTurnId = record.activeLeafTurnId ?? null; - launch = { type: 'continue' }; - } else if (request.parentTurnId === null) { - if (request.expectedGraphRevision !== record.graphRevision) { + return sessions.admitTurn( + request.sessionId, + Effect.gen(function* () { + if (yield* turns.hasOpenOperation(request.sessionId)) { return yield* Effect.fail( - new RequestError({ code: 'conflict', message: 'The conversation graph has moved' }), + new RequestError({ + code: 'busy', + message: 'Another operation is open on this session', + }), ); } - parentTurnId = null; - // Editing "the first prompt" starts fresh only when nothing can precede a root here. - // The session's FIRST root answers that — a later root, relaunched fresh, would read - // the earlier runs as hidden history of its own — while the cut anchors on the active - // lineage's root, whose history holds whatever the hidden prefix is. - const existingTurns = yield* turns.listTurns(request.sessionId); - const firstRoot = existingTurns.find( - (turn) => turn.parentTurnId === null && turn.siblingOrdinal === 1, - ); - const activeRoot = - pathToLeaf( - new Map(existingTurns.map((turn) => [turn.turnId, turn])), - record.activeLeafTurnId, - ).at(0) ?? firstRoot; - const nothingPrecedes = - firstRoot === undefined - ? records.historyId(request.sessionId) === undefined - : !hasHiddenPrefix(record, firstRoot); - if (nothingPrecedes) { - launch = { type: 'fresh' }; - } else { - const forkable = sessions.historyCapabilitiesOf(record.kind).forkAfterTurn === true; - const cut = - forkable && activeRoot !== undefined - ? yield* checkpoints.forkCutBefore(record, activeRoot.turnId) - : undefined; - if (cut === undefined) { + let parentTurnId: TurnId | null; + let launch: TurnLaunch; + if (request.parentTurnId === undefined) { + // Plain send: no guards — targets the current active leaf under the busy rules alone. + parentTurnId = record.activeLeafTurnId ?? null; + launch = { type: 'continue' }; + } else if (request.parentTurnId === null) { + if (request.expectedGraphRevision !== record.graphRevision) { return yield* Effect.fail( new RequestError({ - code: 'unsupported', - message: forkable - ? 'This turn has no provider checkpoint to fork from' - : `${record.kind}: forking from an earlier turn is not supported`, + code: 'conflict', + message: 'The conversation graph has moved', }), ); } - launch = { type: 'fork', cut }; - } - } else { - const existingTurns = yield* turns.listTurns(request.sessionId); - const parent = existingTurns.find((turn) => turn.turnId === request.parentTurnId); - if (!parent) { - return yield* Effect.fail( - new RequestError({ - code: 'not_found', - message: `Unknown turn: ${request.parentTurnId}`, - }), - ); - } - if (parent.state !== 'completed') { - return yield* Effect.fail( - new RequestError({ - code: 'conflict', - message: 'The parent turn has not completed', - }), - ); - } - if (request.expectedGraphRevision !== record.graphRevision) { - return yield* Effect.fail( - new RequestError({ code: 'conflict', message: 'The conversation graph has moved' }), + parentTurnId = null; + // Editing "the first prompt" starts fresh only when nothing can precede a root here. + // The session's FIRST root answers that — a later root, relaunched fresh, would read + // the earlier runs as hidden history of its own — while the cut anchors on the active + // lineage's root, whose history holds whatever the hidden prefix is. + const existingTurns = yield* turns.listTurns(request.sessionId); + const firstRoot = existingTurns.find( + (turn) => turn.parentTurnId === null && turn.siblingOrdinal === 1, ); - } - parentTurnId = request.parentTurnId; - if (request.parentTurnId === record.activeLeafTurnId) { - // Tip-continue on the active lineage: the provider history head IS this leaf. - launch = { type: 'continue' }; - } else { - // A valid checkpoint forks — including at an inactive tip: pi's fork writes a new - // file, and the tip's own history may have grown outside LinkCode (CLI/TUI use), so - // a forking harness never continues a tip blind. Only a harness that cannot fork - // continues a tip by resuming the history its own run wrote to; an interior turn - // is fork-unavailable. - const forkable = sessions.historyCapabilitiesOf(record.kind).forkAfterTurn === true; - const cut = forkable - ? yield* checkpoints.forkCutAfter(record, parent.turnId) - : undefined; - if (cut !== undefined) { + const activeRoot = + pathToLeaf( + new Map(existingTurns.map((turn) => [turn.turnId, turn])), + record.activeLeafTurnId, + ).at(0) ?? firstRoot; + const nothingPrecedes = + firstRoot === undefined + ? records.historyId(request.sessionId) === undefined + : !hasHiddenPrefix(record, firstRoot); + if (nothingPrecedes) { + launch = { type: 'fresh' }; + } else { + const forkable = sessions.historyCapabilitiesOf(record.kind).forkAfterTurn === true; + const cut = + forkable && activeRoot !== undefined + ? yield* checkpoints.forkCutBefore(record, activeRoot.turnId) + : undefined; + if (cut === undefined) { + return yield* Effect.fail( + new RequestError({ + code: 'unsupported', + message: forkable + ? 'This turn has no provider checkpoint to fork from' + : `${record.kind}: forking from an earlier turn is not supported`, + }), + ); + } launch = { type: 'fork', cut }; - } else if (existingTurns.some((turn) => turn.parentTurnId === parent.turnId)) { + } + } else { + const existingTurns = yield* turns.listTurns(request.sessionId); + const parent = existingTurns.find((turn) => turn.turnId === request.parentTurnId); + if (!parent) { + return yield* Effect.fail( + new RequestError({ + code: 'not_found', + message: `Unknown turn: ${request.parentTurnId}`, + }), + ); + } + if (parent.state !== 'completed') { return yield* Effect.fail( new RequestError({ - code: 'unsupported', - message: forkable - ? 'This turn has no provider checkpoint to fork from' - : `${record.kind}: forking from an earlier turn is not supported`, + code: 'conflict', + message: 'The parent turn has not completed', }), ); - } else if (forkable) { + } + if (request.expectedGraphRevision !== record.graphRevision) { return yield* Effect.fail( new RequestError({ - code: 'unsupported', - message: 'This turn has no provider checkpoint to continue from', + code: 'conflict', + message: 'The conversation graph has moved', }), ); + } + parentTurnId = request.parentTurnId; + if (request.parentTurnId === record.activeLeafTurnId) { + // Tip-continue on the active lineage: the provider history head IS this leaf. + launch = { type: 'continue' }; } else { - const historyId = record.runs.find((run) => run.runId === parent.runId)?.historyId; - if (historyId === undefined) { + // A valid checkpoint forks — including at an inactive tip: pi's fork writes a new + // file, and the tip's own history may have grown outside LinkCode (CLI/TUI use), so + // a forking harness never continues a tip blind. Only a harness that cannot fork + // continues a tip by resuming the history its own run wrote to; an interior turn + // is fork-unavailable. + const forkable = sessions.historyCapabilitiesOf(record.kind).forkAfterTurn === true; + const cut = forkable + ? yield* checkpoints.forkCutAfter(record, parent.turnId) + : undefined; + if (cut !== undefined) { + launch = { type: 'fork', cut }; + } else if (existingTurns.some((turn) => turn.parentTurnId === parent.turnId)) { + return yield* Effect.fail( + new RequestError({ + code: 'unsupported', + message: forkable + ? 'This turn has no provider checkpoint to fork from' + : `${record.kind}: forking from an earlier turn is not supported`, + }), + ); + } else if (forkable) { return yield* Effect.fail( new RequestError({ code: 'unsupported', - message: 'This turn has no provider history to continue', + message: 'This turn has no provider checkpoint to continue from', }), ); + } else { + const historyId = record.runs.find( + (run) => run.runId === parent.runId, + )?.historyId; + if (historyId === undefined) { + return yield* Effect.fail( + new RequestError({ + code: 'unsupported', + message: 'This turn has no provider history to continue', + }), + ); + } + launch = { type: 'resume', historyId }; } - launch = { type: 'resume', historyId }; } } - } - const liveRunId = - launch.type === 'continue' ? sessions.liveRunId(request.sessionId) : undefined; - if (liveRunId === undefined && launch.type === 'continue') launch = { type: 'resume' }; - if (request.input.type === 'prompt') { - yield* admitAttachments(record.kind, request.input.blocks); - } - const intent = yield* turns.persistIntent({ - sessionId: request.sessionId, - operationId: request.operationId, - runId: liveRunId ?? mintRunId(), - parentTurnId, - input: request.input, - }); - return { intent, launch }; - }); + const liveRunId = + launch.type === 'continue' ? sessions.liveRunId(request.sessionId) : undefined; + if (liveRunId === undefined && launch.type === 'continue') launch = { type: 'resume' }; + if (request.input.type === 'prompt') { + yield* admitAttachments(record.kind, request.input.blocks); + } + const intent = yield* turns.persistIntent({ + sessionId: request.sessionId, + operationId: request.operationId, + runId: liveRunId ?? mintRunId(), + parentTurnId, + input: request.input, + }); + return { intent, launch }; + }), + ); }), ); } diff --git a/packages/host/engine/src/session/orchestrator.ts b/packages/host/engine/src/session/orchestrator.ts index cf08e3ca7..46be0309b 100644 --- a/packages/host/engine/src/session/orchestrator.ts +++ b/packages/host/engine/src/session/orchestrator.ts @@ -15,7 +15,7 @@ import type { import { userRowMessageId } from '@linkcode/schema'; import type { Transport } from '@linkcode/transport'; import { createWireMessage } from '@linkcode/transport'; -import { Cause, Deferred, Effect, Exit, Scope } from 'effect'; +import { Cause, Deferred, Effect, Exit, Scope, Semaphore } from 'effect'; import type { AgentRuntimeService } from '../agent/runtime-service'; import type { AttachmentIngest } from '../attachment/ingest'; import type { TurnResult } from '../automation/turn-watcher'; @@ -27,6 +27,7 @@ import type { EngineFailure } from '../failure'; import { OperationError, RequestError, toOperationFailure } from '../failure'; import { observeOperation, recordLiveSessions } from '../observability'; import type { ResourceService } from '../resource/service'; +import type { WorktreeService } from '../worktree/worktree-service'; import { LiveSession } from './live-session'; import { SessionEventProcessor } from './session-event-processor'; import { SessionInputDispatcher } from './session-input-dispatcher'; @@ -34,6 +35,7 @@ import type { SessionRecordRegistry } from './session-record-registry'; export class SessionOrchestrator { private readonly sessions = new Map(); + private readonly worktreeGates = new Map(); /** Sessions mid-`delete`: a launch admitted during the delete's own store waits must not install * a live run whose record is about to vanish (and whose journal the final drop would take). */ private readonly deleting = new Set(); @@ -52,6 +54,7 @@ export class SessionOrchestrator { private readonly turns: ConversationTurnService, private readonly journals: ConversationLiveJournals, private readonly ingest: AttachmentIngest, + private readonly worktrees: Pick, private readonly browserTools?: BrowserToolsetFactory, private readonly onRunEnded?: (sessionId: SessionId, runId: RunId) => void, ) { @@ -64,7 +67,63 @@ export class SessionOrchestrator { turns, journals, ); - this.inputs = new SessionInputDispatcher(records, this.events, resources, turns, ingest); + this.inputs = new SessionInputDispatcher( + records, + this.events, + resources, + turns, + ingest, + (sessionId, body) => this.admitTurn(sessionId, body), + ); + } + + /** + * Turn admission on a managed worktree: at most one leaseholder runs a turn, so `body` (the + * caller's admit-and-persist) runs under the worktree's permit after a typed `busy` for any + * co-leaseholder running or holding an admitted turn (its open operation) — two siblings + * admitted concurrently cannot both start one. Sessions without a worktree skip the permit but + * keep the deletion check: `delete` runs outside every caller's critical section, so the record + * a caller validated can be gone by the time this admits. + */ + admitTurn( + sessionId: SessionId, + body: Effect.Effect, + ): Effect.Effect { + const { deleting, records, turns, worktrees } = this; + const isBusy = this.isBusy.bind(this); + const worktree = worktrees.get(sessionId); + const admit = Effect.gen(function* () { + const siblings = worktree === undefined ? [] : worktrees.coLeaseholders(sessionId); + for (let i = 0, len = siblings.length; i < len; i++) { + const sibling = siblings[i]; + if (isBusy(sibling) || (yield* turns.hasOpenOperation(sibling))) { + return yield* Effect.fail( + new RequestError({ + code: 'busy', + message: 'Another session on this worktree is running a turn', + }), + ); + } + } + // Last check before anything durable is written under this session's id. + if (deleting.has(sessionId) || records.get(sessionId) === undefined) { + return yield* Effect.fail( + new RequestError({ code: 'not_found', message: `Unknown session: ${sessionId}` }), + ); + } + return yield* body; + }); + return worktree === undefined + ? admit + : this.worktreeGate(worktree.worktreePath).withPermit(admit); + } + + private worktreeGate(worktreePath: string): Semaphore.Semaphore { + const existing = this.worktreeGates.get(worktreePath); + if (existing) return existing; + const gate = Semaphore.makeUnsafe(1); + this.worktreeGates.set(worktreePath, gate); + return gate; } private get(sessionId: SessionId): LiveSession | undefined { @@ -226,20 +285,28 @@ export class SessionOrchestrator { session.turnInputActive = true; const content: ContentBlock[] = [{ type: 'text', text }]; const { ingest, records, turns } = this; + const admitTurn = this.admitTurn.bind(this); return session.run( Effect.gen({ self: this }, function* () { - if (yield* turns.hasOpenOperation(sessionId)) { - return yield* Effect.fail( - new RequestError({ code: 'busy', message: `Session is busy: ${sessionId}` }), - ); - } - const intent = yield* turns.persistIntent({ + // The worktree gate: an automation turn on a managed worktree is refused while a + // co-leaseholder runs, and its intent persists under the gate like every other turn. + const intent = yield* admitTurn( sessionId, - operationId: mintOperationId(), - runId: session.runId, - parentTurnId: records.get(sessionId)?.activeLeafTurnId ?? null, - input: { type: 'prompt', blocks: yield* ingest.promptBlocks(content) }, - }); + Effect.gen(function* () { + if (yield* turns.hasOpenOperation(sessionId)) { + return yield* Effect.fail( + new RequestError({ code: 'busy', message: `Session is busy: ${sessionId}` }), + ); + } + return yield* turns.persistIntent({ + sessionId, + operationId: mintOperationId(), + runId: session.runId, + parentTurnId: records.get(sessionId)?.activeLeafTurnId ?? null, + input: { type: 'prompt', blocks: yield* ingest.promptBlocks(content) }, + }); + }), + ); const result = yield* Effect.sync(() => { this.events.broadcast( sessionId, diff --git a/packages/host/engine/src/session/session-input-dispatcher.ts b/packages/host/engine/src/session/session-input-dispatcher.ts index 86fe6dab1..bc2b7245e 100644 --- a/packages/host/engine/src/session/session-input-dispatcher.ts +++ b/packages/host/engine/src/session/session-input-dispatcher.ts @@ -19,6 +19,12 @@ import type { LiveSession } from './live-session'; import type { SessionEventProcessor } from './session-event-processor'; import type { SessionRecordRegistry } from './session-record-registry'; +/** The orchestrator's worktree turn gate; see `SessionOrchestrator.admitTurn`. */ +type TurnAdmission = ( + sessionId: SessionId, + body: Effect.Effect, +) => Effect.Effect; + /** Validates and dispatches client input while preserving turn and response state transitions. */ export class SessionInputDispatcher { constructor( @@ -27,6 +33,7 @@ export class SessionInputDispatcher { private readonly resources: ResourceService, private readonly turns: ConversationTurnService, private readonly ingest: AttachmentIngest, + private readonly admitTurn: TurnAdmission, ) {} /** `prepared` is a submit-saga intent already persisted for this dispatch; without one, a @@ -71,11 +78,13 @@ export class SessionInputDispatcher { this.events.rejectInput(sessionId, session, error.message); return Effect.fail(error); } - const { events, ingest, records, resources, turns } = this; + const { admitTurn, events, ingest, records, resources, turns } = this; // Set synchronously, before the first await, so a same-tick second turn input cannot slip // past the gate above while this one is still validating; every failure exit releases it. if (startsTurn) session.turnInputActive = true; - return Effect.gen(function* () { + // Validation and the durable intent. A legacy turn input admits itself here, under the + // worktree gate a saga-prepared intent has already passed. + const prepare = Effect.gen(function* () { // A submit operation in flight owns the session; legacy inputs respect the same admit gate. if (startsTurn && prepared === undefined && (yield* turns.hasOpenOperation(sessionId))) { const error = new RequestError({ @@ -133,7 +142,12 @@ export class SessionInputDispatcher { : input, }); } - const persisted = intent; + return { adapterInput, persisted: intent }; + }); + return Effect.gen(function* () { + const { adapterInput, persisted } = yield* startsTurn && prepared === undefined + ? admitTurn(sessionId, prepare) + : prepare; const persistedTurnId = startsTurn ? nullthrow(persisted, 'turn input without a persisted turn').turn.turnId : undefined; diff --git a/packages/host/engine/src/worktree/worktree-service.ts b/packages/host/engine/src/worktree/worktree-service.ts index 15d44d766..7bd09ffe8 100644 --- a/packages/host/engine/src/worktree/worktree-service.ts +++ b/packages/host/engine/src/worktree/worktree-service.ts @@ -2,7 +2,6 @@ import { createHash } from 'node:crypto'; import { existsSync, readdirSync } from 'node:fs'; import { basename, join, normalize, resolve } from 'node:path'; import type { BranchMode, SessionId, StartOptions, WorktreeRecord } from '@linkcode/schema'; -import { SessionIdSchema } from '@linkcode/schema'; import { Effect, Exit, Semaphore } from 'effect'; import type { EngineFailure } from '../failure'; import { OperationError, RequestError } from '../failure'; @@ -20,15 +19,24 @@ import { switchBranch, } from '../git/worktrees'; import type { WorktreeStore } from './worktree-store'; +import { WorktreeUnavailableError } from './worktree-store'; const RE_UNSAFE_SLUG = /[^\w.-]+/g; const RE_EDGE_DASHES = /^-+|-+$/g; const RE_CHECKED_OUT = /already checked out|already used by worktree/i; const RE_SWITCH_BLOCKED = /would be overwritten|please commit your changes or stash/i; +/** + * Managed git worktrees and the session leases on them. A worktree is provisioned for one session + * and shared by every session forked from it; it is cleaned up when its LAST lease goes, and only + * then — the release marks the row `deleting` durably before any filesystem work, so a fork racing + * the cleanup is refused typed instead of landing on a directory about to vanish. + */ export class WorktreeService { - private readonly bySession = new Map(); + private readonly byPath = new Map(); private readonly byRepoBranch = new Map(); + /** The worktree each leasing session holds, by the record's own `worktreePath`. */ + private readonly leaseBySession = new Map(); private readonly semaphores = new Map(); constructor( @@ -43,16 +51,31 @@ export class WorktreeService { return storeEffect('worktrees.load', 'Failed to load managed worktrees', () => this.store.load(), ).pipe( - Effect.tap((records) => + Effect.tap(({ worktrees }) => Effect.sync(() => { - for (let i = 0, len = records.length; i < len; i++) { - const record = records[i]; - this.bySession.set(record.sessionId, record); - this.byRepoBranch.set(repoBranchKey(record.repoRoot, record.branch), record); - } + for (let i = 0, len = worktrees.length; i < len; i++) this.index(worktrees[i]); }), ), - Effect.andThen(Effect.suspend(() => this.reconcile(durableSessionIds))), + // A lease whose session is gone (deleted while the daemon was down) is swept durably; the + // worktree it held then reconciles below like any other without holders. + Effect.flatMap(({ leases }) => + Effect.forEach( + leases, + (lease) => + durableSessionIds.has(lease.sessionId) + ? Effect.sync(() => { + this.leaseBySession.set(lease.sessionId, lease.worktreePath); + }) + : this.release(lease.sessionId).pipe( + Effect.asVoid, + Effect.catch((error) => + Effect.logWarning('Managed worktree lease sweep deferred', error), + ), + ), + { discard: true }, + ), + ), + Effect.andThen(Effect.suspend(() => this.reconcile())), ); } @@ -81,8 +104,36 @@ export class WorktreeService { }); } + /** Give `sessionId` a hold on the worktree another session already holds — a fork shares its + * source's working tree. Typed `conflict` when the worktree is being removed. */ + acquire(sessionId: SessionId, worktreePath: string): Effect.Effect { + return Effect.tryPromise({ + try: () => this.store.acquireLease(worktreePath, sessionId), + catch: (cause) => + cause instanceof WorktreeUnavailableError + ? new RequestError({ + code: 'conflict', + message: 'The managed worktree is being removed', + }) + : new OperationError({ + subsystem: 'store', + operation: 'worktrees.lease', + publicMessage: 'Failed to lease the managed worktree', + cause, + }), + }).pipe( + Effect.tap(() => + Effect.sync(() => { + this.leaseBySession.set(sessionId, worktreePath); + }), + ), + ); + } + + /** A start in the session's worktree — a resume, or a fork child — needs the directory on disk; + * the lease alone says nothing about that. */ verifyResume(sessionId: SessionId): Effect.Effect { - const record = this.bySession.get(sessionId); + const record = this.get(sessionId); if (!record || existsSync(record.worktreePath)) return Effect.void; return Effect.fail( new RequestError({ @@ -92,22 +143,58 @@ export class WorktreeService { ); } + /** The worktree `sessionId` holds a lease on. */ get(sessionId: SessionId): WorktreeRecord | undefined { - return this.bySession.get(sessionId); + const worktreePath = this.leaseBySession.get(sessionId); + return worktreePath === undefined + ? undefined + : this.byPath.get(normalizeRepoRoot(worktreePath)); + } + + /** The other sessions leasing the worktree `sessionId` holds — the ones whose running turn + * makes this session's next turn `busy`. */ + coLeaseholders(sessionId: SessionId): SessionId[] { + const worktreePath = this.leaseBySession.get(sessionId); + if (worktreePath === undefined) return []; + const others: SessionId[] = []; + for (const [holder, path] of this.leaseBySession) { + if (holder !== sessionId && path === worktreePath) others.push(holder); + } + return others; } hasPath(path: string): boolean { - const key = normalizeRepoRoot(path); - return [...this.bySession.values()].some( - (record) => normalizeRepoRoot(record.worktreePath) === key, - ); + return this.byPath.has(normalizeRepoRoot(path)); } + /** The session is gone: drop its lease, and when it was the last, clean the worktree up. */ cleanupDeletedSession(sessionId: SessionId): Effect.Effect { - const record = this.bySession.get(sessionId); - if (!record) return Effect.void; - return this.semaphore(normalizeRepoRoot(record.repoRoot)).withPermit( - this.cleanupRecord(record), + return this.release(sessionId).pipe( + Effect.flatMap((record) => + record === undefined + ? Effect.void + : this.semaphore(normalizeRepoRoot(record.repoRoot)).withPermit( + this.cleanupRecord(record), + ), + ), + ); + } + + /** Release the session's lease; the worktree comes back only when that lease was the last one + * (the store marked it `deleting`), for the caller to clean up. */ + private release(sessionId: SessionId): Effect.Effect { + return storeEffect('worktrees.release', 'Failed to release the managed worktree', () => + this.store.releaseLease(sessionId), + ).pipe( + Effect.map((released) => { + this.leaseBySession.delete(sessionId); + if (!released?.last) return; + const record = this.byPath.get(normalizeRepoRoot(released.worktreePath)); + if (record === undefined) return; + const deleting: WorktreeRecord = { ...record, state: 'deleting' }; + this.index(deleting); + return deleting; + }), ); } @@ -180,21 +267,24 @@ export class WorktreeService { worktreePath, repoRoot, branch, - sessionId, createdAt: Date.now(), state: 'active', }; + // Two writes, not one transaction: a crash between them leaves an active worktree with no + // lease, which boot reconcile cleans up like any other holder-less worktree. const saved = yield* Effect.exit( storeEffect('worktrees.save', 'Failed to persist managed worktree', () => this.store.save(record), - ), + ).pipe(Effect.andThen(this.acquire(sessionId, worktreePath))), ); if (Exit.isFailure(saved)) { yield* removeWorktreeBestEffort(repoRoot, worktreePath); + yield* storeEffect('worktrees.delete', 'Failed to delete managed worktree', () => + this.store.delete(worktreePath), + ).pipe(Effect.catch(() => Effect.void)); return yield* Effect.failCause(saved.cause); } - this.bySession.set(sessionId, record); - this.byRepoBranch.set(repoBranchKey(repoRoot, branch), record); + this.index(record); yield* this.git.invalidate(options.cwd); yield* this.git.invalidate(worktreePath); return withoutBranch(options, worktreePath); @@ -246,6 +336,8 @@ export class WorktreeService { return semaphore; } + /** Remove a worktree nobody holds: a missing or clean, pushed tree goes with its row; a dirty + * or unpushed one is kept on disk as `orphaned`. */ cleanupRecord(record: WorktreeRecord): Effect.Effect { return Effect.gen({ self: this }, function* () { if (!existsSync(record.worktreePath)) { @@ -277,13 +369,19 @@ export class WorktreeService { }); } - reconcile(sessionIds: ReadonlySet): Effect.Effect { + /** Boot: a worktree with no remaining holder is cleaned up when safe — an `active` one whose + * sessions are gone, a `deleting` one whose cleanup the previous daemon never finished, or an + * `orphaned` one that has become clean and pushed since; a held one whose directory vanished + * is marked orphaned. Leases were swept in `start`. */ + reconcile(): Effect.Effect { return Effect.gen({ self: this }, function* () { - for (const record of this.bySession.values()) { - const hasSession = sessionIds.has(record.sessionId); + const records = Array.from(this.byPath.values()); + for (let i = 0, len = records.length; i < len; i++) { + const record = records[i]; + const held = this.holders(record.worktreePath).length > 0; if (!existsSync(record.worktreePath)) { yield* this.semaphore(normalizeRepoRoot(record.repoRoot)).withPermit( - (hasSession + (held ? this.pruneAdvisory(record.repoRoot).pipe(Effect.andThen(this.markOrphaned(record))) : this.cleanupRecord(record) ).pipe( @@ -292,7 +390,7 @@ export class WorktreeService { ), ), ); - } else if (!hasSession && record.state === 'active') { + } else if (!held) { yield* this.semaphore(normalizeRepoRoot(record.repoRoot)).withPermit( this.cleanupRecord(record).pipe( Effect.catch((error) => @@ -306,6 +404,7 @@ export class WorktreeService { }); } + /** Adopt managed-root directories no row knows: kept as `orphaned`, with no lease. */ scanUnknown(): Effect.Effect { if (!this.root || !existsSync(this.root)) return Effect.void; const root = this.root; @@ -329,7 +428,6 @@ export class WorktreeService { worktreePath: candidate, repoRoot: identity.repoRoot, branch: identity.branch, - sessionId: orphanSessionId(candidate), createdAt: Date.now(), state: 'orphaned', }; @@ -365,16 +463,29 @@ export class WorktreeService { ).pipe( Effect.tap(() => Effect.sync(() => { - this.bySession.delete(record.sessionId); + this.byPath.delete(normalizeRepoRoot(record.worktreePath)); this.byRepoBranch.delete(repoBranchKey(record.repoRoot, record.branch)); + const holders = this.holders(record.worktreePath); + for (let i = 0, len = holders.length; i < len; i++) { + this.leaseBySession.delete(holders[i]); + } }), ), Effect.asVoid, ); } + private holders(worktreePath: string): SessionId[] { + const key = normalizeRepoRoot(worktreePath); + const holders: SessionId[] = []; + for (const [sessionId, path] of this.leaseBySession) { + if (normalizeRepoRoot(path) === key) holders.push(sessionId); + } + return holders; + } + private index(record: WorktreeRecord): void { - this.bySession.set(record.sessionId, record); + this.byPath.set(normalizeRepoRoot(record.worktreePath), record); this.byRepoBranch.set(repoBranchKey(record.repoRoot, record.branch), record); } @@ -409,11 +520,6 @@ function readChildDirectories(path: string): Effect.Effect { ); } -function orphanSessionId(path: string): SessionId { - const digest = createHash('sha256').update(normalizeRepoRoot(path)).digest('hex'); - return SessionIdSchema.parse(`orphan-worktree-${digest}`); -} - function normalizeRepoRoot(path: string): string { const normalized = normalize(resolve(path)); return process.platform === 'win32' ? normalized.toLowerCase() : normalized; diff --git a/packages/host/engine/src/worktree/worktree-store.ts b/packages/host/engine/src/worktree/worktree-store.ts index 328bd1042..c30d8a4bf 100644 --- a/packages/host/engine/src/worktree/worktree-store.ts +++ b/packages/host/engine/src/worktree/worktree-store.ts @@ -1,25 +1,61 @@ -import type { WorktreeRecord } from '@linkcode/schema'; +import type { SessionId, WorktreeLease, WorktreeRecord } from '@linkcode/schema'; +export interface WorktreeStoreSnapshot { + readonly worktrees: WorktreeRecord[]; + readonly leases: WorktreeLease[]; +} + +/** What a released lease was: the worktree it held, and whether it was the last hold on it. */ +export interface WorktreeLeaseRelease { + readonly worktreePath: string; + readonly last: boolean; +} + +/** Rejection from {@link WorktreeStore.acquireLease}: the worktree is gone or `deleting`, so a + * lease on it would name a directory cleanup is about to remove. */ +export class WorktreeUnavailableError extends Error { + constructor(worktreePath: string, options?: ErrorOptions) { + super(`Managed worktree is unavailable: ${worktreePath}`, options); + this.name = 'WorktreeUnavailableError'; + } +} + +/** + * Durable managed-worktree registry: worktree rows and the session leases on them. The daemon + * injects a SQLite implementation on the graph connection — lease release and the `deleting` mark + * MUST be one transaction there; the in-memory default is for tests and embedders. + */ export interface WorktreeStore { - load(): Promise; + load(): Promise; + /** Upsert the worktree row for creation and orphan-marking; leases are untouched. It does not + * police state, so callers never write `active` over a `deleting` row — that only ever holds + * when zero leases remain (the release that set it). */ save(record: WorktreeRecord): Promise; + /** Remove the worktree row and every lease on it. */ delete(worktreePath: string): Promise; + /** Grant `sessionId` its hold; idempotent for the same worktree. Rejects with + * {@link WorktreeUnavailableError} when the worktree is missing or `deleting`. */ + acquireLease(worktreePath: string, sessionId: SessionId): Promise; + /** Drop the session's hold; when it was the last one the worktree is marked `deleting` in the + * same transaction, before any filesystem work. Undefined when the session held none. */ + releaseLease(sessionId: SessionId): Promise; } export class InMemoryWorktreeStore implements WorktreeStore { private readonly records = new Map(); + private readonly leases = new Map(); - load(): Promise { - return Promise.resolve(Array.from(this.records.values(), (record) => structuredClone(record))); + load(): Promise { + return Promise.resolve({ + worktrees: Array.from(this.records.values(), (record) => structuredClone(record)), + leases: Array.from(this.leases.values(), (lease) => structuredClone(lease)), + }); } save(record: WorktreeRecord): Promise { for (const existing of this.records.values()) { if (existing.worktreePath === record.worktreePath) continue; - if ( - existing.sessionId === record.sessionId || - (existing.repoRoot === record.repoRoot && existing.branch === record.branch) - ) { + if (existing.repoRoot === record.repoRoot && existing.branch === record.branch) { return Promise.reject(new Error('worktree already exists')); } } @@ -29,6 +65,38 @@ export class InMemoryWorktreeStore implements WorktreeStore { delete(worktreePath: string): Promise { this.records.delete(worktreePath); + for (const [sessionId, lease] of this.leases) { + if (lease.worktreePath === worktreePath) this.leases.delete(sessionId); + } return Promise.resolve(); } + + acquireLease(worktreePath: string, sessionId: SessionId): Promise { + const record = this.records.get(worktreePath); + if (record === undefined || record.state === 'deleting') { + return Promise.reject(new WorktreeUnavailableError(worktreePath)); + } + const held = this.leases.get(sessionId); + if (held !== undefined && held.worktreePath !== worktreePath) { + return Promise.reject(new Error(`Session ${sessionId} already holds a worktree`)); + } + if (held === undefined) { + this.leases.set(sessionId, { worktreePath, sessionId, createdAt: Date.now() }); + } + return Promise.resolve(); + } + + releaseLease(sessionId: SessionId): Promise { + const lease = this.leases.get(sessionId); + if (lease === undefined) return Promise.resolve(undefined); + this.leases.delete(sessionId); + let remaining = 0; + for (const other of this.leases.values()) { + if (other.worktreePath === lease.worktreePath) remaining += 1; + } + const last = remaining === 0; + const record = this.records.get(lease.worktreePath); + if (last && record !== undefined) record.state = 'deleting'; + return Promise.resolve({ worktreePath: lease.worktreePath, last }); + } } diff --git a/packages/host/engine/tests/integration/engine-worktree.test.ts b/packages/host/engine/tests/integration/engine-worktree.test.ts index bbaed03d0..3fc66d124 100644 --- a/packages/host/engine/tests/integration/engine-worktree.test.ts +++ b/packages/host/engine/tests/integration/engine-worktree.test.ts @@ -2,18 +2,34 @@ import { execFileSync } from 'node:child_process'; import { existsSync, mkdtempSync, rmSync, writeFileSync } from 'node:fs'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; -import type { StartOptions } from '@linkcode/schema'; +import { asHistoryId } from '@linkcode/agent-adapter'; +import type { + AgentHistoryBranchOptions, + AgentHistoryCapabilities, + SessionId, + StartOptions, + TurnId, + WirePayload, +} from '@linkcode/schema'; +import { OperationIdSchema } from '@linkcode/schema'; +import { nullthrow } from 'foxts/guard'; +import { noop } from 'foxts/noop'; import { afterEach, describe, expect, it, vi } from 'vitest'; import { createSessionHarness, FakeAdapter, + settleEngineTasks, startedSessionId, } from '../../src/__tests__/fixtures/session-harness'; +import { InMemoryConversationStore } from '../../src/conversation/conversation-store'; import type { SessionStore } from '../../src/session/session-store'; import { InMemorySessionStore } from '../../src/session/session-store'; import { InMemoryWorkspaceStore } from '../../src/workspace/workspace-store'; import { InMemoryWorktreeStore } from '../../src/worktree/worktree-store'; +const SOURCE_HISTORY = asHistoryId('native-1'); +const CHILD_HISTORY = asHistoryId('native-child'); + const tempRoots: string[] = []; function makeTempDir(): string { @@ -55,6 +71,205 @@ class RejectingStartAdapter extends FakeAdapter { } } +class ForkingAdapter extends FakeAdapter { + override readonly historyCapabilities: AgentHistoryCapabilities = { + list: false, + read: true, + resume: true, + forkAfterTurn: true, + branch: true, + }; + + branchHistory(_opts: AgentHistoryBranchOptions, startOpts: StartOptions): Promise { + this.startedWith = startOpts; + this.emit({ type: 'session-ref', historyId: CHILD_HISTORY }); + return Promise.resolve(); + } +} + +type Harness = ReturnType; + +function worktreeHarness( + makeAdapter: () => FakeAdapter, + stores: { + sessionStore?: InMemorySessionStore; + workspaceStore?: InMemoryWorkspaceStore; + worktreeStore: InMemoryWorktreeStore; + worktreeRoot: string; + conversationStore?: InMemoryConversationStore; + }, +): Harness { + return createSessionHarness( + stores.sessionStore ?? new InMemorySessionStore(), + makeAdapter, + undefined, + undefined, + stores.workspaceStore, + undefined, + { + worktreeStore: stores.worktreeStore, + worktreeRoot: stores.worktreeRoot, + ...(stores.conversationStore && { conversationStore: stores.conversationStore }), + }, + ); +} + +/** Parks the co-leaseholder scan inside the worktree permit, so a delete can land between a + * caller's record check and its durable admit. */ +class GatedConversationStore extends InMemoryConversationStore { + private scan: { sessionId: SessionId; entered: () => void; blocked: Promise } | undefined; + + /** Park the next scan of `sessionId`; resolves once it is parked, returns its release. */ + hold(sessionId: SessionId): { entered: Promise; release: () => void } { + let release = noop; + const blocked = new Promise((resolve) => { + release = resolve; + }); + let entered = noop; + const parked = new Promise((resolve) => { + entered = resolve; + }); + this.scan = { sessionId, entered, blocked }; + return { entered: parked, release }; + } + + override async listOpenOperations(sessionId?: SessionId) { + const scan = this.scan; + if (scan !== undefined && sessionId === scan.sessionId) { + this.scan = undefined; + scan.entered(); + await scan.blocked; + } + return super.listOpenOperations(sessionId); + } +} + +async function startOnWorktree(h: Harness, clientReqId: string, repo: string): Promise { + await h.inject({ + kind: 'session.start', + clientReqId, + opts: { kind: 'claude-code', cwd: repo, branch: { name: 'feature', mode: 'worktree' } }, + }); + return vi.waitFor(() => startedSessionId(h.sent, clientReqId)); +} + +function submitPrompt(h: Harness, clientReqId: string, sessionId: SessionId, text: string) { + return h.inject({ + kind: 'turn.submit', + clientReqId, + sessionId, + operationId: OperationIdSchema.parse(`op-${clientReqId}`), + input: { type: 'prompt', blocks: [{ type: 'text', text }] }, + }); +} + +function submittedTurnId(sent: WirePayload[], replyTo: string): TurnId { + const reply = sent.find( + (payload) => payload.kind === 'turn.submitted' && payload.replyTo === replyTo, + ); + if (reply?.kind !== 'turn.submitted') throw new Error(`no turn.submitted for ${replyTo}`); + return reply.turnId; +} + +/** One settled turn on the source history with a live `ending` checkpoint. */ +async function checkpointedTurn( + h: Harness, + adapter: FakeAdapter, + sessionId: SessionId, + clientReqId: string, +): Promise { + await submitPrompt(h, clientReqId, sessionId, clientReqId); + const turnId = submittedTurnId(h.sent, clientReqId); + adapter.emit({ type: 'session-ref', historyId: SOURCE_HISTORY }); + adapter.emitCheckpoint({ + historyId: SOURCE_HISTORY, + cursor: `cp-${clientReqId}`, + turn: 'ending', + }); + adapter.emit({ type: 'status', status: 'idle' }); + await settleEngineTasks(); + return turnId; +} + +type ForkReply = Extract; + +async function fork( + h: Harness, + clientReqId: string, + sourceSessionId: SessionId, + throughTurnId: TurnId, + expectedGraphRevision: number, +): Promise { + await h.inject({ + kind: 'session.fork', + clientReqId, + sourceSessionId, + throughTurnId, + operationId: OperationIdSchema.parse(`op-${clientReqId}`), + expectedGraphRevision, + }); + return vi.waitFor(() => { + const reply = h.sent.find( + (payload): payload is ForkReply => + (payload.kind === 'session.forked' || payload.kind === 'request.failed') && + payload.replyTo === clientReqId, + ); + return nullthrow(reply, `no fork reply for ${clientReqId}`); + }); +} + +function requestFailed(sent: WirePayload[], replyTo: string) { + return vi.waitFor(() => { + const reply = sent.find( + (payload) => payload.kind === 'request.failed' && payload.replyTo === replyTo, + ); + if (reply?.kind !== 'request.failed') throw new Error(`no request.failed for ${replyTo}`); + return reply; + }); +} + +/** The adapter a fork started, as opposed to the throwaway instances capability lookups mint. */ +function forkedAdapter(h: Harness, source: FakeAdapter): FakeAdapter { + return nullthrow( + h.adapters.find((adapter) => adapter !== source && adapter.startedWith !== null), + 'no forked adapter', + ); +} + +/** A settled turn on a live child running on CHILD_HISTORY, checkpointed so a later turn can fork + * before it. */ +async function childCheckpointedTurn( + h: Harness, + adapter: FakeAdapter, + sessionId: SessionId, + clientReqId: string, +): Promise { + await submitPrompt(h, clientReqId, sessionId, clientReqId); + submittedTurnId(h.sent, clientReqId); + adapter.emitCheckpoint({ historyId: CHILD_HISTORY, cursor: `cp-${clientReqId}`, turn: 'ending' }); + adapter.emit({ type: 'status', status: 'idle' }); + await settleEngineTasks(); +} + +/** The live echo of prompt `text` on `sessionId`, carrying the branch cursor a rewrite hands back. */ +function liveCursor(sent: WirePayload[], sessionId: SessionId, text: string) { + const event = sent + .flatMap((payload) => + payload.kind === 'agent.event' && payload.sessionId === sessionId ? [payload.event] : [], + ) + .findLast( + (candidate) => + candidate.type === 'user-message' && + candidate.branchCursor !== undefined && + candidate.content[0]?.type === 'text' && + candidate.content[0].text === text, + ); + if (event?.type !== 'user-message' || event.branchCursor === undefined) { + throw new Error(`no live prompt echo for ${text}`); + } + return { sourceMessageId: event.messageId, branchCursor: event.branchCursor }; +} + afterEach(() => { const drained = tempRoots.splice(0); for (let i = 0, len = drained.length; i < len; i++) { @@ -90,13 +305,13 @@ describe('engine managed worktree sessions', () => { }, }); const sessionId = await vi.waitFor(() => startedSessionId(h.sent, 'start-delete')); - const [record] = await worktreeStore.load(); + const [record] = (await worktreeStore.load()).worktrees; await h.inject({ kind: 'session.delete', clientReqId: 'delete', sessionId }); await vi.waitFor(() => expect(h.sent).toContainEqual({ kind: 'request.succeeded', replyTo: 'delete' }), ); expect(existsSync(record.worktreePath)).toBe(false); - expect(await worktreeStore.load()).toEqual([]); + expect(await worktreeStore.load()).toEqual({ worktrees: [], leases: [] }); expect((await workspaceStore.load()).some(({ cwd }) => cwd === record.worktreePath)).toBe( false, ); @@ -135,7 +350,7 @@ describe('engine managed worktree sessions', () => { }, }); const sessionId = await vi.waitFor(() => startedSessionId(h.sent, 'start-delete-failure')); - const [record] = await worktreeStore.load(); + const [record] = (await worktreeStore.load()).worktrees; await h.inject({ kind: 'session.delete', clientReqId: 'delete-failure', sessionId }); @@ -148,27 +363,28 @@ describe('engine managed worktree sessions', () => { }), ); expect(existsSync(record.worktreePath)).toBe(true); - expect(await worktreeStore.load()).toMatchObject([{ state: 'active' }]); + expect(await worktreeStore.load()).toMatchObject({ + worktrees: [{ state: 'active' }], + leases: [{ sessionId }], + }); expect(await inner.load()).toHaveLength(1); } finally { await h.engine.stop(); } }); - it('retries orphan cleanup when an already-deleted session is deleted again', async () => { + it('keeps a dirty worktree orphaned and removes it at the next boot once it is clean', async () => { const repo = makeRepo(); const workspaceStore = new InMemoryWorkspaceStore(); const worktreeStore = new InMemoryWorktreeStore(); - const h = createSessionHarness( - new InMemorySessionStore(), - undefined, - undefined, - undefined, + const worktreeRoot = makeTempDir(); + const h = worktreeHarness(() => new FakeAdapter(), { workspaceStore, - undefined, - { worktreeStore, worktreeRoot: makeTempDir() }, - ); + worktreeStore, + worktreeRoot, + }); await h.engine.start(); + let record: { worktreePath: string }; try { await h.inject({ kind: 'session.start', @@ -180,7 +396,7 @@ describe('engine managed worktree sessions', () => { }, }); const sessionId = await vi.waitFor(() => startedSessionId(h.sent, 'start-retry')); - const [record] = await worktreeStore.load(); + [record] = (await worktreeStore.load()).worktrees; const dirtyPath = join(record.worktreePath, 'untracked'); writeFileSync(dirtyPath, 'dirty'); @@ -192,23 +408,29 @@ describe('engine managed worktree sessions', () => { }), ); expect(existsSync(record.worktreePath)).toBe(true); - expect(await worktreeStore.load()).toMatchObject([{ state: 'orphaned' }]); - + expect(await worktreeStore.load()).toMatchObject({ + worktrees: [{ state: 'orphaned' }], + leases: [], + }); rmSync(dirtyPath); - await h.inject({ kind: 'session.delete', clientReqId: 'delete-retry', sessionId }); - await vi.waitFor(() => - expect(h.sent).toContainEqual({ - kind: 'request.succeeded', - replyTo: 'delete-retry', - }), - ); + } finally { + await h.engine.stop(); + } + + const next = worktreeHarness(() => new FakeAdapter(), { + workspaceStore, + worktreeStore, + worktreeRoot, + }); + await next.engine.start(); + try { expect(existsSync(record.worktreePath)).toBe(false); - expect(await worktreeStore.load()).toEqual([]); + expect(await worktreeStore.load()).toEqual({ worktrees: [], leases: [] }); expect((await workspaceStore.load()).some(({ cwd }) => cwd === record.worktreePath)).toBe( false, ); } finally { - await h.engine.stop(); + await next.engine.stop(); } }); @@ -239,9 +461,12 @@ describe('engine managed worktree sessions', () => { }, }); const sessionId = await vi.waitFor(() => startedSessionId(h.sent, 'start')); - const [worktree] = await worktreeStore.load(); + const { + worktrees: [worktree], + leases, + } = await worktreeStore.load(); const [session] = await sessionStore.load(); - expect(worktree.sessionId).toBe(sessionId); + expect(leases).toMatchObject([{ worktreePath: worktree.worktreePath, sessionId }]); expect(session.cwd).toBe(worktree.worktreePath); expect(h.adapters[0].startedWith).toEqual({ kind: 'claude-code', @@ -319,12 +544,303 @@ describe('engine managed worktree sessions', () => { }), ); expect(JSON.stringify(h.sent)).not.toContain('private adapter failure'); - const [worktree] = await worktreeStore.load(); + const { + worktrees: [worktree], + leases, + } = await worktreeStore.load(); const [session] = await sessionStore.load(); - expect(worktree.sessionId).toBe(session.sessionId); + expect(leases).toMatchObject([ + { worktreePath: worktree.worktreePath, sessionId: session.sessionId }, + ]); expect(session.cwd).toBe(worktree.worktreePath); } finally { await h.engine.stop(); } }); }); + +describe('engine managed worktree leases', () => { + it('shares the worktree with a fork child and removes it only after the last lease goes', async () => { + const repo = makeRepo(); + const workspaceStore = new InMemoryWorkspaceStore(); + const worktreeStore = new InMemoryWorktreeStore(); + const h = worktreeHarness(() => new ForkingAdapter(), { + workspaceStore, + worktreeStore, + worktreeRoot: makeTempDir(), + }); + await h.engine.start(); + try { + const sourceId = await startOnWorktree(h, 'start', repo); + const source = nullthrow(h.adapters[0]); + await checkpointedTurn(h, source, sourceId, 't1'); + const secondTurnId = await checkpointedTurn(h, source, sourceId, 't2'); + + const forked = await fork(h, 'fork', sourceId, secondTurnId, 2); + if (forked.kind !== 'session.forked') throw new Error(`fork failed: ${forked.message}`); + const childId = forked.sessionId; + const { + worktrees: [worktree], + leases, + } = await worktreeStore.load(); + expect(leases.map((lease) => lease.sessionId).sort()).toEqual([sourceId, childId].sort()); + expect(new Set(leases.map((lease) => lease.worktreePath))).toEqual( + new Set([worktree.worktreePath]), + ); + expect(forkedAdapter(h, source).startedWith).toMatchObject({ cwd: worktree.worktreePath }); + + await h.inject({ kind: 'session.delete', clientReqId: 'delete-source', sessionId: sourceId }); + await vi.waitFor(() => + expect(h.sent).toContainEqual({ kind: 'request.succeeded', replyTo: 'delete-source' }), + ); + expect(existsSync(worktree.worktreePath)).toBe(true); + expect(await worktreeStore.load()).toMatchObject({ + worktrees: [{ worktreePath: worktree.worktreePath, state: 'active' }], + leases: [{ sessionId: childId }], + }); + expect((await workspaceStore.load()).some(({ cwd }) => cwd === worktree.worktreePath)).toBe( + true, + ); + + await h.inject({ kind: 'session.delete', clientReqId: 'delete-child', sessionId: childId }); + await vi.waitFor(() => + expect(h.sent).toContainEqual({ kind: 'request.succeeded', replyTo: 'delete-child' }), + ); + expect(existsSync(worktree.worktreePath)).toBe(false); + expect(await worktreeStore.load()).toEqual({ worktrees: [], leases: [] }); + expect((await workspaceStore.load()).some(({ cwd }) => cwd === worktree.worktreePath)).toBe( + false, + ); + } finally { + await h.engine.stop(); + } + }); + + it('refuses to fork onto a worktree whose removal has begun and abandons the child', async () => { + const repo = makeRepo(); + const worktreeStore = new InMemoryWorktreeStore(); + const h = worktreeHarness(() => new ForkingAdapter(), { + worktreeStore, + worktreeRoot: makeTempDir(), + }); + await h.engine.start(); + try { + const sourceId = await startOnWorktree(h, 'start', repo); + const source = nullthrow(h.adapters[0]); + const turnId = await checkpointedTurn(h, source, sourceId, 't1'); + // The last lease's release landed durably (another process, or a crash mid-cleanup) while + // this engine still holds the source's view of the worktree. + const [worktree] = (await worktreeStore.load()).worktrees; + await worktreeStore.save({ ...worktree, state: 'deleting' }); + + const forked = await fork(h, 'fork', sourceId, turnId, 1); + expect(forked).toMatchObject({ + kind: 'request.failed', + code: 'conflict', + message: 'The managed worktree is being removed', + }); + expect((await worktreeStore.load()).leases).toMatchObject([{ sessionId: sourceId }]); + expect( + h.sent.filter( + (payload) => payload.kind === 'session.changed' && payload.reason === 'created', + ), + ).toHaveLength(1); + const forkedChild = h.adapters.find( + (adapter) => adapter !== source && adapter.startedWith !== null, + ); + expect(forkedChild).toBeUndefined(); + } finally { + await h.engine.stop(); + } + }); + + it('refuses to fork a session whose managed worktree is missing on disk', async () => { + const repo = makeRepo(); + const worktreeStore = new InMemoryWorktreeStore(); + const h = worktreeHarness(() => new ForkingAdapter(), { + worktreeStore, + worktreeRoot: makeTempDir(), + }); + await h.engine.start(); + try { + const sourceId = await startOnWorktree(h, 'start', repo); + const source = nullthrow(h.adapters[0]); + const turnId = await checkpointedTurn(h, source, sourceId, 't1'); + // The user removed the directory by hand; the lease is still held, so the store would grant + // the child a lease on a path the adapter cannot start in. + const [worktree] = (await worktreeStore.load()).worktrees; + rmSync(worktree.worktreePath, { recursive: true, force: true }); + + const forked = await fork(h, 'fork', sourceId, turnId, 1); + expect(forked).toMatchObject({ + kind: 'request.failed', + code: 'worktree_missing', + message: `The managed worktree is missing at ${worktree.worktreePath}. Restore it or delete this session.`, + }); + expect((await worktreeStore.load()).leases).toMatchObject([{ sessionId: sourceId }]); + const forkedChild = h.adapters.find( + (adapter) => adapter !== source && adapter.startedWith !== null, + ); + expect(forkedChild).toBeUndefined(); + } finally { + await h.engine.stop(); + } + }); + + it('lets only one leaseholder run a turn at a time', async () => { + const repo = makeRepo(); + const worktreeStore = new InMemoryWorktreeStore(); + const h = worktreeHarness(() => new ForkingAdapter(), { + worktreeStore, + worktreeRoot: makeTempDir(), + }); + await h.engine.start(); + try { + const sourceId = await startOnWorktree(h, 'start', repo); + const source = nullthrow(h.adapters[0]); + const turnId = await checkpointedTurn(h, source, sourceId, 't1'); + const forked = await fork(h, 'fork', sourceId, turnId, 1); + if (forked.kind !== 'session.forked') throw new Error(`fork failed: ${forked.message}`); + const childId = forked.sessionId; + + await submitPrompt(h, 'source-turn', sourceId, 'keep going'); + submittedTurnId(h.sent, 'source-turn'); + source.emit({ type: 'status', status: 'running' }); + + await submitPrompt(h, 'child-turn', childId, 'me too'); + expect(await requestFailed(h.sent, 'child-turn')).toMatchObject({ + code: 'busy', + message: 'Another session on this worktree is running a turn', + }); + await h.inject({ + kind: 'agent.input', + clientReqId: 'child-legacy', + sessionId: childId, + input: { type: 'prompt', content: [{ type: 'text', text: 'me too' }] }, + }); + expect(await requestFailed(h.sent, 'child-legacy')).toMatchObject({ code: 'busy' }); + + source.emitCheckpoint({ historyId: SOURCE_HISTORY, cursor: 'cp-3', turn: 'ending' }); + source.emit({ type: 'status', status: 'idle' }); + await settleEngineTasks(); + await submitPrompt(h, 'child-turn-2', childId, 'now'); + await vi.waitFor(() => submittedTurnId(h.sent, 'child-turn-2')); + // The child now holds the worktree's turn: the source is the one refused. + forkedAdapter(h, source).emit({ type: 'status', status: 'running' }); + await submitPrompt(h, 'source-turn-2', sourceId, 'again'); + expect(await requestFailed(h.sent, 'source-turn-2')).toMatchObject({ code: 'busy' }); + } finally { + await h.engine.stop(); + } + }); + + it('persists no turn for a session deleted while its turn waits on the worktree gate', async () => { + const repo = makeRepo(); + const worktreeStore = new InMemoryWorktreeStore(); + const conversationStore = new GatedConversationStore(); + const h = worktreeHarness(() => new ForkingAdapter(), { + worktreeStore, + worktreeRoot: makeTempDir(), + conversationStore, + }); + await h.engine.start(); + try { + const sourceId = await startOnWorktree(h, 'start', repo); + const source = nullthrow(h.adapters[0]); + const turnId = await checkpointedTurn(h, source, sourceId, 't1'); + const forked = await fork(h, 'fork', sourceId, turnId, 1); + if (forked.kind !== 'session.forked') throw new Error(`fork failed: ${forked.message}`); + const childId = forked.sessionId; + + // The child's admission validated its record, then parked on the co-leaseholder scan. + const gate = conversationStore.hold(sourceId); + const submitted = submitPrompt(h, 'child-turn', childId, 'me too'); + await gate.entered; + await h.inject({ kind: 'session.delete', clientReqId: 'delete', sessionId: childId }); + await vi.waitFor(() => + expect(h.sent).toContainEqual({ kind: 'request.succeeded', replyTo: 'delete' }), + ); + gate.release(); + await submitted; + + expect(await requestFailed(h.sent, 'child-turn')).toMatchObject({ + code: 'not_found', + message: `Unknown session: ${childId}`, + }); + expect(await conversationStore.listTurns(childId)).toEqual([]); + expect(await conversationStore.listOpenOperations(childId)).toEqual([]); + } finally { + await h.engine.stop(); + } + }); + + it('refuses to rewrite a prompt on a worktree whose co-leaseholder is running', async () => { + const repo = makeRepo(); + const worktreeStore = new InMemoryWorktreeStore(); + const h = worktreeHarness(() => new ForkingAdapter(), { + worktreeStore, + worktreeRoot: makeTempDir(), + }); + await h.engine.start(); + try { + const sourceId = await startOnWorktree(h, 'start', repo); + const source = nullthrow(h.adapters[0]); + const turnId = await checkpointedTurn(h, source, sourceId, 't1'); + const forked = await fork(h, 'fork', sourceId, turnId, 1); + if (forked.kind !== 'session.forked') throw new Error(`fork failed: ${forked.message}`); + const childId = forked.sessionId; + const child = forkedAdapter(h, source); + + // The child runs two live turns of its own so its second prompt has a bound cut to rewrite. + await childCheckpointedTurn(h, child, childId, 'k1'); + await childCheckpointedTurn(h, child, childId, 'k2'); + const rewriteTarget = liveCursor(h.sent, childId, 'k2'); + + // The source holds the worktree's running turn. + await submitPrompt(h, 'source-run', sourceId, 'keep going'); + submittedTurnId(h.sent, 'source-run'); + source.emit({ type: 'status', status: 'running' }); + + await h.inject({ + kind: 'history.branch', + clientReqId: 'rewrite', + sourceSessionId: childId, + ...rewriteTarget, + content: [{ type: 'text', text: 'k2, edited' }], + }); + expect(await requestFailed(h.sent, 'rewrite')).toMatchObject({ + code: 'busy', + message: 'Another session on this worktree is running a turn', + }); + } finally { + await h.engine.stop(); + } + }); + + it('sweeps a lease whose session is gone at boot and finishes the cleanup', async () => { + const repo = makeRepo(); + const worktreeStore = new InMemoryWorktreeStore(); + const worktreeRoot = makeTempDir(); + const first = worktreeHarness(() => new FakeAdapter(), { worktreeStore, worktreeRoot }); + await first.engine.start(); + let worktreePath: string; + try { + await startOnWorktree(first, 'start', repo); + worktreePath = (await worktreeStore.load()).worktrees[0].worktreePath; + } finally { + await first.engine.stop(); + } + expect(existsSync(worktreePath)).toBe(true); + expect((await worktreeStore.load()).leases).toHaveLength(1); + + // The session store the next boot reads never held that session. + const second = worktreeHarness(() => new FakeAdapter(), { worktreeStore, worktreeRoot }); + await second.engine.start(); + try { + expect(existsSync(worktreePath)).toBe(false); + expect(await worktreeStore.load()).toEqual({ worktrees: [], leases: [] }); + } finally { + await second.engine.stop(); + } + }); +}); diff --git a/packages/host/engine/tests/integration/worktree-service.test.ts b/packages/host/engine/tests/integration/worktree-service.test.ts index 1921b3764..ff23b9f74 100644 --- a/packages/host/engine/tests/integration/worktree-service.test.ts +++ b/packages/host/engine/tests/integration/worktree-service.test.ts @@ -84,7 +84,7 @@ describe('WorktreeService', () => { ), ); expect(result).toEqual({ kind: 'pi', cwd }); - expect(await store.load()).toEqual([]); + expect(await store.load()).toEqual({ worktrees: [], leases: [] }); }); it('rejects worktree mode when the branch is already checked out in the original cwd', async () => { @@ -121,7 +121,7 @@ describe('WorktreeService', () => { ); expect(result).toEqual({ kind: 'pi', cwd }); - expect(await store.load()).toEqual([]); + expect(await store.load()).toEqual({ worktrees: [], leases: [] }); const status = await Effect.runPromise(gitService.getStatus(cwd)); expect(status.isRepo && status.branch).toBe('feature/a'); }); @@ -143,11 +143,11 @@ describe('WorktreeService', () => { expect(result.branch).toBeUndefined(); expect(result.cwd).not.toBe(cwd); expect(existsSync(result.cwd)).toBe(true); - expect((await store.load())[0]).toMatchObject({ - worktreePath: result.cwd, - repoRoot: cwd, - branch: 'feature/a', - state: 'active', + expect(await store.load()).toMatchObject({ + worktrees: [ + { worktreePath: result.cwd, repoRoot: cwd, branch: 'feature/a', state: 'active' }, + ], + leases: [{ worktreePath: result.cwd, sessionId: 'sess-feature' }], }); const status = await Effect.runPromise(gitService.getStatus(result.cwd)); expect(status.isRepo && status.branch).toBe('feature/a'); @@ -221,7 +221,7 @@ describe('WorktreeService', () => { await Effect.runPromise(service.cleanupDeletedSession(id)); expect(existsSync(worktree.cwd)).toBe(false); - expect(await store.load()).toEqual([]); + expect(await store.load()).toEqual({ worktrees: [], leases: [] }); }); it.each([ @@ -250,10 +250,10 @@ describe('WorktreeService', () => { await Effect.runPromise(service.cleanupDeletedSession(id)); expect(existsSync(worktree.cwd)).toBe(true); - expect(await store.load()).toMatchObject([{ state: 'orphaned' }]); + expect(await store.load()).toMatchObject({ worktrees: [{ state: 'orphaned' }], leases: [] }); }); - it('does not update the in-memory ownership state when orphan persistence fails', async () => { + it('leaves the worktree deleting for the next boot when orphan persistence fails', async () => { const inner = new InMemoryWorktreeStore(); let rejectSaves = false; const store: WorktreeStore = { @@ -261,6 +261,8 @@ describe('WorktreeService', () => { save: (record) => rejectSaves ? Promise.reject(new Error('disk unavailable')) : inner.save(record), delete: (path) => inner.delete(path), + acquireLease: (path, sessionId) => inner.acquireLease(path, sessionId), + releaseLease: (sessionId) => inner.releaseLease(sessionId), }; const { service, id, worktree } = await managedWorktree(store); writeFileSync(join(worktree.cwd, 'untracked'), 'dirty'); @@ -269,15 +271,16 @@ describe('WorktreeService', () => { const result = await Effect.runPromiseExit(service.cleanupDeletedSession(id)); expect(result._tag).toBe('Failure'); - expect(service.get(id)?.state).toBe('active'); - expect(await inner.load()).toMatchObject([{ state: 'active' }]); + // The lease release landed durably before the failed inspection: the session holds nothing. + expect(service.get(id)).toBeUndefined(); + expect(await inner.load()).toMatchObject({ worktrees: [{ state: 'deleting' }], leases: [] }); expect(existsSync(worktree.cwd)).toBe(true); }); - it('keeps active ownership when non-force removal fails', async () => { + it('leaves the worktree deleting when non-force removal fails', async () => { const store = new InMemoryWorktreeStore(); const { root, id, worktree } = await managedWorktree(store); - const [record] = await store.load(); + const [record] = (await store.load()).worktrees; await store.save({ ...record, repoRoot: join(temp(), 'missing') }); const restarted = new WorktreeService( store, @@ -290,7 +293,7 @@ describe('WorktreeService', () => { expect(result._tag).toBe('Failure'); expect(existsSync(worktree.cwd)).toBe(true); - expect(await store.load()).toMatchObject([{ state: 'active' }]); + expect(await store.load()).toMatchObject({ worktrees: [{ state: 'deleting' }], leases: [] }); }); it('reconciles active rows without sessions and missing rows on boot', async () => { @@ -303,7 +306,7 @@ describe('WorktreeService', () => { ); await Effect.runPromise(safeRestart.start()); expect(existsSync(safe.worktree.cwd)).toBe(false); - expect(await safeStore.load()).toEqual([]); + expect(await safeStore.load()).toEqual({ worktrees: [], leases: [] }); const missingStore = new InMemoryWorktreeStore(); const missing = await managedWorktree(missingStore); @@ -314,7 +317,7 @@ describe('WorktreeService', () => { await Effect.runPromise(GitService.make([])), ); await Effect.runPromise(missingRestart.start()); - expect(await missingStore.load()).toEqual([]); + expect(await missingStore.load()).toEqual({ worktrees: [], leases: [] }); }); it('retains a missing ownership row when its durable session still exists', async () => { @@ -329,13 +332,16 @@ describe('WorktreeService', () => { await Effect.runPromise(restarted.start(new Set([id]))); - expect(await store.load()).toMatchObject([{ state: 'orphaned' }]); + expect(await store.load()).toMatchObject({ + worktrees: [{ state: 'orphaned' }], + leases: [{ sessionId: id }], + }); await expect(Effect.runPromise(restarted.verifyResume(id))).rejects.toMatchObject({ code: 'worktree_missing', }); }); - it('records unknown linked worktrees as orphaned but ignores ordinary directories', async () => { + it('adopts unknown linked worktrees as orphaned, ignores ordinary directories, and removes an orphan once it is safe', async () => { const cwd = repo(); const root = temp(); const candidate = join(root, 'repo-group', 'feature-leaf'); @@ -352,16 +358,23 @@ describe('WorktreeService', () => { const store = new InMemoryWorktreeStore(); const service = new WorktreeService(store, root, await Effect.runPromise(GitService.make([]))); await Effect.runPromise(service.start()); - await Effect.runPromise(service.start()); expect(existsSync(candidate)).toBe(true); - expect(await store.load()).toMatchObject([ - { - worktreePath: candidate, - repoRoot: cwd, - branch: 'feature/a', - state: 'orphaned', - }, - ]); + expect(await store.load()).toEqual({ + worktrees: [ + expect.objectContaining({ + worktreePath: candidate, + repoRoot: cwd, + branch: 'feature/a', + state: 'orphaned', + }), + ], + leases: [], + }); + + // A holder-less orphan is re-inspected at every boot; clean and pushed, this one goes. + await Effect.runPromise(service.start()); + expect(existsSync(candidate)).toBe(false); + expect(await store.load()).toEqual({ worktrees: [], leases: [] }); }); });