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
14 changes: 14 additions & 0 deletions packages/workflow-executor/src/adapters/server-types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -370,9 +370,23 @@ export const ServerAutomatedInboxAssignmentsResponseSchema = z.object({
assignments: z.array(ServerAutomatedInboxAssignmentSchema),
});

export const SERVER_AUTOMATED_INBOX_READ_FAILURE_REASONS = [
'agent-forbidden',
'agent-unreachable',
'segment-read-failed',
] as const;
export type ServerAutomatedInboxReadFailureReason =
(typeof SERVER_AUTOMATED_INBOX_READ_FAILURE_REASONS)[number];

export interface ServerAutomatedInboxReadFailure {
reason: ServerAutomatedInboxReadFailureReason;
httpStatus?: number;
}

export interface ServerAutomatedInboxSyncRequest {
closed: { recordId: string; stillInSegment: boolean }[];
candidates: string[];
readFailure?: ServerAutomatedInboxReadFailure;
}

export const SERVER_AUTOMATED_INBOX_SYNC_OUTCOMES = [
Expand Down
86 changes: 61 additions & 25 deletions packages/workflow-executor/src/automation-poller.ts
Original file line number Diff line number Diff line change
@@ -1,16 +1,18 @@
import type {
ServerAutomatedInboxAssignment,
ServerAutomatedInboxConfig,
ServerAutomatedInboxReadFailure,
} from './adapters/server-types';
import type { AutomationPort } from './ports/automation-port';
import type { Logger } from './ports/logger-port';
import type { SegmentReaderPort } from './ports/segment-reader-port';

import { AgentHttpError } from '@forestadmin/agent-client';
import { IANAZone } from 'luxon';

import createConsoleLogger from './adapters/console-logger';
import { DEFAULT_STOP_TIMEOUT_S } from './defaults';
import { AutomatedInboxGoneError, extractErrorMessage } from './errors';
import { AgentPortError, AutomatedInboxGoneError, extractErrorMessage } from './errors';
import InFlightRunRegistry from './in-flight-run-registry';

// One membership question per chunk, small enough that a `pk In (...)` stays a query an agent will
Expand All @@ -32,7 +34,7 @@ const MAX_EXCLUDED_RECORDS = 150;

// Every inbox reads the customer's agent several times. Sweeping them all at once piles those reads
// onto the customer's database, and agent-client's ten-second timeout turns the pile-up into inboxes
// that skip their sync.
// whose reads time out.
const MAX_CONCURRENT_INBOX_POLLS = 5;

// The orchestrator's poller lease lives three of these, so a dead holder is replaced within a
Expand Down Expand Up @@ -72,6 +74,53 @@ const isTerminalRun = (runState: string | null | undefined): boolean =>
const isLiveRun = (runState: string | null | undefined): boolean =>
runState != null && LIVE_RUN_STATES.has(runState);

const FORBIDDEN_AGENT_STATUSES: ReadonlySet<number> = new Set([401, 403]);

const UNREACHABLE_AGENT_STATUSES: ReadonlySet<number> = new Set([502, 503, 504]);

// Superagent's own timeout is ECONNABORTED. Any other failure without an HTTP answer (a JWT that
// cannot be signed, a malformed agent URL) is a fault on our side, not an agent to go and restart.
const UNREACHABLE_AGENT_ERROR_CODES: ReadonlySet<string> = new Set([
'ECONNABORTED',
'ECONNREFUSED',
'ECONNRESET',
'EAI_AGAIN',
'EHOSTUNREACH',
'ENETUNREACH',
'ENOTFOUND',
'EPIPE',
'ETIMEDOUT',
]);

// agent-client reports status 0 when a response carries none. 0 is not a real HTTP status, so a
// read that lands there reached no answer worth classifying by code: treat it as unreachable.
const isHttpStatus = (status: number): boolean => status >= 100 && status <= 599;

function classifyReadFailure(error: unknown): ServerAutomatedInboxReadFailure {
if (!(error instanceof AgentPortError)) return { reason: 'segment-read-failed' };

const { cause } = error;

if (!(cause instanceof AgentHttpError)) {
const code = (cause as { code?: unknown })?.code;

return typeof code === 'string' && UNREACHABLE_AGENT_ERROR_CODES.has(code)
? { reason: 'agent-unreachable' }
: { reason: 'segment-read-failed' };
}

const httpStatus = cause.status;

if (!isHttpStatus(httpStatus)) return { reason: 'agent-unreachable' };
if (FORBIDDEN_AGENT_STATUSES.has(httpStatus)) return { reason: 'agent-forbidden', httpStatus };

if (UNREACHABLE_AGENT_STATUSES.has(httpStatus)) {
return { reason: 'agent-unreachable', httpStatus };
}

return { reason: 'segment-read-failed', httpStatus };
}

export type AutomationPollerState = 'idle' | 'running' | 'draining' | 'stopped';

type PaddedPageReason =
Expand All @@ -82,17 +131,12 @@ type PaddedPageReason =
| 'field-without-not-in'
| 'not-in-refused';

/**
* `skipped` is a read that had nothing to ask and so never reached the agent. It is not a
* failure, and it is not proof the agent answers either — which is the distinction the sync
* decision rests on.
*/
interface SegmentRead<T> {
outcome: 'ok' | 'failed' | 'skipped';
items: T[];
paddedPageReason?: PaddedPageReason;
requestedPageSize?: number;
pagesRead?: number;
failure?: ServerAutomatedInboxReadFailure;
}

export interface AutomationPollerConfig {
Expand Down Expand Up @@ -354,22 +398,16 @@ export default class AutomationPoller {
),
]);

// Reaching the agent at all is what the sync attests to. Reporting an empty poll when every
// read failed would tell the orchestrator this inbox is being swept while nothing is, which
// is the one thing a "no sync received" alert must never be lied to about.
if (closed.outcome !== 'ok' && candidates.outcome !== 'ok') {
this.logger('Error', 'Could not reach the agent, reporting nothing for this inbox', {
...logContext,
});

return;
}
// The candidate read is the one that starts runs, so its failure is the one worth fixing.
const readFailure = candidates.failure ?? closed.failure;

// Sent even when both lists are empty: this call is what tells the orchestrator the inbox is
// still being polled, and an inbox with nothing to do must not read as an inbox nobody polls.
// Sent even when both lists are empty or both reads failed: the orchestrator keeps the last
// read failure until a sync comes without one, so a skipped sync would leave the settings
// panel with a stale warning, or none at all.
const results = await this.config.automationPort.sync(config.inboxId, {
closed: closed.items,
candidates: candidates.items,
readFailure,
});

this.logger('Info', 'Automated inbox polled', {
Expand All @@ -380,6 +418,7 @@ export default class AutomationPoller {
candidatePageSize: candidates.requestedPageSize,
candidatePagesRead: candidates.pagesRead,
paddedPageReason: candidates.paddedPageReason,
readFailure: readFailure?.reason,
outcomes: results.reduce<Record<string, number>>(
(counts, { outcome }) => ({ ...counts, [outcome]: (counts[outcome] ?? 0) + 1 }),
{},
Expand Down Expand Up @@ -486,7 +525,7 @@ export default class AutomationPoller {
return false;
});

if (readable.length === 0) return { outcome: 'skipped', items: [] };
if (readable.length === 0) return { items: [] };

const stillInSegment = new Set<string>();

Expand All @@ -502,7 +541,6 @@ export default class AutomationPoller {
}

return {
outcome: 'ok',
items: readable.map(recordId => ({
recordId,
stillInSegment: stillInSegment.has(recordId),
Expand All @@ -528,7 +566,7 @@ export default class AutomationPoller {
error: extractErrorMessage(error),
});

return { outcome: 'failed', items: [] };
return { items: [], failure: classifyReadFailure(error) };
}
}

Expand Down Expand Up @@ -556,7 +594,6 @@ export default class AutomationPoller {

// An agent that ignores an operator it does not know would hand known records back.
return {
outcome: 'ok',
items: page.filter(recordId => !knownSet.has(recordId)),
requestedPageSize: config.maxConcurrentRuns,
};
Expand Down Expand Up @@ -651,7 +688,6 @@ export default class AutomationPoller {
}

return {
outcome: 'ok',
items: [...candidates],
paddedPageReason,
requestedPageSize,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -438,6 +438,25 @@ describe('AgentClientSegmentReader', () => {

await expect(reader.listRecordIds(makeQuery())).rejects.toThrow(AgentPortError);
});

it('should keep the HTTP status of the agent answer on the port error', async () => {
nock(AGENT_URL).get('/forest/orders').query(true).reply(403, {});

await expect(reader.listRecordIds(makeQuery())).rejects.toMatchObject({
cause: expect.objectContaining({ name: 'AgentHttpError', status: 403 }),
});
});

it('should keep the network error code when the agent does not answer', async () => {
nock(AGENT_URL)
.get('/forest/orders')
.query(true)
.replyWithError({ code: 'ECONNREFUSED', message: 'connect ECONNREFUSED' });

await expect(reader.listRecordIds(makeQuery())).rejects.toMatchObject({
cause: expect.objectContaining({ code: 'ECONNREFUSED' }),
});
});
});

describe('field operators', () => {
Expand Down
Loading
Loading