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
62 changes: 46 additions & 16 deletions src/platform/credentials.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,9 @@
import type { Credential, CredentialInfo, CredentialStore } from "@earendil-works/pi-ai";
import type {
AuthOperationOptions,
Credential,
CredentialInfo,
CredentialStore,
} from "@earendil-works/pi-ai";

const SECRET_PREFIX = "provider:";

Expand Down Expand Up @@ -65,21 +70,33 @@ export class PortableCredentialStore implements CredentialStore {
modify(
providerId: string,
fn: (current: Credential | undefined) => Promise<Credential | undefined>,
options?: AuthOperationOptions,
): Promise<Credential | undefined> {
return this.#enqueue(providerId, async () => {
const current = await this.read(providerId);
const next = await fn(current);
if (next === undefined) return current;
await this.#write(providerId, next);
return next;
});
const signal = options?.signal;
return this.#enqueue(
providerId,
async () => {
const current = await this.read(providerId);
signal?.throwIfAborted();
// Once fn starts, finish persistence: a refresh may already have rotated tokens.
const next = await fn(current);
if (next === undefined) return current;
await this.#write(providerId, next);
return next;
},
signal,
);
}

delete(providerId: string): Promise<void> {
return this.#enqueue(providerId, async () => {
this.#memory.delete(providerId);
if (this.#ctx) await this.#ctx.setSecret(`${SECRET_PREFIX}${providerId}`, "");
});
delete(providerId: string, options?: AuthOperationOptions): Promise<void> {
return this.#enqueue(
providerId,
async () => {
this.#memory.delete(providerId);
if (this.#ctx) await this.#ctx.setSecret(`${SECRET_PREFIX}${providerId}`, "");
},
options?.signal,
);
}

async setApiKey(providerId: string, key: string): Promise<void> {
Expand All @@ -98,13 +115,26 @@ export class PortableCredentialStore implements CredentialStore {
}
}

#enqueue<T>(providerId: string, task: () => Promise<T>): Promise<T> {
#enqueue<T>(providerId: string, task: () => Promise<T>, signal?: AbortSignal): Promise<T> {
if (signal?.aborted) return Promise.reject(signal.reason);
const previous = this.#chains.get(providerId) ?? Promise.resolve();
const run = previous.then(task, task);
let onAbort: (() => void) | undefined;
const start = () => {
if (onAbort) signal?.removeEventListener("abort", onAbort);
signal?.throwIfAborted();
return task();
};
const run = previous.then(start, start);
Comment thread
bajrangCoder marked this conversation as resolved.
// Caller cancellation must not release the serialized queue.
this.#chains.set(
providerId,
run.catch(() => undefined),
);
return run;
if (!signal) return run;
return new Promise<T>((resolve, reject) => {
onAbort = () => reject(signal.reason);
signal.addEventListener("abort", onAbort, { once: true });
void run.then(resolve, reject);
});
}
}
Loading