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
Original file line number Diff line number Diff line change
Expand Up @@ -7,15 +7,13 @@ export const CacheTable = React.memo(({ tables }) => {
const deleteCacheTable = useCacheTablesStore((state) => state.delete_cacheTable);
const resetCacheTables = useCacheTablesStore((state) => state.reset);

// ✅ local loading state
const [deletingAll, setDeletingAll] = useState(false);
const [deletingId, setDeletingId] = useState(null); // table id currently deleting

const deleteSingleCache = useCallback(
async (tableId) => {
if (deletingAll || deletingId) return; // avoid concurrent deletes
console.log("Delete cache table:", tableId);
setDeletingId(tableId);
setDeletingId(tableId);
try {
await deleteCacheTable(tableId);
} finally {
Expand All @@ -27,7 +25,6 @@ export const CacheTable = React.memo(({ tables }) => {

const deleteAllCache = useCallback(async () => {
if (deletingAll || deletingId) return;
console.log("Delete all cache tables");
setDeletingAll(true);
try {
await resetCacheTables();
Expand Down
2 changes: 0 additions & 2 deletions reactapp/features/DataStream/components/forecast/dataMenu.js
Original file line number Diff line number Diff line change
Expand Up @@ -150,11 +150,9 @@ const DataMenuControls = React.memo(function DataMenuControls() {
}
// reset();
const cacheKey = getCacheKey(model, date, forecast, cycle, ensemble, vpu, outputFile);
console.log('Generated cache key:', cacheKey);
set_cache_key(cacheKey);

const _prefix = makePrefix(model, date, forecast, cycle, ensemble, vpu, outputFile);
console.log('Generated S3 prefix:', _prefix);
set_prefix(_prefix);
});

Expand Down
1 change: 0 additions & 1 deletion reactapp/features/DataStream/components/map/Mapg.js
Original file line number Diff line number Diff line change
Expand Up @@ -119,7 +119,6 @@ const MainMap = () => {

const deckLayers = useMemo(() => {
if (!isFlowPathsVisible) return EMPTY_LAYERS;
// console.log('Rendering flow paths layer');
const varData = valuesByVar;
const numTimes = timesArr?.length || 0;

Expand Down
1 change: 0 additions & 1 deletion reactapp/features/DataStream/components/menus/CacheMenu.js
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,6 @@ export const CacheMenu = () => {
useEffect(() => {
const fetchCacheTables = async () => {
const files = await getFilesFromCache()
console.log("Fetched cache tables:", files);
set_cacheTables(files);
};
fetchCacheTables();
Expand Down
10 changes: 8 additions & 2 deletions reactapp/features/DataStream/lib/duckdbClient.js
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
import * as duckdb from "@duckdb/duckdb-wasm";

let dbPromise = null;

export function getDuckDB() {
if (!dbPromise) {
dbPromise = (async () => {
Expand All @@ -19,6 +18,11 @@ export function getDuckDB() {
const db = new duckdb.AsyncDuckDB(logger, worker);

await db.instantiate(bundle.mainModule, bundle.pthreadWorker);

await db.open({
accessMode: duckdb.DuckDBAccessMode.READ_WRITE,
opfs: { fileHandling: "auto" },
});

// Optional cleanup
URL.revokeObjectURL(workerUrl);
Expand All @@ -31,9 +35,11 @@ export function getDuckDB() {

export async function getConnection() {
const db = await getDuckDB();
return await db.connect();
const conn = await db.connect();
return conn;
}


// OPTIONAL: wipe all DB state (tables, etc) but keep worker
export async function resetDatabase() {
if (!dbPromise) return;
Expand Down
208 changes: 181 additions & 27 deletions reactapp/features/DataStream/lib/opfsCache.js
Original file line number Diff line number Diff line change
@@ -1,7 +1,13 @@
const CACHE_DIR = "nrds-arrow-cache";
import appAPI from "features/Tethys/services/api/app";
import { tableFromIPC } from "apache-arrow";
import { getNCFiles } from "./s3Utils";
import { DuckDBDataProtocol } from "@duckdb/duckdb-wasm";


const CACHE_DIR = "nrds-cache";
let cacheDirPromise = null;

function formatBytes(bytes, decimals = 2) {
export function formatBytes(bytes, decimals = 2) {
if (bytes === 0) return '0 Bytes';
const k = 1024;
const dm = decimals < 0 ? 0 : decimals;
Expand Down Expand Up @@ -30,40 +36,97 @@ async function getCacheDir() {
}
}

export async function saveArrowToCache(key, buffer) {
// async function saveArrowToCache(url, vpu_gpkg, writable) {
async function saveArrowToCache(url, writable) {
try{
const ncFile = getNCFiles(url);
const buffer = await appAPI.getArrowPerVpu({
ncFile,
});

let dataToWrite;

if (buffer instanceof ArrayBuffer) {
dataToWrite = new Uint8Array(buffer);
} else if (ArrayBuffer.isView(buffer)) {
// covers Uint8Array, DataView, etc.
dataToWrite = new Uint8Array(buffer.buffer, buffer.byteOffset, buffer.byteLength);
} else if (buffer instanceof Blob) {
dataToWrite = buffer;
} else {
console.error("saveArrowToCache: unexpected buffer type", buffer);
throw new Error("saveArrowToCache: expected ArrayBuffer, TypedArray, or Blob");
}

await writable.write(dataToWrite);
await writable.close();
}
catch(error){
console.error("Error fetching Arrow data:", error);
throw error;
}
}

async function cacheParquetToOPFS(url, writable) {
try {
const PARQUETURL=`https://ciroh-community-ngen-datastream.s3.us-east-1.amazonaws.com/${url}`;
const res = await fetch(PARQUETURL, { cache: "no-store" });
if (!res.ok) throw new Error(`Failed to fetch ${url}: ${res.status}`);

// Stream to disk; avoids loading the entire file in memory
if (!res.body) {
const buf = await res.arrayBuffer();
await writable.write(new Uint8Array(buf));
await writable.close();
} else {
// WritableStream from OPFS supports pipeTo in modern browsers

await res.body.pipeTo(writable);
// pipeTo closes the destination by default
}
} catch (err) {
// If pipeTo fails mid-stream, attempt to close to release the lock.
try { await writable.close(); } catch (_) {}
throw err;
}
}

const sqlIdent = (s) => `"${String(s).replace(/"/g, '""')}"`;
const sqlStr = (s) => `'${String(s).replace(/'/g, "''")}'`;

const safeNameForKey = (key) => encodeURIComponent(key);
const tableNameForKey = (key) => String(key).replace(/\.(arrow|parquet)$/i, "");

function isNCFile(key) { return key.endsWith('.nc'); }

function isArrowFile(key) { return key.endsWith('.arrow');}

function isParquetFile(key) { return key.endsWith('.parquet'); }

export async function saveDataToCache(key, url) {
const dir = await getCacheDir();
if (!dir) return; // noop if OPFS unavailable

const safeName = encodeURIComponent(key) + ".arrow";
const safeName = encodeURIComponent(key);
const fileHandle = await dir.getFileHandle(safeName, { create: true });
const writable = await fileHandle.createWritable();

// 🔍 Make sure we always pass a proper binary type to write()
let dataToWrite;

if (buffer instanceof ArrayBuffer) {
dataToWrite = new Uint8Array(buffer);
} else if (ArrayBuffer.isView(buffer)) {
// covers Uint8Array, DataView, etc.
dataToWrite = new Uint8Array(buffer.buffer, buffer.byteOffset, buffer.byteLength);
} else if (buffer instanceof Blob) {
dataToWrite = buffer;
if (isArrowFile(key)) {
await saveArrowToCache(url, writable);
} else {
console.error("saveArrowToCache: unexpected buffer type", buffer);
throw new Error("saveArrowToCache: expected ArrayBuffer, TypedArray, or Blob");
await cacheParquetToOPFS(url, writable);
}

await writable.write(dataToWrite);
await writable.close();
const file = await fileHandle.getFile();
return formatBytes(file.size);
}
function ascii4(u8) {
return String.fromCharCode(...u8);
}

export async function loadArrowFromCache(key) {

export async function loadFromCache(key) {
const dir = await getCacheDir();
if (!dir) return null;

const safeName = encodeURIComponent(key) + ".arrow";
const safeName = encodeURIComponent(key);
try {
const fileHandle = await dir.getFileHandle(safeName);
const file = await fileHandle.getFile();
Expand All @@ -73,6 +136,82 @@ export async function loadArrowFromCache(key) {
}
}


async function doesTableExist(conn, tableName) {
const res = await conn.query(`
SELECT 1
FROM information_schema.tables
WHERE table_schema = 'main'
AND table_name = ${sqlStr(tableName)}
LIMIT 1
`);
return res.toArray().length > 0;
}

// async function createTableFromOPFSParquet({ db, conn, key }) {
// const safeName = encodeURIComponent(key);
// const fileUrl = `opfs://${CACHE_DIR}/${safeName}`;
// const tableName = key.replace(/\.parquet$/i, "");

// await conn.query(`
// CREATE TABLE ${sqlIdent(tableName)} AS
// SELECT * FROM read_parquet(${sqlStr(fileUrl)});
// `);
// }
async function createTableFromOPFSParquet({ conn, key }) {
// 1) Get the OPFS file handle from your cache directory
const cacheDir = await getCacheDir();
const safeName = encodeURIComponent(key);
const fileHandle = await cacheDir.getFileHandle(safeName);

// 2) Register it in DuckDB under some virtual path/name
const duckPath = `${CACHE_DIR}/${safeName}`; // can be any string you like
const bindings = conn.bindings; // This is the AsyncDuckDB instance

await bindings.registerFileHandle(
duckPath,
fileHandle,
DuckDBDataProtocol.BROWSER_FSACCESS,
true
);

// 3) Create table from that registered file name
const tableName = tableNameForKey(key);
await conn.query(`
CREATE TABLE ${sqlIdent(tableName)} AS
SELECT * FROM read_parquet(${sqlStr(duckPath)});
`);
}
async function createTableFromOPFSArrow({ conn, key }) {
const buffer = await loadFromCache(key);
if (!buffer) throw new Error(`Arrow cache missing after save: ${key}`);

const arrowTable = tableFromIPC(new Uint8Array(buffer));
const tableName = tableNameForKey(key);

await conn.insertArrowTable(arrowTable, { name: tableName });
}

export async function createTableFromOPFS({ conn, key, safeName }) {
const tableName = tableNameForKey(key);

if (await doesTableExist(conn, tableName)) {
console.debug(`Table "${tableName}" already exists, skipping.`);
return;
}

if (isArrowFile(key)) {
return createTableFromOPFSArrow({ conn, key });
}
if (isParquetFile(key)) {
return createTableFromOPFSParquet({ conn, key, safeName });
}

throw new Error(`Unsupported file type for key: ${key}`);
}



export async function getFilesFromCache() {
const dir = await getCacheDir();
if (!dir) return null;
Expand All @@ -81,16 +220,30 @@ export async function getFilesFromCache() {
for await (const handle of dir.values()) {
if (handle.kind !== "file") continue;
const file = await handle.getFile();
const id = decodeURIComponent(file.name.replace(".arrow", ""));
const id = decodeURIComponent(file.name);
files.push({id: id, name: id.replaceAll("_", "/"), size: formatBytes(file.size)});
}
return files;
}

export async function statFromCache(key) {
const dir = await getCacheDir();
if (!dir) return null;

const safeName = safeNameForKey(key);
try {
const fileHandle = await dir.getFileHandle(safeName);
const file = await fileHandle.getFile();
return { safeName, sizeBytes: file.size };
} catch {
return null;
}
}

export async function deleteFileFromCache(key) {
const dir = await getCacheDir();
if (!dir) return;
const safeName = encodeURIComponent(key) + ".arrow";
const safeName = encodeURIComponent(key);
try {
await dir.removeEntry(safeName);
return true;
Expand All @@ -109,8 +262,9 @@ export async function clearCache() {
}

export function getCacheKey(model, date, forecast, cycle, ensemble, vpu, outputFile) {
const newOutputFile = isNCFile(outputFile) ? outputFile.replace(".nc", ".arrow") : outputFile;
if (!ensemble){
return `${model}_${date}_${forecast}_${cycle}_${vpu}_${outputFile}`.replace(/\./g,'_').replace(/\//g,'_');
return `${model}_${date}_${forecast}_${cycle}_${vpu}`.replace(/\./g,'_').replace(/\//g,'_') + `_${newOutputFile}`; ;
}
return `${model}_${date}_${forecast}_${cycle}_${ensemble}_${vpu}_${outputFile}`.replace(/\./g,'_').replace(/\//g,'_');
return `${model}_${date}_${forecast}_${cycle}_${ensemble}_${vpu}`.replace(/\./g,'_').replace(/\//g,'_') + `_${newOutputFile}`;
}
Loading
Loading