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
2 changes: 1 addition & 1 deletion docs/guides/MODULE_MAP.md
Original file line number Diff line number Diff line change
Expand Up @@ -228,7 +228,7 @@ Bootstraps UI, initializes store, wires remaining Tier 3 modules. All business l

| Module | Concern |
|--------|---------|
| `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-connection.js` + `websocket-lifecycle.js` + `websocket-watchdog.js` | WebSocket creation, epoch-guarded handshake and bounded jittered reconnect/auth timers, suspendable generation-guarded heartbeat/probe watchdogs, 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 |
Expand Down
6 changes: 5 additions & 1 deletion lib/driver-continuation-pair.js
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
var sessionProvenance = require("./session-provenance");
var splitRoles = require("./session-split-group-roles");

var ACTIVE_RUN_STATES = ["armed", "running", "reviewing", "waiting-worker", "waiting-user", "paused"];

Expand Down Expand Up @@ -37,7 +38,10 @@ function inspect(source, ctx) {
if (!group.pair || group.pair.driverId !== source.localId) {
return { ok: false, error: "Only the configured Driver can continue an attached split pair." };
}
var worker = ctx.sm && ctx.sm.sessions.get(group.pair.workerId);
var roles = splitRoles.normalizePair(group.pair, group.members);
if (!roles.ok || roles.kind === "adhoc") return { ok: false, error: "The attached split pair roles are invalid." };
if (roles.kind === "versioned") return { ok: false, error: "Driver continuation does not yet support a multi-Worker group." };
var worker = ctx.sm && ctx.sm.sessions.get(roles.workerIds[0]);
if (!worker || group.members.indexOf(worker.localId) === -1) return { ok: false, error: "The attached Split Worker is unavailable." };
if ((source.ownerId || null) !== (worker.ownerId || null)) return { ok: false, error: "The attached Split Worker owner changed." };
if (!sessionProvenance.isWorker(worker)) return { ok: false, error: "The attached split partner is not a verified Worker." };
Expand Down
11 changes: 11 additions & 0 deletions lib/multi-worker-feature.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
var ENABLED_GATE = Object.freeze({ feature: "multi-worker-runtime" });

function fromServerConfig(config) {
return config && config.multiWorkerRuntimeEnabled === true ? ENABLED_GATE : null;
}

function isEnabled(gate) {
return gate === ENABLED_GATE;
}

module.exports = { fromServerConfig: fromServerConfig, isEnabled: isEnabled };
22 changes: 14 additions & 8 deletions lib/project-pair-autonomous-stop.js
Original file line number Diff line number Diff line change
@@ -1,15 +1,21 @@
function stopAutonomousWork(session, roles, turnControl, taskControl, workerPermission) {
if (!roles) return false;
var runId = session.autonomousRun && session.autonomousRun.id;
var cancelled = taskControl.cancelAll(roles.worker, runId);
var owned = roles.worker._pairDelegation && roles.worker._pairDelegation.autonomousRunId === runId;
if (!owned && !cancelled) return true;
var workers = roles.workers || [roles.worker];
var cancelled = 0, owned = [];
for (var wi = 0; wi < workers.length; wi++) {
cancelled += taskControl.cancelAll(workers[wi], runId);
if (workers[wi]._pairDelegation && workers[wi]._pairDelegation.autonomousRunId === runId) owned.push(workers[wi]);
}
if (!owned.length && !cancelled) return true;
turnControl.markHumanStop(session);
if (!owned) return true;
taskControl.markInterruption(roles.worker, "user", "The Until complete run was stopped.");
workerPermission.cancelForSession(roles.worker, "The Until complete run was stopped.");
roles.worker.taskStopRequested = true;
if (roles.worker.abortController) roles.worker.abortController.abort();
if (!owned.length) return true;
for (var oi = 0; oi < owned.length; oi++) {
taskControl.markInterruption(owned[oi], "user", "The Until complete run was stopped.");
workerPermission.cancelForSession(owned[oi], "The Until complete run was stopped.");
owned[oi].taskStopRequested = true;
if (owned[oi].abortController) owned[oi].abortController.abort();
}
return true;
}

Expand Down
59 changes: 59 additions & 0 deletions lib/project-pair-close-transaction.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
var structuralClose = require("./project-pair-structural-close");

function attachPairCloseTransaction(ctx) {
function close(args, caller) {
var resolved = ctx.resolvePair(caller, args);
var partner = resolved.partner;
var token = partner._pairDelegation || null;
var interrupted = !!(partner.isProcessing || partner._queryStarting);
var interruption = token && interrupted ? {
source: "driver",
reason: "The Driver closed the Split Worker pair.",
targetTaskId: token.taskId,
requestedAt: Date.now(),
} : null;
var response = token ? ctx.responseText(partner.history || [], token.startIndex) : "";
var failure = token ? ctx.errorSince(partner.history || [], token.startIndex) : null;
var outcome = token ? ctx.taskControl.prepareCompletion(partner, caller, token,
interrupted ? "interrupted" : "completed", response, failure, interruption) : null;
var prepared = token ? ctx.resultCapture.prepare(caller, partner, token, outcome) : { ok: true, prepared: false };
if (!prepared.ok) throw new Error(prepared.error || "could not durably capture the Worker outcome");
var postCommitError = null;
var result = structuralClose.removeSelectedWorker(ctx.store, resolved.group, caller, partner,
ctx.multiWorkerFeature, { afterPersist: function () {
try {
if (token) {
var finalized = ctx.resultCapture.finalizePrepared(token, outcome);
if (!finalized.ok) throw new Error(finalized.error || "prepared close outcome could not be committed");
partner._pairClosing = true;
token.response = response;
token.failure = failure;
token.interrupted = interrupted && !failure;
ctx.taskControl.commitCompletion(partner, caller, token, outcome);
ctx.resultCapture.acceptPrepared(caller, partner, token);
}
} catch (error) { postCommitError = error; }
finally { delete partner._pairClosing; }
return { ok: true };
} });
if (!result.ok) {
if (token && prepared.prepared && !prepared.duplicate) {
var rolledBack = ctx.resultCapture.rollbackPrepared(token, outcome);
if (!rolledBack.ok) throw new Error((result.error || "could not close the Worker pair") +
"; prepared result rollback failed: " + (rolledBack.error || "unknown error"));
}
throw new Error(result.error || "could not close the Worker pair");
}
if (interrupted) {
partner.taskStopRequested = true;
if (partner.abortController) partner.abortController.abort();
}
ctx.workerPermission.cancelForSession(partner, "The Driver closed the Split Worker pair.");
if (postCommitError) partner._pairCloseFinalizationError = postCommitError.message || String(postCommitError);
return { status: "closed", partnerId: partner.localId, interrupted: interrupted,
historyPreserved: true };
}
return { close: close };
}

module.exports = { attachPairCloseTransaction: attachPairCloseTransaction };
18 changes: 18 additions & 0 deletions lib/project-pair-global-stop.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
function stopAllWorkers(session, roles, ctx) {
if (!roles) return false;
for (var i = 0; i < roles.workers.length; i++) {
var worker = roles.workers[i];
if (worker._pairDelegation && worker._pairDelegation.outboxKey && ctx.resultOutbox) {
ctx.resultOutbox.markBlocked(worker._pairDelegation.outboxKey, "The human stopped this Split Worker turn.", "human_stop");
}
ctx.taskControl.markInterruption(worker, "user", "The human stopped this Split Worker turn.");
ctx.workerPermission.cancelForSession(worker, "The human stopped this Split Worker turn.");
worker.taskStopRequested = true;
if (worker.abortController) {
try { worker.abortController.abort(); } catch (e) {}
}
}
return true;
}

module.exports = { stopAllWorkers: stopAllWorkers };
34 changes: 34 additions & 0 deletions lib/project-pair-lifecycle-status.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
var pairUsage = require("./project-pair-usage");

var MAX_TASK_PREVIEW_CHARS = 200;

function clampText(value, max) {
var text = typeof value === "string" ? value.replace(/\s+/g, " ").trim() : "";
if (text.length > max) text = text.slice(0, max - 1) + "…";
return text;
}

function continuityStatus(session) {
var history = session.history || [];
var userTurns = 0;
var errors = 0;
for (var i = 0; i < history.length; i++) {
if (!history[i]) continue;
if (history[i].type === "user_message") userTurns++;
else if (history[i].type === "error") errors++;
}
var lastActivity = typeof session.lastActivity === "number" ? session.lastActivity : null;
return { historyEntries: history.length, userTurns: userTurns, errorEntries: errors,
idleSeconds: lastActivity ? Math.max(0, Math.round((Date.now() - lastActivity) / 1000)) : null };
}

function activityStatus(session) {
var token = session._pairDelegation || null;
return { isProcessing: !!(session.isProcessing || session._queryStarting), delegated: !!token,
currentTask: token ? clampText(token.message, MAX_TASK_PREVIEW_CHARS) : "",
currentTaskId: token && token.taskId || null, interruption: session._pairInterruption || null,
lastTurnInterrupted: !!session._lastTurnInterrupted };
}

module.exports = { activityStatus: activityStatus, clampText: clampText, continuityStatus: continuityStatus,
contextStatus: pairUsage.contextStatus, firstNumber: pairUsage.firstNumber };
Loading
Loading