mirror of
https://github.com/KeygraphHQ/shannon.git
synced 2026-09-14 05:59:06 +02:00
fix(cli): keep shannon logs tailing through a Temporal blip
- End the interactive tail on the log's own terminal marker or Ctrl-C, so a transient Temporal outage no longer aborts the command with exit 1. - Rebuild the memoized Temporal client after a failed poll: a wedged gRPC channel was cached forever, so "retrying…" could never reconnect. - Keep start --follow (CI) bounded — a genuinely dead Temporal still fails the run instead of hanging.
This commit is contained in:
@@ -1,10 +1,12 @@
|
||||
/**
|
||||
* `shannon logs` command — tail a scan's live log.
|
||||
*
|
||||
* The log file is streamed for its content; completion is decided by Temporal (the
|
||||
* workflow's status), so a worker that dies mid-run can't leave the tail hanging. Uses
|
||||
* chokidar for reliable cross-platform file watching and bounded synchronous reads to
|
||||
* prevent duplicate output.
|
||||
* The log file is streamed for its content and ends the tail on its own terminal marker
|
||||
* (`Scan COMPLETED/PARTIAL/FAILED/CANCELLED`) or Ctrl-C. Temporal's workflow status is a backstop
|
||||
* that also closes the tail when a worker dies without writing a marker — for interactive `logs` a
|
||||
* Temporal outage is never fatal (it keeps tailing); only `start --follow` (CI) treats a sustained
|
||||
* outage as a failure. Uses chokidar for reliable cross-platform file watching and bounded
|
||||
* synchronous reads to prevent duplicate output.
|
||||
*/
|
||||
|
||||
import fs from 'node:fs';
|
||||
@@ -127,10 +129,16 @@ export function resolveLogFile(workspaceId: string): string {
|
||||
}
|
||||
|
||||
export interface TailOptions {
|
||||
/** Workflow whose Temporal status decides when the tail stops. Without it, only Ctrl-C ends the tail. */
|
||||
/** Workflow whose Temporal status can end the tail (alongside the file's own terminal marker). */
|
||||
readonly workflowId?: string;
|
||||
/** Called if the tail ends because Temporal became unreachable, with the captured error. */
|
||||
readonly onUnreachable?: (lastError: string) => void;
|
||||
/**
|
||||
* Consecutive Temporal-outage polls before the watch gives up. Interactive `logs` passes
|
||||
* Infinity so a blip never ends the tail (the file marker or Ctrl-C do); `start --follow` (CI)
|
||||
* leaves it bounded so a genuinely dead Temporal fails the run instead of hanging.
|
||||
*/
|
||||
readonly maxConnectFailures?: number;
|
||||
}
|
||||
|
||||
/** Outcome of a tail: whether the streamed log already contained the worker's `Scan FAILED` block. */
|
||||
@@ -139,10 +147,11 @@ export interface TailResult {
|
||||
}
|
||||
|
||||
/**
|
||||
* Stream a scan's log to the terminal until the workflow closes (completion comes from Temporal,
|
||||
* or Ctrl-C). A Temporal outage is warned about and, if sustained, ends the tail with a diagnostic.
|
||||
* Never exits the process: plain `logs` exits; `start --follow` reads the workflow outcome first.
|
||||
* Reports whether the log already showed the failure, so a caller need not print it a second time.
|
||||
* Stream a scan's log to the terminal until the file shows a terminal marker, the workflow closes,
|
||||
* or Ctrl-C. A Temporal outage is warned about; if it reaches `maxConnectFailures` the tail ends
|
||||
* with a diagnostic (bounded for `start --follow`), but interactive `logs` sets that to Infinity so
|
||||
* an outage keeps tailing. Never exits the process itself. Reports whether the log already showed
|
||||
* the failure, so a caller need not print it a second time.
|
||||
*/
|
||||
export function tailUntilComplete(logFile: string, opts: TailOptions = {}): Promise<TailResult> {
|
||||
return new Promise((resolve) => {
|
||||
@@ -186,13 +195,16 @@ export function tailUntilComplete(logFile: string, opts: TailOptions = {}): Prom
|
||||
|
||||
// 1. Output existing content, then stream anything appended. A per-agent file can be created
|
||||
// after the watcher starts, so `add` is handled too and streams it from its first line.
|
||||
// The file's own `Scan COMPLETED/PARTIAL/FAILED/CANCELLED` marker ends the tail on its own —
|
||||
// a Temporal round-trip is a backstop for a worker that dies without writing one, not the
|
||||
// only way to stop.
|
||||
watcher = watch(logFile, { persistent: true });
|
||||
const onFsEvent = (): void => {
|
||||
if (flush() && !opts.workflowId) finish();
|
||||
if (flush()) finish();
|
||||
};
|
||||
watcher.on('change', onFsEvent);
|
||||
watcher.on('add', onFsEvent);
|
||||
if (flush() && !opts.workflowId) {
|
||||
if (flush()) {
|
||||
finish();
|
||||
return;
|
||||
}
|
||||
@@ -200,10 +212,12 @@ export function tailUntilComplete(logFile: string, opts: TailOptions = {}): Prom
|
||||
// 2. Ctrl-C stops watching.
|
||||
process.once('SIGINT', finish);
|
||||
|
||||
// 3. Temporal decides completion. Without a workflow id, the tail relies on Ctrl-C alone.
|
||||
// 3. Temporal backstops completion for a worker that dies without a marker. Without a workflow
|
||||
// id, the tail relies on the file marker or Ctrl-C alone.
|
||||
if (opts.workflowId) {
|
||||
waitForWorkflowClose(opts.workflowId, {
|
||||
signal: controller.signal,
|
||||
...(opts.maxConnectFailures !== undefined ? { maxConnectFailures: opts.maxConnectFailures } : {}),
|
||||
onConnectionTrouble: (lastError) => {
|
||||
if (!done) console.error(`\n⚠ Lost contact with Temporal, retrying… (${lastError})`);
|
||||
},
|
||||
@@ -277,13 +291,13 @@ export interface LogsOptions {
|
||||
|
||||
function tailFileToExit(logFile: string, workflowId: string | undefined, label: string): void {
|
||||
console.error(stdoutIsTerminal() ? `${label}: ${logFile}` : label);
|
||||
let unreachable = false;
|
||||
tailUntilComplete(logFile, {
|
||||
...(workflowId ? { workflowId } : {}),
|
||||
onUnreachable: () => {
|
||||
unreachable = true;
|
||||
},
|
||||
}).finally(() => process.exit(unreachable ? 1 : 0));
|
||||
// Interactive tail: a Temporal outage must never end the session. The file's terminal marker or
|
||||
// Ctrl-C stop it; Temporal stays a soft backstop that reconnects and closes the tail on a
|
||||
// silent worker death, but its unreachability is never fatal here.
|
||||
maxConnectFailures: Number.POSITIVE_INFINITY,
|
||||
}).finally(() => process.exit(0));
|
||||
}
|
||||
|
||||
export function logs(workspaceId: string, options: LogsOptions = {}): void {
|
||||
|
||||
@@ -70,13 +70,31 @@ let clientPromise: Promise<Client> | null = null;
|
||||
|
||||
function getClient(): Promise<Client> {
|
||||
if (!clientPromise) {
|
||||
clientPromise = Connection.connect({ address: ADDRESS }).then(
|
||||
const pending = Connection.connect({ address: ADDRESS }).then(
|
||||
(connection) => new Client({ connection, namespace: NAMESPACE }),
|
||||
);
|
||||
// A rejected connect must not be cached forever: clear the memo so the next call rebuilds
|
||||
// instead of replaying the same failure. Scoped to `pending` so a later successful reconnect
|
||||
// that replaced the memo is left untouched.
|
||||
pending.catch(() => resetClient(pending));
|
||||
clientPromise = pending;
|
||||
}
|
||||
return clientPromise;
|
||||
}
|
||||
|
||||
/**
|
||||
* Drop the memoized client so the next {@link getClient} builds a fresh Connection. The underlying
|
||||
* gRPC channel can wedge such that every reused call fails identically ("Unexpected error while
|
||||
* making gRPC request"), and only a new Connection recovers. Best-effort closes the old channel.
|
||||
* When `only` is given, the memo is cleared only if it still holds that exact promise.
|
||||
*/
|
||||
function resetClient(only?: Promise<Client>): void {
|
||||
if (only !== undefined && clientPromise !== only) return;
|
||||
const previous = clientPromise;
|
||||
clientPromise = null;
|
||||
previous?.then((client) => client.connection.close()).catch(() => {});
|
||||
}
|
||||
|
||||
/** Describe a scan: status, timing, and the agents currently running (from pendingActivities). Null if not found. */
|
||||
export async function describeScan(workflowId: string): Promise<ScanDescription | null> {
|
||||
const client = await getClient();
|
||||
@@ -201,6 +219,9 @@ export async function waitForWorkflowClose(workflowId: string, opts: WatchOption
|
||||
if (err instanceof WorkflowNotFoundError) {
|
||||
return { reason: 'closed' };
|
||||
}
|
||||
// Drop the wedged channel so the next poll dials a fresh one; a cached dead channel would
|
||||
// otherwise fail every retry identically and never recover.
|
||||
resetClient();
|
||||
connectFailures++;
|
||||
lastError = err instanceof Error ? err.message : String(err);
|
||||
if (!warned && connectFailures >= warnAfterFailures) {
|
||||
|
||||
Reference in New Issue
Block a user