Skip to content
Merged
Show file tree
Hide file tree
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
6 changes: 6 additions & 0 deletions docs/guides/MODULE_MAP.md
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,11 @@ Wires all modules, sets up session manager and SDK bridge, dispatches messages.
| `project-pair-lifecycle.js` | Worker management internals: bounded `partner_status`, replacement orchestration, and the per-Driver generation/evaluation ledger. Replacement preflights before mutation and records transactional stage/failure/rollback state through `project-pair-replacement-state.js` |
| `project-pair-usage.js` | Honest Split Worker usage projection. Separates current adapter snapshots, cumulative observed result usage, last-task usage, and compaction observations; unavailable or incomplete evidence stays explicit |
| `project-pair-task-control.js` | Correlated Split Worker task lifecycle: explicit queued follow-ups, exact inspect/cancel, interrupt-and-replace, resume gates, generation/task ids, completion events, and bounded Worker-reported outcome envelopes |
| `project-pair-result-outbox.js` | Durable per-project capture foundation for immutable Worker completion outcomes keyed by owner and stable Driver/Worker origins, task, and generation; atomic persistence and bounded capture retries only, with sender acceptance and recovery deliberately separate |
| `project-pair-result-capture.js` | Bounded completion-capture orchestration: retains the live delegation token across persistence failure, retries capture, and finalizes the existing pair path only after capture succeeds |
| `project-pair-result-delivery.js` | Durable sender acceptance for captured Worker results: persists attempt transitions before SDK delivery, distinguishes transport acceptance from uncertainty, defers query startup, and bounds retries without Worker replay |
| `project-pair-result-pipeline.js` | Wires capture and sender services while keeping coordinator dependencies explicit and bounded |
| `project-pair-result-recovery.js` | Reconciles persisted completed results against stable origins, owner, pair roles, and Worker generation; coalesces event-driven delivery wakes, finalizes accepted results without replay, and serves authorized recovery state/actions |
| `project-pair-autonomous-stop.js` | Run-scoped Split Worker cancellation for Until complete; preserves unrelated queued and active pair work |
| `project-pair-owned-stop.js` | Exact-pair cancellation for a Scheduled Tasks run: establishes the Driver stop barrier, cancels queued Worker work and permissions, and closes/aborts only that pair's live queries |
| `project-pair-message.js` | Non-interrupting live Worker message input. Reports runtime queue acceptance with replay-safe request ids and never claims that the model consumed the message |
Expand Down Expand Up @@ -225,6 +230,7 @@ Bootstraps UI, initializes store, wires remaining Tier 3 modules. All business l
|--------|---------|
| `app-connection.js` + `websocket-lifecycle.js` | WebSocket creation, epoch-guarded handshake/heartbeat/probes, lifecycle-owned bounded jittered reconnect/auth timers, connection status UI, disconnect/restore notifications |
| `app-messages.js` | WebSocket message router (`processMessage`). Dispatches all incoming message types to appropriate handlers |
| `pair-result-status.js` | Per-session pending/blocked/uncertain Split Worker result notice, retained-result view, and explicit duplicate-risk retry confirmation; state is cleared on session switch and never stored locally |
| `message-delivery.js` + `message-delivery-ui.js` | Bounded client message acknowledgement/retry state and context-scoped receipt recovery notices |
| `app-dm.js` | DM mode (open/enter/exit), mate project switching, mate onboarding, DM message rendering, typing indicators |
| `app-home-hub.js` | Home hub rendering, weather, tip rotation, upcoming schedules, project summary |
Expand Down
2 changes: 2 additions & 0 deletions lib/project-connection.js
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,7 @@ function attachConnection(ctx) {
var hydrateSkillCatalog = ctx.hydrateSkillCatalog || function () { return []; };
var warmup = ctx.warmup;
var scheduledMessages = ctx.scheduledMessages;
var pairResultRecovery = ctx.pairResultRecovery;
var sendCursorSharingState = ctx.sendCursorSharingState;

// Adapters are initialized lazily: the first websocket connection into
Expand Down Expand Up @@ -301,6 +302,7 @@ function attachConnection(ctx) {
}
sendTo(ws, { type: "history_done", lastUsage: _lastUsage, lastModelUsage: _lastModelUsage, lastCost: _lastCost, lastStreamInputTokens: _lastStreamInputTokens, contextUsage: active.lastContextUsage || null });
if (scheduledMessages) scheduledMessages.sendState(active, ws);
if (pairResultRecovery) pairResultRecovery.sendState(active, ws);

if (active.isProcessing) {
sendTo(ws, { type: "status", status: "processing" });
Expand Down
86 changes: 86 additions & 0 deletions lib/project-pair-result-capture.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
function attachPairResultCapture(ctx) {
function begin(caller, partner, token, message, projectSlug) {
if (!ctx.resultOutbox) return { ok: true };
var driverOriginId = ctx.resultOutbox.stableOrigin(caller), workerOriginId = ctx.resultOutbox.stableOrigin(partner);
if (!driverOriginId || !workerOriginId) return { ok: false, error: "durable session origins are required before dispatch" };
var started = ctx.resultOutbox.begin({ ownerId: caller.ownerId, projectSlug: projectSlug,
driverOriginId: driverOriginId, workerOriginId: workerOriginId,
taskId: token.taskId, generation: token.generation, driverSessionId: caller.localId,
workerSessionId: partner.localId, historyStartIndex: token.startIndex, message: message,
deliveryRoute: token.wait === true ? "blocking" : "callback" });
if (started.ok) token.outboxKey = started.key;
return started;
}
function finish(caller, partner, token, delivered) {
var group = ctx.store.groupForMember(partner.localId);
if (token._captureFinished || partner._pairDelegation !== token || !group || group.id !== token.groupId || group.members.indexOf(caller.localId) === -1 || group.members.indexOf(partner.localId) === -1 || (group.pair && (group.pair.driverId !== caller.localId || group.pair.workerId !== partner.localId)) || ctx.sm.sessions.get(caller.localId) !== caller || ctx.sm.sessions.get(partner.localId) !== partner || caller.ownerId !== partner.ownerId) return false;
token._captureFinished = true;
try {
if (typeof ctx.onPartnerResult === "function") ctx.onPartnerResult(caller, partner, token);
if (token.outboxKey && ctx.resultOutbox) {
var finalized = ctx.resultOutbox.markFinalized(token.outboxKey);
if (!finalized.ok) throw new Error(finalized.error || "completed result finalization could not be saved");
}
ctx.finishDelegation(group, caller, partner, token);
if (typeof ctx.onStateChange === "function" && token.outboxKey) ctx.onStateChange(caller, ctx.resultOutbox.get(token.outboxKey));
} catch (error) {
token._captureFinished = false;
token._capturePending = true;
token._captureHookError = error.message || String(error);
if (token.outboxKey && ctx.resultOutbox) {
var failed = ctx.resultOutbox.markFinalizationFailure(token.outboxKey, token._captureHookError);
if (typeof ctx.onStateChange === "function") ctx.onStateChange(caller, Object.assign({}, failed.record || ctx.resultOutbox.get(token.outboxKey), {
finalizationState: "failed", finalizationError: token._captureHookError, persistenceError: failed.ok ? null : failed.error,
}));
}
return false;
}
delete token._capturePending;
var resumed = delivered ? true : ctx.resumeDriverWithResult(caller, partner, token);
ctx.drain(caller, partner);
return resumed;
}
function capture(caller, partner, token, outcome) {
if (!ctx.resultOutbox || !token.outboxKey) return finish(caller, partner, token);
if (token._captureInProgress || token._captureComplete) return false;
token._captureInProgress = true;
var saved;
try {
saved = ctx.resultOutbox.capture(token.outboxKey, outcome, function (retryResult) {
token._captureInProgress = false;
if (retryResult && retryResult.ok) {
token._captureComplete = true;
var record = ctx.resultOutbox.get(token.outboxKey);
if (typeof ctx.onStateChange === "function") ctx.onStateChange(caller, record);
if (record && record.deliveryRoute === "callback") deliver(caller, partner, token);
else finish(caller, partner, token);
}
});
} catch (error) {
token._captureInProgress = false;
token._capturePending = true;
token._captureHookError = error.message || String(error);
return false;
}
if (saved.ok) { token._captureInProgress = false; token._captureComplete = true; }
if (saved.ok) {
var current = ctx.resultOutbox.get(token.outboxKey);
if (typeof ctx.onStateChange === "function") ctx.onStateChange(caller, current);
var delivered = current && current.deliveryRoute === "callback" ? deliver(caller, partner, token) : finish(caller, partner, token);
return delivered === false && token._capturePending ? false : delivered;
}
token._capturePending = true;
if (!token._captureFailureReported) {
token._captureFailureReported = true;
try { ctx.sm.sendAndRecord(caller, { type: "error", text: "The Split Worker completed, but its result could not be durably captured yet. The task remains retained for bounded server retry." }); } catch (error) { token._captureNoticeError = error.message || String(error); }
}
return false;
}
function deliver(caller, partner, token) {
if (!ctx.delivery) return finish(caller, partner, token);
return ctx.delivery.wake(caller, partner, token);
}
return { begin: begin, capture: capture, deliver: deliver, finish: finish };
}

module.exports = { attachPairResultCapture: attachPairResultCapture };
139 changes: 139 additions & 0 deletions lib/project-pair-result-delivery.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,139 @@
var RETRY_DELAY_MS = 50;
var DEFAULT_ACCEPTANCE_TIMEOUT_MS = 15000;

function attachPairResultDelivery(ctx) {
function recordMatches(caller, partner, token, record) {
return !!(record && record.key === token.outboxKey && record.ownerId === (caller.ownerId || null) &&
record.driverOriginId === ctx.outbox.stableOrigin(caller) && record.workerOriginId === ctx.outbox.stableOrigin(partner) &&
record.taskId === token.taskId && record.generation === token.generation && record.deliveryRoute === "callback");
}
function valid(caller, partner, token) {
var group = ctx.store.groupForMember(partner.localId), record = ctx.outbox.get(token.outboxKey);
return !!(group && group.id === token.groupId && group.members.indexOf(caller.localId) !== -1 &&
group.members.indexOf(partner.localId) !== -1 && group.pair && group.pair.driverId === caller.localId &&
group.pair.workerId === partner.localId && ctx.sm.sessions.get(caller.localId) === caller &&
ctx.sm.sessions.get(partner.localId) === partner && caller.ownerId === partner.ownerId &&
!caller.destroying && !partner.destroying && !ctx.blockedReason(caller) && partner._pairDelegation === token &&
recordMatches(caller, partner, token, record));
}
function currentAttempt(caller, partner, token, attemptId) {
var record = ctx.outbox.get(token.outboxKey);
return valid(caller, partner, token) && record && record.deliveryState === "attempting" && record.attemptId === attemptId;
}
function textFor(partner, token) {
var failure = token.failure || null, response = token.response || "";
if (token.interrupted) return "[Split Worker execution interrupted] The user interrupted the Split Worker mid-turn. Its work is PARTIAL and unverified — do not treat it as finished. Review what was done and decide next steps with the user.";
return "A Split Worker task delegated through send_to_partner has finished.\n\nOriginal task:\n" + token.message + "\n\n" + (failure ? "Split Worker error:\n" + failure : "Split Worker result:\n" + (response || "(No text response was recorded.)")) + "\n\nReview the result, verify it as needed, and continue the task.";
}
function finish(caller, partner, token) { return ctx.finish(caller, partner, token, true); }
function clearAttemptTimer(token, attemptId) {
if (!token._deliveryTimer || token._deliveryTimer.attemptId !== attemptId) return;
clearTimeout(token._deliveryTimer.timer);
delete token._deliveryTimer;
}
function retry(caller, partner, token, attemptId, reason) {
clearAttemptTimer(token, attemptId);
if (!currentAttempt(caller, partner, token, attemptId)) return { ok: false, stale: true };
var pending = ctx.outbox.markDelivery(token.outboxKey, attemptId, "pending", { deliveryError: reason || null });
if (!pending.ok) {
if (typeof ctx.onStateChange === "function") ctx.onStateChange(caller, ctx.outbox.get(token.outboxKey));
return { ok: false, error: pending.error, stale: pending.stale };
}
if (typeof ctx.onStateChange === "function") ctx.onStateChange(caller, pending.record);
if (pending.record.deliveryAttemptCount < ctx.outbox.maxDeliveryAttempts) {
setTimeout(function () { wake(caller, partner, token); }, RETRY_DELAY_MS * pending.record.deliveryAttemptCount);
}
return { ok: false, pending: true };
}
function uncertain(caller, partner, token, attemptId, reason) {
clearAttemptTimer(token, attemptId);
if (!currentAttempt(caller, partner, token, attemptId)) return false;
token._deliveryUncertain = true;
var saved = ctx.outbox.markDelivery(token.outboxKey, attemptId, "uncertain", { deliveryError: reason || "delivery acceptance is uncertain" });
if (saved.ok && typeof ctx.onStateChange === "function") ctx.onStateChange(caller, saved.record);
return false;
}
function accepted(caller, partner, token, attemptId) {
clearAttemptTimer(token, attemptId);
if (!currentAttempt(caller, partner, token, attemptId)) return { ok: false, stale: true };
var saved = ctx.outbox.markDelivery(token.outboxKey, attemptId, "accepted", { acceptedAt: Date.now() });
if (!saved.ok) {
token._deliveryUncertain = true;
if (typeof ctx.onStateChange === "function") ctx.onStateChange(caller, ctx.outbox.get(token.outboxKey));
return { ok: false, uncertain: true, error: saved.error };
}
token._deliveryAccepted = true;
if (typeof ctx.onStateChange === "function") ctx.onStateChange(caller, saved.record);
return { ok: true, accepted: true, finished: finish(caller, partner, token) !== false };
}
function recordTranscript(caller, token, message, attemptId) {
var outboxRecord = ctx.outbox.get(token.outboxKey);
if (outboxRecord && outboxRecord.transcriptPersisted) return { ok: true };
var matches = function (entry) { return !!(entry && entry.pairResultOutboxKey === token.outboxKey); };
var existing = caller.history && caller.history.some(matches);
var durable = existing && typeof ctx.sm.hasDurableSessionRecord === "function" && ctx.sm.hasDurableSessionRecord(caller, matches);
if (!durable) {
var record = { type: "user_message", text: message, _internal: true, partnerResult: true,
pairResultOutboxKey: token.outboxKey, pairResultAttemptId: attemptId };
var saved;
try {
if (typeof ctx.sm.sendAndRecordDurably !== "function") return { ok: false, error: "Durable Driver transcript recording is unavailable" };
saved = ctx.sm.sendAndRecordDurably(caller, record);
} catch (error) { return { ok: false, error: error.message || String(error) }; }
if (saved !== true) {
if (caller.history && caller.history[caller.history.length - 1] === record) caller.history.pop();
return { ok: false, error: "Driver transcript persistence failed" };
}
}
var marked = ctx.outbox.markTranscript(token.outboxKey, attemptId);
return marked.ok ? { ok: true } : { ok: false, error: marked.error, stale: marked.stale };
}
function wake(caller, partner, token) {
if (!token || token._deliveryUncertain || token._deliveryAccepted) return { ok: false, deferred: true };
if (ctx.blockedReason(caller)) {
var blocked = ctx.outbox.markBlocked(token.outboxKey, ctx.blockedReason(caller), "human_stop");
if (blocked.ok && typeof ctx.onStateChange === "function") ctx.onStateChange(caller, blocked.record);
return { ok: false, blocked: true };
}
if (!valid(caller, partner, token)) return { ok: false, deferred: true };
if (caller._queryStarting) return { ok: false, deferred: true, reason: "query_starting" };
var started = ctx.outbox.beginDelivery(token.outboxKey);
if (!started.ok) return { ok: false, pending: true, error: started.error };
if (started.terminal || started.inFlight || started.blocked || started.exhausted) return { ok: true, terminal: !!started.terminal, exhausted: !!started.exhausted };
var attemptId = started.attemptId, message = textFor(partner, token);
var transcript = recordTranscript(caller, token, message, attemptId);
if (!transcript.ok) return retry(caller, partner, token, attemptId, transcript.error);
var sdk = ctx.getSdk();
if (!sdk || typeof sdk.pushMessage !== "function" || typeof sdk.startQuery !== "function") {
return retry(caller, partner, token, attemptId, "SDK bridge is unavailable");
}
try {
if (caller.queryInstance) {
if (sdk.pushMessage(caller, message) === true) {
return accepted(caller, partner, token, attemptId);
}
return retry(caller, partner, token, attemptId, "Driver query rejected the result");
}
token._deliveryStartupAttemptId = attemptId;
var beforePush = function () { return token._deliveryStartupAttemptId === attemptId && currentAttempt(caller, partner, token, attemptId); };
var onAccepted = function () { accepted(caller, partner, token, attemptId); };
var timeoutMs = Number(ctx.acceptanceTimeoutMs || DEFAULT_ACCEPTANCE_TIMEOUT_MS);
token._deliveryTimer = { attemptId: attemptId, timer: setTimeout(function () {
uncertain(caller, partner, token, attemptId, "delivery acceptance timed out");
}, timeoutMs) };
var startup = sdk.startQuery(caller, message, undefined, ctx.getLinuxUserForSession(caller), beforePush, onAccepted);
Promise.resolve(startup).then(function (result) {
if (!currentAttempt(caller, partner, token, attemptId)) return;
if (result === false) retry(caller, partner, token, attemptId, "Driver query rejected the result");
}, function (error) {
if (currentAttempt(caller, partner, token, attemptId)) uncertain(caller, partner, token, attemptId, error.message || String(error));
});
return { ok: false, pending: true };
} catch (error) {
return uncertain(caller, partner, token, attemptId, error.message || String(error));
}
}
return { wake: wake };
}

module.exports = { attachPairResultDelivery: attachPairResultDelivery };
Loading
Loading