clusternode.java

来自「jsr170接口的java实现。是个apache的开源项目。」· Java 代码 · 共 1,278 行 · 第 1/3 页

JAVA
1,278
字号
            record.writeString(null);            write(record, ntDefs, true);            record.writeChar('\0');            record.update();            setRevision(record.getRevision());            succeeded = true;        } catch (JournalException e) {            String msg = "Unable to create log entry: " + e.getMessage();            log.error(msg);        } catch (Throwable e) {            String msg = "Unexpected error while creating log entry.";            log.error(msg, e);        } finally {            if (!succeeded && record != null) {                record.cancelUpdate();            }        }    }    /**     * {@inheritDoc}     */    public void reregistered(NodeTypeDef ntDef) {        if (status != STARTED) {            log.info("not started: nodetype operation ignored.");            return;        }        Record record = null;        boolean succeeded = false;        try {            record = journal.getProducer(PRODUCER_ID).append();            record.writeString(null);            write(record, ntDef);            record.writeChar('\0');            record.update();            setRevision(record.getRevision());            succeeded = true;        } catch (JournalException e) {            String msg = "Unable to create log entry: " + e.getMessage();            log.error(msg);        } catch (Throwable e) {            String msg = "Unexpected error while creating log entry.";            log.error(msg, e);        } finally {            if (!succeeded && record != null) {                record.cancelUpdate();            }        }    }    /**     * {@inheritDoc}     */    public void unregistered(Collection qnames) {        if (status != STARTED) {            log.info("not started: nodetype operation ignored.");            return;        }        Record record = null;        boolean succeeded = false;        try {            record = journal.getProducer(PRODUCER_ID).append();            record.writeString(null);            write(record, qnames, false);            record.writeChar('\0');            record.update();            setRevision(record.getRevision());            succeeded = true;        } catch (JournalException e) {            String msg = "Unable to create log entry: " + e.getMessage();            log.error(msg);        } catch (Throwable e) {            String msg = "Unexpected error while creating log entry.";            log.error(msg, e);        } finally {            if (!succeeded && record != null) {                record.cancelUpdate();            }        }    }    /**     * {@inheritDoc}     */    public void setListener(NodeTypeEventListener listener) {        nodeTypeListener = listener;    }    /**     * Workspace update channel.     */    class WorkspaceUpdateChannel implements UpdateEventChannel {        /**         * Attribute name used to store record.         */        private static final String ATTRIBUTE_RECORD = "record";        /**         * Workspace name.         */        private final String workspace;        /**         * Create a new instance of this class.         *         * @param workspace workspace name         */        public WorkspaceUpdateChannel(String workspace) {            this.workspace = workspace;        }        /**         * {@inheritDoc}         */        public void updateCreated(Update update) {            if (status != STARTED) {                log.info("not started: update create ignored.");                return;            }            try {                Record record = journal.getProducer(PRODUCER_ID).append();                update.setAttribute(ATTRIBUTE_RECORD, record);            } catch (JournalException e) {                String msg = "Unable to create log entry.";                log.error(msg, e);            } catch (Throwable e) {                String msg = "Unexpected error while creating log entry.";                log.error(msg, e);            }        }        /**         * {@inheritDoc}         */        public void updatePrepared(Update update) {            if (status != STARTED) {                log.info("not started: update prepare ignored.");                return;            }            Record record = (Record) update.getAttribute(ATTRIBUTE_RECORD);            if (record == null) {                String msg = "No record created.";                log.warn(msg);                return;            }            EventStateCollection events = update.getEvents();            ChangeLog changes = update.getChanges();            boolean succeeded = false;            try {                record.writeString(workspace);                write(record, changes, events);                record.writeChar('\0');                succeeded = true;            } catch (JournalException e) {                String msg = "Unable to create log entry: " + e.getMessage();                log.error(msg);            } catch (Throwable e) {                String msg = "Unexpected error while preparing log entry.";                log.error(msg, e);            } finally {                if (!succeeded && record != null) {                    record.cancelUpdate();                    update.setAttribute(ATTRIBUTE_RECORD, null);                }            }        }        /**         * {@inheritDoc}         */        public void updateCommitted(Update update) {            if (status != STARTED) {                log.info("not started: update commit ignored.");                return;            }            Record record = (Record) update.getAttribute(ATTRIBUTE_RECORD);            if (record == null) {                String msg = "No record prepared.";                log.warn(msg);                return;            }            try {                record.update();                setRevision(record.getRevision());                log.info("Appended revision: " + record.getRevision());            } catch (JournalException e) {                String msg = "Unable to commit log entry.";                log.error(msg, e);            } catch (Throwable e) {                String msg = "Unexpected error while committing log entry.";                log.error(msg, e);            } finally {                update.setAttribute(ATTRIBUTE_RECORD, null);            }        }        /**         * {@inheritDoc}         */        public void updateCancelled(Update update) {            if (status != STARTED) {                log.info("not started: update cancel ignored.");                return;            }            Record record = (Record) update.getAttribute(ATTRIBUTE_RECORD);            if (record != null) {                record.cancelUpdate();                update.setAttribute(ATTRIBUTE_RECORD, null);            }        }        /**         * {@inheritDoc}         */        public void setListener(UpdateEventListener listener) {            if (workspace == null) {                versionUpdateListener = listener;            } else {                wspUpdateListeners.remove(workspace);                if (listener != null) {                    wspUpdateListeners.put(workspace, listener);                }            }        }    }    /**     * Workspace lock channel.     */    class WorkspaceLockChannel implements LockEventChannel {        /**         * Workspace name.         */        private final String workspace;        /**         * Create a new instance of this class.         *         * @param workspace workspace name         */        public WorkspaceLockChannel(String workspace) {            this.workspace = workspace;        }        /**         * {@inheritDoc}         */        public ClusterOperation create(NodeId nodeId, boolean deep, String owner) {            if (status != STARTED) {                log.info("not started: lock operation ignored.");                return null;            }            try {                Record record = journal.getProducer(PRODUCER_ID).append();                return new LockOperation(ClusterNode.this, workspace, record,                        nodeId, deep, owner);            } catch (JournalException e) {                String msg = "Unable to create log entry: " + e.getMessage();                log.error(msg);                return null;            } catch (Throwable e) {                String msg = "Unexpected error while creating log entry.";                log.error(msg, e);                return null;            }        }        /**         * {@inheritDoc}         */        public ClusterOperation create(NodeId nodeId) {            if (status != STARTED) {                log.info("not started: unlock operation ignored.");                return null;            }            try {                Record record = journal.getProducer(PRODUCER_ID).append();                return new LockOperation(ClusterNode.this, workspace, record,                        nodeId);            } catch (JournalException e) {                String msg = "Unable to create log entry: " + e.getMessage();                log.error(msg);                return null;            } catch (Throwable e) {                String msg = "Unexpected error while creating log entry.";                log.error(msg, e);                return null;            }        }        /**         * {@inheritDoc}         */        public void setListener(LockEventListener listener) {            wspLockListeners.remove(workspace);            if (listener != null) {                wspLockListeners.put(workspace, listener);            }        }    }    /**     * Invoked when a record starts.     *     * @param workspace workspace, may be <code>null</code>     */    private void start(String workspace) {        this.workspace = workspace;        changeLog = new ChangeLog();        events = new ArrayList();    }    /**     * Process an update operation.     *     * @param operation operation to process     */    private void process(ItemOperation operation) {        operation.apply(changeLog);    }    /**     * Process an event.     *     * @param event event     */    private void process(EventState event) {        events.add(event);    }    /**     * Process a lock operation.     *     * @param nodeId node id     * @param isDeep flag indicating whether lock is deep     * @param owner lock owner     */    private void process(NodeId nodeId, boolean isDeep, String owner) {        LockEventListener listener = (LockEventListener) wspLockListeners.get(workspace);        if (listener == null) {            try {                clusterContext.lockEventsReady(workspace);            } catch (RepositoryException e) {                String msg = "Unable to make lock listener for workspace " +                        workspace + " online: " + e.getMessage();                log.warn(msg);            }            listener = (LockEventListener) wspLockListeners.get(workspace);            if (listener ==  null) {                String msg = "Lock channel unavailable for workspace: " + workspace;                log.error(msg);                return;            }        }        try {            listener.externalLock(nodeId, isDeep, owner);        } catch (RepositoryException e) {            String msg = "Unable to deliver lock event: " + e.getMessage();            log.error(msg);        }    }    /**     * Process an unlock operation.     *     * @param nodeId node id     */    private void process(NodeId nodeId) {        LockEventListener listener = (LockEventListener) wspLockListeners.get(workspace);        if (listener == null) {            try {                clusterContext.lockEventsReady(workspace);            } catch (RepositoryException e) {                String msg = "Unable to make lock listener for workspace " +                        workspace + " online: " + e.getMessage();                log.warn(msg);            }            listener = (LockEventListener) wspLockListeners.get(workspace);            if (listener ==  null) {                String msg = "Lock channel unavailable for workspace: " + workspace;                log.error(msg);                return;            }        }        try {            listener.externalUnlock(nodeId);        } catch (RepositoryException e) {            String msg = "Unable to deliver lock event: " + e.getMessage();            log.error(msg);        }    }    /**     * Process a namespace operation.     *     * @param oldPrefix old prefix. if <code>null</code> this is a fresh mapping     * @param newPrefix new prefix. if <code>null</code> this is an unmap operation     * @param uri uri to map prefix to     */    private void process(String oldPrefix, String newPrefix, String uri) {        if (namespaceListener == null) {            String msg = "Namespace listener unavailable.";            log.error(msg);            return;        }        try {            namespaceListener.externalRemap(oldPrefix, newPrefix, uri);        } catch (RepositoryException e) {            String msg = "Unable to deliver namespace operation: " + e.getMessage();            log.error(msg);        }    }    /**     * Process one or more node type registrations.     *     * @param c collection of node type definitions, if this is a register     *          operation; collection of <code>QName</code>s if this is     *          an unregister operation

⌨️ 快捷键说明

复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?