Skip to content

Commit 4a12097

Browse files
committed
Merge hotfix for large local database startup
2 parents 7c50aef + fa47bfc commit 4a12097

11 files changed

Lines changed: 332 additions & 61 deletions

File tree

‎.github/workflows/ci.yml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ on:
1010
branches:
1111
- main
1212
- "fix/**"
13+
- "hotfix/**"
1314
- "release/**"
1415
- "staging/runtime_session_refactor"
1516
tags:

‎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;
@@ -271,6 +279,8 @@ async function launchGatewayProcess(request: GatewayLaunchRequest): Promise<stri
271279
if (!executable) throw new Error(t("gatewayMissingBinary"));
272280
const targetUrl = stripTrailingSlash(request.targetUrl);
273281
let exited: { code: number | null; signal: NodeJS.Signals | null } | undefined;
282+
const stderrPath = gatewayAutostartLogPath(request.instanceHome, "stderr");
283+
const stderrOffset = fileSize(stderrPath);
274284
const stdio = gatewayLogStdio(request.instanceHome);
275285
const child = spawn(executable, [], {
276286
detached: true,
@@ -283,10 +293,16 @@ async function launchGatewayProcess(request: GatewayLaunchRequest): Promise<stri
283293
exited = { code, signal };
284294
});
285295
child.unref();
286-
const deadline = Date.now() + GATEWAY_START_TIMEOUT_MS;
296+
const deadline = Date.now() + gatewayStartTimeoutMs;
287297
while (Date.now() < deadline) {
288298
if (exited) {
289-
throw new Error(`Gateway exited before becoming healthy (${exitDescription(exited)}).`);
299+
throw new Error(
300+
gatewayLaunchError(
301+
`Gateway exited before becoming healthy (${exitDescription(exited)}).`,
302+
stderrPath,
303+
stderrOffset,
304+
),
305+
);
290306
}
291307
const activeUrl = readActiveGatewayUrl(request.instanceHome);
292308
const candidateUrl = activeUrl ? stripTrailingSlash(activeUrl) : targetUrl;
@@ -300,7 +316,7 @@ async function launchGatewayProcess(request: GatewayLaunchRequest): Promise<stri
300316
await delay(HEALTH_POLL_INTERVAL_MS);
301317
}
302318
stopUnreadyChild(child);
303-
throw new Error(t("gatewayStartTimeout"));
319+
throw new Error(gatewayLaunchError(t("gatewayStartTimeout"), stderrPath, stderrOffset));
304320
}
305321

306322
function gatewayProcessEnv(request: GatewayLaunchRequest): NodeJS.ProcessEnv {
@@ -344,19 +360,64 @@ function exitDescription(exit: { code: number | null; signal: NodeJS.Signals | n
344360

345361
function stopUnreadyChild(child: ChildProcess): void {
346362
try {
347-
if (child.exitCode === null && child.signalCode === null) child.kill();
363+
if (child.exitCode !== null || child.signalCode !== null) return;
364+
if (process.platform === "win32" && child.pid) {
365+
const result = spawnSync("taskkill.exe", ["/PID", String(child.pid), "/T", "/F"], {
366+
stdio: "ignore",
367+
windowsHide: true,
368+
});
369+
if (!result.error && result.status === 0) return;
370+
}
371+
child.kill();
348372
} catch {
349373
// Best effort: the gateway may have already detached or exited.
350374
}
351375
}
352376

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

‎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: 23 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -80,7 +80,21 @@ 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) => {
94+
break Err(error
95+
.context("failed to resynchronize durable Session feed"));
96+
}
97+
}
8498
}
8599
}
86100
Ok(Poll::Ready(None)) if thread_stopping.load(Ordering::SeqCst) => break Ok(()),
@@ -187,6 +201,14 @@ fn reconnect_session_feed(
187201
}
188202
}
189203

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

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

Lines changed: 114 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},
@@ -56,7 +57,9 @@ impl SessionDbService {
5657
if service_is_running() {
5758
return Ok(self.status_payload("running"));
5859
}
59-
if let Some(message) = unreachable_owner_lock_message() {
60+
if let Some(message) = unreachable_owner_lock_message()
61+
&& !terminate_orphaned_session_db_owner()?
62+
{
6063
return Err(anyhow!(message));
6164
}
6265
let service_bin = session_db_binary()
@@ -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,108 @@ 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" => record.parent_process_start_time = value.parse().ok(),
243+
_ => {}
244+
}
245+
}
246+
if record.kind.as_deref() != Some("session_db")
247+
|| record.build_kind.as_deref() != Some(tura_path::build_kind())
248+
|| !record
249+
.home
250+
.as_deref()
251+
.is_some_and(|home| same_path(home, &tura_path::instance_home()))
252+
{
253+
return Ok(false);
254+
}
255+
let (Some(pid), Some(start_time), Some(parent_pid), Some(parent_start_time)) = (
256+
record.pid,
257+
record.process_start_time,
258+
record.parent_pid,
259+
record.parent_process_start_time,
260+
) else {
261+
return Ok(false);
262+
};
263+
let mut system = sysinfo::System::new_all();
264+
system.refresh_processes();
265+
if system
266+
.process(sysinfo::Pid::from_u32(parent_pid))
267+
.is_some_and(|parent| parent.start_time() == parent_start_time)
268+
{
269+
return Ok(false);
270+
}
271+
let Some(process) = system.process(sysinfo::Pid::from_u32(pid)) else {
272+
return Ok(false);
273+
};
274+
if process.start_time() != start_time
275+
|| !process
276+
.name()
277+
.trim_end_matches(".exe")
278+
.eq_ignore_ascii_case("tura_session_db")
279+
|| service_is_running()
280+
{
281+
return Ok(false);
282+
}
283+
if !process.kill() {
284+
return Ok(false);
285+
}
286+
let started = Instant::now();
287+
while started.elapsed() < Duration::from_secs(10) {
288+
system.refresh_processes();
289+
if system.process(sysinfo::Pid::from_u32(pid)).is_none() {
290+
return Ok(true);
291+
}
292+
std::thread::sleep(Duration::from_millis(100));
293+
}
294+
Ok(false)
295+
}
296+
297+
fn current_process_start_time(pid: u32) -> Option<u64> {
298+
let mut system = sysinfo::System::new_all();
299+
system.refresh_processes();
300+
system
301+
.process(sysinfo::Pid::from_u32(pid))
302+
.map(sysinfo::Process::start_time)
303+
}
304+
305+
fn same_path(left: &str, right: &Path) -> bool {
306+
let left = tura_path::normalize_path(Path::new(left));
307+
let right = tura_path::normalize_path(right);
308+
if cfg!(windows) {
309+
left.to_string_lossy().to_lowercase() == right.to_string_lossy().to_lowercase()
310+
} else {
311+
left == right
312+
}
313+
}
314+
203315
fn wait_for_child_exit(child: &mut Child, timeout: Duration) -> bool {
204316
let started = Instant::now();
205317
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)