Repository navigation
Expand file tree
/
Copy pathscraper_run.js
More file actions
73 lines (68 loc) · 2.77 KB
/
Copy pathscraper_run.js
File metadata and controls
73 lines (68 loc) · 2.77 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
'use strict'; /*jslint node:true es9:true*/
import axios from 'axios';
const API_URL = 'https://api.brightdata.com';
const PENDING = ['starting', 'running', 'building'];
const DISCOVER = 'discover_by_';
const enc = encodeURIComponent;
const api_error = (step, e)=>{
const body = e.response?.data;
const msg = typeof body=='string' ? body : JSON.stringify(body??'');
const err = new Error(`${step} failed`
+(e.response ? ` (HTTP ${e.response.status}): ${msg.slice(0, 500)}`
: `: ${e.message}`));
err.step = step;
err.error_code = e.response?.status??e.code;
return err;
};
const call = async(step, req)=>{
try {
return await axios({timeout: 30000, ...req});
} catch(e){
throw api_error(step, e);
}
};
const strip_nulls = data=>JSON.parse(JSON.stringify(data,
(_k, v)=>v==null ? undefined : v));
export function create_scraper_run(opt = {}){
const api_url = opt.api_url||API_URL;
const poll_ms = opt.poll_ms||2000;
const trigger = async({dataset_id, method, input, limit_per_input = 10,
headers})=>
{
const discover = method.startsWith(DISCOVER) && {type: 'discover_new',
discover_by: method.slice(DISCOVER.length), limit_per_input};
const res = await call('trigger', {method: 'POST',
url: `${api_url}/datasets/v3/trigger`, headers,
params: {dataset_id, include_errors: true, ...discover},
data: [].concat(input)});
if (!res.data?.snapshot_id)
throw api_error('trigger', new Error('no snapshot_id returned'));
return res.data.snapshot_id;
};
const progress = async(snapshot_id, headers)=>{
const res = await call('progress', {headers,
url: `${api_url}/datasets/v3/progress/${enc(snapshot_id)}`});
return res.data;
};
const results = async(snapshot_id, headers)=>{
const res = await call('results', {headers, params: {format: 'json'},
url: `${api_url}/datasets/v3/snapshot/${enc(snapshot_id)}`});
if (res.status==202 || PENDING.includes(res.data?.status))
return {snapshot_id, status: 'running'};
return {snapshot_id, status: 'ready', data: strip_nulls(res.data)};
};
const run = async({wait_ms = 45000, headers, ...req})=>{
const deadline = Date.now()+wait_ms;
const h = {...headers, 'x-mcp-dataset-id': req.dataset_id,
'x-mcp-method': req.method};
const snapshot_id = await trigger({...req, headers: h});
for (;;)
{
const res = await results(snapshot_id, h);
if (res.status=='ready' || Date.now()+poll_ms>deadline)
return res;
await new Promise(done=>setTimeout(done, poll_ms));
}
};
return {trigger, progress, results, run};
}