import { BigQuery } from '@google-cloud/bigquery'
import { solanaInstructionDecoder, solanaPortalStream } from '@subsquid/pipes/solana'
import { bigqueryTarget } from '@subsquid/pipes/targets/bigquery'
import * as orcaWhirlpool from './abi/orca_whirlpool/index.js'
const bigquery = new BigQuery({ projectId: 'my-gcp-project' })
await solanaPortalStream({
id: 'orca-swaps',
portal: 'https://portal.sqd.dev/datasets/solana-mainnet',
outputs: solanaInstructionDecoder({
range: { from: '340,000,000' },
programId: orcaWhirlpool.programId,
instructions: { swap: orcaWhirlpool.instructions.swap },
}),
}).pipeTo(
bigqueryTarget({
client: { bigquery },
dataset: 'orca_swaps',
tables: [
{
table: 'swaps',
blockNumberColumn: 'slot',
schema: [
{ name: 'slot', type: 'INT64', mode: 'REQUIRED' },
{ name: 'transaction_index', type: 'INT64', mode: 'REQUIRED' },
{ name: 'instruction_address', type: 'STRING', mode: 'REQUIRED' },
// TIMESTAMP wire format is INT64 microseconds since epoch
{ name: 'block_timestamp', type: 'TIMESTAMP', mode: 'REQUIRED' },
{ name: 'program_id', type: 'STRING', mode: 'REQUIRED' },
],
clusterBy: ['program_id'],
},
],
onData: async ({ store, data }) => {
store.insert(
'swaps',
data.swap.map((d) => ({
slot: d.block.number,
transaction_index: d.rawInstruction.transactionIndex,
instruction_address: d.rawInstruction.instructionAddress.join('.'),
// The Storage Write API does not parse Date/ISO strings; pass microseconds
block_timestamp: d.timestamp ? d.timestamp.getTime() * 1000 : 0,
program_id: d.programId,
})),
)
},
}),
)