From bdf4437ca2d9cf1090d51f342a48cdf978f25a9a Mon Sep 17 00:00:00 2001 From: gojomo Date: Thu, 19 Nov 2009 22:45:07 +0000 Subject: [PATCH] [HER-1546] Springify(5): Update checkpointing to work smoothly with spring-configured crawls core checkpointing system classes: * Checkpointable.java basic interface for checkpoint aware/capable beans * Checkpoint.java one Checkpoint, either in progress or recovery -- a name + storage directory * CheckpointService.java helper bean to enable, trigger checkpoints and recovery * Checkpointer.java(*2), CheckpointInputStream.java, CheckpointRecovery.java DefaultCheckpointRecovery.java, RecoverAction.java (deleted) no longer used --- .../checkpointing/CheckpointInputStream.java | 70 ---- .../checkpointing/CheckpointRecovery.java | 40 --- .../archive/checkpointing/Checkpointable.java | 52 ++- .../archive/checkpointing/Checkpointer.java | 120 ------- .../DefaultCheckpointRecovery.java | 109 ------ .../archive/checkpointing/RecoverAction.java | 38 -- .../archive/crawler/framework/Checkpoint.java | 148 ++++++++ .../crawler/framework/CheckpointService.java | 293 +++++++++++++++ .../crawler/framework/Checkpointer.java | 336 ------------------ 9 files changed, 487 insertions(+), 719 deletions(-) delete mode 100644 commons/src/main/java/org/archive/checkpointing/CheckpointInputStream.java delete mode 100644 commons/src/main/java/org/archive/checkpointing/CheckpointRecovery.java delete mode 100644 commons/src/main/java/org/archive/checkpointing/Checkpointer.java delete mode 100644 commons/src/main/java/org/archive/checkpointing/DefaultCheckpointRecovery.java delete mode 100644 commons/src/main/java/org/archive/checkpointing/RecoverAction.java create mode 100644 engine/src/main/java/org/archive/crawler/framework/Checkpoint.java create mode 100644 engine/src/main/java/org/archive/crawler/framework/CheckpointService.java delete mode 100644 engine/src/main/java/org/archive/crawler/framework/Checkpointer.java diff --git a/commons/src/main/java/org/archive/checkpointing/CheckpointInputStream.java b/commons/src/main/java/org/archive/checkpointing/CheckpointInputStream.java deleted file mode 100644 index f8b31e87..00000000 --- a/commons/src/main/java/org/archive/checkpointing/CheckpointInputStream.java +++ /dev/null @@ -1,70 +0,0 @@ -/* - * This file is part of the Heritrix web crawler (crawler.archive.org). - * - * Licensed to the Internet Archive (IA) by one or more individual - * contributors. - * - * The IA licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.archive.checkpointing; - -import java.io.IOException; -import java.io.InputStream; -import java.io.ObjectInputStream; -import java.net.URI; - - -/** - * Object input stream that provides information useful during checkpoint - * recovery. - * - * @author pjack - */ -public class CheckpointInputStream extends ObjectInputStream -implements CheckpointRecovery { - - - final private CheckpointRecovery recovery; - - - public CheckpointInputStream(InputStream input, - CheckpointRecovery recovery) throws IOException { - super(input); - this.recovery = recovery; - } - - - public String getRecoveredJobName() { - return recovery.getRecoveredJobName(); - } - -// public void setState(Object module, Key key, T value) { -// recovery.setState(module, key, value); -// } - - - public String translatePath(String path) { - return recovery.translatePath(path); - } - - - public URI translateURI(URI uri) { - return recovery.translateURI(uri); - } - - -// public void apply(SingleSheet global) { -// throw new UnsupportedOperationException(); -// } -} diff --git a/commons/src/main/java/org/archive/checkpointing/CheckpointRecovery.java b/commons/src/main/java/org/archive/checkpointing/CheckpointRecovery.java deleted file mode 100644 index c202fd71..00000000 --- a/commons/src/main/java/org/archive/checkpointing/CheckpointRecovery.java +++ /dev/null @@ -1,40 +0,0 @@ -/* - * This file is part of the Heritrix web crawler (crawler.archive.org). - * - * Licensed to the Internet Archive (IA) by one or more individual - * contributors. - * - * The IA licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.archive.checkpointing; - -import java.net.URI; - -/** - * The information about a checkpoint recovery. - * - * @author pjack - */ -public interface CheckpointRecovery { - - String getRecoveredJobName(); - - String translatePath(String path); - - URI translateURI(URI uri); - -// void setState(Object module, Key key, T value); -// -// void apply(SingleSheet global); -} diff --git a/commons/src/main/java/org/archive/checkpointing/Checkpointable.java b/commons/src/main/java/org/archive/checkpointing/Checkpointable.java index 4a4dae84..77908b26 100644 --- a/commons/src/main/java/org/archive/checkpointing/Checkpointable.java +++ b/commons/src/main/java/org/archive/checkpointing/Checkpointable.java @@ -19,19 +19,59 @@ package org.archive.checkpointing; -import java.io.File; import java.io.IOException; -import java.util.List; +import org.archive.crawler.framework.Checkpoint; /** - * @author pjack - * + * Interface for objects that can checkpoint their state, possibly + * but not necessarily into the provided Checkpoint instance, on + * request. + * + * @contributor pjack + * @contributor gojomo */ public interface Checkpointable { + /** + * Note a checkpoint is about to begin. Most beans will ignore, + * but some can use this to wrap up tasks that shouldn't be + * half-done during a checkpoint (and hold a lock to prevent + * tasks from beginning during the checkpoint). + * + * @param checkpointInProgress Checkpoint + */ + void startCheckpoint(Checkpoint checkpointInProgress); - void checkpoint(File dir, List actions) throws IOException; + /** + * Do the actual checkpoint. Beans should ensure any state that + * they would need to recover gets saved in an appropriate place. + * A moderate amount of state may be saved as a JSONObject into + * the Checkpoint (which then keeps it in the checkpoint directory.) + * Larger amounts of state may be stored in a manner private to + * the bean, or via other collaborating beans (such as BdbModule) + * which checkpoint backgin database state. + * + * @param checkpointInProgress Checkpoint + * @throws IOException + */ + void doCheckpoint(Checkpoint checkpointInProgress) throws IOException; - + // + /** + * Cleanup/unlock; need not complete for a checkpoint to be valid. + * + * @param checkpointInProgress Checkpoint + */ + void finishCheckpoint(Checkpoint checkpointInProgress); + + /** + * Used to inform a bean that it should restore its state from + * the given Checkpoint when launched (Lifecycle start()). May be + * autowired or configured after build, but before launch, for + * example via CheckpointService.setRecoveryCheckpointByName(). + * + * @param recoveryCheckpoint Checkpoint + */ + void setRecoveryCheckpoint(Checkpoint recoveryCheckpoint); } diff --git a/commons/src/main/java/org/archive/checkpointing/Checkpointer.java b/commons/src/main/java/org/archive/checkpointing/Checkpointer.java deleted file mode 100644 index 4e2b999c..00000000 --- a/commons/src/main/java/org/archive/checkpointing/Checkpointer.java +++ /dev/null @@ -1,120 +0,0 @@ -/* - * This file is part of the Heritrix web crawler (crawler.archive.org). - * - * Licensed to the Internet Archive (IA) by one or more individual - * contributors. - * - * The IA licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.archive.checkpointing; - -import java.io.File; -import java.io.FileInputStream; -import java.io.FileOutputStream; -import java.io.IOException; -import java.io.ObjectInputStream; -import java.io.ObjectOutputStream; -import java.util.ArrayList; -import java.util.List; - -import org.archive.util.IoUtils; - -/** - * Executes checkpoints and recovers. - * - * @author pjack - */ -public class Checkpointer { - - final public static String ACTIONS_FILE = "actions.serialized"; - - final public static String OBJECT_GRAPH_FILE = "object_graph.serialized"; - - private Checkpointer() { - } - - - public static void checkpoint(/*SheetManager*/Object mgr, File dir) - throws IOException { - List actions = new ArrayList(); -// for (Checkpointable c: mgr.getCheckpointables()) { -// c.checkpoint(dir, actions); -// } - - writeObject(new File(dir, ACTIONS_FILE), actions); - writeObject(new File(dir, OBJECT_GRAPH_FILE), mgr); - } - - - private static void writeObject(File f, Object o) - throws IOException { - ObjectOutputStream oout = null; - try { - oout = new ObjectOutputStream(new FileOutputStream(f)); - oout.writeObject(o); - } finally { - IoUtils.close(oout); - } - - } - - - @SuppressWarnings("unused") - private static List readActions(File dir) - throws IOException { - File actionsFile = new File(dir, ACTIONS_FILE); - ObjectInputStream oinp = null; - try { - oinp = new ObjectInputStream( - new FileInputStream(actionsFile)); - @SuppressWarnings("unchecked") - List actions = (List)oinp.readObject(); - return actions; - } catch (ClassNotFoundException e) { - IOException io = new IOException(); - io.initCause(e); - throw io; - } finally { - IoUtils.close(oinp); - } - } - -// public static SheetManager recover(File dir, CheckpointRecovery recovery) -// throws IOException { -// List actions = readActions(dir); -// for (RecoverAction action: actions) try { -// action.recoverFrom(dir, recovery); -// } catch (Exception e) { -// IOException io = new IOException(); -// io.initCause(e); -// throw io; -// } -// -// CheckpointInputStream cinp = null; -// try { -// File f = new File(dir, OBJECT_GRAPH_FILE); -// cinp = new CheckpointInputStream(new FileInputStream(f), recovery); -// SheetManager mgr = (SheetManager)cinp.readObject(); -// recovery.apply(mgr.getGlobalSheet()); -// return mgr; -// } catch (ClassNotFoundException e) { -// IOException io = new IOException(); -// io.initCause(e); -// throw io; -// }finally { -// IoUtils.close(cinp); -// } -// } - -} diff --git a/commons/src/main/java/org/archive/checkpointing/DefaultCheckpointRecovery.java b/commons/src/main/java/org/archive/checkpointing/DefaultCheckpointRecovery.java deleted file mode 100644 index a92e2b25..00000000 --- a/commons/src/main/java/org/archive/checkpointing/DefaultCheckpointRecovery.java +++ /dev/null @@ -1,109 +0,0 @@ -/* - * This file is part of the Heritrix web crawler (crawler.archive.org). - * - * Licensed to the Internet Archive (IA) by one or more individual - * contributors. - * - * The IA licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.archive.checkpointing; - -import java.net.URI; -import java.util.HashMap; -import java.util.Map; - -/** - * Default implementation. - * - * @author pjack - * - */ -public class DefaultCheckpointRecovery implements CheckpointRecovery { - - -// final private Map> newSettings = -// new IdentityHashMap>(); - - final private Map uriTranslations = new HashMap(); - - final private Map fileTranslations = - new HashMap(); - - final private String name; - - public DefaultCheckpointRecovery(String name) { - this.name = name; - } - - public String getRecoveredJobName() { - return name; - } - - public Map getFileTranslations() { - return fileTranslations; - } - - public Map getURITranslations() { - return uriTranslations; - } - - -// public void setState(Object module, Key key, T value) { -// Map map = newSettings.get(module); -// if (map == null) { -// map = new HashMap(); -// newSettings.put(module, map); -// } -// -// map.put(key, value); -// } - - - public String translatePath(String path) { - Map.Entry match = null; - for (Map.Entry me: fileTranslations.entrySet()) { - if (path.startsWith(me.getKey())) { - if ((match == null) - || (match.getKey().length() < me.getKey().length())) { - match = me; - } - } - } - - if (match == null) { - return path; - } - - int size = match.getKey().length(); - return match.getValue() + path.substring(size); - } - - - public URI translateURI(URI uri) { - URI r = uriTranslations.get(uri); - return r == null ? uri : r; - } - - -// public void apply(SingleSheet global) { -// for (Map.Entry> mod: newSettings.entrySet()) { -// Object module = mod.getKey(); -// for (Map.Entry me: mod.getValue().entrySet()) { -// @SuppressWarnings("unchecked") -// Key k = me.getKey(); -// global.set(module, k, me.getValue()); -// } -// } -// } -} diff --git a/commons/src/main/java/org/archive/checkpointing/RecoverAction.java b/commons/src/main/java/org/archive/checkpointing/RecoverAction.java deleted file mode 100644 index 5ad60aac..00000000 --- a/commons/src/main/java/org/archive/checkpointing/RecoverAction.java +++ /dev/null @@ -1,38 +0,0 @@ -/* - * This file is part of the Heritrix web crawler (crawler.archive.org). - * - * Licensed to the Internet Archive (IA) by one or more individual - * contributors. - * - * The IA licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.archive.checkpointing; - -import java.io.File; -import java.io.Serializable; - - - -/** - * An action to be taken before a checkpoint's object graph is deserialized. - * - * @author pjack - */ -public interface RecoverAction extends Serializable { - - - void recoverFrom(File checkpointDir, CheckpointRecovery recovery) - throws Exception; - -} diff --git a/engine/src/main/java/org/archive/crawler/framework/Checkpoint.java b/engine/src/main/java/org/archive/crawler/framework/Checkpoint.java new file mode 100644 index 00000000..19b207b7 --- /dev/null +++ b/engine/src/main/java/org/archive/crawler/framework/Checkpoint.java @@ -0,0 +1,148 @@ +/* + * This file is part of the Heritrix web crawler (crawler.archive.org). + * + * Licensed to the Internet Archive (IA) by one or more individual + * contributors. + * + * The IA licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.archive.crawler.framework; + +import java.io.File; +import java.io.IOException; +import java.util.logging.Level; +import java.util.logging.Logger; + +import org.apache.commons.io.FileUtils; +import org.archive.spring.ConfigPath; +import org.archive.util.ArchiveUtils; +import org.json.JSONException; +import org.json.JSONObject; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Required; + + +/** + * Represents a single checkpoint, by its name and main store directory. + * + * @contributor gojomo + */ +public class Checkpoint implements InitializingBean { + private final static Logger LOGGER = + Logger.getLogger(Checkpoint.class.getName()); + + String name; + String shortName; + boolean success = false; + + /** + * Checkpoints directory + */ + protected ConfigPath checkpointDir = + new ConfigPath("checkpoint directory",""); + public ConfigPath getCheckpointDir() { + return checkpointDir; + } + @Required // if used in Spring context + public void setCheckpointDir(ConfigPath checkpointsDir) { + this.checkpointDir = checkpointsDir; + } + + protected CheckpointService checkpointService; + public CheckpointService getCheckpointService() { + return this.checkpointService; + } + @Autowired + public void setCheckpointService(CheckpointService checkpointer) { + this.checkpointService = checkpointer; + } + + public Checkpoint() { + } + + /** + * Use immediately after instantiation to fill-in a Checkpoint + * created outside Spring configuration. + * + * @param nextCheckpointNumber + * @param checkpointsDir + */ + public void generateFrom(CheckpointService checkpointer) { + getCheckpointDir().setBase(checkpointer.getCheckpointsDir()); + getCheckpointDir().setPath( + "cp" + + CheckpointService.INDEX_FORMAT.format(checkpointer.getNextCheckpointNumber()) + + "-" + + ArchiveUtils.get14DigitDate()); + getCheckpointDir().getFile().mkdirs(); + afterPropertiesSet(); + } + + public void afterPropertiesSet() { + if(checkpointDir.getBase()==null) { + // if not otherwise set, adopt base from Checkpointer + checkpointDir.setBase(checkpointService.getCheckpointsDir()); + } + name = checkpointDir.getFile().getName(); + shortName = name.substring(name.indexOf("-")); + } + + public void setSuccess(boolean b) { + success = b; + } + public boolean getSuccess() { + return success; + } + + public String getName() { + return name; + } + public String getShortName() { + return shortName; + } + protected void writeValidity(String stamp) { + if(!success) { + return; + } + File valid = new File(checkpointDir.getFile(), CheckpointService.VALIDITY_STAMP_FILENAME); + try { + FileUtils.writeStringToFile(valid, stamp); + } catch (IOException e) { + valid.delete(); + } + } + + public void saveJson(String beanName, JSONObject json) { + try { + File targetFile = new File(getCheckpointDir().getFile(),beanName); + FileUtils.writeStringToFile( + targetFile, + json.toString()); + } catch (IOException e) { + LOGGER.log(Level.SEVERE,"unable to save checkpoint state of "+beanName,e); + setSuccess(false); + } + } + + public JSONObject loadJson(String beanName) { + File sourceFile = new File(getCheckpointDir().getFile(),beanName); + try { + return new JSONObject(FileUtils.readFileToString(sourceFile)); + } catch (JSONException e) { + throw new RuntimeException(e); + } catch (IOException e) { + throw new RuntimeException(e); + } + } +} \ No newline at end of file diff --git a/engine/src/main/java/org/archive/crawler/framework/CheckpointService.java b/engine/src/main/java/org/archive/crawler/framework/CheckpointService.java new file mode 100644 index 00000000..4629f526 --- /dev/null +++ b/engine/src/main/java/org/archive/crawler/framework/CheckpointService.java @@ -0,0 +1,293 @@ +/* + * This file is part of the Heritrix web crawler (crawler.archive.org). + * + * Licensed to the Internet Archive (IA) by one or more individual + * contributors. + * + * The IA licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.archive.crawler.framework; + +import java.io.File; +import java.io.FileFilter; +import java.text.DecimalFormat; +import java.util.Map; +import java.util.Timer; +import java.util.TimerTask; +import java.util.logging.Level; +import java.util.logging.Logger; + +import org.apache.commons.io.filefilter.FileFilterUtils; +import org.apache.commons.lang.ArrayUtils; +import org.archive.checkpointing.Checkpointable; +import org.archive.crawler.reporting.CrawlStatSnapshot; +import org.archive.spring.ConfigPath; +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; +import org.springframework.context.Lifecycle; +import org.springframework.context.support.AbstractApplicationContext; + +/** + * Executes checkpoints, and offers convenience methods for enumerating + * available Checkpoints and injecting a recovery-Checkpoint after + * build and before launch (setRecoveryCheckpointByName). + * + * Offers optional automatic checkpointing at a configurable interval + * in minutes. + * + * @contributor stack + * @contributor gojomo + * @contributor pjack + */ +public class CheckpointService implements Lifecycle, ApplicationContextAware { + private final static Logger LOGGER = + Logger.getLogger(CheckpointService.class.getName()); + + /** Name of file written with timestamp into valid checkpoints */ + public static final String VALIDITY_STAMP_FILENAME = "valid"; + + /** format for serial numbers */ + public static final DecimalFormat INDEX_FORMAT = new DecimalFormat("00000"); + + /** Next overall series checkpoint number */ + protected int nextCheckpointNumber = 1; + + Checkpoint checkpointInProgress; + + CrawlStatSnapshot lastCheckpointSnapshot = null; + + /**Setup in constructor or on a call to recovery */ + protected transient Timer timerThread = null; + + /** + * Checkpoints directory + */ + protected ConfigPath checkpointsDir = + new ConfigPath("checkpoints subdirectory","checkpoints"); + public ConfigPath getCheckpointsDir() { + return checkpointsDir; + } + public void setCheckpointsDir(ConfigPath checkpointsDir) { + this.checkpointsDir = checkpointsDir; + } + + /** + * Period at which to create automatic checkpoints; -1 means + * no auto checkpointing. + */ + int checkpointIntervalMinutes = -1; + public int getCheckpointIntervalMinutes() { + return checkpointIntervalMinutes; + } + public void setCheckpointIntervalMinutes(int checkpointIntervalMinutes) { + this.checkpointIntervalMinutes = checkpointIntervalMinutes; + } + + Checkpoint recoveryCheckpoint; + @Autowired(required=false) + public void setRecoveryCheckpoint(Checkpoint checkpoint) { + this.recoveryCheckpoint = checkpoint; + } + public Checkpoint getRecoveryCheckpoint() { + return this.recoveryCheckpoint; + } + + protected CrawlController controller; + public CrawlController getCrawlController() { + return this.controller; + } + @Autowired + public void setCrawlController(CrawlController controller) { + this.controller = controller; + } + + // ApplicationContextAware implementation, for eventing + AbstractApplicationContext appCtx; + public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { + this.appCtx = (AbstractApplicationContext)applicationContext; + } + + /** + * Create a new Checkpointer + */ + public CheckpointService() { + } + + public void start() { + if (isRunning) { + return; + } + this.isRunning = true; + // Convert period from hours to milliseconds. + long periodMs = getCheckpointIntervalMinutes() * (60 * 1000); + if(periodMs<=0) { + return; + } + TimerTask tt = new TimerTask() { + public void run() { + if (isCheckpointing()) { + LOGGER.info("CheckpointTimerThread skipping checkpoint, " + + "already checkpointing: State: " + + controller.getState()); + return; + } + LOGGER.info("TimerThread request checkpoint"); + requestCrawlCheckpoint(); + } + }; + this.timerThread = new Timer(true); + this.timerThread.schedule(tt, periodMs, periodMs); + LOGGER.info("Installed Checkpoint TimerThread to checkpoint every " + + periodMs + " milliseconds."); + } + + boolean isRunning = false; + public boolean isRunning() { + return isRunning; + } + + + public void stop() { + if (this.timerThread != null) { + LOGGER.info("Cleaned up Checkpoint TimerThread."); + this.timerThread.cancel(); + this.timerThread = null; + } + this.isRunning = false; + } + + /** + * @return Returns the nextCheckpoint index. + */ + public int getNextCheckpointNumber() { + return this.nextCheckpointNumber; + } + + /** + * Run a checkpoint of the crawler + */ + public synchronized String requestCrawlCheckpoint() throws IllegalStateException { + if (isCheckpointing()) { + throw new IllegalStateException("Checkpoint already running."); + } + + // prevent redundant auto-checkpoints when crawler paused + if(controller.isPaused()) { + if (controller.getStatisticsTracker().getSnapshot().sameProgressAs(lastCheckpointSnapshot)) { + LOGGER.info("no progress since last checkpoint; ignoring"); + System.err.println("no progress since last checkpoint; ignoring"); + return null; + } + } + + checkpointInProgress = new Checkpoint(); + checkpointInProgress.generateFrom(this); + + @SuppressWarnings("unchecked") + Map toCheckpoint = appCtx.getBeansOfType(Checkpointable.class); + + try { + // pre (incl. acquire necessary locks) + long startMs = System.currentTimeMillis(); + for(Checkpointable c : toCheckpoint.values()) { + c.startCheckpoint(checkpointInProgress); + } + long duration = System.currentTimeMillis() - startMs; + System.err.println("all startCheckpoint() completed in "+duration+"ms"); + + // flush/write + for(Checkpointable c : toCheckpoint.values()) { + long doMs = System.currentTimeMillis(); + c.doCheckpoint(checkpointInProgress); + long doDuration = System.currentTimeMillis() - doMs; + System.err.println("doCheckpoint() "+c+" in "+doDuration+"ms"); + } + checkpointInProgress.setSuccess(true); + + } catch (Exception e) { + checkpointFailed(e); + } finally { + checkpointInProgress.writeValidity( + controller.getStatisticsTracker().getProgressStamp()); + lastCheckpointSnapshot = controller.getStatisticsTracker().getSnapshot(); + // close (incl. release locks) + for(Checkpointable c : toCheckpoint.values()) { + c.finishCheckpoint(checkpointInProgress); + } + } + + this.nextCheckpointNumber++; + LOGGER.info("finished checkpoint "+checkpointInProgress.getName()); + String nameToReport = checkpointInProgress.getSuccess() ? checkpointInProgress.getName() : null; + this.checkpointInProgress = null; + return nameToReport; + } + + + /** + * @return True if a checkpoint is in progress. + */ + public boolean isCheckpointing() { + return this.checkpointInProgress != null; + } + + /** + * Note that a checkpoint failed + * + * @param e Exception checkpoint failed on. + */ + protected void checkpointFailed(Exception e) { + LOGGER.log(Level.WARNING, " Checkpoint failed", e); + } + + protected void checkpointFailed(final String message) { + LOGGER.warning(message); + } + + public boolean hasAvailableCheckpoints() { + if(getRecoveryCheckpoint()!=null || isRunning()) { + return false; + } + return (getAvailableCheckpointDirectories() != null + && getAvailableCheckpointDirectories().length > 0); + } + + public File[] getAvailableCheckpointDirectories() { + File[] dirs = getCheckpointsDir().getFile().listFiles((FileFilter)FileFilterUtils.directoryFileFilter()); + ArrayUtils.reverse(dirs); + return dirs; + } + + /** + * Given the name of a valid checkpoint subdirectory in the checkpoints + * directory, create a Checkpoint instance, and insert it into all + * Checkpointable beans. + * + * @param selectedCheckpoint + */ + @SuppressWarnings("unchecked") + public void setRecoveryCheckpointByName(String selectedCheckpoint) { + Checkpoint recoveryCheckpoint = new Checkpoint(); + recoveryCheckpoint.getCheckpointDir().setPath(selectedCheckpoint); + recoveryCheckpoint.setCheckpointService(this); + recoveryCheckpoint.afterPropertiesSet(); + setRecoveryCheckpoint(recoveryCheckpoint); + Map toSetRecovery = appCtx.getBeansOfType(Checkpointable.class); + + for(Checkpointable c : toSetRecovery.values()) { + c.setRecoveryCheckpoint(recoveryCheckpoint); + } + } +} //EOC diff --git a/engine/src/main/java/org/archive/crawler/framework/Checkpointer.java b/engine/src/main/java/org/archive/crawler/framework/Checkpointer.java deleted file mode 100644 index 594a2cdb..00000000 --- a/engine/src/main/java/org/archive/crawler/framework/Checkpointer.java +++ /dev/null @@ -1,336 +0,0 @@ -/* - * This file is part of the Heritrix web crawler (crawler.archive.org). - * - * Licensed to the Internet Archive (IA) by one or more individual - * contributors. - * - * The IA licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.archive.crawler.framework; - -import java.io.File; -import java.io.FileOutputStream; -import java.io.IOException; -import java.io.Serializable; -import java.text.DecimalFormat; -import java.util.Timer; -import java.util.TimerTask; -import java.util.logging.Level; -import java.util.logging.Logger; - -import org.archive.util.ArchiveUtils; - -/** - * Runs checkpointing. - * Also keeps history of crawl checkpoints Generally used by CrawlController - * only but also has static utility methods classes that need to participate in - * a checkpoint can use. - * - * @author gojomo - * @author stack - */ -public class Checkpointer implements Serializable { - private static final long serialVersionUID = 7610078446694353173L; - - /** - * Name of file written with timestamp into valid checkpoints. - */ - public static final String VALIDITY_STAMP_FILENAME = "valid"; - - private final static Logger LOGGER = - Logger.getLogger(Checkpointer.class.getName()); - - private static final String DEFAULT_PREFIX = ""; - - /** - * String to prefix any new checkpoint names. - */ - private String checkpointPrefix = DEFAULT_PREFIX; - - /** - * Next overall series checkpoint number. - */ - private int nextCheckpoint = 1; - - /** - * If a checkpoint has begun, its directory under - * checkpointDirectory. - */ - private transient File checkpointInProgressDir = null; - - /** - * If the checkpoint in progress has encountered fatal errors. - */ - private transient boolean checkpointErrors = false; - - /** - * checkpointThread is set if a checkpoint is currently running. - */ - private transient Thread checkpointThread = null; - - private transient CrawlController controller; - - /** - * Setup in constructor or on a call to revovery. - */ - private transient Timer timerThread = null; - - public static final DecimalFormat INDEX_FORMAT = new DecimalFormat("00000"); - - /** - * Create a new CheckpointContext with the given store directory - * @param cc CrawlController instance thats hosting this Checkpointer. - * @param checkpointDir Where to store checkpoint. - */ - public Checkpointer(final CrawlController cc, final File checkpointDir) { - this(cc, DEFAULT_PREFIX); - } - - /** - * Create a new CheckpointContext with the given store directory - * - * @param cc CrawlController instance thats hosting this Checkpointer. - * @param prefix Prefix for checkpoint label. - */ - public Checkpointer(final CrawlController cc, final String prefix) { - super(); - initialize(cc, prefix); - - } - - protected void initialize(final CrawlController cc, final String prefix) { - this.controller = cc; - this.checkpointPrefix = prefix; - // Period is in hours. - int period = cc.getCheckpointerPeriod(); - if (period <= 0) { - return; - } - // Convert period from hours to milliseconds. - long periodMs = period * (1000 * 60 * 60); - TimerTask tt = new TimerTask() { - private CrawlController cController = cc; - public void run() { - if (isCheckpointing()) { - LOGGER.info("CheckpointTimerThread skipping checkpoint, " + - "already checkpointing: State: " + - this.cController.getState()); - return; - } - LOGGER.info("TimerThread request checkpoint"); - this.cController.requestCrawlCheckpoint(); - } - }; - this.timerThread = new Timer(true); - this.timerThread.schedule(tt, periodMs, periodMs); - LOGGER.info("Installed Checkpoint TimerThread to checkpoint every " + - period + " hour(s)."); - } - - void cleanup() { - if (this.timerThread != null) { - LOGGER.info("Cleanedup Checkpoint TimerThread."); - this.timerThread.cancel(); - } - } - - /** - * @return Returns the nextCheckpoint index. - */ - public int getNextCheckpoint() { - return this.nextCheckpoint; - } - - /** - * Run a checkpoint of the crawler. - */ - public void checkpoint() { - String name = "Checkpoint-" + getNextCheckpointName(); - this.checkpointThread = new CheckpointingThread(name); - this.checkpointThread.setDaemon(true); - this.checkpointThread.start(); - } - - /** - * Thread to run the checkpointing. - * @author stack - */ - public class CheckpointingThread extends Thread { - public CheckpointingThread(final String name) { - super(name); - } - - public CrawlController getController() { - return Checkpointer.this.controller; - } - - public void run() { - LOGGER.info("Started"); - // If crawler already paused, don't resume crawling after - // finishing checkpointing. - final boolean alreadyPaused = getController().isPaused() || - getController().isPausing(); - try { - getController().requestCrawlPause(); - // Clear any checkpoint errors. - setCheckpointErrors(false); - if (!waitOnPaused()) { - checkpointFailed("Failed wait for complete pause."); - } else { - createCheckpointInProgressDirectory(); - org.archive.checkpointing.Checkpointer.checkpoint( - //TODO:SPRINGY - null, //getController().getSheetManager(), - checkpointInProgressDir); - } - } catch (Exception e) { - checkpointFailed(e); - } finally { - if (!isCheckpointErrors()) { - writeValidity(); - } - Checkpointer.this.nextCheckpoint++; - clearCheckpointInProgressDirectory(); - LOGGER.info("Finished"); - getController().completePause(); - if (!alreadyPaused) { - getController().requestCrawlResume(); - } - } - } - - private synchronized boolean waitOnPaused() { - // If we're paused we can exit but also exit if the crawl has been - // resumed by the operator. - while(!getController().isPaused() && !getController().isStateRunning()) { - try { - wait(1000 * 3); - } catch (InterruptedException e) { - // May be for us. - } - } - return getController().isPaused(); - } - } - - protected File createCheckpointInProgressDirectory() { - this.checkpointInProgressDir = - new File(Checkpointer.this.controller.getCheckpointsDir().getFile(), - getNextCheckpointName()); - this.checkpointInProgressDir.mkdirs(); - return this.checkpointInProgressDir; - } - - protected void clearCheckpointInProgressDirectory() { - this.checkpointInProgressDir = null; - } - - protected CrawlController getController() { - return this.controller; - } - - /** - * @return next checkpoint name (zero-padding string). - */ - public String getNextCheckpointName() { - return formatCheckpointName(this.checkpointPrefix, this.nextCheckpoint); - } - - public static String formatCheckpointName(final String prefix, - final int index) { - return prefix + INDEX_FORMAT.format(index); - } - - protected void writeValidity() { - File valid = new File(this.checkpointInProgressDir, - VALIDITY_STAMP_FILENAME); - try { - FileOutputStream fos = new FileOutputStream(valid); - fos.write(ArchiveUtils.get14DigitDate().getBytes()); - fos.close(); - } catch (IOException e) { - valid.delete(); - } - } - - /** - * @return Checkpoint directory. Name of the directory is the name of this - * current checkpoint. Null if no checkpoint in progress. - */ - public File getCheckpointInProgressDirectory() { - return this.checkpointInProgressDir; - } - - /** - * @return True if a checkpoint is in progress. - */ - public boolean isCheckpointing() { - return this.checkpointThread != null && this.checkpointThread.isAlive(); - } - - /** - * Note that a checkpoint failed - * - * @param e Exception checkpoint failed on. - */ - protected void checkpointFailed(Exception e) { - LOGGER.log(Level.WARNING, " Checkpoint failed", e); - checkpointFailed(); - } - - protected void checkpointFailed(final String message) { - LOGGER.warning(message); - checkpointFailed(); - } - - protected void checkpointFailed() { - this.checkpointErrors = true; - } - - /** - * @return True if current/last checkpoint failed. - */ - public boolean isCheckpointFailed() { - return this.checkpointErrors; - } - - /** - * @return Return whether this context is at a new crawl, never- - * checkpointed state. - */ - public boolean isAtBeginning() { - return nextCheckpoint == 1; - } - - /** - * Call when recovering from a checkpoint. - * Call this after instance has been revivifyied post-serialization to - * amend counters and directories that effect where checkpoints get stored - * from here on out. - * @param cc CrawlController instance. - */ - public void recover(final CrawlController cc) { - // Prepend the checkpoint name with a little 'r' so we tell apart - // checkpoints made from a recovery. Allow for there being - // multiple 'r' prefixes. - initialize(cc, 'r' + this.checkpointPrefix); - } - - protected boolean isCheckpointErrors() { - return this.checkpointErrors; - } - - protected void setCheckpointErrors(boolean checkpointErrors) { - this.checkpointErrors = checkpointErrors; - } -}