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
87 changes: 84 additions & 3 deletions desktop/src/main/__tests__/ipc-provider-jobs.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@ function setup(
),
cancelFor: vi.fn(async () => undefined),
}
const send = vi.fn()
const deps = {
cli,
state: { workspaceContext: () => "ctx", providerList: () => [] },
Expand All @@ -74,12 +75,12 @@ function setup(
onDrain: async () => undefined,
},
pty: { cancelFor: vi.fn(async () => undefined) },
getMainWindow: () => null,
getMainWindow: () => ({ webContents: { send } }),
providerJobs,
}
// biome-ignore lint/suspicious/noExplicitAny: partial test doubles
registerIpcHandlers(deps as any)
return { providerJobs, cli }
return { providerJobs, cli, send }
}

function invoke(channel: string, args: Record<string, unknown>) {
Expand Down Expand Up @@ -154,7 +155,7 @@ describe("provider job lifecycle over IPC", () => {

it("runs set-source then init on update, clearing the job once", async () => {
const seen: string[][] = []
const { providerJobs } = setup((cliArgs) => {
const { providerJobs, send } = setup((cliArgs) => {
seen.push(cliArgs)
return { lines: [], code: 0 }
})
Expand All @@ -178,6 +179,53 @@ describe("provider job lifecycle over IPC", () => {
expect(providerJobs.get("docker")?.error).toBeTruthy()
})


it("streams update phases and completes the provider job", async () => {
const seen: string[][] = []
const { providerJobs, send } = setup((cliArgs) => {
seen.push(cliArgs)
return { lines: [statusLine(cliArgs[1] === "init" ? "running_init" : "downloading")], code: 0 }
})

const commandId = await invoke("provider_update_streaming", { name: "docker" })
expect(commandId).toEqual(expect.any(String))
expect(providerJobs.get("docker")?.activity).toBe("updating")

await vi.waitFor(() => expect(providerJobs.get("docker")).toBeUndefined())
expect(seen.map((args) => args[1])).toEqual(["set-source", "init"])
expect(send).toHaveBeenCalledWith(
"command-progress",
expect.objectContaining({ commandId, success: true, done: true }),
)
})

it("retains update failure details after a streaming update", async () => {
const { providerJobs } = setup((cliArgs) => ({
lines: [],
code: cliArgs[1] === "init" ? 1 : 0,
}))

await invoke("provider_update_streaming", { name: "docker" })
await vi.waitFor(() => expect(providerJobs.get("docker")?.error).toBeTruthy())
})


it("terminates a streaming update when provider refresh fails", async () => {
const { providerJobs, send } = setup(() => ({ lines: [], code: 0 }))
providerJobs.setRefresh(() => Promise.reject(new Error("refresh boom")))

const commandId = await invoke("provider_update_streaming", { name: "docker" })

await vi.waitFor(() =>
expect(providerJobs.get("docker")?.errorCode).toBe("provider_refresh_failed"),
)
expect(send).toHaveBeenCalledWith(
"command-progress",
expect.objectContaining({ commandId, success: false, done: true }),
)
})


it("does not blame a successful init for a refresh failure afterward", async () => {
const { providerJobs } = setup(() => ({
lines: [statusLine("running_init"), statusLine("ready")],
Expand All @@ -189,4 +237,37 @@ describe("provider job lifecycle over IPC", () => {

expect(providerJobs.get("docker")?.error).not.toBe("refresh boom")
})

it("refreshes and clears a retained provider refresh failure", async () => {
const { providerJobs } = setup(() => ({ lines: [], code: 0 }))
const refresh = vi.fn().mockResolvedValue(undefined)
providerJobs.setRefresh(refresh)
providerJobs.start("docker", "updating")
await providerJobs.finish("docker", {
code: "provider_refresh_failed",
message: "status unavailable",
})

const result = await invoke("provider_refresh_state", { name: "docker" })

expect(result).toEqual({ ok: true })
expect(refresh).toHaveBeenCalledOnce()
expect(providerJobs.get("docker")).toBeUndefined()
})

it("returns failure and retains recovery when provider refresh still fails", async () => {
const { providerJobs } = setup(() => ({ lines: [], code: 0 }))
providerJobs.setRefresh(() => Promise.reject(new Error("still unavailable")))
providerJobs.start("docker", "updating")
await providerJobs.finish("docker", {
code: "provider_refresh_failed",
message: "status unavailable",
})

const result = await invoke("provider_refresh_state", { name: "docker" })

expect(result).toEqual({ ok: false, message: "still unavailable" })
expect(providerJobs.get("docker")?.errorCode).toBe("provider_refresh_failed")
})

})
28 changes: 28 additions & 0 deletions desktop/src/main/__tests__/provider-jobs.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -194,4 +194,32 @@ describe("ProviderJobs", () => {

expect(jobs.get("docker")).toBeUndefined()
})

it("retries refresh for a retained refresh failure", async () => {
const refresh = vi.fn().mockResolvedValue(undefined)
jobs.setRefresh(refresh)
jobs.start("docker", "updating")
await jobs.finish("docker", {
code: "provider_refresh_failed",
message: "status unavailable",
})

await jobs.retryRefresh("docker")

expect(refresh).toHaveBeenCalledOnce()
expect(jobs.get("docker")).toBeUndefined()
})

it("keeps a refresh recovery job when retry fails", async () => {
jobs.setRefresh(() => Promise.reject(new Error("still unavailable")))
jobs.start("docker", "updating")
await jobs.finish("docker", {
code: "provider_refresh_failed",
message: "status unavailable",
})

await expect(jobs.retryRefresh("docker")).rejects.toThrow("still unavailable")
expect(jobs.get("docker")?.errorCode).toBe("provider_refresh_failed")
})

})
117 changes: 114 additions & 3 deletions desktop/src/main/ipc.ts
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,7 @@ type UpdateInfo = {
}

let providerUpdateCache: Record<string, UpdateInfo> = {}
let providerUpdateCacheCheckedAt: string | null = null

const IMAGE_CATALOG_URL =
process.env.DEVSY_IMAGE_CATALOG_URL ??
Expand Down Expand Up @@ -751,6 +752,113 @@ export function registerIpcHandlers(deps: IpcDependencies): {
})
})

ipcMain.handle(
"provider_update_streaming",
async (_event, args: { name: string }) => {
const cmdId = crypto.randomUUID()
const win = deps.getMainWindow()
providerJobs.start(args.name, "updating")

const sendProgress = (message: string, level?: string) => {
const formatted = redactSensitiveText(formatLogLine(message))
providerJobs.appendLog(args.name, formatted)
win?.webContents.send("command-progress", {
commandId: cmdId,
message: formatted,
level,
done: false,
})
}

const runStep = (cliArgs: string[]): Promise<void> =>
new Promise((resolve, reject) => {
cli.runStreaming(
cliArgs,
(line, stream, meta) => {
if (stream === "stdout") {
const envelope = parseCliEnvelope(line)
if (envelope?.kind === "status") {
providerJobs.reportStatus(
args.name,
redactOperationStatus(normalizeOperationStatus(envelope)),
)
return
}
}
sendProgress(line, meta?.level)
},
(code, cliError) => {
if (code === 0) {
resolve()
return
}
reject(
Object.assign(
new Error(
cliError?.message ?? `${cliArgs.join(" ")} exited with ${code}`,
),
{ cliError },
),
)
},
).catch(reject)
})

void (async () => {
let failure: CLIError | undefined
try {
sendProgress("Downloading provider update", "info")
await runStep(["provider", "set-source", args.name, "--use=false"])
sendProgress("Initializing updated provider", "info")
await runStep(["provider", "init", args.name])
} catch (error) {
failure = (error as { cliError?: CLIError }).cliError ?? {
code: "provider_update_failed",
message: errorMessage(error),
}
}
try {
await providerJobs.finish(
args.name,
failure ? redactCLIError(failure) : undefined,
)
} catch (error) {
failure = {
code: "provider_refresh_failed",
message: "The provider updated, but its current state could not be refreshed.",
hint: "Refresh provider status to try again.",
context: { cause: errorMessage(error) },
}
await providerJobs.finish(args.name, redactCLIError(failure))
}
win?.webContents.send("command-progress", {
commandId: cmdId,
message: redactSensitiveText(
formatLogLine(
failure ? "Provider update failed" : "Provider update complete",
failure ? "ERROR" : "INFO",
),
),
level: failure ? "error" : "info",
success: !failure,
cliError: failure ? redactCLIError(failure) : undefined,
done: true,
})
})()

return cmdId
},
)

ipcMain.handle("provider_refresh_state", async (_event, args: { name: string }) => {
try {
await providerJobs.retryRefresh(args.name)
return { ok: true } as const
} catch (error) {
return { ok: false, message: errorMessage(error) } as const
}
})

ipcMain.handle("provider_options", async (_event, args: { name: string }) => {
return cli.run(["provider", "get", args.name])
})
Expand Down Expand Up @@ -825,12 +933,14 @@ export function registerIpcHandlers(deps: IpcDependencies): {
ipcMain.handle("provider_check_updates", async () => {
const out = await computeUpdateChecks()
providerUpdateCache = out
providerUpdateCacheCheckedAt = new Date().toISOString()
return out
})

ipcMain.handle("provider_get_update_cache", async () => {
return providerUpdateCache
})
ipcMain.handle("provider_get_update_cache", async () => ({
updates: providerUpdateCache,
lastCheckedAt: providerUpdateCacheCheckedAt,
}))

ipcMain.handle("image_catalog_get", async () => {
const { cachePath, seedPath } = imageCatalogPaths()
Expand Down Expand Up @@ -1856,6 +1966,7 @@ export function registerIpcHandlers(deps: IpcDependencies): {
void (async () => {
try {
providerUpdateCache = await computeUpdateChecks()
providerUpdateCacheCheckedAt = new Date().toISOString()
} catch {
// Silently swallow background errors.
}
Expand Down
23 changes: 23 additions & 0 deletions desktop/src/main/provider-jobs.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ export interface ProviderJob {
errorCode?: string
errorHint?: string
errorContext?: Record<string, string>
logs?: string[]
}

export class ProviderJobs {
Expand Down Expand Up @@ -64,6 +65,15 @@ export class ProviderJobs {
this.emit()
}


/** Retain recent operation output so failures remain diagnosable after navigation. */
appendLog(name: string, line: string): void {
const job = this.jobs.get(name)
if (!job) return
this.jobs.set(name, { ...job, logs: [...(job.logs ?? []), line].slice(-500) })
this.emit()
}

/** Record the complete current-protocol status event without losing metadata. */
reportStatus(name: string, status: OperationStatus): void {
const job = this.jobs.get(name)
Expand Down Expand Up @@ -140,6 +150,19 @@ export class ProviderJobs {
this.emit()
}


/** Retry only the authoritative provider-state refresh after a completed operation. */
async retryRefresh(name: string): Promise<void> {
const job = this.jobs.get(name)
if (!job || job.errorCode !== "provider_refresh_failed") return
const generation = this.generations.get(name)
await this.refresh?.()
if (this.generations.get(name) !== generation) return
this.jobs.delete(name)
this.generations.delete(name)
this.emit()
}

/**
* Supplies a way to re-read provider state from disk, so a finished job
* isn't cleared before the list reflects what the command just wrote.
Expand Down
Loading
Loading