Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 6 additions & 2 deletions src/state/workflow-state.js
Original file line number Diff line number Diff line change
Expand Up @@ -166,13 +166,17 @@ export async function saveWorkflowSnapshot(
};
let writeErr;
let guarded = false;
// SQLite persistence is part of the chain so drainWriteChain() covers it.
// Otherwise fire-and-forget callers (e.g. actor.subscribe) leave SQLite tmp
// files in flight after drain returns, breaking workspace teardown.
const chain = getWriteChain(workspaceDir)
.then(async () => {
if (await readGuardRunId(statePath, guardRunId)) {
guarded = true;
return;
}
await writeJson(statePath, payload);
await persistSnapshotToSqlite(sqlitePath, payload);
})
.catch((e) => {
writeErr = e;
Expand All @@ -181,7 +185,6 @@ export async function saveWorkflowSnapshot(
await chain;
if (guarded) return null;
if (writeErr) throw writeErr;
await persistSnapshotToSqlite(sqlitePath, payload);
return payload;
}

Expand All @@ -208,13 +211,15 @@ export async function saveWorkflowTerminalState(
};
let writeErr;
let guarded = false;
// SQLite persistence is part of the chain so drainWriteChain() covers it.
const chain = getWriteChain(workspaceDir)
.then(async () => {
if (await readGuardRunId(statePath, guardRunId)) {
guarded = true;
return;
}
await writeJson(statePath, payload);
await persistSnapshotToSqlite(sqlitePath, payload);
})
.catch((e) => {
writeErr = e;
Expand All @@ -223,7 +228,6 @@ export async function saveWorkflowTerminalState(
await chain;
if (guarded) return null;
if (writeErr) throw writeErr;
await persistSnapshotToSqlite(sqlitePath, payload);
return payload;
}

Expand Down
Loading