Skip to content

Commit 3544b18

Browse files
committed
fix: extend TUI startup budget for large local databases
1 parent 81516d7 commit 3544b18

7 files changed

Lines changed: 297 additions & 38 deletions

File tree

‎apps/tui/src/gateway/autostart.ts‎

Lines changed: 70 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,13 @@
1-
import { execFile, spawn, type ChildProcess } from "node:child_process";
2-
import { closeSync, existsSync, mkdirSync, openSync, realpathSync } from "node:fs";
1+
import { execFile, spawn, spawnSync, type ChildProcess } from "node:child_process";
2+
import {
3+
closeSync,
4+
existsSync,
5+
mkdirSync,
6+
openSync,
7+
readSync,
8+
realpathSync,
9+
statSync,
10+
} from "node:fs";
311
import { dirname, join, resolve } from "node:path";
412
import { promisify } from "node:util";
513
import { fileURLToPath } from "node:url";
@@ -19,7 +27,7 @@ import {
1927
type StartupStep = "checking";
2028

2129
const HEALTH_POLL_INTERVAL_MS = 500;
22-
const GATEWAY_START_TIMEOUT_MS = 20_000;
30+
const GATEWAY_START_TIMEOUT_MS = 120_000;
2331

2432
interface GatewayIdentity {
2533
root: string;
@@ -267,6 +275,8 @@ async function launchGatewayProcess(request: GatewayLaunchRequest): Promise<stri
267275
if (!executable) throw new Error(t("gatewayMissingBinary"));
268276
const targetUrl = stripTrailingSlash(request.targetUrl);
269277
let exited: { code: number | null; signal: NodeJS.Signals | null } | undefined;
278+
const stderrPath = gatewayAutostartLogPath(request.instanceHome, "stderr");
279+
const stderrOffset = fileSize(stderrPath);
270280
const stdio = gatewayLogStdio(request.instanceHome);
271281
const child = spawn(executable, [], {
272282
detached: true,
@@ -279,10 +289,16 @@ async function launchGatewayProcess(request: GatewayLaunchRequest): Promise<stri
279289
exited = { code, signal };
280290
});
281291
child.unref();
282-
const deadline = Date.now() + GATEWAY_START_TIMEOUT_MS;
292+
const deadline = Date.now() + gatewayStartTimeoutMs;
283293
while (Date.now() < deadline) {
284294
if (exited) {
285-
throw new Error(`Gateway exited before becoming healthy (${exitDescription(exited)}).`);
295+
throw new Error(
296+
gatewayLaunchError(
297+
`Gateway exited before becoming healthy (${exitDescription(exited)}).`,
298+
stderrPath,
299+
stderrOffset,
300+
),
301+
);
286302
}
287303
const activeUrl = readActiveGatewayUrl(request.instanceHome);
288304
const candidateUrl = activeUrl ? stripTrailingSlash(activeUrl) : targetUrl;
@@ -296,7 +312,7 @@ async function launchGatewayProcess(request: GatewayLaunchRequest): Promise<stri
296312
await delay(HEALTH_POLL_INTERVAL_MS);
297313
}
298314
stopUnreadyChild(child);
299-
throw new Error(t("gatewayStartTimeout"));
315+
throw new Error(gatewayLaunchError(t("gatewayStartTimeout"), stderrPath, stderrOffset));
300316
}
301317

302318
function gatewayProcessEnv(request: GatewayLaunchRequest): NodeJS.ProcessEnv {
@@ -340,19 +356,64 @@ function exitDescription(exit: { code: number | null; signal: NodeJS.Signals | n
340356

341357
function stopUnreadyChild(child: ChildProcess): void {
342358
try {
343-
if (child.exitCode === null && child.signalCode === null) child.kill();
359+
if (child.exitCode !== null || child.signalCode !== null) return;
360+
if (process.platform === "win32" && child.pid) {
361+
const result = spawnSync("taskkill.exe", ["/PID", String(child.pid), "/T", "/F"], {
362+
stdio: "ignore",
363+
windowsHide: true,
364+
});
365+
if (!result.error && result.status === 0) return;
366+
}
367+
child.kill();
344368
} catch {
345369
// Best effort: the gateway may have already detached or exited.
346370
}
347371
}
348372

373+
function gatewayAutostartLogPath(instanceHome: string, stream: "stdout" | "stderr"): string {
374+
return join(instanceHome, ".tura", "logs", `gateway-autostart.${stream}.log`);
375+
}
376+
377+
function fileSize(path: string): number {
378+
try {
379+
return statSync(path).size;
380+
} catch {
381+
return 0;
382+
}
383+
}
384+
385+
function gatewayLaunchError(message: string, stderrPath: string, offset: number): string {
386+
try {
387+
const size = fileSize(stderrPath);
388+
if (size <= offset) return message;
389+
const start = Math.max(offset, size - 8_192);
390+
const buffer = Buffer.alloc(size - start);
391+
const file = openSync(stderrPath, "r");
392+
try {
393+
readSync(file, buffer, 0, buffer.length, start);
394+
} finally {
395+
closeSync(file);
396+
}
397+
const detail = buffer
398+
.toString("utf8")
399+
.split(/\r?\n/u)
400+
.map((line) => line.trim())
401+
.filter(Boolean)
402+
.slice(-8)
403+
.join(" | ");
404+
return detail ? `${message} Gateway stderr: ${detail}` : message;
405+
} catch {
406+
return message;
407+
}
408+
}
409+
349410
function gatewayLogStdio(instanceHome: string): ["ignore", number, number] {
350411
const logDir = join(instanceHome, ".tura", "logs");
351412
mkdirSync(logDir, { recursive: true });
352413
return [
353414
"ignore",
354-
openSync(join(logDir, "gateway-autostart.stdout.log"), "a"),
355-
openSync(join(logDir, "gateway-autostart.stderr.log"), "a"),
415+
openSync(gatewayAutostartLogPath(instanceHome, "stdout"), "a"),
416+
openSync(gatewayAutostartLogPath(instanceHome, "stderr"), "a"),
356417
];
357418
}
358419

‎crates/gateway/src/router_process.rs‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,10 @@ use std::{
2121
time::{Duration, Instant},
2222
};
2323

24-
const ROUTER_HEALTH_TIMEOUT: Duration = Duration::from_secs(20);
24+
// This must exceed the session_db startup budget (30 seconds). Gateway may
25+
// perform two router attempts, while the TUI allows 120 seconds for the whole
26+
// process tree to become healthy.
27+
const ROUTER_HEALTH_TIMEOUT: Duration = Duration::from_secs(45);
2528
const DEFAULT_ROUTER_EXECUTION_TIMEOUT: Duration = Duration::from_secs(35 * 60);
2629
const ROUTER_PROBE_CONNECT_TIMEOUT: Duration = Duration::from_millis(100);
2730
const ROUTER_STARTUP_POLL_INTERVAL: Duration = Duration::from_millis(200);

‎crates/gateway/src/session_feed.rs‎

Lines changed: 20 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -80,7 +80,20 @@ pub fn start_session_feed_tailer() -> Result<(SessionFeedTailer, tokio::sync::on
8080
match subscription.poll_next_entry(SUBSCRIPTION_POLL_INTERVAL) {
8181
Ok(Poll::Ready(Some(entry))) => {
8282
if let Err(error) = reducer.apply(entry) {
83-
break Err(error.context("failed to reduce durable Session feed"));
83+
if !is_session_feed_cursor_gap(&error) {
84+
break Err(error.context("failed to reduce durable Session feed"));
85+
}
86+
match reconnect_session_feed(
87+
&mut reducer,
88+
&thread_stopping,
89+
&thread_cancellation,
90+
) {
91+
Ok(Some(reconnected)) => subscription = reconnected,
92+
Ok(None) => break Ok(()),
93+
Err(error) => break Err(
94+
error.context("failed to resynchronize durable Session feed"),
95+
),
96+
}
8497
}
8598
}
8699
Ok(Poll::Ready(None)) if thread_stopping.load(Ordering::SeqCst) => break Ok(()),
@@ -187,6 +200,12 @@ fn reconnect_session_feed(
187200
}
188201
}
189202

203+
fn is_session_feed_cursor_gap(error: &anyhow::Error) -> bool {
204+
error
205+
.chain()
206+
.any(|cause| cause.to_string().starts_with("session feed cursor gap for "))
207+
}
208+
190209
fn replay_all_sessions(client: &SessionDbClient, reducer: &mut SessionFeedReducer) -> Result<()> {
191210
let mut canonical_session_ids = HashSet::new();
192211
for workspace in client.list_workspaces()? {

‎crates/router/src/services/session_db.rs‎

Lines changed: 116 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,8 @@
66
use anyhow::{anyhow, Result};
77
use serde_json::json;
88
use session_log_contract::client::{
9-
call_service, service_addr_path, service_is_running, unreachable_owner_lock_message,
9+
call_service, service_addr_path, service_is_running, session_db_owner_record_path,
10+
unreachable_owner_lock_message,
1011
};
1112
use std::{
1213
path::{Path, PathBuf},
@@ -57,7 +58,9 @@ impl SessionDbService {
5758
return Ok(self.status_payload("running"));
5859
}
5960
if let Some(message) = unreachable_owner_lock_message() {
60-
return Err(anyhow!(message));
61+
if !terminate_orphaned_session_db_owner()? {
62+
return Err(anyhow!(message));
63+
}
6164
}
6265
let service_bin = session_db_binary()
6366
.ok_or_else(|| anyhow!("session_db service executable tura_session_db not found"))?;
@@ -66,6 +69,13 @@ impl SessionDbService {
6669
let mut command = Command::new(&service_bin);
6770
command
6871
.env("TURA_HOME", tura_path::instance_home())
72+
.env("TURA_ROUTER_PARENT_PID", std::process::id().to_string())
73+
.env(
74+
"TURA_ROUTER_PARENT_START_TIME",
75+
current_process_start_time(std::process::id())
76+
.unwrap_or_default()
77+
.to_string(),
78+
)
6979
.stdin(Stdio::null())
7080
.stdout(Stdio::null())
7181
.stderr(Stdio::null());
@@ -200,6 +210,110 @@ impl SessionDbService {
200210
}
201211
}
202212

213+
#[derive(Default)]
214+
struct SessionDbOwnerRecord {
215+
pid: Option<u32>,
216+
process_start_time: Option<u64>,
217+
kind: Option<String>,
218+
build_kind: Option<String>,
219+
home: Option<String>,
220+
parent_pid: Option<u32>,
221+
parent_process_start_time: Option<u64>,
222+
}
223+
224+
fn terminate_orphaned_session_db_owner() -> Result<bool> {
225+
let raw = match std::fs::read_to_string(session_db_owner_record_path()) {
226+
Ok(raw) => raw,
227+
Err(_) => return Ok(false),
228+
};
229+
let mut record = SessionDbOwnerRecord::default();
230+
for line in raw.lines() {
231+
let Some((key, value)) = line.split_once('=') else {
232+
continue;
233+
};
234+
let value = value.trim();
235+
match key.trim() {
236+
"pid" => record.pid = value.parse().ok(),
237+
"process_start_time" => record.process_start_time = value.parse().ok(),
238+
"kind" => record.kind = Some(value.to_string()),
239+
"build_kind" => record.build_kind = Some(value.to_string()),
240+
"home" => record.home = Some(value.to_string()),
241+
"parent_pid" => record.parent_pid = value.parse().ok(),
242+
"parent_process_start_time" => {
243+
record.parent_process_start_time = value.parse().ok()
244+
}
245+
_ => {}
246+
}
247+
}
248+
if record.kind.as_deref() != Some("session_db")
249+
|| record.build_kind.as_deref() != Some(tura_path::build_kind())
250+
|| !record
251+
.home
252+
.as_deref()
253+
.is_some_and(|home| same_path(home, &tura_path::instance_home()))
254+
{
255+
return Ok(false);
256+
}
257+
let (Some(pid), Some(start_time), Some(parent_pid), Some(parent_start_time)) = (
258+
record.pid,
259+
record.process_start_time,
260+
record.parent_pid,
261+
record.parent_process_start_time,
262+
) else {
263+
return Ok(false);
264+
};
265+
let mut system = sysinfo::System::new_all();
266+
system.refresh_processes();
267+
if system
268+
.process(sysinfo::Pid::from_u32(parent_pid))
269+
.is_some_and(|parent| parent.start_time() == parent_start_time)
270+
{
271+
return Ok(false);
272+
}
273+
let Some(process) = system.process(sysinfo::Pid::from_u32(pid)) else {
274+
return Ok(false);
275+
};
276+
if process.start_time() != start_time
277+
|| !process
278+
.name()
279+
.trim_end_matches(".exe")
280+
.eq_ignore_ascii_case("tura_session_db")
281+
|| service_is_running()
282+
{
283+
return Ok(false);
284+
}
285+
if !process.kill() {
286+
return Ok(false);
287+
}
288+
let started = Instant::now();
289+
while started.elapsed() < Duration::from_secs(10) {
290+
system.refresh_processes();
291+
if system.process(sysinfo::Pid::from_u32(pid)).is_none() {
292+
return Ok(true);
293+
}
294+
std::thread::sleep(Duration::from_millis(100));
295+
}
296+
Ok(false)
297+
}
298+
299+
fn current_process_start_time(pid: u32) -> Option<u64> {
300+
let mut system = sysinfo::System::new_all();
301+
system.refresh_processes();
302+
system
303+
.process(sysinfo::Pid::from_u32(pid))
304+
.map(sysinfo::Process::start_time)
305+
}
306+
307+
fn same_path(left: &str, right: &Path) -> bool {
308+
let left = tura_path::normalize_path(Path::new(left));
309+
let right = tura_path::normalize_path(right);
310+
if cfg!(windows) {
311+
left.to_string_lossy().to_lowercase() == right.to_string_lossy().to_lowercase()
312+
} else {
313+
left == right
314+
}
315+
}
316+
203317
fn wait_for_child_exit(child: &mut Child, timeout: Duration) -> bool {
204318
let started = Instant::now();
205319
while started.elapsed() < timeout {

‎crates/session_log/Cargo.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -51,11 +51,11 @@ lifecycle.workspace = true
5151
rusqlite.workspace = true
5252
serde.workspace = true
5353
serde_json.workspace = true
54+
sysinfo.workspace = true
5455
tracing.workspace = true
5556
tura_path.workspace = true
5657
session_log_contract.workspace = true
5758
uuid.workspace = true
5859

5960
[dev-dependencies]
60-
sysinfo.workspace = true
6161
tempfile.workspace = true

0 commit comments

Comments
 (0)